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:
wenshao 2026-05-16 21:49:08 +08:00
parent fc4d57b6ea
commit db344a4037
2 changed files with 137 additions and 0 deletions

View file

@ -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();

View file

@ -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;
}