session-cold.host.spec.ts 36 KB

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