1
0

request-reconstruction.spec.ts 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474
  1. /**
  2. * Loop-level reconstructability: every request the loop sends is a pure function of the
  3. * session log — messages derive at the step/start boundary and the header is the latest
  4. * request/header snapshot. Each request extends its predecessor unless a logged compaction
  5. * replacement or header change explains the difference.
  6. */
  7. import { describe, expect, it } from 'vitest'
  8. import { Context } from 'cordis'
  9. import LlmService, { LlmError, ReasoningEffortId } from '@deepseek-ai/dsh-llm'
  10. import type { GenerateOptions, LlmModelReasoningInfo, LlmResolvedModelInfo } from '@deepseek-ai/dsh-llm'
  11. import SessionStore, { Session, SessionId, foldRequestHeader } from '@deepseek-ai/dsh-session'
  12. import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
  13. import ToolRegistry, { defineContentToolFixture } from '@deepseek-ai/dsh-tools'
  14. import AgentRegistry, { type Agent } from '@deepseek-ai/dsh-agent'
  15. import AgentLoop from '@deepseek-ai/dsh-agent-loop'
  16. import { MockAdapter, textResponse, toolCallResponse } from './mock-adapter.ts'
  17. async function harness(adapter: MockAdapter, persona = 'stable base') {
  18. const ctx = new Context()
  19. await ctx.plugin(LlmService)
  20. await ctx.plugin(SessionStore)
  21. await ctx.plugin(SystemPrompt, { persona })
  22. await ctx.plugin(ToolRegistry)
  23. await ctx.plugin(AgentRegistry)
  24. await ctx.plugin(AgentLoop, { agents: [] })
  25. ctx.llm.registerAdapter(['mock'], adapter)
  26. return ctx
  27. }
  28. function waitForIdle(ctx: Context, agent: Agent): Promise<void> {
  29. return new Promise((resolve) => {
  30. const dispose = ctx.on('agent/status', (subject, status) => {
  31. if (subject === agent && status === 'idle') {
  32. dispose()
  33. resolve()
  34. }
  35. })
  36. })
  37. }
  38. function send(agent: Agent, text: string) {
  39. agent.followup([{ type: 'text', text }])
  40. }
  41. /** Assert `previous` is a strict value-prefix of `current`. */
  42. function expectPrefixExtension(previous: GenerateOptions, current: GenerateOptions) {
  43. expect(current.messages.length).toBeGreaterThan(previous.messages.length)
  44. expect(current.messages.slice(0, previous.messages.length)).toEqual([...previous.messages])
  45. expect(current.system).toEqual(previous.system)
  46. expect(current.tools).toEqual(previous.tools)
  47. }
  48. function registerEcho(ctx: Context) {
  49. ctx.tools.register(defineContentToolFixture({
  50. name: 'echo',
  51. description: 'echo back',
  52. parameters: { text: { type: 'string' } },
  53. async execute(args) {
  54. return [{ type: 'text', text: `echo: ${String(args.text)}` }]
  55. },
  56. }))
  57. }
  58. describe('request stability across the loop', () => {
  59. it('each step request within a turn append-extends the previous, frozen end to end', async () => {
  60. const adapter = new MockAdapter([
  61. toolCallResponse('c1', 'echo', { text: 'one' }, 'first'),
  62. toolCallResponse('c2', 'echo', { text: 'two' }, 'second'),
  63. textResponse('done'),
  64. ])
  65. const ctx = await harness(adapter)
  66. registerEcho(ctx)
  67. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  68. send(agent, 'go')
  69. await waitForIdle(ctx, agent)
  70. expect(adapter.requests).toHaveLength(3)
  71. expectPrefixExtension(adapter.requests[0]!, adapter.requests[1]!)
  72. expectPrefixExtension(adapter.requests[1]!, adapter.requests[2]!)
  73. for (const request of adapter.requests) {
  74. expect(Object.isFrozen(request)).toBe(true)
  75. expect(Object.isFrozen(request.messages)).toBe(true)
  76. }
  77. // One anchoring header snapshot; no further header events (nothing changed).
  78. const headerEvents = agent.session.events.filter(e => e.type === 'request/header')
  79. expect(headerEvents).toHaveLength(1)
  80. expect(headerEvents[0]?.type === 'request/header' && headerEvents[0].data.reason).toBe('initial')
  81. })
  82. it('a later turn append-extends the previous turn (one conversation, one log)', async () => {
  83. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  84. const ctx = await harness(adapter)
  85. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  86. send(agent, 'first')
  87. await waitForIdle(ctx, agent)
  88. send(agent, 'second')
  89. await waitForIdle(ctx, agent)
  90. expect(adapter.requests).toHaveLength(2)
  91. expectPrefixExtension(adapter.requests[0]!, adapter.requests[1]!)
  92. })
  93. it('logs adapter defaults, supports per-turn effort changes, and restores the effective value', async () => {
  94. const reasoning = {
  95. efforts: [
  96. { id: ReasoningEffortId('high'), name: 'High' },
  97. { id: ReasoningEffortId('max'), name: 'Max' },
  98. ],
  99. defaultEffort: ReasoningEffortId('high'),
  100. }
  101. const adapter = new MockAdapter([textResponse('one'), textResponse('two')], reasoning)
  102. const ctx = await harness(adapter)
  103. const agent = ctx.agentLoop.create(SessionId('effort'), { provider: 'mock', model: 'mock' })
  104. ctx.on('agent/request', async (_agent, turn, _step, _config, _signal, next) => {
  105. const config = await next()
  106. return turn === 2 ? { ...config, reasoningEffort: ReasoningEffortId('max') } : config
  107. })
  108. send(agent, 'first')
  109. await waitForIdle(ctx, agent)
  110. send(agent, 'second')
  111. await waitForIdle(ctx, agent)
  112. expect(adapter.requests.map(request => request.reasoningEffort)).toEqual([
  113. ReasoningEffortId('high'),
  114. ReasoningEffortId('max'),
  115. ])
  116. const headers = agent.session.events.filter(event => event.type === 'request/header')
  117. expect(headers.map(event => event.data.header.config.reasoningEffort)).toEqual([
  118. ReasoningEffortId('high'),
  119. ReasoningEffortId('max'),
  120. ])
  121. expect(headers.map(event => event.data.reason)).toEqual(['initial', 'change'])
  122. const resumedAdapter = new MockAdapter([textResponse('three')], reasoning)
  123. const resumedCtx = await harness(resumedAdapter)
  124. const resumedHandle = await resumedCtx.agents.create({
  125. sessionId: SessionId('effort-resumed'),
  126. seed: structuredClone(agent.session.events),
  127. agentOptions: { provider: 'mock', model: 'mock' },
  128. })
  129. send(resumedHandle.agent, 'third')
  130. await waitForIdle(resumedCtx, resumedHandle.agent)
  131. expect(resumedAdapter.requests[0]?.reasoningEffort).toBe(ReasoningEffortId('max'))
  132. const resumedHeaders = resumedHandle.agent.session.events.filter(event => event.type === 'request/header')
  133. expect(resumedHeaders.at(-1)?.data.header.config.reasoningEffort).toBe(ReasoningEffortId('max'))
  134. expect(resumedHeaders.at(-1)?.data.reason).toBe('resume')
  135. })
  136. it('keeps exact-model resolution, request logging, and dispatch on one adapter registration', async () => {
  137. const ctx = new Context()
  138. await ctx.plugin(LlmService)
  139. await ctx.plugin(SessionStore)
  140. await ctx.plugin(SystemPrompt, { persona: 'stable base' })
  141. await ctx.plugin(ToolRegistry)
  142. await ctx.plugin(AgentRegistry)
  143. await ctx.plugin(AgentLoop, { agents: [] })
  144. const started = Promise.withResolvers<undefined>()
  145. const reasoning = Promise.withResolvers<LlmModelReasoningInfo>()
  146. const first = new class extends MockAdapter {
  147. override async resolveModel(
  148. provider: string,
  149. model: string,
  150. _signal?: AbortSignal,
  151. ): Promise<LlmResolvedModelInfo> {
  152. started.resolve(undefined)
  153. return {
  154. provider,
  155. id: model,
  156. name: model,
  157. reasoning: await reasoning.promise,
  158. }
  159. }
  160. }([textResponse('first')])
  161. const second = new MockAdapter([textResponse('second')], {
  162. efforts: [{ id: ReasoningEffortId('max'), name: 'Max' }],
  163. defaultEffort: ReasoningEffortId('max'),
  164. })
  165. const disposeFirst = ctx.llm.registerAdapter(['mock'], first)
  166. const agent = ctx.agentLoop.create(SessionId('effort-hmr'), { provider: 'mock', model: 'mock' })
  167. send(agent, 'go')
  168. await started.promise
  169. disposeFirst()
  170. ctx.llm.registerAdapter(['mock'], second)
  171. reasoning.resolve({
  172. efforts: [{ id: ReasoningEffortId('high'), name: 'High' }],
  173. defaultEffort: ReasoningEffortId('high'),
  174. })
  175. await waitForIdle(ctx, agent)
  176. expect(first.requests.map(request => request.reasoningEffort)).toEqual([
  177. ReasoningEffortId('high'),
  178. ])
  179. expect(second.requests).toHaveLength(0)
  180. const headers = agent.session.events.filter(event => event.type === 'request/header')
  181. expect(headers.at(-1)?.data.header.config.reasoningEffort).toBe(ReasoningEffortId('high'))
  182. })
  183. it('aborts a blocked reasoning lookup before quiescent disposal completes', async () => {
  184. const started = Promise.withResolvers<AbortSignal>()
  185. const adapter = new class extends MockAdapter {
  186. override resolveModel(
  187. _provider: string,
  188. _model: string,
  189. signal?: AbortSignal,
  190. ): Promise<never> {
  191. if (signal === undefined) return Promise.reject(new Error('missing reasoning signal'))
  192. started.resolve(signal)
  193. return new Promise((_resolve, reject) => {
  194. if (signal.aborted) {
  195. reject(signal.reason instanceof Error ? signal.reason : new Error('reasoning aborted'))
  196. return
  197. }
  198. signal.addEventListener('abort', () => {
  199. reject(signal.reason instanceof Error ? signal.reason : new Error('reasoning aborted'))
  200. }, { once: true })
  201. })
  202. }
  203. }([])
  204. const ctx = await harness(adapter)
  205. const handle = await ctx.agents.create({
  206. sessionId: SessionId('reasoning-dispose'),
  207. agentOptions: { provider: 'mock', model: 'mock' },
  208. })
  209. send(handle.agent, 'go')
  210. const signal = await started.promise
  211. await handle.dispose()
  212. expect(signal.aborted).toBe(true)
  213. expect(handle.agent.status).toBe('disposed')
  214. expect(adapter.requests).toHaveLength(0)
  215. expect(handle.agent.session.events.some(event => event.type === 'request/header')).toBe(false)
  216. })
  217. it.each(['plain error', 'LLM error'] as const)(
  218. 'does not swallow a %s from exact-model resolution',
  219. async (kind) => {
  220. const failure = kind === 'plain error'
  221. ? new Error('reasoning metadata failed')
  222. : new LlmError('unsupported effort', 'UNSUPPORTED_REASONING_EFFORT')
  223. const adapter = new class extends MockAdapter {
  224. override resolveModel(): Promise<never> {
  225. return Promise.reject(failure)
  226. }
  227. }([])
  228. const ctx = await harness(adapter)
  229. const errors: Error[] = []
  230. ctx.on('agent/error', (_agent, _turn, _step, error) => void errors.push(error))
  231. const agent = ctx.agentLoop.create(SessionId(`reasoning-${kind}`), {
  232. provider: 'mock',
  233. model: 'mock',
  234. })
  235. send(agent, 'go')
  236. await waitForIdle(ctx, agent)
  237. expect(errors).toContain(failure)
  238. expect(adapter.requests).toHaveLength(0)
  239. },
  240. )
  241. it('a compaction replace rewrites the resend, and the log explains it', async () => {
  242. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  243. const ctx = await harness(adapter)
  244. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  245. send(agent, 'first')
  246. await waitForIdle(ctx, agent)
  247. // A pre-step listener compacts turn 1's history before turn 2's step —
  248. // the sanctioned surface rewrite, landing OUTSIDE the step.
  249. const preStep = ctx.on('agent/pre-step', () => {
  250. preStep()
  251. const session = agent.session
  252. const nodes = session.surface.nodes
  253. session.append('user/message', {
  254. content: [{ type: 'text', text: '[summary of turn 1]' }],
  255. source: { kind: 'plugin', plugin: 'test-compact' },
  256. }, {
  257. surfaceOp: { op: 'replace', start: nodes[0]!, end: nodes[1]! },
  258. sourceEventSeqs: [nodes[0]!, nodes[1]!],
  259. })
  260. })
  261. send(agent, 'second')
  262. await waitForIdle(ctx, agent)
  263. const second = adapter.requests[1]!
  264. // The rewritten history: summary replaces turn 1's user+assistant pair.
  265. expect(second.messages[0]!.content.some(b => b.type === 'text' && b.text.includes('[summary of turn 1]'))).toBe(true)
  266. // No header event beyond the anchor: the replace is itself in the log.
  267. expect(agent.session.events.filter(e => e.type === 'request/header')).toHaveLength(1)
  268. })
  269. it('a real system-prompt change is a full changed-header snapshot; a stable prompt logs nothing', async () => {
  270. const adapter = new MockAdapter([textResponse('one'), textResponse('two'), textResponse('three')])
  271. const ctx = await harness(adapter)
  272. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  273. send(agent, 'first')
  274. await waitForIdle(ctx, agent)
  275. send(agent, 'second')
  276. await waitForIdle(ctx, agent)
  277. // Identical assembly re-rendered per step is NOT a change.
  278. expect(agent.session.events.filter(e => e.type === 'request/header')).toHaveLength(1)
  279. ctx.systemPrompt.section({ name: 'extra', order: 2, text: 'new guidance' })
  280. send(agent, 'third')
  281. await waitForIdle(ctx, agent)
  282. const snapshots = agent.session.events.filter(e => e.type === 'request/header')
  283. expect(snapshots).toHaveLength(2)
  284. expect(snapshots[1]?.data.reason).toBe('change')
  285. expect(adapter.requests[2]!.system).toContain('new guidance')
  286. // History is preserved across the change — only the header moved.
  287. expect(adapter.requests[2]!.messages.length).toBeGreaterThan(adapter.requests[1]!.messages.length)
  288. })
  289. it('an inject() during the agent/request waterfall joins the NEXT request (the step/start boundary)', async () => {
  290. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  291. const ctx = await harness(adapter)
  292. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  293. let injected = false
  294. ctx.on('agent/request', async (_agent, _turn, _step, _config, _signal, next) => {
  295. if (!injected) {
  296. injected = true
  297. agent.inject([{ type: 'text', text: '[late context]' }], { source: { kind: 'plugin', plugin: 'test' } })
  298. }
  299. return next()
  300. })
  301. send(agent, 'first')
  302. await waitForIdle(ctx, agent)
  303. const first = adapter.requests[0]!
  304. // The inject landed in the log after the boundary: not in THIS request…
  305. expect(first.messages.some(m => m.content.some(b => b.type === 'text' && b.text.includes('[late context]')))).toBe(false)
  306. expect(agent.session.events.some(e => e.type === 'user/message' && e.data.source.kind === 'plugin')).toBe(true)
  307. send(agent, 'second')
  308. await waitForIdle(ctx, agent)
  309. // …but in the next one, at its logged position.
  310. const second = adapter.requests[1]!
  311. expect(second.messages.some(m => m.content.some(b => b.type === 'text' && b.text.includes('[late context]')))).toBe(true)
  312. })
  313. it('a mutation attempt on the frozen request content throws into the step (loud, not silent)', async () => {
  314. const adapter = new MockAdapter([textResponse('one')])
  315. const ctx = await harness(adapter)
  316. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  317. const errors: Error[] = []
  318. ctx.on('agent/error', (_agent, _turn, _step, error) => void errors.push(error))
  319. ctx.on('llm/stream', (options, next) => {
  320. // The historical failure mode this design kills: a listener rewriting
  321. // request content in place. The freeze turns it into a loud error.
  322. options.messages.push({ role: 'user', content: [{ type: 'text', text: 'sneaky' }] })
  323. return next()
  324. })
  325. send(agent, 'go')
  326. await waitForIdle(ctx, agent)
  327. expect(errors).toHaveLength(1)
  328. expect(errors[0]!.message).toMatch(/not extensible|frozen|read only|readonly/i)
  329. })
  330. it('a fresh loop instance over a seeded log anchors with a resume snapshot and stays cache-aligned', async () => {
  331. const adapter = new MockAdapter([textResponse('one')])
  332. const ctx = await harness(adapter)
  333. const agent = ctx.agentLoop.create(SessionId('gen1'), { provider: 'mock', model: 'mock' })
  334. send(agent, 'first')
  335. await waitForIdle(ctx, agent)
  336. // Second generation: a new agent whose session is seeded with the first
  337. // one's full log (the resume/fork path).
  338. const adapter2 = new MockAdapter([textResponse('two')])
  339. const ctx2 = await harness(adapter2)
  340. const handle = await ctx2.agents.create({
  341. sessionId: SessionId('gen2-session'),
  342. seed: [...agent.session.events],
  343. agentOptions: { provider: 'mock', model: 'mock' },
  344. })
  345. const agent2 = handle.agent
  346. send(agent2, 'second')
  347. await waitForIdle(ctx2, agent2)
  348. const snapshots = agent2.session.events.filter(e => e.type === 'request/header')
  349. expect(snapshots).toHaveLength(2)
  350. expect(snapshots[1]?.type === 'request/header' && snapshots[1].data.reason).toBe('resume')
  351. // Identical header across the restart: byte-identical continuation.
  352. expect(adapter2.requests[0]!.system).toEqual(adapter.requests[0]!.system)
  353. expectPrefixExtension(adapter.requests[0]!, adapter2.requests[0]!)
  354. })
  355. it('a delegating listener cannot mutate the seed through next() — the fold stays log-true', async () => {
  356. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  357. const ctx = await harness(adapter)
  358. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  359. ctx.on('agent/request', async (_agent, _turn, _step, _config, _signal, next) => {
  360. const config = await next()
  361. // next() resolves the SAME frozen seed — in-place shaping after
  362. // delegation is unrepresentable, so a "mutate what next() returned"
  363. // listener cannot desync the log from the request (nor reach the
  364. // session's cached header fold, which is deep-cloned away and itself
  365. // frozen).
  366. expect(Object.isFrozen(config)).toBe(true)
  367. expect(() => { (config as { temperature?: number }).temperature = 0.9 }).toThrow(TypeError)
  368. return config
  369. })
  370. send(agent, 'first')
  371. await waitForIdle(ctx, agent)
  372. send(agent, 'second')
  373. await waitForIdle(ctx, agent)
  374. // No changed snapshot was logged (nothing really changed), and the session's own
  375. // fold is immutable state.
  376. expect(agent.session.events.filter(e => e.type === 'request/header')).toHaveLength(1)
  377. expect(Object.isFrozen(agent.session.requestHeader())).toBe(true)
  378. expect(adapter.requests[1]!.temperature).toBeUndefined()
  379. })
  380. it('THEOREM: every request rebuilds byte-equal from the session log alone', async () => {
  381. const adapter = new MockAdapter([
  382. toolCallResponse('c1', 'echo', { text: 'one' }, 'calling'),
  383. textResponse('done'),
  384. textResponse('after change'),
  385. ])
  386. const ctx = await harness(adapter)
  387. registerEcho(ctx)
  388. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  389. send(agent, 'go')
  390. await waitForIdle(ctx, agent)
  391. ctx.systemPrompt.section({ name: 'extra', order: 2, text: 'now with guidance' })
  392. ctx.on('agent/request', async (_agent, _turn, _step, config, _signal, _next) => ({ ...config, temperature: 0.5, maxTokens: 99, stop: ['<END>'] }))
  393. send(agent, 'again')
  394. await waitForIdle(ctx, agent)
  395. expect(adapter.requests).toHaveLength(3)
  396. const events = agent.session.events
  397. const stepStarts = events.filter(e => e.type === 'step/start')
  398. expect(stepStarts).toHaveLength(3)
  399. adapter.requests.forEach((request, index) => {
  400. const stepStart = stepStarts[index]!
  401. // Messages: the derivation over the log prefix strictly before this
  402. // step's step/start — rebuilt here through a completely fresh Session.
  403. const rebuilt = new Session(SessionId(`rebuild-${index}`), structuredClone(events.slice(0, stepStart.seq)))
  404. expect(structuredClone(request.messages)).toEqual(rebuilt.deriveMessages())
  405. // Header: the latest request/header snapshot up to this step's dispatch
  406. // (its header event sits between step/start and the first chunk).
  407. const firstChunk = events.find(e => e.type === 'assistant/chunk' && e.seq > stepStart.seq)!
  408. const header = foldRequestHeader(events.slice(0, firstChunk.seq))!
  409. expect(request.model).toBe(header.config.model)
  410. expect(request.reasoningEffort).toBe(header.config.reasoningEffort)
  411. expect(request.system).toEqual(header.system)
  412. expect(structuredClone(request.tools ?? [])).toEqual(structuredClone(header.tools ?? []))
  413. expect(request.temperature).toBe(header.config.temperature)
  414. expect(request.maxTokens).toBe(header.config.maxTokens)
  415. expect(request.stop).toEqual(header.config.stop)
  416. })
  417. })
  418. })