transport.host.spec.ts 24 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525
  1. import { Context } from '@deepseek-ai/cordis'
  2. import { createScope } from '@deepseek-ai/dsh-scope'
  3. import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
  4. import type { Session, SessionEvent, SessionHeader } from '@deepseek-ai/dsh-session'
  5. import { snapshotSubagentDescriptor } from '@deepseek-ai/dsh-subagent'
  6. import { describe, expect, it, vi } from 'vitest'
  7. import { SessionHistoryController } from '../src/history.ts'
  8. const signal = (): AbortSignal => new AbortController().signal
  9. function append(
  10. session: Session,
  11. type: string,
  12. data: unknown,
  13. options?: { readonly surfaceOp?: unknown; readonly sourceEventSeqs?: readonly number[] },
  14. ): SessionEvent {
  15. return (session.append as unknown as (
  16. eventType: string,
  17. eventData: unknown,
  18. eventOptions?: unknown,
  19. ) => SessionEvent)(type, data, options)
  20. }
  21. function event(type: string, seq: number, data: unknown = {}): SessionEvent {
  22. return { type, seq, time: seq + 1, data } as SessionEvent
  23. }
  24. function cold(
  25. ctx: Context,
  26. header: SessionHeader,
  27. events: readonly SessionEvent[],
  28. ): void {
  29. ctx.provide('sessionPersistence', {
  30. list: () => Promise.resolve([header]),
  31. inspect: () => Promise.resolve({ meta: header, events }),
  32. } as never)
  33. }
  34. interface Deferred<T> {
  35. readonly promise: Promise<T>
  36. resolve(value: T): void
  37. }
  38. function deferred<T>(): Deferred<T> {
  39. let resolve!: (value: T) => void
  40. const promise = new Promise<T>((settle) => { resolve = settle })
  41. return { promise, resolve }
  42. }
  43. async function setup(): Promise<{ ctx: Context; transport: SessionHistoryController }> {
  44. const ctx = new Context()
  45. await ctx.plugin(SessionStore)
  46. const transport = new SessionHistoryController(ctx)
  47. return { ctx, transport }
  48. }
  49. describe('SessionHistoryController', () => {
  50. it('opens at the current cursor and follows later events from an ordinary Session', async () => {
  51. const { ctx, transport } = await setup()
  52. const session = ctx.sessions.create(SessionId('ordinary'), { meta: { cwd: '/workspace' } })
  53. session.append('turn/start', { turn: 1 })
  54. const abort = new AbortController()
  55. const iterator = transport.follow(
  56. { address: { kind: 'session', sessionId: session.id } },
  57. abort.signal,
  58. )[Symbol.asyncIterator]()
  59. expect(await iterator.next()).toMatchObject({ done: false, value: { type: 'opened', cursor: 0 } })
  60. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  61. expect(await iterator.next()).toMatchObject({
  62. done: false,
  63. value: { type: 'event', event: { type: 'turn/end', seq: 1 } },
  64. })
  65. const page = await transport.page(
  66. { address: { kind: 'session', sessionId: session.id }, throughSeq: 1 },
  67. new AbortController().signal,
  68. )
  69. expect(page.events.map(entry => entry.event.seq)).toEqual([0, 1])
  70. abort.abort()
  71. expect(await iterator.next()).toMatchObject({ done: true })
  72. })
  73. it('ends active followers when the owning Controller unloads', async () => {
  74. const ctx = new Context()
  75. await ctx.plugin(SessionStore)
  76. let transport!: SessionHistoryController
  77. const owner = ctx.plugin(Object.assign(
  78. (inner: Context) => { transport = new SessionHistoryController(inner) },
  79. { inject: ['sessions'] },
  80. ))
  81. await owner.await()
  82. const session = ctx.sessions.create(SessionId('controller-unload'), { meta: { cwd: '/workspace' } })
  83. const iterator = transport.follow(
  84. { address: { kind: 'session', sessionId: session.id } },
  85. new AbortController().signal,
  86. )[Symbol.asyncIterator]()
  87. await expect(iterator.next()).resolves.toEqual({
  88. done: false,
  89. value: { type: 'opened', cursor: -1 },
  90. })
  91. const pending = iterator.next()
  92. await owner.dispose()
  93. await expect(pending).resolves.toEqual({ done: true, value: undefined })
  94. await ctx.fiber.dispose()
  95. })
  96. it('resumes from the last applied seq before delivering later live events', async () => {
  97. const { ctx, transport } = await setup()
  98. const session = ctx.sessions.create(SessionId('resume'), { meta: { cwd: '/workspace' } })
  99. session.append('turn/start', { turn: 1 })
  100. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  101. session.append('turn/start', { turn: 2 })
  102. const abort = new AbortController()
  103. const iterator = transport.follow({
  104. address: { kind: 'session', sessionId: session.id },
  105. afterSeq: 0,
  106. }, abort.signal)[Symbol.asyncIterator]()
  107. expect(await iterator.next()).toEqual({ done: false, value: { type: 'opened', cursor: 2 } })
  108. expect(await iterator.next()).toMatchObject({ done: false, value: { type: 'event', event: { seq: 1 } } })
  109. expect(await iterator.next()).toMatchObject({ done: false, value: { type: 'event', event: { seq: 2 } } })
  110. session.append('turn/end', { turn: 2, reason: { kind: 'completed' } })
  111. expect(await iterator.next()).toMatchObject({ done: false, value: { type: 'event', event: { seq: 3 } } })
  112. abort.abort()
  113. expect(await iterator.next()).toMatchObject({ done: true })
  114. })
  115. it('subscribes before a cold read and ignores unrelated and replayed buffered events', async () => {
  116. const { ctx, transport } = await setup()
  117. const sessionId = SessionId('cold-race')
  118. const header = { version: 0, id: sessionId, createdAt: 1, cwd: '/workspace' }
  119. const listed = deferred<readonly SessionHeader[]>()
  120. ctx.provide('sessionPersistence', {
  121. list: () => listed.promise,
  122. inspect: () => Promise.resolve({ meta: header, events: [event('fixture/start', 0)] }),
  123. } as never)
  124. const abort = new AbortController()
  125. const iterator = transport.follow({ address: { kind: 'session', sessionId } }, abort.signal)
  126. [Symbol.asyncIterator]()
  127. const opening = iterator.next()
  128. ctx.emit('session/event', { id: SessionId('unrelated') } as Session, event('fixture/other', 0))
  129. ctx.emit('session/event', { id: sessionId } as Session, event('fixture/start', 0))
  130. listed.resolve([header])
  131. await expect(opening).resolves.toEqual({ done: false, value: { type: 'opened', cursor: 0 } })
  132. const waiting = iterator.next()
  133. abort.abort()
  134. await expect(waiting).resolves.toMatchObject({ done: true })
  135. })
  136. it('bridges the unpublished end-seed boundary when a cold source attaches', async () => {
  137. const ctx = new Context()
  138. await ctx.plugin(SessionStore)
  139. let transport!: SessionHistoryController
  140. let agentCtx!: Context
  141. await ctx.plugin(Object.assign(
  142. (inner: Context) => { transport = new SessionHistoryController(inner) },
  143. { inject: ['sessions'] },
  144. ))
  145. await ctx.plugin(Object.assign(
  146. (inner: Context) => { agentCtx = createScope(inner, { name: 'agent' }).ctx },
  147. { inject: ['sessions'] },
  148. ))
  149. const sessionId = SessionId('cold-attach')
  150. const header = { version: 0, id: sessionId, createdAt: 1, cwd: '/workspace' }
  151. const seed = [event('fixture/start', 0)]
  152. cold(ctx, header, seed)
  153. agentCtx.on('session/created', (session) => {
  154. if (session.id !== sessionId) return
  155. append(session, 'fixture/setup-one', {})
  156. append(session, 'fixture/setup-two', {})
  157. })
  158. const abort = new AbortController()
  159. const iterator = transport.follow({ address: { kind: 'session', sessionId } }, abort.signal)
  160. [Symbol.asyncIterator]()
  161. await expect(iterator.next()).resolves.toEqual({ done: false, value: { type: 'opened', cursor: 0 } })
  162. agentCtx.sessions.create(SessionId('unrelated-created'), { meta: { cwd: '/workspace' } })
  163. const attached = agentCtx.sessions.prepare(sessionId, { meta: header, seed })
  164. agentCtx.sessions.enter(attached)
  165. agentCtx.sessions.announce(attached)
  166. await expect(iterator.next()).resolves.toMatchObject({
  167. done: false,
  168. value: { type: 'event', event: { type: 'session/end-seed', seq: 1 } },
  169. })
  170. await expect(iterator.next()).resolves.toMatchObject({
  171. done: false,
  172. value: { type: 'event', event: { type: 'fixture/setup-one', seq: 2 } },
  173. })
  174. await expect(iterator.next()).resolves.toMatchObject({
  175. done: false,
  176. value: { type: 'event', event: { type: 'fixture/setup-two', seq: 3 } },
  177. })
  178. append(attached, 'fixture/live', {})
  179. await expect(iterator.next()).resolves.toMatchObject({
  180. done: false,
  181. value: { type: 'event', event: { type: 'fixture/live', seq: 4 } },
  182. })
  183. abort.abort()
  184. await expect(iterator.next()).resolves.toMatchObject({ done: true })
  185. })
  186. it('rejects gaps in replayed and live event sequences', async () => {
  187. const replay = await setup()
  188. const replayId = SessionId('replay-gap')
  189. const replayHeader = { version: 0, id: replayId, createdAt: 1, cwd: '/workspace' }
  190. cold(replay.ctx, replayHeader, [event('fixture/start', 0), event('fixture/gap', 2)])
  191. const replayed = replay.transport.follow({
  192. address: { kind: 'session', sessionId: replayId }, afterSeq: -1,
  193. }, signal())[Symbol.asyncIterator]()
  194. await expect(replayed.next()).resolves.toEqual({ done: false, value: { type: 'opened', cursor: 2 } })
  195. await expect(replayed.next()).resolves.toMatchObject({ done: false, value: { event: { seq: 0 } } })
  196. await expect(replayed.next()).rejects.toMatchObject({ failure: { code: 'internal' } })
  197. const live = await setup()
  198. const session = live.ctx.sessions.create(SessionId('live-gap'), { meta: { cwd: '/workspace' } })
  199. append(session, 'fixture/start', {})
  200. live.ctx.provide('agents', { get: () => ({ id: session.id }) } as never)
  201. const followed = live.transport.follow({
  202. address: { kind: 'session', sessionId: session.id },
  203. }, signal())[Symbol.asyncIterator]()
  204. await expect(followed.next()).resolves.toEqual({ done: false, value: { type: 'opened', cursor: 0 } })
  205. live.ctx.emit('session/event', session, event('fixture/gap', 2))
  206. await expect(followed.next()).rejects.toMatchObject({ failure: { code: 'internal' } })
  207. })
  208. it('opens an empty source at cursor -1', async () => {
  209. const { ctx, transport } = await setup()
  210. const session = ctx.sessions.create(SessionId('empty-follow'), { meta: { cwd: '/workspace' } })
  211. const abort = new AbortController()
  212. const iterator = transport.follow({
  213. address: { kind: 'session', sessionId: session.id },
  214. }, abort.signal)[Symbol.asyncIterator]()
  215. await expect(iterator.next()).resolves.toEqual({ done: false, value: { type: 'opened', cursor: -1 } })
  216. await expect(transport.page({
  217. address: { kind: 'session', sessionId: session.id }, throughSeq: -1,
  218. }, signal())).resolves.toMatchObject({ events: [], hasMore: false })
  219. abort.abort()
  220. await expect(iterator.next()).resolves.toMatchObject({ done: true })
  221. })
  222. it('requires the durable parent and mode for a direct subagent address', async () => {
  223. const { ctx, transport } = await setup()
  224. const parentSessionId = SessionId('parent')
  225. const childSessionId = SessionId('child')
  226. ctx.sessions.create(parentSessionId, { meta: { cwd: '/workspace' } })
  227. const child = ctx.sessions.create(childSessionId, {
  228. meta: { cwd: '/workspace', origin: 'subagent', parentSession: parentSessionId },
  229. })
  230. child.append('subagent/descriptor', snapshotSubagentDescriptor({
  231. mode: 'continuable',
  232. provider: 'test',
  233. label: 'child',
  234. }))
  235. const signal = new AbortController().signal
  236. await expect(transport.page({
  237. address: { kind: 'subagent', parentSessionId, childSessionId, mode: 'continuable' },
  238. throughSeq: 0,
  239. }, signal)).resolves.toMatchObject({ events: [{ event: { type: 'subagent/descriptor' } }] })
  240. await expect(transport.page({
  241. address: {
  242. kind: 'subagent',
  243. parentSessionId: SessionId('other-parent'),
  244. childSessionId,
  245. mode: 'continuable',
  246. },
  247. throughSeq: 0,
  248. }, signal)).rejects.toMatchObject({ failure: { code: 'subagent-unauthorized' } })
  249. await expect(transport.page({
  250. address: { kind: 'subagent', parentSessionId, childSessionId, mode: 'one-shot' },
  251. throughSeq: 0,
  252. }, signal)).rejects.toMatchObject({ failure: { code: 'subagent-unauthorized' } })
  253. await expect(transport.page({
  254. address: { kind: 'session', sessionId: childSessionId },
  255. throughSeq: 0,
  256. }, signal)).rejects.toMatchObject({ failure: { code: 'agent-busy' } })
  257. })
  258. it('preserves a cold inspection failure for the Gateway error branch', async () => {
  259. const { ctx, transport } = await setup()
  260. const sessionId = SessionId('corrupt-cold')
  261. const failure = new Error('cold log is corrupt')
  262. const header = { version: 0, id: sessionId, createdAt: 1, cwd: '/workspace' }
  263. ctx.provide('sessionPersistence', {
  264. list: () => Promise.resolve([header]),
  265. inspect: () => Promise.reject(failure),
  266. } as never)
  267. await expect(transport.page({
  268. address: { kind: 'session', sessionId },
  269. throughSeq: -1,
  270. }, new AbortController().signal)).rejects.toBe(failure)
  271. })
  272. it('rejects malformed page and follow cursors at the service boundary', async () => {
  273. const { ctx, transport } = await setup()
  274. const session = ctx.sessions.create(SessionId('validation'), { meta: { cwd: '/workspace' } })
  275. const address = { kind: 'session' as const, sessionId: session.id }
  276. for (const request of [
  277. { address, throughSeq: -2 },
  278. { address, throughSeq: 0.5 },
  279. { address, throughSeq: -1, beforeSeq: -1 },
  280. { address, throughSeq: -1, beforeSeq: 1.5 },
  281. { address, throughSeq: -1, maxMessages: 0 },
  282. { address, throughSeq: -1, maxMessages: 1.5 },
  283. ]) {
  284. await expect(transport.page(request, signal())).rejects.toMatchObject({ failure: { code: 'bad-request' } })
  285. }
  286. await expect(transport.page({ address, throughSeq: 0 }, signal()))
  287. .rejects.toMatchObject({ failure: { code: 'bad-request' } })
  288. const corrupt = await setup()
  289. const corruptId = SessionId('missing-through-seq')
  290. cold(
  291. corrupt.ctx,
  292. { version: 0, id: corruptId, createdAt: 1, cwd: '/workspace' },
  293. [event('fixture/start', 0), event('fixture/gap', 2)],
  294. )
  295. await expect(corrupt.transport.page({
  296. address: { kind: 'session', sessionId: corruptId }, throughSeq: 1,
  297. }, signal())).rejects.toMatchObject({ failure: { code: 'internal' } })
  298. for (const afterSeq of [-2, 0.5]) {
  299. const iterator = transport.follow({ address, afterSeq }, signal())[Symbol.asyncIterator]()
  300. await expect(iterator.next()).rejects.toMatchObject({ failure: { code: 'bad-request' } })
  301. }
  302. const past = transport.follow({ address, afterSeq: 0 }, signal())[Symbol.asyncIterator]()
  303. await expect(past.next()).rejects.toMatchObject({ failure: { code: 'bad-request' } })
  304. })
  305. it('reports missing ordinary and subagent sources without fabricating inspection failures', async () => {
  306. const { ctx, transport } = await setup()
  307. const ordinary = { kind: 'session' as const, sessionId: SessionId('missing') }
  308. await expect(transport.page({ address: ordinary, throughSeq: -1 }, signal()))
  309. .rejects.toMatchObject({ failure: { code: 'internal' } })
  310. ctx.provide('sessionPersistence', {
  311. list: () => Promise.resolve([]),
  312. inspect: () => Promise.reject(new Error('must not inspect')),
  313. } as never)
  314. await expect(transport.page({ address: ordinary, throughSeq: -1 }, signal()))
  315. .rejects.toMatchObject({ failure: { code: 'session-not-found' } })
  316. await expect(transport.page({
  317. address: {
  318. kind: 'subagent',
  319. parentSessionId: SessionId('parent'),
  320. childSessionId: SessionId('missing-child'),
  321. mode: 'continuable',
  322. },
  323. throughSeq: -1,
  324. }, signal())).rejects.toMatchObject({ failure: { code: 'subagent-not-found' } })
  325. })
  326. it('rejects incomplete cold metadata before serving a source', async () => {
  327. const first = await setup()
  328. const sessionId = SessionId('incomplete')
  329. const address = { kind: 'session' as const, sessionId }
  330. first.ctx.provide('sessionPersistence', {
  331. list: () => Promise.resolve([{ version: 0, id: sessionId, createdAt: 1 }]),
  332. inspect: () => Promise.reject(new Error('must not inspect')),
  333. } as never)
  334. await expect(first.transport.page({ address, throughSeq: -1 }, signal()))
  335. .rejects.toMatchObject({ failure: { code: 'session-not-found' } })
  336. const second = await setup()
  337. const listed = { version: 0, id: sessionId, createdAt: 1, cwd: '/workspace' }
  338. second.ctx.provide('sessionPersistence', {
  339. list: () => Promise.resolve([listed]),
  340. inspect: () => Promise.resolve({ meta: { ...listed, cwd: undefined }, events: [] }),
  341. } as never)
  342. await expect(second.transport.page({ address, throughSeq: -1 }, signal()))
  343. .rejects.toMatchObject({ failure: { code: 'session-not-found' } })
  344. })
  345. it('serves cold ordinary history and validates every durable subagent descriptor state', async () => {
  346. const ordinaryBench = await setup()
  347. const ordinaryId = SessionId('cold-ordinary')
  348. const ordinaryHeader = { version: 0, id: ordinaryId, createdAt: 1, cwd: '/workspace' }
  349. cold(ordinaryBench.ctx, ordinaryHeader, [event('turn/start', 0, { turn: 1 })])
  350. await expect(ordinaryBench.transport.page({
  351. address: { kind: 'session', sessionId: ordinaryId },
  352. throughSeq: 0,
  353. }, signal())).resolves.toMatchObject({ events: [{ event: { seq: 0 } }] })
  354. const parentSessionId = SessionId('cold-parent')
  355. const childSessionId = SessionId('cold-child')
  356. const childHeader = {
  357. version: 0,
  358. id: childSessionId,
  359. createdAt: 1,
  360. cwd: '/workspace',
  361. origin: 'subagent' as const,
  362. parentSession: parentSessionId,
  363. }
  364. const childAddress = {
  365. kind: 'subagent' as const,
  366. parentSessionId,
  367. childSessionId,
  368. mode: 'continuable' as const,
  369. }
  370. const missing = await setup()
  371. cold(missing.ctx, childHeader, [])
  372. await expect(missing.transport.page({ address: childAddress, throughSeq: -1 }, signal()))
  373. .rejects.toMatchObject({ failure: { code: 'subagent-catalog-diagnostic', details: { reason: 'unsupported' } } })
  374. const corrupt = await setup()
  375. cold(corrupt.ctx, childHeader, [event('subagent/descriptor', 0, { version: 'bad' })])
  376. await expect(corrupt.transport.page({ address: childAddress, throughSeq: 0 }, signal()))
  377. .rejects.toMatchObject({ failure: { code: 'subagent-catalog-diagnostic', details: { reason: 'corrupt' } } })
  378. const ordinaryChild = await setup()
  379. const { origin: _origin, ...ordinaryChildHeader } = childHeader
  380. cold(ordinaryChild.ctx, ordinaryChildHeader, [])
  381. await expect(ordinaryChild.transport.page({ address: childAddress, throughSeq: -1 }, signal()))
  382. .rejects.toMatchObject({ failure: { code: 'subagent-unauthorized' } })
  383. })
  384. it('uses attached and detached projection cuts and isolates a child projection failure', async () => {
  385. const attached = await setup()
  386. const session = attached.ctx.sessions.create(SessionId('projected'), { meta: { cwd: '/workspace' } })
  387. session.append('turn/start', { turn: 1 })
  388. const snapshot = vi.fn(() => ({ asOfSeq: 0, values: { title: 'attached' } }))
  389. attached.ctx.provide('sessionProjections', { snapshot, restore: vi.fn() } as never)
  390. await expect(attached.transport.page({
  391. address: { kind: 'session', sessionId: session.id },
  392. throughSeq: 0,
  393. }, signal())).resolves.toMatchObject({ projections: { asOfSeq: 0, values: { title: 'attached' } } })
  394. expect(snapshot).toHaveBeenCalledWith(session)
  395. const older = await attached.transport.page({
  396. address: { kind: 'session', sessionId: session.id }, throughSeq: 0, beforeSeq: 1,
  397. }, signal())
  398. expect('projections' in older).toBe(false)
  399. const detached = await setup()
  400. const coldId = SessionId('projected-cold')
  401. const header = { version: 0, id: coldId, createdAt: 1, cwd: '/workspace' }
  402. cold(detached.ctx, header, [event('turn/start', 0, { turn: 1 })])
  403. const restore = vi.fn(() => ({ snapshot: { asOfSeq: 0, values: { title: 'cold' } } }))
  404. detached.ctx.provide('sessionProjections', { snapshot: vi.fn(), restore } as never)
  405. await expect(detached.transport.page({
  406. address: { kind: 'session', sessionId: coldId },
  407. throughSeq: 0,
  408. }, signal())).resolves.toMatchObject({ projections: { values: { title: 'cold' } } })
  409. expect(restore).toHaveBeenCalledWith({}, expect.any(Array), 0)
  410. const failed = await setup()
  411. cold(failed.ctx, header, [event('turn/start', 0, { turn: 1 })])
  412. failed.ctx.provide('sessionProjections', {
  413. snapshot: vi.fn(),
  414. restore: () => { throw new Error('projection failed') },
  415. } as never)
  416. await expect(failed.transport.page({
  417. address: { kind: 'session', sessionId: coldId },
  418. throughSeq: 0,
  419. }, signal())).rejects.toThrow('projection failed')
  420. const child = await setup()
  421. const parentSessionId = SessionId('projection-parent')
  422. const childSessionId = SessionId('projection-child')
  423. const childSession = child.ctx.sessions.create(childSessionId, {
  424. meta: { cwd: '/workspace', origin: 'subagent', parentSession: parentSessionId },
  425. })
  426. childSession.append('subagent/descriptor', snapshotSubagentDescriptor({
  427. mode: 'continuable', provider: 'test', label: 'child',
  428. }))
  429. const warn = vi.spyOn(child.ctx.logger, 'warn').mockImplementation(() => undefined)
  430. child.ctx.provide('sessionProjections', {
  431. snapshot: () => { throw new Error('child projection failed') },
  432. restore: vi.fn(),
  433. } as never)
  434. const page = await child.transport.page({
  435. address: { kind: 'subagent', parentSessionId, childSessionId, mode: 'continuable' },
  436. throughSeq: 0,
  437. }, signal())
  438. expect('projections' in page).toBe(false)
  439. expect(warn).toHaveBeenCalledWith(expect.stringContaining('child projection failed'))
  440. })
  441. it('keeps message-aligned pagination contiguous across replacement provenance', async () => {
  442. const { ctx, transport } = await setup()
  443. const session = ctx.sessions.create(SessionId('pagination'), { meta: { cwd: '/workspace' } })
  444. session.append('turn/start', { turn: 1 })
  445. append(session, 'user/message', { content: [], source: { kind: 'user' } }, { surfaceOp: 'append' })
  446. const firstReply = append(session, 'assistant/message', { turn: 1, step: 1, message: {} }, { surfaceOp: 'append' })
  447. append(session, 'user/message', { content: [], source: { kind: 'user' } }, { surfaceOp: 'append' })
  448. append(session, 'assistant/message', { turn: 1, step: 2, message: {} }, { surfaceOp: 'append' })
  449. const summary = append(session, 'fixture/summary', {})
  450. const replacement = append(session, 'user/message', { content: [], source: { kind: 'plugin' } }, {
  451. surfaceOp: { op: 'replace', start: 1, end: 4 },
  452. sourceEventSeqs: [1, firstReply.seq, 3, 4, summary.seq],
  453. })
  454. const page = await transport.page({
  455. address: { kind: 'session', sessionId: session.id }, throughSeq: replacement.seq, maxMessages: 2,
  456. }, signal())
  457. expect(page.events.map(entry => entry.event.seq)).toEqual([3, 4, 5, replacement.seq])
  458. expect(page.hasMore).toBe(true)
  459. const before = await transport.page({
  460. address: { kind: 'session', sessionId: session.id }, throughSeq: replacement.seq, beforeSeq: 3, maxMessages: 1,
  461. }, signal())
  462. expect(before.events.map(entry => entry.event.seq)).toEqual([2])
  463. })
  464. it('keeps cited source events in the page that owns their appended message', async () => {
  465. const { ctx, transport } = await setup()
  466. const session = ctx.sessions.create(SessionId('pagination-sources'), { meta: { cwd: '/workspace' } })
  467. const source = append(session, 'fixture/source', {})
  468. append(session, 'user/message', { content: [], source: { kind: 'plugin' } }, {
  469. surfaceOp: 'append', sourceEventSeqs: [source.seq],
  470. })
  471. const page = await transport.page({
  472. address: { kind: 'session', sessionId: session.id }, throughSeq: 1, maxMessages: 1,
  473. }, signal())
  474. expect(page.events.map(entry => entry.event.seq)).toEqual([0, 1])
  475. expect(page.hasMore).toBe(false)
  476. })
  477. })