agent.host.spec.ts 19 KB

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