image-offload.spec.ts 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409
  1. /**
  2. * Image offload recovery: an adapter's `IMAGE_OFFLOAD_REQUIRED` failure
  3. * records exact image occurrences and retries without replacing messages.
  4. */
  5. import { afterEach, describe, expect, it } from 'vitest'
  6. import { Context } from '@deepseek-ai/cordis'
  7. import AgentLoop from '@deepseek-ai/dsh-agent-loop'
  8. import type { Agent } from '@deepseek-ai/dsh-agent'
  9. import BasicCompactionEngine from '@deepseek-ai/dsh-compaction-basic'
  10. import TokenMeter from '@deepseek-ai/dsh-token-meter'
  11. import { ImageVariantId } from '@deepseek-ai/dsh-attachment'
  12. import { serializeRequestWithImages } from '@deepseek-ai/dsh-llm-deepseek/src/protocols/chat-completions/serialize.ts'
  13. import { mountAgentLoopTestDependencies } from '@deepseek-ai/dsh-agent-loop-testkit'
  14. import { createAssistantMessage, createToolResultMessage, createUserMessage, IMAGE_OFFLOAD_REQUIRED_CODE, LlmAdapter, LlmError, ToolCallId } from '@deepseek-ai/dsh-llm'
  15. import type { ContentBlock, GenerateOptions, StreamChunk } from '@deepseek-ai/dsh-llm'
  16. import { isReplacementSurfaceEvent, SessionId } from '@deepseek-ai/dsh-session'
  17. import type { Session } from '@deepseek-ai/dsh-session'
  18. import * as offload from '../src/index.ts'
  19. type ScriptEntry = StreamChunk[] | (() => never)
  20. /** Replies one scripted stream per request and declares no retry policy. */
  21. class ScriptedAdapter extends LlmAdapter {
  22. readonly requests: GenerateOptions[] = []
  23. serializeSummary = false
  24. constructor(readonly script: ScriptEntry[]) {
  25. super()
  26. }
  27. async * stream(options: GenerateOptions): AsyncIterable<StreamChunk> {
  28. this.requests.push(options)
  29. if (this.serializeSummary && options.purpose === 'compaction') {
  30. await serializeRequestWithImages(options, {
  31. representation: { kind: 'base64' },
  32. requestImages: new Map([image('first').attachment].map(ref => [ref.attachmentId, {
  33. variantId: ImageVariantId(`sha256:${'b'.repeat(64)}`), attachment: ref,
  34. data: new Uint8Array(ref.bytes), mediaType: ref.mediaType, bytes: ref.bytes,
  35. width: ref.width, height: ref.height, depth: 'uchar', space: 'srgb', hasAlpha: false,
  36. }])),
  37. maxRequestImageBytes: 1,
  38. })
  39. }
  40. const entry = this.script.shift()
  41. if (entry === undefined) throw new Error('script exhausted')
  42. if (typeof entry === 'function') entry()
  43. else yield * entry
  44. }
  45. }
  46. function textResponse(text: string): StreamChunk[] {
  47. return [
  48. { type: 'block-start', index: 0, blockType: 'text' },
  49. { type: 'block-end', index: 0, block: { type: 'text', text } },
  50. { type: 'finish', reason: { kind: 'stop' } },
  51. ]
  52. }
  53. function offloadRequired(offloadImages: number): () => never {
  54. return () => {
  55. throw new LlmError('request images exceed the route budget', IMAGE_OFFLOAD_REQUIRED_CODE, { offloadImages })
  56. }
  57. }
  58. async function harness(adapter: ScriptedAdapter): Promise<Context> {
  59. const ctx = new Context()
  60. contexts.push(ctx)
  61. await mountAgentLoopTestDependencies(ctx)
  62. await ctx.plugin(offload)
  63. await ctx.plugin(AgentLoop, { agents: [] })
  64. ctx.llm.registerAdapter(['mock'], adapter)
  65. return ctx
  66. }
  67. const contexts: Context[] = []
  68. afterEach(async () => {
  69. await Promise.all(contexts.splice(0).map(ctx => ctx.fiber.dispose()))
  70. })
  71. function image(name: string): Extract<ContentBlock, { type: 'image' }> {
  72. return {
  73. type: 'image',
  74. attachment: { attachmentId: `sha256:${'a'.repeat(64)}` as never, name, mediaType: 'image/png', bytes: 1, width: 1, height: 1 },
  75. }
  76. }
  77. function offloadedNames(options: GenerateOptions): string[] {
  78. const names: string[] = []
  79. const visit = (blocks: readonly ContentBlock[]): void => {
  80. for (const block of blocks) {
  81. if (block.type === 'image' && block.offloaded === true) names.push(block.attachment.name ?? '')
  82. if (block.type === 'tool-result') visit(block.content)
  83. }
  84. }
  85. for (const message of options.messages) visit(message.content)
  86. return names
  87. }
  88. /** Surface replacements appended by the recovery, as `[original seq, replacement seq]` pairs. */
  89. function replacements(session: Session): [number, number][] {
  90. return session.snapshotEvents()
  91. .filter(isReplacementSurfaceEvent)
  92. .map(event => [Number(event.sourceEventSeqs?.[0]), Number(event.seq)])
  93. }
  94. function decisions(session: Session) {
  95. return session.snapshotEvents().filter(event => event.type === 'image/offload')
  96. }
  97. async function summaryHarness(script: ScriptEntry[]) {
  98. const adapter = new ScriptedAdapter(script)
  99. const ctx = await harness(adapter)
  100. await ctx.plugin(TokenMeter)
  101. const compact = new BasicCompactionEngine(ctx, { auto: false, summarizationProvider: 'mock', summarizationModel: 'summary' })
  102. const agent = await ctx.agentLoop.create(SessionId('summary-offload'), { provider: 'mock', model: 'mock' })
  103. return { ctx, compact, agent, adapter }
  104. }
  105. async function seedImages(agent: Agent, names: string[]) {
  106. agent.followup(createUserMessage({
  107. content: [{ type: 'text', text: 'conversation details '.repeat(400) }, ...names.map(image)],
  108. source: { kind: 'user' },
  109. }))
  110. await agent.whenIdle()
  111. const nodes = agent.session.surface.nodes
  112. return { start: nodes.at(-2)!, end: nodes.at(-1)! }
  113. }
  114. describe('summary image offload', () => {
  115. it.each([new Error('summary failed'), new LlmError('no count', IMAGE_OFFLOAD_REQUIRED_CODE)])('delegates unhandled summary errors: %s', async (error) => {
  116. const { ctx, agent } = await summaryHarness([])
  117. expect(ctx.waterfall('compaction/summary-error', { session: agent.session, sourceEventSeqs: [], error }, () => false)).toBe(false)
  118. expect(decisions(agent.session)).toEqual([])
  119. })
  120. it('recovers the real summary serializer with a tighter image budget and fresh pricing', async () => {
  121. const { compact, agent, adapter } = await summaryHarness([textResponse('answer'), textResponse('checkpoint')])
  122. const span = await seedImages(agent, ['first', 'second'])
  123. adapter.serializeSummary = true
  124. const result = await compact.compactNow(agent, new AbortController().signal)
  125. expect(result).not.toBeNull()
  126. expect(adapter.requests.map(offloadedNames)).toEqual([[], [], ['first', 'second']])
  127. expect(adapter.requests.slice(1).map(request => request.model)).toEqual(['summary', 'summary'])
  128. const events = agent.session.snapshotEvents()
  129. const types = events.map(event => event.type)
  130. expect(types.filter(type => type === 'compaction/start')).toHaveLength(1)
  131. expect(types.filter(type => type === 'compaction/end')).toHaveLength(1)
  132. const [decision] = decisions(agent.session)
  133. expect(decision?.data.targets).toEqual([{ seq: span.start, imageIndexes: [0, 1] }])
  134. expect(decision!.seq).toBeLessThan(types.indexOf('compaction/summary'))
  135. expect(events[span.start]).not.toHaveProperty('data.content.1.offloaded')
  136. expect(types).not.toContain('compaction/prune')
  137. expect(types).not.toContain('llm/retry')
  138. })
  139. it('offloads only the selected summary span and stops after exhausting its images', async () => {
  140. const { compact, agent, adapter } = await summaryHarness([
  141. textResponse('before'), textResponse('selected'), textResponse('after'),
  142. offloadRequired(1), offloadRequired(1), offloadRequired(1),
  143. ])
  144. await seedImages(agent, ['outside-before'])
  145. const selected = await seedImages(agent, ['first', 'second'])
  146. await seedImages(agent, ['outside-after'])
  147. agent.session.append('turn/start', { turn: 4 })
  148. await expect(compact.compactRegion(selected.start, selected.end, agent)).rejects.toMatchObject({ code: IMAGE_OFFLOAD_REQUIRED_CODE })
  149. expect(adapter.requests.slice(3).map(offloadedNames)).toEqual([[], ['first'], ['first', 'second']])
  150. expect(decisions(agent.session).map(event => event.data.targets)).toEqual([
  151. [{ seq: selected.start, imageIndexes: [0] }], [{ seq: selected.start, imageIndexes: [1] }],
  152. ])
  153. expect(agent.session.surface.replaceGeneration).toBe(0)
  154. const end = agent.session.snapshotEvents().at(-1)
  155. expect(end?.type).toBe('compaction/end')
  156. expect(end?.type === 'compaction/end' && typeof end.data.error).toBe('string')
  157. })
  158. it('preserves omission when a subsequent summary failure is terminal', async () => {
  159. const { compact, agent, adapter } = await summaryHarness([
  160. textResponse('answer'), offloadRequired(1), () => { throw new LlmError('provider outage', 'SERVER') },
  161. ])
  162. await seedImages(agent, ['first'])
  163. await expect(compact.compactNow(agent, new AbortController().signal)).rejects.toMatchObject({ code: 'summary', cause: { code: 'SERVER' } })
  164. expect(adapter.requests).toHaveLength(3)
  165. expect(decisions(agent.session)).toHaveLength(1)
  166. expect(agent.session.surface.replaceGeneration).toBe(0)
  167. })
  168. it('does not record a decision after cancellation during a failed summary', async () => {
  169. const controller = new AbortController()
  170. const reason = new Error('cancel summary')
  171. const { compact, agent, adapter } = await summaryHarness([
  172. textResponse('answer'), () => { controller.abort(reason); return offloadRequired(1)() },
  173. ])
  174. await seedImages(agent, ['first'])
  175. await expect(compact.compactNow(agent, controller.signal)).rejects.toBe(reason)
  176. expect(adapter.requests).toHaveLength(2)
  177. expect(decisions(agent.session)).toEqual([])
  178. })
  179. it('does not retry when cancellation follows a durable omission', async () => {
  180. const { ctx, compact, agent, adapter } = await summaryHarness([textResponse('answer'), offloadRequired(1)])
  181. await seedImages(agent, ['first'])
  182. const controller = new AbortController()
  183. const reason = new Error('cancel after offload')
  184. ctx.on('session/event', (_session, event) => { if (event.type === 'image/offload') controller.abort(reason) })
  185. await expect(compact.compactNow(agent, controller.signal)).rejects.toBe(reason)
  186. expect(adapter.requests).toHaveLength(2)
  187. expect(decisions(agent.session)).toHaveLength(1)
  188. })
  189. it('rejects a concurrently changed selection before recording omission', async () => {
  190. const { compact, agent, adapter } = await summaryHarness([textResponse('answer')])
  191. const selected = await seedImages(agent, ['first'])
  192. adapter.script.push(() => {
  193. agent.session.append('user/message', createUserMessage({ content: [{ type: 'text', text: 'replacement' }], source: { kind: 'user' } }), {
  194. surfaceOp: { op: 'replace', startSeq: selected.start, endSeq: selected.end },
  195. sourceEventSeqs: [selected.start, selected.end],
  196. })
  197. return offloadRequired(1)()
  198. })
  199. await expect(compact.compactNow(agent, new AbortController().signal)).rejects.toMatchObject({ code: 'changed' })
  200. expect(decisions(agent.session)).toEqual([])
  201. })
  202. })
  203. describe('compaction-image-offload', () => {
  204. it('logs one exact image selection and retries without replacing the message', async () => {
  205. const adapter = new ScriptedAdapter([offloadRequired(2), textResponse('sent')])
  206. const ctx = await harness(adapter)
  207. const agent = await ctx.agentLoop.create(SessionId('offload-required'), { provider: 'mock', model: 'mock' })
  208. const delegated: string[] = []
  209. ctx.on('agent/request-error', ({ failure }, next) => {
  210. delegated.push(failure.code)
  211. return next()
  212. })
  213. agent.followup(createUserMessage({
  214. content: [image('a'), { type: 'tool-result', toolCallId: ToolCallId('shot'), content: [image('b')] }, image('c')],
  215. source: { kind: 'user' },
  216. }))
  217. await agent.whenIdle()
  218. expect(adapter.requests).toHaveLength(2)
  219. expect(offloadedNames(adapter.requests[0]!)).toEqual([])
  220. expect(offloadedNames(adapter.requests[1]!)).toEqual(['a', 'b'])
  221. expect(delegated).toEqual([])
  222. const events = agent.session.snapshotEvents()
  223. const types = events.map(event => event.type)
  224. expect(types.filter(type => type === 'llm/retry')).toHaveLength(0)
  225. expect(types.filter(type => type === 'assistant/attempt')).toHaveLength(1)
  226. expect(replacements(agent.session)).toHaveLength(0)
  227. expect(types).not.toContain('compaction/prune')
  228. expect(decisions(agent.session)).toHaveLength(1)
  229. const decision = decisions(agent.session)[0]!
  230. const original = events.find(event => event.type === 'user/message')!.seq
  231. expect(decision.data).toEqual({ targets: [{ seq: original, imageIndexes: [0, 1] }] })
  232. expect(events[original]).toMatchObject({ type: 'user/message', surfaceOp: 'append' })
  233. expect(types.indexOf('assistant/attempt')).toBeLessThan(decision.seq)
  234. expect(decision.seq).toBeLessThan(types.indexOf('assistant/message'))
  235. expect(events[decision.seq + 1]).toMatchObject({ type: 'request/header', data: { reason: 'series' } })
  236. const durable = events[original]!
  237. expect(durable.type === 'user/message' ? durable.data.content[0] : undefined).not.toHaveProperty('offloaded')
  238. })
  239. it('counts the adapter prefix in request order after a surface replacement', async () => {
  240. const adapter = new ScriptedAdapter([offloadRequired(2), textResponse('sent')])
  241. const ctx = await harness(adapter)
  242. const agent = await ctx.agentLoop.create(SessionId('offload-reordered'), { provider: 'mock', model: 'mock' })
  243. // An empty-content assistant node derives no message and carries no occurrence.
  244. agent.session.append('assistant/message', {
  245. turn: 0,
  246. step: 0,
  247. message: createAssistantMessage({ content: [], source: { provider: 'mock', model: 'mock' } }),
  248. stream: [],
  249. }, { surfaceOp: 'append' })
  250. const first = agent.session.append('user/message', createUserMessage({
  251. content: [image('first')], source: { kind: 'user' },
  252. }), { surfaceOp: 'append' })
  253. agent.session.append('user/message', createUserMessage({
  254. content: [image('second')], source: { kind: 'user' },
  255. }), { surfaceOp: 'append' })
  256. agent.session.append('user/message', createUserMessage({
  257. content: [image('replacement')], source: { kind: 'user' },
  258. }), {
  259. surfaceOp: { op: 'replace', startSeq: first.seq, endSeq: first.seq },
  260. sourceEventSeqs: [first.seq],
  261. })
  262. agent.followup(createUserMessage({ content: [{ type: 'text', text: 'send' }], source: { kind: 'user' } }))
  263. await agent.whenIdle()
  264. expect(adapter.requests).toHaveLength(2)
  265. expect(offloadedNames(adapter.requests[1]!)).toEqual(['replacement', 'second'])
  266. expect(replacements(agent.session)).toHaveLength(1)
  267. expect(decisions(agent.session).map(event => event.data.targets)).toEqual([[
  268. { seq: 3, imageIndexes: [0] }, { seq: 2, imageIndexes: [0] },
  269. ]])
  270. })
  271. it('offloads a nested tool-result occurrence and leaves later images untouched', async () => {
  272. const adapter = new ScriptedAdapter([offloadRequired(1), textResponse('sent')])
  273. const ctx = await harness(adapter)
  274. const agent = await ctx.agentLoop.create(SessionId('offload-tool-result'), { provider: 'mock', model: 'mock' })
  275. const callId = ToolCallId('shot')
  276. agent.session.append('turn/start', { turn: 0 })
  277. agent.session.append('assistant/message', {
  278. turn: 0,
  279. step: 1,
  280. message: createAssistantMessage({
  281. content: [{ type: 'tool-call', id: callId, name: 'read_image', arguments: '{}' }],
  282. source: { provider: 'mock', model: 'mock' },
  283. }),
  284. stream: [],
  285. }, { surfaceOp: 'append' })
  286. agent.session.append('tool/call', { turn: 0, step: 1, callId, name: 'read_image', arguments: '{}' })
  287. const result = agent.session.append('tool/result', {
  288. turn: 0,
  289. step: 1,
  290. message: createToolResultMessage({
  291. callId,
  292. content: [
  293. { type: 'tool-result', toolCallId: ToolCallId('empty'), content: [{ type: 'text', text: 'no image' }] },
  294. image('first'),
  295. { type: 'tool-result', toolCallId: ToolCallId('inner'), content: [image('second')] },
  296. ],
  297. isError: false,
  298. }),
  299. }, { surfaceOp: 'append' })
  300. agent.session.append('turn/end', { turn: 0, reason: { kind: 'completed' } })
  301. agent.followup(createUserMessage({ content: [{ type: 'text', text: 'send' }], source: { kind: 'user' } }))
  302. await agent.whenIdle()
  303. expect(offloadedNames(adapter.requests[1]!)).toEqual(['first'])
  304. expect(replacements(agent.session)).toHaveLength(0)
  305. expect(decisions(agent.session)[0]?.data).toEqual({ targets: [{ seq: result.seq, imageIndexes: [0] }] })
  306. expect(agent.session.deriveEventMessage(result)?.source).toEqual(result.data.message.source)
  307. })
  308. it('advances across consecutive failures and preserves the first request snapshot', async () => {
  309. const adapter = new ScriptedAdapter([offloadRequired(1), offloadRequired(1), textResponse('sent')])
  310. const ctx = await harness(adapter)
  311. const agent = await ctx.agentLoop.create(SessionId('offload-repeat'), { provider: 'mock', model: 'mock' })
  312. agent.followup(createUserMessage({
  313. content: [image('first'), image('second'), image('third')], source: { kind: 'user' },
  314. }))
  315. await agent.whenIdle()
  316. expect(adapter.requests.map(offloadedNames)).toEqual([[], ['first'], ['first', 'second']])
  317. expect(decisions(agent.session).map(event => event.data.targets[0]?.imageIndexes)).toEqual([[0], [1]])
  318. expect(agent.session.surface.replaceGeneration).toBe(0)
  319. expect(agent.session.surface.contentGeneration).toBe(2)
  320. })
  321. it('removes the recovery listener when its plugin is disposed', async () => {
  322. const ctx = new Context()
  323. contexts.push(ctx)
  324. await mountAgentLoopTestDependencies(ctx)
  325. const fiber = ctx.plugin(offload)
  326. await fiber
  327. await fiber.dispose()
  328. const adapter = new ScriptedAdapter([offloadRequired(1)])
  329. ctx.llm.registerAdapter(['mock'], adapter)
  330. await ctx.plugin(AgentLoop, { agents: [] })
  331. const agent = await ctx.agentLoop.create(SessionId('offload-unloaded'), { provider: 'mock', model: 'mock' })
  332. agent.followup(createUserMessage({ content: [image('a')], source: { kind: 'user' } }))
  333. await agent.whenIdle()
  334. expect(adapter.requests).toHaveLength(1)
  335. expect(decisions(agent.session)).toEqual([])
  336. expect(ctx.waterfall('compaction/summary-error', {
  337. session: agent.session,
  338. sourceEventSeqs: agent.session.surface.nodes,
  339. error: new LlmError('summary budget exceeded', IMAGE_OFFLOAD_REQUIRED_CODE, { offloadImages: 1 }),
  340. }, () => false)).toBe(false)
  341. expect(decisions(agent.session)).toEqual([])
  342. })
  343. it('leaves every other failure to downstream recovery', async () => {
  344. const adapter = new ScriptedAdapter([() => {
  345. throw new LlmError('provider outage', 'SERVER')
  346. }])
  347. const ctx = await harness(adapter)
  348. const agent = await ctx.agentLoop.create(SessionId('offload-other-failure'), { provider: 'mock', model: 'mock' })
  349. const delegated: string[] = []
  350. ctx.on('agent/request-error', ({ failure }, next) => {
  351. delegated.push(failure.code)
  352. return next()
  353. })
  354. agent.followup(createUserMessage({ content: [image('a')], source: { kind: 'user' } }))
  355. await agent.whenIdle()
  356. expect(delegated).toEqual(['SERVER'])
  357. expect(replacements(agent.session)).toHaveLength(0)
  358. })
  359. it('delegates IMAGE_OFFLOAD_REQUIRED once nothing remains to offload', async () => {
  360. const adapter = new ScriptedAdapter([offloadRequired(1)])
  361. const ctx = await harness(adapter)
  362. const agent = await ctx.agentLoop.create(SessionId('offload-exhausted'), { provider: 'mock', model: 'mock' })
  363. const delegated: string[] = []
  364. ctx.on('agent/request-error', ({ failure }, next) => {
  365. delegated.push(failure.code)
  366. return next()
  367. })
  368. agent.followup(createUserMessage({ content: [{ type: 'text', text: 'no images' }], source: { kind: 'user' } }))
  369. await agent.whenIdle()
  370. expect(delegated).toEqual([IMAGE_OFFLOAD_REQUIRED_CODE])
  371. expect(replacements(agent.session)).toHaveLength(0)
  372. expect(agent.session.snapshotEvents().at(-1)).toMatchObject({ type: 'turn/end', data: { reason: { kind: 'error' } } })
  373. })
  374. })