streams.client.spec.ts 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198
  1. /** Stream scripts, live stream control, cancellation, and the built-in `$events` opening. */
  2. import { describe, expect, it } from 'vitest'
  3. import { RemoteMock, frames, openStream } from '../src/index.ts'
  4. const idle = (): AbortSignal => new AbortController().signal
  5. async function drain(source: AsyncIterable<unknown>): Promise<unknown[]> {
  6. const items: unknown[] = []
  7. for await (const item of source) items.push(item)
  8. return items
  9. }
  10. async function take(source: AsyncIterable<unknown>, count: number): Promise<unknown[]> {
  11. const items: unknown[] = []
  12. for await (const item of source) {
  13. items.push(item)
  14. if (items.length === count) break
  15. }
  16. return items
  17. }
  18. describe('RemoteMock streams', () => {
  19. it('yields frames() then ends, logging the open, and throws on an unmatched endpoint', async () => {
  20. const mock = RemoteMock.create().stream('s/f', frames([{ n: 1 }, { n: 2 }]))
  21. await expect(drain(mock.open('s/f', [{ id: 'a' }], idle()))).resolves.toEqual([{ n: 1 }, { n: 2 }])
  22. expect(mock.log.streams('s/f')).toEqual([{ endpoint: 's/f', args: [{ id: 'a' }], state: 'ended', pushed: 2, seq: 1 }])
  23. expect(mock.log.streams('other')).toEqual([])
  24. expect(() => mock.open('s/g', [], idle())).toThrow('remote-mock: no rule for s/g; registered: $events, s/f')
  25. expect(mock.log.unmatched()).toEqual([{ endpoint: 's/g', mode: 'stream' }])
  26. })
  27. it('keeps openStream() open for pushes, filters by open args, and reports delivery counts', async () => {
  28. const mock = RemoteMock.create().stream('s/f', openStream(['hello']))
  29. const sessionOf = ([request]: readonly unknown[]): string => (request as { sessionId: string }).sessionId
  30. const a = mock.open('s/f', [{ sessionId: 'a' }], idle())
  31. const b = mock.open('s/f', [{ sessionId: 'b' }], idle())
  32. expect(mock.streams.push('s/f', 'only-b', open => sessionOf(open) === 'b')).toBe(1)
  33. expect(mock.streams.push('s/f', 'both')).toBe(2)
  34. expect(mock.streams.fail('s/f', new Error('gone'), open => sessionOf(open) === 'a')).toBe(1)
  35. expect(mock.streams.end('s/f')).toBe(1)
  36. await expect(drain(b)).resolves.toEqual(['hello', 'only-b', 'both'])
  37. await expect(drain(a)).rejects.toThrow('gone')
  38. expect(mock.log.streams().map(entry => [entry.state, entry.pushed])).toEqual([['failed', 2], ['ended', 3]])
  39. expect(mock.streams.push('s/f', 'late')).toBe(0)
  40. })
  41. it('resolves a pending read on push, end, or fail; a second concurrent read is a bug; settling twice is a no-op', async () => {
  42. const mock = RemoteMock.create().stream('s/f', openStream())
  43. const first = mock.open('s/f', [], idle())[Symbol.asyncIterator]()
  44. const pending = first.next()
  45. await expect(first.next()).rejects.toThrow('remote-mock: s/f stream has one consumer')
  46. mock.streams.push('s/f', 'x')
  47. await expect(pending).resolves.toEqual({ value: 'x', done: false })
  48. const ending = first.next()
  49. mock.streams.end('s/f')
  50. await expect(ending).resolves.toEqual({ value: undefined, done: true })
  51. const second = mock.open('s/f', [], idle())[Symbol.asyncIterator]()
  52. const failing = second.next()
  53. mock.streams.fail('s/f', new Error('boom'))
  54. await expect(failing).rejects.toThrow('boom')
  55. await expect(second.next()).rejects.toThrow('boom')
  56. expect(mock.streams.end('s/f')).toBe(0)
  57. const settled = new AbortController()
  58. const ended = mock.open('s/f', [], settled.signal)
  59. mock.streams.end('s/f')
  60. settled.abort()
  61. await expect(drain(ended)).resolves.toEqual([])
  62. expect(mock.log.streams().map(entry => entry.state)).toEqual(['ended', 'failed', 'ended'])
  63. })
  64. it('lets a consumer that returns after the producer ended drop the rest, so drained() settles', async () => {
  65. const mock = RemoteMock.create().stream('s/f', frames(['a', 'b']))
  66. const reader = mock.open('s/f', [], idle())[Symbol.asyncIterator]()
  67. await expect(reader.next()).resolves.toEqual({ value: 'a', done: false })
  68. const pending = mock.streams.drained('s/f')
  69. await expect(reader.return!()).resolves.toEqual({ value: undefined, done: true })
  70. await expect(pending).resolves.toBeUndefined()
  71. expect(mock.log.streams('s/f')).toEqual([{ endpoint: 's/f', args: [], state: 'ended', pushed: 2, seq: 1 }])
  72. expect(mock.streams.push('s/f', 'late')).toBe(0)
  73. })
  74. it('treats consumer abort or return() as cancellation and drops later pushes', async () => {
  75. const mock = RemoteMock.create().stream('s/f', openStream(['queued']))
  76. const controller = new AbortController()
  77. const cancelled = mock.open('s/f', [], controller.signal)[Symbol.asyncIterator]()
  78. await expect(cancelled.next()).resolves.toEqual({ value: 'queued', done: false })
  79. const waiting = cancelled.next()
  80. controller.abort()
  81. await expect(waiting).resolves.toEqual({ value: undefined, done: true })
  82. expect(mock.streams.push('s/f', 'after')).toBe(0)
  83. const aborted = new AbortController()
  84. aborted.abort()
  85. await expect(drain(mock.open('s/f', [], aborted.signal))).resolves.toEqual([])
  86. await expect(take(mock.open('s/f', [], idle()), 1)).resolves.toEqual(['queued'])
  87. expect(mock.log.streams().map(entry => entry.state)).toEqual(['cancelled', 'cancelled', 'cancelled'])
  88. })
  89. it('aborts the stream handle signal when the consumer returns', async () => {
  90. let signal: AbortSignal | undefined
  91. const mock = RemoteMock.create().stream('s/f', (_args, stream) => {
  92. signal = stream.signal
  93. stream.push('first')
  94. })
  95. const reader = mock.open('s/f', [], idle())[Symbol.asyncIterator]()
  96. await expect(reader.next()).resolves.toEqual({ value: 'first', done: false })
  97. await expect(reader.return!()).resolves.toEqual({ value: undefined, done: true })
  98. expect({
  99. state: mock.log.streams('s/f')[0]?.state,
  100. signalAborted: signal?.aborted,
  101. }).toEqual({ state: 'cancelled', signalAborted: true })
  102. })
  103. it('runs script functions with the open args and fails the stream when they throw or reject', async () => {
  104. const mock = RemoteMock.create()
  105. .stream('s/echo', (args, stream) => {
  106. stream.push(args)
  107. stream.end()
  108. stream.end()
  109. stream.fail(new Error('too late'))
  110. })
  111. .stream('s/async', async (_args, stream) => {
  112. await Promise.resolve()
  113. stream.push('later')
  114. stream.end()
  115. })
  116. .stream('s/throws', () => { throw new Error('sync boom') })
  117. .stream('s/rejects', () => Promise.reject(new Error('async boom')))
  118. .stream('s/odd', () => { throw 'string reason' })
  119. await expect(drain(mock.open('s/echo', [{ a: 1 }], idle()))).resolves.toEqual([[{ a: 1 }]])
  120. await expect(drain(mock.open('s/async', [], idle()))).resolves.toEqual(['later'])
  121. await expect(drain(mock.open('s/throws', [], idle()))).rejects.toThrow('sync boom')
  122. await expect(drain(mock.open('s/rejects', [], idle()))).rejects.toThrow('async boom')
  123. await expect(drain(mock.open('s/odd', [], idle()))).rejects.toThrow('string reason')
  124. })
  125. it('drained() settles once the consumer has pulled every push and waits again, or the stream closed', async () => {
  126. const mock = RemoteMock.create().stream('s/f', openStream(['first']))
  127. await expect(mock.streams.drained('s/f')).resolves.toBeUndefined() // nothing open: nothing to drain
  128. const reader = mock.open('s/f', [], idle())[Symbol.asyncIterator]()
  129. const unread = mock.streams.drained('s/f')
  130. let settled = false
  131. void unread.then(() => { settled = true })
  132. await Promise.resolve()
  133. expect(settled).toBe(false) // 'first' is queued and nobody has pulled it
  134. await expect(reader.next()).resolves.toEqual({ value: 'first', done: false })
  135. await Promise.resolve()
  136. expect(settled).toBe(false) // pulled, but the consumer is not waiting for more yet
  137. const waiting = reader.next()
  138. await expect(unread).resolves.toBeUndefined()
  139. mock.streams.push('s/f', 'second')
  140. await expect(waiting).resolves.toEqual({ value: 'second', done: false })
  141. mock.streams.push('s/f', 'third')
  142. mock.streams.end('s/f')
  143. const ended = mock.streams.drained('s/f')
  144. let endedSettled = false
  145. void ended.then(() => { endedSettled = true })
  146. await Promise.resolve()
  147. expect(endedSettled).toBe(false) // ended, but 'third' is still queued
  148. await expect(reader.next()).resolves.toEqual({ value: 'third', done: false })
  149. await expect(ended).resolves.toBeUndefined() // queue empty: a settled stream counts as drained
  150. await expect(reader.next()).resolves.toEqual({ value: undefined, done: true })
  151. await expect(reader.return!()).resolves.toEqual({ value: undefined, done: true }) // returning a settled stream changes nothing
  152. const second = mock.open('s/f', [], idle())[Symbol.asyncIterator]()
  153. await expect(second.next()).resolves.toEqual({ value: 'first', done: false })
  154. const parked = second.next()
  155. await expect(mock.streams.drained('s/f')).resolves.toBeUndefined() // its consumer is waiting
  156. const controller = new AbortController()
  157. const third = mock.open('s/f', [], controller.signal)
  158. const cancelling = mock.streams.drained('s/f', () => true) // third holds 'first' that nobody has read
  159. controller.abort()
  160. await expect(cancelling).resolves.toBeUndefined() // cancellation discards the queue: closed counts as drained
  161. await expect(drain(third)).resolves.toEqual([])
  162. mock.streams.end('s/f')
  163. await expect(parked).resolves.toEqual({ value: undefined, done: true })
  164. })
  165. it('declares a stream without a script: modeOf answers stream and an open is a stream miss', () => {
  166. const mock = RemoteMock.create().load({ streams: ['s/declared'] })
  167. expect(mock.modeOf('s/declared')).toBe('stream')
  168. expect(mock.endpoints()).toEqual(['$events', 's/declared'])
  169. expect(() => mock.open('s/declared', [], idle())).toThrow('remote-mock: no rule for s/declared; registered: $events, s/declared')
  170. expect(mock.log.unmatched()).toEqual([{ endpoint: 's/declared', mode: 'stream' }])
  171. mock.stream('s/declared', frames(['now']))
  172. expect(mock.modeOf('s/declared')).toBe('stream')
  173. })
  174. it('waits for opens with opened(), and answers $events with one ready frame per generation', async () => {
  175. const mock = RemoteMock.create({ host: { home: '/home/me' } })
  176. const second = mock.streams.opened('$events', 2)
  177. const first = mock.open('$events', [{}], idle())
  178. await expect(take(first, 1)).resolves.toEqual([{ type: 'ready', clientId: 'mock-client-1', host: { home: '/home/me' } }])
  179. const again = mock.open('$events', [{}], idle())
  180. await expect(second).resolves.toBeUndefined()
  181. await expect(mock.streams.opened('$events', 1)).resolves.toBeUndefined()
  182. await expect(take(again, 1)).resolves.toEqual([{ type: 'ready', clientId: 'mock-client-2', host: { home: '/home/me' } }])
  183. expect(RemoteMock.create().modeOf('$events')).toBe('stream')
  184. })
  185. })