Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions src/node/services/continuousCompactionJournal.ts
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,7 @@ export async function rebuildContinuousPrefix(
),
workspaceId,
messagesWithSentinel: addInterruptedSentinel(prepared.providerRequestMessages),
replayReceiptMessages: prepared.activeContextMessages,
postCompactionAttachments: journal.postCompactionAttachments,
deferLoadingToolNames: deferred,
});
Expand Down
1 change: 1 addition & 0 deletions src/node/services/continuousCompactionSummary.ts
Original file line number Diff line number Diff line change
Expand Up @@ -154,6 +154,7 @@ export async function summarizeContinuousCompaction(args: {
);
const messages = await prepareMessagesForProvider({
messagesWithSentinel: addInterruptedSentinel(prepared.providerRequestMessages),
replayReceiptMessages: prepared.activeContextMessages,
Comment thread
ThomasK33 marked this conversation as resolved.
effectiveAgentId: "compact",
toolNamesForSentinel: [],
postCompactionAttachments: null,
Expand Down
22 changes: 21 additions & 1 deletion src/node/services/historyService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Comment thread
ThomasK33 marked this conversation as resolved.

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 (
Expand Down
10 changes: 9 additions & 1 deletion src/node/services/messagePipeline.ts
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,12 @@ export interface PrepareMessagesOptions {
* become the same tool_reference output the live turn sent.
*/
deferLoadingToolNames?: ReadonlySet<string>;
/**
* 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[];
}

/**
Expand Down Expand Up @@ -122,6 +128,7 @@ export async function prepareMessagesForProvider(
anthropicCacheTtl,
workspaceId,
deferLoadingToolNames,
replayReceiptMessages,
} = opts;

// --- XumMessage-level transforms ---
Expand Down Expand Up @@ -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;

Expand Down
138 changes: 137 additions & 1 deletion src/node/services/streamManager.continuousCompaction.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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"];
Expand Down Expand Up @@ -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",
Expand Down
Loading
Loading