agent.host.spec.ts 17 KB

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