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
|