request-recovery.spec.ts 24 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602
  1. import { describe, expect, it, vi } from 'vitest'
  2. import { Context } from 'cordis'
  3. import LlmService, {
  4. CallId,
  5. CONTEXT_WINDOW_EXCEEDED_CODE,
  6. HarnessError,
  7. LlmAdapter,
  8. LlmError,
  9. ProviderRequestId,
  10. } from '@deepseek-ai/dsh-llm'
  11. import type { GenerateOptions, LlmFailure, StreamChunk } from '@deepseek-ai/dsh-llm'
  12. import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
  13. import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
  14. import ToolRegistry, { defineTool } from '@deepseek-ai/dsh-tools'
  15. import type { PostToolDecision } from '@deepseek-ai/dsh-tools'
  16. import AgentRegistry from '@deepseek-ai/dsh-agent'
  17. import type { Agent } from '@deepseek-ai/dsh-agent'
  18. import AgentLoop from '@deepseek-ai/dsh-agent-loop'
  19. import { maxTokensResponse, textResponse, toolCallResponse } from './mock-adapter.ts'
  20. class FailureScriptAdapter extends LlmAdapter {
  21. requests: GenerateOptions[] = []
  22. constructor(private readonly entries: (Error | StreamChunk[])[]) {
  23. super()
  24. }
  25. async * stream(options: GenerateOptions): AsyncIterable<StreamChunk> {
  26. this.requests.push(options)
  27. const entry = this.entries.shift()
  28. if (entry === undefined) throw new Error('failure script exhausted')
  29. if (entry instanceof Error) throw entry
  30. yield* entry
  31. }
  32. }
  33. class IteratorConstructionFailureAdapter extends LlmAdapter {
  34. stream(_options: GenerateOptions): AsyncIterable<StreamChunk> {
  35. return {
  36. [Symbol.asyncIterator](): AsyncIterator<StreamChunk> {
  37. throw new LlmError('iterator construction failed', 'ITERATOR_CONSTRUCTION')
  38. },
  39. }
  40. }
  41. }
  42. class SynchronousDispatchFailureAdapter extends LlmAdapter {
  43. constructor(private readonly error: Error) {
  44. super()
  45. }
  46. stream(_options: GenerateOptions): AsyncIterable<StreamChunk> {
  47. throw this.error
  48. }
  49. }
  50. class IteratorResultGetterFailureAdapter extends LlmAdapter {
  51. constructor(
  52. private readonly field: 'done' | 'value',
  53. private readonly error: Error,
  54. ) {
  55. super()
  56. }
  57. stream(_options: GenerateOptions): AsyncIterable<StreamChunk> {
  58. const result = this.field === 'done' ? {} : { done: false }
  59. Object.defineProperty(result, this.field, { get: () => { throw this.error } })
  60. return {
  61. [Symbol.asyncIterator](): AsyncIterator<StreamChunk> {
  62. return { next: () => Promise.resolve(result as unknown as IteratorResult<StreamChunk>) }
  63. },
  64. }
  65. }
  66. }
  67. const streamListenerFailureCases: readonly [string, (ctx: Context) => void][] = [
  68. ['synchronous listener throw', (ctx) => {
  69. ctx.on('llm/stream', () => { throw new Error('synchronous stream listener failed') })
  70. }],
  71. ['invalid listener iterable', (ctx) => {
  72. ctx.on('llm/stream', () => ({}) as AsyncIterable<StreamChunk>)
  73. }],
  74. ['listener wrapper iteration failure', (ctx) => {
  75. ctx.on('llm/stream', (_options, next) => (async function * () {
  76. for await (const chunk of next()) {
  77. yield chunk
  78. throw new Error('stream listener wrapper failed')
  79. }
  80. })())
  81. }],
  82. ]
  83. async function harness(adapter?: LlmAdapter): Promise<Context> {
  84. const ctx = new Context()
  85. await ctx.plugin(LlmService)
  86. await ctx.plugin(SessionStore)
  87. await ctx.plugin(SystemPrompt)
  88. await ctx.plugin(ToolRegistry)
  89. await ctx.plugin(AgentRegistry)
  90. await ctx.plugin(AgentLoop, { agents: [] })
  91. if (adapter) ctx.llm.registerAdapter(['mock'], adapter)
  92. return ctx
  93. }
  94. function waitForIdle(ctx: Context, agent: Agent): Promise<void> {
  95. return new Promise((resolve) => {
  96. const dispose = ctx.on('agent/status', (subject, status) => {
  97. if (subject === agent && status === 'idle') {
  98. dispose()
  99. resolve()
  100. }
  101. })
  102. })
  103. }
  104. function send(agent: Agent): void {
  105. agent.send([{ type: 'text', text: 'go' }])
  106. }
  107. function contextError(message = 'context too large'): LlmError {
  108. return new LlmError(message, CONTEXT_WINDOW_EXCEEDED_CODE)
  109. }
  110. describe('agent post-step and request-error lifecycle', () => {
  111. it('fires post-step after results, buffered context, and steering but before step/end', async () => {
  112. const twoCalls: StreamChunk[] = [
  113. { type: 'block-start', index: 0, blockType: 'tool-call' },
  114. { type: 'block-end', index: 0, block: { type: 'tool-call', id: CallId('call-1'), name: 'work', arguments: '{}' } },
  115. { type: 'block-start', index: 1, blockType: 'tool-call' },
  116. { type: 'block-end', index: 1, block: { type: 'tool-call', id: CallId('call-2'), name: 'work', arguments: '{}' } },
  117. { type: 'usage', usage: { inputTokens: 10, outputTokens: 5 } },
  118. { type: 'finish', reason: { kind: 'tool-calls' } },
  119. ]
  120. const adapter = new FailureScriptAdapter([twoCalls, textResponse('done')])
  121. const ctx = await harness(adapter)
  122. ctx.tools.register(defineTool({
  123. name: 'work',
  124. description: 'do work',
  125. parameters: {},
  126. async execute(_args, exec) {
  127. if (exec.callId === CallId('call-2')) {
  128. exec.agent?.steer([{ type: 'text', text: 'steered' }], { source: { kind: 'plugin', plugin: 'test' } })
  129. }
  130. return [{ type: 'text', text: 'worked' }]
  131. },
  132. }))
  133. ctx.on('tools/post-execute', async (exec, _result): Promise<PostToolDecision> => ({
  134. kind: 'accept',
  135. additionalContexts: [{
  136. content: [{ type: 'text', text: `context for ${exec.callId}` }],
  137. source: { kind: 'plugin', plugin: 'test' },
  138. }],
  139. }))
  140. const agent = ctx.agentLoop.create(SessionId('post-step-order'), { provider: 'mock', model: 'mock' })
  141. const order: string[] = []
  142. ctx.on('session/event', (_session, event) => {
  143. if (
  144. event.type === 'assistant/message' || event.type === 'tool/call'
  145. || event.type === 'tool/result' || event.type === 'context/message'
  146. || event.type === 'steering/message' || event.type === 'step/end'
  147. ) {
  148. if (!('step' in event.data) || event.data.step === 1) order.push(event.type)
  149. }
  150. })
  151. ctx.on('agent/post-step', (subject, turn, step, signal) => {
  152. if (subject !== agent || step !== 1) return
  153. expect({ turn, step, aborted: signal.aborted }).toEqual({ turn: 1, step: 1, aborted: false })
  154. subject.inject([{ type: 'text', text: 'listener mutation' }], { source: { kind: 'plugin', plugin: 'post-step' } })
  155. order.push('agent/post-step')
  156. })
  157. send(agent)
  158. await waitForIdle(ctx, agent)
  159. expect(order).toEqual([
  160. 'assistant/message',
  161. 'tool/call',
  162. 'tool/result',
  163. 'tool/call',
  164. 'tool/result',
  165. 'context/message',
  166. 'context/message',
  167. 'steering/message',
  168. 'context/message',
  169. 'agent/post-step',
  170. 'step/end',
  171. ])
  172. })
  173. it('fires post-step for max-tokens and lets cancellation override that success', async () => {
  174. const adapter = new FailureScriptAdapter([maxTokensResponse('partial')])
  175. const ctx = await harness(adapter)
  176. const agent = ctx.agentLoop.create(SessionId('cancel-post-step-max-tokens'), { provider: 'mock', model: 'mock' })
  177. let entered!: () => void
  178. const postStepEntered = new Promise<void>((resolve) => { entered = resolve })
  179. ctx.on('agent/post-step', async (_agent, turn, step, signal) => {
  180. expect({ turn, step }).toEqual({ turn: 1, step: 1 })
  181. entered()
  182. await new Promise<void>((resolve) => {
  183. signal.addEventListener('abort', () => { resolve() }, { once: true })
  184. })
  185. })
  186. send(agent)
  187. const idle = waitForIdle(ctx, agent)
  188. await postStepEntered
  189. agent.cancel('cancelled during max-tokens post-step')
  190. await idle
  191. expect(agent.session.events.find(event => event.type === 'assistant/message')).toMatchObject({
  192. data: { usage: { inputTokens: 10, outputTokens: 7 } },
  193. })
  194. expect(agent.session.events.at(-1)).toMatchObject({
  195. type: 'turn/end',
  196. data: { reason: { kind: 'aborted', reason: 'cancelled during max-tokens post-step' } },
  197. })
  198. })
  199. it('closes the successful step as disposed when disposal lands during post-step', async () => {
  200. const adapter = new FailureScriptAdapter([
  201. toolCallResponse('dispose-call', 'work', {}),
  202. textResponse('must not continue'),
  203. ])
  204. const ctx = await harness(adapter)
  205. ctx.tools.register(defineTool({
  206. name: 'work',
  207. description: 'do work',
  208. parameters: {},
  209. async execute() { return [{ type: 'text', text: 'worked' }] },
  210. }))
  211. const agent = ctx.agentLoop.create(SessionId('dispose-post-step'), { provider: 'mock', model: 'mock' })
  212. let entered!: () => void
  213. const postStepEntered = new Promise<void>((resolve) => { entered = resolve })
  214. ctx.on('agent/post-step', async (_agent, turn, step, signal) => {
  215. expect({ turn, step }).toEqual({ turn: 1, step: 1 })
  216. entered()
  217. await new Promise<void>((resolve) => {
  218. signal.addEventListener('abort', () => { resolve() }, { once: true })
  219. })
  220. })
  221. send(agent)
  222. await postStepEntered
  223. await ctx.fiber.dispose()
  224. expect(adapter.requests).toHaveLength(1)
  225. const boundaries = agent.session.events.filter(event =>
  226. event.type === 'step/start' || event.type === 'step/end',
  227. )
  228. expect(boundaries.map(event => event.type)).toEqual(['step/start', 'step/end'])
  229. expect(boundaries.map(event => event.data)).toEqual([
  230. { turn: 1, step: 1 },
  231. { turn: 1, step: 1 },
  232. ])
  233. expect(agent.session.events.at(-1)).toMatchObject({
  234. type: 'turn/end',
  235. data: { reason: { kind: 'disposed' } },
  236. })
  237. })
  238. it.each([
  239. ['thrown', contextError()],
  240. ['in-band', [{ type: 'finish', reason: { kind: 'error', failure: { message: 'too large', code: CONTEXT_WINDOW_EXCEEDED_CODE, status: 400 } } }] satisfies StreamChunk[]],
  241. ] as const)('recovers a %s request failure in a new reconstructable step', async (_style, failure) => {
  242. const adapter = new FailureScriptAdapter([failure, textResponse('recovered')])
  243. const ctx = await harness(adapter)
  244. const agent = ctx.agentLoop.create(SessionId(`recover-${_style}`), { provider: 'mock', model: 'mock' })
  245. const attempts: number[] = []
  246. ctx.on('agent/request-error', async (subject, turn, step, error, facts, history) => {
  247. expect(subject).toBe(agent)
  248. expect({ turn, step, code: error.code }).toEqual({ turn: 1, step: 1, code: CONTEXT_WINDOW_EXCEEDED_CODE })
  249. expect(facts.code).toBe(CONTEXT_WINDOW_EXCEEDED_CODE)
  250. attempts.push(history.length)
  251. subject.session.append('context/message', {
  252. content: [{ type: 'text', text: 'RECOVERY SURFACE MUTATION' }],
  253. source: { kind: 'plugin', plugin: 'test-recovery' },
  254. }, { surfaceOp: 'append' })
  255. return { action: 'retry' }
  256. })
  257. send(agent)
  258. await waitForIdle(ctx, agent)
  259. expect(attempts).toEqual([0])
  260. expect(adapter.requests).toHaveLength(2)
  261. expect(JSON.stringify(adapter.requests[1]!.messages)).toContain('RECOVERY SURFACE MUTATION')
  262. const starts = agent.session.events.filter(event => event.type === 'step/start')
  263. const ends = agent.session.events.filter(event => event.type === 'step/end')
  264. expect(starts.map(event => event.data.step)).toEqual([1, 2])
  265. expect(ends.map(event => event.data.step)).toEqual([1, 2])
  266. const recovery = agent.session.events.find(event => event.type === 'context/message')!
  267. expect(ends[0]!.seq).toBeLessThan(recovery.seq)
  268. expect(recovery.seq).toBeLessThan(starts[1]!.seq)
  269. })
  270. it.each(streamListenerFailureCases)('does not offer %s to request recovery', async (_name, install) => {
  271. const ctx = await harness(new FailureScriptAdapter([textResponse('unused')]))
  272. const agent = ctx.agentLoop.create(SessionId(`stream-plugin-${_name.replaceAll(' ', '-')}`), { provider: 'mock', model: 'mock' })
  273. let recoveries = 0
  274. install(ctx)
  275. ctx.on('agent/request-error', async (_agent, _turn, _step, _error, _failure, _history, _signal, next) => {
  276. recoveries += 1
  277. return next()
  278. })
  279. send(agent)
  280. await waitForIdle(ctx, agent)
  281. expect(recoveries).toBe(0)
  282. expect(agent.session.events.at(-1)).toMatchObject({ type: 'turn/end', data: { reason: { kind: 'error' } } })
  283. })
  284. it('does not offer a nested model-call failure as the outer request failure', async () => {
  285. const outer = new FailureScriptAdapter([textResponse('outer adapter must not run')])
  286. const nested = new FailureScriptAdapter([contextError('nested overflow')])
  287. const ctx = await harness(outer)
  288. ctx.llm.registerAdapter(['nested'], nested)
  289. ctx.on('llm/stream', (options, next) => {
  290. if (options.provider !== 'mock') return next()
  291. return (async function* () {
  292. yield* ctx.llm.stream({
  293. provider: 'nested',
  294. model: 'nested',
  295. messages: [],
  296. ...options.signal === undefined ? {} : { signal: options.signal },
  297. })
  298. yield* next()
  299. })()
  300. })
  301. const agent = ctx.agentLoop.create(SessionId('nested-stream-not-recoverable'), { provider: 'mock', model: 'mock' })
  302. let recoveries = 0
  303. ctx.on('agent/request-error', async (_agent, _turn, _step, _error, _failure, _history, _signal, next) => {
  304. recoveries += 1
  305. return next()
  306. })
  307. send(agent)
  308. await waitForIdle(ctx, agent)
  309. expect(nested.requests).toHaveLength(1)
  310. expect(outer.requests).toHaveLength(0)
  311. expect(recoveries).toBe(0)
  312. expect(agent.session.events.at(-1)).toMatchObject({
  313. type: 'turn/end',
  314. data: { reason: { kind: 'error', message: 'nested overflow', code: CONTEXT_WINDOW_EXCEEDED_CODE } },
  315. })
  316. })
  317. it.each(['prompt-submit', 'prompt-assembly', 'pre-step', 'request'] as const)(
  318. 'does not offer %s middleware failures to request recovery',
  319. async (boundary) => {
  320. const adapter = new FailureScriptAdapter([textResponse('unused')])
  321. const ctx = await harness(adapter)
  322. if (boundary === 'prompt-submit') {
  323. ctx.on('agent/prompt-submit', () => { throw new Error('prompt submit failed') })
  324. } else if (boundary === 'prompt-assembly') {
  325. ctx.on('system-prompt/assemble', () => { throw new Error('prompt assembly failed') })
  326. } else if (boundary === 'pre-step') {
  327. ctx.on('agent/pre-step', () => { throw new Error('pre-step failed') })
  328. } else {
  329. ctx.on('agent/request', () => { throw new Error('request middleware failed') })
  330. }
  331. const agent = ctx.agentLoop.create(SessionId(`${boundary}-not-recoverable`), { provider: 'mock', model: 'mock' })
  332. let recoveries = 0
  333. ctx.on('agent/request-error', async (_agent, _turn, _step, _error, _failure, _history, _signal, next) => {
  334. recoveries += 1
  335. return next()
  336. })
  337. send(agent)
  338. await waitForIdle(ctx, agent)
  339. expect(recoveries).toBe(0)
  340. expect(adapter.requests).toHaveLength(0)
  341. expect(agent.session.events.at(-1)).toMatchObject({ type: 'turn/end', data: { reason: { kind: 'error' } } })
  342. },
  343. )
  344. it('does not offer result, tool, or post-step plugin failures to request recovery', async () => {
  345. for (const failure of ['result', 'tool', 'post-step'] as const) {
  346. const adapter = new FailureScriptAdapter([
  347. failure === 'tool' ? toolCallResponse(`call-${failure}`, 'work', {}) : textResponse('done'),
  348. ...(failure === 'tool' ? [textResponse('done')] : []),
  349. ])
  350. const ctx = await harness(adapter)
  351. if (failure === 'result') ctx.on('agent/step-result', () => { throw new Error('result failed') })
  352. if (failure === 'post-step') ctx.on('agent/post-step', () => { throw new Error('post-step failed') })
  353. if (failure === 'tool') {
  354. vi.spyOn(ctx.tools, 'execute').mockRejectedValue(new Error('tool service failed'))
  355. }
  356. const agent = ctx.agentLoop.create(SessionId(`${failure}-not-recoverable`), { provider: 'mock', model: 'mock' })
  357. let recoveries = 0
  358. ctx.on('agent/request-error', async (_agent, _turn, _step, _error, _failure, _history, _signal, next) => {
  359. recoveries += 1
  360. return next()
  361. })
  362. send(agent)
  363. await waitForIdle(ctx, agent)
  364. expect(recoveries, failure).toBe(0)
  365. }
  366. })
  367. it.each([
  368. ['synchronous dispatch', (error: Error) => new SynchronousDispatchFailureAdapter(error)],
  369. ['done getter', (error: Error) => new IteratorResultGetterFailureAdapter('done', error)],
  370. ['value getter', (error: Error) => new IteratorResultGetterFailureAdapter('value', error)],
  371. ] as const)('preserves original Error identity for adapter %s', async (_name, makeAdapter) => {
  372. const original = contextError(`${_name} overflow`)
  373. const ctx = await harness(makeAdapter(original))
  374. const agent = ctx.agentLoop.create(SessionId(`identity-${_name.replaceAll(' ', '-')}`), { provider: 'mock', model: 'mock' })
  375. let seen: Error | undefined
  376. ctx.on('agent/request-error', async (_agent, _turn, _step, error, _failure, _history, _signal, next) => {
  377. seen = error
  378. return next()
  379. })
  380. send(agent)
  381. await waitForIdle(ctx, agent)
  382. expect(seen).toBe(original)
  383. })
  384. it('keeps an adapter error with a hostile message accessor on the recovery path', async () => {
  385. const original = Object.defineProperty(new HarnessError('provider failed', 'SERVER'), 'message', {
  386. get() { throw new Error('SDK message accessor trap') },
  387. })
  388. const ctx = await harness(new SynchronousDispatchFailureAdapter(original))
  389. const agent = ctx.agentLoop.create(SessionId('hostile-message-recovery'), { provider: 'mock', model: 'mock' })
  390. let seenError: Error | undefined
  391. let seenFailure: LlmFailure | undefined
  392. ctx.on('agent/request-error', async (_agent, _turn, _step, error, failure, _history, _signal, next) => {
  393. seenError = error
  394. seenFailure = failure
  395. return next()
  396. })
  397. send(agent)
  398. await waitForIdle(ctx, agent)
  399. expect(seenError).toBe(original)
  400. expect(seenFailure).toEqual({ message: 'LLM adapter failed', code: 'SERVER' })
  401. expect(agent.session.events.at(-1)).toMatchObject({
  402. type: 'turn/end',
  403. data: { reason: { kind: 'error', failure: { message: 'LLM adapter failed', code: 'SERVER' } } },
  404. })
  405. })
  406. it('passes structured facts beside the original Error and records its cause chain on exhaustion', async () => {
  407. const original = new LlmError('provider busy', 'RATE_LIMIT', {
  408. cause: new Error('upstream connection reset'),
  409. status: 429,
  410. providerRetryAfterMs: 2_000,
  411. requestId: ProviderRequestId('req-9'),
  412. })
  413. Object.freeze(original)
  414. const ctx = await harness(new SynchronousDispatchFailureAdapter(original))
  415. const agent = ctx.agentLoop.create(SessionId('structured-request-failure'), { provider: 'mock', model: 'mock' })
  416. let seenError: Error | undefined
  417. let seenFailure: LlmFailure | undefined
  418. let seenHistory: readonly LlmFailure[] | undefined
  419. ctx.on('agent/request-error', async (
  420. _agent, _turn, _step, error, failure, history, _signal, next,
  421. ) => {
  422. seenError = error
  423. seenFailure = failure
  424. seenHistory = history
  425. return next()
  426. })
  427. send(agent)
  428. await waitForIdle(ctx, agent)
  429. expect(seenError).toBe(original)
  430. expect(seenFailure).toEqual({
  431. message: 'provider busy',
  432. code: 'RATE_LIMIT',
  433. status: 429,
  434. providerRetryAfterMs: 2_000,
  435. requestId: ProviderRequestId('req-9'),
  436. })
  437. expect(seenHistory).toEqual([])
  438. expect(Object.isFrozen(seenHistory)).toBe(true)
  439. expect(agent.session.events.at(-1)).toMatchObject({
  440. type: 'turn/end',
  441. data: {
  442. reason: {
  443. kind: 'error',
  444. step: 1,
  445. failure: {
  446. message: 'provider busy: upstream connection reset',
  447. code: 'RATE_LIMIT',
  448. status: 429,
  449. providerRetryAfterMs: 2_000,
  450. requestId: ProviderRequestId('req-9'),
  451. },
  452. },
  453. },
  454. })
  455. })
  456. it('classifies iterator construction and explicit NO_ADAPTER as model-request failures', async () => {
  457. for (const scenario of ['iterator', 'no-adapter'] as const) {
  458. const ctx = scenario === 'iterator' ? await harness(new IteratorConstructionFailureAdapter()) : await harness()
  459. const agent = ctx.agentLoop.create(SessionId(`request-boundary-${scenario}`), { provider: 'mock', model: 'mock' })
  460. let seen = ''
  461. ctx.on('agent/request-error', async (_agent, _turn, _step, error, _failure, _history, _signal, next) => {
  462. seen = error.code ?? ''
  463. return next()
  464. })
  465. send(agent)
  466. await waitForIdle(ctx, agent)
  467. expect(seen).toBe(scenario === 'iterator' ? 'ITERATOR_CONSTRUCTION' : 'NO_ADAPTER')
  468. }
  469. })
  470. it('tracks consecutive retry attempts and resets after a successful request', async () => {
  471. const capped = new FailureScriptAdapter([contextError('first overflow'), contextError('second overflow')])
  472. const cappedCtx = await harness(capped)
  473. const cappedAgent = cappedCtx.agentLoop.create(SessionId('retry-cap'), { provider: 'mock', model: 'mock' })
  474. const cappedHistories: string[][] = []
  475. cappedCtx.on('agent/request-error', async (
  476. _agent, _turn, _step, _error, _failure, history, _signal, next,
  477. ) => {
  478. const codes = history.map(entry => entry.code)
  479. cappedHistories.push(codes)
  480. return codes.length < 1 ? { action: 'retry' } : next()
  481. })
  482. send(cappedAgent)
  483. await waitForIdle(cappedCtx, cappedAgent)
  484. expect(cappedHistories).toEqual([[], [CONTEXT_WINDOW_EXCEEDED_CODE]])
  485. const reset = new FailureScriptAdapter([
  486. contextError('first overflow'),
  487. toolCallResponse('retry-reset-call', 'work', {}),
  488. contextError('later overflow'),
  489. ])
  490. const resetCtx = await harness(reset)
  491. resetCtx.tools.register(defineTool({
  492. name: 'work',
  493. description: 'continue',
  494. parameters: {},
  495. async execute() { return [{ type: 'text', text: 'worked' }] },
  496. }))
  497. const resetAgent = resetCtx.agentLoop.create(SessionId('retry-reset'), { provider: 'mock', model: 'mock' })
  498. const resetHistories: { step: number; codes: string[] }[] = []
  499. resetCtx.on('agent/request-error', async (
  500. _agent, _turn, step, _error, _failure, history, _signal, next,
  501. ) => {
  502. resetHistories.push({ step, codes: history.map(entry => entry.code) })
  503. return resetHistories.length === 1 ? { action: 'retry' } : next()
  504. })
  505. send(resetAgent)
  506. await waitForIdle(resetCtx, resetAgent)
  507. expect(resetHistories).toEqual([{ step: 1, codes: [] }, { step: 3, codes: [] }])
  508. })
  509. it('preserves the original provider error when recovery throws', async () => {
  510. const adapter = new FailureScriptAdapter([contextError('original overflow')])
  511. const ctx = await harness(adapter)
  512. const agent = ctx.agentLoop.create(SessionId('recovery-throws'), { provider: 'mock', model: 'mock' })
  513. ctx.on('agent/request-error', () => { throw new Error('recovery exploded') })
  514. send(agent)
  515. await waitForIdle(ctx, agent)
  516. expect(agent.session.events.at(-1)).toMatchObject({
  517. type: 'turn/end',
  518. data: { reason: { kind: 'error', failure: { message: 'original overflow', code: CONTEXT_WINDOW_EXCEEDED_CODE } } },
  519. })
  520. })
  521. it.each(['cancel', 'dispose'] as const)('keeps %s live through request recovery', async (action) => {
  522. const adapter = new FailureScriptAdapter([contextError()])
  523. const ctx = await harness(adapter)
  524. const agent = ctx.agentLoop.create(SessionId(`${action}-recovery`), { provider: 'mock', model: 'mock' })
  525. let entered!: () => void
  526. const recoveryEntered = new Promise<void>((resolve) => { entered = resolve })
  527. ctx.on('agent/request-error', async (_agent, _turn, _step, _error, _failure, _history, signal) => {
  528. entered()
  529. await new Promise<void>((resolve) => {
  530. signal.addEventListener('abort', () => { resolve() }, { once: true })
  531. })
  532. return { action: 'retry' }
  533. })
  534. send(agent)
  535. const idle = waitForIdle(ctx, agent)
  536. await recoveryEntered
  537. if (action === 'cancel') {
  538. agent.cancel('cancelled during recovery')
  539. await idle
  540. } else {
  541. await ctx.fiber.dispose()
  542. }
  543. expect(adapter.requests).toHaveLength(1)
  544. expect(agent.session.events.at(-1)).toMatchObject({
  545. type: 'turn/end',
  546. data: { reason: action === 'cancel' ? { kind: 'aborted', reason: 'cancelled during recovery' } : { kind: 'disposed' } },
  547. })
  548. })
  549. })