feat(ai): auto-marked signing anthropic-messages proxies from 400

The anthropic-messages transport now self-heals on the first '400 Invalid signature in thinking block': demote every unsigned thinking block in the current request, retry once, and pin the (baseUrl, modelId) as signing in the provider session state so subsequent turns skip the round-trip. The retried assistant message surfaces disabledFeatures: ['unsigned-thinking-replay'] so UIs can nudge the user toward the permanent 'compat.replayUnsignedThinking: false' setting in models.yml.\n\nFixes #4297
This commit is contained in:
roboomp
2026-07-02 11:54:47 +00:00
parent 989ac98a6b
commit c8d6fb5580
3 changed files with 327 additions and 2 deletions
+1 -1
View File
@@ -4,7 +4,7 @@
### Added
- Added an actionable remediation hint on the Anthropic `400 Invalid `signature` in `thinking` block` error for custom `anthropic-messages` providers that have not overridden `compat.replayUnsignedThinking`, naming the provider and the exact `models.yml` knob to flip. ([#4297](https://github.com/can1357/oh-my-pi/issues/4297))
- Added a runtime signing-endpoint auto-detect on `anthropic-messages`: when an unmarked custom proxy returns `400 Invalid `signature` in `thinking` block`, the transport demotes every unsigned thinking block in the request, retries once, and pins the (baseUrl, modelId) as signing in the provider session state so subsequent turns skip the round-trip. The successful assistant message surfaces `disabledFeatures: ["unsigned-thinking-replay"]` so UIs can prompt the user to persist the change with `compat.replayUnsignedThinking: false` in `models.yml`. Includes an actionable remediation hint on the raw `400` when the auto-retry can't run. ([#4297](https://github.com/can1357/oh-my-pi/issues/4297))
## [16.3.1] - 2026-07-02
+59 -1
View File
@@ -329,15 +329,26 @@ const ANTHROPIC_PROVIDER_SESSION_STATE_KEY = "anthropic-messages";
type AnthropicProviderSessionState = ProviderSessionState & {
strictToolsDisabled: boolean;
fastModeDisabled: boolean;
/**
* Runtime-learned: this endpoint returned `400 Invalid signature in
* thinking block` for a replayed unsigned thinking block, so it must be
* treated as a signing proxy from now on. All subsequent requests demote
* unsigned thinking to text for this (baseUrl, modelId), same behavior as
* an explicit `compat.replayUnsignedThinking: false`. Cleared on session
* close.
*/
replayUnsignedThinkingDisabled: boolean;
};
function createAnthropicProviderSessionState(): AnthropicProviderSessionState {
const state: AnthropicProviderSessionState = {
strictToolsDisabled: false,
fastModeDisabled: false,
replayUnsignedThinkingDisabled: false,
close: () => {
state.strictToolsDisabled = false;
state.fastModeDisabled = false;
state.replayUnsignedThinkingDisabled = false;
},
};
return state;
@@ -1724,6 +1735,7 @@ const streamAnthropicOnce = (
let disableStrictTools =
(providerSessionState?.strictToolsDisabled ?? false) || (model.compat?.disableStrictTools ?? false);
let dropFastMode = providerSessionState?.fastModeDisabled ?? false;
let forceDemoteUnsignedThinking = providerSessionState?.replayUnsignedThinkingDisabled ?? false;
const mergedCallerHeaders = mergeHeaders(model.headers, options?.headers);
const umansGatewayWebSearchHeader = getUmansWebSearchHeader(model, mergedCallerHeaders);
@@ -1834,6 +1846,7 @@ const streamAnthropicOnce = (
options,
disableStrictTools,
umansGatewayWebSearchHeader !== undefined,
forceDemoteUnsignedThinking,
);
if (disableStrictTools) {
dropAnthropicStrictTools(nextParams);
@@ -2403,6 +2416,39 @@ const streamAnthropicOnce = (
firstTokenTime = undefined;
continue;
}
if (
!forceDemoteUnsignedThinking &&
firstTokenTime === undefined &&
!streamedReplayUnsafeContent &&
isInvalidThinkingSignatureError(
streamFailure instanceof Error ? streamFailure.message : String(streamFailure),
)
) {
logger.warn(
"anthropic: signing proxy detected (Invalid signature in thinking block), demoting unsigned thinking and retrying",
{
provider: model.provider,
model: model.id,
baseUrl,
error: streamFailure instanceof Error ? streamFailure.message : String(streamFailure),
},
);
if (providerSessionState) {
providerSessionState.replayUnsignedThinkingDisabled = true;
}
forceDemoteUnsignedThinking = true;
params = await prepareParams();
providerRetryAttempt = 0;
output.content.length = 0;
output.model = model.id;
output.responseId = undefined;
output.errorMessage = undefined;
output.providerPayload = undefined;
output.usage = createEmptyUsage(copilotDynamicHeaders?.premiumRequests);
output.stopReason = "stop";
firstTokenTime = undefined;
continue;
}
if (
!dropFastMode &&
model.provider === "anthropic" &&
@@ -2479,6 +2525,9 @@ const streamAnthropicOnce = (
if (dropFastMode && model.provider === "anthropic" && options?.serviceTier === "priority") {
output.disabledFeatures = [...(output.disabledFeatures ?? []), "priority"];
}
if (forceDemoteUnsignedThinking && model.compat.replayUnsignedThinking) {
output.disabledFeatures = [...(output.disabledFeatures ?? []), "unsigned-thinking-replay"];
}
stream.push({ type: "done", reason: output.stopReason, message: output });
stream.end();
} catch (error) {
@@ -3051,7 +3100,16 @@ function buildParams(
options?: AnthropicOptions,
disableStrictTools = false,
useUmansGatewayWebSearch = false,
forceDemoteUnsignedThinking = false,
): MessageCreateParamsStreaming {
// A session-scoped auto-demote (learned from a live signing 400) clones the
// resolved compat with `replayUnsignedThinking: false` so every subsequent
// downstream read (convertAnthropicMessages, transformMessages) sees the
// demoted default without mutating the shared `model` reference.
const effectiveModel =
forceDemoteUnsignedThinking && model.compat.replayUnsignedThinking
? { ...model, compat: { ...model.compat, replayUnsignedThinking: false } }
: model;
const { cacheControl } = getCacheControl(model, options?.cacheRetention, isOAuthToken);
// Pre-compute system blocks so they occupy the right slot in the serialized body.
@@ -3177,7 +3235,7 @@ function buildParams(
// metadata → max_tokens → thinking → context_management → output_config → stream.
const params: MessageCreateParamsStreaming = {
model: options?.requestModelId ?? model.requestModelId ?? model.id,
messages: convertAnthropicMessages(context.messages, model, isOAuthToken, {
messages: convertAnthropicMessages(context.messages, effectiveModel, isOAuthToken, {
serverSideFallbackEnabled: !!options?.fallbacks?.length,
}),
...(systemBlocks && { system: systemBlocks }),
@@ -0,0 +1,267 @@
import { afterEach, describe, expect, it, vi } from "bun:test";
import { streamAnthropic } from "@oh-my-pi/pi-ai/providers/anthropic";
import { AnthropicMessages } from "@oh-my-pi/pi-ai/providers/anthropic-client";
import type {
AssistantMessage,
AssistantMessageEvent,
Context,
Message,
Model,
ProviderSessionState,
} from "@oh-my-pi/pi-ai/types";
import { buildModel } from "@oh-my-pi/pi-catalog/build";
/**
* Regression for #4297 — the anthropic-messages transport auto-heals the very
* first `400 Invalid signature in thinking block` from an unmarked custom
* signing proxy: demote every unsigned thinking block in the request, retry
* once, and pin the (baseUrl, modelId) as signing in the session state so
* subsequent turns skip the demotion round-trip.
*/
const model: Model<"anthropic-messages"> = buildModel({
id: "cf-anthropic/claude-opus-4-8",
name: "Claude Opus 4.8 via cloudflared",
api: "anthropic-messages",
provider: "cf-anthropic",
baseUrl: "https://opencode.cloudflare.dev/anthropic",
reasoning: true,
input: ["text"],
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
contextWindow: 200_000,
maxTokens: 8_192,
});
const priorTurnContext: Context = {
messages: [
{ role: "user", content: "Summarize README", timestamp: 0 },
{
role: "assistant",
content: [
{ type: "thinking", thinking: "Read the file, then summarise.", thinkingSignature: "" },
{ type: "text", text: "The README covers the CLI." },
],
api: "anthropic-messages",
provider: "cf-anthropic",
model: "cf-anthropic/claude-opus-4-8",
usage: {
input: 0,
output: 0,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 0,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
},
stopReason: "stop",
timestamp: 0,
} satisfies AssistantMessage,
{ role: "user", content: "Translate to French.", timestamp: 0 },
] satisfies Message[],
};
function createSignatureRejection(): Error {
const error = new Error(
'400 {"type":"error","error":{"type":"invalid_request_error","message":"messages.1.content.0: Invalid `signature` in `thinking` block"},"request_id":"req_test"}',
);
Object.assign(error, { status: 400 });
return error;
}
interface AnthropicWireBlock {
type: string;
thinking?: string;
text?: string;
signature?: string;
}
interface AnthropicWireMessage {
role: string;
content: AnthropicWireBlock[] | string;
}
interface CapturedRequestPayload {
messages?: AnthropicWireMessage[];
}
function extractPriorAssistantBlocks(params: unknown): AnthropicWireBlock[] {
if (!params || typeof params !== "object" || !("messages" in params)) return [];
const { messages } = params as CapturedRequestPayload;
if (!Array.isArray(messages)) return [];
for (const msg of messages) {
if (msg.role !== "assistant") continue;
if (typeof msg.content === "string") continue;
return msg.content;
}
return [];
}
const successEvents = [
{
type: "message_start",
message: {
id: "msg_ok",
usage: {
input_tokens: 12,
output_tokens: 0,
cache_read_input_tokens: 0,
cache_creation_input_tokens: 0,
},
},
},
{ type: "content_block_start", index: 0, content_block: { type: "text", text: "" } },
{ type: "content_block_delta", index: 0, delta: { type: "text_delta", text: "Bonjour." } },
{ type: "content_block_stop", index: 0 },
{
type: "message_delta",
delta: { stop_reason: "end_turn" },
usage: {
input_tokens: 12,
output_tokens: 4,
cache_read_input_tokens: 0,
cache_creation_input_tokens: 0,
},
},
{ type: "message_stop" },
] as const;
function successRequest() {
const response = new Response(null, { status: 200, headers: { "request-id": "req_ok" } });
return {
async withResponse() {
return {
data: (async function* () {
for (const event of successEvents) {
yield event;
}
})(),
response,
request_id: response.headers.get("request-id"),
};
},
};
}
function readReplayUnsignedThinkingDisabled(map: Map<string, ProviderSessionState>): boolean | undefined {
for (const [key, value] of map) {
if (!key.startsWith("anthropic-messages")) continue;
if (typeof value !== "object" || value === null) continue;
if (!("replayUnsignedThinkingDisabled" in value)) continue;
const flag = value.replayUnsignedThinkingDisabled;
return typeof flag === "boolean" ? flag : undefined;
}
return undefined;
}
describe("#4297 anthropic-messages runtime signing auto-mark", () => {
afterEach(() => {
vi.restoreAllMocks();
});
it("demotes unsigned thinking, retries, and pins the session on the first signing 400", async () => {
const providerSessionState = new Map<string, ProviderSessionState>();
const capturedPayloads: unknown[] = [];
let attempt = 0;
vi.spyOn(AnthropicMessages.prototype, "create").mockImplementation((params: unknown) => {
attempt += 1;
capturedPayloads.push(params);
if (attempt === 1) {
return {
async withResponse() {
throw createSignatureRejection();
},
} as never;
}
return successRequest() as never;
});
const stream = streamAnthropic(model, priorTurnContext, {
apiKey: "sk-ant-test",
providerSessionState,
});
const events: AssistantMessageEvent[] = [];
for await (const event of stream) {
events.push(event);
}
const result = await stream.result();
expect(attempt).toBe(2);
expect(result.stopReason).toBe("stop");
expect(result.errorMessage).toBeUndefined();
const firstAttemptBlocks = extractPriorAssistantBlocks(capturedPayloads[0]);
const firstThinking = firstAttemptBlocks.find(block => block.type === "thinking");
expect(firstThinking?.signature).toBe("");
expect(firstThinking?.thinking).toBe("Read the file, then summarise.");
const retryBlocks = extractPriorAssistantBlocks(capturedPayloads[1]);
expect(retryBlocks.find(block => block.type === "thinking")).toBeUndefined();
const demotedText = retryBlocks.find(block => block.type === "text");
expect(demotedText?.text).toContain("Read the file, then summarise.");
expect(readReplayUnsignedThinkingDisabled(providerSessionState)).toBe(true);
expect(result.disabledFeatures).toContain("unsigned-thinking-replay");
});
it("pre-demotes unsigned thinking on subsequent turns once the session is pinned", async () => {
const providerSessionState = new Map<string, ProviderSessionState>();
const capturedPayloads: unknown[] = [];
vi.spyOn(AnthropicMessages.prototype, "create").mockImplementation((params: unknown) => {
capturedPayloads.push(params);
return successRequest() as never;
});
// Seed the session state as though a prior turn had already auto-marked
// the endpoint. This mirrors the shape produced by the runtime retry so
// subsequent turns never repeat the 400 round-trip.
providerSessionState.set(`anthropic-messages:${model.baseUrl}\u0000${model.id}`, {
close: () => {},
strictToolsDisabled: false,
fastModeDisabled: false,
replayUnsignedThinkingDisabled: true,
} as ProviderSessionState);
const stream = streamAnthropic(model, priorTurnContext, {
apiKey: "sk-ant-test",
providerSessionState,
});
for await (const _ of stream) {
/* drain */
}
const result = await stream.result();
expect(result.stopReason).toBe("stop");
expect(capturedPayloads.length).toBe(1);
const blocks = extractPriorAssistantBlocks(capturedPayloads[0]);
expect(blocks.find(block => block.type === "thinking")).toBeUndefined();
expect(blocks.find(block => block.type === "text")?.text).toContain("Read the file, then summarise.");
expect(result.disabledFeatures).toContain("unsigned-thinking-replay");
});
it("does not auto-mark on unrelated Anthropic invalid_request_error 400s", async () => {
const providerSessionState = new Map<string, ProviderSessionState>();
let attempt = 0;
vi.spyOn(AnthropicMessages.prototype, "create").mockImplementation(() => {
attempt += 1;
return {
async withResponse() {
const error = new Error(
'400 {"type":"error","error":{"type":"invalid_request_error","message":"Some other validation failure"},"request_id":"req_test"}',
);
Object.assign(error, { status: 400 });
throw error;
},
} as never;
});
const stream = streamAnthropic(model, priorTurnContext, {
apiKey: "sk-ant-test",
providerSessionState,
});
for await (const _ of stream) {
/* drain */
}
const result = await stream.result();
expect(attempt).toBe(1);
expect(result.stopReason).toBe("error");
expect(result.errorMessage).toContain("Some other validation failure");
expect(readReplayUnsignedThinkingDisabled(providerSessionState)).toBe(false);
});
});