mirror of
https://github.com/QwenLM/qwen-code.git
synced 2026-08-31 10:16:57 +00:00
fix(core): roll back partial assistant push on retryable mid-stream errors
Regression surfaced by @yiliang114 on PR #4176: a stream attempt that yields a `functionCall` and then throws a retryable error (e.g. `StreamContentError` with a 429 payload, or `InvalidStreamError`) triggers the partial-assistant-turn push in `processStreamResponse`, but the outer `sendMessageStream` retry loop catches the error and issues a fresh attempt. The failed-attempt `model[functionCall]` was left in history, and the successful retry's response landed as a SECOND consecutive `model` entry — invalid user/model alternation AND the failed-attempt `tool_use` is orphan on the wire (no matching tool_result). The very wedge this PR is meant to escape. Track the index of the pushed partial in a new `pendingPartialAssistantTurnIndex` field, reset it on every `sendMessageStream` entry, and pop it before each retry-and-continue path in the catch block (rate-limit, reactive-compression, transient stream anomaly, content-validation retry). Paths that `break` (unretryable) keep the partial — the caller will see it as part of the error surface and the existing repair / dedup machinery in `useGeminiStream.handleCompletedTools` handles it correctly. The pop is defensive (index-bounds + role check) so a hypothetical intervening `setHistory` / `truncateHistory` can't cause an out-of- bounds splice. Adds a regression test in `geminiChat.test.ts` that reproduces @yiliang114's exact shape: yield `functionCall`, throw `StreamContentError(429)`, second attempt yields `Success after retry` plain text + STOP. Asserts the final history is exactly `[user, model(success text)]` — no leading failed-attempt model turn, no orphan `functionCall` anywhere. All 92 geminiChat, 133 client, 92 useGeminiStream tests pass; full core suite green; tsc clean.
This commit is contained in:
parent
fc4d57b6ea
commit
db344a4037
2 changed files with 137 additions and 0 deletions
|
|
@ -2382,6 +2382,87 @@ describe('GeminiChat', async () => {
|
|||
}
|
||||
});
|
||||
|
||||
it('rolls back the partial assistant turn when a retryable error fires after a tool_use chunk', async () => {
|
||||
// Regression for the case @yiliang114 reproduced on #4176: stream
|
||||
// attempt 1 yields a `functionCall` (which triggers the partial-
|
||||
// history push in `processStreamResponse`), then the SAME stream
|
||||
// throws a retryable error (e.g. a TPM 429 `StreamContentError`).
|
||||
// The outer retry loop must drop the partial before issuing the
|
||||
// retry — otherwise the retry's response lands as a SECOND
|
||||
// consecutive `model` entry and the failed-attempt `tool_use`
|
||||
// becomes orphan on the wire (invalid alternation +
|
||||
// tool_use_id-with-no-matching-tool_use 400).
|
||||
vi.useFakeTimers();
|
||||
try {
|
||||
const tpmError = new StreamContentError(
|
||||
'{"error":{"code":"429","message":"Throttling: TPM(1/1)"}}',
|
||||
);
|
||||
const failingStream = (async function* () {
|
||||
yield {
|
||||
candidates: [
|
||||
{
|
||||
content: {
|
||||
parts: [
|
||||
{
|
||||
functionCall: {
|
||||
id: 'call_failed_retry_attempt',
|
||||
name: 'read_file',
|
||||
args: { path: '/tmp/a.txt' },
|
||||
},
|
||||
},
|
||||
],
|
||||
},
|
||||
},
|
||||
],
|
||||
} as unknown as GenerateContentResponse;
|
||||
throw tpmError;
|
||||
})();
|
||||
const successStream = (async function* () {
|
||||
yield {
|
||||
candidates: [
|
||||
{
|
||||
content: { parts: [{ text: 'Success after retry' }] },
|
||||
finishReason: 'STOP',
|
||||
},
|
||||
],
|
||||
} as unknown as GenerateContentResponse;
|
||||
})();
|
||||
vi.mocked(mockContentGenerator.generateContentStream)
|
||||
.mockResolvedValueOnce(failingStream)
|
||||
.mockResolvedValueOnce(successStream);
|
||||
|
||||
const stream = await chat.sendMessageStream(
|
||||
'test-model',
|
||||
{ message: 'test' },
|
||||
'prompt-rollback-on-retry',
|
||||
);
|
||||
const iterator = stream[Symbol.asyncIterator]();
|
||||
// Advance through the rate-limit RETRY + delay, drain all events.
|
||||
for (;;) {
|
||||
const next = iterator.next();
|
||||
await vi.advanceTimersByTimeAsync(60_000);
|
||||
const r = await next;
|
||||
if (r.done) break;
|
||||
}
|
||||
|
||||
const history = chat.getHistory();
|
||||
// History must NOT contain the failed attempt's partial
|
||||
// model[functionCall]. Expected shape: [user, model(success
|
||||
// text)] — exactly two entries, alternation intact.
|
||||
expect(history.length).toBe(2);
|
||||
expect(history[0]!.role).toBe('user');
|
||||
expect(history[1]!.role).toBe('model');
|
||||
const successText = history[1]!.parts!.find((p) => p.text)?.text;
|
||||
expect(successText).toBe('Success after retry');
|
||||
// Defensively: NO functionCall anywhere in history.
|
||||
expect(history.some((h) => h.parts?.some((p) => p.functionCall))).toBe(
|
||||
false,
|
||||
);
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
});
|
||||
|
||||
it('should retry on TPM throttling StreamContentError with initial delay', async () => {
|
||||
vi.useFakeTimers();
|
||||
|
||||
|
|
|
|||
|
|
@ -553,6 +553,24 @@ export class GeminiChat {
|
|||
*/
|
||||
private hasFailedCompressionAttempt = false;
|
||||
|
||||
/**
|
||||
* Index into `this.history` of the model turn that `processStreamResponse`
|
||||
* persisted on the CURRENT in-flight attempt's mid-stream error. `null` if
|
||||
* no partial has been pushed (the common case).
|
||||
*
|
||||
* The retry loop in `sendMessageStream` reads this on every catch to roll
|
||||
* the partial back BEFORE retrying — without that pop, a retryable
|
||||
* mid-stream error (rate limit, transient stream anomaly) leaves the
|
||||
* failed attempt's `model[functionCall]` in history, and the successful
|
||||
* retry's response lands as a SECOND consecutive model turn (invalid
|
||||
* user/model alternation, plus the failed-attempt tool_use is orphan on
|
||||
* the wire — the very wedge this whole subsystem is meant to escape).
|
||||
*
|
||||
* Reset to `null` on every `sendMessageStream` entry so a marker left
|
||||
* over from a prior unretryable break doesn't bleed into the next send.
|
||||
*/
|
||||
private pendingPartialAssistantTurnIndex: number | null = null;
|
||||
|
||||
/**
|
||||
* Creates a new GeminiChat instance.
|
||||
*
|
||||
|
|
@ -730,6 +748,12 @@ export class GeminiChat {
|
|||
});
|
||||
this.sendPromise = streamDonePromise;
|
||||
|
||||
// Clear any partial-push marker left over from a prior unretryable
|
||||
// break path — the marker is per-send; carrying it across sends
|
||||
// would let the next send's retry catch wrongly pop a now-valid
|
||||
// model entry sitting at the stale index.
|
||||
this.pendingPartialAssistantTurnIndex = null;
|
||||
|
||||
let compressionInfo: ChatCompressionInfo;
|
||||
let requestContents: Content[];
|
||||
let userContentAdded = false;
|
||||
|
|
@ -855,12 +879,35 @@ export class GeminiChat {
|
|||
} catch (error) {
|
||||
lastError = error;
|
||||
|
||||
// If `processStreamResponse` persisted a partial assistant turn
|
||||
// (mid-stream error after a `functionCall` was already yielded),
|
||||
// every retry-and-continue path below must drop that turn first.
|
||||
// Otherwise a successful retry's response lands AFTER the stale
|
||||
// failed-attempt model turn — two consecutive `model` entries
|
||||
// with an orphan tool_use in the first, re-triggering the
|
||||
// "tool_use_id ... corresponding tool_use" 400 this fix is
|
||||
// supposed to escape. Paths that `break` (unretryable) keep
|
||||
// the partial — the caller will see it as part of the error
|
||||
// surface.
|
||||
const popPartialIfPushed = () => {
|
||||
const idx = self.pendingPartialAssistantTurnIndex;
|
||||
if (idx === null) return;
|
||||
if (
|
||||
self.history.length > idx &&
|
||||
self.history[idx]?.role === 'model'
|
||||
) {
|
||||
self.history.splice(idx, 1);
|
||||
}
|
||||
self.pendingPartialAssistantTurnIndex = null;
|
||||
};
|
||||
|
||||
// Handle rate-limit / throttling errors returned as stream content.
|
||||
// These arrive as StreamContentError with finish_reason="error_finish"
|
||||
// from the pipeline, containing the throttling message in the content.
|
||||
// Covers TPM throttling, GLM rate limits, and other provider throttling.
|
||||
const isRateLimit = isRateLimitError(error, extraRetryErrorCodes);
|
||||
if (isRateLimit && rateLimitRetryCount < maxRateLimitRetries) {
|
||||
popPartialIfPushed();
|
||||
rateLimitRetryCount++;
|
||||
const delayMs = getRateLimitRetryDelayMs(rateLimitRetryCount, {
|
||||
...RATE_LIMIT_RETRY_OPTIONS,
|
||||
|
|
@ -935,6 +982,7 @@ export class GeminiChat {
|
|||
reactiveInfo.compressionStatus ===
|
||||
CompressionStatus.COMPRESSED
|
||||
) {
|
||||
popPartialIfPushed();
|
||||
requestContents = self.getHistory(true);
|
||||
debugLogger.info(
|
||||
`Reactive compression succeeded: ` +
|
||||
|
|
@ -991,6 +1039,7 @@ export class GeminiChat {
|
|||
isTransientStreamError &&
|
||||
invalidStreamRetryCount < INVALID_STREAM_RETRY_CONFIG.maxRetries
|
||||
) {
|
||||
popPartialIfPushed();
|
||||
invalidStreamRetryCount++;
|
||||
const delayMs =
|
||||
INVALID_STREAM_RETRY_CONFIG.initialDelayMs *
|
||||
|
|
@ -1024,6 +1073,7 @@ export class GeminiChat {
|
|||
const isContentError = error instanceof InvalidStreamError;
|
||||
if (isContentError) {
|
||||
if (attempt < INVALID_CONTENT_RETRY_OPTIONS.maxAttempts - 1) {
|
||||
popPartialIfPushed();
|
||||
logContentRetry(
|
||||
self.config,
|
||||
new ContentRetryEvent(
|
||||
|
|
@ -1586,6 +1636,12 @@ export class GeminiChat {
|
|||
...consolidatedHistoryParts,
|
||||
],
|
||||
});
|
||||
// Track the pushed turn so the outer sendMessageStream retry loop
|
||||
// can roll it back if it decides to retry the same send. Without
|
||||
// this, a successful retry would leave the failed attempt's
|
||||
// partial `model[functionCall]` as a stale leading model turn in
|
||||
// front of the retry's real response.
|
||||
this.pendingPartialAssistantTurnIndex = this.history.length - 1;
|
||||
}
|
||||
throw streamError;
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue