stream-server.host.spec.ts 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283
  1. import { once } from 'node:events'
  2. import { createServer, type Server } from 'node:http'
  3. import { afterEach, describe, expect, it, vi } from 'vitest'
  4. import WebSocket from 'ws'
  5. import {
  6. RemoteStreamMuxServer,
  7. type RemoteStreamFailureMapper,
  8. type RemoteStreamOpener,
  9. } from '../src/stream-server.ts'
  10. interface RunningMux {
  11. readonly http: Server
  12. readonly mux: RemoteStreamMuxServer
  13. readonly url: string
  14. }
  15. const running = new Set<RunningMux>()
  16. afterEach(async () => {
  17. await Promise.all([...running].map(async (entry) => {
  18. running.delete(entry)
  19. await entry.mux.close().catch(() => undefined)
  20. await closeHttp(entry.http)
  21. }))
  22. })
  23. describe('Remote stream mux server carrier lifecycle', () => {
  24. it('sends WebSocket Ping control frames without application messages', async () => {
  25. const entry = await startMux(async (_endpoint, _payload, signal) => waitForAbort(signal), 20)
  26. const client = await connect(entry.url)
  27. const serverSocket = acceptedSocket(entry.mux)
  28. const messages = vi.fn()
  29. client.on('message', messages)
  30. const ping = once(client, 'ping')
  31. const pong = once(serverSocket, 'pong')
  32. expect((await ping)[0]).toEqual(Buffer.alloc(0))
  33. expect((await pong)[0]).toEqual(Buffer.alloc(0))
  34. expect(messages).not.toHaveBeenCalled()
  35. const closingPing = vi.spyOn(serverSocket, 'ping')
  36. client.pause()
  37. serverSocket.close()
  38. expect(serverSocket.readyState).toBe(WebSocket.CLOSING)
  39. await new Promise<void>((resolve) => { setTimeout(resolve, 25) })
  40. expect(closingPing).not.toHaveBeenCalled()
  41. const closed = once(client, 'close')
  42. client.resume()
  43. await closed
  44. })
  45. it('terminates a socket that does not answer the previous heartbeat', async () => {
  46. const entry = await startMux(async (_endpoint, _payload, signal) => waitForAbort(signal), 20)
  47. const client = await connect(entry.url)
  48. const serverSocket = acceptedSocket(entry.mux)
  49. serverSocket.removeAllListeners('pong')
  50. const terminated = vi.spyOn(serverSocket, 'terminate')
  51. const closed = once(client, 'close')
  52. await vi.waitFor(() => { expect(terminated).toHaveBeenCalledOnce() })
  53. await closed
  54. })
  55. it('rejects binary, malformed, and duplicate logical-stream messages', async () => {
  56. const entry = await startMux(async (_endpoint, _payload, signal) => waitForAbort(signal))
  57. const binary = await connect(entry.url)
  58. const binaryClosed = once(binary, 'close')
  59. binary.send(Buffer.from('{}'))
  60. const binaryEvent = await binaryClosed
  61. expect(binaryEvent[0]).toBe(1003)
  62. const malformed = await connect(entry.url)
  63. const malformedClosed = once(malformed, 'close')
  64. malformed.send('not json')
  65. const malformedEvent = await malformedClosed
  66. expect(malformedEvent[0]).toBe(1008)
  67. expect(String(malformedEvent[1])).toBe('invalid Remote stream request')
  68. const duplicate = await connect(entry.url)
  69. const longId = 'same'.repeat(100)
  70. duplicate.send(openFrame(longId))
  71. duplicate.send(openFrame(longId))
  72. const duplicateEvent = await once(duplicate, 'close')
  73. expect(duplicateEvent[0]).toBe(1008)
  74. expect(String(duplicateEvent[1])).toBe('invalid Remote stream request')
  75. const noInput = await connect(entry.url)
  76. noInput.send(openFrame('no-input'))
  77. noInput.send(JSON.stringify({ type: 'input', streamId: 'no-input', value: 'unexpected' }))
  78. const noInputEvent = await once(noInput, 'close')
  79. expect(noInputEvent[0]).toBe(1008)
  80. expect(String(noInputEvent[1])).toBe('invalid Remote stream request')
  81. })
  82. it('accepts all ws text representations and terminates a carrier error', async () => {
  83. const entry = await startMux(async (_endpoint, _payload, signal) => waitForAbort(signal))
  84. const client = await connect(entry.url)
  85. const serverSocket = acceptedSocket(entry.mux)
  86. const cancel = JSON.stringify({ type: 'cancel', streamId: 'absent' })
  87. serverSocket.emit('message', [Buffer.from(cancel)], false)
  88. serverSocket.emit('message', Uint8Array.from(Buffer.from(cancel)).buffer, false)
  89. const closed = once(client, 'close')
  90. serverSocket.emit('error', new Error('fixture carrier failure'))
  91. await closed
  92. })
  93. it('does not send an end frame after clean source cancellation', async () => {
  94. let opened!: () => void
  95. const didOpen = new Promise<void>((resolve) => { opened = resolve })
  96. let returned!: () => void
  97. const didReturn = new Promise<void>((resolve) => { returned = resolve })
  98. const entry = await startMux(async (_endpoint, _payload, signal) => {
  99. opened()
  100. return cleanlyCancelled(signal, returned)
  101. })
  102. const client = await connect(entry.url)
  103. const frames: unknown[] = []
  104. client.on('message', (data) => {
  105. if (!Buffer.isBuffer(data)) throw new TypeError('fixture expected a Buffer frame')
  106. frames.push(JSON.parse(data.toString('utf8')) as unknown)
  107. })
  108. client.send(openFrame('cancelled'))
  109. await didOpen
  110. client.send(JSON.stringify({ type: 'cancel', streamId: 'cancelled' }))
  111. await didReturn
  112. await new Promise<void>((resolve) => { setImmediate(resolve) })
  113. expect(frames).toEqual([])
  114. client.close()
  115. await once(client, 'close')
  116. })
  117. it('closes the carrier when ws reports an item write failure', async () => {
  118. let release!: () => void
  119. const released = new Promise<void>((resolve) => { release = resolve })
  120. let opened!: () => void
  121. const didOpen = new Promise<void>((resolve) => { opened = resolve })
  122. const entry = await startMux(async () => delayedItem(released, opened))
  123. const client = await connect(entry.url)
  124. client.send(openFrame('write-failure'))
  125. await didOpen
  126. const serverSocket = acceptedSocket(entry.mux)
  127. const mutable = serverSocket as unknown as {
  128. send(data: unknown, callback: (error?: Error) => void): void
  129. }
  130. mutable.send = (_data, callback): void => {
  131. callback(new Error('fixture ws write failure'))
  132. }
  133. const closed = once(client, 'close')
  134. release()
  135. const closeEvent = await closed
  136. expect(closeEvent[0]).toBe(1011)
  137. expect(String(closeEvent[1])).toBe('Remote stream failure could not be delivered')
  138. })
  139. it('contains an item produced after its socket closes', async () => {
  140. let release!: () => void
  141. const released = new Promise<void>((resolve) => { release = resolve })
  142. let opened!: () => void
  143. const didOpen = new Promise<void>((resolve) => { opened = resolve })
  144. let returned!: () => void
  145. const didReturn = new Promise<void>((resolve) => { returned = resolve })
  146. const entry = await startMux(async () => delayedItem(released, opened, returned))
  147. const client = await connect(entry.url)
  148. client.send(openFrame('late-item'))
  149. await didOpen
  150. const serverSocket = acceptedSocket(entry.mux)
  151. client.close()
  152. await once(client, 'close')
  153. await vi.waitFor(() => { expect(serverSocket.readyState).toBe(WebSocket.CLOSED) })
  154. release()
  155. await didReturn
  156. })
  157. it('terminates active sockets on close and reports a repeated close', async () => {
  158. let opened!: () => void
  159. const didOpen = new Promise<void>((resolve) => { opened = resolve })
  160. let returned!: () => void
  161. const didReturn = new Promise<void>((resolve) => { returned = resolve })
  162. const entry = await startMux(async (_endpoint, _payload, signal) => {
  163. opened()
  164. return cleanlyCancelled(signal, returned)
  165. })
  166. const client = await connect(entry.url)
  167. client.send(openFrame('active'))
  168. await didOpen
  169. const closed = once(client, 'close')
  170. await entry.mux.close()
  171. running.delete(entry)
  172. await closed
  173. await didReturn
  174. await expect(entry.mux.close()).rejects.toThrow()
  175. await closeHttp(entry.http)
  176. })
  177. })
  178. const mapFailure: RemoteStreamFailureMapper = error => ({
  179. code: 'internal',
  180. message: error instanceof Error ? error.message : String(error),
  181. details: {},
  182. })
  183. async function startMux(open: RemoteStreamOpener, heartbeatIntervalMs = 2_000): Promise<RunningMux> {
  184. const mux = new RemoteStreamMuxServer(open, mapFailure, heartbeatIntervalMs)
  185. const http = createServer()
  186. http.on('upgrade', (request, socket, head) => { mux.handleUpgrade(request, socket, head) })
  187. await new Promise<void>((resolve, reject) => {
  188. http.once('error', reject)
  189. http.listen(0, '127.0.0.1', () => {
  190. http.off('error', reject)
  191. resolve()
  192. })
  193. })
  194. const address = http.address()
  195. if (address === null || typeof address === 'string') throw new Error('fixture HTTP server has no TCP port')
  196. const entry = { http, mux, url: `ws://127.0.0.1:${String(address.port)}` }
  197. running.add(entry)
  198. return entry
  199. }
  200. async function connect(url: string): Promise<WebSocket> {
  201. const socket = new WebSocket(url)
  202. await once(socket, 'open')
  203. return socket
  204. }
  205. function acceptedSocket(mux: RemoteStreamMuxServer): WebSocket {
  206. const exposed = mux as unknown as { server: { clients: Set<WebSocket> } }
  207. const socket = [...exposed.server.clients][0]
  208. if (socket === undefined) throw new Error('fixture mux has no accepted socket')
  209. return socket
  210. }
  211. function openFrame(streamId: string): string {
  212. return JSON.stringify({ type: 'open', streamId, endpoint: 'fixture/follow', payload: {} })
  213. }
  214. async function *waitForAbort(signal: AbortSignal): AsyncIterable<never> {
  215. await new Promise<void>((resolve) => {
  216. if (signal.aborted) resolve()
  217. else signal.addEventListener('abort', () => { resolve() }, { once: true })
  218. })
  219. }
  220. async function *cleanlyCancelled(signal: AbortSignal, returned: () => void): AsyncIterable<never> {
  221. try {
  222. await new Promise<void>((resolve) => {
  223. if (signal.aborted) resolve()
  224. else signal.addEventListener('abort', () => { resolve() }, { once: true })
  225. })
  226. } finally {
  227. returned()
  228. }
  229. }
  230. async function *delayedItem(
  231. released: Promise<void>,
  232. opened: () => void,
  233. returned: () => void = () => {},
  234. ): AsyncIterable<string> {
  235. try {
  236. opened()
  237. await released
  238. yield 'item'
  239. } finally {
  240. returned()
  241. }
  242. }
  243. async function closeHttp(server: Server): Promise<void> {
  244. if (!server.listening) return
  245. await new Promise<void>((resolve, reject) => {
  246. server.close((error) => {
  247. if (error === undefined) resolve()
  248. else reject(error)
  249. })
  250. })
  251. }