api-proxy-cold.spec.ts 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489
  1. /**
  2. * Cold-session and degenerate-composition paths of the host ApiProxy:
  3. * metadata-only listing, Agent-free history reads, subagent ownership
  4. * isolation, and prompt failure mapping.
  5. */
  6. import { mkdtempSync, writeFileSync, utimesSync } from 'node:fs'
  7. import { tmpdir } from 'node:os'
  8. import { join } from 'node:path'
  9. import { describe, expect, it, vi } from 'vitest'
  10. import { Context } from 'cordis'
  11. import SessionStore from '@deepseek-ai/dsh-session'
  12. import AgentRegistry from '@deepseek-ai/dsh-agent'
  13. import { MessageId } from '@deepseek-ai/dsh-llm'
  14. import type { Agent } from '@deepseek-ai/dsh-agent'
  15. import UserInteractionService from '@deepseek-ai/dsh-user-interaction'
  16. import type { SessionEvent, SessionHeader, SessionId } from '@deepseek-ai/dsh-session'
  17. import {
  18. PersistenceCoordinator,
  19. SessionPersistenceRevision,
  20. type PersistenceBackend,
  21. type StoredPrefix,
  22. } from '@deepseek-ai/dsh-session-persistence'
  23. import type { RpcRequest } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  24. import { RpcId } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  25. import { createApiProxy } from '@deepseek-ai/dsh-host-apiproxy'
  26. const sid = (id: string): SessionId => id as SessionId
  27. let nextRpc = 1
  28. function request<P>(payload: P): RpcRequest<P> {
  29. return { rpcId: RpcId(`cold-${String(nextRpc++)}`), payload }
  30. }
  31. function header(id: string, createdAt: number, extra: Partial<SessionHeader> = {}): SessionHeader {
  32. return { version: 0, id: sid(id), createdAt, cwd: '/proj', ...extra }
  33. }
  34. describe('sessions.list cold merge', () => {
  35. it('summarizes unattached sessions: log mtime, locate-less and vanished-log createdAt fallbacks, lineage', async () => {
  36. const ctx = new Context()
  37. await ctx.plugin(SessionStore)
  38. await ctx.plugin(UserInteractionService)
  39. const root = mkdtempSync(join(tmpdir(), 'dsh-cold-'))
  40. const logPath = join(root, 'a.log')
  41. writeFileSync(logPath, 'log-bytes')
  42. utimesSync(logPath, 5000, 5000) // mtime 5_000_000 ms — newer than every createdAt below
  43. const metas = [
  44. header('session-a', 1000),
  45. header('session-b', 2000, { parentSession: sid('session-parent'), origin: 'subagent' }),
  46. header('session-c', 1500),
  47. ]
  48. // Structural fake of the persistence face list() consumes: list + locate.
  49. // locate: a real per-session file (mtime wins), a backend without one
  50. // (SQLite shape → createdAt), and a path whose file vanished (stat ENOENT
  51. // → createdAt).
  52. ctx.provide('sessionPersistence', {
  53. list: () => Promise.resolve(metas),
  54. locate: (meta: SessionHeader) => {
  55. if (meta.id === sid('session-a')) return { kind: 'jsonl', path: logPath }
  56. if (meta.id === sid('session-c')) return { kind: 'jsonl', path: join(root, 'vanished.log') }
  57. return undefined
  58. },
  59. })
  60. const api = createApiProxy(ctx, { defaultTarget: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp', workspaceRoot: '/tmp' })
  61. const response = await api.sessions.list(request({}))
  62. expect(response.result.ok).toBe(true)
  63. if (!response.result.ok) throw new Error('unreachable')
  64. const items = response.result.value.items
  65. expect(items.map(item => item.sessionId)).toEqual(['session-a', 'session-b', 'session-c'])
  66. const [a, b, c] = items
  67. expect(a?.updatedAt).toBeCloseTo(5_000_000, -3)
  68. expect(a?.running).toBe(false)
  69. // Cold summaries are never blank: lazy persistence keeps never-appended
  70. // sessions out of list(), so a listed session necessarily has events.
  71. expect(items.every(item => !item.blank)).toBe(true)
  72. expect(a?.cwd).toBe('/proj')
  73. expect(a?.parentSessionId).toBeUndefined()
  74. expect(b?.updatedAt).toBe(2000)
  75. expect(b?.parentSessionId).toBe('session-parent')
  76. expect(b?.origin).toBe('subagent')
  77. expect(c?.updatedAt).toBe(1500)
  78. })
  79. })
  80. describe('attached updatedAt excludes end-seed', () => {
  81. it('reports the last real work, not the pickup, so a resumed-untouched session does not float', async () => {
  82. const ctx = new Context()
  83. await ctx.plugin(SessionStore)
  84. await ctx.plugin(UserInteractionService)
  85. await ctx.plugin(AgentRegistry)
  86. const api = createApiProxy(ctx, { defaultTarget: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp', workspaceRoot: '/tmp' })
  87. // Old work, resumed just now: the log tail would report the pickup.
  88. const worked = 1_000_000
  89. const resumed = ctx.sessions.create(sid('resumed-untouched'), {
  90. seed: [
  91. { type: 'turn/start', seq: 0, time: worked, data: { turn: 1 } },
  92. { type: 'turn/end', seq: 1, time: worked, data: { turn: 1, reason: { kind: 'completed' } } },
  93. ],
  94. meta: { cwd: '/proj', createdAt: 500 },
  95. })
  96. ctx.agents.register({ id: resumed.id, session: resumed, status: 'idle', ctx } as Agent)
  97. const boundary = resumed.events.at(-1)
  98. expect(boundary?.type).toBe('session/end-seed')
  99. expect(boundary?.time).toBeGreaterThan(worked)
  100. const listed = await api.sessions.list(request({}))
  101. if (!listed.result.ok) throw new Error('list failed')
  102. const summary = listed.result.value.items.find(item => item.sessionId === 'resumed-untouched')
  103. expect(summary?.updatedAt).toBe(worked)
  104. // Real work appended after end-seed does move it.
  105. resumed.append('turn/start', { turn: 2 })
  106. const after = await api.sessions.list(request({}))
  107. if (!after.result.ok) throw new Error('list failed')
  108. const moved = after.result.value.items.find(item => item.sessionId === 'resumed-untouched')
  109. expect(moved?.updatedAt).toBeGreaterThan(worked)
  110. })
  111. })
  112. describe('cold history recovery view', () => {
  113. it('shows in-memory interruption repair without activating the session', async () => {
  114. const ctx = new Context()
  115. await ctx.plugin(SessionStore)
  116. await ctx.plugin(UserInteractionService)
  117. const sessionId = sid('session-interrupted')
  118. const meta = header(sessionId, 1000)
  119. const stored: StoredPrefix<never> = {
  120. meta,
  121. events: [{ type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } }],
  122. revision: SessionPersistenceRevision('history-recovery-test:1'),
  123. }
  124. const backend: PersistenceBackend<never> = {
  125. name: 'history-recovery-test',
  126. loadStored: id => Promise.resolve(id === sessionId ? structuredClone(stored) : undefined),
  127. readStoredRevision: id => Promise.resolve(
  128. id === sessionId ? SessionPersistenceRevision('history-recovery-test:1') : undefined,
  129. ),
  130. appendBatch: () => Promise.resolve(),
  131. commitRepair: () => Promise.resolve(),
  132. list: () => Promise.resolve([structuredClone(meta)]),
  133. }
  134. const coordinator = new PersistenceCoordinator(ctx, backend)
  135. ctx.provide('sessionPersistence', {
  136. list: (signal?: AbortSignal) => backend.list(signal),
  137. inspect: (id: SessionId, signal?: AbortSignal) => coordinator.inspect(id, signal),
  138. locate: () => undefined,
  139. } as never)
  140. const api = createApiProxy(ctx, { defaultTarget: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp', workspaceRoot: '/tmp' })
  141. const history = await api.sessions.history(request({ sessionId, beforeSeq: 2, maxMessages: 10 }))
  142. if (!history.result.ok) throw new Error('history failed')
  143. expect(history.result.value.events.map(entry => entry.event)).toMatchInlineSnapshot(`
  144. [
  145. {
  146. "data": {
  147. "turn": 1,
  148. },
  149. "seq": 0,
  150. "time": 1,
  151. "type": "turn/start",
  152. },
  153. {
  154. "data": {
  155. "reason": {
  156. "kind": "interrupted",
  157. },
  158. "turn": 1,
  159. },
  160. "seq": 1,
  161. "time": 1,
  162. "type": "turn/end",
  163. },
  164. ]
  165. `)
  166. expect(ctx.sessions.get(sessionId)).toBeUndefined()
  167. await ctx.fiber.dispose()
  168. })
  169. })
  170. describe('subagent ownership fence', () => {
  171. it('reads a cold child without an Agent and rejects generic resume or adoption', async () => {
  172. const ctx = new Context()
  173. await ctx.plugin(SessionStore)
  174. await ctx.plugin(AgentRegistry)
  175. await ctx.plugin(UserInteractionService)
  176. const sessionId = sid('session-child')
  177. const meta = header('session-child', 1000, {
  178. parentSession: sid('session-parent'),
  179. seedLength: 0,
  180. origin: 'subagent',
  181. })
  182. const events = [
  183. { type: 'turn/start', seq: 0, time: 1, data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } } },
  184. {
  185. type: 'user/message',
  186. seq: 1,
  187. time: 2,
  188. data: { content: [{ type: 'text', text: 'work' }], source: { kind: 'user' } },
  189. surfaceOp: 'append',
  190. },
  191. {
  192. type: 'subagent/descriptor',
  193. seq: 2,
  194. time: 3,
  195. data: { version: 2, mode: 'continuable', provider: 'spawn', label: 'child' },
  196. },
  197. { type: 'turn/end', seq: 3, time: 4, data: { turn: 1, reason: { kind: 'completed' } } },
  198. ] as SessionEvent[]
  199. const inspect = vi.fn(() => Promise.resolve({ meta, events }))
  200. ctx.provide('sessionPersistence', {
  201. list: () => Promise.resolve([meta]),
  202. inspect,
  203. locate: () => undefined,
  204. } as never)
  205. const resume = vi.spyOn(ctx.agents, 'resume')
  206. const api = createApiProxy(ctx, { defaultTarget: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp', workspaceRoot: '/tmp' })
  207. const history = await api.sessions.history(request({ sessionId }))
  208. expect(history.result.ok).toBe(true)
  209. if (history.result.ok) {
  210. expect(history.result.value.events.map(entry => entry.event.type)).toEqual(events.map(event => event.type))
  211. }
  212. expect(ctx.agents.get(sessionId)).toBeUndefined()
  213. const prompt = await api.sessions.prompt(request({
  214. sessionId,
  215. mode: 'queue',
  216. content: [{ type: 'text', text: 'follow up' }],
  217. }))
  218. expect(prompt.result.ok).toBe(false)
  219. if (!prompt.result.ok) {
  220. expect(prompt.result.error).toMatchObject({
  221. code: 'agent-busy',
  222. details: { reason: 'use subagent delivery for this child session' },
  223. })
  224. }
  225. const create = await api.sessions.create(request({ sessionId, cwd: '/proj' }))
  226. expect(create.result.ok).toBe(false)
  227. if (!create.result.ok) expect(create.result.error.code).toBe('agent-busy')
  228. expect(resume).not.toHaveBeenCalled()
  229. expect(ctx.agents.get(sessionId)).toBeUndefined()
  230. expect(inspect).toHaveBeenCalledTimes(3)
  231. })
  232. it('no longer treats a descriptor-only cold child without origin as subagent-owned', async () => {
  233. const ctx = new Context()
  234. await ctx.plugin(SessionStore)
  235. await ctx.plugin(AgentRegistry)
  236. await ctx.plugin(UserInteractionService)
  237. const sessionId = sid('session-legacy-child')
  238. const meta = header('session-legacy-child', 1000, {
  239. parentSession: sid('session-parent'),
  240. seedLength: 0,
  241. })
  242. const events = [
  243. {
  244. type: 'subagent/descriptor',
  245. seq: 0,
  246. time: 1,
  247. data: { version: 2, mode: 'continuable', provider: 'spawn', label: 'child' },
  248. },
  249. ] as SessionEvent[]
  250. ctx.provide('sessionPersistence', {
  251. list: () => Promise.resolve([meta]),
  252. inspect: () => Promise.resolve({ meta, events }),
  253. locate: () => undefined,
  254. } as never)
  255. // Pre-#1569 stores classify a child only through the descriptor event and
  256. // carry no header `origin`; the pre-release decision stops recognizing
  257. // them, so the ownership fence lets generic resume reach the registry
  258. // instead of answering `agent-busy`.
  259. const resume = vi.spyOn(ctx.agents, 'resume')
  260. .mockRejectedValue(new Error('registry unavailable in this bench'))
  261. const api = createApiProxy(ctx, { defaultTarget: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp', workspaceRoot: '/tmp' })
  262. const prompt = await api.sessions.prompt(request({
  263. sessionId,
  264. mode: 'queue',
  265. content: [{ type: 'text', text: 'follow up' }],
  266. }))
  267. expect(resume).toHaveBeenCalledTimes(1)
  268. expect(prompt.result.ok).toBe(false)
  269. if (!prompt.result.ok) expect(prompt.result.error.code).toBe('internal')
  270. })
  271. it('rejects origin-marked and runtime-owned live children from generic controls', async () => {
  272. const ctx = new Context()
  273. await ctx.plugin(SessionStore)
  274. await ctx.plugin(AgentRegistry)
  275. await ctx.plugin(UserInteractionService)
  276. const parentSession = ctx.sessions.create(sid('session-parent'), { meta: { cwd: '/proj' } })
  277. const parent = { id: parentSession.id, session: parentSession, status: 'idle', ctx } as Agent
  278. ctx.agents.register(parent)
  279. const originSession = ctx.sessions.create(sid('session-origin-child'), {
  280. meta: { cwd: '/proj', parentSession: parent.id, origin: 'subagent' },
  281. })
  282. const cancel = vi.fn()
  283. const updateInbox = vi.fn(() => 'applied' as const)
  284. const originChild = {
  285. id: originSession.id,
  286. session: originSession,
  287. status: 'idle',
  288. ctx,
  289. cancel,
  290. updateInbox,
  291. } as unknown as Agent
  292. ctx.agents.register(originChild)
  293. const startingSession = ctx.sessions.create(sid('session-starting-child'), {
  294. meta: { cwd: '/proj', parentSession: parent.id },
  295. })
  296. const startingChild = { id: startingSession.id, session: startingSession, status: 'idle', ctx } as Agent
  297. ctx.agents.enter(startingChild, parent)
  298. const api = createApiProxy(ctx, { defaultTarget: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp', workspaceRoot: '/tmp' })
  299. const stopped = await api.sessions.cancel(request({ sessionId: originChild.id }))
  300. expect(stopped.result.ok).toBe(false)
  301. if (!stopped.result.ok) expect(stopped.result.error.code).toBe('agent-busy')
  302. expect(cancel).not.toHaveBeenCalled()
  303. const queued = await api.sessions.updateQueue(request({
  304. sessionId: originChild.id,
  305. itemId: MessageId('queued-item'),
  306. action: { kind: 'remove' },
  307. }))
  308. expect(queued.result.ok).toBe(false)
  309. if (!queued.result.ok) expect(queued.result.error.code).toBe('agent-busy')
  310. expect(updateInbox).not.toHaveBeenCalled()
  311. const models = await api.sessions.models(request({ sessionId: startingChild.id }))
  312. expect(models.result.ok).toBe(false)
  313. if (!models.result.ok) expect(models.result.error.code).toBe('agent-busy')
  314. const create = await api.sessions.create(request({ sessionId: originChild.id, cwd: '/proj' }))
  315. expect(create.result.ok).toBe(false)
  316. if (!create.result.ok) expect(create.result.error.code).toBe('agent-busy')
  317. const history = await api.sessions.history(request({ sessionId: originChild.id }))
  318. expect(history.result.ok).toBe(true)
  319. expect(ctx.agents.get(originChild.id)).toBe(originChild)
  320. })
  321. it('does not classify an ordinary fork from an inherited ancestor descriptor', async () => {
  322. const ctx = new Context()
  323. await ctx.plugin(SessionStore)
  324. await ctx.plugin(AgentRegistry)
  325. await ctx.plugin(UserInteractionService)
  326. const session = ctx.sessions.create(sid('session-ordinary-fork'), {
  327. seed: [{
  328. type: 'subagent/descriptor',
  329. seq: 0,
  330. time: 1,
  331. data: { version: 2, mode: 'continuable', provider: 'spawn', label: 'ancestor' },
  332. }],
  333. meta: { cwd: '/proj', parentSession: sid('session-source'), seedLength: 1 },
  334. })
  335. const followup = vi.fn()
  336. const agent = { id: session.id, session, status: 'idle', ctx, followup } as unknown as Agent
  337. ctx.agents.register(agent)
  338. const api = createApiProxy(ctx, { defaultTarget: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp', workspaceRoot: '/tmp' })
  339. const response = await api.sessions.prompt(request({
  340. sessionId: agent.id,
  341. mode: 'queue',
  342. content: [{ type: 'text', text: 'ordinary work' }],
  343. }))
  344. expect(response.result.ok).toBe(true)
  345. expect(followup).toHaveBeenCalledOnce()
  346. })
  347. })
  348. describe('degenerate composition (no persistence, no factory)', () => {
  349. it('list skips the cold merge and history reports missing persistence as internal', async () => {
  350. const ctx = new Context()
  351. await ctx.plugin(SessionStore)
  352. await ctx.plugin(AgentRegistry)
  353. await ctx.plugin(UserInteractionService)
  354. const api = createApiProxy(ctx, { defaultTarget: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp', workspaceRoot: '/tmp' })
  355. const listed = await api.sessions.list(request({}))
  356. expect(listed.result.ok).toBe(true)
  357. if (listed.result.ok) expect(listed.result.value.items).toEqual([])
  358. // No persistence means cold history cannot inspect a transcript.
  359. const response = await api.sessions.history(request({ sessionId: sid('session-ghost') }))
  360. expect(response.result.ok).toBe(false)
  361. if (!response.result.ok) {
  362. expect(response.result.error.code).toBe('internal')
  363. expect(response.result.error.message).toMatch(/history unavailable for session "session-ghost"/)
  364. }
  365. })
  366. it('maps a persistence catalog miss to session-not-found without inspection', async () => {
  367. const ctx = new Context()
  368. await ctx.plugin(SessionStore)
  369. await ctx.plugin(AgentRegistry)
  370. await ctx.plugin(UserInteractionService)
  371. const inspect = vi.fn()
  372. ctx.provide('sessionPersistence', {
  373. list: () => Promise.resolve([]),
  374. inspect,
  375. } as never)
  376. const api = createApiProxy(ctx, { defaultTarget: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp', workspaceRoot: '/tmp' })
  377. const response = await api.sessions.history(request({ sessionId: sid('session-missing') }))
  378. expect(response.result.ok).toBe(false)
  379. if (!response.result.ok) expect(response.result.error.code).toBe('session-not-found')
  380. expect(inspect).not.toHaveBeenCalled()
  381. })
  382. })
  383. describe('sessions.prompt synchronous rejection', () => {
  384. it('maps a synchronous send throw (disposed/invalid input) to agent-busy with the reason attached', async () => {
  385. const ctx = new Context()
  386. await ctx.plugin(SessionStore)
  387. await ctx.plugin(AgentRegistry)
  388. await ctx.plugin(UserInteractionService)
  389. const session = ctx.sessions.create(sid('session-throwing'))
  390. // A live structural stub whose delivery verbs throw synchronously, the
  391. // shape a disposed loop presents at this seam.
  392. ctx.agents.register({
  393. id: session.id,
  394. session,
  395. status: 'idle',
  396. ctx,
  397. followup: () => { throw new Error('agent "session-throwing" lifecycle disposed') },
  398. steer: () => { throw new Error('agent "session-throwing" lifecycle disposed') },
  399. } as unknown as Agent)
  400. const api = createApiProxy(ctx, { defaultTarget: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp', workspaceRoot: '/tmp' })
  401. for (const mode of ['queue', 'steer'] as const) {
  402. const response = await api.sessions.prompt(request({
  403. sessionId: session.id, mode, content: [{ type: 'text' as const, text: 'x' }],
  404. }))
  405. expect(response.result.ok).toBe(false)
  406. if (!response.result.ok) {
  407. expect(response.result.error.code).toBe('agent-busy')
  408. expect(response.result.error.message).toBe('prompt rejected')
  409. expect(response.result.error.details).toEqual({
  410. reason: 'Error: agent "session-throwing" lifecycle disposed',
  411. })
  412. }
  413. }
  414. })
  415. it('classifies a raced cold-resume ID collision as agent-busy', async () => {
  416. const ctx = new Context()
  417. await ctx.plugin(SessionStore)
  418. await ctx.plugin(AgentRegistry)
  419. await ctx.plugin(UserInteractionService)
  420. const sessionId = sid('race-resume')
  421. const meta: SessionHeader = header('race-resume', 1000)
  422. ctx.provide('sessionPersistence', {
  423. list: () => Promise.resolve([meta]),
  424. inspect: () => Promise.resolve({ meta, events: [] as SessionEvent[] }),
  425. locate: () => undefined,
  426. } as never)
  427. // The raced winner: a live parent-owned subagent publishes the identity
  428. // while the generic cold resume is in flight, so the resume collides.
  429. const parentSession = ctx.sessions.create(sid('race-parent'), { meta: { cwd: '/proj' } })
  430. const parent = { id: parentSession.id, session: parentSession, status: 'idle', ctx } as Agent
  431. ctx.agents.register(parent)
  432. const childSession = ctx.sessions.create(sessionId, {
  433. meta: { cwd: '/proj', parentSession: parent.id, origin: 'subagent' },
  434. })
  435. const child = { id: sessionId, session: childSession, status: 'idle', ctx } as unknown as Agent
  436. vi.spyOn(ctx.agents, 'resume').mockImplementationOnce(async () => {
  437. // The parent's `enter()` wins the identity between the pre-resume
  438. // re-check and publication; the generic resume then collides.
  439. ctx.agents.register(child)
  440. throw new Error('session id already published')
  441. })
  442. const api = createApiProxy(ctx, { defaultTarget: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp', workspaceRoot: '/tmp' })
  443. const models = await api.sessions.models(request({ sessionId }))
  444. expect(models.result.ok).toBe(false)
  445. if (!models.result.ok) {
  446. expect(models.result.error).toMatchObject({
  447. code: 'agent-busy',
  448. details: { reason: 'use subagent delivery for this child session' },
  449. })
  450. }
  451. })
  452. })