stream-server.host.spec.ts 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312
  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('requires two missed heartbeats before terminating an unresponsive socket', 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 once(client, 'ping')
  53. await once(client, 'ping')
  54. expect(terminated).not.toHaveBeenCalled()
  55. await vi.waitFor(() => { expect(terminated).toHaveBeenCalledOnce() })
  56. await closed
  57. })
  58. it('keeps the socket when a delayed Pong arrives before the final check', async () => {
  59. const entry = await startMux(async (_endpoint, _payload, signal) => waitForAbort(signal), 20)
  60. const client = await connect(entry.url, false)
  61. const serverSocket = acceptedSocket(entry.mux)
  62. const terminated = vi.spyOn(serverSocket, 'terminate')
  63. let finalCheck: (() => void) | undefined
  64. const immediate = vi.spyOn(globalThis, 'setImmediate').mockImplementation((callback) => {
  65. finalCheck = callback
  66. return 0 as unknown as NodeJS.Immediate
  67. })
  68. try {
  69. await once(client, 'ping')
  70. await once(client, 'ping')
  71. await vi.waitFor(() => { expect(finalCheck).toBeDefined() })
  72. serverSocket.emit('pong', Buffer.alloc(0))
  73. finalCheck?.()
  74. expect(terminated).not.toHaveBeenCalled()
  75. } finally {
  76. immediate.mockRestore()
  77. const closed = once(client, 'close')
  78. client.close()
  79. await closed
  80. }
  81. })
  82. it('rejects binary, malformed, and duplicate logical-stream messages', async () => {
  83. const entry = await startMux(async (_endpoint, _payload, signal) => waitForAbort(signal))
  84. const binary = await connect(entry.url)
  85. const binaryClosed = once(binary, 'close')
  86. binary.send(Buffer.from('{}'))
  87. const binaryEvent = await binaryClosed
  88. expect(binaryEvent[0]).toBe(1003)
  89. const malformed = await connect(entry.url)
  90. const malformedClosed = once(malformed, 'close')
  91. malformed.send('not json')
  92. const malformedEvent = await malformedClosed
  93. expect(malformedEvent[0]).toBe(1008)
  94. expect(String(malformedEvent[1])).toBe('invalid Remote stream request')
  95. const duplicate = await connect(entry.url)
  96. const longId = 'same'.repeat(100)
  97. duplicate.send(openFrame(longId))
  98. duplicate.send(openFrame(longId))
  99. const duplicateEvent = await once(duplicate, 'close')
  100. expect(duplicateEvent[0]).toBe(1008)
  101. expect(String(duplicateEvent[1])).toBe('invalid Remote stream request')
  102. const noInput = await connect(entry.url)
  103. noInput.send(openFrame('no-input'))
  104. noInput.send(JSON.stringify({ type: 'input', streamId: 'no-input', value: 'unexpected' }))
  105. const noInputEvent = await once(noInput, 'close')
  106. expect(noInputEvent[0]).toBe(1008)
  107. expect(String(noInputEvent[1])).toBe('invalid Remote stream request')
  108. })
  109. it('accepts all ws text representations and terminates a carrier error', async () => {
  110. const entry = await startMux(async (_endpoint, _payload, signal) => waitForAbort(signal))
  111. const client = await connect(entry.url)
  112. const serverSocket = acceptedSocket(entry.mux)
  113. const cancel = JSON.stringify({ type: 'cancel', streamId: 'absent' })
  114. serverSocket.emit('message', [Buffer.from(cancel)], false)
  115. serverSocket.emit('message', Uint8Array.from(Buffer.from(cancel)).buffer, false)
  116. const closed = once(client, 'close')
  117. serverSocket.emit('error', new Error('fixture carrier failure'))
  118. await closed
  119. })
  120. it('does not send an end frame after clean source cancellation', async () => {
  121. let opened!: () => void
  122. const didOpen = new Promise<void>((resolve) => { opened = resolve })
  123. let returned!: () => void
  124. const didReturn = new Promise<void>((resolve) => { returned = resolve })
  125. const entry = await startMux(async (_endpoint, _payload, signal) => {
  126. opened()
  127. return cleanlyCancelled(signal, returned)
  128. })
  129. const client = await connect(entry.url)
  130. const frames: unknown[] = []
  131. client.on('message', (data) => {
  132. if (!Buffer.isBuffer(data)) throw new TypeError('fixture expected a Buffer frame')
  133. frames.push(JSON.parse(data.toString('utf8')) as unknown)
  134. })
  135. client.send(openFrame('cancelled'))
  136. await didOpen
  137. client.send(JSON.stringify({ type: 'cancel', streamId: 'cancelled' }))
  138. await didReturn
  139. await new Promise<void>((resolve) => { setImmediate(resolve) })
  140. expect(frames).toEqual([])
  141. client.close()
  142. await once(client, 'close')
  143. })
  144. it('closes the carrier when ws reports an item write failure', async () => {
  145. let release!: () => void
  146. const released = new Promise<void>((resolve) => { release = resolve })
  147. let opened!: () => void
  148. const didOpen = new Promise<void>((resolve) => { opened = resolve })
  149. const entry = await startMux(async () => delayedItem(released, opened))
  150. const client = await connect(entry.url)
  151. client.send(openFrame('write-failure'))
  152. await didOpen
  153. const serverSocket = acceptedSocket(entry.mux)
  154. const mutable = serverSocket as unknown as {
  155. send(data: unknown, callback: (error?: Error) => void): void
  156. }
  157. mutable.send = (_data, callback): void => {
  158. callback(new Error('fixture ws write failure'))
  159. }
  160. const closed = once(client, 'close')
  161. release()
  162. const closeEvent = await closed
  163. expect(closeEvent[0]).toBe(1011)
  164. expect(String(closeEvent[1])).toBe('Remote stream failure could not be delivered')
  165. })
  166. it('contains an item produced after its socket closes', async () => {
  167. let release!: () => void
  168. const released = new Promise<void>((resolve) => { release = resolve })
  169. let opened!: () => void
  170. const didOpen = new Promise<void>((resolve) => { opened = resolve })
  171. let returned!: () => void
  172. const didReturn = new Promise<void>((resolve) => { returned = resolve })
  173. const entry = await startMux(async () => delayedItem(released, opened, returned))
  174. const client = await connect(entry.url)
  175. client.send(openFrame('late-item'))
  176. await didOpen
  177. const serverSocket = acceptedSocket(entry.mux)
  178. client.close()
  179. await once(client, 'close')
  180. await vi.waitFor(() => { expect(serverSocket.readyState).toBe(WebSocket.CLOSED) })
  181. release()
  182. await didReturn
  183. })
  184. it('terminates active sockets on close and reports a repeated close', async () => {
  185. let opened!: () => void
  186. const didOpen = new Promise<void>((resolve) => { opened = resolve })
  187. let returned!: () => void
  188. const didReturn = new Promise<void>((resolve) => { returned = resolve })
  189. const entry = await startMux(async (_endpoint, _payload, signal) => {
  190. opened()
  191. return cleanlyCancelled(signal, returned)
  192. })
  193. const client = await connect(entry.url)
  194. client.send(openFrame('active'))
  195. await didOpen
  196. const closed = once(client, 'close')
  197. await entry.mux.close()
  198. running.delete(entry)
  199. await closed
  200. await didReturn
  201. await expect(entry.mux.close()).rejects.toThrow()
  202. await closeHttp(entry.http)
  203. })
  204. })
  205. const mapFailure: RemoteStreamFailureMapper = error => ({
  206. code: 'internal',
  207. message: error instanceof Error ? error.message : String(error),
  208. details: {},
  209. })
  210. async function startMux(open: RemoteStreamOpener, heartbeatIntervalMs = 2_000): Promise<RunningMux> {
  211. const mux = new RemoteStreamMuxServer(open, mapFailure, heartbeatIntervalMs)
  212. const http = createServer()
  213. http.on('upgrade', (request, socket, head) => { mux.handleUpgrade(request, socket, head) })
  214. await new Promise<void>((resolve, reject) => {
  215. http.once('error', reject)
  216. http.listen(0, '127.0.0.1', () => {
  217. http.off('error', reject)
  218. resolve()
  219. })
  220. })
  221. const address = http.address()
  222. if (address === null || typeof address === 'string') throw new Error('fixture HTTP server has no TCP port')
  223. const entry = { http, mux, url: `ws://127.0.0.1:${String(address.port)}` }
  224. running.add(entry)
  225. return entry
  226. }
  227. async function connect(url: string, autoPong = true): Promise<WebSocket> {
  228. const socket = new WebSocket(url, { autoPong })
  229. await once(socket, 'open')
  230. return socket
  231. }
  232. function acceptedSocket(mux: RemoteStreamMuxServer): WebSocket {
  233. const exposed = mux as unknown as { server: { clients: Set<WebSocket> } }
  234. const socket = [...exposed.server.clients][0]
  235. if (socket === undefined) throw new Error('fixture mux has no accepted socket')
  236. return socket
  237. }
  238. function openFrame(streamId: string): string {
  239. return JSON.stringify({ type: 'open', streamId, endpoint: 'fixture/follow', payload: {} })
  240. }
  241. async function *waitForAbort(signal: AbortSignal): AsyncIterable<never> {
  242. await new Promise<void>((resolve) => {
  243. if (signal.aborted) resolve()
  244. else signal.addEventListener('abort', () => { resolve() }, { once: true })
  245. })
  246. }
  247. async function *cleanlyCancelled(signal: AbortSignal, returned: () => void): AsyncIterable<never> {
  248. try {
  249. await new Promise<void>((resolve) => {
  250. if (signal.aborted) resolve()
  251. else signal.addEventListener('abort', () => { resolve() }, { once: true })
  252. })
  253. } finally {
  254. returned()
  255. }
  256. }
  257. async function *delayedItem(
  258. released: Promise<void>,
  259. opened: () => void,
  260. returned: () => void = () => {},
  261. ): AsyncIterable<string> {
  262. try {
  263. opened()
  264. await released
  265. yield 'item'
  266. } finally {
  267. returned()
  268. }
  269. }
  270. async function closeHttp(server: Server): Promise<void> {
  271. if (!server.listening) return
  272. await new Promise<void>((resolve, reject) => {
  273. server.close((error) => {
  274. if (error === undefined) resolve()
  275. else reject(error)
  276. })
  277. })
  278. }