transport.host.spec.ts 28 KB

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