session.client.spec.ts 47 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102
  1. /** Session object lifecycle, event-window transport, commands, and resync behavior. */
  2. import { afterEach, describe, expect, it, vi } from 'vitest'
  3. import type { SessionEvent } from '@deepseek-ai/dsh-session/types'
  4. import type { SessionId } from '@deepseek-ai/dsh-api-remotes/client'
  5. import { RemoteStreamCarrierError } from '@deepseek-ai/dsh-api-gateway/client'
  6. import { RemoteError } from '@deepseek-ai/dsh-typert-protocol'
  7. import { JUMP_PAGE_MESSAGES, Session, type SessionOptions } from '../src/client/sessions/session.ts'
  8. import { FakeApiClient, deferred, err, fakeRemote, ok } from './fake-api.client.ts'
  9. import { entries, ev, historyValue, plainTurn } from './event-script.client.ts'
  10. const SID = 'fk-s1' as SessionId
  11. const PARENT = 'fk-parent' as SessionId
  12. afterEach(() => {
  13. vi.unstubAllGlobals()
  14. })
  15. function makeSession(
  16. api = new FakeApiClient(),
  17. options: SessionOptions = {},
  18. ): { api: FakeApiClient; session: Session } {
  19. return { api, session: new Session(SID, fakeRemote(api), options) }
  20. }
  21. function follow(
  22. api: FakeApiClient,
  23. event: SessionEvent,
  24. ): Promise<void> {
  25. return api.pushFollow(SID, {
  26. type: 'event',
  27. event: event as never,
  28. })
  29. }
  30. function windowEntries(session: Session) {
  31. return session.eventSource.getSnapshot().entries
  32. }
  33. function eventSeqs(session: Session): number[] {
  34. return windowEntries(session).map(entry => entry.event.seq)
  35. }
  36. function histResponse(events: SessionEvent[], hasMore = false) {
  37. return Promise.resolve(ok(historyValue(events, hasMore)))
  38. }
  39. describe('Session file upload', () => {
  40. it('uses the background body carrier, reports progress, and validates its receipt', async () => {
  41. const progress = vi.fn()
  42. const post = vi.fn(async (request: {
  43. path: string
  44. body: Blob
  45. headers?: Readonly<Record<string, string>>
  46. signal?: AbortSignal
  47. onProgress?: (progress: { loaded: number; total?: number }) => void
  48. }) => {
  49. request.onProgress?.({ loaded: 2, total: 4 })
  50. return {
  51. status: 200,
  52. body: JSON.stringify({
  53. ok: true,
  54. value: {
  55. receiptId: 'receipt-1',
  56. file: { attachmentId: 'file-1', name: 'notes & refs.pdf', bytes: 4 },
  57. },
  58. }),
  59. }
  60. })
  61. const { api, session } = makeSession(undefined, { backgroundUploads: { post } })
  62. const abort = new AbortController()
  63. const file = new Blob([Uint8Array.of(1, 2, 3, 4)])
  64. await expect(session.uploadFile(file, 'notes & refs.pdf', abort.signal, progress)).resolves.toEqual({
  65. ok: true,
  66. value: {
  67. receiptId: 'receipt-1',
  68. file: { attachmentId: 'file-1', name: 'notes & refs.pdf', bytes: 4 },
  69. },
  70. })
  71. expect(post).toHaveBeenCalledWith(expect.objectContaining({
  72. path: '/api/session/uploadFileBinary?sessionId=fk-s1&name=notes+%26+refs.pdf',
  73. body: file,
  74. headers: { 'content-type': 'application/octet-stream' },
  75. signal: abort.signal,
  76. }))
  77. expect(progress).toHaveBeenCalledWith({ loaded: 2, total: 4 })
  78. expect(api.callsOf('session.uploadFile')).toEqual([])
  79. })
  80. it('preserves a background business failure and supports unnamed files without observers', async () => {
  81. const post = vi.fn((_request: { readonly body: Blob }) => Promise.resolve({
  82. status: 200,
  83. body: JSON.stringify({
  84. ok: false,
  85. error: { code: 'session/attachment-invalid', message: 'denied', details: { reason: 'NOPE' } },
  86. }),
  87. }))
  88. const { session } = makeSession(undefined, { backgroundUploads: { post } })
  89. await expect(session.uploadFile(new Blob([]))).resolves.toMatchObject({
  90. ok: false,
  91. error: { code: 'session/attachment-invalid', message: 'denied', details: { reason: 'NOPE' } },
  92. })
  93. const request = post.mock.calls[0]?.[0]
  94. expect(request).toMatchObject({
  95. path: '/api/session/uploadFileBinary?sessionId=fk-s1',
  96. headers: { 'content-type': 'application/octet-stream' },
  97. })
  98. expect(request?.body).toBeInstanceOf(Blob)
  99. })
  100. it('keeps byte and Blob fallback uploads on the generated Remote carrier', async () => {
  101. const { api, session } = makeSession()
  102. await expect(session.uploadFile(Uint8Array.of(0, 0, 0), 'bytes.bin')).resolves.toMatchObject({ ok: true })
  103. await expect(session.uploadFile(new Blob([Uint8Array.of(1)]))).resolves.toMatchObject({ ok: true })
  104. expect(api.callsOf('session.uploadFile')).toEqual([
  105. { sessionId: SID, data: 'AAAA', name: 'bytes.bin' },
  106. { sessionId: SID, data: 'AQ==' },
  107. ])
  108. })
  109. it('folds non-200 and malformed background responses into transport failures', async () => {
  110. const bodies: unknown[] = [
  111. null,
  112. { ok: 'yes' },
  113. { ok: false, error: null },
  114. { ok: false, error: { code: 1, message: 'x', details: {} } },
  115. { ok: false, error: { code: 'x', message: 1, details: {} } },
  116. { ok: false, error: { code: 'x', message: 'x', details: null } },
  117. { ok: true, value: null },
  118. { ok: true, value: { receiptId: 1, file: {} } },
  119. { ok: true, value: { receiptId: 'r', file: null } },
  120. { ok: true, value: { receiptId: 'r', file: { attachmentId: 1, name: 'x', bytes: 1 } } },
  121. { ok: true, value: { receiptId: 'r', file: { attachmentId: 'a', name: 1, bytes: 1 } } },
  122. { ok: true, value: { receiptId: 'r', file: { attachmentId: 'a', name: 'x', bytes: '1' } } },
  123. { ok: true, value: { receiptId: 'r', file: { attachmentId: 'a', name: 'x', bytes: 1.5 } } },
  124. { ok: true, value: { receiptId: 'r', file: { attachmentId: 'a', name: 'x', bytes: -1 } } },
  125. ]
  126. for (const body of bodies) {
  127. const { session } = makeSession(undefined, {
  128. backgroundUploads: { post: () => Promise.resolve({ status: 200, body: JSON.stringify(body) }) },
  129. })
  130. await expect(session.uploadFile(new Blob([])))
  131. .rejects.toThrow(/file upload transport returned an invalid/)
  132. }
  133. const { session } = makeSession(undefined, {
  134. backgroundUploads: { post: () => Promise.resolve({ status: 503, body: 'unavailable' }) },
  135. })
  136. await expect(session.uploadFile(new Blob([])))
  137. .rejects.toThrow('file upload transport failed with HTTP 503')
  138. })
  139. it('refuses a direct file upload for a continuable subagent before either carrier runs', async () => {
  140. const post = vi.fn()
  141. const api = new FakeApiClient()
  142. const session = new Session(SID, fakeRemote(api), {
  143. address: { parentSessionId: PARENT, childSessionId: SID, mode: 'continuable' },
  144. backgroundUploads: { post },
  145. })
  146. await expect(session.uploadFile(new Blob([]))).resolves.toMatchObject({
  147. ok: false,
  148. error: { code: 'subagent/attachment-invalid', details: { reason: 'SUBAGENT_FILE_UNSUPPORTED' } },
  149. })
  150. expect(post).not.toHaveBeenCalled()
  151. expect(api.callsOf('session.uploadFile')).toEqual([])
  152. })
  153. })
  154. describe('Session open', () => {
  155. it('keeps a bare Session blank until an authoritative lifecycle signal arrives', () => {
  156. const { session } = makeSession()
  157. expect(session.getSnapshot()).toMatchObject({ blank: true, promptAttempted: false, running: false })
  158. session.handleRunning(true)
  159. expect(session.getSnapshot()).toMatchObject({ blank: false, running: true })
  160. })
  161. it('installs the tail page: cold → loading → open with window and nodes in place', async () => {
  162. const { api, session } = makeSession()
  163. const page = plainTurn(10, 3, '问', '答')
  164. api.onHistory = () => histResponse(page, true)
  165. expect(session.getSnapshot().openState).toBe('cold')
  166. const opening = session.open()
  167. expect(session.getSnapshot().openState).toBe('loading')
  168. await opening
  169. const snapshot = session.getSnapshot()
  170. expect(snapshot.openState).toBe('open')
  171. expect(snapshot.hasMore).toBe(true)
  172. expect(eventSeqs(session)).toEqual([10, 11, 12, 13, 14, 15])
  173. expect(session.eventSource.getSnapshot().change).toMatchObject({ kind: 'replace' })
  174. })
  175. it('is idempotent: concurrent opens share one follow, reopening when open is a no-op', async () => {
  176. const { api, session } = makeSession()
  177. await Promise.all([session.open(), session.open()])
  178. await session.open()
  179. expect(api.callsOf('session.follow')).toHaveLength(1)
  180. expect(api.callsOf('session.history')).toEqual([])
  181. })
  182. it('lands an error result in openState=error with the Remote failure kept', async () => {
  183. const { api, session } = makeSession()
  184. api.onHistory = () => Promise.resolve(err(new RemoteError('session/not-found', 'gone', { sessionId: SID })))
  185. await session.open()
  186. const snapshot = session.getSnapshot()
  187. expect(snapshot.openState).toBe('error')
  188. expect(snapshot.openError?.code).toBe('session/not-found')
  189. })
  190. it('lands exhausted carrier retries in openState=error as gateway/internal', async () => {
  191. const { api, session } = makeSession()
  192. // Two consecutive carrier losses before any opening is accepted exhaust the
  193. // Gateway's retry budget; the escaping failure crosses the stream boundary marked.
  194. api.onHistory = () => Promise.reject(new RemoteStreamCarrierError('history carrier down'))
  195. await session.open()
  196. expect(session.getSnapshot().openState).toBe('error')
  197. expect(session.getSnapshot().openError).toMatchObject({
  198. code: 'gateway/internal', message: 'history carrier down',
  199. })
  200. expect(api.followStarts).toHaveLength(2)
  201. })
  202. it('lands a packed live record in openState=error instead of crashing the stream loop', async () => {
  203. const { api, session } = makeSession()
  204. api.onHistory = () => histResponse(plainTurn(0, 0, 'a', 'b'))
  205. await session.open()
  206. expect(session.getSnapshot().openState).toBe('open')
  207. // The live tail may carry only events; a packed record breaks that contract.
  208. await api.pushFollow(SID, {
  209. type: 'chunks',
  210. event: {
  211. type: 'chunkrow/text-chunks',
  212. seq: 6,
  213. time: 6,
  214. data: { turn: 1, step: 1, index: 0, texts: ['a'], dt: [] },
  215. },
  216. } as never)
  217. await vi.waitFor(() => { expect(session.getSnapshot().openState).toBe('error') })
  218. expect(session.getSnapshot().openError).toMatchObject({
  219. code: 'gateway/internal', message: 'session live stream emitted a packed history record',
  220. })
  221. })
  222. it('lands a Gateway-marked stream failure in openState=error', async () => {
  223. const { api, session } = makeSession()
  224. api.onHistory = () => Promise.reject(new Error('socket died'))
  225. await session.open()
  226. expect(session.getSnapshot().openState).toBe('error')
  227. expect(session.getSnapshot().openError).toMatchObject({ code: 'gateway/internal', message: 'socket died' })
  228. })
  229. it('stitches live frames arriving while history is pending, dropping the page overlap', async () => {
  230. const { api, session } = makeSession()
  231. const gate = deferred<Awaited<ReturnType<FakeApiClient['onHistory']>>>()
  232. api.onHistory = () => gate.promise
  233. const opening = session.open()
  234. // Three live frames land while the opening snapshot is pending; seq 15 overlaps its tail.
  235. const page = plainTurn(10, 0, '早', '安')
  236. const deliveries = [
  237. follow(api, ev.turnStart(15, 1)),
  238. follow(api, ev.user(16, '插进来的')),
  239. ]
  240. gate.resolve(ok({
  241. records: entries(page) as never[],
  242. hasMore: false,
  243. modelSelection: { provider: 'deepseek-official', model: 'deepseek-v4-flash' },
  244. }))
  245. await Promise.all([opening, ...deliveries])
  246. const seqs = eventSeqs(session)
  247. // Overlapping seq-15 frame (== page tail turn/end) was dropped; 16 appended once.
  248. expect(seqs).toEqual([10, 11, 12, 13, 14, 15, 16])
  249. })
  250. })
  251. describe('live event path', () => {
  252. async function opened(events: SessionEvent[] = plainTurn(0, 0, 'a', 'b')) {
  253. const { api, session } = makeSession()
  254. api.onHistory = () => histResponse(events)
  255. await session.open()
  256. return { api, session }
  257. }
  258. it('drops replayed frames at or below the window tail', async () => {
  259. const { api, session } = await opened()
  260. const before = session.eventSource.getSnapshot()
  261. await follow(api, ev.user(3, '重放'))
  262. expect(session.eventSource.getSnapshot()).toBe(before)
  263. })
  264. it('keeps the authoritative host blank bit across unrelated log events', async () => {
  265. const { api, session } = await opened([])
  266. session.handleBlank(true)
  267. await Promise.all([
  268. follow(api, ev.commandRun(0, 'cmd-perm', 'permission', ' danger-full-access')),
  269. follow(api, ev.commandDone(1, 'cmd-perm', 'success', 'preset danger-full-access')),
  270. ])
  271. const snapshot = session.getSnapshot()
  272. expect(eventSeqs(session)).toEqual([0, 1])
  273. expect(snapshot.blank).toBe(true)
  274. })
  275. it('repairs a seq gap by repulling the tail page instead of appending a hole', async () => {
  276. const { api, session } = await opened(plainTurn(0, 0, 'a', 'b')) // tail seq = 5
  277. const repaired = [...plainTurn(0, 0, 'a', 'b'), ...plainTurn(6, 1, 'c', 'd')]
  278. api.onHistory = () => histResponse(repaired)
  279. // seq 9 with tail 5 → gap; the event detours to the buffer and one history refetch fires.
  280. await follow(api, ev.assistant(9, 1, 'd'))
  281. await vi.waitFor(() => {
  282. expect(api.callsOf('session.history')).toHaveLength(1)
  283. })
  284. await vi.waitFor(() => {
  285. expect(eventSeqs(session)).toEqual(
  286. repaired.filter(event => event.seq <= 9).map(event => event.seq),
  287. )
  288. })
  289. })
  290. })
  291. describe('paging', () => {
  292. it('prepends an older page and keeps seq continuity', async () => {
  293. const older = plainTurn(0, 0, '旧问', '旧答')
  294. const newer = plainTurn(6, 1, '新问', '新答')
  295. const { api, session } = makeSession()
  296. api.onHistory = payload => payload.beforeSeq === undefined
  297. ? histResponse(newer, true)
  298. : histResponse(older, false)
  299. await session.open()
  300. await session.loadOlder()
  301. const snapshot = session.getSnapshot()
  302. expect(api.callsOf('session.follow')).toHaveLength(1)
  303. expect(api.callsOf('session.history')).toMatchObject([
  304. { sessionId: SID, throughSeq: 11, beforeSeq: 6 },
  305. ])
  306. expect(snapshot.hasMore).toBe(false)
  307. expect(eventSeqs(session)).toEqual([...older, ...newer].map(event => event.seq))
  308. })
  309. it('installs a page without interpreting business replacement metadata', async () => {
  310. const { api, session } = makeSession()
  311. api.onHistory = () => histResponse([
  312. ev.compactSummary(80, '窗外范围的摘要', 3, 40),
  313. ev.compactCheckpoint(81, 80, 3, 40),
  314. ev.user(82, '压缩后的新问题'),
  315. ], true)
  316. const errorSpy = vi.spyOn(console, 'error').mockImplementation(() => undefined)
  317. try {
  318. await session.open()
  319. const snapshot = session.getSnapshot()
  320. expect(snapshot.openState).toBe('open')
  321. expect(eventSeqs(session)).toEqual([80, 81, 82])
  322. expect(errorSpy).not.toHaveBeenCalled()
  323. } finally {
  324. errorSpy.mockRestore()
  325. }
  326. })
  327. it('drops a discontinuous older page fail-soft (window unchanged, hasMore cleared)', async () => {
  328. const { api, session } = makeSession()
  329. api.onHistory = payload => payload.beforeSeq === undefined
  330. ? histResponse(plainTurn(10, 1, '新', '页'), true)
  331. : histResponse(plainTurn(0, 0, '断', '层'), true) // tail seq 5, but baseSeq is 10 → hole
  332. const errorSpy = vi.spyOn(console, 'error').mockImplementation(() => undefined)
  333. try {
  334. await session.open()
  335. const windowBefore = session.eventSource.getSnapshot()
  336. await session.loadOlder()
  337. const snapshot = session.getSnapshot()
  338. expect(session.eventSource.getSnapshot().entries).toEqual(windowBefore.entries)
  339. expect(snapshot.hasMore).toBe(false)
  340. } finally {
  341. errorSpy.mockRestore()
  342. }
  343. })
  344. it('loadThrough pages repeatedly until the window covers the target seq', async () => {
  345. const oldest = plainTurn(0, 0, '最旧问', '最旧答')
  346. const middle = plainTurn(6, 1, '中问', '中答')
  347. const newest = plainTurn(12, 2, '新问', '新答')
  348. const { api, session } = makeSession()
  349. api.onHistory = (payload) => {
  350. if (payload.beforeSeq === undefined) return histResponse(newest, true)
  351. return payload.beforeSeq === 12 ? histResponse(middle, true) : histResponse(oldest, false)
  352. }
  353. await session.open()
  354. const gate = deferred<Awaited<ReturnType<FakeApiClient['onHistory']>>>()
  355. api.onHistory = (payload) => {
  356. api.onHistory = payload2 => payload2.beforeSeq === 12 ? histResponse(middle, true) : histResponse(oldest, false)
  357. void payload
  358. return gate.promise
  359. }
  360. const jump = session.loadThrough(0)
  361. expect(session.getSnapshot().loadingOlder).toBe(true)
  362. gate.resolve(ok(historyValue(middle, true)))
  363. await jump
  364. const snapshot = session.getSnapshot()
  365. expect(snapshot.loadingOlder).toBe(false)
  366. expect(eventSeqs(session)).toEqual([...oldest, ...middle, ...newest].map(event => event.seq))
  367. expect(api.callsOf('session.history')).toMatchObject([
  368. { beforeSeq: 12, maxMessages: JUMP_PAGE_MESSAGES },
  369. { beforeSeq: 6, maxMessages: JUMP_PAGE_MESSAGES },
  370. ])
  371. })
  372. it('loadThrough is a no-op when the window already covers the target or the session is not open', async () => {
  373. const { api, session } = makeSession()
  374. await session.loadThrough(0) // cold: no-op
  375. expect(api.calls).toEqual([])
  376. api.onHistory = () => histResponse(plainTurn(6, 1, 'x', 'y'), true)
  377. await session.open()
  378. const calls = api.calls.length
  379. await session.loadThrough(6) // baseSeq is already 6
  380. await session.loadThrough(9) // inside the window
  381. expect(api.calls.length).toBe(calls)
  382. })
  383. it('loadThrough retargets a running jump to the lowest requested seq and shares its completion', async () => {
  384. const oldest = plainTurn(0, 0, 'a', 'b')
  385. const middle = plainTurn(6, 1, 'c', 'd')
  386. const { api, session } = makeSession()
  387. api.onHistory = () => histResponse(plainTurn(12, 2, 'e', 'f'), true)
  388. await session.open()
  389. const gate = deferred<Awaited<ReturnType<FakeApiClient['onHistory']>>>()
  390. api.onHistory = () => {
  391. api.onHistory = () => histResponse(oldest, false)
  392. return gate.promise
  393. }
  394. const first = session.loadThrough(6)
  395. const second = session.loadThrough(0)
  396. gate.resolve(ok(historyValue(middle, true)))
  397. await Promise.all([first, second])
  398. expect(eventSeqs(session)).toEqual([...oldest, ...middle].map(event => event.seq).concat([12, 13, 14, 15, 16, 17]))
  399. expect(api.callsOf('session.history')).toHaveLength(2)
  400. })
  401. it('loadThrough refused by a busy pager leaves no target behind for later jumps', async () => {
  402. const middle = plainTurn(6, 1, 'c', 'd')
  403. const { api, session } = makeSession()
  404. api.onHistory = () => histResponse(plainTurn(12, 2, 'e', 'f'), true)
  405. await session.open()
  406. // A plain single-page pull holds the busy flag while the jump is refused.
  407. const gate = deferred<Awaited<ReturnType<FakeApiClient['onHistory']>>>()
  408. api.onHistory = () => gate.promise
  409. const older = session.loadOlder()
  410. await session.loadThrough(0) // refused: must not park seq 0 anywhere
  411. gate.resolve(ok(historyValue(middle, true)))
  412. await older
  413. // A later jump to a nearer seq pages exactly to it — a leaked 0 target
  414. // would keep pulling three-event pages all the way to the head.
  415. api.onHistory = (payload) => {
  416. const start = ((payload as { beforeSeq?: number }).beforeSeq ?? 0) - 3
  417. return histResponse(
  418. [ev.user(start, `u${String(start)}`), ev.user(start + 1, `u${String(start + 1)}`), ev.user(start + 2, `u${String(start + 2)}`)],
  419. start > 0,
  420. )
  421. }
  422. await session.loadThrough(4)
  423. // Covered at seq 3 (≤ 4) after one page; a leaked 0 target would add a
  424. // third call at beforeSeq 3 and pull the head to 0.
  425. expect(api.callsOf('session.history').map(call => (call as { beforeSeq?: number }).beforeSeq))
  426. .toEqual([12, 6])
  427. expect(eventSeqs(session)[0]).toBe(3)
  428. })
  429. it('loadThrough stops paging when the event stream generation moves mid-loop', async () => {
  430. const { api, session } = makeSession()
  431. api.onHistory = () => histResponse(plainTurn(12, 2, 'x', 'y'), true)
  432. await session.open()
  433. const gate = deferred<Awaited<ReturnType<FakeApiClient['onHistory']>>>()
  434. api.onHistory = () => gate.promise
  435. const jump = session.loadThrough(0)
  436. // The address is rebuilt while the first page is in flight.
  437. api.onHistory = () => histResponse(plainTurn(12, 2, 'x', 'y'), true)
  438. const rebuilt = session.resync()
  439. gate.resolve(ok(historyValue(plainTurn(6, 1, 'c', 'd'), true)))
  440. await jump
  441. await rebuilt
  442. // The stale loop must not page the new generation toward its old target:
  443. // history calls are the gated page and the resync tail only.
  444. expect(api.callsOf('session.history')).toHaveLength(1)
  445. expect(session.getSnapshot().loadingOlder).toBe(false)
  446. })
  447. it('loadThrough stops on a page that makes no progress instead of looping', async () => {
  448. const { api, session } = makeSession()
  449. api.onHistory = payload => payload.beforeSeq === undefined
  450. ? histResponse(plainTurn(12, 2, 'x', 'y'), true)
  451. : histResponse([], true) // empty page still claiming more history
  452. await session.open()
  453. await session.loadThrough(0)
  454. expect(session.getSnapshot().loadingOlder).toBe(false)
  455. expect(api.callsOf('session.history')).toHaveLength(1)
  456. })
  457. it('loadThrough fails soft on a thrown page and clears its busy state', async () => {
  458. const { api, session } = makeSession()
  459. api.onHistory = () => histResponse(plainTurn(12, 2, 'x', 'y'), true)
  460. await session.open()
  461. api.onHistory = () => Promise.reject(new Error('page wire down'))
  462. const errorSpy = vi.spyOn(console, 'error').mockImplementation(() => undefined)
  463. try {
  464. await session.loadThrough(0)
  465. expect(errorSpy).toHaveBeenCalled()
  466. expect(session.getSnapshot().loadingOlder).toBe(false)
  467. } finally {
  468. errorSpy.mockRestore()
  469. }
  470. })
  471. it('ignores loadOlder while one is in flight (single request)', async () => {
  472. const { api, session } = makeSession()
  473. api.onHistory = () => histResponse(plainTurn(6, 1, 'x', 'y'), true)
  474. await session.open()
  475. const gate = deferred<Awaited<ReturnType<FakeApiClient['onHistory']>>>()
  476. api.onHistory = () => gate.promise
  477. const first = session.loadOlder()
  478. const second = session.loadOlder()
  479. gate.resolve(ok({
  480. records: entries(plainTurn(0, 0, 'a', 'b')) as never[],
  481. hasMore: false,
  482. modelSelection: { provider: 'deepseek-official', model: 'deepseek-v4-flash' },
  483. }))
  484. await Promise.all([first, second])
  485. expect(api.callsOf('session.follow')).toHaveLength(1)
  486. expect(api.callsOf('session.history')).toHaveLength(1)
  487. })
  488. })
  489. describe('prompt and cancel errors', () => {
  490. it('routes an addressed child through non-activating history, continuation prompt, and interrupt only', async () => {
  491. const api = new FakeApiClient()
  492. const session = new Session(SID, fakeRemote(api), {
  493. address: { parentSessionId: PARENT, childSessionId: SID, mode: 'continuable' },
  494. parentAvailable: true,
  495. })
  496. await session.open()
  497. const prompted = await session.prompt([{ type: 'text', text: '继续' }], 'queue')
  498. const cancelled = await session.cancel()
  499. expect(prompted).toEqual({ ok: true, value: { accepted: true } })
  500. expect(cancelled).toEqual({ ok: true, value: { accepted: true } })
  501. expect(api.callsOf('session.follow')).toEqual([
  502. {
  503. address: {
  504. kind: 'subagent', parentSessionId: PARENT, childSessionId: SID, mode: 'continuable',
  505. },
  506. maxMessages: 50,
  507. },
  508. ])
  509. expect(api.callsOf('subagent.history')).toEqual([])
  510. expect(api.callsOf('subagents.prompt')).toEqual([
  511. {
  512. requestId: expect.any(String) as unknown as string,
  513. parentSessionId: PARENT, childSessionId: SID,
  514. mode: 'continuable',
  515. content: [{ type: 'text', text: '继续' }],
  516. clientTimeZone: new Intl.DateTimeFormat().resolvedOptions().timeZone,
  517. },
  518. ])
  519. expect(api.callsOf('subagents.interruptByParent')).toEqual([
  520. { childSessionId: SID, parentSessionId: PARENT, mode: 'continuable' },
  521. ])
  522. expect(api.callsOf('session.history')).toEqual([])
  523. expect(api.callsOf('session.prompt')).toEqual([])
  524. expect(api.callsOf('session.cancel')).toEqual([])
  525. // A successful interrupt leaves no stop error behind.
  526. expect(session.getSnapshot().promptError).toBeNull()
  527. expect(session.getSnapshot().subagent).toEqual({
  528. address: { parentSessionId: PARENT, childSessionId: SID, mode: 'continuable' },
  529. parentAvailable: true,
  530. })
  531. })
  532. it('forwards continuation image parts to the subagent prompt Remote unstripped', async () => {
  533. const api = new FakeApiClient()
  534. const session = new Session(SID, fakeRemote(api), {
  535. address: { parentSessionId: PARENT, childSessionId: SID, mode: 'continuable' },
  536. parentAvailable: true,
  537. })
  538. await session.open()
  539. const content = [
  540. { type: 'text' as const, text: '看这张图' },
  541. { type: 'image' as const, mediaType: 'image/png' as const, data: 'aGk=', name: 'shot.png' },
  542. ]
  543. const prompted = await session.prompt(content, 'queue')
  544. expect(prompted).toEqual({ ok: true, value: { accepted: true } })
  545. expect(api.callsOf('subagents.prompt')).toEqual([
  546. {
  547. requestId: expect.any(String) as unknown as string,
  548. parentSessionId: PARENT, childSessionId: SID,
  549. mode: 'continuable',
  550. content,
  551. clientTimeZone: new Intl.DateTimeFormat().resolvedOptions().timeZone,
  552. },
  553. ])
  554. expect(session.getSnapshot().promptError).toBeNull()
  555. })
  556. it('lands an interrupt business failure in promptError with op=stop', async () => {
  557. const api = new FakeApiClient()
  558. api.onSubagentInterrupt = () => Promise.resolve(err(new RemoteError('subagent/unauthorized', 'nope', { childSessionId: SID })))
  559. const session = new Session(SID, fakeRemote(api), {
  560. address: { parentSessionId: PARENT, childSessionId: SID, mode: 'continuable' },
  561. parentAvailable: true,
  562. })
  563. await session.open()
  564. const cancelled = await session.cancel()
  565. expect(cancelled).toMatchObject({ ok: false, error: { code: 'subagent/unauthorized' } })
  566. expect(session.getSnapshot().promptError).toMatchObject({
  567. op: 'stop', error: { code: 'subagent/unauthorized' },
  568. })
  569. })
  570. it('rejects staged files instead of dropping them from subagent continuations', async () => {
  571. const api = new FakeApiClient()
  572. const session = new Session(SID, fakeRemote(api), {
  573. address: { parentSessionId: PARENT, childSessionId: SID, mode: 'continuable' },
  574. parentAvailable: true,
  575. })
  576. await session.open()
  577. const prompted = await session.prompt([
  578. { type: 'file', receiptId: 'receipt' as never },
  579. { type: 'text', text: '继续' },
  580. ], 'queue')
  581. expect(prompted).toMatchObject({
  582. ok: false,
  583. error: {
  584. code: 'subagent/attachment-invalid',
  585. details: { reason: 'SUBAGENT_FILE_UNSUPPORTED' },
  586. },
  587. })
  588. expect(api.callsOf('subagents.prompt')).toEqual([])
  589. })
  590. it('sends a one-shot address to the Host under the continuable marker', async () => {
  591. const api = new FakeApiClient()
  592. api.onSubagentPrompt = () => Promise.resolve(err(new RemoteError(
  593. 'subagent/not-resumable', 'subagent cannot be resumed', { childSessionId: SID },
  594. )))
  595. const session = new Session(SID, fakeRemote(api), {
  596. address: { parentSessionId: PARENT, childSessionId: SID, mode: 'one-shot' },
  597. })
  598. await session.open()
  599. const prompted = await session.prompt([{ type: 'text', text: '继续' }], 'queue')
  600. const cancelled = await session.cancel()
  601. // The Host reads the durable descriptor; the wire marker stays 'continuable'.
  602. expect(prompted).toMatchObject({ ok: false, error: { code: 'subagent/not-resumable' } })
  603. expect(cancelled).toEqual({ ok: true, value: { accepted: true } })
  604. expect(api.callsOf('subagents.prompt')).toMatchObject([
  605. { parentSessionId: PARENT, childSessionId: SID, mode: 'continuable' },
  606. ])
  607. expect(api.callsOf('subagents.interruptByParent')).toEqual([
  608. { childSessionId: SID, parentSessionId: PARENT, mode: 'continuable' },
  609. ])
  610. expect(api.callsOf('session.follow')).toEqual([
  611. {
  612. address: {
  613. kind: 'subagent', parentSessionId: PARENT, childSessionId: SID, mode: 'one-shot',
  614. },
  615. maxMessages: 50,
  616. },
  617. ])
  618. expect(api.callsOf('subagent.history')).toEqual([])
  619. expect(api.callsOf('session.cancel')).toEqual([])
  620. })
  621. it('delivers an image continuation to the Host without narrowing its upload parts', async () => {
  622. const api = new FakeApiClient()
  623. const session = new Session(SID, fakeRemote(api), {
  624. address: { parentSessionId: PARENT, childSessionId: SID, mode: 'continuable' },
  625. })
  626. await session.open()
  627. const prompted = await session.prompt(
  628. [{ type: 'text', text: '看图' }, { type: 'image', mediaType: 'image/png', data: 'AA==' }],
  629. 'queue',
  630. )
  631. expect(prompted).toEqual({ ok: true, value: { accepted: true } })
  632. expect(api.callsOf('subagents.prompt')).toMatchObject([
  633. { content: [{ type: 'text' }, { type: 'image', mediaType: 'image/png', data: 'AA==' }] },
  634. ])
  635. })
  636. it('publishes the first-prompt lifecycle synchronously before the Remote settles', async () => {
  637. const { api, session } = makeSession()
  638. session.handleBlank(true)
  639. expect(session.getSnapshot()).toMatchObject({
  640. blank: true, promptAttempted: false, awaitingFirstTurn: false,
  641. })
  642. const inFlight = session.prompt([{ type: 'text', text: '要发的' }], 'queue')
  643. expect(session.getSnapshot()).toMatchObject({
  644. blank: true, promptAttempted: true, awaitingFirstTurn: true,
  645. })
  646. const result = await inFlight
  647. expect(result.ok).toBe(true)
  648. expect(session.getSnapshot()).toMatchObject({
  649. blank: false, promptAttempted: true, awaitingFirstTurn: true,
  650. })
  651. expect(api.callsOf('session.prompt')).toMatchObject([{
  652. sessionId: SID,
  653. mode: 'queue',
  654. content: [{ type: 'text', text: '要发的' }],
  655. clientTimeZone: new Intl.DateTimeFormat().resolvedOptions().timeZone,
  656. }])
  657. session.handleRunning(true)
  658. expect(session.getSnapshot()).toMatchObject({ running: true, awaitingFirstTurn: false })
  659. })
  660. it('keeps the attempted-first-prompt state when the Host rejects the prompt', async () => {
  661. const { api, session } = makeSession()
  662. session.handleBlank(true)
  663. api.onPrompt = () => Promise.resolve(err(new RemoteError('session/agent-busy', 'busy', { reason: 'x' })))
  664. const result = await session.prompt([{ type: 'text', text: '失败的' }], 'queue')
  665. expect(result.ok).toBe(false)
  666. expect(session.getSnapshot().promptError).toMatchObject({ op: 'send', error: { code: 'session/agent-busy' } })
  667. expect(session.getSnapshot()).toMatchObject({
  668. blank: true, promptAttempted: true, awaitingFirstTurn: true,
  669. })
  670. })
  671. it('propagates a non-Remote throw raised while cancelling', async () => {
  672. const { api, session } = makeSession()
  673. api.onCancel = () => Promise.reject(new Error('cancel transport down'))
  674. await expect(session.cancel()).rejects.toThrow('cancel transport down')
  675. expect(session.getSnapshot().promptError).toBeNull()
  676. })
  677. it('reads session-authorized attachment bytes and keeps the opaque id on the wire', async () => {
  678. const { api, session } = makeSession()
  679. const result = await session.readAttachment('attachment-1' as never)
  680. expect(result).toEqual({
  681. ok: true,
  682. value: {
  683. attachment: { attachmentId: 'a', mediaType: 'image/png', bytes: 1, width: 1, height: 1 },
  684. data: Uint8Array.of(0),
  685. },
  686. })
  687. expect(api.callsOf('session.attachment')).toEqual([{
  688. sessionId: SID, attachmentId: 'attachment-1',
  689. }])
  690. })
  691. })
  692. describe('rename', () => {
  693. it('settles the title projection cell from the unary response (higher-seq-wins vs the push frame)', async () => {
  694. const { api, session } = makeSession()
  695. api.onRename = () => Promise.resolve(ok({ title: '正名', seq: 7 }))
  696. const result = await session.rename(' 正名 ')
  697. expect(result).toMatchObject({ ok: true, value: { title: '正名', seq: 7 } })
  698. expect(api.callsOf('session.rename')).toMatchObject([{ sessionId: SID, title: ' 正名 ' }])
  699. expect(session.projections.faceOf('title').getSnapshot()).toBe('正名')
  700. // A stale lower-seq apply (the push-frame path routes into this same
  701. // store) must not roll the settled value back.
  702. session.projections.apply('title', '旧名', 3)
  703. expect(session.projections.faceOf('title').getSnapshot()).toBe('正名')
  704. })
  705. it('returns the business error untouched and folds a transport throw to internal', async () => {
  706. const { api, session } = makeSession()
  707. api.onRename = () => Promise.resolve(err(new RemoteError('session/title-invalid', 'empty', { sessionId: SID })))
  708. const rejected = await session.rename(' ')
  709. expect(rejected).toMatchObject({ ok: false, error: { code: 'session/title-invalid' } })
  710. expect(session.projections.faceOf('title').getSnapshot()).toBeUndefined()
  711. api.onRename = () => Promise.reject(new Error('rename transport down'))
  712. await expect(session.rename('x')).rejects.toThrow('rename transport down')
  713. })
  714. })
  715. describe('remaining branches', () => {
  716. it('propagates a non-Remote throw raised while prompting', async () => {
  717. const { api, session } = makeSession()
  718. api.onPrompt = () => Promise.reject(new Error('prompt wire down'))
  719. await expect(session.prompt([{ type: 'text', text: 'x' }], 'queue')).rejects.toThrow('prompt wire down')
  720. expect(session.getSnapshot().promptError).toBeNull()
  721. })
  722. it('cancel business error also lands op=stop promptError', async () => {
  723. const { api, session } = makeSession()
  724. api.onCancel = () => Promise.resolve(err(new RemoteError('session/agent-busy', 'nope', { reason: 'r' })))
  725. await session.cancel()
  726. expect(session.getSnapshot().promptError).toMatchObject({ op: 'stop', error: { code: 'session/agent-busy' } })
  727. })
  728. it('loadOlder guards: not-open/no-hasMore no-op, err result kept window, empty page updates hasMore, throw fail-soft', async () => {
  729. const { api, session } = makeSession()
  730. await session.loadOlder() // cold: no-op, zero calls
  731. expect(api.calls).toEqual([])
  732. api.onHistory = () => histResponse(plainTurn(6, 1, 'x', 'y'), true)
  733. await session.open()
  734. // err result: window unchanged
  735. api.onHistory = () => Promise.resolve(err(new RemoteError('gateway/internal', 'x', {})))
  736. await session.loadOlder()
  737. expect(eventSeqs(session)).toHaveLength(6)
  738. expect(session.getSnapshot().hasMore).toBe(true)
  739. // empty page: hasMore adopts the response
  740. api.onHistory = () => histResponse([], false)
  741. await session.loadOlder()
  742. expect(session.getSnapshot().hasMore).toBe(false)
  743. // hasMore false now: further loadOlder is a guard no-op
  744. const calls = api.calls.length
  745. await session.loadOlder()
  746. expect(api.calls.length).toBe(calls)
  747. // throw path: fail-soft with console.error
  748. const errorSpy = vi.spyOn(console, 'error').mockImplementation(() => undefined)
  749. try {
  750. await session.resync()
  751. api.onHistory = () => histResponse(plainTurn(6, 1, 'x', 'y'), true)
  752. await session.resync()
  753. api.onHistory = () => Promise.reject(new Error('page wire down'))
  754. await session.loadOlder()
  755. expect(errorSpy).toHaveBeenCalled()
  756. expect(session.getSnapshot().loadingOlder).toBe(false)
  757. } finally {
  758. errorSpy.mockRestore()
  759. }
  760. })
  761. it('subscribe delivers snapshot-change notifications and unsubscribes', async () => {
  762. const { api, session } = makeSession()
  763. api.onHistory = () => histResponse(plainTurn(0, 0, 'a', 'b'))
  764. let notified = 0
  765. const unsubscribe = session.subscribe(() => { notified++ })
  766. await session.open()
  767. await new Promise(resolve => setTimeout(resolve, 0))
  768. expect(notified).toBeGreaterThan(0)
  769. const seen = notified
  770. unsubscribe()
  771. session.handleRunning(true) // any snapshot mutation; the listener must stay silent
  772. await new Promise(resolve => setTimeout(resolve, 0))
  773. expect(notified).toBe(seen)
  774. })
  775. it('rejects an opening page that does not end at the opening cursor', async () => {
  776. const { api, session } = makeSession()
  777. let call = 0
  778. api.onHistory = () => {
  779. call++
  780. return histResponse(plainTurn(0, 0, 'a', 'b'))
  781. }
  782. api.followCursor = 11
  783. await session.open()
  784. expect(call).toBe(1)
  785. const snapshot = session.getSnapshot()
  786. expect(snapshot.openState).toBe('error')
  787. expect(snapshot.openError).toMatchObject({
  788. code: 'gateway/internal', message: 'session event stream page did not end at its requested cursor',
  789. })
  790. expect(eventSeqs(session)).toEqual([])
  791. })
  792. it('deduplicates repeated running flips and records removal', () => {
  793. const { session } = makeSession()
  794. const before = session.getSnapshot()
  795. session.handleRunning(false) // already false: dedup branch
  796. expect(session.getSnapshot()).toBe(before)
  797. session.handleRemoved()
  798. expect(session.getSnapshot().removed).toBe(true)
  799. })
  800. it('drops live events while cold/error (no window upkeep)', async () => {
  801. const { api, session } = makeSession()
  802. await follow(api, ev.user(0, '冷态帧'))
  803. expect(eventSeqs(session)).toEqual([])
  804. api.onHistory = () => Promise.resolve(err(new RemoteError('gateway/internal', 'x', {})))
  805. await session.open()
  806. await follow(api, ev.user(0, '错态帧'))
  807. expect(eventSeqs(session)).toEqual([])
  808. })
  809. it('preserves a Host-reported failure that terminates the live source', async () => {
  810. const { api, session } = makeSession()
  811. api.onHistory = () => histResponse(plainTurn(0, 0, 'a', 'b'))
  812. await session.open()
  813. const failure = new RemoteError('session/not-found', 'session disappeared', { sessionId: SID })
  814. api.failStreams(failure)
  815. await vi.waitFor(() => { expect(session.getSnapshot().openState).toBe('error') })
  816. expect(session.getSnapshot().openError).toMatchObject({
  817. code: failure.code, message: failure.message, details: failure.details,
  818. })
  819. })
  820. it('coalesces queued gap frames behind one repair and exposes a failed repair', async () => {
  821. const { api, session } = makeSession()
  822. api.onHistory = () => histResponse(plainTurn(0, 0, 'a', 'b'))
  823. await session.open()
  824. const gate = deferred<Awaited<ReturnType<FakeApiClient['onHistory']>>>()
  825. let repairs = 0
  826. api.onHistory = () => {
  827. repairs++
  828. return gate.promise
  829. }
  830. const deliveries = Promise.all([
  831. follow(api, ev.user(9, '洞一')),
  832. follow(api, ev.user(10, '洞二')),
  833. ])
  834. await vi.waitFor(() => { expect(repairs).toBe(1) })
  835. gate.reject(new RemoteError('gateway/internal', 'repair wire down', {}))
  836. await deliveries
  837. await vi.waitFor(() => { expect(session.getSnapshot().openState).toBe('error') })
  838. expect(session.getSnapshot().openError).toMatchObject({ code: 'gateway/internal', message: 'repair wire down' })
  839. expect(eventSeqs(session)).toHaveLength(6)
  840. })
  841. it('doOpen transport throw of a stale generation is swallowed (generation guard in catch)', async () => {
  842. const { api, session } = makeSession()
  843. const stale = deferred<Awaited<ReturnType<FakeApiClient['onHistory']>>>()
  844. api.onHistory = () => stale.promise
  845. const opening = session.open()
  846. api.onHistory = () => histResponse(plainTurn(0, 0, 'a', 'b'))
  847. const resynced = session.resync()
  848. stale.reject(new Error('stale wire'))
  849. await Promise.all([opening, resynced])
  850. expect(session.getSnapshot().openState).toBe('open') // stale catch did not write error
  851. })
  852. it('drops a stale doOpen whose history resolved successfully after resync superseded it', async () => {
  853. const { api, session } = makeSession()
  854. const stale = deferred<Awaited<ReturnType<FakeApiClient['onHistory']>>>()
  855. api.onHistory = () => stale.promise
  856. const opening = session.open()
  857. api.onHistory = () => histResponse(plainTurn(6, 1, '新', '代'))
  858. const resynced = session.resync()
  859. stale.resolve(ok({
  860. records: entries(plainTurn(0, 0, '旧', '代')) as never[],
  861. hasMore: false,
  862. modelSelection: { provider: 'deepseek-official', model: 'stale' },
  863. })) // success, but its generation is gone
  864. await Promise.all([opening, resynced])
  865. expect(eventSeqs(session)).toEqual(plainTurn(6, 1, '新', '代').map(event => event.seq))
  866. })
  867. it('drops a gap repair superseded by a full resync while its pull was in flight', async () => {
  868. const { api, session } = makeSession()
  869. api.onHistory = () => histResponse(plainTurn(0, 0, 'a', 'b'))
  870. await session.open()
  871. const repairPull = deferred<Awaited<ReturnType<FakeApiClient['onHistory']>>>()
  872. api.onHistory = () => repairPull.promise
  873. const delivery = follow(api, ev.user(9, '洞'))
  874. await vi.waitFor(() => { expect(api.callsOf('session.history')).toHaveLength(1) })
  875. api.onHistory = () => histResponse(plainTurn(6, 1, 'c', 'd'))
  876. const resynced = session.resync() // bumps the generation
  877. repairPull.resolve(ok({
  878. records: entries(plainTurn(0, 0, '旧', '页')) as never[],
  879. hasMore: false,
  880. modelSelection: { provider: 'deepseek-official', model: 'stale' },
  881. })) // repair result: stale, dropped
  882. await Promise.all([delivery, resynced])
  883. expect(eventSeqs(session)).toEqual(plainTurn(6, 1, 'c', 'd').map(event => event.seq))
  884. })
  885. it('successful cancel leaves no promptError', async () => {
  886. const { api, session } = makeSession()
  887. api.onHistory = () => histResponse(plainTurn(0, 0, 'a', 'b'))
  888. await session.open()
  889. const result = await session.cancel()
  890. expect(result.ok).toBe(true)
  891. expect(session.getSnapshot().promptError).toBeNull()
  892. })
  893. it('dispose is a reserved no-op on resident instances', async () => {
  894. const { session } = makeSession()
  895. await expect(session.dispose()).resolves.toBeUndefined()
  896. })
  897. it('carries raw history and follow events through the event feed', async () => {
  898. const { api, session } = makeSession()
  899. const historyCall = ev.toolCall(6, 1, 'h1', 'bash', '{"cmd":"pwd"}')
  900. const historyResult = ev.toolResult(7, 1, 'h1', 'done')
  901. api.onHistory = () => Promise.resolve(ok({
  902. records: [
  903. ...entries(plainTurn(0, 0, 'a', 'b')),
  904. { type: 'event', event: historyCall },
  905. { type: 'event', event: historyResult },
  906. ] as never[],
  907. hasMore: false,
  908. modelSelection: { provider: 'deepseek-official', model: 'deepseek-v4-flash' },
  909. }))
  910. await session.open()
  911. expect(windowEntries(session).slice(-2)).toEqual([
  912. { type: 'event', event: historyCall },
  913. { type: 'event', event: historyResult },
  914. ])
  915. const liveCall = ev.toolCall(8, 2, 'l1', 'write', '{"file_path":"a.ts"}')
  916. await follow(api, liveCall)
  917. expect(windowEntries(session).at(-1)).toEqual({ type: 'event', event: liveCall })
  918. const liveResult = ev.toolResult(9, 2, 'l1', 'ok')
  919. await follow(api, liveResult)
  920. expect(windowEntries(session).at(-1)).toEqual({ type: 'event', event: liveResult })
  921. })
  922. })
  923. describe('resync', () => {
  924. it('keeps the old feed until the reconnect snapshot, then repairs queued live gaps', async () => {
  925. const { api, session } = makeSession()
  926. api.onHistory = () => histResponse(plainTurn(0, 0, '旧', '窗'))
  927. await session.open()
  928. const oldWindow = session.eventSource.getSnapshot()
  929. const replacement = deferred<Awaited<ReturnType<FakeApiClient['onHistory']>>>()
  930. api.followCursor = 15
  931. api.onHistory = () => replacement.promise
  932. const publications: ReturnType<Session['eventSource']['getSnapshot']>[] = []
  933. const off = session.eventSource.subscribe(() => {
  934. publications.push(session.eventSource.getSnapshot())
  935. })
  936. const syncing = session.resync()
  937. await vi.waitFor(() => { expect(api.callsOf('session.follow')).toHaveLength(2) })
  938. expect(session.eventSource.getSnapshot()).toBe(oldWindow)
  939. expect(publications).toEqual([])
  940. api.onHistory = () => histResponse([
  941. ...plainTurn(10, 2, '终', '页'),
  942. ev.user(16, '后到低位'),
  943. ev.user(17, '后到高位'),
  944. ])
  945. const liveDeliveries = Promise.all([
  946. follow(api, ev.user(17, '后到高位')),
  947. follow(api, ev.user(16, '后到低位')),
  948. ])
  949. expect(session.eventSource.getSnapshot()).toBe(oldWindow)
  950. replacement.resolve(ok({
  951. records: entries(plainTurn(10, 2, '终', '页')) as never[],
  952. hasMore: false,
  953. modelSelection: { provider: 'deepseek-official', model: 'deepseek-v4-flash' },
  954. }))
  955. await Promise.all([syncing, liveDeliveries])
  956. await vi.waitFor(() => {
  957. expect(eventSeqs(session)).toEqual([10, 11, 12, 13, 14, 15, 16, 17])
  958. })
  959. expect(publications).toHaveLength(2)
  960. expect(publications.map(snapshot => snapshot.change.kind)).toEqual(['replace', 'replace'])
  961. expect(publications[0]?.entries.map(entry => entry.event.seq)).toEqual([10, 11, 12, 13, 14, 15])
  962. expect(publications[1]?.entries.map(entry => entry.event.seq)).toEqual([10, 11, 12, 13, 14, 15, 16, 17])
  963. off()
  964. })
  965. it('rebuilds the window without clearing control state; cold instances no-op', async () => {
  966. const { api, session } = makeSession()
  967. api.onHistory = () => histResponse(plainTurn(0, 0, 'a', 'b'))
  968. await session.open()
  969. session.handleRunning(true)
  970. session.handleAgentError('still visible')
  971. api.onHistory = () => histResponse([...plainTurn(0, 0, 'a', 'b'), ...plainTurn(6, 1, 'c', 'd')])
  972. await session.resync()
  973. const snapshot = session.getSnapshot()
  974. expect(snapshot.openState).toBe('open')
  975. expect(snapshot.running).toBe(true)
  976. expect(snapshot.lastAgentError).toBe('still visible')
  977. expect(eventSeqs(session)).toHaveLength(12)
  978. const cold = makeSession()
  979. await cold.session.resync()
  980. expect(cold.api.calls).toEqual([]) // never opened: no traffic
  981. })
  982. it('drops a stale in-flight open superseded by resync (generation guard)', async () => {
  983. const { api, session } = makeSession()
  984. const stale = deferred<Awaited<ReturnType<FakeApiClient['onHistory']>>>()
  985. api.onHistory = () => stale.promise
  986. const firstOpen = session.open()
  987. api.onHistory = () => histResponse(plainTurn(6, 1, '新', '代'))
  988. const resynced = session.resync()
  989. stale.reject(new Error('dead connection')) // the doomed pre-disconnect request fails late
  990. await firstOpen
  991. await resynced
  992. const snapshot = session.getSnapshot()
  993. expect(snapshot.openState).toBe('open') // stale failure did not settle the fresh generation into error
  994. expect(eventSeqs(session)).toEqual(plainTurn(6, 1, '新', '代').map(event => event.seq))
  995. })
  996. })
  997. describe('snapshot ownership', () => {
  998. it('publishes event-window appends without changing an unrelated Session snapshot', async () => {
  999. const { api, session } = makeSession()
  1000. api.onHistory = () => histResponse(plainTurn(0, 0, '稳', '定'))
  1001. await session.open()
  1002. const sessionBefore = session.getSnapshot()
  1003. const windowBefore = session.eventSource.getSnapshot()
  1004. const firstEntry = windowBefore.entries[0]
  1005. await follow(api, ev.user(6, '追加'))
  1006. const windowAfter = session.eventSource.getSnapshot()
  1007. expect(session.getSnapshot()).toBe(sessionBefore)
  1008. expect(windowAfter).not.toBe(windowBefore)
  1009. expect(windowAfter.entries[0]).toBe(firstEntry)
  1010. expect(windowAfter.change).toMatchObject({ kind: 'append' })
  1011. })
  1012. })