agent.host.spec.ts 17 KB

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