transport.client.spec.ts 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295
  1. import { describe, expect, it, vi } from 'vitest'
  2. import {
  3. RemoteStream,
  4. RemoteStreamCarrierError,
  5. RemoteStreamError,
  6. type RemoteStreamOptions,
  7. } from '@deepseek-ai/dsh-api-gateway/client'
  8. import type { RemoteResult } from '@deepseek-ai/dsh-typert-protocol'
  9. import {
  10. createSessionControlStream,
  11. SessionEventStream,
  12. sessionStreamFailure,
  13. type SessionJournalChange,
  14. type SessionRemote,
  15. } from '../src/client/index.ts'
  16. import type {
  17. SessionAddress,
  18. SessionControlFrame,
  19. SessionEventEntry,
  20. SessionFollowFrame,
  21. SessionFollowRequest,
  22. SessionPage,
  23. SessionPageRequest,
  24. } from '../src/types.ts'
  25. type SessionTransportRemote = Pick<SessionRemote, 'control' | 'follow' | 'page'>
  26. const ADDRESS: SessionAddress = { kind: 'session', sessionId: 'session-1' as never }
  27. const AVAILABLE_CONNECTION = {
  28. hostDescription: {
  29. getSnapshot: () => ({
  30. version: 'fixture', cwd: '/fixture', attachedSessions: 0, home: '/home/fixture', canOpenPath: true,
  31. }),
  32. subscribe: () => () => {},
  33. },
  34. }
  35. function entry(seq: number): SessionEventEntry {
  36. return { event: { type: 'turn/start', seq, time: seq, data: { turn: seq } } }
  37. }
  38. function page(events: readonly SessionEventEntry[], hasMore = false): SessionPage {
  39. return { events, hasMore }
  40. }
  41. function sessionClient(remote: SessionTransportRemote) {
  42. return {
  43. session: remote as SessionRemote,
  44. $stream: <Item>(options: RemoteStreamOptions<Item>) => (
  45. new RemoteStream(AVAILABLE_CONNECTION, options)
  46. ),
  47. }
  48. }
  49. interface FollowGeneration {
  50. readonly frames: readonly SessionFollowFrame[]
  51. readonly terminal?: Error
  52. readonly hold?: boolean
  53. readonly waitAfterFrames?: Promise<void>
  54. }
  55. class ScriptedSessionRemote implements SessionTransportRemote {
  56. readonly followRequests: SessionFollowRequest[] = []
  57. readonly pageRequests: SessionPageRequest[] = []
  58. readonly signals: AbortSignal[] = []
  59. constructor(
  60. private readonly generations: FollowGeneration[],
  61. private readonly pages: RemoteResult<SessionPage>[],
  62. private readonly controlFrames: readonly SessionControlFrame[] = [],
  63. private readonly holdControl = true,
  64. ) {}
  65. async *follow(request: SessionFollowRequest, signal = new AbortController().signal): AsyncIterable<SessionFollowFrame> {
  66. const generation = this.generations.shift()
  67. if (generation === undefined) throw new Error('no scripted Session generation')
  68. this.followRequests.push(request)
  69. this.signals.push(signal)
  70. for (const frame of generation.frames) yield frame
  71. await generation.waitAfterFrames
  72. if (generation.terminal !== undefined) throw generation.terminal
  73. if (generation.hold === true && !signal.aborted) {
  74. await new Promise<void>((resolve) => {
  75. signal.addEventListener('abort', () => { resolve() }, { once: true })
  76. })
  77. }
  78. }
  79. page(request: SessionPageRequest): Promise<RemoteResult<SessionPage>> {
  80. this.pageRequests.push(request)
  81. const result = this.pages.shift()
  82. if (result === undefined) throw new Error('no scripted Session page')
  83. return Promise.resolve(result)
  84. }
  85. async *control(signal = new AbortController().signal): AsyncIterable<SessionControlFrame> {
  86. for (const frame of this.controlFrames) yield frame
  87. if (this.holdControl && !signal.aborted) {
  88. await new Promise<void>((resolve) => {
  89. signal.addEventListener('abort', () => { resolve() }, { once: true })
  90. })
  91. }
  92. }
  93. }
  94. describe('Session Client stream adapters', () => {
  95. it('binds an event journal to one address and publishes replace, append, and prepend changes', async () => {
  96. const remote = new ScriptedSessionRemote(
  97. [{
  98. frames: [
  99. { type: 'opened', cursor: 3 },
  100. { type: 'event', ...entry(3) },
  101. { type: 'event', ...entry(4) },
  102. ],
  103. hold: true,
  104. }],
  105. [
  106. { ok: true, value: page([entry(2), entry(3)], true) },
  107. { ok: true, value: page([entry(0), entry(1)], false) },
  108. ],
  109. )
  110. const changes: SessionJournalChange[] = []
  111. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  112. publish: (change) => { changes.push(change) },
  113. failed: vi.fn(),
  114. })
  115. await stream.open({ maxMessages: 50 })
  116. await vi.waitFor(() => { expect(changes).toHaveLength(2) })
  117. await stream.prepend({ beforeSeq: 2, maxMessages: 50 })
  118. expect(remote.followRequests).toEqual([{ address: ADDRESS }])
  119. expect(remote.pageRequests).toEqual([
  120. { address: ADDRESS, throughSeq: 3, maxMessages: 50 },
  121. { address: ADDRESS, throughSeq: 4, beforeSeq: 2, maxMessages: 50 },
  122. ])
  123. expect(changes).toMatchObject([
  124. { type: 'replace', entries: [entry(2), entry(3)], hasMore: true },
  125. { type: 'append', entry: entry(4) },
  126. { type: 'prepend', entries: [entry(0), entry(1)], hasMore: false },
  127. ])
  128. await stream.dispose()
  129. expect(remote.signals[0]?.aborted).toBe(true)
  130. })
  131. it('resumes after the applied cursor and repairs through the addressed tail page', async () => {
  132. const lost = new RemoteStreamCarrierError('lost')
  133. const remote = new ScriptedSessionRemote(
  134. [
  135. {
  136. frames: [{ type: 'opened', cursor: 1 }, { type: 'event', ...entry(2) }],
  137. terminal: lost,
  138. },
  139. { frames: [{ type: 'opened', cursor: 4 }], hold: true },
  140. ],
  141. [
  142. { ok: true, value: page([entry(0), entry(1)]) },
  143. { ok: true, value: page([entry(0), entry(1), entry(2), entry(3), entry(4)]) },
  144. ],
  145. )
  146. const changes: SessionJournalChange[] = []
  147. const carrierFailed = vi.fn()
  148. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  149. publish: (change) => { changes.push(change) },
  150. carrierFailed,
  151. failed: vi.fn(),
  152. })
  153. await stream.open({ maxMessages: 50 })
  154. await vi.waitFor(() => { expect(remote.followRequests).toHaveLength(2) })
  155. expect(remote.followRequests).toEqual([
  156. { address: ADDRESS },
  157. { address: ADDRESS, afterSeq: 2 },
  158. ])
  159. expect(remote.pageRequests).toEqual([
  160. { address: ADDRESS, throughSeq: 1, maxMessages: 50 },
  161. { address: ADDRESS, throughSeq: 4, maxMessages: 50 },
  162. ])
  163. expect(changes.map(change => change.type)).toEqual(['replace', 'append', 'replace'])
  164. expect(carrierFailed).toHaveBeenCalledWith(lost)
  165. await stream.dispose()
  166. })
  167. it('repairs a resumed event stream without an optional message limit', async () => {
  168. const finish = Promise.withResolvers<undefined>()
  169. const remote = new ScriptedSessionRemote(
  170. [
  171. {
  172. frames: [{ type: 'opened', cursor: 0 }],
  173. waitAfterFrames: finish.promise,
  174. terminal: new RemoteStreamCarrierError('lost'),
  175. },
  176. { frames: [{ type: 'opened', cursor: 1 }], hold: true },
  177. ],
  178. [
  179. { ok: true, value: page([entry(0)]) },
  180. { ok: true, value: page([entry(0), entry(1)]) },
  181. ],
  182. )
  183. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  184. publish: vi.fn(),
  185. failed: vi.fn(),
  186. })
  187. await stream.open({})
  188. finish.resolve(undefined)
  189. await vi.waitFor(() => { expect(remote.pageRequests).toHaveLength(2) })
  190. expect(remote.pageRequests).toEqual([
  191. { address: ADDRESS, throughSeq: 0 },
  192. { address: ADDRESS, throughSeq: 1 },
  193. ])
  194. await stream.dispose()
  195. })
  196. it('turns a page failure into a typed stream failure and closes follow', async () => {
  197. const failure = { code: 'session-not-found', message: 'missing', details: { sessionId: 'session-1' } } as const
  198. const remote = new ScriptedSessionRemote(
  199. [{ frames: [{ type: 'opened', cursor: -1 }], hold: true }],
  200. [{ ok: false, error: failure }],
  201. )
  202. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  203. publish: vi.fn(),
  204. failed: vi.fn(),
  205. })
  206. await expect(stream.open({})).rejects.toBeInstanceOf(RemoteStreamError)
  207. await expect(stream.open({})).rejects.toThrow('already opened')
  208. expect(sessionStreamFailure(new RemoteStreamError(failure.code, failure.message, failure.details)))
  209. .toEqual(failure)
  210. expect(sessionStreamFailure(new Error('local'))).toBeUndefined()
  211. expect(remote.signals[0]?.aborted).toBe(true)
  212. expect(remote.pageRequests).toEqual([{ address: ADDRESS, throughSeq: -1 }])
  213. })
  214. it('maps the Host-wide control baseline and deltas into one snapshot stream', async () => {
  215. const baseline: SessionControlFrame = {
  216. type: 'baseline',
  217. value: { queues: {}, jobs: {}, projections: {} },
  218. }
  219. const update: SessionControlFrame = {
  220. type: 'queue', sessionId: 'session-1' as never, items: [],
  221. }
  222. const remote = new ScriptedSessionRemote([], [], [baseline, update])
  223. const accept = vi.fn<(frame: SessionControlFrame) => void>()
  224. const stream = createSessionControlStream(sessionClient(remote), {
  225. accept,
  226. failed: vi.fn(),
  227. })
  228. stream.start()
  229. stream.start()
  230. await vi.waitFor(() => { expect(accept).toHaveBeenCalledTimes(2) })
  231. expect(accept.mock.calls.map(([frame]) => frame)).toEqual([baseline, update])
  232. await stream.dispose()
  233. await stream.dispose()
  234. })
  235. it('classifies control streams that end before and after their opening baseline', async () => {
  236. const beforeFailed = vi.fn()
  237. const before = createSessionControlStream(
  238. sessionClient(new ScriptedSessionRemote([], [], [], false)),
  239. { accept: vi.fn(), failed: beforeFailed },
  240. )
  241. before.start()
  242. await vi.waitFor(() => { expect(beforeFailed).toHaveBeenCalledOnce() })
  243. expect(beforeFailed.mock.calls[0]?.[0]).toMatchObject({
  244. message: 'session control stream ended before its opening snapshot',
  245. })
  246. await before.dispose()
  247. const baseline: SessionControlFrame = {
  248. type: 'baseline',
  249. value: { queues: {}, jobs: {}, projections: {} },
  250. }
  251. const carrierFailed = vi.fn()
  252. const failed = vi.fn()
  253. const afterRemote = new ScriptedSessionRemote([], [], [baseline], false)
  254. const after = createSessionControlStream(sessionClient(afterRemote), {
  255. accept: vi.fn(),
  256. carrierFailed: (error) => {
  257. carrierFailed(error)
  258. void after.dispose()
  259. },
  260. failed,
  261. })
  262. after.start()
  263. await vi.waitFor(() => { expect(carrierFailed).toHaveBeenCalledOnce() })
  264. expect(carrierFailed.mock.calls[0]?.[0]).toMatchObject({
  265. message: 'session control stream ended without a terminal result',
  266. })
  267. expect(failed).not.toHaveBeenCalled()
  268. await after.dispose()
  269. })
  270. })