compact-loop-repro.spec.ts 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365
  1. import { describe, expect, it } from 'vitest'
  2. import { Context } from 'cordis'
  3. import { toolPairingBalancedAfter, toolPairingBalancedBefore } from '@deepseek-ai/dsh-compact'
  4. import { CONTEXT_WINDOW_EXCEEDED_CODE, LlmError } from '@deepseek-ai/dsh-llm'
  5. import type { ContentBlock, GenerateOptions, StreamChunk } from '@deepseek-ai/dsh-llm'
  6. import { CallId, LlmAdapter } from '@deepseek-ai/dsh-llm'
  7. import { defineTool } from '@deepseek-ai/dsh-tools'
  8. import type { Agent } from '@deepseek-ai/dsh-agent'
  9. import AgentLoop from '@deepseek-ai/dsh-agent-loop'
  10. import { mountAgentLoopTestDependencies } from '@deepseek-ai/dsh-agent-loop-testkit'
  11. import * as Invariants from '@deepseek-ai/dsh-invariants'
  12. import { BasicCompactService } from '@deepseek-ai/dsh-compact-basic'
  13. import TokenMeterService from '@deepseek-ai/dsh-token-meter'
  14. import * as LlmRetry from '@deepseek-ai/dsh-llm-retry'
  15. import { SessionId, type SurfaceEvent } from '@deepseek-ai/dsh-session'
  16. /**
  17. * CBR-001 regression through the real loop. A replacement checkpoint has a high
  18. * log seq at the surface head and carries no tool pair, so both adjacent cuts
  19. * must be safe and re-compacting that checkpoint alone must succeed. This pins
  20. * surface-position semantics rather than raw-log scanning.
  21. */
  22. class ReproCompactService extends BasicCompactService {
  23. override async summarize(): Promise<{ summary: ContentBlock[]; provider: string; model: string }> {
  24. return {
  25. summary: [{ type: 'text', text: 'CHECKPOINT SUMMARY' }],
  26. provider: 'mock',
  27. model: 'stub',
  28. }
  29. }
  30. }
  31. /** Each call emits one tool-call until exhausted, then a final text answer. */
  32. class StepwiseToolAdapter extends LlmAdapter {
  33. calls = 0
  34. constructor(private toolSteps: number) {
  35. super()
  36. }
  37. override resolveModelContext(): Promise<{ contextWindow: number }> {
  38. return Promise.resolve({ contextWindow: 400 })
  39. }
  40. async * stream(_options: GenerateOptions): AsyncIterable<StreamChunk> {
  41. const n = this.calls
  42. this.calls += 1
  43. if (n < this.toolSteps) {
  44. const id = CallId(`c${n}`)
  45. const args = `{"i":${n}}`
  46. yield { type: 'block-start', index: 0, blockType: 'text' }
  47. yield { type: 'block-end', index: 0, block: { type: 'text', text: `step ${n}` } }
  48. yield { type: 'block-start', index: 1, blockType: 'tool-call' }
  49. yield { type: 'block-end', index: 1, block: { type: 'tool-call', id, name: 'work', arguments: args } }
  50. yield { type: 'finish', reason: { kind: 'tool-calls' } }
  51. return
  52. }
  53. yield { type: 'block-start', index: 0, blockType: 'text' }
  54. yield { type: 'block-end', index: 0, block: { type: 'text', text: 'all done' } }
  55. yield { type: 'finish', reason: { kind: 'stop' } }
  56. }
  57. }
  58. /** First conversation request overflows, then the rebuilt retry succeeds. */
  59. class OverflowRecoveryAdapter extends LlmAdapter {
  60. readonly conversationRequests: GenerateOptions[] = []
  61. readonly summaryRequests: GenerateOptions[] = []
  62. constructor(
  63. private readonly delivery: 'thrown' | 'in-band',
  64. private readonly transientAfterOverflow = false,
  65. ) {
  66. super()
  67. }
  68. override resolveModelContext(): Promise<{ contextWindow: number }> {
  69. return Promise.resolve({ contextWindow: 128 })
  70. }
  71. override async * stream(options: GenerateOptions): AsyncIterable<StreamChunk> {
  72. if (options.system?.includes('You are a compaction engine')) {
  73. this.summaryRequests.push(options)
  74. yield { type: 'block-start', index: 0, blockType: 'text' }
  75. yield { type: 'block-end', index: 0, block: { type: 'text', text: 'RECOVERY CHECKPOINT' } }
  76. yield { type: 'finish', reason: { kind: 'stop' } }
  77. return
  78. }
  79. this.conversationRequests.push(options)
  80. if (this.conversationRequests.length === 1) {
  81. if (this.delivery === 'thrown') {
  82. throw new LlmError('request too large for model context', CONTEXT_WINDOW_EXCEEDED_CODE)
  83. }
  84. yield {
  85. type: 'finish',
  86. reason: {
  87. kind: 'error',
  88. failure: {
  89. message: 'request too large for model context',
  90. code: CONTEXT_WINDOW_EXCEEDED_CODE,
  91. },
  92. },
  93. }
  94. return
  95. }
  96. if (this.transientAfterOverflow && this.conversationRequests.length === 2) {
  97. throw new LlmError('temporary provider outage', 'SERVER')
  98. }
  99. yield { type: 'block-start', index: 0, blockType: 'text' }
  100. yield { type: 'block-end', index: 0, block: { type: 'text', text: 'recovered' } }
  101. yield { type: 'finish', reason: { kind: 'stop' } }
  102. }
  103. }
  104. async function harness(toolSteps: number): Promise<{ ctx: Context; compact: ReproCompactService }> {
  105. const ctx = new Context()
  106. await mountAgentLoopTestDependencies(ctx)
  107. await ctx.plugin(Invariants)
  108. await ctx.plugin(AgentLoop, { agents: [] })
  109. await ctx.plugin(TokenMeterService)
  110. ctx.llm.registerAdapter(['mock'], new StepwiseToolAdapter(toolSteps))
  111. ctx.tools.register(defineTool({
  112. name: 'work',
  113. description: 'does work',
  114. parameters: { i: { type: 'number' } },
  115. async execute() {
  116. return [{ type: 'text', text: 'work result' }]
  117. },
  118. }))
  119. // Small window so several tool steps cross the threshold and compaction
  120. // fires within the runaway turn after enough history can shrink.
  121. const compact = new ReproCompactService(ctx, {
  122. auto: true,
  123. thresholdRatio: 0.5,
  124. retainTokens: 50,
  125. maxTokens: 8192,
  126. compactionRetries: 1,
  127. })
  128. return { ctx, compact }
  129. }
  130. function waitForIdle(ctx: Context, agent: Agent): Promise<void> {
  131. return new Promise((resolve) => {
  132. const dispose = ctx.on('agent/status', (subject, status) => {
  133. if (subject === agent && status === 'idle') {
  134. dispose()
  135. resolve()
  136. }
  137. })
  138. })
  139. }
  140. function seedOverflowHistory(agent: Agent): void {
  141. for (let turn = 1; turn <= 2; turn += 1) {
  142. const sentinel = turn === 1 ? 'OLD HISTORY SENTINEL' : 'RECENT HISTORY'
  143. agent.session.append('turn/start', {
  144. turn,
  145. trigger: { kind: 'message', source: { kind: 'user' } },
  146. })
  147. agent.session.append('user/message', {
  148. content: [{ type: 'text', text: `${sentinel} ${'old context '.repeat(200)}` }],
  149. source: { kind: 'user' },
  150. }, { surfaceOp: 'append' })
  151. agent.session.append('step/start', { turn, step: 1 })
  152. agent.session.append('assistant/message', {
  153. provenance: { provider: 'mock', model: 'mock' },
  154. turn,
  155. step: 1,
  156. content: [{ type: 'text', text: `historical response ${turn} ${'detail '.repeat(200)}` }],
  157. }, { surfaceOp: 'append' })
  158. agent.session.append('step/end', { turn, step: 1 })
  159. agent.session.append('turn/end', { turn, reason: { kind: 'completed' } })
  160. }
  161. }
  162. describe('CBR-001: a real-loop checkpoint is a valid boundary on both sides', () => {
  163. it('uses the model actually routed by agent/request for post-step pressure', async () => {
  164. const { ctx } = await harness(8)
  165. ctx.on('agent/request', async (_agent, _turn, _step, config) => ({ ...config, provider: 'mock', model: 'mock' }))
  166. try {
  167. const agent = ctx.agentLoop.create(SessionId('routed-pressure'), {
  168. provider: 'unconfigured-agent-fallback',
  169. model: 'unconfigured-agent-fallback',
  170. })
  171. agent.send([{ type: 'text', text: 'do a routed multi-step task' }])
  172. await waitForIdle(ctx, agent)
  173. expect(agent.session.requestHeader()?.config.model).toBe('mock')
  174. expect(agent.session.events.some(event => event.type === 'compact/summary')).toBe(true)
  175. expect(agent.session.events.at(-1)).toMatchObject({
  176. type: 'turn/end',
  177. data: { reason: { kind: 'completed' } },
  178. })
  179. } finally {
  180. await ctx.fiber.dispose()
  181. }
  182. })
  183. it('runs automatic pressure after the current tool result and before step/end', async () => {
  184. const { ctx } = await harness(8)
  185. try {
  186. const agent = ctx.agentLoop.create(SessionId('post-step-order'), { provider: 'mock', model: 'mock' })
  187. agent.send([{ type: 'text', text: 'do tool work' }])
  188. await waitForIdle(ctx, agent)
  189. const events = [...agent.session.events]
  190. const compactStart = events.find(event => event.type === 'compact/start')
  191. expect(compactStart).toBeDefined()
  192. const precedingResult = events.findLast(event =>
  193. event.type === 'tool/result' && event.seq < compactStart!.seq,
  194. )
  195. if (precedingResult?.type !== 'tool/result') throw new Error('expected a durable tool result before compaction')
  196. const stepEnd = events.find(event =>
  197. event.type === 'step/end'
  198. && event.data.step === precedingResult.data.step
  199. && event.seq > compactStart!.seq,
  200. )
  201. expect(precedingResult.seq).toBeLessThan(compactStart!.seq)
  202. expect(compactStart!.seq).toBeLessThan(stepEnd!.seq)
  203. } finally {
  204. await ctx.fiber.dispose()
  205. }
  206. })
  207. it('the head checkpoint the loop lands is a balanced cut on both sides', async () => {
  208. const { ctx } = await harness(8)
  209. try {
  210. const agent = ctx.agentLoop.create(SessionId('repro'), { provider: 'mock', model: 'mock' })
  211. agent.send([{ type: 'text', text: 'do a long multi-step task' }])
  212. await waitForIdle(ctx, agent)
  213. const events = [...agent.session.events]
  214. // A compaction ran: at least one checkpoint landed on the surface.
  215. const checkpoints = events.filter(
  216. (e): e is SurfaceEvent =>
  217. e.type === 'user/message'
  218. && typeof (e as SurfaceEvent).surfaceOp === 'object',
  219. )
  220. expect(checkpoints.length).toBeGreaterThan(0)
  221. // High log position does not make a text-only checkpoint mid-step; both
  222. // its start and end cuts are balanced in surface order.
  223. const nodes = agent.session.surface.nodes
  224. for (const cp of checkpoints) {
  225. const index = nodes.indexOf(cp.seq)
  226. if (index === -1) continue // shadowed by a later checkpoint — no longer an edge.
  227. expect(toolPairingBalancedBefore(agent.session, cp.seq),
  228. `checkpoint seq ${cp.seq} must be a balanced region START`).toBe(true)
  229. expect(toolPairingBalancedAfter(agent.session, cp.seq),
  230. `checkpoint seq ${cp.seq} must be a balanced region END`).toBe(true)
  231. }
  232. } finally {
  233. await ctx.fiber.dispose()
  234. }
  235. })
  236. })
  237. describe('context-overflow recovery across the real loop and compact-basic', () => {
  238. it.each(['thrown', 'in-band'] as const)(
  239. 'force-compacts a %s overflow between failed and retry steps',
  240. async (delivery) => {
  241. const ctx = new Context()
  242. const adapter = new OverflowRecoveryAdapter(delivery)
  243. await mountAgentLoopTestDependencies(ctx)
  244. await ctx.plugin(Invariants)
  245. await ctx.plugin(AgentLoop, { agents: [] })
  246. await ctx.plugin(TokenMeterService)
  247. ctx.llm.registerAdapter(['mock'], adapter)
  248. ctx.on('agent/request', async (_agent, _turn, _step, config) => ({ ...config, provider: 'mock', model: 'mock' }))
  249. await ctx.plugin(BasicCompactService, {
  250. thresholdRatio: 1,
  251. retainTokens: 100,
  252. maxTokens: 64,
  253. compactionRetries: 0,
  254. maxOverflowRetries: 1,
  255. })
  256. try {
  257. const agent = ctx.agentLoop.create(SessionId(`overflow-${delivery}`), {
  258. provider: 'unconfigured-agent-fallback',
  259. model: 'unconfigured-agent-fallback',
  260. })
  261. seedOverflowHistory(agent)
  262. agent.send([{ type: 'text', text: 'continue from history' }])
  263. await agent.whenIdle()
  264. expect(adapter.conversationRequests).toHaveLength(2)
  265. expect(adapter.summaryRequests).toHaveLength(1)
  266. expect(JSON.stringify(adapter.conversationRequests[0]!.messages)).toContain('OLD HISTORY SENTINEL')
  267. const retry = JSON.stringify(adapter.conversationRequests[1]!.messages)
  268. expect(retry).toContain('RECOVERY CHECKPOINT')
  269. expect(retry).not.toContain('OLD HISTORY SENTINEL')
  270. const events = [...agent.session.events]
  271. const failedEnd = events.find(event =>
  272. event.type === 'step/end' && event.data.turn === 3 && event.data.step === 1,
  273. )!
  274. const retryStart = events.find(event =>
  275. event.type === 'step/start' && event.data.turn === 3 && event.data.step === 2,
  276. )!
  277. const compaction = events.filter(event =>
  278. event.type === 'compact/start'
  279. || event.type === 'compact/summary'
  280. || event.type === 'compact/end',
  281. )
  282. expect(compaction.map(event => event.type)).toEqual([
  283. 'compact/start',
  284. 'compact/summary',
  285. 'compact/end',
  286. ])
  287. expect(compaction.every(event => event.seq > failedEnd.seq && event.seq < retryStart.seq)).toBe(true)
  288. expect(events.at(-1)).toMatchObject({
  289. type: 'turn/end',
  290. data: { reason: { kind: 'completed' } },
  291. })
  292. } finally {
  293. await ctx.fiber.dispose()
  294. }
  295. },
  296. )
  297. it('keeps context-overflow and transient retry budgets independent in one sequence', async () => {
  298. const ctx = new Context()
  299. const adapter = new OverflowRecoveryAdapter('thrown', true)
  300. await mountAgentLoopTestDependencies(ctx)
  301. await ctx.plugin(Invariants)
  302. await ctx.plugin(LlmRetry, {
  303. maxTransientRetries: 1,
  304. initialDelayMs: 1,
  305. maxDelayMs: 1,
  306. jitterRatio: 0,
  307. })
  308. await ctx.plugin(AgentLoop, { agents: [] })
  309. await ctx.plugin(TokenMeterService)
  310. ctx.llm.registerAdapter(['mock'], adapter)
  311. await ctx.plugin(BasicCompactService, {
  312. thresholdRatio: 1,
  313. retainTokens: 100,
  314. maxTokens: 64,
  315. compactionRetries: 0,
  316. maxOverflowRetries: 1,
  317. })
  318. try {
  319. const agent = ctx.agentLoop.create(SessionId('alternating-recovery'), { provider: 'mock', model: 'mock' })
  320. seedOverflowHistory(agent)
  321. agent.send([{ type: 'text', text: 'continue from history' }])
  322. await agent.whenIdle()
  323. expect(adapter.conversationRequests).toHaveLength(3)
  324. expect(adapter.summaryRequests).toHaveLength(1)
  325. expect(agent.session.events.filter(event => event.type === 'llm/retry').map(event => event.data))
  326. .toEqual([expect.objectContaining({ step: 2, retry: 1, failure: { message: 'temporary provider outage', code: 'SERVER' } })])
  327. expect(agent.session.events.filter(event => event.type === 'step/start').slice(-3).map(event => event.data.step))
  328. .toEqual([1, 2, 3])
  329. expect(agent.session.events.at(-1)).toMatchObject({
  330. type: 'turn/end',
  331. data: { reason: { kind: 'completed' } },
  332. })
  333. } finally {
  334. await ctx.fiber.dispose()
  335. }
  336. })
  337. })