agent.host.spec.ts 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452
  1. import { mkdtempSync, writeFileSync } from 'node:fs'
  2. import { tmpdir } from 'node:os'
  3. import { join } from 'node:path'
  4. import { Context } from '@deepseek-ai/cordis'
  5. import AgentRegistry from '@deepseek-ai/dsh-agent'
  6. import type { Agent } from '@deepseek-ai/dsh-agent'
  7. import { agentPresetProjectionDefinition } from '@deepseek-ai/dsh-agent-presets'
  8. import SessionStore, { SessionId, SessionLogOffset, SessionSeq } from '@deepseek-ai/dsh-session'
  9. import type { SessionEvent, SessionHeader } from '@deepseek-ai/dsh-session'
  10. import type { SessionObservation } from '@deepseek-ai/dsh-session-query'
  11. import type { SessionInspection } from '@deepseek-ai/dsh-session-persistence'
  12. import TypertRegistry from '@deepseek-ai/dsh-typert-registry'
  13. import { afterEach, describe, expect, it, vi } from 'vitest'
  14. import {
  15. ApiSessionAgentController,
  16. ApiSessionCwdConflict,
  17. ApiSessionNotFound,
  18. ApiSessionSubagentOwnership,
  19. inspectApiSession,
  20. } from '../src/agent.ts'
  21. import { installModelSelectionProjection } from '../src/model-selection-projection.ts'
  22. import { installSessionReadTestServices, testSessionPersistence } from './test-remote.ts'
  23. const roots: Context[] = []
  24. afterEach(async () => {
  25. await Promise.all(roots.splice(0).map(ctx => ctx.fiber.dispose()))
  26. })
  27. async function harness(): Promise<{ ctx: Context; agents: ApiSessionAgentController }> {
  28. const ctx = new Context()
  29. roots.push(ctx)
  30. await ctx.plugin(TypertRegistry)
  31. await ctx.plugin(SessionStore)
  32. await ctx.plugin(AgentRegistry)
  33. installSessionReadTestServices(ctx)
  34. ctx.sessionProjections.register(agentPresetProjectionDefinition)
  35. installModelSelectionProjection(ctx)
  36. ctx.provide('agentDefaultModel', {
  37. currentSelection: () => ({ provider: 'fixture', model: 'fixture-model' }),
  38. saveSelection: () => Promise.resolve(),
  39. } as never)
  40. return { ctx, agents: new ApiSessionAgentController(ctx) }
  41. }
  42. function header(id: string, cwd: string | null = '/workspace'): SessionHeader {
  43. return {
  44. version: 0,
  45. id: SessionId(id),
  46. createdAt: 1,
  47. isSeeded: false,
  48. ...(cwd === null ? {} : { cwd }),
  49. }
  50. }
  51. function unseededInspection(
  52. meta: SessionHeader,
  53. events: readonly SessionEvent[] = [],
  54. ): SessionInspection {
  55. if (meta.isSeeded) throw new Error('seeded inspection fixtures require an explicit inherited cut')
  56. return { meta, inheritedEventCount: SessionLogOffset(0), events }
  57. }
  58. function providePersistence(ctx: Context, persistence: Record<string, unknown>): () => void {
  59. return ctx.provide('sessionPersistence', testSessionPersistence(ctx, persistence) as never)
  60. }
  61. function agent(ctx: Context, meta: SessionHeader): Agent {
  62. const session = ctx.sessions.create(meta.id, { meta })
  63. return { id: meta.id, session, status: 'idle', ctx } as Agent
  64. }
  65. function unpublishedAgent(ctx: Context, meta: SessionHeader): Agent {
  66. return {
  67. id: meta.id,
  68. session: { id: meta.id, header: meta, events: [] },
  69. status: 'idle',
  70. ctx,
  71. } as unknown as Agent
  72. }
  73. describe('ApiSession identity failures', () => {
  74. it('describes cwd conflicts with and without a recorded cwd', () => {
  75. expect(new ApiSessionCwdConflict(SessionId('missing-cwd'), '/wanted', undefined).message)
  76. .toContain('records no cwd')
  77. expect(new ApiSessionCwdConflict(SessionId('wrong-cwd'), '/wanted', '/existing').message)
  78. .toContain('belongs to "/existing"')
  79. })
  80. it('maps absent and cwd-less point observations to not found', async () => {
  81. const ctx = new Context()
  82. roots.push(ctx)
  83. await ctx.plugin(SessionStore)
  84. installSessionReadTestServices(ctx)
  85. await expect(inspectApiSession(ctx, SessionId('missing')))
  86. .rejects.toBeInstanceOf(ApiSessionNotFound)
  87. const inspect = vi.fn(() => Promise.resolve(undefined))
  88. const disposeMissing = providePersistence(ctx, {
  89. list: () => Promise.resolve([]),
  90. inspect,
  91. })
  92. await expect(inspectApiSession(ctx, SessionId('missing'))).rejects.toBeInstanceOf(ApiSessionNotFound)
  93. expect(inspect).toHaveBeenCalledOnce()
  94. disposeMissing()
  95. const listed = header('cwd-less-catalog', null)
  96. const disposeListed = providePersistence(ctx, {
  97. list: () => Promise.resolve([listed]),
  98. inspect: () => Promise.resolve(unseededInspection(listed)),
  99. })
  100. await expect(inspectApiSession(ctx, listed.id)).rejects.toBeInstanceOf(ApiSessionNotFound)
  101. disposeListed()
  102. const catalog = header('cwd-less-inspect')
  103. const inspected = header('cwd-less-inspect', null)
  104. providePersistence(ctx, {
  105. list: () => Promise.resolve([catalog]),
  106. inspect: () => Promise.resolve(unseededInspection(inspected)),
  107. })
  108. await expect(inspectApiSession(ctx, catalog.id)).rejects.toBeInstanceOf(ApiSessionNotFound)
  109. })
  110. it('forwards an explicit inspection signal', async () => {
  111. const ctx = new Context()
  112. roots.push(ctx)
  113. await ctx.plugin(SessionStore)
  114. installSessionReadTestServices(ctx)
  115. const meta = header('signalled-inspection')
  116. const inspect = vi.fn(() => Promise.resolve(unseededInspection(meta)))
  117. providePersistence(ctx, { inspect })
  118. const signal = new AbortController().signal
  119. await expect(inspectApiSession(ctx, meta.id, signal)).resolves.toEqual(unseededInspection(meta))
  120. expect(inspect).toHaveBeenCalledWith(meta.id, signal)
  121. })
  122. })
  123. describe('ApiSession Agent lookup and recovery', () => {
  124. it('resumes directly from a retained observation and rejects an invalid observed header', async () => {
  125. const { ctx, agents } = await harness()
  126. const meta = header('observed-resume')
  127. const resumed = unpublishedAgent(ctx, meta)
  128. const resume = vi.spyOn(ctx.agents, 'resume').mockResolvedValue({
  129. agent: resumed,
  130. dispose: () => Promise.resolve(),
  131. })
  132. const observed = {
  133. source: 'prepared',
  134. header: meta,
  135. inheritedEventCount: SessionLogOffset(0),
  136. events: [],
  137. cursor: -1,
  138. projections: { asOfSeq: -1, values: {} },
  139. retain: vi.fn(),
  140. [Symbol.dispose]: vi.fn(),
  141. } as unknown as SessionObservation
  142. await expect(agents.resolveObservedAgent(observed)).resolves.toEqual({ agent: resumed })
  143. expect(resume).toHaveBeenCalledWith(expect.objectContaining({ resumeSessionId: meta.id }))
  144. const invalid = {
  145. ...observed,
  146. header: header('observed-without-cwd', null),
  147. } as SessionObservation
  148. await expect(agents.resolveObservedAgent(invalid)).resolves.toMatchObject({
  149. error: { code: 'session/not-found' },
  150. })
  151. })
  152. it('projects live Agent contexts and maps missing cold identities through Typert lookup failures', async () => {
  153. const { ctx } = await harness()
  154. const live = agent(ctx, header('live'))
  155. ctx.agents.register(live)
  156. providePersistence(ctx, {
  157. list: () => Promise.resolve([]),
  158. inspect: vi.fn(),
  159. })
  160. const host = ctx.typert.contexts.getHost('agent')
  161. if (host === undefined) throw new Error('Agent Context resolver was not registered')
  162. await expect(host.resolve(live.id)).resolves.toBe(live.ctx)
  163. await expect(host.resolve(SessionId('missing'))).rejects.toMatchObject({ code: 'session/not-found' })
  164. })
  165. it('returns raced ordinary Agents and ownership failures after resume throws', async () => {
  166. const ordinary = await harness()
  167. const ordinaryMeta = header('ordinary-race')
  168. providePersistence(ordinary.ctx, {
  169. list: () => Promise.resolve([ordinaryMeta]),
  170. inspect: () => Promise.resolve(unseededInspection(ordinaryMeta)),
  171. })
  172. const winner = agent(ordinary.ctx, ordinaryMeta)
  173. vi.spyOn(ordinary.ctx.agents, 'resume').mockImplementation(async () => {
  174. ordinary.ctx.agents.register(winner)
  175. throw new Error('raced publication')
  176. })
  177. await expect(ordinary.agents.resolveAgent(ordinaryMeta.id)).resolves.toEqual({ agent: winner })
  178. const child = await harness()
  179. const childMeta = header('child-race')
  180. providePersistence(child.ctx, {
  181. list: () => Promise.resolve([childMeta]),
  182. inspect: () => Promise.resolve(unseededInspection(childMeta)),
  183. })
  184. vi.spyOn(child.ctx.agents, 'resume').mockImplementation(async () => {
  185. child.ctx.sessions.create(childMeta.id, {
  186. meta: { ...childMeta, parentSession: SessionId('parent'), origin: 'subagent' },
  187. })
  188. throw new Error('raced child publication')
  189. })
  190. await expect(child.agents.resolveAgent(childMeta.id)).resolves.toMatchObject({
  191. error: { code: 'session/agent-busy' },
  192. })
  193. })
  194. it('reports not-found and ordinary resume failures without fabricating an Agent', async () => {
  195. const missing = await harness()
  196. providePersistence(missing.ctx, {
  197. list: () => Promise.resolve([]),
  198. inspect: vi.fn(),
  199. })
  200. await expect(missing.agents.resolveAgent(SessionId('missing'))).resolves.toMatchObject({
  201. error: { code: 'session/not-found' },
  202. })
  203. const failed = await harness()
  204. const meta = header('failed')
  205. providePersistence(failed.ctx, {
  206. list: () => Promise.resolve([meta]),
  207. inspect: () => Promise.resolve(unseededInspection(meta)),
  208. })
  209. vi.spyOn(failed.ctx.agents, 'resume').mockRejectedValue(new Error('factory unavailable'))
  210. await expect(failed.agents.resolveAgent(meta.id)).resolves.toMatchObject({
  211. error: { code: 'gateway/internal', message: expect.stringContaining('factory unavailable') as string },
  212. })
  213. })
  214. it('requires projected observations before activation', async () => {
  215. const { agents } = await harness()
  216. const meta = header('unprojected-observation')
  217. const observed = {
  218. source: 'prepared',
  219. header: meta,
  220. inheritedEventCount: SessionLogOffset(0),
  221. events: [],
  222. cursor: -1,
  223. retain: vi.fn(),
  224. [Symbol.dispose]: vi.fn(),
  225. } as unknown as SessionObservation
  226. expect(() => agents.presetForObservation(observed)).toThrow(
  227. 'Agent activation requires a projected Session observation',
  228. )
  229. })
  230. })
  231. describe('ApiSession model selection', () => {
  232. it('requires the model-selection projection', async () => {
  233. const { ctx, agents } = await harness()
  234. const live = agent(ctx, header('missing-model-projection'))
  235. vi.spyOn(ctx.sessionProjections, 'stateOf').mockReturnValue(undefined)
  236. expect(() => agents.selectionFor(live)).toThrow('required modelSelection projection')
  237. })
  238. it('reads a reasoning-free request and consumes only the exact pending selection', async () => {
  239. const { ctx, agents } = await harness()
  240. const logged = agent(ctx, header('logged-model'))
  241. logged.session.append('request/header', {
  242. header: { config: { provider: 'logged-provider', model: 'logged-model' } },
  243. reason: 'initial',
  244. })
  245. expect(agents.selectionFor(logged).current).toEqual({
  246. provider: 'logged-provider',
  247. model: 'logged-model',
  248. })
  249. const pending = agent(ctx, header('pending-model'))
  250. const selection = agents.selectionFor(pending)
  251. agents.selectForNextRequest(pending, {
  252. provider: 'selected-provider',
  253. model: 'selected-model',
  254. reasoningEffort: 'high' as never,
  255. })
  256. expect(selection.current).toMatchObject({
  257. provider: 'selected-provider', model: 'selected-model', reasoningEffort: 'high',
  258. })
  259. expect(agents.consumeSelection(pending, 'other-provider', 'selected-model', 'high')).toBe(false)
  260. expect(agents.consumeSelection(pending, 'selected-provider', 'other-model', 'high')).toBe(false)
  261. expect(agents.consumeSelection(pending, 'selected-provider', 'selected-model', 'low')).toBe(false)
  262. expect(agents.consumeSelection(pending, 'selected-provider', 'selected-model', 'high')).toBe(true)
  263. expect(selection.current).toEqual({ provider: 'fixture', model: 'fixture-model' })
  264. const untouched = agent(ctx, header('uninstalled-model'))
  265. expect(agents.consumeSelection(untouched, 'fixture', 'fixture-model', undefined)).toBe(false)
  266. })
  267. })
  268. describe('ApiSession create or adoption', () => {
  269. it('shares one in-flight creation between concurrent callers', async () => {
  270. const { ctx, agents } = await harness()
  271. const cwd = mkdtempSync(join(tmpdir(), 'dsh-session-controller-concurrent-'))
  272. const meta = header('concurrent-create', cwd)
  273. const created = unpublishedAgent(ctx, meta)
  274. let release!: () => void
  275. const gate = new Promise<void>((resolve) => { release = resolve })
  276. const create = vi.spyOn(ctx.agents, 'create').mockImplementation(async () => {
  277. await gate
  278. return { agent: created, dispose: () => Promise.resolve() }
  279. })
  280. const first = agents.ensureSession(meta.id, cwd, false)
  281. const second = agents.ensureSession(meta.id, cwd, false)
  282. release()
  283. await expect(Promise.all([first, second])).resolves.toEqual([created, created])
  284. expect(create).toHaveBeenCalledOnce()
  285. })
  286. it('accepts a raced ordinary creation and rejects a raced attached child', async () => {
  287. const ordinary = await harness()
  288. const cwd = mkdtempSync(join(tmpdir(), 'dsh-session-controller-create-'))
  289. const ordinaryMeta = header('create-race', cwd)
  290. const winner = agent(ordinary.ctx, ordinaryMeta)
  291. vi.spyOn(ordinary.ctx.agents, 'create').mockImplementation(async () => {
  292. ordinary.ctx.agents.register(winner)
  293. throw new Error('raced creation')
  294. })
  295. await expect(ordinary.agents.ensureSession(ordinaryMeta.id, cwd, false))
  296. .resolves.toBe(winner)
  297. const child = await harness()
  298. const childCwd = mkdtempSync(join(tmpdir(), 'dsh-session-controller-child-'))
  299. const childId = SessionId('create-child-race')
  300. vi.spyOn(child.ctx.agents, 'create').mockImplementation(async () => {
  301. child.ctx.sessions.create(childId, {
  302. meta: { cwd: childCwd, parentSession: SessionId('parent'), origin: 'subagent' },
  303. })
  304. throw new Error('raced child creation')
  305. })
  306. await expect(child.agents.ensureSession(childId, childCwd, false))
  307. .rejects.toBeInstanceOf(ApiSessionSubagentOwnership)
  308. })
  309. it('validates ownership and cwd on the Agent returned by creation', async () => {
  310. const child = await harness()
  311. const childCwd = mkdtempSync(join(tmpdir(), 'dsh-session-controller-returned-child-'))
  312. const childMeta = {
  313. ...header('returned-child', childCwd),
  314. parentSession: SessionId('parent'),
  315. origin: 'subagent' as const,
  316. }
  317. const childAgent = unpublishedAgent(child.ctx, childMeta)
  318. vi.spyOn(child.ctx.agents, 'create').mockResolvedValue({
  319. agent: childAgent,
  320. dispose: () => Promise.resolve(),
  321. })
  322. await expect(child.agents.ensureSession(childMeta.id, childCwd, false))
  323. .rejects.toBeInstanceOf(ApiSessionSubagentOwnership)
  324. const wrong = await harness()
  325. const requestedCwd = mkdtempSync(join(tmpdir(), 'dsh-session-controller-wrong-cwd-'))
  326. const wrongAgent = unpublishedAgent(wrong.ctx, header('wrong-returned-cwd', '/other'))
  327. vi.spyOn(wrong.ctx.agents, 'create').mockResolvedValue({
  328. agent: wrongAgent,
  329. dispose: () => Promise.resolve(),
  330. })
  331. await expect(wrong.agents.ensureSession(wrongAgent.id, requestedCwd, false))
  332. .rejects.toBeInstanceOf(ApiSessionCwdConflict)
  333. })
  334. it('resumes a matching persisted identity and preserves its selected preset', async () => {
  335. const { ctx, agents } = await harness()
  336. const meta = { ...header('stored'), agentPreset: 'minimal' }
  337. const events = [{
  338. type: 'agent-preset/selected',
  339. seq: SessionSeq(0),
  340. time: 1,
  341. data: { agentPreset: 'minimal' },
  342. }] as SessionEvent[]
  343. providePersistence(ctx, {
  344. list: () => Promise.resolve([meta]),
  345. inspect: () => Promise.resolve(unseededInspection(meta, events)),
  346. })
  347. ctx.provide('agentPresets', {
  348. resolve: (id?: string) => Promise.resolve({ id: id ?? 'minimal' }),
  349. mount: () => Promise.resolve(),
  350. } as never)
  351. const resumedSession = ctx.sessions.prepare(meta.id, {
  352. seed: structuredClone(events),
  353. meta: structuredClone(meta),
  354. inheritedEventCount: SessionLogOffset(0),
  355. seedSource: 'persistence',
  356. })
  357. const resumed = {
  358. id: meta.id,
  359. session: resumedSession,
  360. status: 'idle',
  361. ctx,
  362. } as unknown as Agent
  363. const resume = vi.spyOn(ctx.agents, 'resume').mockResolvedValue({
  364. agent: resumed,
  365. dispose: () => Promise.resolve(),
  366. })
  367. await expect(agents.ensureSession(meta.id, '/workspace', true, 'minimal')).resolves.toBe(resumed)
  368. expect(resume).toHaveBeenCalledWith(expect.objectContaining({ resumeSessionId: meta.id }))
  369. })
  370. it('rejects an ownership race before resume and a persisted cwd conflict', async () => {
  371. const child = await harness()
  372. const childMeta = header('resume-child-race')
  373. providePersistence(child.ctx, {
  374. list: () => Promise.resolve([childMeta]),
  375. inspect: () => Promise.resolve(unseededInspection(childMeta)),
  376. })
  377. child.ctx.provide('agentPresets', {
  378. resolve: () => {
  379. child.ctx.sessions.create(childMeta.id, {
  380. meta: { ...childMeta, parentSession: SessionId('parent'), origin: 'subagent' },
  381. })
  382. return Promise.resolve({ id: 'standard' })
  383. },
  384. mount: () => Promise.resolve(),
  385. } as never)
  386. await expect(child.agents.resolveAgent(childMeta.id)).resolves.toMatchObject({
  387. error: { code: 'session/agent-busy' },
  388. })
  389. const conflict = await harness()
  390. const stored = header('stored-cwd-conflict', '/stored')
  391. providePersistence(conflict.ctx, {
  392. list: () => Promise.resolve([stored]),
  393. inspect: () => Promise.resolve(unseededInspection(stored)),
  394. })
  395. await expect(conflict.agents.ensureSession(stored.id, '/requested', true))
  396. .rejects.toBeInstanceOf(ApiSessionCwdConflict)
  397. })
  398. it('surfaces directory creation failure and rejects setup without a scoped Agent', async () => {
  399. const { agents } = await harness()
  400. const parent = mkdtempSync(join(tmpdir(), 'dsh-session-controller-file-'))
  401. const file = join(parent, 'file')
  402. writeFileSync(file, 'not a directory')
  403. await expect(agents.ensureSession(SessionId('mkdir-failure'), join(file, 'child'), false))
  404. .rejects.toThrow('failed to ensure project directory')
  405. const composition = await agents.composeAgent(undefined)
  406. expect(() => composition.setup(new Context())).toThrow('Agent setup has no scoped Agent')
  407. })
  408. })