transport.client.spec.ts 13 KB

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