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:
@@ -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
@@ -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." }],
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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[] = [
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user