api-proxy-cold.spec.ts 35 KB

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