Retry a dead provider item inline instead of ending the turn
Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
This commit is contained in:
co-authored by
Sisyphus
parent
4f6ed43fed
commit
a4edbcc71e
+50
-1
@@ -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<ApprovalRequest, 'matchedPattern' | 'suggestedPattern' | 'repeated'>;
|
||||
|
||||
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<string, number>();
|
||||
/** 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<AgentEvent, { type: 'tool-output' }>[] = [];
|
||||
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<string, ApprovalContext>();
|
||||
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;
|
||||
|
||||
+128
-1
@@ -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);
|
||||
});
|
||||
|
||||
|
||||
Reference in New Issue
Block a user