fix: fixed tool-call recovery, stream parsing, and image normalization flows

- Fixed tool result validation to reject invalid blocks and append explicit error diagnostics.
- Handled aborted and leaked tool calls by returning abort results and dropping partial leaked output.
- Fixed Anthropic streaming by enforcing strict SSE parsing and content-block lifecycle checks.
- Fixed image handling by normalizing model-context inputs and preserving images on resize failure.
This commit is contained in:
can1357
2026-06-08 14:40:18 +02:00
parent 1e5017bf7d
commit 0890b2be61
9 changed files with 799 additions and 389 deletions
+9
View File
@@ -1,10 +1,19 @@
# Changelog
## [Unreleased]
### Removed
- Removed the `maxToolCallsPerTurn` option from `AgentOptions` and `AgentLoopConfig`, so assistant turns are no longer capped after a configured number of completed tool calls
### Fixed
- Fixed tool result parsing to mark assistant tool outputs with unsupported content block shapes as errors and include a diagnostic text block
- Fixed GPT-5 Harmony leakage handling by recovering valid leaked tool calls when possible and discarding leaked partial assistant output before retrying
- Fixed tool-call cancellation handling so aborted tools are marked aborted with an explicit reason and do not report generic errors
- Fixed tool-call completion so assistant messages on abort keep only completed tool-call blocks and continue processing tool calls when a length stop still included results
- Fixed runs that stopped with reason `length` after returning tool results so execution continues to handle additional tool calls
## [15.10.3] - 2026-06-08
### Added
+292 -130
View File
@@ -20,6 +20,7 @@ import {
createHarmonyAuditEvent,
detectHarmonyLeakInAssistantMessage,
extractHarmonyRemoved,
recoverHarmonyToolCall,
type HarmonyDetection,
type HarmonyRecoveredToolCall,
isHarmonyLeakMitigationTarget,
@@ -32,7 +33,6 @@ import {
finishChatSpan,
finishExecuteToolSpan,
finishInvokeAgentSpan,
fireOnRunEnd,
PiGenAIAttr,
recordSkippedTool,
resolveTelemetry,
@@ -68,6 +68,76 @@ class HarmonyLeakInterruption extends Error {
}
}
type AssistantContentBlock = AssistantMessage["content"][number];
type AssistantToolCallBlock = Extract<AssistantContentBlock, { type: "toolCall" }>;
type CloneableRecord = Record<string, unknown>;
function cloneUnknown(value: unknown): unknown {
if (Array.isArray(value)) return value.map(cloneUnknown);
if (!value || typeof value !== "object") return value;
const source = value as CloneableRecord;
const out: CloneableRecord = {};
for (const [key, child] of Object.entries(source)) {
out[key] = cloneUnknown(child);
}
return out;
}
function cloneToolArguments(args: AssistantToolCallBlock["arguments"]): AssistantToolCallBlock["arguments"] {
return cloneUnknown(args) as AssistantToolCallBlock["arguments"];
}
function snapshotAssistantContentBlock(block: AssistantContentBlock): AssistantContentBlock {
switch (block.type) {
case "text":
return { ...block };
case "thinking":
return { ...block };
case "redactedThinking":
return { ...block };
case "toolCall":
return { ...block, arguments: cloneToolArguments(block.arguments) };
}
}
function snapshotAssistantMessage(message: AssistantMessage): AssistantMessage {
return {
...message,
content: message.content.map(snapshotAssistantContentBlock),
usage: {
...message.usage,
cost: { ...message.usage.cost },
},
disabledFeatures: message.disabledFeatures ? [...message.disabledFeatures] : undefined,
};
}
function snapshotAssistantMessageEvent(event: AssistantMessageEvent): AssistantMessageEvent {
switch (event.type) {
case "start":
return { ...event, partial: snapshotAssistantMessage(event.partial) };
case "text_start":
case "text_delta":
case "text_end":
case "thinking_start":
case "thinking_delta":
case "thinking_end":
case "toolcall_start":
case "toolcall_delta":
return { ...event, partial: snapshotAssistantMessage(event.partial) };
case "toolcall_end":
return {
...event,
toolCall: snapshotAssistantContentBlock(event.toolCall) as AssistantToolCallBlock,
partial: snapshotAssistantMessage(event.partial),
};
case "done":
return { ...event, message: snapshotAssistantMessage(event.message) };
case "error":
return { ...event, error: snapshotAssistantMessage(event.error) };
}
}
/**
* Normalize a value coming back from `tool.execute()` (or its streaming partial-update callback)
* into a structurally valid {@link AgentToolResult}.
@@ -77,7 +147,7 @@ class HarmonyLeakInterruption extends Error {
* (missing `content` array → crash on reload). We coerce at the single boundary where untyped
* results enter the agent loop, so every downstream consumer can rely on the type.
*/
function coerceToolResult(raw: unknown): { result: AgentToolResult<any>; malformed: boolean } {
function coerceToolResult(raw: unknown): { result: AgentToolResult<unknown>; malformed: boolean } {
const rawObj = raw && typeof raw === "object" ? (raw as Record<string, unknown>) : null;
const rawContent = rawObj?.content;
const details = rawObj && "details" in rawObj ? rawObj.details : {};
@@ -98,8 +168,12 @@ function coerceToolResult(raw: unknown): { result: AgentToolResult<any>; malform
}
const content: AgentToolResult["content"] = [];
let invalidBlocks = 0;
for (const block of rawContent) {
if (!block || typeof block !== "object" || !("type" in block)) continue;
if (!block || typeof block !== "object" || !("type" in block)) {
invalidBlocks++;
continue;
}
if (block.type === "text" && typeof (block as { text?: unknown }).text === "string") {
content.push({ type: "text", text: sanitizeText((block as { text: string }).text) });
} else if (
@@ -108,9 +182,20 @@ function coerceToolResult(raw: unknown): { result: AgentToolResult<any>; malform
typeof (block as { mimeType?: unknown }).mimeType === "string"
) {
content.push(block as { type: "image"; data: string; mimeType: string });
} else {
invalidBlocks++;
}
}
return { result: { content, details, ...(explicitError ? { isError: true } : {}) }, malformed: false };
if (invalidBlocks > 0) {
content.push({
type: "text",
text: `Tool returned an invalid result: ${invalidBlocks} content block${invalidBlocks === 1 ? "" : "s"} had an unsupported shape.`,
});
}
return {
result: { content, details, ...(explicitError || invalidBlocks > 0 ? { isError: true } : {}) },
malformed: invalidBlocks > 0,
};
}
/**
@@ -176,7 +261,7 @@ export function agentLoopContinue(
(async () => {
const newMessages: AgentMessage[] = [];
const currentContext: AgentContext = { ...context };
const currentContext: AgentContext = { ...context, messages: [...context.messages] };
stream.push({ type: "agent_start" });
stream.push({ type: "turn_start" });
@@ -211,9 +296,6 @@ function buildAgentEndEvent(
): Extract<AgentEvent, { type: "agent_end" }> {
if (!telemetry) return { type: "agent_end", messages };
const snapshot = telemetry.collector.snapshot({ stepCount });
if (telemetry.collector.markRunEnded()) {
fireOnRunEnd(telemetry, snapshot.summary, snapshot.coverage);
}
return { type: "agent_end", messages, telemetry: snapshot.summary, coverage: snapshot.coverage };
}
@@ -313,22 +395,26 @@ function normalizeMessagesForProvider(
return messages;
}
let changed = false;
const normalized = messages.map(message => {
let hasThinking = false;
for (const message of messages) {
if (message.role !== "assistant" || !Array.isArray(message.content)) continue;
for (const block of message.content) {
if (block.type === "thinking") {
hasThinking = true;
break;
}
}
if (hasThinking) break;
}
if (!hasThinking) return messages;
return messages.map(message => {
if (message.role !== "assistant" || !Array.isArray(message.content)) {
return message;
}
const filtered = message.content.filter(block => block.type !== "thinking");
if (filtered.length === message.content.length) {
return message;
}
changed = true;
return { ...message, content: filtered };
return filtered.length === message.content.length ? message : { ...message, content: filtered };
});
return changed ? normalized : messages;
}
export const INTENT_FIELD = "_i";
@@ -552,6 +638,12 @@ async function runLoopBody(
continue;
}
}
if (recovered) {
message = snapshotAssistantMessage(message);
currentContext.messages.push(message);
stream.push({ type: "message_start", message: snapshotAssistantMessage(message) });
stream.push({ type: "message_end", message: snapshotAssistantMessage(message) });
}
newMessages.push(message);
let steeringMessagesFromExecution: AgentMessage[] | undefined;
@@ -640,6 +732,9 @@ async function runLoopBody(
status: "skipped",
});
}
if (message.stopReason === "length" && toolResults.length > 0) {
hasMoreToolCalls = true;
}
}
stream.push({ type: "turn_end", message, toolResults });
@@ -816,8 +911,27 @@ async function streamAssistantResponse(
let partialMessage: AssistantMessage | null = null;
let addedPartial = false;
const completedToolCallIds = new Set<string>();
const responseIterator = response[Symbol.asyncIterator]();
const finishAbortedStream = async (): Promise<AssistantMessage> => {
try {
await responseIterator.return?.();
} catch {
// Provider cancellation failures cannot change the committed aborted message.
}
const aborted = emitAbortedAssistantMessage(
partialMessage,
addedPartial,
completedToolCallIds,
context,
config,
stream,
requestSignal,
);
await finishChat(aborted);
return aborted;
};
// Set up a single abort race: register the abort listener once for the whole
// stream and reuse the same race promise for every iterator.next() instead of
@@ -826,16 +940,7 @@ async function streamAssistantResponse(
let detachAbortListener: (() => void) | undefined;
if (requestSignal) {
if (requestSignal.aborted) {
const aborted = emitAbortedAssistantMessage(
partialMessage,
addedPartial,
context,
config,
stream,
requestSignal,
);
await finishChat(aborted);
return aborted;
return await finishAbortedStream();
}
const { promise, resolve } = Promise.withResolvers<typeof ABORTED>();
const onAbort = () => resolve(ABORTED);
@@ -850,37 +955,51 @@ async function streamAssistantResponse(
if (abortRacePromise) {
const result = await Promise.race([responseIterator.next(), abortRacePromise]);
if (result === ABORTED) {
responseIterator.return?.()?.catch(() => {});
const aborted = emitAbortedAssistantMessage(
partialMessage,
addedPartial,
context,
config,
stream,
requestSignal,
);
await finishChat(aborted);
return aborted;
return await finishAbortedStream();
}
next = result;
} else {
next = await responseIterator.next();
}
if (requestSignal?.aborted) {
const aborted = emitAbortedAssistantMessage(
partialMessage,
addedPartial,
context,
config,
stream,
requestSignal,
);
await finishChat(aborted);
return aborted;
}
if (next.done) break;
const event = next.value;
if (event.type === "done" || event.type === "error") {
let finalMessage = retainCompletedToolCalls(await response.result(), completedToolCallIds);
if (harmonyMitigationEnabled) {
const detection = detectHarmonyLeakInAssistantMessage(finalMessage);
if (detection) {
const recovered = recoverHarmonyToolCall(finalMessage, detection);
const removed = recovered?.removed ?? extractHarmonyRemoved(finalMessage, detection);
if (addedPartial) {
emitDiscardedHarmonyPartial(
partialMessage,
stream,
`Discarded after GPT-5 Harmony protocol leakage (${signalListLabel(detection.signals)})`,
);
context.messages.pop();
addedPartial = false;
}
throw new HarmonyLeakInterruption(detection, removed, recovered);
}
}
finalMessage = snapshotAssistantMessage(finalMessage);
if (addedPartial) {
context.messages[context.messages.length - 1] = finalMessage;
} else {
context.messages.push(finalMessage);
}
if (!addedPartial) {
stream.push({ type: "message_start", message: snapshotAssistantMessage(finalMessage) });
}
stream.push({ type: "message_end", message: snapshotAssistantMessage(finalMessage) });
await finishChat(finalMessage);
return finalMessage;
}
if (requestSignal?.aborted) {
return await finishAbortedStream();
}
// Yield to the event loop periodically to prevent busy-wait
// when the LLM is streaming chunks faster than the loop can rest.
await yieldIfDue();
@@ -890,7 +1009,7 @@ async function streamAssistantResponse(
partialMessage = event.partial;
context.messages.push(partialMessage);
addedPartial = true;
stream.push({ type: "message_start", message: { ...partialMessage } });
stream.push({ type: "message_start", message: snapshotAssistantMessage(partialMessage) });
break;
case "text_start":
@@ -903,63 +1022,48 @@ async function streamAssistantResponse(
case "toolcall_delta":
case "toolcall_end":
if (partialMessage) {
if (event.type === "toolcall_end") {
completedToolCallIds.add(event.toolCall.id);
}
partialMessage = event.partial;
context.messages[context.messages.length - 1] = partialMessage;
config.onAssistantMessageEvent?.(partialMessage, event);
if (signal?.aborted) {
continue;
}
stream.push({
type: "message_update",
assistantMessageEvent: event,
message: { ...partialMessage },
assistantMessageEvent: snapshotAssistantMessageEvent(event),
message: snapshotAssistantMessage(partialMessage),
});
}
break;
case "done":
case "error": {
const finalMessage = await response.result();
if (harmonyMitigationEnabled) {
const detection = detectHarmonyLeakInAssistantMessage(finalMessage);
if (detection) {
const removed = extractHarmonyRemoved(finalMessage, detection);
if (addedPartial) {
context.messages.pop();
addedPartial = false;
}
throw new HarmonyLeakInterruption(detection, removed);
}
}
if (addedPartial) {
context.messages[context.messages.length - 1] = finalMessage;
} else {
context.messages.push(finalMessage);
}
if (!addedPartial) {
stream.push({ type: "message_start", message: { ...finalMessage } });
}
stream.push({ type: "message_end", message: finalMessage });
await finishChat(finalMessage);
return finalMessage;
}
}
}
} finally {
detachAbortListener?.();
}
const trailing = await response.result();
let trailing = await response.result();
if (harmonyMitigationEnabled) {
const detection = detectHarmonyLeakInAssistantMessage(trailing);
if (detection) {
const recovered = recoverHarmonyToolCall(trailing, detection);
const removed = recovered?.removed ?? extractHarmonyRemoved(trailing, detection);
if (addedPartial) {
emitDiscardedHarmonyPartial(
partialMessage,
stream,
`Discarded after GPT-5 Harmony protocol leakage (${signalListLabel(detection.signals)})`,
);
context.messages.pop();
addedPartial = false;
}
throw new HarmonyLeakInterruption(detection, extractHarmonyRemoved(trailing, detection));
throw new HarmonyLeakInterruption(detection, removed, recovered);
}
}
trailing = snapshotAssistantMessage(trailing);
if (addedPartial) {
context.messages[context.messages.length - 1] = trailing;
stream.push({ type: "message_end", message: snapshotAssistantMessage(trailing) });
}
await finishChat(trailing);
return trailing;
});
@@ -973,6 +1077,30 @@ async function streamAssistantResponse(
}
}
function retainCompletedToolCalls(message: AssistantMessage, completedToolCallIds: ReadonlySet<string>): AssistantMessage {
if (message.stopReason !== "error" && message.stopReason !== "aborted") return message;
let changed = false;
const content = message.content.filter(block => {
if (block.type !== "toolCall") return true;
const keep = completedToolCallIds.has(block.id);
if (!keep) changed = true;
return keep;
});
return changed ? { ...message, content } : message;
}
function emitDiscardedHarmonyPartial(
partialMessage: AssistantMessage | null,
stream: EventStream<AgentEvent, AgentMessage[]>,
errorMessage: string,
): void {
if (!partialMessage) return;
stream.push({
type: "message_end",
message: snapshotAssistantMessage({ ...partialMessage, stopReason: "error", errorMessage }),
});
}
/** Resolve the human-readable reason an abort carried. A caller that aborts via
* `AbortController.abort(reason)` with a string or a non-`AbortError` `Error`
* (e.g. the coding agent's user-interrupt label) gets that text surfaced on the
@@ -991,39 +1119,45 @@ export function abortReasonText(signal: AbortSignal | undefined): string {
function emitAbortedAssistantMessage(
partialMessage: AssistantMessage | null,
addedPartial: boolean,
completedToolCallIds: ReadonlySet<string>,
context: AgentContext,
config: AgentLoopConfig,
stream: EventStream<AgentEvent, AgentMessage[]>,
requestSignal: AbortSignal | undefined,
): AssistantMessage {
const errorMessage = abortReasonText(requestSignal);
const abortedMessage: AssistantMessage = partialMessage
? { ...partialMessage, stopReason: "aborted", errorMessage }
: {
role: "assistant",
content: [],
api: config.model.api,
provider: config.model.provider,
model: config.model.id,
usage: {
input: 0,
output: 0,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 0,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
},
stopReason: "aborted",
errorMessage,
timestamp: Date.now(),
};
const abortedMessage = snapshotAssistantMessage(
retainCompletedToolCalls(
partialMessage
? { ...partialMessage, stopReason: "aborted", errorMessage }
: {
role: "assistant",
content: [],
api: config.model.api,
provider: config.model.provider,
model: config.model.id,
usage: {
input: 0,
output: 0,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 0,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
},
stopReason: "aborted",
errorMessage,
timestamp: Date.now(),
},
completedToolCallIds,
),
);
if (addedPartial) {
context.messages[context.messages.length - 1] = abortedMessage;
} else {
context.messages.push(abortedMessage);
stream.push({ type: "message_start", message: { ...abortedMessage } });
stream.push({ type: "message_start", message: snapshotAssistantMessage(abortedMessage) });
}
stream.push({ type: "message_end", message: abortedMessage });
stream.push({ type: "message_end", message: snapshotAssistantMessage(abortedMessage) });
return abortedMessage;
}
@@ -1061,7 +1195,7 @@ async function executeToolCalls(
: steeringAbortController.signal;
const interruptState = { triggered: false };
let steeringMessages: AgentMessage[] | undefined;
let steeringCheck: Promise<void> | null = null;
let steeringCheckTail: Promise<void> = Promise.resolve();
const records = toolCalls.map(toolCall => ({
toolCall,
@@ -1085,21 +1219,17 @@ async function executeToolCalls(
if (!shouldInterruptImmediately || !getSteeringMessages || interruptState.triggered) {
return;
}
if (steeringCheck) {
await steeringCheck;
return;
}
steeringCheck = (async () => {
const check = steeringCheckTail.then(async () => {
if (interruptState.triggered) return;
const steering = await getSteeringMessages();
if (steering.length > 0) {
steeringMessages = steering;
interruptState.triggered = true;
steeringAbortController.abort();
}
})().finally(() => {
steeringCheck = null;
});
await steeringCheck;
steeringCheckTail = check.catch(() => {});
await check;
};
const emitToolResult = (record: (typeof records)[number], result: AgentToolResult<any>, isError: boolean): void => {
@@ -1171,6 +1301,16 @@ async function executeToolCalls(
}
}
record.args = argsForExecution;
if (toolSignal.aborted) {
record.skipped = true;
recordSkippedTool(telemetry, {
toolCallId: toolCall.id,
toolName: toolCall.name,
status: "aborted",
});
emitToolResult(record, createToolSignalAbortedResult(toolSignal), true);
return;
}
record.started = true;
stream.push({
type: "tool_execution_start",
@@ -1198,6 +1338,11 @@ async function executeToolCalls(
await runInActiveSpan(toolSpan, async () => {
try {
if (!tool) throw new Error(`Tool ${toolCall.name} not found`);
if (toolSignal.aborted) {
result = createToolSignalAbortedResult(toolSignal);
isError = true;
return;
}
let effectiveArgs: Record<string, unknown>;
try {
@@ -1224,8 +1369,15 @@ async function executeToolCalls(
throw new ToolCallBlockedError(beforeResult.reason);
}
}
// Reflect post-hook args so emitted tool results / afterToolCall see what actually executed.
record.args = effectiveArgs;
if (toolSignal.aborted) {
result = createToolSignalAbortedResult(toolSignal);
isError = true;
return;
}
const executionArgs = transformToolCallArguments
? transformToolCallArguments(effectiveArgs, toolCall.name)
: effectiveArgs;
record.args = executionArgs;
const toolContext = getToolContext
? getToolContext({
@@ -1237,14 +1389,14 @@ async function executeToolCalls(
: undefined;
const rawResult = await tool.execute(
toolCall.id,
transformToolCallArguments ? transformToolCallArguments(effectiveArgs, toolCall.name) : effectiveArgs,
executionArgs,
toolSignal,
partialResult => {
stream.push({
type: "tool_execution_update",
toolCallId: toolCall.id,
toolName: toolCall.name,
args: effectiveArgs,
args: executionArgs,
partialResult: coerceToolResult(partialResult).result,
});
},
@@ -1262,7 +1414,7 @@ async function executeToolCalls(
isError = true;
}
if (afterToolCall) {
if (afterToolCall && !toolSignal.aborted) {
try {
const after = await afterToolCall(
{
@@ -1295,6 +1447,7 @@ async function executeToolCalls(
});
const interrupted = interruptState.triggered;
const abortedDuringExecution = toolSignal.aborted && isError;
if (interrupted) {
record.skipped = true;
emitToolResult(record, createSkippedToolResult(), true);
@@ -1305,13 +1458,14 @@ async function executeToolCalls(
const firstTextBlock = result.content?.[0];
const errorMessageForSpan =
caughtError === undefined && isError && firstTextBlock?.type === "text" ? firstTextBlock.text : undefined;
const status = interrupted
? "aborted"
: caughtError instanceof ToolCallBlockedError
? "blocked"
: isError
? "error"
: "ok";
const status =
interrupted || abortedDuringExecution
? "aborted"
: caughtError instanceof ToolCallBlockedError
? "blocked"
: isError
? "error"
: "ok";
finishExecuteToolSpan(telemetry, toolSpan, {
result,
isError,
@@ -1417,6 +1571,14 @@ function createAbortedToolResult(
return toolResultMessage;
}
function createToolSignalAbortedResult(signal: AbortSignal): AgentToolResult<unknown> {
const reason = abortReasonText(signal);
return {
content: [{ type: "text", text: `Tool was not executed because the run was aborted: ${reason}.` }],
details: {},
};
}
function createSkippedToolResult(): AgentToolResult<any> {
return {
content: [{ type: "text", text: "Skipped due to queued user message." }],
+7 -3
View File
@@ -1,13 +1,17 @@
# Changelog
## [Unreleased]
### Changed
- Changed Anthropic retry handling to avoid retrying 4xx responses other than 408 and 429
### Fixed
- Fixed raw Anthropic SSE handling by parsing event frames with strict JSON parsing and matching event-type validation, surfacing malformed frames as stream errors instead of repairing them
- Fixed Anthropic stream envelope handling to reject duplicate `content_block_start` indexes and block deltas/stops for unopened blocks, preventing malformed envelope states from producing partial output
- Fixed Anthropic image conversion to normalize `image/jpg` to `image/jpeg` and emit a placeholder for unsupported image MIME types
- Fixed Anthropic thinking request preparation by clamping `max_tokens` to provider/model limits and adjusting thinking budgets to a valid value
- Fixed the Anthropic stream parser shipping a truncated tool call as a completed turn. When a transport drop cut the SSE stream mid-`tool_use` and a transparent reconnect spliced a fresh message envelope onto the same stream, the duplicate `message_start` was deduped but the orphaned tool block — which never received its `content_block_stop` — survived in the assistant message with its seed `{}` (or partially-parsed) arguments. The terminal stop signal from the reconnect then let it flow through as a normal tool call, so e.g. a `read` dispatched with `{}` failed downstream validation (`path: expected string, received undefined`). The parser now treats any tool block left open at stream end as a truncated envelope and routes it through the existing retry/error path instead of emitting bogus arguments.
### Fixed
- Fixed the Zhipu Coding Plan login prompt advertising a misleading `sk-...` placeholder. Zhipu API keys are formatted `<id>.<secret>` (no `sk-` prefix), so the placeholder now matches the actual format instead of suggesting the wrong shape. ([#2106](https://github.com/can1357/oh-my-pi/issues/2106))
## [15.10.4] - 2026-06-08
+355 -213
View File
@@ -57,7 +57,7 @@ import { AssistantMessageEventStream } from "../utils/event-stream";
import { isFoundryEnabled } from "../utils/foundry";
import { finalizeErrorMessage, type RawHttpRequestDump, rewriteCopilotError } from "../utils/http-inspector";
import { getStreamFirstEventTimeoutMs, getStreamIdleTimeoutMs, iterateWithIdleTimeout } from "../utils/idle-iterator";
import { parseJsonWithRepair, parseStreamingJson, parseStreamingJsonThrottled } from "../utils/json-parse";
import { parseStreamingJsonThrottled } from "../utils/json-parse";
import { parseGitHubCopilotApiKey } from "../utils/oauth/github-copilot";
import { notifyProviderResponse } from "../utils/provider-response";
import { isCopilotTransientModelError } from "../utils/retry";
@@ -257,6 +257,26 @@ export function buildAnthropicHeaders(options: AnthropicHeaderOptions): Record<s
}
type AnthropicCacheControl = NonNullable<TextBlockParam["cache_control"]>;
type AnthropicImageMediaType = "image/jpeg" | "image/png" | "image/gif" | "image/webp";
function normalizeAnthropicImageMediaType(mimeType: string): AnthropicImageMediaType | undefined {
const normalized = mimeType.trim().toLowerCase();
if (normalized === "image/jpg") return "image/jpeg";
if (
normalized === "image/jpeg" ||
normalized === "image/png" ||
normalized === "image/gif" ||
normalized === "image/webp"
) {
return normalized;
}
return undefined;
}
function cloneAnthropicCacheControl(cacheControl: AnthropicCacheControl): AnthropicCacheControl {
return { ...cacheControl };
}
type AnthropicOutputConfig = NonNullable<MessageCreateParamsStreaming["output_config"]>;
@@ -750,42 +770,67 @@ function convertContentBlocks(
type: "image";
source: {
type: "base64";
media_type: "image/jpeg" | "image/png" | "image/gif" | "image/webp";
media_type: AnthropicImageMediaType;
data: string;
};
}
> {
const textBlocks = content
.filter((block): block is TextContent => block.type === "text")
.map(block => block.text.toWellFormed())
.filter(text => text.trim().length > 0);
const imageBlocks = content.filter((block): block is ImageContent => block.type === "image");
const omittedImages = !supportsImages && imageBlocks.length > 0;
if (imageBlocks.length === 0 || !supportsImages) {
if (omittedImages) {
textBlocks.push(NON_VISION_IMAGE_PLACEHOLDER);
}
return textBlocks.join("\n").toWellFormed();
}
const blocks: Array<
| { type: "text"; text: string }
| {
type: "image";
source: {
type: "base64";
media_type: AnthropicImageMediaType;
data: string;
};
}
> = [];
let sawText = false;
let sawImage = false;
const blocks = [
...textBlocks.map(text => ({
type: "text" as const,
text,
})),
...imageBlocks.map(block => ({
type: "image" as const,
for (const block of content) {
if (block.type === "text") {
const text = block.text.toWellFormed();
if (text.trim().length === 0) continue;
sawText = true;
blocks.push({ type: "text", text });
continue;
}
if (!supportsImages) {
blocks.push({ type: "text", text: NON_VISION_IMAGE_PLACEHOLDER });
continue;
}
const mediaType = normalizeAnthropicImageMediaType(block.mimeType);
if (!mediaType) {
blocks.push({ type: "text", text: `[unsupported image: ${block.mimeType}]` });
continue;
}
sawImage = true;
blocks.push({
type: "image",
source: {
type: "base64" as const,
media_type: block.mimeType as "image/jpeg" | "image/png" | "image/gif" | "image/webp",
type: "base64",
media_type: mediaType,
data: block.data,
},
})),
];
});
}
if (!textBlocks.length) {
if (!supportsImages) {
return blocks
.filter((block): block is { type: "text"; text: string } => block.type === "text")
.map(block => block.text)
.join("\n")
.toWellFormed();
}
if (sawImage && !sawText) {
blocks.unshift({
type: "text" as const,
type: "text",
text: "(see attached image)",
});
}
@@ -887,6 +932,16 @@ type FoundryTlsOptions = {
key?: string;
};
const foundryTlsOptionsCache = new Map<string, FoundryTlsOptions | undefined>();
function foundryTlsOptionsCacheKey(): string {
return JSON.stringify([
$env.NODE_EXTRA_CA_CERTS ?? null,
$env.CLAUDE_CODE_CLIENT_CERT ?? null,
$env.CLAUDE_CODE_CLIENT_KEY ?? null,
]);
}
function resolveAnthropicBaseUrl(model: Model<"anthropic-messages">, apiKey?: string): string | undefined {
if (model.provider === "github-copilot") {
return normalizeAnthropicBaseUrl(resolveGitHubCopilotBaseUrl(model.baseUrl, apiKey) ?? model.baseUrl);
@@ -975,6 +1030,9 @@ function resolveFoundryTlsOptions(model: Model<"anthropic-messages">): FoundryTl
if (model.provider !== "anthropic") return undefined;
if (!isFoundryEnabled()) return undefined;
const cacheKey = foundryTlsOptionsCacheKey();
if (foundryTlsOptionsCache.has(cacheKey)) return foundryTlsOptionsCache.get(cacheKey);
const ca = resolvePemValue($env.NODE_EXTRA_CA_CERTS, "NODE_EXTRA_CA_CERTS");
const cert = resolvePemValue($env.CLAUDE_CODE_CLIENT_CERT, "CLAUDE_CODE_CLIENT_CERT");
const key = resolvePemValue($env.CLAUDE_CODE_CLIENT_KEY, "CLAUDE_CODE_CLIENT_KEY");
@@ -987,7 +1045,9 @@ function resolveFoundryTlsOptions(model: Model<"anthropic-messages">): FoundryTl
if (ca) options.ca = [...tls.rootCertificates, ca];
if (cert) options.cert = cert;
if (key) options.key = key;
return Object.keys(options).length > 0 ? options : undefined;
const resolved = Object.keys(options).length > 0 ? options : undefined;
foundryTlsOptionsCache.set(cacheKey, resolved);
return resolved;
}
function buildClaudeCodeTlsFetchOptions(
@@ -1037,14 +1097,8 @@ const ANTHROPIC_MESSAGE_EVENTS: ReadonlySet<string> = new Set([
]);
/**
* Anthropic keepalive `ping` events carry no message content, but they prove the
* upstream connection is alive during long server-side gaps (extended thinking,
* slow tool execution). They are normally dropped before reaching the consumer;
* we instead surface them as lightweight markers so the idle watchdog
* (`iterateWithIdleTimeout`) resets its deadline on every ping. Without this, a
* connection that is demonstrably still streaming pings still trips
* "Anthropic stream stalled while waiting for the next event". The message-event
* branches in `streamAnthropic` match none of these markers, so they are ignored.
* Iterate over Anthropic SSE events from a raw Response, preserving ping events
* for liveness and rejecting malformed complete event envelopes.
*/
type RawMessagePingEvent = { type: "ping" };
type AnthropicStreamEvent = RawMessageStreamEvent | RawMessagePingEvent;
@@ -1079,7 +1133,10 @@ async function* iterateAnthropicEvents(
}
try {
const event = parseJsonWithRepair<RawMessageStreamEvent>(sse.data);
const event = JSON.parse(sse.data) as RawMessageStreamEvent;
if (event.type !== sse.event) {
throw new Error(`event type ${event.type} does not match SSE event ${sse.event}`);
}
if (event.type === "message_start") {
sawMessageStart = true;
} else if (event.type === "message_stop") {
@@ -1094,7 +1151,7 @@ async function* iterateAnthropicEvents(
}
}
if (sawMessageStart && !sawMessageEnd) {
if (sawMessageStart && !sawMessageEnd && !signal?.aborted) {
throw createAnthropicStreamEnvelopeError("stream ended before message_stop");
}
}
@@ -1177,17 +1234,12 @@ function getAnthropicCompat(
const PROVIDER_MAX_RETRIES = 3;
const PROVIDER_BASE_DELAY_MS = 2000;
/**
* Check if an error from the Anthropic SDK is a rate-limit/transient error that
* should be retried before any content has been emitted.
*
* Includes malformed JSON stream-envelope parse errors seen from some
* Anthropic-compatible proxy endpoints.
*/
/** Transient stream corruption errors where the response was truncated mid-JSON. */
function isTransientStreamParseError(error: unknown): boolean {
if (!(error instanceof Error)) return false;
return /json parse error|unterminated string|unexpected end of json input/i.test(error.message);
return /unterminated string|unexpected end of json input|unexpected end of data|unexpected eof|end of file|eof while parsing|truncated/i.test(
error.message,
);
}
const ANTHROPIC_STREAM_ENVELOPE_ERROR_PREFIX = "Anthropic stream envelope error:";
@@ -1235,6 +1287,8 @@ export function isProviderRetryableError(error: unknown, provider?: string): boo
// `streamSimple` a/b/c policy), so surface them immediately instead of
// burning the retry budget here.
if (isUsageLimitError(error.message)) return false;
const status = extractHttpStatusFromError(error);
if (status !== undefined && status >= 400 && status < 500 && status !== 408 && status !== 429) return false;
const msg = error.message.toLowerCase();
if (
isUnexpectedSocketCloseMessage(msg) ||
@@ -1268,13 +1322,12 @@ export type AnthropicUsageLike = {
/**
* Capture Anthropic's optional cache-creation TTL breakdown and server-tool-use
* counters into the harness Usage shape. Only sets fields that were reported, so
* a `message_delta` that omits `cache_creation` does not clobber the breakdown
* established at `message_start`.
* counters into the harness Usage shape. Omitted/null fields are no-ops; explicit
* zero-valued objects clear prior extras from earlier stream usage snapshots.
*/
export function applyAnthropicUsageExtras(usage: Usage, source: AnthropicUsageLike): void {
const cacheCreation = source.cache_creation;
if (cacheCreation) {
if (cacheCreation != null) {
const fiveMinute = cacheCreation.ephemeral_5m_input_tokens ?? 0;
const oneHour = cacheCreation.ephemeral_1h_input_tokens ?? 0;
if (fiveMinute > 0 || oneHour > 0) {
@@ -1282,10 +1335,12 @@ export function applyAnthropicUsageExtras(usage: Usage, source: AnthropicUsageLi
...(fiveMinute > 0 ? { ephemeral5m: fiveMinute } : {}),
...(oneHour > 0 ? { ephemeral1h: oneHour } : {}),
};
} else {
delete usage.cttl;
}
}
const serverToolUse = source.server_tool_use;
if (serverToolUse) {
if (serverToolUse != null) {
const webSearch = serverToolUse.web_search_requests ?? 0;
const webFetch = serverToolUse.web_fetch_requests ?? 0;
if (webSearch > 0 || webFetch > 0) {
@@ -1293,6 +1348,8 @@ export function applyAnthropicUsageExtras(usage: Usage, source: AnthropicUsageLi
...(webSearch > 0 ? { webSearch } : {}),
...(webFetch > 0 ? { webFetch } : {}),
};
} else {
delete usage.server;
}
}
}
@@ -1466,6 +1523,11 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = (
let sawEvent = false;
let sawMessageStart = false;
let sawTerminalEnvelope = false;
let sawMessageStop = false;
const openBlocks = new Map<
number,
{ contentIndex: number; kind: "text" | "thinking" | "redactedThinking" | "toolCall" }
>();
const timedAnthropicStream = iterateWithIdleTimeout(anthropicStream, {
idleTimeoutMs,
@@ -1508,6 +1570,12 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = (
}
if (event.type === "content_block_start") {
if (sawTerminalEnvelope) {
throw createAnthropicStreamEnvelopeError(`received ${event.type} after terminal stop signal`);
}
if (openBlocks.has(event.index)) {
throw createAnthropicStreamEnvelopeError(`duplicate content_block_start index ${event.index}`);
}
if (!firstTokenTime) firstTokenTime = Date.now();
if (event.content_block.type === "text") {
streamedReplayUnsafeContent = true;
@@ -1517,12 +1585,15 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = (
index: event.index,
};
output.content.push(block);
const contentIndex = output.content.length - 1;
openBlocks.set(event.index, { contentIndex, kind: "text" });
stream.push({
type: "text_start",
contentIndex: output.content.length - 1,
contentIndex,
partial: output,
});
} else if (event.content_block.type === "thinking") {
streamedReplayUnsafeContent = true;
const block: Block = {
type: "thinking",
thinking: "",
@@ -1530,18 +1601,25 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = (
index: event.index,
};
output.content.push(block);
const contentIndex = output.content.length - 1;
openBlocks.set(event.index, { contentIndex, kind: "thinking" });
stream.push({
type: "thinking_start",
contentIndex: output.content.length - 1,
contentIndex,
partial: output,
});
} else if (event.content_block.type === "redacted_thinking") {
streamedReplayUnsafeContent = true;
const block: Block = {
type: "redactedThinking",
data: event.content_block.data,
index: event.index,
};
output.content.push(block);
openBlocks.set(event.index, {
contentIndex: output.content.length - 1,
kind: "redactedThinking",
});
} else if (event.content_block.type === "tool_use") {
streamedReplayUnsafeContent = true;
const block: Block = {
@@ -1555,92 +1633,115 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = (
index: event.index,
};
output.content.push(block);
const contentIndex = output.content.length - 1;
openBlocks.set(event.index, { contentIndex, kind: "toolCall" });
stream.push({
type: "toolcall_start",
contentIndex: output.content.length - 1,
contentIndex,
partial: output,
});
}
} else if (event.type === "content_block_delta") {
if (sawTerminalEnvelope) {
throw createAnthropicStreamEnvelopeError(`received ${event.type} after terminal stop signal`);
}
const openBlock = openBlocks.get(event.index);
if (!openBlock) {
throw createAnthropicStreamEnvelopeError(`received content_block_delta for unopened index ${event.index}`);
}
const block = blocks[openBlock.contentIndex];
if (event.delta.type === "text_delta") {
const index = blocks.findIndex(b => b.index === event.index);
const block = blocks[index];
if (block && block.type === "text") {
block.text += event.delta.text;
stream.push({
type: "text_delta",
contentIndex: index,
delta: event.delta.text,
partial: output,
});
if (openBlock.kind !== "text" || block?.type !== "text") {
throw createAnthropicStreamEnvelopeError(`received text_delta for ${openBlock.kind} block`);
}
streamedReplayUnsafeContent = true;
block.text += event.delta.text;
stream.push({
type: "text_delta",
contentIndex: openBlock.contentIndex,
delta: event.delta.text,
partial: output,
});
} else if (event.delta.type === "thinking_delta") {
const index = blocks.findIndex(b => b.index === event.index);
const block = blocks[index];
if (block && block.type === "thinking") {
block.thinking += event.delta.thinking;
stream.push({
type: "thinking_delta",
contentIndex: index,
delta: event.delta.thinking,
partial: output,
});
if (openBlock.kind !== "thinking" || block?.type !== "thinking") {
throw createAnthropicStreamEnvelopeError(`received thinking_delta for ${openBlock.kind} block`);
}
streamedReplayUnsafeContent = true;
block.thinking += event.delta.thinking;
stream.push({
type: "thinking_delta",
contentIndex: openBlock.contentIndex,
delta: event.delta.thinking,
partial: output,
});
} else if (event.delta.type === "input_json_delta") {
const index = blocks.findIndex(b => b.index === event.index);
const block = blocks[index];
if (block && block.type === "toolCall") {
block.partialJson += event.delta.partial_json;
const throttled = parseStreamingJsonThrottled(block.partialJson, block.lastParseLen ?? 0);
if (throttled) {
block.arguments = throttled.value;
block.lastParseLen = throttled.parsedLen;
}
stream.push({
type: "toolcall_delta",
contentIndex: index,
delta: event.delta.partial_json,
partial: output,
});
if (openBlock.kind !== "toolCall" || block?.type !== "toolCall") {
throw createAnthropicStreamEnvelopeError(`received input_json_delta for ${openBlock.kind} block`);
}
streamedReplayUnsafeContent = true;
block.partialJson += event.delta.partial_json;
const throttled = parseStreamingJsonThrottled(block.partialJson, block.lastParseLen ?? 0);
if (throttled) {
block.arguments = throttled.value;
block.lastParseLen = throttled.parsedLen;
}
stream.push({
type: "toolcall_delta",
contentIndex: openBlock.contentIndex,
delta: event.delta.partial_json,
partial: output,
});
} else if (event.delta.type === "signature_delta") {
const index = blocks.findIndex(b => b.index === event.index);
const block = blocks[index];
if (block && block.type === "thinking") {
block.thinkingSignature = block.thinkingSignature || "";
block.thinkingSignature += event.delta.signature;
if (openBlock.kind !== "thinking" || block?.type !== "thinking") {
throw createAnthropicStreamEnvelopeError(`received signature_delta for ${openBlock.kind} block`);
}
streamedReplayUnsafeContent = true;
block.thinkingSignature = block.thinkingSignature || "";
block.thinkingSignature += event.delta.signature;
}
} else if (event.type === "content_block_stop") {
const index = blocks.findIndex(b => b.index === event.index);
const block = blocks[index];
if (block) {
delete (block as { index?: number }).index;
if (block.type === "text") {
stream.push({
type: "text_end",
contentIndex: index,
content: block.text,
partial: output,
});
} else if (block.type === "thinking") {
stream.push({
type: "thinking_end",
contentIndex: index,
content: block.thinking,
partial: output,
});
} else if (block.type === "toolCall") {
block.arguments = parseStreamingJson(block.partialJson);
delete (block as { partialJson?: string }).partialJson;
delete (block as { lastParseLen?: number }).lastParseLen;
stream.push({
type: "toolcall_end",
contentIndex: index,
toolCall: block,
partial: output,
});
}
if (sawTerminalEnvelope) {
throw createAnthropicStreamEnvelopeError(`received ${event.type} after terminal stop signal`);
}
const openBlock = openBlocks.get(event.index);
if (!openBlock) {
throw createAnthropicStreamEnvelopeError(`received content_block_stop for unopened index ${event.index}`);
}
const block = blocks[openBlock.contentIndex];
if (!block || block.type !== openBlock.kind) {
throw createAnthropicStreamEnvelopeError(`content_block_stop kind mismatch for index ${event.index}`);
}
openBlocks.delete(event.index);
delete (block as { index?: number }).index;
if (block.type === "text") {
streamedReplayUnsafeContent = true;
stream.push({
type: "text_end",
contentIndex: openBlock.contentIndex,
content: block.text,
partial: output,
});
} else if (block.type === "thinking") {
streamedReplayUnsafeContent = true;
stream.push({
type: "thinking_end",
contentIndex: openBlock.contentIndex,
content: block.thinking,
partial: output,
});
} else if (block.type === "toolCall") {
streamedReplayUnsafeContent = true;
const finalJson =
block.partialJson.length > 0 ? block.partialJson : JSON.stringify(block.arguments ?? {});
block.arguments = JSON.parse(finalJson) as ToolCall["arguments"];
delete (block as { partialJson?: string }).partialJson;
delete (block as { lastParseLen?: number }).lastParseLen;
stream.push({
type: "toolcall_end",
contentIndex: openBlock.contentIndex,
toolCall: block,
partial: output,
});
}
} else if (event.type === "message_delta") {
const rawStopReason = event.delta.stop_reason;
@@ -1683,6 +1784,7 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = (
calculateCost(model, output.usage);
} else if (event.type === "message_stop") {
sawTerminalEnvelope = true;
sawMessageStop = true;
}
}
@@ -1696,25 +1798,18 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = (
if (!sawEvent || !sawMessageStart) {
throw createAnthropicStreamEnvelopeError("stream ended before message_start");
}
if (!sawTerminalEnvelope) {
throw createAnthropicStreamEnvelopeError("stream ended before terminal stop signal");
if (!sawMessageStop) {
throw createAnthropicStreamEnvelopeError("stream ended before message_stop");
}
// An open tool_use block — one that never received its
// `content_block_stop` — means the stream was truncated mid-tool-call.
// In practice this is a transport drop that a transparent reconnect
// splices back together: the reconnect's `message_start` is deduped
// above, yet the orphaned block survives with its seed `{}` (or a
// partial-parse) arguments. Emitting it would dispatch a tool call the
// model never finished generating. Surface it as a truncated envelope so
// the existing retry/error path engages instead of shipping bogus args.
if (
blocks.some(
block =>
block.type === "toolCall" && (block as { partialJson?: string }).partialJson !== undefined,
)
) {
throw createAnthropicStreamEnvelopeError("stream ended with an unterminated tool_use block");
if (openBlocks.size > 0) {
const firstOpenBlock = openBlocks.entries().next().value;
if (firstOpenBlock) {
const [openIndex, openBlock] = firstOpenBlock;
throw createAnthropicStreamEnvelopeError(
`stream ended with an unterminated ${openBlock.kind} block at index ${openIndex}`,
);
}
throw createAnthropicStreamEnvelopeError("stream ended with an unterminated content block");
}
if (output.stopReason === "aborted" || output.stopReason === "error") {
@@ -1849,12 +1944,11 @@ function applyClaudeCodeSystemCache(
blocks: AnthropicSystemBlock[],
cacheControl: AnthropicCacheControl | undefined,
): number {
if (!cacheControl || blocks.length <= 2) return 0;
blocks[2] = { ...blocks[2], cache_control: cacheControl };
if (blocks.length === 3) return 1;
if (!cacheControl || blocks.length === 0) return 0;
const lastIndex = blocks.length - 1;
blocks[lastIndex] = { ...blocks[lastIndex], cache_control: cacheControl };
return 2;
if (blocks[lastIndex].cache_control != null) return 0;
blocks[lastIndex] = { ...blocks[lastIndex], cache_control: cloneAnthropicCacheControl(cacheControl) };
return 1;
}
export function buildAnthropicSystemBlocks(
@@ -1891,8 +1985,8 @@ export function buildAnthropicSystemBlocks(
blocks.push({ type: "text", text: prompt });
}
const lastIndex = blocks.length - 1;
if (cacheControl && lastIndex >= 0) {
blocks[lastIndex] = { ...blocks[lastIndex], cache_control: cacheControl };
if (cacheControl && lastIndex >= 0 && blocks[lastIndex].cache_control == null) {
blocks[lastIndex] = { ...blocks[lastIndex], cache_control: cloneAnthropicCacheControl(cacheControl) };
}
return blocks.length > 0 ? blocks : undefined;
}
@@ -2009,7 +2103,6 @@ export function buildAnthropicClientOptions(args: AnthropicClientOptionsArgs): A
...(tlsFetchOptions ? { fetchOptions: tlsFetchOptions } : {}),
};
}
// OpenCode Zen's Anthropic-compatible gateway accepts bearer auth only;
// leaving apiKey set lets the client add X-Api-Key, which upstream Alibaba rejects.
if (model.provider === "opencode-zen") {
@@ -2025,9 +2118,16 @@ export function buildAnthropicClientOptions(args: AnthropicClientOptionsArgs): A
};
}
const authorizationHeader = getHeaderCaseInsensitive(defaultHeaders, "Authorization");
const shouldSuppressClientApiKey =
!oauthToken &&
!isAnthropicApiBaseUrl(baseUrl) &&
typeof authorizationHeader === "string" &&
/^Bearer\s+/i.test(authorizationHeader);
return {
isOAuthToken: oauthToken,
apiKey: oauthToken ? null : apiKey,
apiKey: oauthToken || shouldSuppressClientApiKey ? null : apiKey,
authToken: oauthToken ? apiKey : undefined,
baseURL: baseUrl,
maxRetries: 5,
@@ -2052,6 +2152,7 @@ function disableThinkingIfToolChoiceForced(params: MessageCreateParamsStreaming)
if (toolChoice.type !== "any" && toolChoice.type !== "tool") return;
delete params.thinking;
delete params.context_management;
const outputConfig = params.output_config as AnthropicOutputConfig | undefined;
if (!outputConfig) return;
@@ -2068,11 +2169,20 @@ function ensureMaxTokensForThinking(params: MessageCreateParamsStreaming, model:
const budgetTokens = thinking.budget_tokens ?? 0;
if (budgetTokens <= 0) return;
const maxTokens = params.max_tokens ?? 0;
const requiredMaxTokens = budgetTokens + OUTPUT_FALLBACK_BUFFER;
if (maxTokens < requiredMaxTokens) {
params.max_tokens = Math.min(requiredMaxTokens, model.maxTokens);
const maxAllowedTokens = Math.min(CLAUDE_CODE_MAX_OUTPUT_TOKENS, model.maxTokens);
const currentMaxTokens = Math.min(params.max_tokens ?? maxAllowedTokens, maxAllowedTokens);
const raisedMaxTokens = Math.min(Math.max(currentMaxTokens, budgetTokens + OUTPUT_FALLBACK_BUFFER), maxAllowedTokens);
params.max_tokens = raisedMaxTokens;
if (budgetTokens + OUTPUT_FALLBACK_BUFFER <= raisedMaxTokens) return;
const clampedBudget = raisedMaxTokens - OUTPUT_FALLBACK_BUFFER;
if (clampedBudget <= 0) {
throw new Error(
`Anthropic thinking budget requires max_tokens greater than ${OUTPUT_FALLBACK_BUFFER}; got ${raisedMaxTokens}`,
);
}
thinking.budget_tokens = clampedBudget;
}
type CacheControlBlock = {
@@ -2082,39 +2192,35 @@ type CacheControlBlock = {
function applyCacheControlToLastBlock<T extends CacheControlBlock>(
blocks: T[],
cacheControl: AnthropicCacheControl,
): void {
if (blocks.length === 0) return;
): boolean {
if (blocks.length === 0) return false;
const lastIndex = blocks.length - 1;
blocks[lastIndex] = { ...blocks[lastIndex], cache_control: cacheControl };
if (blocks[lastIndex].cache_control != null) return false;
blocks[lastIndex] = { ...blocks[lastIndex], cache_control: cloneAnthropicCacheControl(cacheControl) };
return true;
}
function applyCacheControlToLastTextBlock(
blocks: Array<ContentBlockParam & CacheControlBlock>,
cacheControl: AnthropicCacheControl,
): void {
if (blocks.length === 0) return;
): boolean {
if (blocks.length === 0) return false;
for (let i = blocks.length - 1; i >= 0; i--) {
if (blocks[i].type === "text") {
blocks[i] = { ...blocks[i], cache_control: cacheControl };
return;
if (blocks[i].cache_control != null) return false;
blocks[i] = { ...blocks[i], cache_control: cloneAnthropicCacheControl(cacheControl) };
return true;
}
}
applyCacheControlToLastBlock(blocks, cacheControl);
return false;
}
function applyPromptCaching(params: MessageCreateParamsStreaming, cacheControl?: AnthropicCacheControl): void {
if (!cacheControl) return;
// Skip if cache_control breakpoints were already placed externally on messages.
for (const message of params.messages) {
if (Array.isArray(message.content)) {
if ((message.content as Array<ContentBlockParam & CacheControlBlock>).some(b => b.cache_control != null))
return;
}
}
const MAX_CACHE_BREAKPOINTS = 4;
let cacheBreakpointsUsed = 0;
let cacheBreakpointsUsed = countCacheControlBreakpoints(params);
if (cacheBreakpointsUsed >= MAX_CACHE_BREAKPOINTS) return;
let isCCLayout = false;
if (params.system && Array.isArray(params.system) && params.system.length > 0) {
@@ -2122,9 +2228,12 @@ function applyPromptCaching(params: MessageCreateParamsStreaming, cacheControl?:
params.system.length >= 3 &&
(params.system[0] as { text?: string }).text?.startsWith(CLAUDE_BILLING_HEADER_PREFIX) === true;
if (isCCLayout) {
cacheBreakpointsUsed += applyClaudeCodeSystemCache(params.system as AnthropicSystemBlock[], cacheControl);
} else {
applyCacheControlToLastBlock(params.system, cacheControl);
const placed = Math.min(
MAX_CACHE_BREAKPOINTS - cacheBreakpointsUsed,
applyClaudeCodeSystemCache(params.system as AnthropicSystemBlock[], cacheControl),
);
cacheBreakpointsUsed += placed;
} else if (applyCacheControlToLastBlock(params.system, cacheControl)) {
cacheBreakpointsUsed++;
}
}
@@ -2137,14 +2246,19 @@ function applyPromptCaching(params: MessageCreateParamsStreaming, cacheControl?:
const message = params.messages[i];
if (!message) continue;
if (typeof message.content === "string") {
message.content = [{ type: "text", text: message.content, cache_control: cacheControl }];
message.content = [
{ type: "text", text: message.content, cache_control: cloneAnthropicCacheControl(cacheControl) },
];
cacheBreakpointsUsed++;
} else if (Array.isArray(message.content) && message.content.length > 0) {
applyCacheControlToLastTextBlock(
message.content as Array<ContentBlockParam & CacheControlBlock>,
cacheControl,
);
cacheBreakpointsUsed++;
if (
applyCacheControlToLastTextBlock(
message.content as Array<ContentBlockParam & CacheControlBlock>,
cacheControl,
)
) {
cacheBreakpointsUsed++;
}
}
}
}
@@ -2157,7 +2271,9 @@ function normalizeCacheControlBlockTtl(block: CacheControlBlock, seenFiveMinute:
return;
}
if (seenFiveMinute.value) {
delete cacheControl.ttl;
const normalized = cloneAnthropicCacheControl(cacheControl);
delete normalized.ttl;
block.cache_control = normalized;
}
}
@@ -2322,7 +2438,7 @@ function buildParams(
});
// Pre-compute tools.
let tools: ReturnType<typeof convertTools> | undefined;
let tools: AnthropicWireTool[] | undefined;
if (context.tools) {
tools = convertTools(
context.tools,
@@ -2385,11 +2501,11 @@ function buildParams(
// metadata → max_tokens → thinking → context_management → output_config → stream.
const params: MessageCreateParamsStreaming = {
model: model.id,
messages: convertAnthropicMessages(context.messages, model, isOAuthToken),
messages: convertAnthropicMessages(context.messages, model, isOAuthToken, baseUrl),
...(systemBlocks && { system: systemBlocks }),
...(tools !== undefined && { tools }),
...(metadata && { metadata }),
max_tokens: Math.min(CLAUDE_CODE_MAX_OUTPUT_TOKENS, options?.maxTokens || model.maxTokens),
max_tokens: Math.min(CLAUDE_CODE_MAX_OUTPUT_TOKENS, model.maxTokens, options?.maxTokens || model.maxTokens),
...(thinking && { thinking }),
...(contextManagement && { context_management: contextManagement }),
...(outputConfig && { output_config: outputConfig }),
@@ -2397,8 +2513,8 @@ function buildParams(
};
// Opus 4.7+ rejects non-default sampling parameters with 400 error.
const allowSamplingParams = !hasOpus47ApiRestrictions(model.id);
if (allowSamplingParams && options?.temperature !== undefined && !options?.thinkingEnabled) {
const allowSamplingParams = !hasOpus47ApiRestrictions(model.id) && !params.thinking;
if (allowSamplingParams && options?.temperature !== undefined) {
params.temperature = options.temperature;
}
if (allowSamplingParams && options?.topP !== undefined) {
@@ -2476,9 +2592,8 @@ function isZaiAnthropicEndpoint(model: Model<"anthropic-messages">): boolean {
* arguments (#2005). Known non-signing hosts are also preserved for
* compatibility.
*/
function shouldReplayUnsignedThinking(model: Model<"anthropic-messages">): boolean {
function shouldReplayUnsignedThinking(model: Model<"anthropic-messages">, baseUrl: string | undefined): boolean {
if (model.provider === "zai" || model.provider === "deepseek") return true;
const baseUrl = model.baseUrl;
if (baseUrl) {
try {
const hostname = new URL(baseUrl).hostname.toLowerCase();
@@ -2514,12 +2629,13 @@ export function convertAnthropicMessages(
messages: Message[],
model: Model<"anthropic-messages">,
isOAuthToken: boolean,
baseUrl = resolveAnthropicBaseUrl(model),
): AnthropicMessageParam[] {
const params: AnthropicMessageParam[] = [];
// Indices of params emitted from `developer` messages. After the main pass,
// the ones whose placement satisfies Anthropic's mid-conversation rules are
// upgraded from the `user` role to the authoritative `system` role.
const developerParamIndices: number[] = [];
const params: AnthropicMessageParam[] = [];
const transformedMessages = transformMessages(messages, model, normalizeToolCallId);
@@ -2578,7 +2694,7 @@ export function convertAnthropicMessages(
}
if (block.thinking.trim().length === 0) continue;
if (!block.thinkingSignature || block.thinkingSignature.trim().length === 0) {
if (shouldReplayUnsignedThinking(model)) {
if (shouldReplayUnsignedThinking(model, baseUrl)) {
blocks.push({
type: "thinking",
thinking: block.thinking.toWellFormed(),
@@ -2728,6 +2844,7 @@ function isJsonSchemaArrayNode(schema: Record<string, unknown>): boolean {
const t = schema.type;
if (t === "array") return true;
if (Array.isArray(t) && t.includes("array") && !t.includes("object")) return true;
if (schema.items !== undefined || Array.isArray(schema.prefixItems)) return true;
return false;
}
@@ -2754,6 +2871,14 @@ function pickAnthropicScalarType(type: unknown): string | undefined {
}
return undefined;
}
function pickAnthropicEffectiveScalarType(schema: Record<string, unknown>): string | undefined {
const explicit = pickAnthropicScalarType(schema.type);
if (explicit) return explicit;
if (isRecord(schema.properties)) return "object";
if (schema.items !== undefined || Array.isArray(schema.prefixItems)) return "array";
return undefined;
}
function anthropicPerTypeKeep(scalarType: string | undefined): Set<string> | undefined {
switch (scalarType) {
@@ -2768,14 +2893,6 @@ function anthropicPerTypeKeep(scalarType: string | undefined): Set<string> | und
}
}
/**
* Per-schema-object memoization slot for the normalized Anthropic tool form. We stamp
* the result onto the host via a `Symbol` property (mirroring `utils/schema/stamps.ts`)
* instead of using a `WeakMap`: it's a single hidden-class slot, so warm reads are
* direct property access and write-once cycles resolve to the in-progress result.
*/
const kAnthropicToolNormal = Symbol("pi.schema.anthropic.toolNormal");
/**
* Normalize a JSON Schema node for Anthropic tool `input_schema`.
*
@@ -2796,20 +2913,20 @@ const kAnthropicToolNormal = Symbol("pi.schema.anthropic.toolNormal");
* pass downstream demotes those shapes to non-strict instead of fabricating a closed
* object, so callers like the resolve tool keep working open-map semantics.
*/
export function normalizeAnthropicToolSchema(schema: unknown): unknown {
if (Array.isArray(schema)) return schema.map(entry => normalizeAnthropicToolSchema(entry));
function normalizeAnthropicToolSchemaNode(
schema: unknown,
cache: WeakMap<Record<string, unknown>, Record<string, unknown>>,
): unknown {
if (Array.isArray(schema)) return schema.map(entry => normalizeAnthropicToolSchemaNode(entry, cache));
if (!isRecord(schema)) return schema;
const slot = schema as Record<symbol, Record<string, unknown> | undefined>;
const existing = slot[kAnthropicToolNormal];
const existing = cache.get(schema);
if (existing !== undefined) return existing;
const result: Record<string, unknown> = {};
// Pre-stamp before recursion so cyclic schemas resolve to the in-progress object
// (mirrors the WeakMap-set-before-recurse pattern the original implementation used).
Object.defineProperty(schema, kAnthropicToolNormal, { value: result, writable: true, configurable: true });
cache.set(schema, result);
const scalarType = pickAnthropicScalarType(schema.type);
const scalarType = pickAnthropicEffectiveScalarType(schema);
const perTypeKeep = anthropicPerTypeKeep(scalarType);
const spill: Array<[string, unknown]> = [];
@@ -2848,12 +2965,12 @@ export function normalizeAnthropicToolSchema(schema: unknown): unknown {
const sourceProperties = result.properties as Record<string, unknown>;
for (const propName in sourceProperties) {
if (!Object.hasOwn(sourceProperties, propName)) continue;
normalizedProperties[propName] = normalizeAnthropicToolSchema(sourceProperties[propName]);
normalizedProperties[propName] = normalizeAnthropicToolSchemaNode(sourceProperties[propName], cache);
}
result.properties = normalizedProperties;
}
if (isRecord(result.additionalProperties)) {
const normalized = normalizeAnthropicToolSchema(result.additionalProperties);
const normalized = normalizeAnthropicToolSchemaNode(result.additionalProperties, cache);
if (isRecord(normalized) && Object.keys(normalized).length === 0) {
result.additionalProperties = true;
} else {
@@ -2861,17 +2978,17 @@ export function normalizeAnthropicToolSchema(schema: unknown): unknown {
}
}
if (Array.isArray(result.items)) {
result.items = result.items.map(item => normalizeAnthropicToolSchema(item));
result.items = result.items.map(item => normalizeAnthropicToolSchemaNode(item, cache));
} else if (isRecord(result.items)) {
result.items = normalizeAnthropicToolSchema(result.items);
result.items = normalizeAnthropicToolSchemaNode(result.items, cache);
}
if (Array.isArray(result.prefixItems)) {
result.prefixItems = result.prefixItems.map(item => normalizeAnthropicToolSchema(item));
result.prefixItems = result.prefixItems.map(item => normalizeAnthropicToolSchemaNode(item, cache));
}
for (const key of COMBINATOR_KEYS) {
const variants = result[key];
if (Array.isArray(variants)) {
result[key] = variants.map(variant => normalizeAnthropicToolSchema(variant));
result[key] = variants.map(variant => normalizeAnthropicToolSchemaNode(variant, cache));
}
}
for (const defsKey of ["$defs", "definitions"] as const) {
@@ -2881,7 +2998,7 @@ export function normalizeAnthropicToolSchema(schema: unknown): unknown {
const sourceDefs = definitions as Record<string, unknown>;
for (const name in sourceDefs) {
if (!Object.hasOwn(sourceDefs, name)) continue;
normalizedDefs[name] = normalizeAnthropicToolSchema(sourceDefs[name]);
normalizedDefs[name] = normalizeAnthropicToolSchemaNode(sourceDefs[name], cache);
}
result[defsKey] = normalizedDefs;
}
@@ -2890,6 +3007,10 @@ export function normalizeAnthropicToolSchema(schema: unknown): unknown {
return result;
}
export function normalizeAnthropicToolSchema(schema: unknown): unknown {
return normalizeAnthropicToolSchemaNode(schema, new WeakMap());
}
type AnthropicToolSchemaPlan = {
inputSchema: AnthropicToolInputSchema;
strict: boolean;
@@ -2910,6 +3031,25 @@ function hasNullVariant(schema: Record<string, unknown>): boolean {
if (Array.isArray(schema.type) && schema.type.includes("null")) return true;
return Array.isArray(schema.anyOf) && schema.anyOf.some(variant => isRecord(variant) && variant.type === "null");
}
function hasAnthropicSchemaDefiningKeyword(schema: Record<string, unknown>): boolean {
if (
schema.type !== undefined ||
schema.properties !== undefined ||
schema.additionalProperties !== undefined ||
schema.items !== undefined ||
schema.prefixItems !== undefined ||
schema.enum !== undefined ||
schema.const !== undefined ||
schema.$ref !== undefined
) {
return true;
}
for (const key of COMBINATOR_KEYS) {
if (schema[key] !== undefined) return true;
}
return schema.$defs !== undefined || schema.definitions !== undefined;
}
function makeAnthropicNullableSchema(schema: unknown, budget: AnthropicStrictBudget): unknown | undefined {
if (isRecord(schema)) {
@@ -2948,6 +3088,8 @@ function normalizeAnthropicStrictSchemaNode(
const cached = cache.get(schema);
if (cached) return cached;
if (!hasAnthropicSchemaDefiningKeyword(schema)) return undefined;
// Strict tool use only supports closed objects. Open maps stay available on
// the non-strict schema plan instead of producing an Anthropic 400.
if (isJsonSchemaObjectNode(schema) && schema.additionalProperties !== false) {
@@ -457,11 +457,11 @@ describe("anthropic stream envelope handling", () => {
expect(attempt).toBe(1);
expect(countEvents(events, "toolcall_start")).toBe(1);
expect(countEvents(events, "toolcall_delta")).toBe(1);
expect(countEvents(events, "toolcall_end")).toBe(1);
expect(countEvents(events, "toolcall_end")).toBe(0);
expect(countEvents(events, "error")).toBe(1);
expect(countEvents(events, "done")).toBe(0);
expect(result.stopReason).toBe("error");
expect(result.errorMessage).toContain("stream ended before terminal stop signal");
expect(result.errorMessage).toContain("Unterminated string");
const toolCall = result.content[0];
expect(toolCall?.type).toBe("toolCall");
@@ -494,7 +494,7 @@ describe("anthropic stream envelope handling", () => {
expect(countEvents(events, "done")).toBe(0);
expect(countEvents(events, "error")).toBe(1);
expect(result.stopReason).toBe("error");
expect(result.errorMessage).toContain("unterminated tool_use block");
expect(result.errorMessage).toContain("unterminated toolCall block");
});
it("parses raw SSE directly so unknown events do not fail Anthropic streams", async () => {
vi.spyOn(AnthropicMessages.prototype, "create").mockImplementation(
@@ -541,7 +541,7 @@ describe("anthropic stream envelope handling", () => {
expect(result.content).toEqual([{ type: "text", text: "partial" }]);
});
it("repairs malformed JSON in raw SSE event data before parsing", async () => {
it("surfaces malformed raw SSE event JSON instead of repairing protocol frames", async () => {
const malformedTextDelta =
'{"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"line\\qbreak"}}';
const successEvents = createTextSuccessEvents("unused");
@@ -556,13 +556,16 @@ describe("anthropic stream envelope handling", () => {
vi.spyOn(AnthropicMessages.prototype, "create").mockImplementation(() => createRawSseRequest(frames) as never);
const stream = streamAnthropic(model, context, { apiKey: "sk-ant-test" });
for await (const _ of stream) {
// drain stream
const events: AssistantMessageEvent[] = [];
for await (const event of stream) {
events.push(event);
}
const result = await stream.result();
expect(result.stopReason).toBe("stop");
expect(result.content).toEqual([{ type: "text", text: "line\\qbreak" }]);
expect(countEvents(events, "error")).toBe(1);
expect(countEvents(events, "done")).toBe(0);
expect(result.stopReason).toBe("error");
expect(result.errorMessage).toContain("Could not parse Anthropic SSE event content_block_delta");
});
it("surfaces a refusal fallback message when stop_details is null", async () => {
const refusalEvents: MockAnthropicEvent[] = [
+6 -4
View File
@@ -1,6 +1,7 @@
# Changelog
## [Unreleased]
### Added
- Added support for paste marker highlighting with accent styling (`[Paste #N, +X lines]`/`[Paste #N, Y chars]`) in the prompt editor, matching the visual treatment of image references
@@ -8,10 +9,15 @@
### Changed
- Normalized image content before it enters model context so attached images are downscaled and preprocessed for prompts, steering messages, follow-ups, and custom agent messages
- Changed image marker format to include pixel dimensions when available (`[Image #N, WxH]`), falling back to bare `[Image #N]` when header cannot be decoded
- Changed the prompt editor to highlight large-paste placeholders (`[Paste #N, +X lines]`/`[Paste #N, Y chars]`) with the same accent styling as image references (bold, no hyperlink), and to delete image/paste markers atomically: a single backspace or forward-delete removes the whole marker instead of leaving a broken `[Paste #N, +X lines` behind.
- Browser tool helpers (`tab.*`) are now individually tracked and time-bounded: when a `run` cell hits its budget, the timeout error names the still-running helper(s) and how long each has been stalled (e.g. `... (stalled on tab.screenshot({ selector: ".x" }) (29.9s))`) instead of the opaque `Browser code execution timed out after 30000ms`. Page-coupled helpers that should resolve quickly (`observe`, `screenshot`, `extract`) also fail fast with a named per-op error at `min(cellBudget, 20s)`, leaving budget for the rest of the cell, rather than silently consuming the whole budget.
### Removed
- Removed the special Anthropic `claude-opus-4-8` tool-call batch cap; sessions no longer abort an in-flight provider stream after a fixed number of completed tool calls.
### Fixed
- Fixed task-row shimmer timing so every running description starts its highlight on the first character together and reaches the last character together, regardless of text length.
@@ -23,10 +29,6 @@
- Fixed edit tool result previews to show only current-file lines and collapse long inserted blocks instead of echoing removed content.
- Fixed `generateDiffString` to omit the mid-skip `...` placeholder between two nearby edits, conveying the elided gap via the jump in line numbers instead (consistent with how leading/trailing context skips already render). The placeholder row was indistinguishable from a genuine `...` context line and wasted a row in compact previews.
### Removed
- Removed the special Anthropic `claude-opus-4-8` tool-call batch cap; sessions no longer abort an in-flight provider stream after a fixed number of completed tool calls.
## [15.10.4] - 2026-06-08
### Added
@@ -215,6 +215,7 @@ import { parseCommandArgs } from "../utils/command-args";
import { type EditMode, resolveEditMode } from "../utils/edit-mode";
import { resolveFileDisplayMode } from "../utils/file-display-mode";
import { extractFileMentions, generateFileMentionMessages } from "../utils/file-mentions";
import { normalizeModelContextImages } from "../utils/image-loading";
import { buildNamedToolChoice } from "../utils/tool-choice";
import type { AuthStorage } from "./auth-storage";
import type { ClientBridge, ClientBridgePermissionOption, ClientBridgePermissionOutcome } from "./client-bridge";
@@ -4272,6 +4273,24 @@ export class AgentSession {
};
}
async #normalizeMessageContentImages(
content: string | (TextContent | ImageContent)[],
): Promise<string | (TextContent | ImageContent)[]> {
if (typeof content === "string") return content;
const images = content.filter((part): part is ImageContent => part.type === "image");
if (images.length === 0) return content;
const normalizedImages = await normalizeModelContextImages(images);
if (!normalizedImages) return content;
let imageIndex = 0;
return content.map(part => (part.type === "image" ? normalizedImages[imageIndex++]! : part));
}
async #normalizeAgentMessageImages<T extends AgentMessage>(message: T): Promise<T> {
const content = await this.#normalizeMessageContentImages(message.content);
if (content === message.content) return message;
return { ...message, content } as T;
}
/**
* Send a prompt to the agent.
* - Handles extension commands (registered via pi.registerCommand) immediately, even during streaming
@@ -4371,10 +4390,11 @@ export class AgentSession {
const hasPendingUserDirective = this.#toolChoiceQueue.inspect().includes("user-force");
const eagerTodoPrelude =
!options?.synthetic && !hasPendingUserDirective ? this.#createEagerTodoPrelude(expandedText) : undefined;
const normalizedImages = await normalizeModelContextImages(options?.images);
const userContent: (TextContent | ImageContent)[] = [{ type: "text", text: expandedText }];
if (options?.images) {
userContent.push(...options.images);
if (normalizedImages) {
userContent.push(...normalizedImages);
}
const promptAttribution = options?.attribution ?? (options?.synthetic ? "agent" : "user");
@@ -4391,6 +4411,7 @@ export class AgentSession {
try {
await this.#promptWithMessage(message, expandedText, {
...options,
images: normalizedImages,
prependMessages: eagerTodoPrelude ? [eagerTodoPrelude.message] : undefined,
appendMessages: keywordNotices.length > 0 ? keywordNotices : undefined,
});
@@ -4533,7 +4554,9 @@ export class AgentSession {
useHashLines: resolveFileDisplayMode(this).hashLines,
snapshotStore: getFileSnapshotStore(this),
});
messages.push(...fileMentionMessages);
for (const fileMentionMessage of fileMentionMessages) {
messages.push(await this.#normalizeAgentMessageImages(fileMentionMessage));
}
}
const beforeAgentStartSystemPrompt = await this.#buildSystemPromptForAgentStart(expandedText);
@@ -4549,15 +4572,17 @@ export class AgentSession {
const promptAttribution: "user" | "agent" | undefined =
"attribution" in message ? message.attribution : undefined;
for (const msg of result.messages) {
messages.push({
role: "custom",
customType: msg.customType,
content: msg.content,
display: msg.display,
details: msg.details,
attribution: msg.attribution ?? promptAttribution ?? (message.role === "user" ? "user" : "agent"),
timestamp: Date.now(),
});
messages.push(
await this.#normalizeAgentMessageImages({
role: "custom",
customType: msg.customType,
content: msg.content,
display: msg.display,
details: msg.details,
attribution: msg.attribution ?? promptAttribution ?? (message.role === "user" ? "user" : "agent"),
timestamp: Date.now(),
}),
);
}
}
@@ -4765,11 +4790,12 @@ export class AgentSession {
* Internal: Queue a steering message (already expanded, no extension command check).
*/
async #queueSteer(text: string, images?: ImageContent[]): Promise<void> {
const normalizedImages = await normalizeModelContextImages(images);
const displayText = text || (images && images.length > 0 ? "[Image]" : "");
this.#steeringMessages.push({ text: displayText });
const content: (TextContent | ImageContent)[] = [{ type: "text", text }];
if (images && images.length > 0) {
content.push(...images);
if (normalizedImages && normalizedImages.length > 0) {
content.push(...normalizedImages);
}
this.agent.steer({
role: "user",
@@ -4784,11 +4810,12 @@ export class AgentSession {
* Internal: Queue a follow-up message (already expanded, no extension command check).
*/
async #queueFollowUp(text: string, images?: ImageContent[]): Promise<void> {
const normalizedImages = await normalizeModelContextImages(images);
const displayText = text || (images && images.length > 0 ? "[Image]" : "");
this.#followUpMessages.push({ text: displayText });
const content: (TextContent | ImageContent)[] = [{ type: "text", text }];
if (images && images.length > 0) {
content.push(...images);
if (normalizedImages && normalizedImages.length > 0) {
content.push(...normalizedImages);
}
this.agent.followUp({
role: "user",
@@ -4932,16 +4959,17 @@ export class AgentSession {
attribution: message.attribution ?? "agent",
timestamp: Date.now(),
};
const normalizedAppMessage = await this.#normalizeAgentMessageImages(appMessage);
if (this.isStreaming) {
if (options?.deliverAs === "nextTurn") {
this.#queueHiddenNextTurnMessage(appMessage, options?.triggerTurn ?? false);
this.#queueHiddenNextTurnMessage(normalizedAppMessage, options?.triggerTurn ?? false);
return;
}
if (options?.deliverAs === "followUp") {
this.agent.followUp(appMessage);
this.agent.followUp(normalizedAppMessage);
} else {
this.agent.steer(appMessage);
this.agent.steer(normalizedAppMessage);
}
return;
}
@@ -4949,16 +4977,16 @@ export class AgentSession {
if (options?.deliverAs === "nextTurn") {
if (options?.triggerTurn) {
if (this.#clientBridge?.deferAgentInitiatedTurns && !this.#allowAcpAgentInitiatedTurns) {
this.#queueHiddenNextTurnMessage(appMessage, false);
this.#queueHiddenNextTurnMessage(normalizedAppMessage, false);
return;
}
await this.agent.prompt(appMessage);
await this.agent.prompt(normalizedAppMessage);
return;
}
this.agent.appendMessage(appMessage);
this.agent.appendMessage(normalizedAppMessage);
this.sessionManager.appendCustomMessageEntry(
message.customType,
message.content,
normalizedAppMessage.customType,
normalizedAppMessage.content,
message.display,
message.details,
message.attribution ?? "agent",
@@ -4968,17 +4996,17 @@ export class AgentSession {
if (options?.triggerTurn) {
if (this.#clientBridge?.deferAgentInitiatedTurns && !this.#allowAcpAgentInitiatedTurns) {
this.#queueHiddenNextTurnMessage(appMessage, false);
this.#queueHiddenNextTurnMessage(normalizedAppMessage, false);
return;
}
await this.agent.prompt(appMessage);
await this.agent.prompt(normalizedAppMessage);
return;
}
this.agent.appendMessage(appMessage);
this.agent.appendMessage(normalizedAppMessage);
this.sessionManager.appendCustomMessageEntry(
message.customType,
message.content,
normalizedAppMessage.customType,
normalizedAppMessage.content,
message.display,
message.details,
message.attribution ?? "agent",
@@ -2,7 +2,7 @@ import * as fs from "node:fs/promises";
import type { ImageContent } from "@oh-my-pi/pi-ai";
import { formatBytes, readImageMetadata, SUPPORTED_IMAGE_MIME_TYPES } from "@oh-my-pi/pi-utils";
import { resolveReadPath } from "../tools/path-utils";
import { formatDimensionNote, resizeImage } from "./image-resize";
import { formatDimensionNote, resizeImage, type ImageResizeOptions } from "./image-resize";
export const MAX_IMAGE_INPUT_BYTES = 20 * 1024 * 1024;
export const SUPPORTED_INPUT_IMAGE_MIME_TYPES = SUPPORTED_IMAGE_MIME_TYPES;
@@ -50,6 +50,36 @@ export async function ensureSupportedImageInput(image: ImageContent): Promise<Im
}
}
export interface NormalizeModelContextImagesOptions {
resize?: ImageResizeOptions;
}
/**
* Normalize image blocks before they enter agent/model context. This keeps
* provider request construction from having to resize an unbounded batch of
* large images on the streaming hot path. Images are processed sequentially on
* purpose: `resizeImage` may fan out multiple encoders for one image, so the
* outer image batch must stay bounded.
*/
export async function normalizeModelContextImages(
images: ImageContent[] | undefined,
options?: NormalizeModelContextImagesOptions,
): Promise<ImageContent[] | undefined> {
if (!images || images.length === 0) return undefined;
const normalized: ImageContent[] = [];
for (const image of images) {
try {
const resized = await resizeImage(image, options?.resize);
normalized.push({ type: "image", data: resized.data, mimeType: resized.mimeType });
} catch {
// Preserve existing caller behavior for decode/resize failures: keep the
// user's image block rather than dropping it from the turn.
normalized.push(image);
}
}
return normalized;
}
export async function loadImageInput(options: LoadImageInputOptions): Promise<LoadedImageInput | null> {
const maxBytes = options.maxBytes ?? MAX_IMAGE_INPUT_BYTES;
const resolvedPath = options.resolvedPath ?? resolveReadPath(options.path, options.cwd);
@@ -1,5 +1,5 @@
import { describe, expect, test } from "bun:test";
import { ensureSupportedImageInput } from "../src/utils/image-loading";
import { ensureSupportedImageInput, normalizeModelContextImages } from "../src/utils/image-loading";
// 1x1 red PNG (69 bytes). Bun.Image sniffs format from bytes, so we can pass
// this with a non-supported MIME type and the conversion path runs over the
@@ -7,6 +7,17 @@ import { ensureSupportedImageInput } from "../src/utils/image-loading";
const RED_1X1_PNG_BASE64 =
"iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAIAAACQd1PeAAAADElEQVR4nGP4z8AAAAMBAQDJ/pLvAAAAAElFTkSuQmCC";
async function makeRedPng(width: number, height: number): Promise<string> {
const seed = Buffer.from(RED_1X1_PNG_BASE64, "base64");
const upscaled = await new Bun.Image(seed).resize(width, height, { filter: "nearest" }).png().bytes();
return Buffer.from(upscaled).toBase64();
}
async function dimensions(image: { data: string }): Promise<{ width: number; height: number }> {
const metadata = await new Bun.Image(Buffer.from(image.data, "base64")).metadata();
return { width: metadata.width, height: metadata.height };
}
describe("ensureSupportedImageInput", () => {
test("passes supported mime types through unchanged", async () => {
const input = { type: "image" as const, data: RED_1X1_PNG_BASE64, mimeType: "image/png" };
@@ -41,3 +52,22 @@ describe("ensureSupportedImageInput", () => {
expect(result).toBeNull();
});
});
describe("normalizeModelContextImages", () => {
test("downscales multiple large images before model context", async () => {
const wide = { type: "image" as const, data: await makeRedPng(2000, 1500), mimeType: "image/png" };
const tall = { type: "image" as const, data: await makeRedPng(1200, 2200), mimeType: "image/png" };
const result = await normalizeModelContextImages([wide, tall]);
expect(result).toHaveLength(2);
expect(result?.[0]?.type).toBe("image");
expect(result?.[1]?.type).toBe("image");
const wideDims = await dimensions(result![0]!);
const tallDims = await dimensions(result![1]!);
expect(wideDims.width).toBeLessThanOrEqual(1568);
expect(wideDims.height).toBeLessThanOrEqual(1568);
expect(tallDims.width).toBeLessThanOrEqual(1568);
expect(tallDims.height).toBeLessThanOrEqual(1568);
});
});