transport.host.spec.ts 28 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644
  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 type { SessionObservation } from '@deepseek-ai/dsh-session-query'
  6. import { snapshotSubagentDescriptor } from '@deepseek-ai/dsh-subagent'
  7. import { subagentIdentityProjectionDefinition } from '@deepseek-ai/dsh-subagent/src/projection.ts'
  8. import { describe, expect, it, vi } from 'vitest'
  9. import { SessionHistoryController } from '../src/history.ts'
  10. import { installSessionReadTestServices, testSessionPersistence } from './test-remote.ts'
  11. const signal = (): AbortSignal => new AbortController().signal
  12. function append(
  13. session: Session,
  14. type: string,
  15. data: unknown,
  16. options?: { readonly surfaceOp?: unknown; readonly sourceEventSeqs?: readonly number[] },
  17. ): SessionEvent {
  18. return (session.append as unknown as (
  19. eventType: string,
  20. eventData: unknown,
  21. eventOptions?: unknown,
  22. ) => SessionEvent)(type, data, options)
  23. }
  24. function event(type: string, seq: number, data: unknown = {}): SessionEvent {
  25. return {
  26. type,
  27. seq,
  28. time: seq + 1,
  29. data,
  30. ...type.startsWith('fixture/') ? { ignorable: true } : {},
  31. } as SessionEvent
  32. }
  33. function cold(
  34. ctx: Context,
  35. header: SessionHeader,
  36. events: readonly SessionEvent[],
  37. ): void {
  38. ctx.provide('sessionPersistence', testSessionPersistence(ctx, {
  39. list: () => Promise.resolve([header]),
  40. inspect: () => Promise.resolve({ meta: header, events }),
  41. }) as never)
  42. }
  43. interface Deferred<T> {
  44. readonly promise: Promise<T>
  45. resolve(value: T): void
  46. }
  47. function deferred<T>(): Deferred<T> {
  48. let resolve!: (value: T) => void
  49. const promise = new Promise<T>((settle) => { resolve = settle })
  50. return { promise, resolve }
  51. }
  52. async function setup(): Promise<{ ctx: Context; transport: SessionHistoryController }> {
  53. const ctx = new Context()
  54. await ctx.plugin(SessionStore)
  55. installSessionReadTestServices(ctx)
  56. ctx.sessionProjections.register(subagentIdentityProjectionDefinition)
  57. const transport = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  58. return { ctx, transport }
  59. }
  60. describe('SessionHistoryController', () => {
  61. it('opens at the current cursor and follows later events from an ordinary Session', async () => {
  62. const { ctx, transport } = await setup()
  63. const session = ctx.sessions.create(SessionId('ordinary'), { meta: { cwd: '/workspace' } })
  64. session.append('turn/start', { turn: 1 })
  65. const abort = new AbortController()
  66. const iterator = transport.follow(
  67. { address: { kind: 'session', sessionId: session.id } },
  68. abort.signal,
  69. )[Symbol.asyncIterator]()
  70. expect(await iterator.next()).toMatchObject({ done: false, value: { type: 'snapshot', cursor: 0 } })
  71. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  72. expect(await iterator.next()).toMatchObject({
  73. done: false,
  74. value: { type: 'event', event: { type: 'turn/end', seq: 1 } },
  75. })
  76. const page = await transport.page(
  77. { address: { kind: 'session', sessionId: session.id }, throughSeq: 1 },
  78. new AbortController().signal,
  79. )
  80. expect(page.events.map(entry => entry.event.seq)).toEqual([0, 1])
  81. abort.abort()
  82. expect(await iterator.next()).toMatchObject({ done: true })
  83. })
  84. it('ends active followers when the owning Controller unloads', async () => {
  85. const ctx = new Context()
  86. await ctx.plugin(SessionStore)
  87. installSessionReadTestServices(ctx)
  88. let transport!: SessionHistoryController
  89. const owner = ctx.plugin(Object.assign(
  90. (inner: Context) => {
  91. transport = new SessionHistoryController(inner, (observation) => { observation[Symbol.dispose]() })
  92. },
  93. { inject: ['sessions', 'sessionQuery'] },
  94. ))
  95. await owner.await()
  96. const session = ctx.sessions.create(SessionId('controller-unload'), { meta: { cwd: '/workspace' } })
  97. const iterator = transport.follow(
  98. { address: { kind: 'session', sessionId: session.id } },
  99. new AbortController().signal,
  100. )[Symbol.asyncIterator]()
  101. await expect(iterator.next()).resolves.toMatchObject({
  102. done: false,
  103. value: { type: 'snapshot', cursor: -1 },
  104. })
  105. const pending = iterator.next()
  106. await owner.dispose()
  107. await expect(pending).resolves.toEqual({ done: true, value: undefined })
  108. await ctx.fiber.dispose()
  109. })
  110. it('reconnects with a complete replacement snapshot before later live events', async () => {
  111. const { ctx, transport } = await setup()
  112. const session = ctx.sessions.create(SessionId('resume'), { meta: { cwd: '/workspace' } })
  113. session.append('turn/start', { turn: 1 })
  114. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  115. session.append('turn/start', { turn: 2 })
  116. const abort = new AbortController()
  117. const iterator = transport.follow({
  118. address: { kind: 'session', sessionId: session.id },
  119. }, abort.signal)[Symbol.asyncIterator]()
  120. expect(await iterator.next()).toMatchObject({
  121. done: false,
  122. value: {
  123. type: 'snapshot',
  124. cursor: 2,
  125. events: [{ event: { seq: 0 } }, { event: { seq: 1 } }, { event: { seq: 2 } }],
  126. },
  127. })
  128. session.append('turn/end', { turn: 2, reason: { kind: 'completed' } })
  129. expect(await iterator.next()).toMatchObject({ done: false, value: { type: 'event', event: { seq: 3 } } })
  130. abort.abort()
  131. expect(await iterator.next()).toMatchObject({ done: true })
  132. })
  133. it('subscribes before a cold read and ignores unrelated and replayed buffered events', async () => {
  134. const { ctx, transport } = await setup()
  135. const sessionId = SessionId('cold-race')
  136. const header = { version: 0, id: sessionId, createdAt: 1, cwd: '/workspace' }
  137. const inspected = deferred<{ meta: SessionHeader; events: readonly SessionEvent[] }>()
  138. ctx.provide('sessionPersistence', testSessionPersistence(ctx, {
  139. inspect: () => inspected.promise,
  140. }) as never)
  141. const abort = new AbortController()
  142. const iterator = transport.follow({ address: { kind: 'session', sessionId } }, abort.signal)
  143. [Symbol.asyncIterator]()
  144. const opening = iterator.next()
  145. ctx.emit('session/event', {
  146. id: SessionId('unrelated'), events: [event('fixture/other', 0)],
  147. } as unknown as Session, event('fixture/other', 0))
  148. ctx.emit('session/event', {
  149. id: sessionId, events: [event('fixture/start', 0)],
  150. } as unknown as Session, event('fixture/start', 0))
  151. inspected.resolve({ meta: header, events: [event('fixture/start', 0)] })
  152. await expect(opening).resolves.toMatchObject({ done: false, value: { type: 'snapshot', cursor: 0 } })
  153. const waiting = iterator.next()
  154. abort.abort()
  155. await expect(waiting).resolves.toMatchObject({ done: true })
  156. })
  157. it('buffers creation while the opening observation is unresolved', async () => {
  158. const ctx = new Context()
  159. await ctx.plugin(SessionStore)
  160. const sessionId = SessionId('created-during-observation')
  161. const header = { version: 0, id: sessionId, createdAt: 1, cwd: '/workspace' }
  162. const observed = deferred<SessionObservation>()
  163. ctx.provide('sessionQuery', { observeSession: () => observed.promise } as never)
  164. const transport = new SessionHistoryController(ctx, vi.fn())
  165. const abort = new AbortController()
  166. const iterator = transport.follow({ address: { kind: 'session', sessionId } }, abort.signal)
  167. [Symbol.asyncIterator]()
  168. const opening = iterator.next()
  169. const attached = ctx.sessions.create(sessionId, { meta: header, seed: [event('fixture/seed', 0)] })
  170. observed.resolve({
  171. source: 'live',
  172. header: attached.header,
  173. events: attached.events,
  174. cursor: attached.seq - 1,
  175. projections: { asOfSeq: attached.seq - 1, values: {} },
  176. retain: vi.fn(),
  177. [Symbol.dispose]: vi.fn(),
  178. } as unknown as SessionObservation)
  179. await expect(opening).resolves.toMatchObject({
  180. done: false,
  181. value: {
  182. type: 'snapshot', cursor: 1, events: [{ event: { seq: 0 } }, { event: { seq: 1 } }],
  183. },
  184. })
  185. expect(attached.id).toBe(sessionId)
  186. abort.abort()
  187. await expect(iterator.next()).resolves.toMatchObject({ done: true })
  188. })
  189. it('bridges the unpublished end-seed boundary when a cold source attaches', async () => {
  190. const ctx = new Context()
  191. await ctx.plugin(SessionStore)
  192. installSessionReadTestServices(ctx)
  193. let transport!: SessionHistoryController
  194. let agentCtx!: Context
  195. await ctx.plugin(Object.assign(
  196. (inner: Context) => {
  197. transport = new SessionHistoryController(inner, (observation) => { observation[Symbol.dispose]() })
  198. },
  199. { inject: ['sessions', 'sessionQuery'] },
  200. ))
  201. await ctx.plugin(Object.assign(
  202. (inner: Context) => { agentCtx = createScope(inner, { name: 'agent' }).ctx },
  203. { inject: ['sessions'] },
  204. ))
  205. const sessionId = SessionId('cold-attach')
  206. const header = { version: 0, id: sessionId, createdAt: 1, cwd: '/workspace' }
  207. const seed = [event('fixture/start', 0)]
  208. cold(ctx, header, seed)
  209. agentCtx.on('session/created', (session) => {
  210. if (session.id !== sessionId) return
  211. append(session, 'fixture/setup-one', {})
  212. append(session, 'fixture/setup-two', {})
  213. })
  214. const abort = new AbortController()
  215. const iterator = transport.follow({ address: { kind: 'session', sessionId } }, abort.signal)
  216. [Symbol.asyncIterator]()
  217. await expect(iterator.next()).resolves.toMatchObject({ done: false, value: { type: 'snapshot', cursor: 0 } })
  218. agentCtx.sessions.create(SessionId('unrelated-created'), { meta: { cwd: '/workspace' } })
  219. const attached = agentCtx.sessions.prepare(sessionId, { meta: header, seed })
  220. agentCtx.sessions.enter(attached)
  221. agentCtx.sessions.announce(attached)
  222. await expect(iterator.next()).resolves.toMatchObject({
  223. done: false,
  224. value: { type: 'event', event: { type: 'session/end-seed', seq: 1 } },
  225. })
  226. await expect(iterator.next()).resolves.toMatchObject({
  227. done: false,
  228. value: { type: 'event', event: { type: 'fixture/setup-one', seq: 2 } },
  229. })
  230. await expect(iterator.next()).resolves.toMatchObject({
  231. done: false,
  232. value: { type: 'event', event: { type: 'fixture/setup-two', seq: 3 } },
  233. })
  234. append(attached, 'fixture/live', {})
  235. await expect(iterator.next()).resolves.toMatchObject({
  236. done: false,
  237. value: { type: 'event', event: { type: 'fixture/live', seq: 4 } },
  238. })
  239. abort.abort()
  240. await expect(iterator.next()).resolves.toMatchObject({ done: true })
  241. })
  242. it('rejects gaps in replayed and live event sequences', async () => {
  243. const replay = await setup()
  244. const replayId = SessionId('replay-gap')
  245. const replayHeader = { version: 0, id: replayId, createdAt: 1, cwd: '/workspace' }
  246. cold(replay.ctx, replayHeader, [event('fixture/start', 0), event('fixture/gap', 2)])
  247. const replayed = replay.transport.follow({
  248. address: { kind: 'session', sessionId: replayId },
  249. }, signal())[Symbol.asyncIterator]()
  250. await expect(replayed.next()).rejects.toMatchObject({ code: 'SESSION_QUERY_CORRUPT_SESSION' })
  251. const live = await setup()
  252. const session = live.ctx.sessions.create(SessionId('live-gap'), { meta: { cwd: '/workspace' } })
  253. append(session, 'fixture/start', {})
  254. live.ctx.provide('agents', { get: () => ({ id: session.id }) } as never)
  255. const followed = live.transport.follow({
  256. address: { kind: 'session', sessionId: session.id },
  257. }, signal())[Symbol.asyncIterator]()
  258. await expect(followed.next()).resolves.toMatchObject({ done: false, value: { type: 'snapshot', cursor: 0 } })
  259. const skipped = event('fixture/skipped', 1)
  260. const gap = event('fixture/gap', 2)
  261. live.ctx.emit('session/event', {
  262. id: session.id,
  263. events: [event('fixture/start', 0), skipped, gap],
  264. } as unknown as Session, gap)
  265. await expect(followed.next()).rejects.toMatchObject({ failure: { code: 'internal' } })
  266. })
  267. it('opens an empty source at cursor -1', async () => {
  268. const { ctx, transport } = await setup()
  269. const session = ctx.sessions.create(SessionId('empty-follow'), { meta: { cwd: '/workspace' } })
  270. const abort = new AbortController()
  271. const iterator = transport.follow({
  272. address: { kind: 'session', sessionId: session.id },
  273. }, abort.signal)[Symbol.asyncIterator]()
  274. await expect(iterator.next()).resolves.toMatchObject({ done: false, value: { type: 'snapshot', cursor: -1 } })
  275. await expect(transport.page({
  276. address: { kind: 'session', sessionId: session.id }, throughSeq: -1,
  277. }, signal())).resolves.toMatchObject({ events: [], hasMore: false })
  278. abort.abort()
  279. await expect(iterator.next()).resolves.toMatchObject({ done: true })
  280. })
  281. it('publishes an empty projection baseline when the query has no registry', async () => {
  282. const ctx = new Context()
  283. await ctx.plugin(SessionStore)
  284. const sessionId = SessionId('projectionless-follow')
  285. const meta = { version: 0, id: sessionId, createdAt: 1, cwd: '/workspace' }
  286. ctx.provide('sessionQuery', {
  287. observeSession: () => Promise.resolve({
  288. source: 'live', header: meta, events: [], cursor: -1,
  289. retain: vi.fn(), [Symbol.dispose]: vi.fn(),
  290. } satisfies SessionObservation),
  291. } as never)
  292. const history = new SessionHistoryController(ctx, vi.fn())
  293. const abort = new AbortController()
  294. const iterator = history.follow({ address: { kind: 'session', sessionId } }, abort.signal)
  295. [Symbol.asyncIterator]()
  296. await expect(iterator.next()).resolves.toMatchObject({
  297. value: { type: 'snapshot', projections: { asOfSeq: -1, values: {} } },
  298. })
  299. abort.abort()
  300. await expect(iterator.next()).resolves.toMatchObject({ done: true })
  301. await ctx.fiber.dispose()
  302. })
  303. it('disposes a retained promotion when background activation rejects synchronously', async () => {
  304. const ctx = new Context()
  305. await ctx.plugin(SessionStore)
  306. const sessionId = SessionId('promotion-failure')
  307. const meta = { version: 0, id: sessionId, createdAt: 1, cwd: '/workspace' }
  308. const disposePromotion = vi.fn()
  309. const promotion = {
  310. source: 'prepared', header: meta, events: [], cursor: -1,
  311. projections: { asOfSeq: -1, values: {} },
  312. retain: vi.fn(), [Symbol.dispose]: disposePromotion,
  313. } as unknown as SessionObservation
  314. const source = {
  315. ...promotion,
  316. retain: () => promotion,
  317. [Symbol.dispose]: vi.fn(),
  318. } as SessionObservation
  319. ctx.provide('sessionQuery', {
  320. observeSession: () => Promise.resolve(source),
  321. } as never)
  322. const history = new SessionHistoryController(ctx, () => { throw new Error('activation failed') })
  323. const iterator = history.follow({ address: { kind: 'session', sessionId } }, signal())
  324. [Symbol.asyncIterator]()
  325. await expect(iterator.next()).resolves.toMatchObject({ value: { type: 'snapshot' } })
  326. await expect(iterator.next()).rejects.toThrow('activation failed')
  327. expect(disposePromotion).toHaveBeenCalledOnce()
  328. await ctx.fiber.dispose()
  329. })
  330. it('requires the durable parent and mode for a direct subagent address', async () => {
  331. const { ctx, transport } = await setup()
  332. const parentSessionId = SessionId('parent')
  333. const childSessionId = SessionId('child')
  334. ctx.sessions.create(parentSessionId, { meta: { cwd: '/workspace' } })
  335. const child = ctx.sessions.create(childSessionId, {
  336. meta: { cwd: '/workspace', origin: 'subagent', parentSession: parentSessionId },
  337. })
  338. child.append('subagent/descriptor', snapshotSubagentDescriptor({
  339. mode: 'continuable',
  340. provider: 'test',
  341. label: 'child',
  342. }))
  343. const signal = new AbortController().signal
  344. await expect(transport.page({
  345. address: { kind: 'subagent', parentSessionId, childSessionId, mode: 'continuable' },
  346. throughSeq: 0,
  347. }, signal)).resolves.toMatchObject({ events: [{ event: { type: 'subagent/descriptor' } }] })
  348. await expect(transport.page({
  349. address: {
  350. kind: 'subagent',
  351. parentSessionId: SessionId('other-parent'),
  352. childSessionId,
  353. mode: 'continuable',
  354. },
  355. throughSeq: 0,
  356. }, signal)).rejects.toMatchObject({ failure: { code: 'subagent-unauthorized' } })
  357. await expect(transport.page({
  358. address: { kind: 'subagent', parentSessionId, childSessionId, mode: 'one-shot' },
  359. throughSeq: 0,
  360. }, signal)).rejects.toMatchObject({ failure: { code: 'subagent-unauthorized' } })
  361. await expect(transport.page({
  362. address: { kind: 'session', sessionId: childSessionId },
  363. throughSeq: 0,
  364. }, signal)).rejects.toMatchObject({ failure: { code: 'agent-busy' } })
  365. })
  366. it('preserves a cold inspection failure for the Gateway error branch', async () => {
  367. const { ctx, transport } = await setup()
  368. const sessionId = SessionId('corrupt-cold')
  369. const failure = new Error('cold log is corrupt')
  370. const header = { version: 0, id: sessionId, createdAt: 1, cwd: '/workspace' }
  371. ctx.provide('sessionPersistence', testSessionPersistence(ctx, {
  372. list: () => Promise.resolve([header]),
  373. inspect: () => Promise.reject(failure),
  374. }) as never)
  375. await expect(transport.page({
  376. address: { kind: 'session', sessionId },
  377. throughSeq: -1,
  378. }, new AbortController().signal)).rejects.toMatchObject({
  379. code: 'SESSION_QUERY_PERSISTENCE_FAILED',
  380. cause: failure,
  381. })
  382. })
  383. it('rejects malformed page and follow cursors at the service boundary', async () => {
  384. const { ctx, transport } = await setup()
  385. const session = ctx.sessions.create(SessionId('validation'), { meta: { cwd: '/workspace' } })
  386. const address = { kind: 'session' as const, sessionId: session.id }
  387. for (const request of [
  388. { address, throughSeq: -2 },
  389. { address, throughSeq: 0.5 },
  390. { address, throughSeq: -1, beforeSeq: -1 },
  391. { address, throughSeq: -1, beforeSeq: 1.5 },
  392. { address, throughSeq: -1, maxMessages: 0 },
  393. { address, throughSeq: -1, maxMessages: 1.5 },
  394. ]) {
  395. await expect(transport.page(request, signal())).rejects.toMatchObject({ failure: { code: 'bad-request' } })
  396. }
  397. await expect(transport.page({ address, throughSeq: 0 }, signal()))
  398. .rejects.toMatchObject({ failure: { code: 'bad-request' } })
  399. const corrupt = await setup()
  400. const corruptId = SessionId('missing-through-seq')
  401. cold(
  402. corrupt.ctx,
  403. { version: 0, id: corruptId, createdAt: 1, cwd: '/workspace' },
  404. [event('fixture/start', 0), event('fixture/gap', 2)],
  405. )
  406. await expect(corrupt.transport.page({
  407. address: { kind: 'session', sessionId: corruptId }, throughSeq: 1,
  408. }, signal())).rejects.toMatchObject({ code: 'SESSION_QUERY_CORRUPT_SESSION' })
  409. for (const maxMessages of [0, 0.5]) {
  410. const iterator = transport.follow({ address, maxMessages }, signal())[Symbol.asyncIterator]()
  411. await expect(iterator.next()).rejects.toMatchObject({ failure: { code: 'bad-request' } })
  412. }
  413. })
  414. it('reports missing ordinary and subagent sources without fabricating inspection failures', async () => {
  415. const { ctx, transport } = await setup()
  416. const ordinary = { kind: 'session' as const, sessionId: SessionId('missing') }
  417. await expect(transport.page({ address: ordinary, throughSeq: -1 }, signal()))
  418. .rejects.toMatchObject({ failure: { code: 'session-not-found' } })
  419. const inspect = vi.fn(() => Promise.resolve(undefined))
  420. ctx.provide('sessionPersistence', testSessionPersistence(ctx, {
  421. list: () => Promise.resolve([]),
  422. inspect,
  423. }) as never)
  424. await expect(transport.page({ address: ordinary, throughSeq: -1 }, signal()))
  425. .rejects.toMatchObject({ failure: { code: 'session-not-found' } })
  426. await expect(transport.page({
  427. address: {
  428. kind: 'subagent',
  429. parentSessionId: SessionId('parent'),
  430. childSessionId: SessionId('missing-child'),
  431. mode: 'continuable',
  432. },
  433. throughSeq: -1,
  434. }, signal())).rejects.toMatchObject({ failure: { code: 'subagent-not-found' } })
  435. expect(inspect).toHaveBeenCalledTimes(2)
  436. })
  437. it('rejects incomplete cold metadata before serving a source', async () => {
  438. const first = await setup()
  439. const sessionId = SessionId('incomplete')
  440. const address = { kind: 'session' as const, sessionId }
  441. const firstHeader = { version: 0, id: sessionId, createdAt: 1 }
  442. first.ctx.provide('sessionPersistence', testSessionPersistence(first.ctx, {
  443. list: () => Promise.resolve([firstHeader]),
  444. inspect: () => Promise.resolve({ meta: firstHeader, events: [] }),
  445. }) as never)
  446. await expect(first.transport.page({ address, throughSeq: -1 }, signal()))
  447. .rejects.toMatchObject({ failure: { code: 'session-not-found' } })
  448. const second = await setup()
  449. const listed = { version: 0, id: sessionId, createdAt: 1, cwd: '/workspace' }
  450. const inspected = { version: 0, id: sessionId, createdAt: 1 }
  451. second.ctx.provide('sessionPersistence', testSessionPersistence(second.ctx, {
  452. list: () => Promise.resolve([listed]),
  453. inspect: () => Promise.resolve({ meta: inspected, events: [] }),
  454. }) as never)
  455. await expect(second.transport.page({ address, throughSeq: -1 }, signal()))
  456. .rejects.toMatchObject({ failure: { code: 'session-not-found' } })
  457. })
  458. it('serves cold ordinary history and validates every durable subagent descriptor state', async () => {
  459. const ordinaryBench = await setup()
  460. const ordinaryId = SessionId('cold-ordinary')
  461. const ordinaryHeader = { version: 0, id: ordinaryId, createdAt: 1, cwd: '/workspace' }
  462. cold(ordinaryBench.ctx, ordinaryHeader, [event('turn/start', 0, { turn: 1 })])
  463. await expect(ordinaryBench.transport.page({
  464. address: { kind: 'session', sessionId: ordinaryId },
  465. throughSeq: 0,
  466. }, signal())).resolves.toMatchObject({ events: [{ event: { seq: 0 } }] })
  467. const parentSessionId = SessionId('cold-parent')
  468. const childSessionId = SessionId('cold-child')
  469. const childHeader = {
  470. version: 0,
  471. id: childSessionId,
  472. createdAt: 1,
  473. cwd: '/workspace',
  474. origin: 'subagent' as const,
  475. parentSession: parentSessionId,
  476. }
  477. const childAddress = {
  478. kind: 'subagent' as const,
  479. parentSessionId,
  480. childSessionId,
  481. mode: 'continuable' as const,
  482. }
  483. const missing = await setup()
  484. cold(missing.ctx, childHeader, [])
  485. await expect(missing.transport.page({ address: childAddress, throughSeq: -1 }, signal()))
  486. .rejects.toMatchObject({ failure: { code: 'subagent-catalog-diagnostic', details: { reason: 'corrupt' } } })
  487. const corrupt = await setup()
  488. cold(corrupt.ctx, childHeader, [event('subagent/descriptor', 0, { version: 'bad' })])
  489. await expect(corrupt.transport.page({ address: childAddress, throughSeq: 0 }, signal()))
  490. .rejects.toMatchObject({ failure: { code: 'subagent-catalog-diagnostic', details: { reason: 'corrupt' } } })
  491. const ordinaryChild = await setup()
  492. const { origin: _origin, ...ordinaryChildHeader } = childHeader
  493. cold(ordinaryChild.ctx, ordinaryChildHeader, [])
  494. await expect(ordinaryChild.transport.page({ address: childAddress, throughSeq: -1 }, signal()))
  495. .rejects.toMatchObject({ failure: { code: 'subagent-unauthorized' } })
  496. })
  497. it('reports an unavailable descriptor when an observed child has no projection value', async () => {
  498. const ctx = new Context()
  499. await ctx.plugin(SessionStore)
  500. const parentSessionId = SessionId('missing-projection-parent')
  501. const childSessionId = SessionId('missing-projection-child')
  502. const meta: SessionHeader = {
  503. version: 0,
  504. id: childSessionId,
  505. createdAt: 1,
  506. cwd: '/workspace',
  507. origin: 'subagent',
  508. parentSession: parentSessionId,
  509. }
  510. ctx.provide('sessionQuery', {
  511. observeSession: () => Promise.resolve({
  512. source: 'live', header: meta, events: [], cursor: -1,
  513. projections: { asOfSeq: -1, values: {} },
  514. retain: vi.fn(), [Symbol.dispose]: vi.fn(),
  515. } as unknown as SessionObservation),
  516. } as never)
  517. const history = new SessionHistoryController(ctx, vi.fn())
  518. await expect(history.page({
  519. address: { kind: 'subagent', parentSessionId, childSessionId, mode: 'continuable' },
  520. throughSeq: -1,
  521. }, signal())).rejects.toMatchObject({
  522. failure: { code: 'subagent-catalog-diagnostic', details: { reason: 'unsupported' } },
  523. })
  524. await ctx.fiber.dispose()
  525. })
  526. it('keeps pages projection-free and computes projections only for child authorization', async () => {
  527. const ordinary = await setup()
  528. const session = ordinary.ctx.sessions.create(SessionId('projected'), { meta: { cwd: '/workspace' } })
  529. session.append('turn/start', { turn: 1 })
  530. const ordinarySnapshot = vi.spyOn(ordinary.ctx.sessionProjections, 'snapshot')
  531. const ordinaryPage = await ordinary.transport.page({
  532. address: { kind: 'session', sessionId: session.id },
  533. throughSeq: 0,
  534. }, signal())
  535. expect('projections' in ordinaryPage).toBe(false)
  536. expect(ordinarySnapshot).not.toHaveBeenCalled()
  537. const child = await setup()
  538. const parentSessionId = SessionId('projection-parent')
  539. const childSessionId = SessionId('projection-child')
  540. const childSession = child.ctx.sessions.create(childSessionId, {
  541. meta: { cwd: '/workspace', origin: 'subagent', parentSession: parentSessionId },
  542. })
  543. childSession.append('subagent/descriptor', snapshotSubagentDescriptor({
  544. mode: 'continuable', provider: 'test', label: 'child',
  545. }))
  546. const childSnapshot = vi.spyOn(child.ctx.sessionProjections, 'snapshot')
  547. const page = await child.transport.page({
  548. address: { kind: 'subagent', parentSessionId, childSessionId, mode: 'continuable' },
  549. throughSeq: 0,
  550. }, signal())
  551. expect('projections' in page).toBe(false)
  552. expect(childSnapshot).toHaveBeenCalledWith(childSession)
  553. })
  554. it('keeps message-aligned pagination contiguous across replacement provenance', async () => {
  555. const { ctx, transport } = await setup()
  556. const session = ctx.sessions.create(SessionId('pagination'), { meta: { cwd: '/workspace' } })
  557. session.append('turn/start', { turn: 1 })
  558. append(session, 'user/message', { content: [], source: { kind: 'user' } }, { surfaceOp: 'append' })
  559. const firstReply = append(session, 'assistant/message', { turn: 1, step: 1, message: {} }, { surfaceOp: 'append' })
  560. append(session, 'user/message', { content: [], source: { kind: 'user' } }, { surfaceOp: 'append' })
  561. append(session, 'assistant/message', { turn: 1, step: 2, message: {} }, { surfaceOp: 'append' })
  562. const summary = append(session, 'fixture/summary', {})
  563. const replacement = append(session, 'user/message', { content: [], source: { kind: 'plugin' } }, {
  564. surfaceOp: { op: 'replace', start: 1, end: 4 },
  565. sourceEventSeqs: [1, firstReply.seq, 3, 4, summary.seq],
  566. })
  567. const page = await transport.page({
  568. address: { kind: 'session', sessionId: session.id }, throughSeq: replacement.seq, maxMessages: 2,
  569. }, signal())
  570. expect(page.events.map(entry => entry.event.seq)).toEqual([3, 4, 5, replacement.seq])
  571. expect(page.hasMore).toBe(true)
  572. const before = await transport.page({
  573. address: { kind: 'session', sessionId: session.id }, throughSeq: replacement.seq, beforeSeq: 3, maxMessages: 1,
  574. }, signal())
  575. expect(before.events.map(entry => entry.event.seq)).toEqual([2])
  576. })
  577. it('keeps cited source events in the page that owns their appended message', async () => {
  578. const { ctx, transport } = await setup()
  579. const session = ctx.sessions.create(SessionId('pagination-sources'), { meta: { cwd: '/workspace' } })
  580. const source = append(session, 'fixture/source', {})
  581. append(session, 'user/message', { content: [], source: { kind: 'plugin' } }, {
  582. surfaceOp: 'append', sourceEventSeqs: [source.seq],
  583. })
  584. const page = await transport.page({
  585. address: { kind: 'session', sessionId: session.id }, throughSeq: 1, maxMessages: 1,
  586. }, signal())
  587. expect(page.events.map(entry => entry.event.seq)).toEqual([0, 1])
  588. expect(page.hasMore).toBe(false)
  589. })
  590. })