transport.client.spec.ts 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565
  1. import { describe, expect, it, vi } from 'vitest'
  2. import {
  3. RemoteStream,
  4. RemoteStreamCarrierError,
  5. type RemoteStreamOptions,
  6. } from '@deepseek-ai/dsh-api-gateway/client'
  7. import { RemoteError } from '@deepseek-ai/dsh-typert-protocol'
  8. import { LlmAttemptId } from '@deepseek-ai/dsh-llm'
  9. import { SESSION_FORMAT_VERSION } from '@deepseek-ai/dsh-session/types'
  10. import type { RemoteResult } from '@deepseek-ai/dsh-typert-protocol'
  11. import {
  12. createSessionControlStream,
  13. SessionEventStream,
  14. type SessionJournalChange,
  15. type SessionRemote,
  16. } from '../src/client/index.ts'
  17. import type { SessionRemotes } from '../src/client/sessions/remotes.ts'
  18. import type {
  19. SessionAddress,
  20. SessionAssistantStreamBaseline,
  21. SessionAssistantStreamFrame,
  22. SessionControlFrame,
  23. SessionEventEntry,
  24. SessionFollowFrame,
  25. SessionFollowRequest,
  26. SessionHistoryRecord,
  27. SessionPage,
  28. SessionPageRequest,
  29. } from '../src/types.ts'
  30. type SessionTransportRemote = Pick<SessionRemote, 'control' | 'follow' | 'page'>
  31. const ADDRESS: SessionAddress = { kind: 'session', sessionId: 'session-1' as never }
  32. const AVAILABLE_CONNECTION = {
  33. generation: {
  34. getSnapshot: () => ({ id: 1, host: { home: '/home/fixture' } }),
  35. subscribe: () => () => {},
  36. },
  37. }
  38. function entry(seq: number): SessionEventEntry {
  39. return { type: 'event', event: { type: 'turn/start', seq, time: seq, data: { turn: seq } } }
  40. }
  41. function page(records: readonly SessionHistoryRecord[], hasMore = false): SessionPage {
  42. return { records, hasMore }
  43. }
  44. function snapshot(
  45. cursor: number,
  46. records: readonly SessionHistoryRecord[],
  47. hasMore = false,
  48. assistantStream: SessionAssistantStreamBaseline = { revision: 0 },
  49. ): SessionFollowFrame {
  50. return {
  51. type: 'snapshot',
  52. header: {
  53. version: SESSION_FORMAT_VERSION,
  54. id: ADDRESS.kind === 'session' ? ADDRESS.sessionId : ADDRESS.childSessionId,
  55. createdAt: 0,
  56. isSeeded: false,
  57. },
  58. cursor,
  59. records,
  60. hasMore,
  61. projections: { asOfSeq: cursor, values: {} },
  62. assistantStream,
  63. }
  64. }
  65. function assistantFrame(frame: SessionAssistantStreamFrame): SessionFollowFrame {
  66. return { type: 'assistant-stream', frame }
  67. }
  68. function sessionClient(remote: SessionTransportRemote): SessionRemotes {
  69. return {
  70. session: remote as SessionRemote,
  71. $stream: <Item>(options: RemoteStreamOptions<Item>) => (
  72. new RemoteStream(AVAILABLE_CONNECTION, options)
  73. ),
  74. commands: { execute: () => Promise.reject(new Error('stream tests never run commands')) },
  75. subagents: {
  76. list: () => Promise.reject(new Error('stream tests never read the subagent catalog')),
  77. prompt: () => Promise.reject(new Error('stream tests never prompt a subagent')),
  78. interruptByParent: () => Promise.reject(new Error('stream tests never interrupt a subagent')),
  79. },
  80. }
  81. }
  82. interface FollowGeneration {
  83. readonly frames: readonly SessionFollowFrame[]
  84. readonly terminal?: Error
  85. readonly hold?: boolean
  86. readonly waitAfterFrames?: Promise<void>
  87. }
  88. class ScriptedSessionRemote implements SessionTransportRemote {
  89. readonly followRequests: SessionFollowRequest[] = []
  90. readonly pageRequests: SessionPageRequest[] = []
  91. readonly signals: AbortSignal[] = []
  92. constructor(
  93. private readonly generations: FollowGeneration[],
  94. private readonly pages: RemoteResult<SessionPage>[],
  95. private readonly controlFrames: readonly SessionControlFrame[] = [],
  96. private readonly holdControl = true,
  97. ) {}
  98. async *follow(request: SessionFollowRequest, signal = new AbortController().signal): AsyncIterable<SessionFollowFrame> {
  99. const generation = this.generations.shift()
  100. if (generation === undefined) throw new Error('no scripted Session generation')
  101. this.followRequests.push(request)
  102. this.signals.push(signal)
  103. for (const frame of generation.frames) yield frame
  104. await generation.waitAfterFrames
  105. if (generation.terminal !== undefined) throw generation.terminal
  106. if (generation.hold === true && !signal.aborted) {
  107. await new Promise<void>((resolve) => {
  108. signal.addEventListener('abort', () => { resolve() }, { once: true })
  109. })
  110. }
  111. }
  112. page(request: SessionPageRequest): Promise<RemoteResult<SessionPage>> {
  113. this.pageRequests.push(request)
  114. const result = this.pages.shift()
  115. if (result === undefined) throw new Error('no scripted Session page')
  116. return Promise.resolve(result)
  117. }
  118. async *control(signal = new AbortController().signal): AsyncIterable<SessionControlFrame> {
  119. for (const frame of this.controlFrames) yield frame
  120. if (this.holdControl && !signal.aborted) {
  121. await new Promise<void>((resolve) => {
  122. signal.addEventListener('abort', () => { resolve() }, { once: true })
  123. })
  124. }
  125. }
  126. }
  127. describe('Session Client stream adapters', () => {
  128. it('opts into assistant notifications and publishes the reconnect baseline plus live frame', async () => {
  129. const attemptId = LlmAttemptId('transport-attempt')
  130. const baseline: SessionAssistantStreamBaseline = {
  131. revision: 2,
  132. activeAttempt: {
  133. attemptId,
  134. startedAfterSeq: -1,
  135. turn: 1,
  136. step: 1,
  137. nextIndex: 1,
  138. stream: [{ type: 'text-chunks', time0: 0, index: 0, dt: [], texts: ['a'] }],
  139. },
  140. }
  141. const frame: SessionAssistantStreamFrame = {
  142. type: 'chunk', attemptId, revision: 3, index: 1,
  143. time: 1, chunk: { type: 'text-delta', index: 0, text: 'b' },
  144. }
  145. const remote = new ScriptedSessionRemote(
  146. [{ frames: [snapshot(0, [entry(0)], false, baseline), assistantFrame(frame)], hold: true }],
  147. [],
  148. )
  149. const changes: SessionJournalChange[] = []
  150. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  151. publish: (change) => { changes.push(change) },
  152. failed: vi.fn(),
  153. })
  154. await stream.open({})
  155. await vi.waitFor(() => { expect(changes).toHaveLength(2) })
  156. expect(remote.followRequests).toEqual([{ address: ADDRESS, assistantStream: true }])
  157. expect(changes).toMatchObject([
  158. { type: 'replace', page: { assistantStream: baseline } },
  159. { type: 'assistant-stream', frame },
  160. ])
  161. await stream.dispose()
  162. })
  163. it('rejects an opted-in opening that omits its Assistant baseline', async () => {
  164. const remote = new ScriptedSessionRemote([{
  165. frames: [{
  166. type: 'snapshot',
  167. header: {
  168. version: SESSION_FORMAT_VERSION,
  169. id: ADDRESS.sessionId,
  170. createdAt: 0,
  171. isSeeded: false,
  172. },
  173. cursor: -1,
  174. records: [],
  175. hasMore: false,
  176. projections: { asOfSeq: -1, values: {} },
  177. }],
  178. }], [])
  179. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  180. publish: vi.fn(),
  181. failed: vi.fn(),
  182. })
  183. try {
  184. await expect(stream.open({})).rejects.toMatchObject({
  185. code: 'gateway/internal',
  186. message: 'session assistant stream omitted its opted-in opening baseline',
  187. })
  188. } finally {
  189. await stream.dispose()
  190. }
  191. })
  192. it('rejects an Assistant frame that arrives before the opening baseline', async () => {
  193. const remote = new ScriptedSessionRemote([{
  194. frames: [assistantFrame({
  195. type: 'start', attemptId: LlmAttemptId('pre-opening-attempt'),
  196. revision: 1, startedAfterSeq: -1, turn: 1, step: 1,
  197. })],
  198. }], [])
  199. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  200. publish: vi.fn(),
  201. failed: vi.fn(),
  202. })
  203. try {
  204. await expect(stream.open({})).rejects.toMatchObject({
  205. code: 'gateway/internal',
  206. message: 'session event stream emitted an entry before its opening cursor',
  207. })
  208. expect(remote.followRequests).toEqual([{ address: ADDRESS, assistantStream: true }])
  209. } finally {
  210. await stream.dispose()
  211. }
  212. })
  213. it('rebaselines after a transient assistant revision gap without advancing the durable cursor', async () => {
  214. const attemptId = LlmAttemptId('gapped-attempt')
  215. const start: SessionAssistantStreamFrame = {
  216. type: 'start', attemptId, revision: 1, startedAfterSeq: -1,
  217. turn: 1, step: 1,
  218. }
  219. const gap: SessionAssistantStreamFrame = {
  220. type: 'chunk', attemptId, revision: 3, index: 0,
  221. time: 1, chunk: { type: 'text-delta', index: 0, text: 'lost predecessor' },
  222. }
  223. const replacement: SessionAssistantStreamBaseline = {
  224. revision: 3,
  225. activeAttempt: {
  226. attemptId,
  227. startedAfterSeq: -1,
  228. turn: 1,
  229. step: 1,
  230. nextIndex: 1,
  231. stream: [{ type: 'text-chunks', time0: 1, index: 0, dt: [], texts: ['lost predecessor'] }],
  232. },
  233. }
  234. const remote = new ScriptedSessionRemote([
  235. {
  236. frames: [snapshot(0, [entry(0)]), assistantFrame(start), assistantFrame(gap)],
  237. },
  238. { frames: [snapshot(0, [entry(0)], false, replacement)], hold: true },
  239. ], [])
  240. const changes: SessionJournalChange[] = []
  241. const carrierFailed = vi.fn()
  242. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  243. publish: (change) => { changes.push(change) },
  244. carrierFailed,
  245. failed: vi.fn(),
  246. })
  247. await stream.open({})
  248. await vi.waitFor(() => { expect(remote.followRequests).toHaveLength(2) })
  249. expect(changes.map(change => change.type)).toEqual([
  250. 'replace', 'assistant-stream', 'replace',
  251. ])
  252. expect(changes.at(-1)).toMatchObject({
  253. type: 'replace', page: { assistantStream: replacement },
  254. })
  255. expect(remote.pageRequests).toEqual([])
  256. expect(carrierFailed).toHaveBeenCalledWith(expect.objectContaining({
  257. message: 'session assistant stream skipped revision 2',
  258. }))
  259. await stream.dispose()
  260. })
  261. it('rebaselines when a replacement Agent lifecycle restarts at revision one', async () => {
  262. const attemptId = LlmAttemptId('replacement-lifecycle-attempt')
  263. const previous: SessionAssistantStreamBaseline = {
  264. revision: 2,
  265. activeAttempt: {
  266. attemptId,
  267. startedAfterSeq: -1,
  268. turn: 1,
  269. step: 1,
  270. nextIndex: 1,
  271. stream: [{ type: 'text-chunks', time0: 1, index: 0, dt: [], texts: ['old'] }],
  272. },
  273. }
  274. const replacementStart: SessionAssistantStreamFrame = {
  275. type: 'start', attemptId, revision: 1, startedAfterSeq: -1,
  276. turn: 2, step: 1,
  277. }
  278. const replacement: SessionAssistantStreamBaseline = {
  279. revision: 1,
  280. activeAttempt: {
  281. attemptId,
  282. startedAfterSeq: -1,
  283. turn: 2,
  284. step: 1,
  285. nextIndex: 0,
  286. stream: [],
  287. },
  288. }
  289. const remote = new ScriptedSessionRemote([
  290. {
  291. frames: [snapshot(0, [entry(0)], false, previous), assistantFrame(replacementStart)],
  292. },
  293. { frames: [snapshot(0, [entry(0)], false, replacement)], hold: true },
  294. ], [])
  295. const changes: SessionJournalChange[] = []
  296. const carrierFailed = vi.fn()
  297. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  298. publish: (change) => { changes.push(change) },
  299. carrierFailed,
  300. failed: vi.fn(),
  301. })
  302. try {
  303. await stream.open({})
  304. await vi.waitFor(() => { expect(remote.followRequests).toHaveLength(2) })
  305. expect(changes).toMatchObject([
  306. { type: 'replace', page: { assistantStream: previous } },
  307. { type: 'replace', page: { assistantStream: replacement } },
  308. ])
  309. expect(carrierFailed).toHaveBeenCalledWith(expect.objectContaining({
  310. message: 'session assistant stream skipped revision 3',
  311. }))
  312. } finally {
  313. await stream.dispose()
  314. }
  315. })
  316. it('validates one scalar current-event range before publishing Client entries', async () => {
  317. const remote = new ScriptedSessionRemote(
  318. [{ frames: [snapshot(2, [entry(0), entry(1), entry(2)]), entry(3)], hold: true }],
  319. [],
  320. )
  321. const changes: SessionJournalChange[] = []
  322. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  323. publish: (change) => { changes.push(change) },
  324. failed: vi.fn(),
  325. })
  326. await stream.open({})
  327. await vi.waitFor(() => { expect(changes).toHaveLength(2) })
  328. expect(changes[0]).toMatchObject({
  329. type: 'replace',
  330. entries: [
  331. entry(0),
  332. entry(1),
  333. entry(2),
  334. ],
  335. })
  336. expect(changes[1]).toEqual({ type: 'append', entry: entry(3) })
  337. await stream.dispose()
  338. })
  339. it('binds an event journal to one address and publishes replace, append, and prepend changes', async () => {
  340. const remote = new ScriptedSessionRemote(
  341. [{
  342. frames: [
  343. snapshot(3, [entry(2), entry(3)], true),
  344. entry(3),
  345. entry(4),
  346. ],
  347. hold: true,
  348. }],
  349. [
  350. { ok: true, value: page([entry(0), entry(1)], false) },
  351. ],
  352. )
  353. const changes: SessionJournalChange[] = []
  354. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  355. publish: (change) => { changes.push(change) },
  356. failed: vi.fn(),
  357. })
  358. await stream.open({ maxMessages: 50 })
  359. await vi.waitFor(() => { expect(changes).toHaveLength(2) })
  360. await stream.prepend({ beforeSeq: 2, maxMessages: 50 })
  361. expect(remote.followRequests).toEqual([{
  362. address: ADDRESS, assistantStream: true, maxMessages: 50,
  363. }])
  364. expect(remote.pageRequests).toEqual([
  365. { address: ADDRESS, throughSeq: 4, beforeSeq: 2, maxMessages: 50 },
  366. ])
  367. expect(changes).toMatchObject([
  368. { type: 'replace', entries: [entry(2), entry(3)], hasMore: true },
  369. { type: 'append', entry: entry(4) },
  370. { type: 'prepend', entries: [entry(0), entry(1)], hasMore: false },
  371. ])
  372. await stream.dispose()
  373. expect(remote.signals[0]?.aborted).toBe(true)
  374. })
  375. it('replaces the retained window from each reconnect snapshot', async () => {
  376. const lost = new RemoteStreamCarrierError('lost')
  377. const remote = new ScriptedSessionRemote(
  378. [
  379. {
  380. frames: [snapshot(1, [entry(0), entry(1)]), entry(2)],
  381. terminal: lost,
  382. },
  383. { frames: [snapshot(4, [entry(0), entry(1), entry(2), entry(3), entry(4)])], hold: true },
  384. ],
  385. [],
  386. )
  387. const changes: SessionJournalChange[] = []
  388. const carrierFailed = vi.fn()
  389. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  390. publish: (change) => { changes.push(change) },
  391. carrierFailed,
  392. failed: vi.fn(),
  393. })
  394. await stream.open({ maxMessages: 50 })
  395. await vi.waitFor(() => { expect(remote.followRequests).toHaveLength(2) })
  396. expect(remote.followRequests).toEqual([
  397. { address: ADDRESS, assistantStream: true, maxMessages: 50 },
  398. { address: ADDRESS, assistantStream: true, maxMessages: 50 },
  399. ])
  400. expect(remote.pageRequests).toEqual([])
  401. expect(changes.map(change => change.type)).toEqual(['replace', 'append', 'replace'])
  402. expect(carrierFailed).toHaveBeenCalledWith(lost)
  403. await stream.dispose()
  404. })
  405. it('repairs a resumed event stream without an optional message limit', async () => {
  406. const finish = Promise.withResolvers<undefined>()
  407. const remote = new ScriptedSessionRemote(
  408. [
  409. {
  410. frames: [snapshot(0, [entry(0)])],
  411. waitAfterFrames: finish.promise,
  412. terminal: new RemoteStreamCarrierError('lost'),
  413. },
  414. { frames: [snapshot(1, [entry(0), entry(1)])], hold: true },
  415. ],
  416. [],
  417. )
  418. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  419. publish: vi.fn(),
  420. failed: vi.fn(),
  421. })
  422. await stream.open({})
  423. finish.resolve(undefined)
  424. await vi.waitFor(() => { expect(remote.followRequests).toHaveLength(2) })
  425. expect(remote.followRequests).toEqual([
  426. { address: ADDRESS, assistantStream: true },
  427. { address: ADDRESS, assistantStream: true },
  428. ])
  429. expect(remote.pageRequests).toEqual([])
  430. await stream.dispose()
  431. })
  432. it('repairs a live gap without adding an absent message limit', async () => {
  433. const remote = new ScriptedSessionRemote(
  434. [{ frames: [snapshot(0, [entry(0)]), entry(2)], hold: true }],
  435. [{ ok: true, value: page([entry(0), entry(1), entry(2)]) }],
  436. )
  437. const changes: SessionJournalChange[] = []
  438. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  439. publish: (change) => { changes.push(change) },
  440. failed: vi.fn(),
  441. })
  442. await stream.open({})
  443. await vi.waitFor(() => { expect(changes).toHaveLength(2) })
  444. expect(remote.pageRequests).toEqual([{ address: ADDRESS, throughSeq: 2 }])
  445. await stream.dispose()
  446. })
  447. it('turns a pagination failure into a typed stream failure', async () => {
  448. const failure = new RemoteError('session/not-found', 'missing', { sessionId: 'session-1' as never })
  449. const remote = new ScriptedSessionRemote(
  450. [{ frames: [snapshot(-1, [])], hold: true }],
  451. [{ ok: false, error: failure }],
  452. )
  453. const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
  454. publish: vi.fn(),
  455. failed: vi.fn(),
  456. })
  457. await stream.open({})
  458. await expect(stream.prepend({})).rejects.toMatchObject({ code: 'session/not-found' })
  459. await expect(stream.open({})).rejects.toThrow('already opened')
  460. expect(remote.signals[0]?.aborted).toBe(false)
  461. expect(remote.pageRequests).toEqual([{ address: ADDRESS, throughSeq: -1 }])
  462. await stream.dispose()
  463. expect(remote.signals[0]?.aborted).toBe(true)
  464. })
  465. it('maps the Host-wide control baseline and deltas into one snapshot stream', async () => {
  466. const baseline: SessionControlFrame = {
  467. type: 'baseline',
  468. value: { queues: {}, jobs: {}, projections: {} },
  469. }
  470. const update: SessionControlFrame = {
  471. type: 'queue', sessionId: 'session-1' as never, items: [],
  472. }
  473. const remote = new ScriptedSessionRemote([], [], [baseline, update])
  474. const accept = vi.fn<(frame: SessionControlFrame) => void>()
  475. const stream = createSessionControlStream(sessionClient(remote), {
  476. accept,
  477. failed: vi.fn(),
  478. })
  479. stream.start()
  480. stream.start()
  481. await vi.waitFor(() => { expect(accept).toHaveBeenCalledTimes(2) })
  482. expect(accept.mock.calls.map(([frame]) => frame)).toEqual([baseline, update])
  483. await stream.dispose()
  484. await stream.dispose()
  485. })
  486. it('classifies control streams that end before and after their opening baseline', async () => {
  487. const beforeFailed = vi.fn()
  488. const before = createSessionControlStream(
  489. sessionClient(new ScriptedSessionRemote([], [], [], false)),
  490. { accept: vi.fn(), failed: beforeFailed },
  491. )
  492. before.start()
  493. await vi.waitFor(() => { expect(beforeFailed).toHaveBeenCalledOnce() })
  494. expect(beforeFailed.mock.calls[0]?.[0]).toMatchObject({
  495. message: 'session control stream ended before its opening snapshot',
  496. })
  497. await before.dispose()
  498. const baseline: SessionControlFrame = {
  499. type: 'baseline',
  500. value: { queues: {}, jobs: {}, projections: {} },
  501. }
  502. const carrierFailed = vi.fn()
  503. const failed = vi.fn()
  504. const afterRemote = new ScriptedSessionRemote([], [], [baseline], false)
  505. const after = createSessionControlStream(sessionClient(afterRemote), {
  506. accept: vi.fn(),
  507. carrierFailed: (error) => {
  508. carrierFailed(error)
  509. void after.dispose()
  510. },
  511. failed,
  512. })
  513. after.start()
  514. await vi.waitFor(() => { expect(carrierFailed).toHaveBeenCalledOnce() })
  515. expect(carrierFailed.mock.calls[0]?.[0]).toMatchObject({
  516. message: 'session control stream ended without a terminal result',
  517. })
  518. expect(failed).not.toHaveBeenCalled()
  519. await after.dispose()
  520. })
  521. })