diff --git a/src/session.ts b/src/session.ts index 5b8c96e..628d342 100644 --- a/src/session.ts +++ b/src/session.ts @@ -2,6 +2,7 @@ import { isStepCount, generateText, streamText, + APICallError, type LanguageModel, type ModelMessage, type ToolApprovalResponse, @@ -15,7 +16,7 @@ import { Notebook, type NotebookState } from './notebook'; import { Permissions, type PermissionConfig } from './permission'; import type { PluginHost } from './plugins'; import { systemPrompt } from './prompt'; -import { pruneToFit } from './prune'; +import { detachProviderItems, pruneToFit } from './prune'; import { createSkillTool, renderSkills, type Skill } from './skills'; import { disabledToolNames, onBashOutput, tools as builtinTools, type ToolSetName } from './tools'; @@ -96,6 +97,17 @@ const REPEAT_LIMIT = 3; const callKey = (toolName: string, input: unknown) => `${toolName}:${JSON.stringify(input ?? null)}`; +/** + * The provider rejected an `item_reference` because it no longer holds that item: + * 404 "Item with id 'msg_...' not found". Retrying the same history repeats it, so + * this is the one failure that is worth answering by rewriting the history. + */ +const isStaleItemError = (error: unknown): boolean => + APICallError.isInstance(error) && /item with id '[^']*' not found/i.test(error.message); + +const STALE_ITEM_NOTICE = + 'The provider no longer had part of this session stored. Re-sent the history inline and carried on.'; + type ApprovalContext = Pick; export class Session { @@ -109,6 +121,8 @@ export class Session { private readonly permissions: Permissions; /** Calls seen this turn, for the repeat guard. Cleared per turn, not per step. */ private readonly seen = new Map(); + /** One stale-item repair per turn, so a repeating 404 cannot loop the run. */ + private staleItemsRepaired = false; private controller: AbortController | undefined; constructor(private readonly opts: SessionOptions) { @@ -331,6 +345,7 @@ export class Session { // Per turn, not per step: a tool called once in each of three steps is the // loop this guards against. this.seen.clear(); + this.staleItemsRepaired = false; const outputs: Extract[] = []; onBashOutput(({ toolCallId, chunk }) => { @@ -346,6 +361,19 @@ export class Session { } } + /** + * Rewrites the history so nothing points at provider-side storage, once per turn. + * + * The 404 repeats for every reference in the request, and a repair that could run + * twice would retry a request that cannot be made to work. + */ + private repairStaleItems(): boolean { + if (this.staleItemsRepaired) return false; + this.staleItemsRepaired = true; + this.replace(detachProviderItems(this.messages)); + return true; + } + private async *run( signal: AbortSignal, threshold: number, @@ -360,6 +388,8 @@ export class Session { const guardNotices: string[] = []; const why = new Map(); let sawError = false; + let delivered = false; + let staleRetry = false; const result = streamText({ model: this.model, @@ -405,14 +435,17 @@ export class Session { while (guardNotices.length > 0) yield { type: 'notice', text: guardNotices.shift()! }; switch (part.type) { case 'text-delta': + delivered = true; yield { type: 'text', text: part.text }; break; case 'reasoning-delta': + delivered = true; yield { type: 'reasoning', text: part.text }; break; case 'tool-input-start': // Arrives before the arguments finish streaming, so the UI can name // the tool while the model is still writing its input. + delivered = true; yield { type: 'tool-start', id: part.id, name: part.toolName }; break; case 'tool-call': @@ -448,6 +481,14 @@ export class Session { yield { type: 'done' }; return; case 'error': + // A stale item is rejected before generation starts, so nothing has + // been said yet and the request can be rebuilt. Once output is on + // screen it cannot be unsent, and a retry would repeat it. + if (!delivered && isStaleItemError(part.error) && this.repairStaleItems()) { + yield { type: 'notice', text: STALE_ITEM_NOTICE }; + staleRetry = true; + break; + } sawError = true; yield { type: 'error', error: part.error }; break; @@ -460,10 +501,18 @@ export class Session { yield { type: 'done' }; return; } + if (!delivered && isStaleItemError(error) && this.repairStaleItems()) { + yield { type: 'notice', text: STALE_ITEM_NOTICE }; + continue; + } yield { type: 'error', error }; return; } + // The history was rewritten under this run, so its promise-shaped results + // describe a request that no longer stands. Run again rather than read them. + if (staleRetry) continue; + // A stream that ended in an error has no response messages or usage to // await; touching them would throw NoOutputGeneratedError. if (sawError) return; diff --git a/test/compact.test.ts b/test/compact.test.ts index 59366f1..0ec6693 100644 --- a/test/compact.test.ts +++ b/test/compact.test.ts @@ -1,7 +1,7 @@ import { expect, test } from 'bun:test'; import { MockLanguageModelV4, simulateReadableStream } from 'ai/test'; import type { LanguageModelV4CallOptions, LanguageModelV4StreamPart } from '@ai-sdk/provider'; -import type { ModelMessage } from 'ai'; +import { APICallError, type ModelMessage } from 'ai'; import { mkdtempSync, rmSync } from 'node:fs'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; @@ -326,3 +326,130 @@ test('a compacted turn sends no assistant item reference whose reasoning was pru } } }), 20_000); + +test('a compacted turn inlines a plain assistant item instead of referencing remote storage', async () => { + const seen: LanguageModelV4CallOptions[] = []; + const session = new Session({ + messages: [ + { role: 'user', content: `earlier question ${'x'.repeat(4000)}` }, + { + role: 'assistant', + content: [{ type: 'text', text: 'remembered answer', providerOptions: { openai: { itemId: 'msg_stale' } } }], + }, + ], + compactThreshold: 100, + model: new MockLanguageModelV4({ + doStream: async (o) => { + seen.push(o); + return stream(text('ok')); + }, + }), + askApproval: async () => 'deny', + }); + + const events: AgentEvent[] = []; + for await (const event of session.send('continue')) events.push(event); + + expect(events.some((event) => event.type === 'compacted')).toBe(true); + expect(JSON.stringify(seen[0]?.prompt)).not.toContain('msg_stale'); + expect(JSON.stringify(seen[0]?.prompt)).toContain('remembered answer'); +}); + +/** + * The failure this guards against: an item reference the provider no longer holds is + * rejected with 404 "Item with id 'msg_...' not found", and every retry of the same + * history is rejected the same way, so resuming a session could never get going. + */ +const staleItem = (id: string) => + new APICallError({ + message: `Item with id '${id}' not found.`, + url: 'https://api.openai.com/v1/responses', + requestBodyValues: {}, + statusCode: 404, + isRetryable: false, + }); + +const withStoredItem = (): ModelMessage[] => [ + { role: 'user', content: 'earlier question' }, + { + role: 'assistant', + content: [{ type: 'text', text: 'remembered answer', providerOptions: { openai: { itemId: 'msg_gone' } } }], + }, +]; + +test('a rejected stale item reference is retried inline rather than ending the turn', async () => { + const seen: LanguageModelV4CallOptions[] = []; + let call = 0; + const session = new Session({ + messages: withStoredItem(), + model: new MockLanguageModelV4({ + doStream: async (o) => { + seen.push(o); + if (call++ === 0) throw staleItem('msg_gone'); + return stream(text('picked up where we left off')); + }, + }), + askApproval: async () => 'deny', + maxRetries: 0, + }); + + const events: AgentEvent[] = []; + for await (const event of session.send('continue')) events.push(event); + + expect(events.map((e) => e.type)).toEqual(['notice', 'text', 'done']); + expect(call).toBe(2); + expect(JSON.stringify(seen[0]?.prompt)).toContain('msg_gone'); + // The retry carries the same content with nothing pointing at provider storage. + expect(JSON.stringify(seen[1]?.prompt)).not.toContain('msg_gone'); + expect(JSON.stringify(seen[1]?.prompt)).toContain('remembered answer'); + // Repaired in place, so a save or a later turn cannot resend the dead reference. + expect(JSON.stringify(session.messages)).not.toContain('msg_gone'); +}); + +test('a stale item rejection that survives the repair is reported once, not looped', async () => { + let call = 0; + const session = new Session({ + messages: withStoredItem(), + model: new MockLanguageModelV4({ + doStream: async () => { + call++; + throw staleItem('msg_gone'); + }, + }), + askApproval: async () => 'deny', + maxRetries: 0, + }); + + const events: AgentEvent[] = []; + for await (const event of session.send('continue')) events.push(event); + + expect(events.map((e) => e.type)).toEqual(['notice', 'error']); + expect(call).toBe(2); +}); + +test('a stale item arriving after text was streamed is reported, not silently repeated', async () => { + let call = 0; + const session = new Session({ + messages: withStoredItem(), + model: new MockLanguageModelV4({ + doStream: async () => { + call++; + return stream([ + { type: 'text-start', id: '0' }, + { type: 'text-delta', id: '0', delta: 'half an answer' }, + { type: 'error', error: staleItem('msg_gone') }, + ]); + }, + }), + askApproval: async () => 'deny', + maxRetries: 0, + }); + + const events: AgentEvent[] = []; + for await (const event of session.send('continue')) events.push(event); + + // Retrying here would deliver "half an answer" twice. + expect(events.map((e) => e.type)).toEqual(['text', 'error']); + expect(call).toBe(1); +}); +