diff --git a/packages/agent/CHANGELOG.md b/packages/agent/CHANGELOG.md index 687e1cb1f..9d0472231 100644 --- a/packages/agent/CHANGELOG.md +++ b/packages/agent/CHANGELOG.md @@ -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 diff --git a/packages/agent/src/agent-loop.ts b/packages/agent/src/agent-loop.ts index 052b72afb..fc859bb0c 100644 --- a/packages/agent/src/agent-loop.ts +++ b/packages/agent/src/agent-loop.ts @@ -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; +type CloneableRecord = Record; + +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; malformed: boolean } { +function coerceToolResult(raw: unknown): { result: AgentToolResult; malformed: boolean } { const rawObj = raw && typeof raw === "object" ? (raw as Record) : null; const rawContent = rawObj?.content; const details = rawObj && "details" in rawObj ? rawObj.details : {}; @@ -98,8 +168,12 @@ function coerceToolResult(raw: unknown): { result: AgentToolResult; 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; 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 { 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(); const responseIterator = response[Symbol.asyncIterator](); + const finishAbortedStream = async (): Promise => { + 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(); 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): 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, + 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, context: AgentContext, config: AgentLoopConfig, stream: EventStream, 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 | null = null; + let steeringCheckTail: Promise = 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, 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; 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 { + const reason = abortReasonText(signal); + return { + content: [{ type: "text", text: `Tool was not executed because the run was aborted: ${reason}.` }], + details: {}, + }; +} + function createSkippedToolResult(): AgentToolResult { return { content: [{ type: "text", text: "Skipped due to queued user message." }], diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index b67b895bc..272a04f91 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -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 `.` (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 diff --git a/packages/ai/src/providers/anthropic.ts b/packages/ai/src/providers/anthropic.ts index 20a5fad31..652a854cd 100644 --- a/packages/ai/src/providers/anthropic.ts +++ b/packages/ai/src/providers/anthropic.ts @@ -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; +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; @@ -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(); + +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 = 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(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( 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, 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).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, - cacheControl, - ); - cacheBreakpointsUsed++; + if ( + applyCacheControlToLastTextBlock( + message.content as Array, + 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 | 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): 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 | 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 | undefined { switch (scalarType) { @@ -2768,14 +2893,6 @@ function anthropicPerTypeKeep(scalarType: string | undefined): Set | 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>, +): unknown { + if (Array.isArray(schema)) return schema.map(entry => normalizeAnthropicToolSchemaNode(entry, cache)); if (!isRecord(schema)) return schema; - const slot = schema as Record | undefined>; - const existing = slot[kAnthropicToolNormal]; + const existing = cache.get(schema); if (existing !== undefined) return existing; const result: Record = {}; - // 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; 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; 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): 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): 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) { diff --git a/packages/ai/test/anthropic-stream-envelope.test.ts b/packages/ai/test/anthropic-stream-envelope.test.ts index 42ea9a377..a10fc8df6 100644 --- a/packages/ai/test/anthropic-stream-envelope.test.ts +++ b/packages/ai/test/anthropic-stream-envelope.test.ts @@ -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[] = [ diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index a3ba00a7d..3411f6528 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -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 diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index 7d56f4228..f4ea08104 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -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 { + 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(message: T): Promise { + 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 { + 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 { + 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", diff --git a/packages/coding-agent/src/utils/image-loading.ts b/packages/coding-agent/src/utils/image-loading.ts index 9f172188d..4be0a534e 100644 --- a/packages/coding-agent/src/utils/image-loading.ts +++ b/packages/coding-agent/src/utils/image-loading.ts @@ -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 { + 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 { const maxBytes = options.maxBytes ?? MAX_IMAGE_INPUT_BYTES; const resolvedPath = options.resolvedPath ?? resolveReadPath(options.path, options.cwd); diff --git a/packages/coding-agent/test/image-input-normalization.test.ts b/packages/coding-agent/test/image-input-normalization.test.ts index 8e4a0ca2a..39bdddd99 100644 --- a/packages/coding-agent/test/image-input-normalization.test.ts +++ b/packages/coding-agent/test/image-input-normalization.test.ts @@ -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 { + 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); + }); +});