mirror of
https://github.com/QwenLM/qwen-code.git
synced 2026-08-31 10:16:57 +00:00
fix(core): defer chat-recording flush until partial-turn rollback decision
Yiliang114's PR #4176 follow-up: `popPartialIfPushed()` rolls a failed-attempt's `model[functionCall]` out of in-memory `this.history` before retrying, but the same partial turn was already appended to the chat-recording JSONL by `recordAssistantTurn()`. After a successful retry, live history was clean while the durable transcript still carried two assistant entries — the failed attempt and the success. `--resume` then rehydrated the failed `model[functionCall]` and the resumed model context picked up a tool_use the live session correctly discarded. Fix: the recording call on the stream-error + hasToolCall path is now deferred. processStreamResponse stashes the would-be record on a new `pendingPartialAssistantRecord` field instead of immediately appending. The retry loop's `popPartialIfPushed` clears the stash alongside the in-memory splice — no flush, no leak. Once the retry loop has settled (success-break: stash already null; unretryable-break: partial survived), a flush hook right after the for-loop appends the stashed record to JSONL so the transcript matches whatever lives in `this.history`. The success path (streamError === null) still records immediately as before — only the at-risk-of-rollback path is deferred. Lifecycle is paired one-to-one with `pendingPartialAssistantTurnIndex`: set together, popped together, flushed together, and reset in every history-replacement method (sendMessageStream entry, setHistory, clearHistory, truncateHistory, addHistory, stripThoughtsFromHistory) so an exotic path can't leak a stale stash into a future send's flush. Two regression tests in geminiChat.test.ts: 1. "rolls back the chat-recording entry too when the retry succeeds" — failing-then-succeeding stream sequence asserts `recordAssistantTurn` is called exactly once with the success text, no functionCall part anywhere in the recorded message. Verified the test fails (mock called twice) when the deferral branch is bypassed. 2. "flushes the chat-recording entry on the unretryable break path" — synthetic non-retryable error after a tool_use chunk asserts the partial IS recorded so the JSONL stays aligned with the in-memory partial that survives. Verified the test fails (mock called zero times) when the post-loop flush is removed. Tests: 104/104 geminiChat (+2 new), 133/133 client, 94/94 useGeminiStream. tsc + eslint + prettier clean.
This commit is contained in:
parent
08216d9484
commit
2dbfc4e3bb
2 changed files with 274 additions and 8 deletions
|
|
@ -2524,6 +2524,188 @@ describe('GeminiChat', async () => {
|
|||
}
|
||||
});
|
||||
|
||||
it('rolls back the chat-recording entry too when the retry succeeds (yiliang114 PR #4176 follow-up)', async () => {
|
||||
// The in-memory rollback test above asserts `this.history` ends
|
||||
// clean after a retry-success. This test asserts the same about
|
||||
// chat-recording JSONL: the failed attempt's `recordAssistantTurn`
|
||||
// call must NOT have been flushed, so `--resume` won't rehydrate
|
||||
// a model[functionCall] turn the live session correctly discarded.
|
||||
// Without the deferred-flush stash + popPartialIfPushed clear,
|
||||
// `recordAssistantTurn` was called twice (once for the partial,
|
||||
// once for the success) and only the in-memory pop fixed live
|
||||
// history; the durable transcript stayed corrupt.
|
||||
vi.useFakeTimers();
|
||||
try {
|
||||
const recordAssistantTurn = vi.fn();
|
||||
const chatWithRecording = new GeminiChat(
|
||||
mockConfig,
|
||||
config,
|
||||
[],
|
||||
{
|
||||
recordAssistantTurn,
|
||||
recordChatCompression: vi.fn(),
|
||||
} as unknown as ConstructorParameters<typeof GeminiChat>[3],
|
||||
uiTelemetryService,
|
||||
);
|
||||
|
||||
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_recording',
|
||||
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 chatWithRecording.sendMessageStream(
|
||||
'test-model',
|
||||
{ message: 'test' },
|
||||
'prompt-recording-rollback',
|
||||
);
|
||||
const iterator = stream[Symbol.asyncIterator]();
|
||||
for (;;) {
|
||||
const next = iterator.next();
|
||||
await vi.advanceTimersByTimeAsync(60_000);
|
||||
const r = await next;
|
||||
if (r.done) break;
|
||||
}
|
||||
|
||||
// Exactly one recording: the successful retry's text turn.
|
||||
// The failed attempt's partial functionCall must have been
|
||||
// discarded by `popPartialIfPushed` clearing the deferred-flush
|
||||
// stash, never reaching the JSONL.
|
||||
expect(recordAssistantTurn).toHaveBeenCalledTimes(1);
|
||||
const recordedMessage = recordAssistantTurn.mock.calls[0]![0]
|
||||
?.message as Array<{ text?: string; functionCall?: unknown }>;
|
||||
const recordedText = recordedMessage.find((p) => p.text)?.text;
|
||||
expect(recordedText).toBe('Success after retry');
|
||||
// No functionCall part anywhere in the recorded turn.
|
||||
expect(recordedMessage.some((p) => p.functionCall)).toBe(false);
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
});
|
||||
|
||||
it('flushes the chat-recording entry on the unretryable break path (kept partial → durable JSONL)', async () => {
|
||||
// Counterpart to the rollback test: when the retry budget is
|
||||
// exhausted (or the error is unretryable from the start), the
|
||||
// partial assistant turn IS kept in `this.history` — and the
|
||||
// chat-recording JSONL must match. Without the deferred-flush
|
||||
// path firing at the rethrow site, the JSONL silently drops a
|
||||
// partial that's still in live history, and the orphan-tool_use
|
||||
// repair pass at session-load has no dangling functionCall to
|
||||
// close → `--resume` first send 400s with the very wedge the
|
||||
// repair was supposed to escape.
|
||||
vi.useFakeTimers();
|
||||
try {
|
||||
const recordAssistantTurn = vi.fn();
|
||||
const chatWithRecording = new GeminiChat(
|
||||
mockConfig,
|
||||
config,
|
||||
[],
|
||||
{
|
||||
recordAssistantTurn,
|
||||
recordChatCompression: vi.fn(),
|
||||
} as unknown as ConstructorParameters<typeof GeminiChat>[3],
|
||||
uiTelemetryService,
|
||||
);
|
||||
|
||||
// Unretryable: a non-rate-limit, non-InvalidStream error after
|
||||
// a tool_use chunk lands. The catch block falls through to
|
||||
// `break` with the partial kept in memory.
|
||||
const failingStream = (async function* () {
|
||||
yield {
|
||||
candidates: [
|
||||
{
|
||||
content: {
|
||||
parts: [
|
||||
{
|
||||
functionCall: {
|
||||
id: 'call_unretryable_kept',
|
||||
name: 'read_file',
|
||||
args: { path: '/tmp/k.txt' },
|
||||
},
|
||||
},
|
||||
],
|
||||
},
|
||||
},
|
||||
],
|
||||
} as unknown as GenerateContentResponse;
|
||||
throw new Error('synthetic unretryable mid-stream failure');
|
||||
})();
|
||||
vi.mocked(
|
||||
mockContentGenerator.generateContentStream,
|
||||
).mockResolvedValueOnce(failingStream);
|
||||
|
||||
const stream = await chatWithRecording.sendMessageStream(
|
||||
'test-model',
|
||||
{ message: 'test' },
|
||||
'prompt-recording-flush-on-break',
|
||||
);
|
||||
const iterator = stream[Symbol.asyncIterator]();
|
||||
await expect(
|
||||
(async () => {
|
||||
for (;;) {
|
||||
const r = await iterator.next();
|
||||
if (r.done) return;
|
||||
}
|
||||
})(),
|
||||
).rejects.toThrow(/synthetic unretryable/);
|
||||
|
||||
// In-memory: partial is kept (the wedge-recovery contract that
|
||||
// the rest of this PR's machinery relies on).
|
||||
const history = chatWithRecording.getHistory();
|
||||
const lastModelTurn = history.findLast((h) => h.role === 'model');
|
||||
expect(
|
||||
lastModelTurn?.parts?.some(
|
||||
(p) => p.functionCall?.id === 'call_unretryable_kept',
|
||||
),
|
||||
).toBe(true);
|
||||
|
||||
// JSONL: must contain the same partial turn so `--resume` sees
|
||||
// a transcript that matches live history. Exactly one record
|
||||
// (no success retry happened on this path).
|
||||
expect(recordAssistantTurn).toHaveBeenCalledTimes(1);
|
||||
const recordedMessage = recordAssistantTurn.mock.calls[0]![0]
|
||||
?.message as Array<{ functionCall?: { id?: string } }>;
|
||||
expect(
|
||||
recordedMessage.some(
|
||||
(p) => p.functionCall?.id === 'call_unretryable_kept',
|
||||
),
|
||||
).toBe(true);
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
});
|
||||
|
||||
it('should retry on TPM throttling StreamContentError with initial delay', async () => {
|
||||
vi.useFakeTimers();
|
||||
|
||||
|
|
|
|||
|
|
@ -583,6 +583,32 @@ export class GeminiChat {
|
|||
*/
|
||||
private pendingPartialAssistantTurnIndex: number | null = null;
|
||||
|
||||
/**
|
||||
* Deferred-flush record for the partial assistant turn pushed on a
|
||||
* mid-stream error. Stashed instead of immediately appended to the
|
||||
* chat-recording JSONL so a subsequent retry-success can roll back the
|
||||
* persisted record alongside the in-memory pop. Without deferral, the
|
||||
* JSONL transcript keeps the failed attempt even though live history
|
||||
* dropped it — `--resume` rehydrates the failed `model[functionCall]`
|
||||
* turn and the resumed model context picks up a tool_use the in-session
|
||||
* run intentionally discarded (yiliang114 repro on PR #4176).
|
||||
*
|
||||
* Lifecycle is paired with `pendingPartialAssistantTurnIndex`:
|
||||
* - Set together on stream-error + hasToolCall + hasContent.
|
||||
* - Cleared together by `popPartialIfPushed` when the retry loop
|
||||
* rolls the partial back.
|
||||
* - Flushed together to JSONL after the retry loop exits if the
|
||||
* partial survived (unretryable break path → about to throw to
|
||||
* the caller). At that point the partial is durable in memory and
|
||||
* the recording must match.
|
||||
* - Reset on `sendMessageStream` entry / setHistory / clearHistory /
|
||||
* truncateHistory / addHistory / stripThoughtsFromHistory so a leak
|
||||
* from any exotic path can't bleed into a future send.
|
||||
*/
|
||||
private pendingPartialAssistantRecord:
|
||||
| Parameters<ChatRecordingService['recordAssistantTurn']>[0]
|
||||
| null = null;
|
||||
|
||||
/**
|
||||
* Heap-pressure compaction is process-wide pressure applied per chat. If one
|
||||
* heap-triggered attempt cannot reduce history, briefly back off this chat
|
||||
|
|
@ -826,8 +852,13 @@ export class GeminiChat {
|
|||
// 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.
|
||||
// model entry sitting at the stale index. The deferred-record
|
||||
// stash gets the same per-send reset for the same reason: a
|
||||
// leftover from a prior unretryable break would otherwise get
|
||||
// appended to JSONL by THIS send's retry-loop flush, attaching
|
||||
// someone else's failed turn to this conversation.
|
||||
this.pendingPartialAssistantTurnIndex = null;
|
||||
this.pendingPartialAssistantRecord = null;
|
||||
|
||||
let compressionInfo: ChatCompressionInfo;
|
||||
let requestContents: Content[];
|
||||
|
|
@ -974,6 +1005,13 @@ export class GeminiChat {
|
|||
self.history.splice(idx, 1);
|
||||
}
|
||||
self.pendingPartialAssistantTurnIndex = null;
|
||||
// Discard the deferred chat-recording record alongside the
|
||||
// in-memory pop so the JSONL transcript also drops the
|
||||
// failed attempt. (Paired with the stash in
|
||||
// processStreamResponse — see the field-level comment on
|
||||
// `pendingPartialAssistantRecord` for the failure mode this
|
||||
// fixes.)
|
||||
self.pendingPartialAssistantRecord = null;
|
||||
};
|
||||
|
||||
// Handle rate-limit / throttling errors returned as stream content.
|
||||
|
|
@ -1169,6 +1207,23 @@ export class GeminiChat {
|
|||
}
|
||||
}
|
||||
|
||||
// The retry loop has settled: any partial that was rolled back
|
||||
// had its stash cleared by `popPartialIfPushed`; any partial that
|
||||
// survived (success break with no partial set, or unretryable
|
||||
// break with the partial kept) is now durable in memory, so the
|
||||
// deferred chat-recording append must finally land on disk.
|
||||
// Without this flush the unretryable-break path persists the
|
||||
// partial in `this.history` but the JSONL transcript silently
|
||||
// drops it — `--resume` then loads a truncated transcript that
|
||||
// doesn't match the live session shape, and the orphan-tool_use
|
||||
// repair pass at session-load has nothing to repair.
|
||||
if (self.pendingPartialAssistantRecord) {
|
||||
self.chatRecordingService?.recordAssistantTurn(
|
||||
self.pendingPartialAssistantRecord,
|
||||
);
|
||||
self.pendingPartialAssistantRecord = null;
|
||||
}
|
||||
|
||||
// Max output tokens escalation: if the retry loop succeeded with
|
||||
// the capped default (8K) but hit MAX_TOKENS, retry once at the
|
||||
// model's full output limit. This ensures models with large output
|
||||
|
|
@ -1459,8 +1514,11 @@ export class GeminiChat {
|
|||
// shows up at that index in a future send (defense-in-depth — the
|
||||
// helper also bounds-checks, but a stale marker that happens to
|
||||
// line up with a real model turn could otherwise pop the wrong
|
||||
// entry).
|
||||
// entry). Drop the deferred-record stash for the same reason: a
|
||||
// later flush would otherwise append a turn that doesn't match
|
||||
// the (now-empty) live history.
|
||||
this.pendingPartialAssistantTurnIndex = null;
|
||||
this.pendingPartialAssistantRecord = null;
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -1478,8 +1536,12 @@ export class GeminiChat {
|
|||
// clearHistory: any future addHistory variant that splices into
|
||||
// the middle (instead of plain push) would shift indices and a
|
||||
// stale marker could splice the wrong entry. Belt-and-suspenders;
|
||||
// the next sendMessageStream entry also clears it.
|
||||
// the next sendMessageStream entry also clears it. The
|
||||
// deferred-record stash is paired with the marker — keeping it
|
||||
// around past an external addHistory could let a subsequent
|
||||
// retry-loop flush land a stale partial in JSONL.
|
||||
this.pendingPartialAssistantTurnIndex = null;
|
||||
this.pendingPartialAssistantRecord = null;
|
||||
}
|
||||
|
||||
setHistory(history: Content[]): void {
|
||||
|
|
@ -1489,8 +1551,10 @@ export class GeminiChat {
|
|||
// marker MUST be cleared — otherwise `popPartialIfPushed` could find
|
||||
// a model turn at the stale index in the replacement history and
|
||||
// splice an entry that has nothing to do with the original partial
|
||||
// push, corrupting the conversation.
|
||||
// push, corrupting the conversation. Drop the paired deferred-record
|
||||
// stash too: its referent (the model turn at the old index) is gone.
|
||||
this.pendingPartialAssistantTurnIndex = null;
|
||||
this.pendingPartialAssistantRecord = null;
|
||||
}
|
||||
|
||||
truncateHistory(keepCount: number): void {
|
||||
|
|
@ -1500,8 +1564,10 @@ export class GeminiChat {
|
|||
// the marker rather than try to fix it up — it's per-send and
|
||||
// ephemeral, so losing it across a truncate is safe (the
|
||||
// sendMessageStream that pushed it has already finished or will
|
||||
// start fresh on the next call).
|
||||
// start fresh on the next call). Drop the paired deferred-record
|
||||
// stash for the same reason.
|
||||
this.pendingPartialAssistantTurnIndex = null;
|
||||
this.pendingPartialAssistantRecord = null;
|
||||
}
|
||||
|
||||
stripThoughtsFromHistory(): void {
|
||||
|
|
@ -1510,8 +1576,11 @@ export class GeminiChat {
|
|||
.filter((content): content is Content => content !== null);
|
||||
// Filter+map replaces `this.history` with a new array, so any pending
|
||||
// partial-push marker is now indexed against an array that no longer
|
||||
// exists. Clear it for the same reason setHistory does.
|
||||
// exists. Clear it for the same reason setHistory does. The deferred
|
||||
// chat-recording stash is paired with the marker — drop it too so a
|
||||
// later flush can't land a turn that doesn't exist in live history.
|
||||
this.pendingPartialAssistantTurnIndex = null;
|
||||
this.pendingPartialAssistantRecord = null;
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -1712,7 +1781,7 @@ export class GeminiChat {
|
|||
) {
|
||||
const contextWindowSize =
|
||||
this.config.getContentGeneratorConfig()?.contextWindowSize;
|
||||
this.chatRecordingService?.recordAssistantTurn({
|
||||
const recordArgs = {
|
||||
model,
|
||||
message: [
|
||||
...(thoughtContentPart ? [thoughtContentPart] : []),
|
||||
|
|
@ -1730,7 +1799,22 @@ export class GeminiChat {
|
|||
],
|
||||
tokens: usageMetadata,
|
||||
contextWindowSize,
|
||||
});
|
||||
};
|
||||
if (streamError !== null) {
|
||||
// Stream-error + tool-use partial: defer the JSONL append until
|
||||
// the outer retry loop decides whether to roll back this attempt.
|
||||
// If the same send retries successfully, popPartialIfPushed clears
|
||||
// this stash and the failed attempt never lands on disk; if the
|
||||
// retry path doesn't apply (unretryable break), the stash is
|
||||
// flushed at the rethrow site so JSONL stays aligned with the
|
||||
// partial that survives in-memory. Without this, retry-success
|
||||
// leaves a failed `model[functionCall]` durable in JSONL and
|
||||
// `--resume` rehydrates a turn the live session correctly
|
||||
// discarded (yiliang114 PR #4176 repro).
|
||||
this.pendingPartialAssistantRecord = recordArgs;
|
||||
} else {
|
||||
this.chatRecordingService?.recordAssistantTurn(recordArgs);
|
||||
}
|
||||
}
|
||||
|
||||
// Mid-stream failure recovery: if the upstream stream threw (typical on
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue