compact-loop-repro.spec.ts 12 KB

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