transport.client.spec.ts 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321
  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 snapshot(
  42. cursor: number,
  43. events: readonly SessionEventEntry[],
  44. hasMore = false,
  45. ): SessionFollowFrame {
  46. return {
  47. type: 'snapshot',
  48. header: {
  49. version: 0,
  50. id: ADDRESS.kind === 'session' ? ADDRESS.sessionId : ADDRESS.childSessionId,
  51. createdAt: 0,
  52. },
  53. cursor,
  54. events,
  55. hasMore,
  56. projections: { asOfSeq: cursor, values: {} },
  57. }
  58. }
  59. function sessionClient(remote: SessionTransportRemote) {
  60. return {
  61. session: remote as SessionRemote,
  62. $stream: <Item>(options: RemoteStreamOptions<Item>) => (
  63. new RemoteStream(AVAILABLE_CONNECTION, options)
  64. ),
  65. }
  66. }
  67. interface FollowGeneration {
  68. readonly frames: readonly SessionFollowFrame[]
  69. readonly terminal?: Error
  70. readonly hold?: boolean
  71. readonly waitAfterFrames?: Promise<void>
  72. }
  73. class ScriptedSessionRemote implements SessionTransportRemote {
  74. readonly followRequests: SessionFollowRequest[] = []
  75. readonly pageRequests: SessionPageRequest[] = []
  76. readonly signals: AbortSignal[] = []
  77. constructor(
  78. private readonly generations: FollowGeneration[],
  79. private readonly pages: RemoteResult<SessionPage>[],
  80. private readonly controlFrames: readonly SessionControlFrame[] = [],
  81. private readonly holdControl = true,
  82. ) {}
  83. async *follow(request: SessionFollowRequest, signal = new AbortController().signal): AsyncIterable<SessionFollowFrame> {
  84. const generation = this.generations.shift()
  85. if (generation === undefined) throw new Error('no scripted Session generation')
  86. this.followRequests.push(request)
  87. this.signals.push(signal)
  88. for (const frame of generation.frames) yield frame
  89. await generation.waitAfterFrames
  90. if (generation.terminal !== undefined) throw generation.terminal
  91. if (generation.hold === true && !signal.aborted) {
  92. await new Promise<void>((resolve) => {
  93. signal.addEventListener('abort', () => { resolve() }, { once: true })
  94. })
  95. }
  96. }
  97. page(request: SessionPageRequest): Promise<RemoteResult<SessionPage>> {
  98. this.pageRequests.push(request)
  99. const result = this.pages.shift()
  100. if (result === undefined) throw new Error('no scripted Session page')
  101. return Promise.resolve(result)
  102. }
  103. async *control(signal = new AbortController().signal): AsyncIterable<SessionControlFrame> {
  104. for (const frame of this.controlFrames) yield frame
  105. if (this.holdControl && !signal.aborted) {
  106. await new Promise<void>((resolve) => {
  107. signal.addEventListener('abort', () => { resolve() }, { once: true })
  108. })
  109. }
  110. }
  111. }
  112. describe('Session Client stream adapters', () => {
  113. it('binds an event journal to one address and publishes replace, append, and prepend changes', async () => {
  114. const remote = new ScriptedSessionRemote(
  115. [{
  116. frames: [
  117. snapshot(3, [entry(2), entry(3)], true),
  118. { type: 'event', ...entry(3) },
  119. { type: 'event', ...entry(4) },
  120. ],
  121. hold: true,
  122. }],
  123. [
  124. { ok: true, value: page([entry(0), entry(1)], false) },
  125. ],
  126. )
  127. const changes: SessionJournalChange[] = []
  128. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  129. publish: (change) => { changes.push(change) },
  130. failed: vi.fn(),
  131. })
  132. await stream.open({ maxMessages: 50 })
  133. await vi.waitFor(() => { expect(changes).toHaveLength(2) })
  134. await stream.prepend({ beforeSeq: 2, maxMessages: 50 })
  135. expect(remote.followRequests).toEqual([{ address: ADDRESS, maxMessages: 50 }])
  136. expect(remote.pageRequests).toEqual([
  137. { address: ADDRESS, throughSeq: 4, beforeSeq: 2, maxMessages: 50 },
  138. ])
  139. expect(changes).toMatchObject([
  140. { type: 'replace', entries: [entry(2), entry(3)], hasMore: true },
  141. { type: 'append', entry: entry(4) },
  142. { type: 'prepend', entries: [entry(0), entry(1)], hasMore: false },
  143. ])
  144. await stream.dispose()
  145. expect(remote.signals[0]?.aborted).toBe(true)
  146. })
  147. it('replaces the retained window from each reconnect snapshot', async () => {
  148. const lost = new RemoteStreamCarrierError('lost')
  149. const remote = new ScriptedSessionRemote(
  150. [
  151. {
  152. frames: [snapshot(1, [entry(0), entry(1)]), { type: 'event', ...entry(2) }],
  153. terminal: lost,
  154. },
  155. { frames: [snapshot(4, [entry(0), entry(1), entry(2), entry(3), entry(4)])], hold: true },
  156. ],
  157. [],
  158. )
  159. const changes: SessionJournalChange[] = []
  160. const carrierFailed = vi.fn()
  161. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  162. publish: (change) => { changes.push(change) },
  163. carrierFailed,
  164. failed: vi.fn(),
  165. })
  166. await stream.open({ maxMessages: 50 })
  167. await vi.waitFor(() => { expect(remote.followRequests).toHaveLength(2) })
  168. expect(remote.followRequests).toEqual([
  169. { address: ADDRESS, maxMessages: 50 },
  170. { address: ADDRESS, maxMessages: 50 },
  171. ])
  172. expect(remote.pageRequests).toEqual([])
  173. expect(changes.map(change => change.type)).toEqual(['replace', 'append', 'replace'])
  174. expect(carrierFailed).toHaveBeenCalledWith(lost)
  175. await stream.dispose()
  176. })
  177. it('repairs a resumed event stream without an optional message limit', async () => {
  178. const finish = Promise.withResolvers<undefined>()
  179. const remote = new ScriptedSessionRemote(
  180. [
  181. {
  182. frames: [snapshot(0, [entry(0)])],
  183. waitAfterFrames: finish.promise,
  184. terminal: new RemoteStreamCarrierError('lost'),
  185. },
  186. { frames: [snapshot(1, [entry(0), entry(1)])], hold: true },
  187. ],
  188. [],
  189. )
  190. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  191. publish: vi.fn(),
  192. failed: vi.fn(),
  193. })
  194. await stream.open({})
  195. finish.resolve(undefined)
  196. await vi.waitFor(() => { expect(remote.followRequests).toHaveLength(2) })
  197. expect(remote.followRequests).toEqual([{ address: ADDRESS }, { address: ADDRESS }])
  198. expect(remote.pageRequests).toEqual([])
  199. await stream.dispose()
  200. })
  201. it('repairs a live gap without adding an absent message limit', async () => {
  202. const remote = new ScriptedSessionRemote(
  203. [{ frames: [snapshot(0, [entry(0)]), { type: 'event', ...entry(2) }], hold: true }],
  204. [{ ok: true, value: page([entry(0), entry(1), entry(2)]) }],
  205. )
  206. const changes: SessionJournalChange[] = []
  207. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  208. publish: (change) => { changes.push(change) },
  209. failed: vi.fn(),
  210. })
  211. await stream.open({})
  212. await vi.waitFor(() => { expect(changes).toHaveLength(2) })
  213. expect(remote.pageRequests).toEqual([{ address: ADDRESS, throughSeq: 2 }])
  214. await stream.dispose()
  215. })
  216. it('turns a pagination failure into a typed stream failure', async () => {
  217. const failure = { code: 'session-not-found', message: 'missing', details: { sessionId: 'session-1' } } as const
  218. const remote = new ScriptedSessionRemote(
  219. [{ frames: [snapshot(-1, [])], hold: true }],
  220. [{ ok: false, error: failure }],
  221. )
  222. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  223. publish: vi.fn(),
  224. failed: vi.fn(),
  225. })
  226. await stream.open({})
  227. await expect(stream.prepend({})).rejects.toBeInstanceOf(RemoteStreamError)
  228. await expect(stream.open({})).rejects.toThrow('already opened')
  229. expect(sessionStreamFailure(new RemoteStreamError(failure.code, failure.message, failure.details)))
  230. .toEqual(failure)
  231. expect(sessionStreamFailure(new Error('local'))).toBeUndefined()
  232. expect(remote.signals[0]?.aborted).toBe(false)
  233. expect(remote.pageRequests).toEqual([{ address: ADDRESS, throughSeq: -1 }])
  234. await stream.dispose()
  235. expect(remote.signals[0]?.aborted).toBe(true)
  236. })
  237. it('maps the Host-wide control baseline and deltas into one snapshot stream', async () => {
  238. const baseline: SessionControlFrame = {
  239. type: 'baseline',
  240. value: { queues: {}, jobs: {}, projections: {} },
  241. }
  242. const update: SessionControlFrame = {
  243. type: 'queue', sessionId: 'session-1' as never, items: [],
  244. }
  245. const remote = new ScriptedSessionRemote([], [], [baseline, update])
  246. const accept = vi.fn<(frame: SessionControlFrame) => void>()
  247. const stream = createSessionControlStream(sessionClient(remote), {
  248. accept,
  249. failed: vi.fn(),
  250. })
  251. stream.start()
  252. stream.start()
  253. await vi.waitFor(() => { expect(accept).toHaveBeenCalledTimes(2) })
  254. expect(accept.mock.calls.map(([frame]) => frame)).toEqual([baseline, update])
  255. await stream.dispose()
  256. await stream.dispose()
  257. })
  258. it('classifies control streams that end before and after their opening baseline', async () => {
  259. const beforeFailed = vi.fn()
  260. const before = createSessionControlStream(
  261. sessionClient(new ScriptedSessionRemote([], [], [], false)),
  262. { accept: vi.fn(), failed: beforeFailed },
  263. )
  264. before.start()
  265. await vi.waitFor(() => { expect(beforeFailed).toHaveBeenCalledOnce() })
  266. expect(beforeFailed.mock.calls[0]?.[0]).toMatchObject({
  267. message: 'session control stream ended before its opening snapshot',
  268. })
  269. await before.dispose()
  270. const baseline: SessionControlFrame = {
  271. type: 'baseline',
  272. value: { queues: {}, jobs: {}, projections: {} },
  273. }
  274. const carrierFailed = vi.fn()
  275. const failed = vi.fn()
  276. const afterRemote = new ScriptedSessionRemote([], [], [baseline], false)
  277. const after = createSessionControlStream(sessionClient(afterRemote), {
  278. accept: vi.fn(),
  279. carrierFailed: (error) => {
  280. carrierFailed(error)
  281. void after.dispose()
  282. },
  283. failed,
  284. })
  285. after.start()
  286. await vi.waitFor(() => { expect(carrierFailed).toHaveBeenCalledOnce() })
  287. expect(carrierFailed.mock.calls[0]?.[0]).toMatchObject({
  288. message: 'session control stream ended without a terminal result',
  289. })
  290. expect(failed).not.toHaveBeenCalled()
  291. await after.dispose()
  292. })
  293. })