session-cold.host.spec.ts 38 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961
  1. /**
  2. * Cold-session and degenerate-composition paths of the Session Controller:
  3. * metadata-only listing, Agent-free history reads, subagent ownership
  4. * isolation, and prompt failure mapping.
  5. */
  6. import { describe, expect, it, vi } from 'vitest'
  7. import { mkdtempSync, writeFileSync } from 'node:fs'
  8. import { tmpdir } from 'node:os'
  9. import { join } from 'node:path'
  10. import { Context } from '@deepseek-ai/cordis'
  11. import SessionStore, { SESSION_FORMAT_VERSION, SessionLogOffset, SessionSeq } from '@deepseek-ai/dsh-session'
  12. import AgentRegistry from '@deepseek-ai/dsh-agent'
  13. import { SessionHistoryController } from '@deepseek-ai/dsh-api-session-controller/src/history.ts'
  14. import { subagentIdentityProjectionDefinition } from '@deepseek-ai/dsh-subagent/src/projection.ts'
  15. import TypertRegistry from '@deepseek-ai/dsh-typert-registry'
  16. import { createUserMessage, MessageId } from '@deepseek-ai/dsh-llm'
  17. import { snapshotSubagentDescriptor } from '@deepseek-ai/dsh-subagent'
  18. import type { Agent } from '@deepseek-ai/dsh-agent'
  19. import type { Session, SessionEvent, SessionHeader, SessionId } from '@deepseek-ai/dsh-session'
  20. import type { SessionPromptRequest, SessionRequestId } from '../src/types.ts'
  21. import {
  22. PersistenceCoordinator,
  23. SessionPersistenceRevision,
  24. type PersistenceBackend,
  25. type StoredPrefix,
  26. } from '@deepseek-ai/dsh-session-persistence'
  27. import { ApiSessionList } from '../src/list.ts'
  28. import {
  29. createSessionTestRemote,
  30. currentSessionListing,
  31. installSessionReadTestServices,
  32. testSessionPersistence,
  33. } from './test-remote.ts'
  34. const sid = (id: string): SessionId => id as SessionId
  35. function request<P>(payload: P): P {
  36. return payload
  37. }
  38. let nextRequestId = 1
  39. function promptRequest(
  40. payload: Omit<SessionPromptRequest, 'requestId'>,
  41. ): SessionPromptRequest {
  42. return {
  43. ...payload,
  44. requestId: `cold-${String(nextRequestId++)}` as SessionRequestId,
  45. }
  46. }
  47. function header(id: string, createdAt: number, extra: Partial<SessionHeader> = {}): SessionHeader {
  48. return { version: SESSION_FORMAT_VERSION, id: sid(id), createdAt, cwd: '/proj', isSeeded: false, ...extra }
  49. }
  50. function providePersistence(ctx: Context, persistence: Record<string, unknown>): () => void {
  51. return ctx.provide('sessionPersistence', testSessionPersistence(ctx, persistence) as never)
  52. }
  53. describe('sessions.list cold merge', () => {
  54. it('fully observes only small possibly-blank artifacts and treats unavailable probes as visible', async () => {
  55. const ctx = new Context()
  56. await ctx.plugin(SessionStore)
  57. const root = mkdtempSync(join(tmpdir(), 'dsh-cold-'))
  58. const smallPath = join(root, 'small.log')
  59. const historicalPath = join(root, 'historical.log')
  60. const largePath = join(root, 'large.log')
  61. writeFileSync(smallPath, 'x'.repeat(1024))
  62. writeFileSync(historicalPath, 'historical')
  63. writeFileSync(largePath, 'x'.repeat(1025))
  64. const metas: SessionHeader[] = [
  65. header('small-blank', 100),
  66. header('small-conversation', 200),
  67. header('large-unknown', 300),
  68. header('cached-nonblank', 400),
  69. header('seeded-cold', 450, { isSeeded: true }),
  70. header('locationless', 500, { parentSession: sid('session-parent'), origin: 'subagent' }),
  71. header('vanished', 600),
  72. header('read-failure', 700),
  73. { version: SESSION_FORMAT_VERSION, id: sid('missing-cwd'), createdAt: 800, isSeeded: false },
  74. header('migration-small-blank', 900),
  75. ]
  76. const inspect = vi.fn(async (id: SessionId) => {
  77. if (id === sid('small-blank')) {
  78. return {
  79. meta: metas[0]!,
  80. events: [{ type: 'session/end-seed', seq: SessionSeq(0), time: 700, data: {} }] satisfies SessionEvent[],
  81. }
  82. }
  83. if (id === sid('small-conversation')) {
  84. return {
  85. meta: metas[1]!,
  86. events: [
  87. { type: 'turn/start', seq: SessionSeq(0), time: 800, data: { turn: 1 } },
  88. {
  89. type: 'user/message', seq: SessionSeq(1), time: 1200,
  90. data: createUserMessage({ content: [{ type: 'text', text: 'worked' }], source: { kind: 'user' } }),
  91. surfaceOp: 'append',
  92. },
  93. ] satisfies SessionEvent[],
  94. }
  95. }
  96. if (id === sid('migration-small-blank')) {
  97. return {
  98. meta: metas[9]!,
  99. events: [{ type: 'session/end-seed', seq: 0, time: 900, data: {} }] as SessionEvent[],
  100. }
  101. }
  102. if (id === sid('read-failure')) throw new Error('simulated read failure')
  103. throw new Error(`unexpected cold read: ${id}`)
  104. })
  105. providePersistence(ctx, {
  106. list: () => Promise.resolve([
  107. ...metas.slice(0, -1).map(currentSessionListing),
  108. {
  109. status: 'migration-required' as const,
  110. storedVersion: 0,
  111. targetVersion: SESSION_FORMAT_VERSION,
  112. header: metas[9]!,
  113. location: { kind: 'jsonl', path: historicalPath },
  114. },
  115. ]),
  116. locate: (meta: SessionHeader) => {
  117. if (meta.id === sid('large-unknown') || meta.id === sid('seeded-cold')) {
  118. return { kind: 'jsonl', path: largePath }
  119. }
  120. if (meta.id === sid('migration-small-blank')) {
  121. return { kind: 'jsonl', path: join(root, 'missing-current.log') }
  122. }
  123. if (meta.id === sid('locationless')) return undefined
  124. if (meta.id === sid('vanished')) return { kind: 'jsonl', path: join(root, 'vanished.log') }
  125. return { kind: 'jsonl', path: smallPath }
  126. },
  127. inspect,
  128. })
  129. ctx.provide('sessionProjectionCache', {
  130. cachedSnapshot: (meta: SessionHeader) => {
  131. if (meta.id === sid('seeded-cold')) throw new Error('seeded cold listing must not guess a body cut')
  132. if (meta.id === sid('small-blank')) {
  133. return { asOfSeq: SessionSeq(0), values: { sessionListMetadata: { blank: true, lastPromptAt: null } } }
  134. }
  135. if (meta.id === sid('small-conversation')) {
  136. return { asOfSeq: SessionSeq(0), values: { sessionListMetadata: { blank: true, lastPromptAt: 900 } } }
  137. }
  138. if (meta.id === sid('cached-nonblank')) {
  139. return { asOfSeq: SessionSeq(1), values: { sessionListMetadata: { blank: false, lastPromptAt: 1000 } } }
  140. }
  141. return undefined
  142. },
  143. hydratePrepared: (session: Session, events: readonly SessionEvent[]) =>
  144. ctx.sessionProjections.hydrate(session, {}, events, SessionLogOffset(0)),
  145. } as never)
  146. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  147. const response = await remote.list(request({}))
  148. expect(response.ok).toBe(true)
  149. if (!response.ok) throw new Error('unreachable')
  150. const byId = Object.fromEntries(response.value.items.map(item => [item.sessionId, item]))
  151. expect(byId['small-blank']).toMatchObject({ blank: true, updatedAt: 100, running: false })
  152. expect(byId['small-conversation']).toMatchObject({ blank: false, updatedAt: 1200 })
  153. expect(byId['large-unknown']).toMatchObject({ blank: false, updatedAt: 300 })
  154. expect(byId['cached-nonblank']).toMatchObject({ blank: false, updatedAt: 1000 })
  155. expect(byId['seeded-cold']).toMatchObject({ blank: false, updatedAt: 450 })
  156. expect(byId['locationless']).toMatchObject({
  157. blank: false,
  158. updatedAt: 500,
  159. parentSessionId: 'session-parent',
  160. origin: 'subagent',
  161. })
  162. expect(byId['vanished']).toMatchObject({ blank: false, updatedAt: 600 })
  163. expect(byId['read-failure']).toMatchObject({ blank: false, updatedAt: 700 })
  164. expect(byId['missing-cwd']).toBeUndefined()
  165. expect(byId['migration-small-blank']).toMatchObject({ blank: true, updatedAt: 900 })
  166. expect(inspect).toHaveBeenCalledTimes(4)
  167. expect(inspect.mock.calls.map(([id]) => id)).toEqual(expect.arrayContaining([
  168. sid('small-blank'),
  169. sid('small-conversation'),
  170. sid('read-failure'),
  171. sid('migration-small-blank'),
  172. ]))
  173. })
  174. it('can disable bounded cold observations without hiding cold Sessions', async () => {
  175. const ctx = new Context()
  176. await ctx.plugin(SessionStore)
  177. const meta = header('probe-disabled', 100)
  178. const inspect = vi.fn()
  179. providePersistence(ctx, {
  180. list: () => Promise.resolve([currentSessionListing(meta)]),
  181. locate: () => ({ kind: 'jsonl', path: '/not-read' }),
  182. inspect,
  183. })
  184. const remote = createSessionTestRemote(ctx, {
  185. defaultModelSelection: () => ({ provider: 'p', model: 'm' }),
  186. cwd: '/tmp',
  187. coldBlankProbeMaxBytes: 0,
  188. })
  189. const response = await remote.list(request({}))
  190. if (!response.ok) throw new Error('unreachable')
  191. expect(response.value.items).toEqual([
  192. expect.objectContaining({ sessionId: meta.id, blank: false, updatedAt: meta.createdAt }),
  193. ])
  194. expect(inspect).not.toHaveBeenCalled()
  195. })
  196. it('prefers a live row attached during the query without folding its seed', async () => {
  197. const ctx = new Context()
  198. await ctx.plugin(SessionStore)
  199. await ctx.plugin(AgentRegistry)
  200. const meta = header('attached-during-list', 100)
  201. const started = Promise.withResolvers<undefined>()
  202. const release = Promise.withResolvers<undefined>()
  203. providePersistence(ctx, {
  204. list: async () => {
  205. started.resolve(undefined)
  206. await release.promise
  207. return [currentSessionListing(meta)]
  208. },
  209. })
  210. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  211. const listing = remote.list(request({}))
  212. await started.promise
  213. const session = ctx.sessions.create(meta.id, {
  214. seed: [
  215. { type: 'turn/start', seq: SessionSeq(0), time: 200, data: { turn: 1 } },
  216. {
  217. type: 'user/message', seq: SessionSeq(1), time: 300,
  218. data: createUserMessage({ content: [{ type: 'text', text: 'live' }], source: { kind: 'user' } }),
  219. surfaceOp: 'append',
  220. },
  221. ],
  222. meta: {
  223. ...meta.cwd === undefined ? {} : { cwd: meta.cwd },
  224. createdAt: meta.createdAt,
  225. },
  226. })
  227. ctx.agents.register({ id: session.id, session, status: 'running', ctx } as Agent)
  228. release.resolve(undefined)
  229. const response = await listing
  230. if (!response.ok) throw new Error('list failed')
  231. expect(response.value.items).toEqual([
  232. expect.objectContaining({
  233. sessionId: meta.id,
  234. blank: false,
  235. running: true,
  236. updatedAt: 100,
  237. }),
  238. ])
  239. })
  240. it('prefers a Session that attaches during its bounded cold observation', async () => {
  241. const ctx = new Context()
  242. await ctx.plugin(SessionStore)
  243. await ctx.plugin(AgentRegistry)
  244. const root = mkdtempSync(join(tmpdir(), 'dsh-cold-race-'))
  245. const path = join(root, 'small.log')
  246. writeFileSync(path, 'small')
  247. const meta = header('attached-during-probe', 100)
  248. providePersistence(ctx, {
  249. list: () => Promise.resolve([currentSessionListing(meta)]),
  250. locate: () => ({ kind: 'jsonl', path }),
  251. inspect: () => {
  252. const session = ctx.sessions.create(meta.id, {
  253. meta,
  254. seed: [{ type: 'turn/start', seq: SessionSeq(0), time: 200, data: { turn: 1 } }],
  255. })
  256. ctx.agents.register({ id: session.id, session, status: 'running', ctx } as Agent)
  257. return Promise.resolve({
  258. meta,
  259. inheritedEventCount: SessionLogOffset(0),
  260. events: [],
  261. })
  262. },
  263. })
  264. const remote = createSessionTestRemote(ctx, {
  265. defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp',
  266. })
  267. const response = await remote.list(request({}))
  268. if (!response.ok) throw new Error('list failed')
  269. expect(response.value.items).toEqual([
  270. expect.objectContaining({ sessionId: meta.id, running: true, blank: false }),
  271. ])
  272. await ctx.fiber.dispose()
  273. })
  274. it('propagates a cold location failure instead of returning a partial list', async () => {
  275. const ctx = new Context()
  276. await ctx.plugin(SessionStore)
  277. const meta = header('broken-cache', 100)
  278. providePersistence(ctx, {
  279. list: () => Promise.resolve([currentSessionListing(meta)]),
  280. locate: () => { throw new Error('location failed') },
  281. })
  282. const remote = createSessionTestRemote(ctx, {
  283. defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp',
  284. })
  285. await expect(remote.list(request({}))).resolves.toMatchObject({
  286. ok: false,
  287. error: { message: expect.stringContaining('location failed') as string },
  288. })
  289. await ctx.fiber.dispose()
  290. })
  291. it('supports an unsignalled probe whose observation has no projection registry', async () => {
  292. const ctx = new Context()
  293. await ctx.plugin(SessionStore)
  294. await ctx.plugin(AgentRegistry)
  295. installSessionReadTestServices(ctx)
  296. const root = mkdtempSync(join(tmpdir(), 'dsh-cold-unprojected-'))
  297. const path = join(root, 'small.log')
  298. writeFileSync(path, 'small')
  299. const meta = header('unprojected-small', 100)
  300. ctx.provide('sessionPersistence', {
  301. list: () => Promise.resolve([currentSessionListing(meta)]),
  302. locate: () => ({ kind: 'jsonl', path }),
  303. } as never)
  304. vi.spyOn(ctx.sessionQuery, 'listSessions').mockResolvedValue([{
  305. header: meta, live: false, persisted: true,
  306. }])
  307. vi.spyOn(ctx.sessionQuery, 'observeSession').mockResolvedValue({
  308. source: 'prepared',
  309. header: meta,
  310. inheritedEventCount: SessionLogOffset(0),
  311. events: [],
  312. cursor: -1,
  313. retain: vi.fn(), [Symbol.dispose]: vi.fn(),
  314. })
  315. const list = new ApiSessionList(ctx, 1024)
  316. await expect(list.list()).resolves.toEqual([
  317. expect.objectContaining({ sessionId: meta.id, blank: false }),
  318. ])
  319. await ctx.fiber.dispose()
  320. })
  321. })
  322. describe('attached updatedAt tracks human prompts', () => {
  323. it('ignores pickup and non-prompt work after the latest human message', async () => {
  324. const ctx = new Context()
  325. await ctx.plugin(SessionStore)
  326. await ctx.plugin(AgentRegistry)
  327. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  328. await new Promise(resolve => setTimeout(resolve, 0))
  329. // Old work, resumed just now: the log tail would report the pickup.
  330. const worked = 1_000_000
  331. const resumed = ctx.sessions.create(sid('resumed-untouched'), {
  332. seed: [
  333. { type: 'turn/start', seq: SessionSeq(0), time: worked, data: { turn: 1 } },
  334. {
  335. type: 'user/message', seq: SessionSeq(1), time: worked,
  336. data: createUserMessage({ content: [{ type: 'text', text: 'worked' }], source: { kind: 'user' } }),
  337. surfaceOp: 'append',
  338. },
  339. { type: 'turn/end', seq: SessionSeq(2), time: worked + 1, data: { turn: 1, reason: { kind: 'completed' } } },
  340. ],
  341. meta: { cwd: '/proj', createdAt: 500 },
  342. })
  343. ctx.agents.register({ id: resumed.id, session: resumed, status: 'idle', ctx } as Agent)
  344. const boundary = resumed.snapshotEvents().at(-1)
  345. expect(boundary?.type).toBe('session/end-seed')
  346. expect(boundary?.time).toBeGreaterThan(worked)
  347. const listed = await remote.list(request({}))
  348. if (!listed.ok) throw new Error('list failed')
  349. const summary = listed.value.items.find(item => item.sessionId === 'resumed-untouched')
  350. expect(summary?.updatedAt).toBe(500)
  351. // A lifecycle boundary is not a human update.
  352. resumed.append('turn/start', { turn: 2 })
  353. const afterBoundary = await remote.list(request({}))
  354. if (!afterBoundary.ok) throw new Error('list failed')
  355. expect(afterBoundary.value.items.find(item => item.sessionId === 'resumed-untouched')?.updatedAt)
  356. .toBe(worked)
  357. const prompt = resumed.append('user/message', createUserMessage({
  358. content: [{ type: 'text', text: 'new prompt' }],
  359. source: { kind: 'user' },
  360. }), { surfaceOp: 'append' })
  361. const after = await remote.list(request({}))
  362. if (!after.ok) throw new Error('list failed')
  363. const moved = after.value.items.find(item => item.sessionId === 'resumed-untouched')
  364. expect(moved?.updatedAt).toBe(prompt.time)
  365. })
  366. })
  367. describe('cold history recovery view', () => {
  368. it('shows in-memory interruption repair without activating the session', async () => {
  369. const ctx = new Context()
  370. await ctx.plugin(SessionStore)
  371. const sessionId = sid('session-interrupted')
  372. const meta = header(sessionId, 1000)
  373. const stored: StoredPrefix<never> = {
  374. meta,
  375. inheritedEventCount: SessionLogOffset(0),
  376. events: [{ type: 'turn/start', seq: SessionSeq(0), time: 1, data: { turn: 1 } }],
  377. revision: SessionPersistenceRevision('history-recovery-test:1'),
  378. }
  379. const backend: PersistenceBackend<never> = {
  380. name: 'history-recovery-test',
  381. loadStored: id => Promise.resolve(id === sessionId ? structuredClone(stored) : undefined),
  382. readStoredRevision: id => Promise.resolve(
  383. id === sessionId ? SessionPersistenceRevision('history-recovery-test:1') : undefined,
  384. ),
  385. appendBatch: () => Promise.resolve(),
  386. commitRepair: () => Promise.resolve(),
  387. list: () => Promise.resolve([currentSessionListing(structuredClone(meta))]),
  388. }
  389. const coordinator = new PersistenceCoordinator(ctx, backend)
  390. providePersistence(ctx, {
  391. list: (signal?: AbortSignal) => backend.list(signal),
  392. inspect: (id: SessionId, signal?: AbortSignal) => coordinator.inspect(id, signal),
  393. borrowSession: (id: SessionId, signal?: AbortSignal) => coordinator.borrowSession(id, signal),
  394. locate: () => undefined,
  395. })
  396. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  397. const history = await remote.page({
  398. address: { kind: 'session', sessionId },
  399. throughSeq: 1,
  400. beforeSeq: 2,
  401. maxMessages: 10,
  402. })
  403. if (!history.ok) throw new Error('history failed')
  404. expect(history.value.records.map(record => record.event)).toMatchInlineSnapshot(`
  405. [
  406. {
  407. "data": {
  408. "turn": 1,
  409. },
  410. "seq": 0,
  411. "time": 1,
  412. "type": "turn/start",
  413. },
  414. {
  415. "data": {
  416. "reason": {
  417. "kind": "interrupted",
  418. },
  419. "turn": 1,
  420. },
  421. "seq": 1,
  422. "time": 1,
  423. "type": "turn/end",
  424. },
  425. ]
  426. `)
  427. expect(ctx.sessions.get(sessionId)).toBeUndefined()
  428. await ctx.fiber.dispose()
  429. })
  430. })
  431. describe('Remote Agent and Session lookup policy', () => {
  432. it('deduplicates a cold resume across Agent and Session parameters', async () => {
  433. const ctx = new Context()
  434. await ctx.plugin(TypertRegistry)
  435. await ctx.plugin(SessionStore)
  436. await ctx.plugin(AgentRegistry)
  437. const sessionId = sid('session-remote-cold')
  438. const meta = header(sessionId, 1000)
  439. const inspect = vi.fn(() => Promise.resolve({
  440. meta,
  441. inheritedEventCount: SessionLogOffset(0),
  442. events: [] as SessionEvent[],
  443. }))
  444. providePersistence(ctx, {
  445. list: () => Promise.resolve([currentSessionListing(meta)]),
  446. inspect,
  447. locate: () => undefined,
  448. })
  449. const resumedSession = { id: sessionId, header: meta, events: [] } as unknown as import('@deepseek-ai/dsh-session').Session
  450. const resumedAgent = { id: sessionId, session: resumedSession, status: 'idle', ctx } as Agent
  451. const release = Promise.withResolvers<undefined>()
  452. const resume = vi.spyOn(ctx.agents, 'resume').mockImplementation(async () => {
  453. await release.promise
  454. return { agent: resumedAgent, dispose: () => Promise.resolve() }
  455. })
  456. const defaultAgentLookup = ctx.typert.lookups.get('agent')
  457. const defaultSessionLookup = ctx.typert.lookups.get('session')
  458. createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  459. await vi.waitFor(() => {
  460. expect(ctx.typert.lookups.get('agent')).not.toBe(defaultAgentLookup)
  461. expect(ctx.typert.lookups.get('session')).not.toBe(defaultSessionLookup)
  462. })
  463. const agentLookup = ctx.typert.lookups.get('agent')
  464. const sessionLookup = ctx.typert.lookups.get('session')
  465. if (agentLookup === undefined || sessionLookup === undefined) throw new Error('core lookup providers were not mounted')
  466. const resolvedAgent = Promise.resolve(agentLookup.resolve(sessionId))
  467. const resolvedSession = Promise.resolve(sessionLookup.resolve(sessionId))
  468. await vi.waitFor(() => { expect(resume).toHaveBeenCalledOnce() })
  469. release.resolve(undefined)
  470. await expect(resolvedAgent).resolves.toBe(resumedAgent)
  471. await expect(resolvedSession).resolves.toBe(resumedSession)
  472. expect(inspect).toHaveBeenCalledOnce()
  473. })
  474. it('preserves the subagent ownership fence for cold and live Remote lookups', async () => {
  475. const ctx = new Context()
  476. await ctx.plugin(TypertRegistry)
  477. await ctx.plugin(SessionStore)
  478. await ctx.plugin(AgentRegistry)
  479. const coldId = sid('session-remote-cold-child')
  480. const coldMeta = header(coldId, 1000, {
  481. parentSession: sid('session-parent'),
  482. origin: 'subagent',
  483. })
  484. const inspect = vi.fn(() => Promise.resolve({
  485. meta: coldMeta,
  486. inheritedEventCount: SessionLogOffset(0),
  487. events: [] as SessionEvent[],
  488. }))
  489. providePersistence(ctx, {
  490. list: () => Promise.resolve([currentSessionListing(coldMeta)]),
  491. inspect,
  492. locate: () => undefined,
  493. })
  494. const liveSession = ctx.sessions.create(sid('session-remote-live-child'), {
  495. meta: { cwd: '/proj', parentSession: sid('session-parent'), origin: 'subagent' },
  496. })
  497. const liveAgent = { id: liveSession.id, session: liveSession, status: 'idle', ctx } as Agent
  498. ctx.agents.register(liveAgent)
  499. const resume = vi.spyOn(ctx.agents, 'resume')
  500. const defaultAgentLookup = ctx.typert.lookups.get('agent')
  501. const defaultSessionLookup = ctx.typert.lookups.get('session')
  502. createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  503. await vi.waitFor(() => {
  504. expect(ctx.typert.lookups.get('agent')).not.toBe(defaultAgentLookup)
  505. expect(ctx.typert.lookups.get('session')).not.toBe(defaultSessionLookup)
  506. })
  507. const agentLookup = ctx.typert.lookups.get('agent')
  508. const sessionLookup = ctx.typert.lookups.get('session')
  509. if (agentLookup === undefined || sessionLookup === undefined) throw new Error('core lookup providers were not mounted')
  510. const ownershipFailure = {
  511. code: 'session/agent-busy',
  512. details: { reason: 'use subagent delivery for this child session' },
  513. }
  514. const coldFailure = Promise.resolve(agentLookup.resolve(coldId))
  515. const liveFailure = Promise.resolve(sessionLookup.resolve(liveSession.id))
  516. await expect(coldFailure).rejects.toMatchObject(ownershipFailure)
  517. await expect(liveFailure).rejects.toMatchObject(ownershipFailure)
  518. expect(resume).not.toHaveBeenCalled()
  519. expect(inspect).toHaveBeenCalledOnce()
  520. })
  521. })
  522. describe('subagent ownership fence', () => {
  523. it('reads a cold child without an Agent and rejects generic resume or adoption', async () => {
  524. const ctx = new Context()
  525. await ctx.plugin(SessionStore)
  526. await ctx.plugin(AgentRegistry)
  527. const sessionId = sid('session-child')
  528. const meta = header('session-child', 1000, {
  529. parentSession: sid('session-parent'),
  530. isSeeded: true,
  531. origin: 'subagent',
  532. })
  533. const events = [
  534. {
  535. type: 'turn/start',
  536. seq: SessionSeq(0),
  537. time: 1,
  538. data: {
  539. turn: 1,
  540. trigger: { kind: 'message', source: { kind: 'user' } },
  541. } as SessionEvent<'turn/start'>['data'],
  542. },
  543. {
  544. type: 'user/message',
  545. seq: SessionSeq(1),
  546. time: 2,
  547. data: createUserMessage({ content: [{ type: 'text', text: 'work' }], source: { kind: 'user' } }),
  548. surfaceOp: 'append',
  549. },
  550. {
  551. type: 'subagent/descriptor',
  552. seq: SessionSeq(2),
  553. time: 3,
  554. data: snapshotSubagentDescriptor({
  555. mode: 'continuable',
  556. provider: 'spawn',
  557. label: 'child',
  558. }),
  559. },
  560. { type: 'turn/end', seq: SessionSeq(3), time: 4, data: { turn: 1, reason: { kind: 'completed' } } },
  561. ] satisfies SessionEvent[]
  562. const inspect = vi.fn(() => Promise.resolve({
  563. meta,
  564. inheritedEventCount: SessionLogOffset(0),
  565. events,
  566. }))
  567. providePersistence(ctx, {
  568. list: () => Promise.resolve([currentSessionListing(meta)]),
  569. inspect,
  570. locate: () => undefined,
  571. })
  572. const resume = vi.spyOn(ctx.agents, 'resume')
  573. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  574. ctx.sessionProjections.register(subagentIdentityProjectionDefinition)
  575. const history = await new SessionHistoryController(
  576. ctx,
  577. (observation) => { observation[Symbol.dispose]() },
  578. ).page({
  579. address: {
  580. kind: 'subagent',
  581. parentSessionId: meta.parentSession as SessionId,
  582. childSessionId: sessionId,
  583. mode: 'continuable',
  584. },
  585. throughSeq: 3,
  586. }, new AbortController().signal)
  587. expect(history.records.map(record => record.event.type))
  588. .toEqual(events.map(event => event.type))
  589. expect(ctx.agents.get(sessionId)).toBeUndefined()
  590. const prompt = await remote.prompt(promptRequest({
  591. sessionId,
  592. mode: 'queue',
  593. content: [{ type: 'text', text: 'follow up' }],
  594. }))
  595. expect(prompt.ok).toBe(false)
  596. if (!prompt.ok) {
  597. expect(prompt.error).toMatchObject({
  598. code: 'session/agent-busy',
  599. details: { reason: 'use subagent delivery for this child session' },
  600. })
  601. }
  602. const create = await remote.create(request({ sessionId, cwd: '/proj' }))
  603. expect(create.ok).toBe(false)
  604. if (!create.ok) expect(create.error.code).toBe('session/agent-busy')
  605. expect(resume).not.toHaveBeenCalled()
  606. expect(ctx.agents.get(sessionId)).toBeUndefined()
  607. expect(inspect).toHaveBeenCalledTimes(3)
  608. })
  609. it('no longer treats a descriptor-only cold child without origin as subagent-owned', async () => {
  610. const ctx = new Context()
  611. await ctx.plugin(SessionStore)
  612. await ctx.plugin(AgentRegistry)
  613. const sessionId = sid('session-legacy-child')
  614. const meta = header('session-legacy-child', 1000, {
  615. parentSession: sid('session-parent'),
  616. isSeeded: true,
  617. })
  618. const events = [
  619. {
  620. type: 'subagent/descriptor',
  621. seq: SessionSeq(0),
  622. time: 1,
  623. data: { version: 2, mode: 'continuable', provider: 'spawn', label: 'child' },
  624. },
  625. ] satisfies SessionEvent[]
  626. providePersistence(ctx, {
  627. list: () => Promise.resolve([currentSessionListing(meta)]),
  628. inspect: () => Promise.resolve({
  629. meta,
  630. inheritedEventCount: SessionLogOffset(0),
  631. events,
  632. }),
  633. locate: () => undefined,
  634. })
  635. // Stores whose headers predate `origin` classify a child only through the
  636. // descriptor event; the pre-release decision stops recognizing them, so
  637. // the ownership fence lets generic resume reach the registry instead of
  638. // answering `agent-busy`.
  639. const resume = vi.spyOn(ctx.agents, 'resume')
  640. .mockRejectedValue(new Error('registry unavailable in this bench'))
  641. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  642. const prompt = await remote.prompt(promptRequest({
  643. sessionId,
  644. mode: 'queue',
  645. content: [{ type: 'text', text: 'follow up' }],
  646. }))
  647. expect(resume).toHaveBeenCalledTimes(1)
  648. expect(prompt.ok).toBe(false)
  649. if (!prompt.ok) expect(prompt.error.code).toBe('gateway/internal')
  650. })
  651. it('rejects origin-marked and runtime-owned live children from generic controls', async () => {
  652. const ctx = new Context()
  653. await ctx.plugin(SessionStore)
  654. await ctx.plugin(AgentRegistry)
  655. const parentSession = ctx.sessions.create(sid('session-parent'), { meta: { cwd: '/proj' } })
  656. const parent = { id: parentSession.id, session: parentSession, status: 'idle', ctx } as Agent
  657. ctx.agents.register(parent)
  658. const originSession = ctx.sessions.create(sid('session-origin-child'), {
  659. meta: { cwd: '/proj', parentSession: parent.id, origin: 'subagent' },
  660. })
  661. const cancel = vi.fn()
  662. const updateInbox = vi.fn(() => 'applied' as const)
  663. const originChild = {
  664. id: originSession.id,
  665. session: originSession,
  666. status: 'idle',
  667. ctx,
  668. cancel,
  669. updateInbox,
  670. } as unknown as Agent
  671. ctx.agents.register(originChild)
  672. const startingSession = ctx.sessions.create(sid('session-starting-child'), {
  673. meta: { cwd: '/proj', parentSession: parent.id },
  674. })
  675. const startingChild = { id: startingSession.id, session: startingSession, status: 'idle', ctx } as Agent
  676. ctx.agents.enter(startingChild, parent)
  677. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  678. const stopped = await remote.cancel(request({ sessionId: originChild.id }))
  679. expect(stopped.ok).toBe(false)
  680. if (!stopped.ok) expect(stopped.error.code).toBe('session/agent-busy')
  681. expect(cancel).not.toHaveBeenCalled()
  682. const queued = await remote.updateQueue(request({
  683. sessionId: originChild.id,
  684. itemId: MessageId('queued-item'),
  685. action: { kind: 'remove' },
  686. }))
  687. expect(queued.ok).toBe(false)
  688. if (!queued.ok) expect(queued.error.code).toBe('session/agent-busy')
  689. expect(updateInbox).not.toHaveBeenCalled()
  690. const selection = await remote.selectModel(request({
  691. sessionId: startingChild.id,
  692. provider: 'p',
  693. model: 'm',
  694. }))
  695. expect(selection.ok).toBe(false)
  696. if (!selection.ok) expect(selection.error.code).toBe('session/agent-busy')
  697. const create = await remote.create(request({ sessionId: originChild.id, cwd: '/proj' }))
  698. expect(create.ok).toBe(false)
  699. if (!create.ok) expect(create.error.code).toBe('session/agent-busy')
  700. expect(ctx.agents.get(originChild.id)).toBe(originChild)
  701. })
  702. it('does not classify an ordinary fork from an inherited ancestor descriptor', async () => {
  703. const ctx = new Context()
  704. await ctx.plugin(SessionStore)
  705. await ctx.plugin(AgentRegistry)
  706. const session = ctx.sessions.create(sid('session-ordinary-fork'), {
  707. seed: [{
  708. type: 'subagent/descriptor',
  709. seq: SessionSeq(0),
  710. time: 1,
  711. data: { version: 2, mode: 'continuable', provider: 'spawn', label: 'ancestor' },
  712. }],
  713. inheritedEventCount: SessionLogOffset(1),
  714. meta: { cwd: '/proj', parentSession: sid('session-source'), isSeeded: true },
  715. })
  716. const followup = vi.fn()
  717. const agent = { id: session.id, session, status: 'idle', ctx, followup } as unknown as Agent
  718. ctx.agents.register(agent)
  719. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  720. const response = await remote.prompt(promptRequest({
  721. sessionId: agent.id,
  722. mode: 'queue',
  723. content: [{ type: 'text', text: 'ordinary work' }],
  724. }))
  725. expect(response.ok).toBe(true)
  726. expect(followup).toHaveBeenCalledOnce()
  727. })
  728. it('canonicalizes a supplied browser zone on the exact prompt and rejects invalid names', async () => {
  729. const ctx = new Context()
  730. await ctx.plugin(SessionStore)
  731. await ctx.plugin(AgentRegistry)
  732. const session = ctx.sessions.create(sid('session-browser-zone'), { meta: { cwd: '/proj' } })
  733. const followup = vi.fn()
  734. const agent = { id: session.id, session, status: 'idle', ctx, followup } as unknown as Agent
  735. ctx.agents.register(agent)
  736. const remote = createSessionTestRemote(ctx, {
  737. defaultModelSelection: () => ({ provider: 'p', model: 'm' }),
  738. cwd: '/tmp',
  739. })
  740. const alias = 'US/Pacific'
  741. const canonical = new Intl.DateTimeFormat('en-US', { timeZone: alias })
  742. .resolvedOptions().timeZone
  743. const zonedRequest = promptRequest({
  744. sessionId: agent.id,
  745. mode: 'queue' as const,
  746. content: [{ type: 'text' as const, text: 'zoned work' }],
  747. clientTimeZone: alias,
  748. })
  749. await expect(remote.prompt(zonedRequest)).resolves.toMatchObject({ ok: true })
  750. expect(followup).toHaveBeenNthCalledWith(1, expect.objectContaining({
  751. source: { kind: 'user', rpcId: zonedRequest.requestId, clientTimeZone: canonical },
  752. }))
  753. const utcRequest = promptRequest({
  754. sessionId: agent.id,
  755. mode: 'queue' as const,
  756. content: [{ type: 'text' as const, text: 'UTC work' }],
  757. clientTimeZone: 'UTC',
  758. })
  759. await expect(remote.prompt(utcRequest)).resolves.toMatchObject({ ok: true })
  760. expect(followup).toHaveBeenNthCalledWith(2, expect.objectContaining({
  761. source: { kind: 'user', rpcId: utcRequest.requestId, clientTimeZone: 'UTC' },
  762. }))
  763. const unzonedRequest = promptRequest({
  764. sessionId: agent.id,
  765. mode: 'queue' as const,
  766. content: [{ type: 'text' as const, text: 'headless work' }],
  767. })
  768. await expect(remote.prompt(unzonedRequest)).resolves.toMatchObject({ ok: true })
  769. expect(followup).toHaveBeenNthCalledWith(3, expect.objectContaining({
  770. source: { kind: 'user', rpcId: unzonedRequest.requestId },
  771. }))
  772. for (const clientTimeZone of ['', ' UTC', 'CST', 'Not/A_Real_Zone']) {
  773. const invalid = await remote.prompt(promptRequest({
  774. sessionId: agent.id,
  775. mode: 'queue' as const,
  776. content: [{ type: 'text' as const, text: 'invalid zone' }],
  777. clientTimeZone,
  778. }))
  779. expect(invalid).toMatchObject({
  780. ok: false,
  781. error: {
  782. code: 'session/invalid-time-zone',
  783. message: 'clientTimeZone must be UTC or a valid IANA Area/Location name',
  784. details: { value: clientTimeZone },
  785. },
  786. })
  787. }
  788. expect(followup).toHaveBeenCalledTimes(3)
  789. })
  790. })
  791. describe('degenerate composition (no persistence, no factory)', () => {
  792. it('lists no cold rows and reports an absent point source as not found', async () => {
  793. const ctx = new Context()
  794. await ctx.plugin(SessionStore)
  795. await ctx.plugin(AgentRegistry)
  796. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  797. const listed = await remote.list(request({}))
  798. expect(listed.ok).toBe(true)
  799. if (listed.ok) expect(listed.value.items).toEqual([])
  800. // No persistence means cold history cannot inspect a transcript.
  801. const response = await remote.page({
  802. address: { kind: 'session', sessionId: sid('session-ghost') },
  803. throughSeq: -1,
  804. })
  805. expect(response.ok).toBe(false)
  806. if (!response.ok) {
  807. expect(response.error.code).toBe('session/not-found')
  808. }
  809. })
  810. it('maps a missing direct persistence read to session-not-found', async () => {
  811. const ctx = new Context()
  812. await ctx.plugin(SessionStore)
  813. await ctx.plugin(AgentRegistry)
  814. const inspect = vi.fn()
  815. providePersistence(ctx, {
  816. list: () => Promise.resolve([]),
  817. inspect,
  818. })
  819. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  820. const response = await remote.page({
  821. address: { kind: 'session', sessionId: sid('session-missing') },
  822. throughSeq: -1,
  823. })
  824. expect(response.ok).toBe(false)
  825. if (!response.ok) expect(response.error.code).toBe('session/not-found')
  826. expect(inspect).toHaveBeenCalledOnce()
  827. })
  828. })
  829. describe('sessions.prompt synchronous rejection', () => {
  830. it('maps a synchronous send throw (disposed/invalid input) to agent-busy with the reason attached', async () => {
  831. const ctx = new Context()
  832. await ctx.plugin(SessionStore)
  833. await ctx.plugin(AgentRegistry)
  834. const session = ctx.sessions.create(sid('session-throwing'))
  835. // A live structural stub whose delivery verbs throw synchronously, the
  836. // shape a disposed loop presents at this gateway boundary.
  837. ctx.agents.register({
  838. id: session.id,
  839. session,
  840. status: 'idle',
  841. ctx,
  842. followup: () => { throw new Error('agent "session-throwing" lifecycle disposed') },
  843. steer: () => { throw new Error('agent "session-throwing" lifecycle disposed') },
  844. } as unknown as Agent)
  845. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  846. for (const mode of ['queue', 'steer'] as const) {
  847. const response = await remote.prompt(promptRequest({
  848. sessionId: session.id, mode, content: [{ type: 'text' as const, text: 'x' }],
  849. }))
  850. expect(response.ok).toBe(false)
  851. if (!response.ok) {
  852. expect(response.error.code).toBe('session/agent-busy')
  853. expect(response.error.message).toBe('prompt rejected')
  854. expect(response.error.details).toEqual({
  855. reason: 'Error: agent "session-throwing" lifecycle disposed',
  856. })
  857. }
  858. }
  859. })
  860. it('classifies a raced cold-resume ID collision as agent-busy', async () => {
  861. const ctx = new Context()
  862. await ctx.plugin(SessionStore)
  863. await ctx.plugin(AgentRegistry)
  864. const sessionId = sid('race-resume')
  865. const meta: SessionHeader = header('race-resume', 1000)
  866. providePersistence(ctx, {
  867. list: () => Promise.resolve([currentSessionListing(meta)]),
  868. inspect: () => Promise.resolve({
  869. meta,
  870. inheritedEventCount: SessionLogOffset(0),
  871. events: [] as SessionEvent[],
  872. }),
  873. locate: () => undefined,
  874. })
  875. // The raced winner: a live parent-owned subagent publishes the identity
  876. // while the generic cold resume is in flight, so the resume collides.
  877. const parentSession = ctx.sessions.create(sid('race-parent'), { meta: { cwd: '/proj' } })
  878. const parent = { id: parentSession.id, session: parentSession, status: 'idle', ctx } as Agent
  879. ctx.agents.register(parent)
  880. const childSession = ctx.sessions.create(sessionId, {
  881. meta: { cwd: '/proj', parentSession: parent.id, origin: 'subagent' },
  882. })
  883. const child = { id: sessionId, session: childSession, status: 'idle', ctx } as unknown as Agent
  884. vi.spyOn(ctx.agents, 'resume').mockImplementationOnce(async () => {
  885. // The parent's `enter()` wins the identity between the pre-resume
  886. // re-check and publication; the generic resume then collides.
  887. ctx.agents.register(child)
  888. throw new Error('session id already published')
  889. })
  890. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  891. const selection = await remote.selectModel(request({ sessionId, provider: 'p', model: 'm' }))
  892. expect(selection.ok).toBe(false)
  893. if (!selection.ok) {
  894. expect(selection.error).toMatchObject({
  895. code: 'session/agent-busy',
  896. details: { reason: 'use subagent delivery for this child session' },
  897. })
  898. }
  899. })
  900. })