transport.client.spec.ts 13 KB

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