diff --git a/src/node/services/continuousCompactionJournal.ts b/src/node/services/continuousCompactionJournal.ts index 1373a765860..fd40ce1ce12 100644 --- a/src/node/services/continuousCompactionJournal.ts +++ b/src/node/services/continuousCompactionJournal.ts @@ -82,6 +82,7 @@ export async function rebuildContinuousPrefix( ), workspaceId, messagesWithSentinel: addInterruptedSentinel(prepared.providerRequestMessages), + replayReceiptMessages: prepared.activeContextMessages, postCompactionAttachments: journal.postCompactionAttachments, deferLoadingToolNames: deferred, }); diff --git a/src/node/services/continuousCompactionSummary.ts b/src/node/services/continuousCompactionSummary.ts index 91d762303d1..cd33f276ac7 100644 --- a/src/node/services/continuousCompactionSummary.ts +++ b/src/node/services/continuousCompactionSummary.ts @@ -154,6 +154,7 @@ export async function summarizeContinuousCompaction(args: { ); const messages = await prepareMessagesForProvider({ messagesWithSentinel: addInterruptedSentinel(prepared.providerRequestMessages), + replayReceiptMessages: prepared.activeContextMessages, effectiveAgentId: "compact", toolNamesForSentinel: [], postCompactionAttachments: null, diff --git a/src/node/services/historyService.ts b/src/node/services/historyService.ts index 53dc9855d8c..f665624a448 100644 --- a/src/node/services/historyService.ts +++ b/src/node/services/historyService.ts @@ -3543,13 +3543,33 @@ export class HistoryService { hadErrorMetadata && !commitWorthy && !hasDurableRefusalMetadata && - (existingMessage.parts?.length ?? 0) === 0; + (existingMessage.parts?.length ?? 0) === 0 && + // A row that already holds the replay receipt keeps it (a crash after the receipt + // write, before this partial was deleted, #5886). + existingMessage.metadata?.anthropicThinkingReplay !== "off"; + + // #5886: the Anthropic thinking-repair receipt (MuxMetadata.anthropicThinkingReplay) + // must outlive a turn that produced no output, or the next turn replays the removed + // thinking and pays one more 400. Add only the receipt to the stored row: its parts + // stay as they are, and an empty row never reaches the provider or shows as a reply. + const keepsReplayReceipt = + !shouldCommit && + partial.metadata?.anthropicThinkingReplay === "off" && + existingMessage.metadata?.anthropicThinkingReplay !== "off"; if (shouldCommit) { const updateResult = await this.updateHistoryUnderWriteLock(workspaceId, partial); if (!updateResult.success) { return updateResult; } + } else if (keepsReplayReceipt) { + const updateResult = await this.updateHistoryUnderWriteLock(workspaceId, { + ...existingMessage, + metadata: { ...existingMessage.metadata, anthropicThinkingReplay: "off" }, + }); + if (!updateResult.success) { + return updateResult; + } } else if (shouldDeleteErroredPlaceholder) { const deleteMessageResult = await this.deleteMessageUnderWriteLock(workspaceId, partial.id); if ( diff --git a/src/node/services/messagePipeline.ts b/src/node/services/messagePipeline.ts index 81a7f68d446..ca304d0215f 100644 --- a/src/node/services/messagePipeline.ts +++ b/src/node/services/messagePipeline.ts @@ -77,6 +77,12 @@ export interface PrepareMessagesOptions { * become the same tool_reference output the live turn sent. */ deferLoadingToolNames?: ReadonlySet; + /** + * Rows searched for the Anthropic thinking-repair receipt; defaults to messagesWithSentinel. + * The empty-row filter drops a repaired turn that produced no output, so callers that + * filter pass the rows from before it (#5886). + */ + replayReceiptMessages?: MuxMessage[]; } /** @@ -122,6 +128,7 @@ export async function prepareMessagesForProvider( anthropicCacheTtl, workspaceId, deferLoadingToolNames, + replayReceiptMessages, } = opts; // --- XumMessage-level transforms --- @@ -228,7 +235,8 @@ export async function prepareMessagesForProvider( // context segment. Same strip as the one-request repair (stripReasoningReplay), which also // covers adaptive thinking, where transformModelMessages ignores anthropicStripReasoning. const segmentMessages = - providerForMessages === "anthropic" && hasAnthropicReplayReceipt(messagesWithSentinel) + providerForMessages === "anthropic" && + hasAnthropicReplayReceipt(replayReceiptMessages ?? messagesWithSentinel) ? stripReasoningReplay(transformedMessages, "anthropic") : transformedMessages; diff --git a/src/node/services/streamManager.continuousCompaction.test.ts b/src/node/services/streamManager.continuousCompaction.test.ts index ae0c479f9e8..896a34dc153 100644 --- a/src/node/services/streamManager.continuousCompaction.test.ts +++ b/src/node/services/streamManager.continuousCompaction.test.ts @@ -13,7 +13,7 @@ import { historyWriteLockPath, removeSessionDirUnderMemoryLocks } from "./worksp import { createAnthropic } from "@ai-sdk/anthropic"; import { readFile, writeFile } from "node:fs/promises"; import assert from "@/common/utils/assert"; -import { createMuxMessage } from "@/common/types/message"; +import { createMuxMessage, type MuxMessage } from "@/common/types/message"; import { ContinuousCompactionJournalSchema, type ContinuousCompactionJournal, @@ -933,6 +933,50 @@ describe("continuous prefix prepareStep and journal", () => { } ); + it("prefix replay honors a thinking-repair receipt on an empty kept row (#5886)", async () => { + const journal = journalFixture(); + journal.preparation.modelString = "anthropic:claude-opus-5-5"; + journal.preparation.effectiveThinkingLevel = "high"; + const signed = createMuxMessage("signed", "assistant", "", undefined, [ + { + type: "reasoning", + text: "removed thinking", + providerOptions: { anthropic: { signature: "sig-removed" } }, + }, + { type: "text", text: "kept answer" }, + ]); + // A repaired turn that failed before output: commitPartial kept only its receipt. + const receiptOnly: MuxMessage = { + id: "repaired-no-output", + role: "assistant", + metadata: { anthropicThinkingReplay: "off" }, + parts: [], + }; + journal.prefixSourceRows = [ + journal.boundary, + createMuxMessage("user-1", "user", "first"), + signed, + createMuxMessage("user-2", "user", "second"), + receiptOnly, + createMuxMessage("user-3", "user", "third"), + ]; + const { deferLoadingToolNames: _deferred, ...preparation } = journal.preparation; + const expected = await assemblePromptPayload({ + ...preparation, + workspaceId, + history: journal.prefixSourceRows, + systemMessage: "", + postCompactionAttachments: journal.postCompactionAttachments, + }); + const actual = (await rebuildContinuousPrefix(journal, workspaceId)).filter( + (message) => message.role !== "system" + ); + expect(actual).toEqual(expected.messages.filter((message) => message.role !== "system")); + const serialized = JSON.stringify(actual); + expect(serialized).toContain("kept answer"); + expect(serialized).not.toContain("removed thinking"); + }); + it("prefix replay keeps native tool search results as tool references (#5262)", async () => { const journal = journalFixture(); journal.preparation.deferLoadingToolNames = ["slack_send_message"]; @@ -1180,6 +1224,98 @@ describe("continuous prefix prepareStep and journal", () => { }); } + it("a fallback over a consumed prefix sends none of that prefix's thinking (#5886)", async () => { + const { swap } = await setup(); + swap.journal.liveTailCopySpec.partIndex = 0; + // The swapped prefix predates the turn's thinking-repair receipt. + const systemCount = swap.prefix.filter((message) => message.role === "system").length; + swap.prefix = [ + ...swap.prefix.slice(0, systemCount), + { + role: "assistant", + content: [ + { + type: "reasoning", + text: "removed prefix thinking", + providerOptions: { anthropic: { signature: "sig-removed" } }, + }, + { type: "text", text: "kept prefix answer" }, + ], + }, + ...swap.prefix.slice(systemCount), + ]; + const nextModel = "anthropic:fallback-model"; + const rebuilt = createMuxMessage("live", "assistant", "", { stepStartPartIndices: [0] }); + rebuilt.parts = [ + { type: "text", text: "retained step" }, + { + type: "dynamic-tool", + toolCallId: "keep", + toolName: "bash", + state: "output-available", + input: {}, + output: { success: true }, + }, + { type: "text", text: "refused response after swap" }, + ]; + const { deferLoadingToolNames: _deferred, ...preparation } = swap.journal.preparation; + const payload = await assemblePromptPayload({ + ...preparation, + modelString: nextModel, + systemMessage: "Fresh fallback system", + workspaceId, + history: [createMuxMessage("prompt", "user", "original request"), rebuilt], + }); + let firstFallbackStep: ai.ModelMessage[] | undefined; + const turn = await startLiveTurn({ + requestOptions: { + initialMetadata: { anthropicThinkingReplay: "off" }, + modelFallback: { + chain: [nextModel], + prepare: () => + Promise.resolve({ + success: true as const, + data: { + model, + modelString: nextModel, + messages: payload.messages, + system: payload.system, + tools: payload.tools, + thinkingLevel: "off" as const, + }, + }), + }, + }, + attempts: [ + async function* (options, manager) { + expect(manager.setPrefixSwap(workspaceId, swap)).toBe(true); + await prepareStepForTests(options, originalMessages); + yield { type: "start-step" }; + yield { type: "text-delta", text: "retained step" }; + yield { type: "finish-step", usage: { inputTokens: 1, outputTokens: 1, totalTokens: 2 } }; + yield { type: "finish", finishReason: "content-filter" }; + }, + async function* (options) { + // The messages the fallback's first provider request actually carries. + const step = await prepareStepForTests(options, options.messages ?? [], 0); + firstFallbackStep = step?.messages ?? options.messages; + yield* answer(); + }, + ], + }); + await turn.completion; + + expect(turn.calls).toHaveLength(2); + // The fallback request carries the swapped prefix as recorded (thinking included): + // prepareStep's prefix-swap guard (#5086) removes it before the provider call. + const requested = JSON.stringify(turn.calls[1]?.messages); + expect(requested).toContain("kept prefix answer"); + expect(requested).toContain("removed prefix thinking"); + const sent = JSON.stringify(firstFallbackStep); + expect(sent).toContain("kept prefix answer"); + expect(sent).not.toContain("removed prefix thinking"); + }); + for (const family of ["anthropic", "openai"]) { for (const mode of [ "pending", diff --git a/src/node/services/streamManager.errorRecovery.test.ts b/src/node/services/streamManager.errorRecovery.test.ts index 9825197938e..7c806aff725 100644 --- a/src/node/services/streamManager.errorRecovery.test.ts +++ b/src/node/services/streamManager.errorRecovery.test.ts @@ -1,5 +1,5 @@ import { describe, test, expect, spyOn } from "bun:test"; -import type { TurnEngineEvent } from "./streamManager"; +import type { TurnEngineEvent, TurnExecutionOptions } from "./streamManager"; import * as aiSdk from "ai"; import { APICallError, @@ -12,7 +12,8 @@ import { createAnthropic } from "@ai-sdk/anthropic"; import { createOpenAI } from "@ai-sdk/openai"; import { createStreamManagerForTests, fakeStreamText } from "./streamManager.testHarness"; import { OPENAI_RESPONSES_BASE_URL_HINT } from "./utils/openAIResponsesBaseUrlHint"; -import type { MuxMetadata } from "@/common/types/message"; +import { createMuxMessage, type MuxMessage, type MuxMetadata } from "@/common/types/message"; +import { assemblePromptPayload } from "./turnContextAssembler"; import { log } from "./log"; import { installStreamManagerTestHistory, @@ -23,7 +24,9 @@ import { appendPartialAssistantForTests, createStreamResultForTests, prepareStepForTests, + REFUSAL_FINISH, } from "./streamManager.suite.testHarness"; +import { Ok } from "@/common/types/result"; installStreamManagerTestHistory(); @@ -134,6 +137,7 @@ function createRecoveryHarness() { initialMetadata?: Partial; /** SDK-reported total usage for every attempt of this turn. */ streamUsage?: unknown; + requestOptions?: Partial; }) { const historySequence = input.historySequence ?? 1; const messageId = `${input.workspaceId}-${historySequence}`; @@ -154,6 +158,7 @@ function createRecoveryHarness() { providerOptions: input.providerOptions, ...(input.initialMetadata != null ? { initialMetadata: input.initialMetadata } : {}), providedRuntimeTempDir: "", + ...input.requestOptions, }) ); if (!result.success) throw new Error(`Expected stream to start: ${JSON.stringify(result)}`); @@ -1390,6 +1395,238 @@ describe("StreamManager - Anthropic thinking signature recovery", () => { expect(errorPartial?.metadata?.anthropicThinkingReplay).toBe("off"); }); + // #5886: the receipt must survive a turn that produced no output, or the next turn + // puts the removed thinking back and pays one more 400. + describe("receipt on a turn without output", () => { + const signedEarlierTurn = (): MuxMessage => + createMuxMessage("earlier-assistant", "assistant", "", { historySequence: 1 }, [ + { + type: "reasoning", + text: "earlier thinking", + providerOptions: { anthropic: { signature: "sig-bound-to-old-prefix" } }, + }, + { type: "text", text: "earlier answer" }, + ]); + async function seedEarlierTurn(workspaceId: string) { + for (const message of [ + createMuxMessage("earlier-user", "user", "earlier", { historySequence: 0 }), + signedEarlierTurn(), + createMuxMessage("now-user", "user", "now", { historySequence: 2 }), + ]) { + const appended = await historyService.appendToHistory(workspaceId, message); + if (!appended.success) throw new Error(appended.error); + } + } + /** Thinking parts the next Anthropic turn would replay from the committed history. */ + async function nextTurnThinking(workspaceId: string) { + const history = await historyService.getHistoryFromLatestBoundary(workspaceId); + if (!history.success) throw new Error(history.error); + const payload = await assemblePromptPayload({ + history: [ + ...history.data, + createMuxMessage("next-user", "user", "next", { historySequence: 4 }), + ], + systemMessage: "system", + modelString: "anthropic:claude-opus-5-5", + providerForMessages: "anthropic", + effectiveThinkingLevel: "medium", + effectiveAgentId: "exec", + toolNamesForSentinel: [], + workspaceId, + }); + const serialized = JSON.stringify(payload.messages); + // The earlier answer itself must stay: only its thinking goes. + expect(serialized).toContain("earlier answer"); + return payload.messages.flatMap((message) => + message.role === "assistant" && Array.isArray(message.content) + ? message.content.filter((part) => part.type === "reasoning") + : [] + ); + } + + test("a failed step-0 retry keeps the receipt through commitPartial", async () => { + const workspaceId = "anthropic-receipt-failed-retry"; + await seedEarlierTurn(workspaceId); + const harness = createRecoveryHarness(); + const { model } = scriptedAnthropicModel([rejection, rejection]); + + await harness.run({ + workspaceId, + historySequence: 3, + model, + messages: messages(), + attempts: [realSdkAttempt, realSdkAttempt], + }); + expect(harness.errors()).toHaveLength(1); + const committed = await historyService.commitPartial(workspaceId); + if (!committed.success) throw new Error(committed.error); + + expect(await nextTurnThinking(workspaceId)).toEqual([]); + }); + + test("a partial left by a crash after the repair keeps the receipt", async () => { + const workspaceId = "anthropic-receipt-crash"; + await seedEarlierTurn(workspaceId); + await appendPartialAssistantForTests(workspaceId, "crashed-turn", 3); + // What flushPartialWrite left on disk before the retry request: no parts, no error. + const written = await historyService.writePartial(workspaceId, { + id: "crashed-turn", + role: "assistant", + metadata: { historySequence: 3, partial: true, anthropicThinkingReplay: "off" }, + parts: [], + }); + if (!written.success) throw new Error(written.error); + const committed = await historyService.commitPartial(workspaceId); + if (!committed.success) throw new Error(committed.error); + + expect(await nextTurnThinking(workspaceId)).toEqual([]); + }); + + test("a crash after the receipt reached the row keeps it on the next commit", async () => { + const workspaceId = "anthropic-receipt-crash-after-row"; + await seedEarlierTurn(workspaceId); + // commitPartial already wrote the receipt to the row, then crashed before + // deleting the errored partial. + const placeholder = await historyService.appendToHistory(workspaceId, { + id: "repaired-turn", + role: "assistant", + metadata: { historySequence: 3, partial: true, anthropicThinkingReplay: "off" }, + parts: [], + }); + if (!placeholder.success) throw new Error(placeholder.error); + const written = await historyService.writePartial(workspaceId, { + id: "repaired-turn", + role: "assistant", + metadata: { + historySequence: 3, + partial: true, + anthropicThinkingReplay: "off", + error: "rejected again", + errorType: "reasoning_rejected", + }, + parts: [], + }); + if (!written.success) throw new Error(written.error); + const committed = await historyService.commitPartial(workspaceId); + if (!committed.success) throw new Error(committed.error); + + expect(await nextTurnThinking(workspaceId)).toEqual([]); + }); + + test("without a receipt, an empty failed turn still leaves no row", async () => { + const workspaceId = "anthropic-no-receipt-failed-turn"; + await seedEarlierTurn(workspaceId); + await appendPartialAssistantForTests(workspaceId, "failed-turn", 3); + const written = await historyService.writePartial(workspaceId, { + id: "failed-turn", + role: "assistant", + metadata: { historySequence: 3, partial: true, error: "boom", errorType: "unknown" }, + parts: [], + }); + if (!written.success) throw new Error(written.error); + const committed = await historyService.commitPartial(workspaceId); + if (!committed.success) throw new Error(committed.error); + + const history = await historyService.getHistoryFromLatestBoundary(workspaceId); + if (!history.success) throw new Error(history.error); + expect(history.data.map((message) => message.id)).not.toContain("failed-turn"); + expect(await nextTurnThinking(workspaceId)).toHaveLength(1); + }); + }); + + test("a step-0 thinking rebuild after the repair sends no thinking", async () => { + // #5886: the rebuild reads history, which does not hold this turn's receipt yet. + const thinkingOverrideState: NonNullable = { + applied: "medium", + }; + const harness = createRecoveryHarness(); + const { model, requestBlockTypes } = scriptedAnthropicModel([ + () => { + // The user changes the thinking level while the rejected request is in flight. + thinkingOverrideState.pending = "high"; + return rejection(); + }, + success, + ]); + const preparedSdkAttempt: Attempt = async function* (options) { + const prepared = await prepareStepForTests(options, options.messages ?? [], 0); + yield* aiSdk.streamText({ + model: options.model, + messages: prepared?.messages ?? options.messages ?? [], + maxRetries: 0, + }).fullStream; + }; + + await harness.run({ + workspaceId: "anthropic-receipt-rebuild", + model, + messages: messages(), + attempts: [preparedSdkAttempt, preparedSdkAttempt], + providerOptions: {}, + requestOptions: { + thinkingOverrideState, + rebuildProviderOptionsForThinkingLevel: (level) => ({ + effectiveLevel: level, + providerOptions: {}, + }), + // History-based rebuild: still holds the earlier signed thinking. + rebuildFirstStepForThinkingLevel: () => Promise.resolve(messages()), + }, + }); + + expect(harness.streamEnds()).toHaveLength(1); + expect(thinkingOverrideState.applied).toBe("high"); + expect(requestBlockTypes).toHaveLength(2); + expect(requestBlockTypes[1]).not.toContain("thinking"); + expect(requestBlockTypes[1]).toContain("tool_use"); + }); + + test("a model fallback after the repair sends no thinking", async () => { + // #5886: prepare() rebuilds the fallback request from history, which does not hold + // this turn's receipt yet. + const harness = createRecoveryHarness(); + const { model } = scriptedAnthropicModel([rejection]); + const refusedRetry: Attempt = async function* () { + await Promise.resolve(); + yield { type: "start-step" }; + yield { type: "finish-step", usage: TEST_USAGE, finishReason: "content-filter" }; + yield REFUSAL_FINISH; + }; + const fallbackModel = createTestLanguageModel("claude-sonnet-5", "anthropic.messages"); + + const { calls } = await harness.run({ + workspaceId: "anthropic-receipt-fallback", + model, + modelString: "anthropic:claude-opus-5-5", + messages: messages(), + attempts: [realSdkAttempt, refusedRetry, textAttempt("fallback answer")], + requestOptions: { + modelFallback: { + chain: ["anthropic:claude-sonnet-5"], + prepare: (modelString) => + Promise.resolve( + Ok({ + model: fallbackModel, + modelString, + // History-based rebuild: still holds the earlier signed thinking. + messages: messages(), + system: "system", + tools: undefined, + thinkingLevel: "medium" as const, + }) + ), + }, + }, + }); + + expect(calls).toHaveLength(3); + expect(calls[2]?.model).toBe(fallbackModel); + const fallbackParts = (calls[2]?.messages ?? []).flatMap((message) => + message.role === "assistant" && Array.isArray(message.content) ? message.content : [] + ); + expect(fallbackParts.map((part) => part.type)).toEqual(["tool-call"]); + }); + test("logs prefix-binding thinking drops for the step at info, without paths", async () => { const harness = createRecoveryHarness(); const [messageStart, ...rest] = successEvents; diff --git a/src/node/services/streamManager.ts b/src/node/services/streamManager.ts index a5045961277..9f65ba68e3d 100644 --- a/src/node/services/streamManager.ts +++ b/src/node/services/streamManager.ts @@ -465,6 +465,8 @@ interface StepMessageTracker { * put the removed blocks back (preserved thinking). */ anthropicThinkingStripped?: boolean; + /** True once this turn holds the Anthropic thinking-repair receipt (reads initialMetadata). */ + hasAnthropicReplayReceipt?: () => boolean; } interface StreamRequestConfig { stopCause?: StreamStopCause; @@ -3158,9 +3160,15 @@ export class StreamManager { appliedLevel, thinkingOverride ); + // #5886: the rebuild reads history, which does not hold this turn's + // thinking-repair receipt yet. Keep the removed thinking out. + const receiptApplied = + stepTracker?.hasAnthropicReplayReceipt?.() === true + ? stripReasoningReplay(rebuilt, "anthropic") + : rebuilt; // Same per-step transforms the construction-time messages receive. rebuiltFirstStepMessages = dedupeNativeToolReferences( - await transformStepMessages(rebuilt) + await transformStepMessages(receiptApplied) ); if (stepTracker) { stepTracker.latestMessages = rebuiltFirstStepMessages; @@ -3423,6 +3431,9 @@ export class StreamManager { cumulativeUsage: { inputTokens: 0, outputTokens: 0, totalTokens: 0 }, cumulativeProviderMetadata: undefined, }; + // Single source for same-turn rebuilds: the receipt lives on initialMetadata (#5886). + stepTracker.hasAnthropicReplayReceipt = () => + streamInfo.initialMetadata?.anthropicThinkingReplay === "off"; // Mid-turn thinking override: route applied levels into this stream's // metadata (partials, stream-end, final assistant message). Wired before @@ -4181,7 +4192,11 @@ export class StreamManager { const nextRequest = this.buildStreamRequestConfig({ model: prepared.data.model, modelString: prepared.data.modelString, - messages: prepared.data.messages, + // #5886: prepare() rebuilds from history, which does not hold this turn's + // thinking-repair receipt yet. Keep the removed thinking out. + messages: streamInfo.stepTracker.hasAnthropicReplayReceipt?.() + ? stripReasoningReplay(prepared.data.messages, "anthropic") + : prepared.data.messages, system: prepared.data.system, tools: prepared.data.tools, providerOptions: prepared.data.providerOptions, @@ -5954,9 +5969,8 @@ export class StreamManager { // later block, so the next turns must keep it out (messagePipeline reads this // receipt). initialMetadata feeds the partial, error partial and final row alike // (as the model-fallback record does). Persist before the retry request goes out so - // a crash mid-retry cannot forget the strip. A step-0 repair has no parts yet, and - // commitPartial drops an empty partial, so a crash or a second rejection before any - // output loses the receipt: the next turn then pays one more retry (fail-safe). + // a crash mid-retry cannot forget the strip. A step-0 repair has no parts yet: + // commitPartial keeps the receipt on the empty row (#5886). streamInfo.initialMetadata = { ...streamInfo.initialMetadata, anthropicThinkingReplay: "off", diff --git a/src/node/services/turnContextAssembler.ts b/src/node/services/turnContextAssembler.ts index 11c6b267edc..444aa3a58bc 100644 --- a/src/node/services/turnContextAssembler.ts +++ b/src/node/services/turnContextAssembler.ts @@ -250,6 +250,7 @@ export async function assemblePromptPayload( ? tagUserRowsWithHistoryItemIds(prepared.providerRequestMessages) : prepared.providerRequestMessages ), + replayReceiptMessages: prepared.activeContextMessages, effectiveAgentId: options.effectiveAgentId, toolNamesForSentinel: options.toolNamesForSentinel, planContentForTransition: options.planContentForTransition,