session-cold.host.spec.ts 38 KB

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