All files / src/adapters/socket.io index.ts

0% Statements 0/34
0% Branches 0/23
0% Functions 0/10
0% Lines 0/34

Press n or j to go to the next uncovered block, b, p or k for the previous block.

1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106                                                                                                                                                                                                                   
import { Server } from 'socket.io'
import Adapter from '../../lib/adapter.js'
import GleeQuoreMessage from '../../lib/message.js'
 
class SocketIOAdapter extends Adapter {
  private server: Server
 
  name(): string {
    return 'Socket.IO adapter'
  }
 
  async connect(): Promise<this> {
    return this._connect()
  }
 
  async send(message: GleeQuoreMessage): Promise<void> {
    return this._send(message)
  }
 
  async _connect(): Promise<this> {
    const config = await this.resolveProtocolConfig('ws')
    const websocketOptions = config?.server
    const serverUrl: URL = new URL(this.serverUrlExpanded)
    const asyncapiServerPort: number = serverUrl.port
      ? Number(serverUrl.port)
      : 80
    const optionsPort: number = websocketOptions?.port
    const port: number = optionsPort || asyncapiServerPort
 
    const serverOptions: { [key: string]: any } = {
      path: serverUrl.pathname || '/',
      serveClient: false,
      transports: ['websocket'],
    }
 
    if (websocketOptions.httpServer) {
      const server = websocketOptions.httpServer
      Iif (!optionsPort && String(server.address().port) !== String(port)) {
        console.error(
          `Your custom HTTP server is listening on port ${
            server.address().port
          } but your AsyncAPI file says it must listen on ${port}. Please fix the inconsistency.`
        )
        process.exit(1)
      }
      this.server = new Server(server, serverOptions)
    } else {
      this.server = new Server({
        ...serverOptions,
        ...{
          cors: {
            origin: true,
          },
        },
      })
    }
 
    this.server.on('connect', (socket) => {
      this.emit('server:ready', {
        name: this.name(),
        adapter: this,
        connection: socket,
        channels: this.channelNames,
      })
 
      socket.onAny((eventName, payload) => {
        const msg = this._createMessage(eventName, payload)
        this.emit('message', msg, socket)
      })
    })
 
    Iif (!websocketOptions.httpServer) {
      this.server.listen(port)
    }
    return this
  }
 
  async _send(message: GleeQuoreMessage): Promise<void> {
    if (message.broadcast) {
      this.app.syncCluster(message)
 
      this.connections
        .filter(({ channels }) => channels.includes(message.channel))
        .forEach((connection) => {
          connection.getRaw().emit(message.channel, message.payload)
        })
    } else {
      Iif (!message.connection) {
        throw new Error(
          'There is no Socket.IO connection to send the message yet.'
        )
      }
      message.connection.getRaw().emit(message.channel, message.payload)
    }
  }
 
  _createMessage(eventName: string, payload: any) {
    return new GleeQuoreMessage({
      payload,
      channel: eventName,
    })
  }
}
 
export default SocketIOAdapter