transport.host.spec.ts 29 KB

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