request-recovery.spec.ts 24 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605
  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, { defineContentToolFixture } 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.followup([{ 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(defineContentToolFixture({
  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. // Injected context is a plugin-sourced user/message; the direct human
  144. // prompt (user source) stays untracked as before.
  145. const isInjected = event.type === 'user/message' && event.data.source.kind !== 'user'
  146. if (
  147. event.type === 'assistant/message' || event.type === 'tool/call'
  148. || event.type === 'tool/result' || isInjected
  149. || event.type === 'steering/message' || event.type === 'step/end'
  150. ) {
  151. if (!('step' in event.data) || event.data.step === 1) order.push(isInjected ? 'context/message' : event.type)
  152. }
  153. })
  154. ctx.on('agent/post-step', (subject, turn, step, signal) => {
  155. if (subject !== agent || step !== 1) return
  156. expect({ turn, step, aborted: signal.aborted }).toEqual({ turn: 1, step: 1, aborted: false })
  157. subject.inject([{ type: 'text', text: 'listener mutation' }], { source: { kind: 'plugin', plugin: 'post-step' } })
  158. order.push('agent/post-step')
  159. })
  160. send(agent)
  161. await waitForIdle(ctx, agent)
  162. expect(order).toEqual([
  163. 'assistant/message',
  164. 'tool/call',
  165. 'tool/result',
  166. 'tool/call',
  167. 'tool/result',
  168. 'context/message',
  169. 'context/message',
  170. 'steering/message',
  171. 'context/message',
  172. 'agent/post-step',
  173. 'step/end',
  174. ])
  175. })
  176. it('fires post-step for max-tokens and lets cancellation override that success', async () => {
  177. const adapter = new FailureScriptAdapter([maxTokensResponse('partial')])
  178. const ctx = await harness(adapter)
  179. const agent = ctx.agentLoop.create(SessionId('cancel-post-step-max-tokens'), { provider: 'mock', model: 'mock' })
  180. let entered!: () => void
  181. const postStepEntered = new Promise<void>((resolve) => { entered = resolve })
  182. ctx.on('agent/post-step', async (_agent, turn, step, signal) => {
  183. expect({ turn, step }).toEqual({ turn: 1, step: 1 })
  184. entered()
  185. await new Promise<void>((resolve) => {
  186. signal.addEventListener('abort', () => { resolve() }, { once: true })
  187. })
  188. })
  189. send(agent)
  190. const idle = waitForIdle(ctx, agent)
  191. await postStepEntered
  192. agent.cancel({ kind: 'user' })
  193. await idle
  194. expect(agent.session.events.find(event => event.type === 'assistant/message')).toMatchObject({
  195. data: { usage: { inputTokens: 10, outputTokens: 7 } },
  196. })
  197. expect(agent.session.events.at(-1)).toMatchObject({
  198. type: 'turn/end',
  199. data: { reason: { kind: 'aborted' } },
  200. })
  201. })
  202. it('closes the successful step as disposed when disposal lands during post-step', async () => {
  203. const adapter = new FailureScriptAdapter([
  204. toolCallResponse('dispose-call', 'work', {}),
  205. textResponse('must not continue'),
  206. ])
  207. const ctx = await harness(adapter)
  208. ctx.tools.register(defineContentToolFixture({
  209. name: 'work',
  210. description: 'do work',
  211. parameters: {},
  212. async execute() { return [{ type: 'text', text: 'worked' }] },
  213. }))
  214. const agent = ctx.agentLoop.create(SessionId('dispose-post-step'), { provider: 'mock', model: 'mock' })
  215. let entered!: () => void
  216. const postStepEntered = new Promise<void>((resolve) => { entered = resolve })
  217. ctx.on('agent/post-step', async (_agent, turn, step, signal) => {
  218. expect({ turn, step }).toEqual({ turn: 1, step: 1 })
  219. entered()
  220. await new Promise<void>((resolve) => {
  221. signal.addEventListener('abort', () => { resolve() }, { once: true })
  222. })
  223. })
  224. send(agent)
  225. await postStepEntered
  226. await ctx.fiber.dispose()
  227. expect(adapter.requests).toHaveLength(1)
  228. const boundaries = agent.session.events.filter(event =>
  229. event.type === 'step/start' || event.type === 'step/end',
  230. )
  231. expect(boundaries.map(event => event.type)).toEqual(['step/start', 'step/end'])
  232. expect(boundaries.map(event => event.data)).toEqual([
  233. { turn: 1, step: 1 },
  234. { turn: 1, step: 1 },
  235. ])
  236. expect(agent.session.events.at(-1)).toMatchObject({
  237. type: 'turn/end',
  238. data: { reason: { kind: 'disposed' } },
  239. })
  240. })
  241. it.each([
  242. ['thrown', contextError()],
  243. ['in-band', [{ type: 'finish', reason: { kind: 'error', failure: { message: 'too large', code: CONTEXT_WINDOW_EXCEEDED_CODE, status: 400 } } }] satisfies StreamChunk[]],
  244. ] as const)('recovers a %s request failure in a new reconstructable step', async (_style, failure) => {
  245. const adapter = new FailureScriptAdapter([failure, textResponse('recovered')])
  246. const ctx = await harness(adapter)
  247. const agent = ctx.agentLoop.create(SessionId(`recover-${_style}`), { provider: 'mock', model: 'mock' })
  248. const attempts: number[] = []
  249. ctx.on('agent/request-error', async (subject, turn, step, error, facts, history) => {
  250. expect(subject).toBe(agent)
  251. expect({ turn, step, code: error.code }).toEqual({ turn: 1, step: 1, code: CONTEXT_WINDOW_EXCEEDED_CODE })
  252. expect(facts.code).toBe(CONTEXT_WINDOW_EXCEEDED_CODE)
  253. attempts.push(history.length)
  254. subject.session.append('user/message', {
  255. content: [{ type: 'text', text: 'RECOVERY SURFACE MUTATION' }],
  256. source: { kind: 'plugin', plugin: 'test-recovery' },
  257. }, { surfaceOp: 'append' })
  258. return { action: 'retry' }
  259. })
  260. send(agent)
  261. await waitForIdle(ctx, agent)
  262. expect(attempts).toEqual([0])
  263. expect(adapter.requests).toHaveLength(2)
  264. expect(JSON.stringify(adapter.requests[1]!.messages)).toContain('RECOVERY SURFACE MUTATION')
  265. const starts = agent.session.events.filter(event => event.type === 'step/start')
  266. const ends = agent.session.events.filter(event => event.type === 'step/end')
  267. expect(starts.map(event => event.data.step)).toEqual([1, 2])
  268. expect(ends.map(event => event.data.step)).toEqual([1, 2])
  269. const recovery = agent.session.events.find(event => event.type === 'user/message' && event.data.source.kind === 'plugin')!
  270. expect(ends[0]!.seq).toBeLessThan(recovery.seq)
  271. expect(recovery.seq).toBeLessThan(starts[1]!.seq)
  272. })
  273. it.each(streamListenerFailureCases)('does not offer %s to request recovery', async (_name, install) => {
  274. const ctx = await harness(new FailureScriptAdapter([textResponse('unused')]))
  275. const agent = ctx.agentLoop.create(SessionId(`stream-plugin-${_name.replaceAll(' ', '-')}`), { provider: 'mock', model: 'mock' })
  276. let recoveries = 0
  277. install(ctx)
  278. ctx.on('agent/request-error', async (_agent, _turn, _step, _error, _failure, _history, _signal, next) => {
  279. recoveries += 1
  280. return next()
  281. })
  282. send(agent)
  283. await waitForIdle(ctx, agent)
  284. expect(recoveries).toBe(0)
  285. expect(agent.session.events.at(-1)).toMatchObject({ type: 'turn/end', data: { reason: { kind: 'error' } } })
  286. })
  287. it('does not offer a nested model-call failure as the outer request failure', async () => {
  288. const outer = new FailureScriptAdapter([textResponse('outer adapter must not run')])
  289. const nested = new FailureScriptAdapter([contextError('nested overflow')])
  290. const ctx = await harness(outer)
  291. ctx.llm.registerAdapter(['nested'], nested)
  292. ctx.on('llm/stream', (options, next) => {
  293. if (options.provider !== 'mock') return next()
  294. return (async function* () {
  295. yield* ctx.llm.stream({
  296. provider: 'nested',
  297. model: 'nested',
  298. messages: [],
  299. ...options.signal === undefined ? {} : { signal: options.signal },
  300. })
  301. yield* next()
  302. })()
  303. })
  304. const agent = ctx.agentLoop.create(SessionId('nested-stream-not-recoverable'), { provider: 'mock', model: 'mock' })
  305. let recoveries = 0
  306. ctx.on('agent/request-error', async (_agent, _turn, _step, _error, _failure, _history, _signal, next) => {
  307. recoveries += 1
  308. return next()
  309. })
  310. send(agent)
  311. await waitForIdle(ctx, agent)
  312. expect(nested.requests).toHaveLength(1)
  313. expect(outer.requests).toHaveLength(0)
  314. expect(recoveries).toBe(0)
  315. expect(agent.session.events.at(-1)).toMatchObject({
  316. type: 'turn/end',
  317. data: { reason: { kind: 'error', message: 'nested overflow', code: CONTEXT_WINDOW_EXCEEDED_CODE } },
  318. })
  319. })
  320. it.each(['prompt-submit', 'prompt-assembly', 'pre-step', 'request'] as const)(
  321. 'does not offer %s middleware failures to request recovery',
  322. async (boundary) => {
  323. const adapter = new FailureScriptAdapter([textResponse('unused')])
  324. const ctx = await harness(adapter)
  325. if (boundary === 'prompt-submit') {
  326. ctx.on('agent/prompt-submit', () => { throw new Error('prompt submit failed') })
  327. } else if (boundary === 'prompt-assembly') {
  328. ctx.on('system-prompt/assemble', () => { throw new Error('prompt assembly failed') })
  329. } else if (boundary === 'pre-step') {
  330. ctx.on('agent/pre-step', () => { throw new Error('pre-step failed') })
  331. } else {
  332. ctx.on('agent/request', () => { throw new Error('request middleware failed') })
  333. }
  334. const agent = ctx.agentLoop.create(SessionId(`${boundary}-not-recoverable`), { provider: 'mock', model: 'mock' })
  335. let recoveries = 0
  336. ctx.on('agent/request-error', async (_agent, _turn, _step, _error, _failure, _history, _signal, next) => {
  337. recoveries += 1
  338. return next()
  339. })
  340. send(agent)
  341. await waitForIdle(ctx, agent)
  342. expect(recoveries).toBe(0)
  343. expect(adapter.requests).toHaveLength(0)
  344. expect(agent.session.events.at(-1)).toMatchObject({ type: 'turn/end', data: { reason: { kind: 'error' } } })
  345. },
  346. )
  347. it('does not offer result, tool, or post-step plugin failures to request recovery', async () => {
  348. for (const failure of ['result', 'tool', 'post-step'] as const) {
  349. const adapter = new FailureScriptAdapter([
  350. failure === 'tool' ? toolCallResponse(`call-${failure}`, 'work', {}) : textResponse('done'),
  351. ...(failure === 'tool' ? [textResponse('done')] : []),
  352. ])
  353. const ctx = await harness(adapter)
  354. if (failure === 'result') ctx.on('agent/step-result', () => { throw new Error('result failed') })
  355. if (failure === 'post-step') ctx.on('agent/post-step', () => { throw new Error('post-step failed') })
  356. if (failure === 'tool') {
  357. vi.spyOn(ctx.tools, 'execute').mockRejectedValue(new Error('tool service failed'))
  358. }
  359. const agent = ctx.agentLoop.create(SessionId(`${failure}-not-recoverable`), { provider: 'mock', model: 'mock' })
  360. let recoveries = 0
  361. ctx.on('agent/request-error', async (_agent, _turn, _step, _error, _failure, _history, _signal, next) => {
  362. recoveries += 1
  363. return next()
  364. })
  365. send(agent)
  366. await waitForIdle(ctx, agent)
  367. expect(recoveries, failure).toBe(0)
  368. }
  369. })
  370. it.each([
  371. ['synchronous dispatch', (error: Error) => new SynchronousDispatchFailureAdapter(error)],
  372. ['done getter', (error: Error) => new IteratorResultGetterFailureAdapter('done', error)],
  373. ['value getter', (error: Error) => new IteratorResultGetterFailureAdapter('value', error)],
  374. ] as const)('preserves original Error identity for adapter %s', async (_name, makeAdapter) => {
  375. const original = contextError(`${_name} overflow`)
  376. const ctx = await harness(makeAdapter(original))
  377. const agent = ctx.agentLoop.create(SessionId(`identity-${_name.replaceAll(' ', '-')}`), { provider: 'mock', model: 'mock' })
  378. let seen: Error | undefined
  379. ctx.on('agent/request-error', async (_agent, _turn, _step, error, _failure, _history, _signal, next) => {
  380. seen = error
  381. return next()
  382. })
  383. send(agent)
  384. await waitForIdle(ctx, agent)
  385. expect(seen).toBe(original)
  386. })
  387. it('keeps an adapter error with a hostile message accessor on the recovery path', async () => {
  388. const original = Object.defineProperty(new HarnessError('provider failed', 'SERVER'), 'message', {
  389. get() { throw new Error('SDK message accessor trap') },
  390. })
  391. const ctx = await harness(new SynchronousDispatchFailureAdapter(original))
  392. const agent = ctx.agentLoop.create(SessionId('hostile-message-recovery'), { provider: 'mock', model: 'mock' })
  393. let seenError: Error | undefined
  394. let seenFailure: LlmFailure | undefined
  395. ctx.on('agent/request-error', async (_agent, _turn, _step, error, failure, _history, _signal, next) => {
  396. seenError = error
  397. seenFailure = failure
  398. return next()
  399. })
  400. send(agent)
  401. await waitForIdle(ctx, agent)
  402. expect(seenError).toBe(original)
  403. expect(seenFailure).toEqual({ message: 'LLM adapter failed', code: 'SERVER' })
  404. expect(agent.session.events.at(-1)).toMatchObject({
  405. type: 'turn/end',
  406. data: { reason: { kind: 'error', failure: { message: 'LLM adapter failed', code: 'SERVER' } } },
  407. })
  408. })
  409. it('passes structured facts beside the original Error and records its cause chain on exhaustion', async () => {
  410. const original = new LlmError('provider busy', 'RATE_LIMIT', {
  411. cause: new Error('upstream connection reset'),
  412. status: 429,
  413. providerRetryAfterMs: 2_000,
  414. requestId: ProviderRequestId('req-9'),
  415. })
  416. Object.freeze(original)
  417. const ctx = await harness(new SynchronousDispatchFailureAdapter(original))
  418. const agent = ctx.agentLoop.create(SessionId('structured-request-failure'), { provider: 'mock', model: 'mock' })
  419. let seenError: Error | undefined
  420. let seenFailure: LlmFailure | undefined
  421. let seenHistory: readonly LlmFailure[] | undefined
  422. ctx.on('agent/request-error', async (
  423. _agent, _turn, _step, error, failure, history, _signal, next,
  424. ) => {
  425. seenError = error
  426. seenFailure = failure
  427. seenHistory = history
  428. return next()
  429. })
  430. send(agent)
  431. await waitForIdle(ctx, agent)
  432. expect(seenError).toBe(original)
  433. expect(seenFailure).toEqual({
  434. message: 'provider busy',
  435. code: 'RATE_LIMIT',
  436. status: 429,
  437. providerRetryAfterMs: 2_000,
  438. requestId: ProviderRequestId('req-9'),
  439. })
  440. expect(seenHistory).toEqual([])
  441. expect(Object.isFrozen(seenHistory)).toBe(true)
  442. expect(agent.session.events.at(-1)).toMatchObject({
  443. type: 'turn/end',
  444. data: {
  445. reason: {
  446. kind: 'error',
  447. step: 1,
  448. failure: {
  449. message: 'provider busy: upstream connection reset',
  450. code: 'RATE_LIMIT',
  451. status: 429,
  452. providerRetryAfterMs: 2_000,
  453. requestId: ProviderRequestId('req-9'),
  454. },
  455. },
  456. },
  457. })
  458. })
  459. it('classifies iterator construction and explicit NO_ADAPTER as model-request failures', async () => {
  460. for (const scenario of ['iterator', 'no-adapter'] as const) {
  461. const ctx = scenario === 'iterator' ? await harness(new IteratorConstructionFailureAdapter()) : await harness()
  462. const agent = ctx.agentLoop.create(SessionId(`request-boundary-${scenario}`), { provider: 'mock', model: 'mock' })
  463. let seen = ''
  464. ctx.on('agent/request-error', async (_agent, _turn, _step, error, _failure, _history, _signal, next) => {
  465. seen = error.code ?? ''
  466. return next()
  467. })
  468. send(agent)
  469. await waitForIdle(ctx, agent)
  470. expect(seen).toBe(scenario === 'iterator' ? 'ITERATOR_CONSTRUCTION' : 'NO_ADAPTER')
  471. }
  472. })
  473. it('tracks consecutive retry attempts and resets after a successful request', async () => {
  474. const capped = new FailureScriptAdapter([contextError('first overflow'), contextError('second overflow')])
  475. const cappedCtx = await harness(capped)
  476. const cappedAgent = cappedCtx.agentLoop.create(SessionId('retry-cap'), { provider: 'mock', model: 'mock' })
  477. const cappedHistories: string[][] = []
  478. cappedCtx.on('agent/request-error', async (
  479. _agent, _turn, _step, _error, _failure, history, _signal, next,
  480. ) => {
  481. const codes = history.map(entry => entry.code)
  482. cappedHistories.push(codes)
  483. return codes.length < 1 ? { action: 'retry' } : next()
  484. })
  485. send(cappedAgent)
  486. await waitForIdle(cappedCtx, cappedAgent)
  487. expect(cappedHistories).toEqual([[], [CONTEXT_WINDOW_EXCEEDED_CODE]])
  488. const reset = new FailureScriptAdapter([
  489. contextError('first overflow'),
  490. toolCallResponse('retry-reset-call', 'work', {}),
  491. contextError('later overflow'),
  492. ])
  493. const resetCtx = await harness(reset)
  494. resetCtx.tools.register(defineContentToolFixture({
  495. name: 'work',
  496. description: 'continue',
  497. parameters: {},
  498. async execute() { return [{ type: 'text', text: 'worked' }] },
  499. }))
  500. const resetAgent = resetCtx.agentLoop.create(SessionId('retry-reset'), { provider: 'mock', model: 'mock' })
  501. const resetHistories: { step: number; codes: string[] }[] = []
  502. resetCtx.on('agent/request-error', async (
  503. _agent, _turn, step, _error, _failure, history, _signal, next,
  504. ) => {
  505. resetHistories.push({ step, codes: history.map(entry => entry.code) })
  506. return resetHistories.length === 1 ? { action: 'retry' } : next()
  507. })
  508. send(resetAgent)
  509. await waitForIdle(resetCtx, resetAgent)
  510. expect(resetHistories).toEqual([{ step: 1, codes: [] }, { step: 3, codes: [] }])
  511. })
  512. it('preserves the original provider error when recovery throws', async () => {
  513. const adapter = new FailureScriptAdapter([contextError('original overflow')])
  514. const ctx = await harness(adapter)
  515. const agent = ctx.agentLoop.create(SessionId('recovery-throws'), { provider: 'mock', model: 'mock' })
  516. ctx.on('agent/request-error', () => { throw new Error('recovery exploded') })
  517. send(agent)
  518. await waitForIdle(ctx, agent)
  519. expect(agent.session.events.at(-1)).toMatchObject({
  520. type: 'turn/end',
  521. data: { reason: { kind: 'error', failure: { message: 'original overflow', code: CONTEXT_WINDOW_EXCEEDED_CODE } } },
  522. })
  523. })
  524. it.each(['cancel', 'dispose'] as const)('keeps %s live through request recovery', async (action) => {
  525. const adapter = new FailureScriptAdapter([contextError()])
  526. const ctx = await harness(adapter)
  527. const agent = ctx.agentLoop.create(SessionId(`${action}-recovery`), { provider: 'mock', model: 'mock' })
  528. let entered!: () => void
  529. const recoveryEntered = new Promise<void>((resolve) => { entered = resolve })
  530. ctx.on('agent/request-error', async (_agent, _turn, _step, _error, _failure, _history, signal) => {
  531. entered()
  532. await new Promise<void>((resolve) => {
  533. signal.addEventListener('abort', () => { resolve() }, { once: true })
  534. })
  535. return { action: 'retry' }
  536. })
  537. send(agent)
  538. const idle = waitForIdle(ctx, agent)
  539. await recoveryEntered
  540. if (action === 'cancel') {
  541. agent.cancel({ kind: 'user' })
  542. await idle
  543. } else {
  544. await ctx.fiber.dispose()
  545. }
  546. expect(adapter.requests).toHaveLength(1)
  547. expect(agent.session.events.at(-1)).toMatchObject({
  548. type: 'turn/end',
  549. data: { reason: action === 'cancel' ? { kind: 'aborted' } : { kind: 'disposed' } },
  550. })
  551. })
  552. })