transport.host.spec.ts 28 KB

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