agent.host.spec.ts 18 KB

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