stream-server.host.spec.ts 8.5 KB

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