-> split header and data
header, _, b64data = url.partition(",")
media_type = (
header.split(";")[0].replace("data:", "")
or "image/jpeg"
)
anthropic_parts.append(
{
"type": "image",
"source": {
"type": "base64",
"media_type": media_type,
"data": b64data,
},
}
)
else:
# Remote URL -- Anthropic supports url source type natively.
# See: https://docs.anthropic.com/en/docs/build-with-claude/vision#url-based-images
anthropic_parts.append(
{
"type": "image",
"source": {
"type": "url",
"url": url,
},
}
)
elif part.get("type") == "input_document":
# `input_document` is Studio's normalised content type
# for PDFs / docs. The frontend sends either
# `{type:"input_document", file_data:"data:application/pdf;base64,..."}`
# or `{type:"input_document", file_url:"https://..."}`,
# plus optional `filename` and `media_type`.
# Translate to Anthropic's native `document` block.
url = part.get("file_url") or ""
data_uri = part.get("file_data") or ""
title = part.get("filename")
# Treat any "data:" URI with no actual base64
# payload (`data:application/pdf;base64,` or
# whitespace-only) as missing so the file_url
# branch below can take over. Matches the
# OpenAI-side fallback so a malformed inline
# payload + valid remote URL still attaches.
data_uri_valid = False
b64data = ""
header = ""
if data_uri.startswith("data:"):
header, _, b64data = data_uri.partition(",")
data_uri_valid = bool(b64data.strip())
if data_uri_valid:
media_type = (
part.get("media_type")
or header.split(";")[0].replace("data:", "")
or "application/pdf"
)
doc_block: dict[str, Any] = {
"type": "document",
"source": {
"type": "base64",
"media_type": media_type,
"data": b64data,
},
# Opt into Anthropic's natural-citation
# pipeline; without this no citations_delta
# events fire. See
# https://platform.claude.com/docs/en/build-with-claude/citations
"citations": {"enabled": True},
}
if title:
doc_block["title"] = title
anthropic_parts.append(doc_block)
elif url:
doc_block = {
"type": "document",
"source": {
"type": "url",
"url": url,
},
"citations": {"enabled": True},
}
if title:
doc_block["title"] = title
anthropic_parts.append(doc_block)
# Skip whole-message append when nothing usable survived.
# An empty content array (e.g. user dropped only an unparseable
# `input_document`) would 400 the Anthropic API with
# "messages.N.content: at least one block is required".
if anthropic_parts:
filtered.append({"role": msg["role"], "content": anthropic_parts})
else:
filtered.append(msg)
# Claude 4.7 family removed temperature / top_p / top_k entirely.
# The earlier guard only handled top_k; temperature is now also
# rejected with 400 "temperature is deprecated for this model".
# Latch the match once and reuse it everywhere temperature or
# top_k would otherwise be set — including the thinking-mode
# override below, which used to force temperature=1.
sampling_removed = bool(_ANTHROPIC_4_7_SAMPLING_REMOVED.match(model))
body: dict[str, Any] = {
"model": model,
"messages": filtered,
"max_tokens": max_tokens or 1024, # required by Anthropic
"stream": True,
}
if not sampling_removed:
body["temperature"] = temperature
if top_k is not None and top_k > 0 and not sampling_removed:
body["top_k"] = top_k
# Anthropic only caches a prefix when at least one cache_control
# marker is attached to it — the frontend defaults
# enable_prompt_caching to True for Anthropic, so treat `None` the
# same as True here (callers that don't set the flag still get
# caching). Pass False explicitly to opt out.
prompt_caching_enabled = enable_prompt_caching is not False
# Anthropic accepts an optional `ttl` on each cache_control marker
# (default is the 5m ephemeral pool; set "1h" to land in the 1h
# pool instead). Per the prompt-caching docs, 1h cache writes are
# billed at 2x base input vs 1.25x for 5m, but reads are 0.1x for
# both. The 1h pool is the right pick when conversations span
# multiple short bursts more than 5 minutes apart -- the read
# discount makes up for the 1.6x write premium after a single
# additional hit. Anything other than the known TTL strings is
# dropped to avoid sending a malformed marker.
#
# The `extended-cache-ttl-2025-04-11` beta header that originally
# gated 1h TTL has been promoted to GA: as of 2026-05 the live
# API accepts `ttl: "1h"` without any beta opt-in. Verified
# against api.anthropic.com on claude-opus-4-7 (status 200 +
# `ephemeral_1h_input_tokens` populated). The test below pins
# the contract by asserting the header is NOT on the wire so a
# future regression that reintroduces the gate would surface
# before users see a 400.
cache_marker: dict[str, Any] = {"type": "ephemeral"}
if prompt_cache_ttl in ("5m", "1h"):
cache_marker["ttl"] = prompt_cache_ttl
if system:
if prompt_caching_enabled:
# System block is the most stable prefix across turns, so
# it gets its own breakpoint. Skipped when system is
# empty — there's nothing to cache, and an empty marker
# is a no-op.
body["system"] = [
{
"type": "text",
"text": system,
"cache_control": dict(cache_marker),
}
]
else:
body["system"] = system
if prompt_caching_enabled and filtered:
# Second breakpoint at the end of the conversation. Anthropic
# caches the longest matching prefix up to a cache_control
# marker; placing one on the latest message means turn N+1
# rehydrates everything up through turn N from cache instead
# of recomputing it. This is what makes caching actually work
# when the system prompt is empty or shorter than Anthropic's
# ~1024-token cache floor — the conversation history carries
# the bulk of the input tokens. Anthropic allows up to 4
# breakpoints per request; we use at most 2 (system + tail).
last_msg = filtered[-1]
content = last_msg.get("content")
if isinstance(content, str):
last_msg["content"] = [
{
"type": "text",
"text": content,
"cache_control": dict(cache_marker),
}
]
elif isinstance(content, list) and content:
# Don't mutate the caller's list. Rebuild the tail with
# cache_control attached to the final block so an
# upstream image-bearing turn still cleanly slots into
# the cache as part of the conversational prefix.
head = list(content[:-1])
tail = content[-1]
if isinstance(tail, dict):
head.append({**tail, "cache_control": dict(cache_marker)})
else:
head.append(tail)
last_msg["content"] = head
thinking_spec = _anthropic_thinking_spec(model)
allowed_efforts = (
thinking_spec.efforts
if thinking_spec
else ("none", "low", "medium", "high")
)
effort = reasoning_effort if reasoning_effort in allowed_efforts else None
# Claude 4.6 Opus/Sonnet accept top-tier adaptive effort as "max" only;
# "xhigh" is rejected (supported on Claude 4.7). Map our shared "xhigh"
# semantic to "max" for 4.6 outbound requests while still accepting
# both in ``allowed_efforts`` for persisted / cross-provider UI state.
if effort == "xhigh" and model.startswith(
("claude-opus-4-6", "claude-sonnet-4-6")
):
effort = "max"
if effort is None:
if enable_thinking is False:
effort = "none"
elif enable_thinking is True:
effort = "medium"
# Normalize one semantic Thinking control into Anthropic's two model-era
# APIs: adaptive effort on Claude 4.6/4.7, manual budget_tokens on 4.5.
if effort and effort != "none":
# Anthropic rejects top_k whenever thinking is enabled.
body.pop("top_k", None)
# Earlier families (4.5/4.6) require temperature=1 when
# thinking is enabled and forbid top_p in the same request:
# "temperature and top_p cannot both be specified for this
# model. Please use only one."
# On Claude 4.7, temperature was removed entirely — sending
# any value (including 1) returns 400 — so skip the override
# there and let the model use its default sampling.
if not sampling_removed:
body["temperature"] = 1
body.pop("top_p", None)
if thinking_spec and thinking_spec.kind == "adaptive":
# `display` defaults to "omitted" on Claude Opus 4.7 (per the
# adaptive-thinking docs) — without an explicit opt-in the
# API emits an empty thinking block plus a signature_delta,
# so our SSE handler would surface a stray
# and the reasoning panel would stay blank. Force
# "summarized" so 4.7 streams thinking_delta events like
# 4.6 does. On 4.6 / Sonnet 4.6 this is the default, so
# setting it explicitly is harmless.
body["thinking"] = {"type": "adaptive", "display": "summarized"}
# Per the Messages API reference, the effort knob for
# adaptive thinking lives under `output_config.effort` —
# NOT as a top-level field. Sending `effort: ...` directly
# produces a 400 "effort: Extra inputs are not permitted".
# Allowed values: low | medium | high | xhigh | max. See:
# https://platform.claude.com/docs/en/api/messages
body["output_config"] = {"effort": effort}
elif thinking_spec and thinking_spec.kind == "manual":
budget_tokens = {"low": 1024, "medium": 2048, "high": 4096}[effort]
body["thinking"] = {
"type": "enabled",
"budget_tokens": budget_tokens,
}
# Anthropic requires max_tokens to be strictly greater than
# thinking.budget_tokens on the manual-thinking path.
if body.get("max_tokens", 0) <= budget_tokens:
body["max_tokens"] = budget_tokens + 1024
# Anthropic server-side web_search — see
# https://platform.claude.com/docs/en/agents-and-tools/tool-use/web-search-tool
# The tool type is date-pinned per model family. Newer Opus /
# Sonnet 4.6 + 4.7 accept `web_search_20260209` with dynamic
# filtering (Claude writes code to filter results before they
# reach context); everything else uses `web_search_20250305`.
# `_anthropic_web_search_version` picks the right one. Anthropic
# dispatches search calls server-side, returning server_tool_use
# + web_search_tool_result blocks in the SSE stream, plus
# url-citation annotations on text deltas. We translate all of
# that into our local _toolEvent shape so the chat UI renders
# web_search exactly like OpenAI's path.
if enabled_tools and "web_search" in enabled_tools:
anthropic_tools = list(body.get("tools") or [])
anthropic_tools.append(
{
"type": _anthropic_web_search_version(model),
"name": "web_search",
"max_uses": 5,
}
)
body["tools"] = anthropic_tools
# Anthropic server-side web_fetch reads a single URL (text/PDF)
# and returns a `web_fetch_tool_result` document block. Opt in
# via `enabled_tools=["web_fetch"]`; no beta header required.
# `_anthropic_web_fetch_version` picks `web_fetch_20260209`
# (dynamic filtering) for Opus 4.6/4.7 + Sonnet 4.6, falling
# back to `web_fetch_20250910` elsewhere; mismatched variants
# return 400 so the per-model picker is required.
web_fetch_enabled = bool(enabled_tools and "web_fetch" in enabled_tools)
if web_fetch_enabled:
anthropic_tools = list(body.get("tools") or [])
anthropic_tools.append(
{
"type": _anthropic_web_fetch_version(model),
"name": "web_fetch",
"max_uses": 5,
}
)
body["tools"] = anthropic_tools
# Anthropic server-side code execution — see
# https://platform.claude.com/docs/en/agents-and-tools/tool-use/code-execution-tool
# The tool type is date-pinned per model family.
# `_anthropic_code_execution_version` picks `code_execution_20260120`
# for Opus 4.5+ / Sonnet 4.5+ / Opus 4.7 / Sonnet 4.6 (adds REPL
# state persistence + programmatic tool calling) and falls back
# to `code_execution_20250825` everywhere else. Both versions
# run Python + bash + str_replace file edits inside a 5 GB
# sandboxed container per request, with no internet access, and
# both are unlocked by the same `code-execution-2025-08-25`
# `anthropic-beta` header set further down. On the SSE stream
# Anthropic emits two sub-tool names -- `bash_code_execution`
# and `text_editor_code_execution` -- wrapped in the standard
# server_tool_use / *_tool_result block shape.
# v1 wires the tool only; file uploads (container_upload
# content blocks and generated-file retrieval via the Files
# API) are a deliberate follow-up.
code_execution_enabled = bool(
enabled_tools and "code_execution" in enabled_tools
)
if code_execution_enabled:
anthropic_tools = list(body.get("tools") or [])
anthropic_tools.append(
{
"type": _anthropic_code_execution_version(model),
"name": "code_execution",
}
)
body["tools"] = anthropic_tools
# Reuse the prior turn's container so filesystem state
# (files written, packages installed, variables set)
# persists across turns of the same thread. Anthropic
# exposes the container id on the Message object's
# top-level `container.id`; on the SSE stream we latch it
# off `message_start.message.container.id` further down
# and emit a `container_ready` _toolEvent so the chat
# adapter persists it on the thread record. A stale id
# (container expired / not found) surfaces as a 4xx
# below, where we emit `container_invalidated` and let
# the next turn fall back to auto-create.
if anthropic_code_exec_container_id:
body["container"] = anthropic_code_exec_container_id
# Server-side context compaction — see
# https://platform.claude.com/docs/en/build-with-claude/compaction
# Beta as of `compact-2026-01-12`. When `compaction_threshold` is
# provided AND the model accepts compaction (Opus 4.6+ / 4.7,
# Sonnet 4.6, Mythos preview), attach
# `context_management.edits[{type:"compact_20260112", trigger:
# {type:"input_tokens", value:N}}]` to the body. Anthropic runs
# the compaction step server-side once the rendered prompt
# crosses the threshold and replies with a top-level
# `context_management` block plus `usage.iterations[]` so we can
# account per-iteration. Below-min thresholds get clamped up to
# 50K so the request doesn't 400.
compaction_active = (
compaction_threshold is not None
and compaction_threshold > 0
and _anthropic_supports_compaction(model)
)
if compaction_active and compaction_threshold is not None:
trigger_value = max(
int(compaction_threshold),
_ANTHROPIC_COMPACTION_MIN,
)
body["context_management"] = {
"edits": [
{
"type": _ANTHROPIC_COMPACTION_TYPE,
"trigger": {
"type": "input_tokens",
"value": trigger_value,
},
}
]
}
# fast_mode is Opus 4.6/4.7 only; silently drop elsewhere.
# Incompatible with the Priority service_tier (frontend gate
# prevents both at once; backend lets Anthropic 400 if combined).
fast_mode_active = bool(fast_mode) and _anthropic_supports_fast_mode(model)
if fast_mode_active:
body["speed"] = "fast"
url = f"{self.base_url}/messages"
completion_id = f"chatcmpl-anthropic-{model.replace('/', '-')}"
# Log the outgoing config keys (not the messages themselves) so we
# can prove which thinking/effort fields actually reached the wire.
# If Anthropic skips reasoning despite a configured effort, this
# tells us whether we sent the field or dropped it on the floor.
logger.info(
"Anthropic request shape (model=%s, has_thinking=%s, thinking=%s, "
"output_config=%s, temperature=%s, has_top_p=%s, has_top_k=%s, "
"max_tokens=%s)",
model,
"thinking" in body,
body.get("thinking"),
body.get("output_config"),
body.get("temperature"),
"top_p" in body,
"top_k" in body,
body.get("max_tokens"),
)
# Translate Anthropic stop reasons onto the OpenAI chat-completions
# `finish_reason` vocabulary. `pause_turn` maps to None so the
# adapter does NOT emit a finish_reason chunk: pause_turn means
# Claude paused a long server-tool turn (web_search / web_fetch)
# and will continue once the user (or our retry) sends back the
# partial assistant message. Forwarding it as "stop" makes the
# OpenAI client think the answer is done and truncates the
# rendered message. `refusal` maps to "content_filter" as the
# nearest semantic match. See
# https://platform.claude.com/docs/en/api/messages#response-stop-reason
_finish_reason_map: dict[str, Optional[str]] = {
"end_turn": "stop",
"max_tokens": "length",
"stop_sequence": "stop",
"tool_use": "tool_calls",
"refusal": "content_filter",
"pause_turn": None,
}
logger.info("Proxying Anthropic Messages API to %s (model=%s)", url, model)
request_headers = self._auth_headers()
# Anthropic accepts comma-separated beta features in a single
# `anthropic-beta` header. Merge our flags onto whatever the
# registry's extra_headers contributed (currently nothing on
# the beta axis, just anthropic-version) so future betas
# added at the registry level keep working.
existing_beta = request_headers.get("anthropic-beta", "").strip()
beta_parts = (
[p.strip() for p in existing_beta.split(",") if p.strip()]
if existing_beta
else []
)
if code_execution_enabled and _ANTHROPIC_CODE_EXECUTION_BETA not in beta_parts:
beta_parts.append(_ANTHROPIC_CODE_EXECUTION_BETA)
if compaction_active and _ANTHROPIC_COMPACTION_BETA not in beta_parts:
beta_parts.append(_ANTHROPIC_COMPACTION_BETA)
if fast_mode_active and _ANTHROPIC_FAST_MODE_BETA not in beta_parts:
beta_parts.append(_ANTHROPIC_FAST_MODE_BETA)
if beta_parts:
request_headers["anthropic-beta"] = ",".join(beta_parts)
try:
async with _http_client.stream(
"POST",
url,
json = body,
headers = request_headers,
timeout = self._stream_timeout,
) as response:
if response.status_code != 200:
error_body = await response.aread()
error_text = error_body.decode("utf-8", errors = "replace")
logger.error(
"Anthropic returned %d: %s",
response.status_code,
error_text[:500],
)
# Stale container detection (mirror of the OpenAI
# path). When we sent a `container` field and the
# response is 4xx with any hint that the id is
# expired / missing, emit container_invalidated so
# the chat adapter clears the stored id and the
# next turn falls back to auto-create.
if (
anthropic_code_exec_container_id
and 400 <= response.status_code < 500
):
lowered = error_text.lower()
if "container" in lowered and (
"expired" in lowered
or "not_found" in lowered
or "not found" in lowered
or "no such container" in lowered
or "invalid" in lowered
):
yield (
f"data: "
f"{_json.dumps({'id': completion_id, 'object': 'chat.completion.chunk', 'choices': [{'index': 0, 'delta': {}, 'finish_reason': None}], '_toolEvent': {'type': 'container_invalidated'}})}"
)
yield _error_sse_line(
response.status_code, error_text, self.provider_type
)
return
# NOTE: same manual __anext__ loop as stream_chat_completion — see comment there.
lines_gen = response.aiter_lines().__aiter__()
thinking_open = False
# Diagnostic counters for the next time the user reports
# "no thinking content" — distinguishes "Anthropic never sent
# thinking_delta" from "frontend didn't render the chunks".
event_counts: dict[str, int] = {}
# web_search state. Anthropic emits the query inside an
# `input_json_delta` stream on a `server_tool_use` content
# block, then a separate `web_search_tool_result` block
# with the URL list. Unlike OpenAI we get per-call results
# directly, so each tool card carries its own citations.
# `current_server_tool_use`: {id, name, partial_json_buffer}
# `current_result_block`: {tool_use_id, results}
# Both go to None when the matching content_block_stop fires.
current_server_tool_use: Optional[dict[str, Any]] = None
current_result_block: Optional[dict[str, Any]] = None
web_search_calls: dict[str, dict[str, Any]] = {}
# code_execution state. Anthropic's
# `code_execution_20250825` tool emits the same
# server_tool_use → *_tool_result block shape as
# web_search, but the server_tool_use carries one of
# two sub-tool names (`bash_code_execution` or
# `text_editor_code_execution`) and the result block
# type matches (`bash_code_execution_tool_result` /
# `text_editor_code_execution_tool_result`). Kept
# parallel to web_search state so the two paths don't
# collide when both pills are on in the same turn.
current_code_exec_use: Optional[dict[str, Any]] = None
current_code_exec_result: Optional[dict[str, Any]] = None
code_execution_calls: dict[str, dict[str, Any]] = {}
# web_fetch state. Same server_tool_use → *_tool_result
# block shape as web_search but the server_tool_use
# carries name="web_fetch" and the result block is
# `web_fetch_tool_result` with content.type=
# `web_fetch_result` (success) or `web_fetch_tool_error`
# (failure). Kept separate from web_search state so a
# turn that uses both does not collide.
current_web_fetch_use: Optional[dict[str, Any]] = None
current_web_fetch_result: Optional[dict[str, Any]] = None
web_fetch_calls: dict[str, dict[str, Any]] = {}
# Compaction state. Server-side compaction emits a
# `{type:"compaction", content:"..."}` content block
# whenever it runs. The summary text can land on the
# start event AND/OR via text_delta events on the same
# block (Anthropic's wire format is permissive here).
# Accumulate in `current_compaction["content"]` and emit
# on content_block_stop so the chat-adapter can persist
# it onto the assistant message for round-tripping on
# the next turn.
current_compaction: Optional[dict[str, Any]] = None
compaction_blocks_seen = 0
# Document citations from ``citations_delta`` events.
# Deduped by type-specific anchor key; inline [N] is
# injected after each cited run, and the full list is
# forwarded as a synthetic document_citations tool_event
# on message_stop for the Sources panel.
document_citations: list[dict[str, Any]] = []
# Counts surfaced in the final log line so reports of
# "Code execution did nothing" can be triaged at a
# glance. generated_files_count is interesting for the
# future Files API PR — when bash creates files inside
# the container, they show up as file_id entries on
# bash_code_execution_result.content, and v1 drops
# them. Track the count so we know how often it would
# have mattered.
code_execution_generated_files = 0
# Container id captured from `message_start.message.container.id`
# when code_execution is enabled. Emit a `container_ready`
# _toolEvent on first sight so the chat adapter persists it
# on the thread record. Only emitted when the value differs
# from the inbound id — no churn on reuse.
latched_container_id: Optional[str] = None
container_id_emitted = False
# Cache usage tracking. message_start carries the input
# accounting (incl. cache_creation_input_tokens and
# cache_read_input_tokens); message_delta carries cumulative
# output_tokens. Both are surfaced in the "stream complete"
# log so prompt caching can be verified per-request without
# opening the Anthropic dashboard.
last_usage: dict[str, Any] = {}
def _content_chunk(text: str) -> str:
chunk = {
"id": completion_id,
"object": "chat.completion.chunk",
"choices": [
{
"index": 0,
"delta": {"content": text},
"finish_reason": None,
}
],
}
return f"data: {_json.dumps(chunk)}"
def _emit_tool_event(payload: dict[str, Any]) -> str:
chunk = {
"id": completion_id,
"object": "chat.completion.chunk",
"choices": [
{
"index": 0,
"delta": {},
"finish_reason": None,
}
],
"_toolEvent": payload,
}
return f"data: {_json.dumps(chunk)}"
def _format_web_search_results(
results: list[Any],
) -> str:
blocks: list[str] = []
for r in results:
if not isinstance(r, dict):
continue
if r.get("type") != "web_search_result":
continue
url = r.get("url", "")
title = r.get("title") or url
if not url:
continue
blocks.append(f"Title: {title}\nURL: {url}")
return "\n---\n".join(blocks)
def _format_web_fetch_result(inner: dict[str, Any]) -> str:
"""Render a `web_fetch_tool_result.content` payload
as the Title / URL / snippet block CodeExecutionToolUI
and parseSourcesFromResult already expect from the
web_search path.
Success shape (text):
{type: web_fetch_result, url, retrieved_at,
content: {type: document, source: {type: text,
media_type, data}, title?}}
Success shape (pdf): source.type=base64 + media_type=
application/pdf. We do not surface the base64
bytes; the title + url is enough for the source
pill, and the model still sees the document
contents on its side.
Error shape: {type: web_fetch_tool_error, error_code}.
"""
inner_type = inner.get("type") or ""
if inner_type == "web_fetch_tool_error":
return f"Error: {inner.get('error_code', 'unknown')}"
url = inner.get("url", "")
document = inner.get("content") or {}
title = ""
snippet = ""
if isinstance(document, dict):
title = document.get("title") or ""
source = document.get("source") or {}
if isinstance(source, dict):
media_type = source.get("media_type") or ""
data = source.get("data") or ""
# Inline a short text preview so the source
# pill carries usable context; skip for PDFs
# since the body is base64-encoded.
if (
media_type.startswith("text/")
and isinstance(data, str)
and data
):
snippet = data[:240].strip()
# Frontend parseSourcesFromResult only emits a source
# pill when both `Title:` and `URL:` are present, so
# fall back to the URL when Anthropic omits the
# document title (matches the web_search formatter).
if not title and url:
title = url
parts: list[str] = []
if title:
parts.append(f"Title: {title}")
if url:
parts.append(f"URL: {url}")
if snippet:
parts.append(f"Snippet: {snippet}")
return "\n".join(parts) if parts else "(fetch complete)"
def _format_code_execution_result(
inner: dict[str, Any],
) -> str:
"""Render an Anthropic code-execution result block as
the preformatted text payload the frontend's
CodeExecutionToolUI displays inside a . Handles
bash, text_editor (view/create/str_replace), and the
matching error variants.
"""
inner_type = inner.get("type") or ""
if inner_type.endswith("_error"):
return f"Error: {inner.get('error_code', 'unknown')}"
if inner_type == "bash_code_execution_result":
stdout = inner.get("stdout") or ""
stderr = inner.get("stderr") or ""
return_code = inner.get("return_code")
parts: list[str] = []
if stdout:
parts.append(stdout)
if stderr:
parts.append(f"--- stderr ---\n{stderr}")
if isinstance(return_code, int) and return_code != 0:
parts.append(f"return_code: {return_code}")
return "\n".join(parts) if parts else "(no output)"
if inner_type == "text_editor_code_execution_result":
# view: file content; create: is_file_update flag;
# str_replace: diff `lines` list. The matching
# server_tool_use carries the command + path, but
# that's encoded into the tool_start arguments
# already — here we only format the result body.
if "lines" in inner and isinstance(inner.get("lines"), list):
return "\n".join(str(line) for line in inner["lines"])
if "is_file_update" in inner:
return (
"Updated" if inner.get("is_file_update") else "Created"
)
content_field = inner.get("content")
if isinstance(content_field, str):
return content_field
return "(file operation complete)"
return "(code execution complete)"
try:
while True:
try:
line = await lines_gen.__anext__()
except StopAsyncIteration:
break
if not line or line.startswith("event:"):
continue
if not line.startswith("data:"):
continue
data_str = line[len("data:") :].strip()
if not data_str:
continue
try:
event = _json.loads(data_str)
except _json.JSONDecodeError:
continue
event_type = event.get("type")
if event_type == "content_block_delta":
delta_kind = (event.get("delta") or {}).get("type")
key = f"{event_type}:{delta_kind}"
else:
key = event_type or ""
event_counts[key] = event_counts.get(key, 0) + 1
# message_start carries the input-side usage block
# including cache_creation_input_tokens and
# cache_read_input_tokens. message_delta updates
# output_tokens (and may overwrite the input fields
# with final values). Merge both into last_usage.
if event_type == "message_start":
start_usage = (event.get("message") or {}).get("usage")
if isinstance(start_usage, dict):
last_usage.update(start_usage)
if event_type == "content_block_start":
content_block = event.get("content_block") or {}
block_type = content_block.get("type")
block_name = content_block.get("name")
if (
block_type == "server_tool_use"
and block_name == "web_search"
):
tool_use_id = content_block.get("id", "") or (
f"ws_{len(web_search_calls)}"
)
current_server_tool_use = {
"id": tool_use_id,
"buffer": "",
}
web_search_calls[tool_use_id] = {
"query": "",
"results": [],
}
elif block_type == "web_search_tool_result":
tool_use_id = content_block.get("tool_use_id", "")
# Anthropic sometimes ships the full results
# list on the start event; sometimes deltas
# follow. Capture whatever is present and
# finalize on content_block_stop.
content = content_block.get("content") or []
current_result_block = {
"tool_use_id": tool_use_id,
"results": list(content)
if isinstance(content, list)
else [],
}
elif (
block_type == "server_tool_use"
and block_name == "web_fetch"
):
tool_use_id = content_block.get("id", "") or (
f"wf_{len(web_fetch_calls)}"
)
current_web_fetch_use = {
"id": tool_use_id,
"buffer": "",
}
web_fetch_calls[tool_use_id] = {
"url": "",
"result": None,
}
elif block_type == "web_fetch_tool_result":
tool_use_id = content_block.get("tool_use_id", "")
inner = content_block.get("content") or {}
current_web_fetch_result = {
"tool_use_id": tool_use_id,
"inner": inner if isinstance(inner, dict) else {},
}
elif block_type == "server_tool_use" and block_name in (
"bash_code_execution",
"text_editor_code_execution",
):
tool_use_id = content_block.get("id", "") or (
f"ce_{len(code_execution_calls)}"
)
kind = (
"bash"
if block_name == "bash_code_execution"
else "text_editor"
)
current_code_exec_use = {
"id": tool_use_id,
"kind": kind,
"buffer": "",
}
code_execution_calls[tool_use_id] = {
"kind": kind,
"arguments": {},
"result": None,
}
elif block_type in (
"bash_code_execution_tool_result",
"text_editor_code_execution_tool_result",
):
# Anthropic ships the full result content
# on the start event for code-exec result
# blocks (unlike web_search, which can
# split across deltas). Capture it and
# finalize on content_block_stop so the
# ordering matches the web_search path.
tool_use_id = content_block.get("tool_use_id", "")
inner = content_block.get("content") or {}
current_code_exec_result = {
"tool_use_id": tool_use_id,
"inner": inner if isinstance(inner, dict) else {},
}
elif block_type == "compaction":
# Server-side compaction emits a `compaction`
# content block on the assistant message.
# Anthropic may include the summary text on
# this start event AND/OR stream it via
# text_delta events on the same block. See
# https://platform.claude.com/docs/en/build-with-claude/compaction
# Capture either form; finalize and emit
# on content_block_stop. The chat-adapter
# persists the block onto the assistant
# message so the next turn's request
# carries it back -- Anthropic then skips
# re-compaction from scratch.
seed = content_block.get("content") or ""
current_compaction = {
"content": seed if isinstance(seed, str) else "",
}
elif event_type == "content_block_delta":
delta = event.get("delta", {})
delta_type = delta.get("type")
if delta_type == "thinking_delta":
# Anthropic streams extended-thinking content as
# thinking_delta events on a separate content
# block. Wrap as inline ... so
# the frontend's parseAssistantContent lifts it
# into the reasoning panel — same pattern as
# the OpenAI Responses path.
thinking_text = delta.get("thinking", "")
if thinking_text:
if not thinking_open:
thinking_text = f"{thinking_text}"
thinking_open = True
yield _content_chunk(thinking_text)
elif delta_type == "text_delta":
text = delta.get("text", "")
# text_deltas inside a compaction block
# carry the summary chunks; route them
# into the compaction buffer and DON'T
# yield them to the user-visible stream
# -- the summary is opaque internal
# state, not assistant prose.
if current_compaction is not None:
if text:
current_compaction["content"] += text
else:
# First text after a thinking block closes the
# tag we opened above. Anthropic emits
# a content_block_stop between blocks, but
# closing on the text_delta transition is more
# forgiving if events arrive out of order.
if thinking_open:
yield _content_chunk("")
thinking_open = False
if text:
yield _content_chunk(text)
# web_search citations: web_search_tool_result.
# User-doc citations: citations_delta below.
elif delta_type == "citations_delta":
# One citation per event; collapse onto a
# numbered footnote list and inject [N]
# inline. See
# https://platform.claude.com/docs/en/build-with-claude/citations
cit = delta.get("citation")
if isinstance(cit, dict):
key = _anthropic_citation_key(cit)
idx_for_marker: Optional[int] = None
for idx, existing in enumerate(
document_citations, start = 1
):
if existing.get("_key") == key:
idx_for_marker = idx
break
if idx_for_marker is None:
document_citations.append({**cit, "_key": key})
idx_for_marker = len(document_citations)
yield _content_chunk(f"[{idx_for_marker}]")
elif delta_type == "input_json_delta":
# Streamed partial_json carrying tool inputs
# — the search query for web_search, or the
# command/path/etc. for code execution.
# Route to whichever buffer is open. The two
# state slots are exclusive in practice
# (Anthropic doesn't interleave tool input
# streams), but checking both keeps the
# dispatch robust if that ever changes.
partial = delta.get("partial_json", "")
if current_server_tool_use is not None:
current_server_tool_use["buffer"] += partial
elif current_code_exec_use is not None:
current_code_exec_use["buffer"] += partial
elif current_web_fetch_use is not None:
current_web_fetch_use["buffer"] += partial
# signature_delta and any other delta types are
# intentionally skipped — they carry trust /
# verification metadata, not user-visible content.
elif event_type == "content_block_stop":
if current_server_tool_use is not None:
# End of the server_tool_use block — parse the
# accumulated input_json into a query and
# emit tool_start. The matching tool_end fires
# later when the web_search_tool_result block
# closes with the actual results.
buffer = current_server_tool_use["buffer"]
query = ""
if buffer:
try:
parsed = _json.loads(buffer)
if isinstance(parsed, dict):
q = parsed.get("query", "")
if isinstance(q, str):
query = q
except Exception:
query = ""
tool_use_id = current_server_tool_use["id"]
if tool_use_id in web_search_calls:
web_search_calls[tool_use_id]["query"] = query
yield _emit_tool_event(
{
"type": "tool_start",
"tool_name": "web_search",
"tool_call_id": tool_use_id,
"arguments": (
{"query": query} if query else {}
),
}
)
current_server_tool_use = None
elif current_result_block is not None:
# End of a web_search_tool_result — emit
# tool_end carrying the search results as
# Title:/URL: blocks. parseSourcesFromResult
# on the frontend lifts these into source
# pills at message tail.
tool_use_id = current_result_block["tool_use_id"]
results = current_result_block["results"]
if tool_use_id in web_search_calls:
web_search_calls[tool_use_id]["results"] = results
result_text = _format_web_search_results(results)
yield _emit_tool_event(
{
"type": "tool_end",
"tool_call_id": tool_use_id,
"result": (result_text or "(search complete)"),
}
)
current_result_block = None
elif current_code_exec_use is not None:
# End of a code-execution server_tool_use —
# parse the buffered input_json into a
# {command, path, ...} dict and emit
# tool_start. The matching tool_end fires
# on the result block's content_block_stop.
buffer = current_code_exec_use["buffer"]
parsed_args: dict[str, Any] = {}
if buffer:
try:
parsed_obj = _json.loads(buffer)
if isinstance(parsed_obj, dict):
parsed_args = parsed_obj
except Exception:
parsed_args = {}
tool_use_id = current_code_exec_use["id"]
kind = current_code_exec_use["kind"]
emit_args = {"kind": kind, **parsed_args}
if tool_use_id in code_execution_calls:
code_execution_calls[tool_use_id]["arguments"] = (
emit_args
)
yield _emit_tool_event(
{
"type": "tool_start",
"tool_name": "code_execution",
"tool_call_id": tool_use_id,
"arguments": emit_args,
}
)
current_code_exec_use = None
elif current_compaction is not None:
# End of a compaction block. Emit it as a
# synthetic tool_event so the chat-adapter
# can persist the {type:"compaction",
# content:"..."} payload onto the
# assistant message. The next turn's
# request body forwards the content_part
# verbatim and Anthropic recognises it
# as the prior compaction state.
compaction_blocks_seen += 1
yield _emit_tool_event(
{
"type": "compaction_block",
"content": current_compaction["content"],
}
)
current_compaction = None
elif current_code_exec_result is not None:
# End of a code-execution result block —
# format the inner result into the text
# payload CodeExecutionToolUI renders.
tool_use_id = current_code_exec_result["tool_use_id"]
inner = current_code_exec_result["inner"]
# Track generated-file count for the
# follow-up Files API PR. v1 drops them.
if isinstance(inner, dict):
file_blocks = inner.get("content")
if isinstance(file_blocks, list):
for entry in file_blocks:
if isinstance(entry, dict) and entry.get(
"file_id"
):
code_execution_generated_files += 1
result_text = _format_code_execution_result(
inner if isinstance(inner, dict) else {}
)
if tool_use_id in code_execution_calls:
code_execution_calls[tool_use_id]["result"] = (
result_text
)
yield _emit_tool_event(
{
"type": "tool_end",
"tool_call_id": tool_use_id,
"result": result_text,
}
)
current_code_exec_result = None
elif current_web_fetch_use is not None:
# End of the web_fetch server_tool_use —
# parse the buffered input_json into the
# URL the model asked Anthropic to fetch
# and emit tool_start. The matching
# tool_end fires on the result block's
# content_block_stop just below.
buffer = current_web_fetch_use["buffer"]
url = ""
if buffer:
try:
parsed = _json.loads(buffer)
if isinstance(parsed, dict):
probe = parsed.get("url", "")
if isinstance(probe, str):
url = probe
except Exception:
logger.debug(
"Failed to parse web_fetch input_json",
buffer = buffer,
)
url = ""
tool_use_id = current_web_fetch_use["id"]
if tool_use_id in web_fetch_calls:
web_fetch_calls[tool_use_id]["url"] = url
yield _emit_tool_event(
{
"type": "tool_start",
"tool_name": "web_fetch",
"tool_call_id": tool_use_id,
"arguments": ({"url": url} if url else {}),
}
)
current_web_fetch_use = None
elif current_web_fetch_result is not None:
# End of the web_fetch_tool_result —
# format Title / URL / snippet for the
# frontend source pill and emit tool_end.
# `inner` is sanitised to a dict at the
# matching content_block_start, and the
# formatter always returns a non-empty
# string (defaulting to "(fetch complete)"
# when no fields are present), so no
# extra fallback is needed here.
tool_use_id = current_web_fetch_result["tool_use_id"]
result_text = _format_web_fetch_result(
current_web_fetch_result["inner"]
)
if tool_use_id in web_fetch_calls:
web_fetch_calls[tool_use_id]["result"] = result_text
yield _emit_tool_event(
{
"type": "tool_end",
"tool_call_id": tool_use_id,
"result": result_text,
}
)
current_web_fetch_result = None
elif thinking_open:
# Close the tag when the thinking block
# ends, in case no text_delta follows (e.g.
# display=omitted on Claude 4.7, or thinking-
# only turns).
yield _content_chunk("")
thinking_open = False
elif event_type == "message_delta":
delta_usage = event.get("usage")
if isinstance(delta_usage, dict):
last_usage.update(delta_usage)
# When a fresh compaction has run, Anthropic
# publishes per-iteration token counts in
# `usage.iterations[]`. The top-level
# input_tokens / output_tokens only cover the
# `message` iteration, NOT the compaction
# passes — billing has to sum the whole
# array. See
# https://platform.claude.com/docs/en/build-with-claude/compaction
# Fold the compaction iterations into
# `compaction_input_tokens` / `compaction_output_tokens`
# so the cost surface can add them without
# re-walking the array (and so the closing
# log line names the figures).
iterations = delta_usage.get("iterations")
if isinstance(iterations, list):
c_in = 0
c_out = 0
for it in iterations:
if (
isinstance(it, dict)
and it.get("type") == "compaction"
):
c_in += int(it.get("input_tokens") or 0)
c_out += int(it.get("output_tokens") or 0)
if c_in or c_out:
last_usage["compaction_input_tokens"] = c_in
last_usage["compaction_output_tokens"] = c_out
# Anthropic reports the code_execution container
# id on `message_delta.delta.container.{id,
# expires_at}` (NOT on message_start — at start
# the container hasn't been provisioned yet).
# Latch on first sight and emit container_ready
# only when the value differs from the inbound
# id, so steady-state reuse doesn't re-write
# the same id to the thread record every turn.
delta_obj = event.get("delta") or {}
container_obj = delta_obj.get("container")
if (
isinstance(container_obj, dict)
and latched_container_id is None
):
probe = container_obj.get("id")
if isinstance(probe, str) and probe:
latched_container_id = probe
if (
latched_container_id
and not container_id_emitted
and latched_container_id
!= anthropic_code_exec_container_id
):
yield _emit_tool_event(
{
"type": "container_ready",
"container_id": latched_container_id,
}
)
container_id_emitted = True
stop_reason = event.get("delta", {}).get("stop_reason")
if stop_reason:
if thinking_open:
yield _content_chunk("")
thinking_open = False
# `pause_turn` is in-progress, not terminal:
# the SSE stream still ends with [DONE] via
# message_stop but we skip emitting a
# finish_reason="stop" chunk that would
# truncate the rendered message in the UI.
mapped = _finish_reason_map.get(stop_reason, "stop")
# Streaming refusal: emit a visible notice
# plus an out-of-band _toolEvent so the
# frontend can prune the refused turn.
# The mapped finish_reason is
# "content_filter" per OpenAI spec.
# https://platform.claude.com/docs/en/test-and-evaluate/strengthen-guardrails/handle-streaming-refusals
if stop_reason == "refusal":
logger.warning(
"Anthropic refusal stop_reason (model=%s)",
model,
)
# Drop signal rides _toolEvent (not
# text) so assistant content cannot
# spoof a context reset.
yield _content_chunk(
"\n\n_The response was stopped by "
"Anthropic's safety classifier. Edit "
"or remove the previous turn and try "
"again._"
)
yield _emit_tool_event(
{"type": "anthropic_refusal"}
)
if mapped is not None:
chunk = {
"id": completion_id,
"object": "chat.completion.chunk",
"choices": [
{
"index": 0,
"delta": {},
"finish_reason": mapped,
}
],
}
yield f"data: {_json.dumps(chunk)}"
elif event_type == "message_stop":
if thinking_open:
yield _content_chunk("")
thinking_open = False
# Forward document_citations so the Sources
# panel can render the inline [N] footnotes.
# ``cited_text`` is truncated server-side to
# keep SSE bytes bounded on long spans.
if document_citations:
clean_cits = []
for c in document_citations:
entry = {k: v for k, v in c.items() if k != "_key"}
cited = entry.get("cited_text")
if (
isinstance(cited, str)
and len(cited) > _CITED_TEXT_MAX_LEN
):
entry["cited_text"] = (
cited[:_CITED_TEXT_MAX_LEN] + "…"
)
clean_cits.append(entry)
yield _emit_tool_event(
{
"type": "document_citations",
"citations": clean_cits,
}
)
# Final include_usage-style chunk so callers can
# see cache_creation / cache_read without
# scraping the server log.
usage_line = _build_usage_chunk(
completion_id,
"anthropic",
last_usage,
)
if usage_line:
yield usage_line
yield "data: [DONE]"
await (
response.aclose()
) # set PoolByteStream._closed=True FIRST
break
except GeneratorExit:
await response.aclose() # set PoolByteStream._closed=True FIRST
await lines_gen.aclose() # now safe — aclose() is a no-op
raise
finally:
# Surface per-event-type counts + web_search summary so
# reports of "no reasoning panel content" / "Search
# didn't do anything" can be triaged at a glance.
web_search_requested = bool(
enabled_tools and "web_search" in enabled_tools
)
web_search_invocations = len(web_search_calls)
total_results = sum(
len(sc.get("results") or []) for sc in web_search_calls.values()
)
queries = [
sc["query"]
for sc in web_search_calls.values()
if sc.get("query")
]
# cache_read_input_tokens > 0 on turn N proves the
# cache_control marker on the system block is doing
# its job — turn 1 will show cache_creation > 0
# instead. cache_creation tokens are billed at a
# small premium; cache_read tokens are billed at a
# discount.
code_execution_invocations = len(code_execution_calls)
code_execution_results = sum(
1
for c in code_execution_calls.values()
if c.get("result") is not None
)
web_fetch_requested = web_fetch_enabled
web_fetch_invocations = len(web_fetch_calls)
web_fetch_urls = [
wf["url"] for wf in web_fetch_calls.values() if wf.get("url")
]
logger.info(
"Anthropic stream complete (model=%s, "
"web_search_requested=%s, web_search_invocations=%s, "
"results=%s, queries=%s, "
"web_fetch_requested=%s, web_fetch_invocations=%s, "
"web_fetch_urls=%s, "
"code_execution_requested=%s, "
"code_execution_invocations=%s, "
"code_execution_results=%s, "
"code_execution_generated_files=%s, "
"container_id_in=%s, container_id_out=%s, "
"input_tokens=%s, output_tokens=%s, "
"cache_creation_input_tokens=%s, "
"cache_read_input_tokens=%s, "
"compaction_input_tokens=%s, "
"compaction_output_tokens=%s, "
"compaction_blocks_seen=%s, events=%s)",
model,
web_search_requested,
web_search_invocations,
total_results,
queries,
web_fetch_requested,
web_fetch_invocations,
web_fetch_urls,
code_execution_enabled,
code_execution_invocations,
code_execution_results,
code_execution_generated_files,
anthropic_code_exec_container_id,
latched_container_id,
last_usage.get("input_tokens"),
last_usage.get("output_tokens"),
last_usage.get("cache_creation_input_tokens"),
last_usage.get("cache_read_input_tokens"),
last_usage.get("compaction_input_tokens"),
last_usage.get("compaction_output_tokens"),
compaction_blocks_seen,
event_counts,
)
await response.aclose()
await lines_gen.aclose()
except httpx.ConnectError as exc:
logger.error("Connection error to %s: %s", self.provider_type, exc)
yield _error_sse_line(
502,
f"Failed to connect to {self.provider_type}: {exc}",
self.provider_type,
)
except httpx.ReadTimeout as exc:
logger.error("Read timeout from %s: %s", self.provider_type, exc)
yield _error_sse_line(
504,
f"Timeout waiting for {self.provider_type} response",
self.provider_type,
)
except httpx.HTTPError as exc:
logger.error("HTTP error from %s: %s", self.provider_type, exc)
yield _error_sse_line(
502,
f"Error communicating with {self.provider_type}: {exc}",
self.provider_type,
)
async def _stream_openai_responses(
self,
messages: list[dict[str, Any]],
model: str,
temperature: float,
top_p: float,
max_tokens: Optional[int],
enable_thinking: Optional[bool],
reasoning_effort: Optional[str],
enabled_tools: Optional[list[str]] = None,
enable_prompt_caching: Optional[bool] = None,
openai_code_exec_container_id: Optional[str] = None,
compaction_threshold: Optional[int] = None,
) -> AsyncGenerator[str, None]:
"""
Call OpenAI's /v1/responses endpoint and translate its SSE stream back
into OpenAI Chat Completions chunk format.
The Responses API uses a different request shape (``input`` instead of
``messages``, ``instructions`` for system prompts, ``max_output_tokens``
for the budget) and emits event-typed SSE frames (e.g.
``response.output_text.delta``) rather than chat-completion chunks.
``presence_penalty`` / ``top_k`` are not part of the Responses contract
and are dropped here intentionally.
"""
import json as _json
is_openai_cloud = _is_openai_family_cloud(self.base_url)
image_generation_requested = bool(
enabled_tools and "image_generation" in enabled_tools and is_openai_cloud
)
# Split system messages out into a single `instructions` string and
# translate user/assistant messages into the Responses input shape.
instructions_parts: list[str] = []
input_items: list[dict[str, Any]] = []
openai_replay_items: list[dict[str, Any]] = []
previous_response_id: Optional[str] = None
for msg in messages:
role = msg.get("role")
content = msg.get("content", "")
if role == "system":
if isinstance(content, str):
if content:
instructions_parts.append(content)
elif isinstance(content, list):
for part in content:
if part.get("type") == "text" and part.get("text"):
instructions_parts.append(part["text"])
continue
if isinstance(content, str):
input_items.append({"role": role, "content": content})
continue
if isinstance(content, list):
translated_parts: list[dict[str, Any]] = []
used_previous_response_id = False
for part in content:
part_type = part.get("type")
if part_type == "text":
translated_parts.append(
{"type": "input_text", "text": part.get("text", "")}
)
elif part_type == "image_url":
url = part.get("image_url", {}).get("url", "")
if url:
# Responses takes image_url as a flat string (both
# https:// URLs and data: URLs are accepted).
translated_parts.append(
{"type": "input_image", "image_url": url}
)
elif (
part_type == "reasoning"
and role == "assistant"
and image_generation_requested
):
replay_item = _sanitize_openai_reasoning_replay_item(part)
if replay_item:
openai_replay_items.append(replay_item)
elif (
part_type == "image_generation_call"
and role == "assistant"
and image_generation_requested
):
response_id = (
part.get("response_id")
or part.get("openai_response_id")
or part.get("previous_response_id")
)
call_id = part.get("id") or part.get("image_generation_call_id")
if isinstance(call_id, str) and call_id:
if isinstance(response_id, str) and response_id:
previous_response_id = response_id
input_items = []
translated_parts = []
used_previous_response_id = True
else:
previous_response_id = None
openai_replay_items.append(
{"type": "image_generation_call", "id": call_id}
)
elif part_type == "input_document":
# OpenAI Responses accepts PDFs / docs as
# `{type:"input_file", file_data:"data:application/pdf;base64,..."}`
# or `{type:"input_file", file_url:"https://..."}`,
# with optional `filename`. See
# https://developers.openai.com/api/docs/guides/images-vision
# Map Studio's normalised `input_document` shape
# straight onto Responses' `input_file`.
file_url = part.get("file_url")
file_data = part.get("file_data")
filename = part.get("filename")
# Mirror the Anthropic-side guard: any "data:" URI
# without an actual base64 payload (`data:application/pdf;base64,`
# or whitespace-only) would otherwise be forwarded
# to OpenAI as `file_data=""`, which 400s the whole
# turn. Treat such payloads as missing AND fall
# back to file_url if one is also present, so a
# recoverable remote URL doesn't get discarded in
# favour of a malformed inline payload.
file_data_valid = bool(
isinstance(file_data, str)
and file_data
and (
not file_data.startswith("data:")
or file_data.partition(",")[2].strip()
)
)
block: dict[str, Any] = {"type": "input_file"}
if file_data_valid:
block["file_data"] = file_data
elif file_url:
block["file_url"] = file_url
else:
continue
if filename:
block["filename"] = filename
translated_parts.append(block)
if translated_parts and not used_previous_response_id:
input_items.append({"role": role, "content": translated_parts})
if previous_response_id:
# OpenAI's documented multi-turn image generation path can use
# `previous_response_id` to carry the prior generated image and
# paired reasoning state. Prefer that over manual item replay when
# we captured the response id; keep replay below as a fallback for
# older stored turns that only have an image_generation_call id.
openai_replay_items = []
elif (
_openai_image_replay_requires_reasoning(model)
and reasoning_effort != "none"
and enable_thinking is not False
):
filtered_replay_items: list[dict[str, Any]] = []
has_reasoning_replay = False
dropped_image_replay_without_reasoning = False
for item in openai_replay_items:
if item.get("type") == "reasoning":
has_reasoning_replay = True
filtered_replay_items.append(item)
elif item.get("type") == "image_generation_call":
if has_reasoning_replay:
filtered_replay_items.append(item)
else:
dropped_image_replay_without_reasoning = True
else:
filtered_replay_items.append(item)
openai_replay_items = filtered_replay_items
if dropped_image_replay_without_reasoning:
yield _error_sse_line(
400,
"OpenAI image edit reference is missing paired reasoning state. "
"Regenerate the image, then retry the edit.",
self.provider_type,
)
return
image_generation_has_reference = bool(
previous_response_id
or any(
isinstance(item, dict) and item.get("type") == "image_generation_call"
for item in openai_replay_items
)
)
if openai_replay_items:
insert_at = len(input_items)
for index in range(len(input_items) - 1, -1, -1):
if input_items[index].get("role") == "user":
insert_at = index
break
input_items[insert_at:insert_at] = openai_replay_items
# NOTE: gpt-5.x / o3 / gpt-4.5 are reasoning-class models. They reject
# temperature and top_p with `Unsupported parameter` 400s on
# /v1/responses (and on /v1/chat/completions for the same families).
# The PROVIDER_REGISTRY['openai'] model_id_allowlist already scopes
# the picker to those families, so we never need to send sampling
# knobs here. ``reasoning.effort`` defaults to "medium" server-side
# if omitted — surface it in a future commit if a knob is wanted.
del temperature, top_p # explicit drop — params are accepted for
# API symmetry with the other stream methods but not forwarded.
body: dict[str, Any] = {
"model": model,
"input": input_items,
"stream": True,
}
if previous_response_id:
body["previous_response_id"] = previous_response_id
# `summary: "auto"` is what makes /v1/responses emit reasoning
# summary events — without it OpenAI returns no thinking text on
# most reasoning models, the SSE handler has no …
# to wrap, and the chat reasoning panel stays blank. Always pair
# an explicit effort with summary except for the explicit "off"
# case (effort: "none"), where summaries are pointless.
summary_unsupported = bool(
_OPENAI_REASONING_SUMMARY_UNSUPPORTED.match(model.strip().lower())
)
if reasoning_effort in (
"minimal",
"low",
"medium",
"high",
"max",
"xhigh",
):
body["reasoning"] = {"effort": reasoning_effort}
if not summary_unsupported:
body["reasoning"]["summary"] = "auto"
elif reasoning_effort == "none" or enable_thinking is False:
body["reasoning"] = {"effort": "none"}
elif enable_thinking is True:
body["reasoning"] = {"effort": "medium"}
if not summary_unsupported:
body["reasoning"]["summary"] = "auto"
if instructions_parts:
body["instructions"] = "\n\n".join(instructions_parts)
if max_tokens is not None:
body["max_output_tokens"] = max_tokens
# Prompt caching on /v1/responses is automatic and free, but the
# default in-memory policy only survives ~5-10 min of inactivity
# (up to ~1 hr). Opt into the 24-hour retention policy so a chat
# left idle overnight still hits the cache on the next turn.
# Pricing is identical to in_memory per OpenAI's docs.
#
# Gated on the base URL because ollama / llama.cpp / "custom"
# presets all collapse to provider_type="openai" in
# toExternalBackendProviderType, so they also land in this
# helper. Those servers expose /v1/responses-shaped routes in
# some configurations but don't implement
# prompt_cache_retention — sending the field unconditionally
# would 400 them. Match the public OpenAI host strictly so the
# field only goes to OpenAI cloud. Studio's openai model picker
# is registry-scoped to gpt-5.x / o3 / gpt-4.5, all of which
# accept this parameter (gpt-5.5+ already defaults to "24h" and
# rejects "in_memory", so it's a safe no-op there).
# OpenAI-family cloud: api.openai.com OR Azure OpenAI Foundry
# (*.openai.azure.com). Both expose the same Responses-API
# extensions used below -- prompt_cache_retention,
# context_management compaction, container shell tool -- so
# treat them uniformly. Non-cloud OpenAI-compatible servers
# (ollama / llama.cpp / vLLM / "custom" preset) hit /v1/responses
# without these extensions and would 400 on the unknown body
# fields, so they intentionally fall outside this gate.
if is_openai_cloud and enable_prompt_caching is not False:
body["prompt_cache_retention"] = "24h"
# OpenAI server-side context compaction — see
# https://developers.openai.com/api/docs/guides/compaction
# When `compaction_threshold` is provided on a cloud OpenAI
# request, attach `context_management: [{type:"compaction",
# compact_threshold:N}]` so the API runs server-side
# compaction when the rendered prompt crosses the threshold.
# No beta header is required; no dated version pin. The field
# is silently dropped for non-cloud backends because ollama /
# llama.cpp / "custom" presets land in this helper and would
# 400 on an unknown body field.
if (
is_openai_cloud
and compaction_threshold is not None
and compaction_threshold > 0
):
body["context_management"] = [
{
"type": "compaction",
"compact_threshold": int(compaction_threshold),
}
]
# OpenAI server-side tools — see
# https://developers.openai.com/api/docs/guides/tools
# https://developers.openai.com/api/docs/guides/tools-shell
# The frontend's Search/Code buttons map to the unified
# enabled_tools shorthand; translate that into the Responses-API
# tool schema. Other built-in tools (file_search,
# code_interpreter, image_generation, computer_use_preview) can
# be added with the same pattern when we surface their toggles.
code_execution_enabled_openai = bool(
enabled_tools and "code_execution" in enabled_tools and is_openai_cloud
)
# OpenAI's image_generation tool is a Responses-API server tool.
# See https://developers.openai.com/api/docs/guides/tools-image-generation
# The model picks size / quality / background server-side and
# delegates rendering to a gpt-image-* family model; the result
# comes back inline as an `image_generation_call` output item
# with a base64 image. Available on every gpt-5.x family member
# plus gpt-4.1 / gpt-4o / o3 per the docs; restrict to cloud
# OpenAI because the local llama.cpp / ollama backends don't
# implement it and would 400.
image_generation_enabled_openai = image_generation_requested
def _openai_image_generation_tool() -> dict[str, Any]:
tool: dict[str, Any] = {"type": "image_generation"}
if image_generation_has_reference:
# OpenAI's Responses image tool defaults to `auto`. For
# Studio's explicit follow-up edit flow, force edit mode so
# the provider uses the previous response / call id as image
# context instead of treating the text as a fresh generation.
tool["action"] = "edit"
return tool
if enabled_tools:
tools_array: list[dict[str, Any]] = []
if "web_search" in enabled_tools:
tools_array.append({"type": "web_search"})
if code_execution_enabled_openai:
# `container_auto` lets OpenAI auto-create a fresh
# container per request; we capture the resulting
# container_id off the SSE stream and the chat-adapter
# persists it onto the thread record. Subsequent turns
# in the same thread pass it back as
# `openai_code_exec_container_id`, which we translate to
# `container_reference` here so the model sees
# filesystem state from prior turns. Container expires
# after ~20 min of inactivity per OpenAI's default
# policy — a stale id 400s, the chat-adapter clears it
# via container_invalidated, and the next turn falls
# back to auto-create.
shell_env: dict[str, Any]
if openai_code_exec_container_id:
shell_env = {
"type": "container_reference",
"container_id": openai_code_exec_container_id,
}
else:
shell_env = {"type": "container_auto"}
tools_array.append({"type": "shell", "environment": shell_env})
if image_generation_enabled_openai:
tools_array.append(_openai_image_generation_tool())
if tools_array:
body["tools"] = tools_array
url = f"{self.base_url}/responses"
completion_id = f"chatcmpl-openai-{model.replace('/', '-')}"
logger.info("Proxying OpenAI Responses API to %s (model=%s)", url, model)
def _build_body(container_id_for_this_attempt: Optional[str]) -> dict[str, Any]:
"""Snapshot of the request body. Called once for the initial
attempt and again with ``None`` for the post-expiry retry.
Returns a fresh dict so the retry doesn't share state with the
first attempt.
"""
attempt_body = dict(body)
if enabled_tools:
tools_array_attempt: list[dict[str, Any]] = []
if "web_search" in enabled_tools:
tools_array_attempt.append({"type": "web_search"})
if code_execution_enabled_openai:
if container_id_for_this_attempt:
env_attempt: dict[str, Any] = {
"type": "container_reference",
"container_id": container_id_for_this_attempt,
}
else:
env_attempt = {"type": "container_auto"}
tools_array_attempt.append(
{"type": "shell", "environment": env_attempt}
)
if image_generation_enabled_openai:
tools_array_attempt.append(_openai_image_generation_tool())
if tools_array_attempt:
attempt_body["tools"] = tools_array_attempt
else:
attempt_body.pop("tools", None)
return attempt_body
def _is_openai_container_expired_error(error_text: str) -> bool:
"""Match the substring patterns OpenAI uses for expired / missing
code-exec containers. There's no official error code in the public
docs, so we substring-match a small set.
"""
lowered = error_text.lower()
if "container" not in lowered:
return False
return (
"expired" in lowered
or "not_found" in lowered
or "not found" in lowered
or "no such container" in lowered
)
try:
retried = False
attempt_container_id = openai_code_exec_container_id
while True:
attempt_body = _build_body(attempt_container_id)
async with _http_client.stream(
"POST",
url,
json = attempt_body,
headers = self._auth_headers(),
timeout = self._stream_timeout,
) as response:
if response.status_code != 200:
error_body = await response.aread()
error_text = error_body.decode("utf-8", errors = "replace")
logger.error(
"OpenAI Responses returned %d: %s",
response.status_code,
error_text[:500],
)
expired_container_4xx = (
attempt_container_id
and 400 <= response.status_code < 500
and _is_openai_container_expired_error(error_text)
)
if expired_container_4xx and not retried:
yield (
f"data: "
f"{_json.dumps({'id': completion_id, 'object': 'chat.completion.chunk', 'choices': [{'index': 0, 'delta': {}, 'finish_reason': None}], '_toolEvent': {'type': 'container_invalidated'}})}"
)
retried = True
attempt_container_id = None
continue
yield _error_sse_line(
response.status_code, error_text, self.provider_type
)
return
# NOTE: same manual __anext__ loop as stream_chat_completion —
# see comment there for the GeneratorExit / aclose ordering.
lines_gen = response.aiter_lines().__aiter__()
done_emitted = False
reasoning_open = False
reasoning_emitted = False
# Latched from response.completed / response.incomplete so
# the final log can surface input_tokens_details.cached_tokens —
# the field that proves prompt_cache_retention="24h" is
# actually hitting OpenAI's cache instead of recomputing
# the prefix every turn.
last_usage: Optional[dict[str, Any]] = None
# Per-call state for OpenAI's server-side web_search tool. Mapped
# back into our local _toolEvent shape so the existing chat-UI
# renderer surfaces web_search the same way it does for local
# tool calls: a "Searching…" tool-call card, then a `tool_end`
# carrying citations formatted as
# Title: …\nURL: …\nSnippet: …\n---\n…
# blocks (which the frontend's parseSourcesFromResult lifts
# into source content parts at end of stream).
# web_search_calls preserves insertion order so we can apply
# the aggregated citation list onto the *last* call's
# tool_end — that's the one the frontend's source-pill
# extraction reads (parseSourcesFromResult flatMaps every
# web_search result, so a single non-empty result is enough
# to surface all sources at message tail).
# OpenAI emits url_citation annotations on text deltas, not
# per call — there's no wire field linking a citation back
# to a specific search invocation. Hence the shared list.
# web_search_calls: { item_id -> {query} }
web_search_calls: dict[str, dict[str, Any]] = {}
all_url_citations: list[dict[str, Any]] = []
# Shell-tool (code execution) state. OpenAI emits
# `shell_call` items (model requesting a command list)
# paired with `shell_call_output` items (execution
# results). We mirror the Anthropic code-execution UX
# by emitting one `_toolEvent` tool_start per
# shell_call and one tool_end per shell_call_output;
# they're linked via `shell_call_output.call_id`
# matching `shell_call.id`. Items are independent of
# web_search (different keyed map).
# shell_calls: { call_id -> {commands, output} }
shell_calls: dict[str, dict[str, Any]] = {}
# Container id captured from the response stream. When
# it differs from the inbound id, emit a synthetic
# `container_ready` _toolEvent so the frontend can
# persist it onto the thread record for the next turn.
# Where OpenAI surfaces it is documented loosely; we
# probe two known fields (response.container_id on
# response.completed, item.environment.container_id on
# shell_call output items) and latch the first one we
# see.
latched_container_id: Optional[str] = None
container_id_emitted = False
current_openai_response_id: Optional[str] = None
last_openai_reasoning_replay_item: Optional[dict[str, Any]] = None
openai_reasoning_replay_items: dict[str, dict[str, Any]] = {}
image_generation_calls_started: set[str] = set()
# Buffer for a citation marker straddling two delta events;
# prepended onto the next delta. See _split_pending_citation_tail.
pending_marker_tail: str = ""
# Segments deferred while their markers reference unseen
# source_ids; held in arrival order so output never
# leapfrogs an earlier deferred segment. Flushed on
# annotation events and force-flushed at end-of-stream
# with leftover private-use codepoints stripped.
pending_citation_segments: list[str] = []
def _record_openai_response_id(payload: dict[str, Any]) -> None:
nonlocal current_openai_response_id
response_obj = payload.get("response")
candidates: list[Any] = []
if isinstance(response_obj, dict):
candidates.append(response_obj.get("id"))
candidates.append(payload.get("response_id"))
for candidate in candidates:
if isinstance(candidate, str) and candidate:
current_openai_response_id = candidate
return
def _drain_pending_segments(force: bool) -> str:
"""Re-attempt resolution on buffered segments in order.
Stops at the first still-unresolved segment unless
``force`` (end-of-stream), where lingering markers are stripped."""
out: list[str] = []
while pending_citation_segments:
seg = pending_citation_segments[0]
rewritten, unresolved = _rewrite_citation_markers_partial(
seg,
all_url_citations,
)
if unresolved and not force:
pending_citation_segments[0] = rewritten
break
if unresolved and force:
rewritten = _replace_openai_citation_markers(
rewritten,
all_url_citations,
)
pending_citation_segments.pop(0)
if rewritten:
out.append(rewritten)
return "".join(out)
def _flush_pending_marker_tail(tail: str) -> str:
"""Render any leftover citation tail at end-of-stream.
Unterminated tails drop (no annotation to bind to). If the
close byte arrived concatenated, rewrite then scrub any
residual private-use bytes and any orphan ``cite``
literal so the renderer never sees raw markup. url_citations
are aggregated separately and applied to web_search tool_end.
"""
if not tail:
return ""
if _OPENAI_CITE_STOP not in tail:
# Unterminated: drop the whole tail, otherwise the
# residual ``cite`` would leak as plain text.
return ""
rendered = _replace_openai_citation_markers(
tail, all_url_citations
)
# Scrub residual private-use bytes (e.g. a partial opener).
for ch in ("", "", ""):
rendered = rendered.replace(ch, "")
# Drop any orphan ``cite`` literal -- meaningless
# without its closing byte and matching url_citation.
rendered = re.sub(r"^cite\S*", "", rendered)
return rendered
def _emit_tool_event(payload: dict[str, Any]) -> str:
chunk = {
"id": completion_id,
"object": "chat.completion.chunk",
"choices": [
{
"index": 0,
"delta": {},
"finish_reason": None,
}
],
"_toolEvent": payload,
}
return f"data: {_json.dumps(chunk)}"
def _format_shell_output(output: Any) -> str:
"""Render an OpenAI `shell_call_output.output` list
as the preformatted text payload the frontend's
CodeExecutionToolUI displays inside a . Each
entry has stdout/stderr/outcome — concatenate them
with a separator block per entry and append
`return_code` / `(timeout)` annotations only when
they convey information beyond "succeeded".
"""
if not isinstance(output, list):
return ""
parts: list[str] = []
for entry in output:
if not isinstance(entry, dict):
continue
stdout = entry.get("stdout") or ""
stderr = entry.get("stderr") or ""
outcome = entry.get("outcome") or {}
chunk_parts: list[str] = []
if stdout:
chunk_parts.append(stdout)
if stderr:
chunk_parts.append(f"--- stderr ---\n{stderr}")
if isinstance(outcome, dict):
outcome_type = outcome.get("type")
if outcome_type == "exit":
exit_code = outcome.get("exit_code")
if isinstance(exit_code, int) and exit_code != 0:
chunk_parts.append(f"return_code: {exit_code}")
elif outcome_type == "timeout":
chunk_parts.append("(timeout)")
if chunk_parts:
parts.append("\n".join(chunk_parts))
return (
"\n--- next command ---\n".join(parts)
if parts
else "(no output)"
)
def _record_url_citation(payload: dict[str, Any]) -> None:
"""Append a url_citation onto the shared all_url_citations
list. Dedup by URL — the same URL can be cited many
times under different ``source_id`` aliases (one per
span/locator), so collect every alias we see onto
the matching entry's ``source_ids`` list. The
delta-text rewriter resolves any of those aliases
back to this entry's URL. The id may live under
``source_id``, ``id``, or ``locator`` across the
Responses API revisions."""
if payload.get("type") != "url_citation":
return
url = payload.get("url", "")
if not url:
return
source_id = (
payload.get("source_id")
or payload.get("id")
or payload.get("locator")
or ""
)
# Single pass: either backfill aliases onto an
# existing URL entry (and return) or fall through
# to append a fresh one.
for c in all_url_citations:
if c["url"] != url:
continue
if source_id:
aliases = c.setdefault("source_ids", [])
if source_id not in aliases:
aliases.append(source_id)
return
title = payload.get("title") or url
snippet = payload.get("snippet") or payload.get("quote") or ""
all_url_citations.append(
{
"url": url,
"title": title,
"snippet": snippet,
"source_ids": [source_id] if source_id else [],
}
)
def _record_openai_reasoning_replay_item(
payload: Any,
) -> Optional[dict[str, Any]]:
if not isinstance(payload, dict):
return None
item_id = payload.get("id") or payload.get("item_id")
if not isinstance(item_id, str) or not item_id:
return None
existing = openai_reasoning_replay_items.setdefault(
item_id,
{
"type": "reasoning",
"id": item_id,
"summary": [],
"status": "completed",
},
)
if payload.get("type") == "reasoning":
sanitized = _sanitize_openai_reasoning_replay_item(payload)
if sanitized:
existing.update(sanitized)
return existing
summary_text = ""
part = payload.get("part")
if (
isinstance(part, dict)
and part.get("type") == "summary_text"
):
text = part.get("text")
if isinstance(text, str):
summary_text = text
elif (
payload.get("type")
== "response.reasoning_summary_text.done"
):
text = payload.get("text")
if isinstance(text, str):
summary_text = text
if summary_text:
summary_index = payload.get("summary_index")
summary = existing.setdefault("summary", [])
if isinstance(summary, list):
summary_part = {
"type": "summary_text",
"text": summary_text,
}
if (
isinstance(summary_index, int)
and summary_index >= 0
):
while len(summary) <= summary_index:
summary.append(
{"type": "summary_text", "text": ""}
)
summary[summary_index] = summary_part
else:
summary.append(summary_part)
return existing
def _image_generation_arguments(
prompt: str,
raw_item_id: Any,
) -> dict[str, Any]:
arguments: dict[str, Any] = {"kind": "image", "prompt": prompt}
if isinstance(raw_item_id, str) and raw_item_id:
arguments["openai_image_generation_call_id"] = raw_item_id
if current_openai_response_id:
arguments["openai_response_id"] = current_openai_response_id
if last_openai_reasoning_replay_item:
arguments["openai_reasoning_item"] = (
last_openai_reasoning_replay_item
)
return arguments
def _extract_reasoning_text(payload: Any) -> str:
if payload is None:
return ""
if isinstance(payload, str):
return payload
if isinstance(payload, list):
out: list[str] = []
for item in payload:
text = _extract_reasoning_text(item)
if text:
out.append(text)
return "".join(out)
if isinstance(payload, dict):
# OpenAI responses may carry reasoning summaries in
# different envelope fields across event variants.
for key in ("text", "delta", "content", "summary"):
if key in payload:
text = _extract_reasoning_text(payload.get(key))
if text:
return text
if payload.get("type") == "summary_text":
return _extract_reasoning_text(payload.get("text"))
return ""
def _chunk_with_text(text: str) -> str:
chunk = {
"id": completion_id,
"object": "chat.completion.chunk",
"choices": [
{
"index": 0,
"delta": {"content": text},
"finish_reason": None,
}
],
}
return f"data: {_json.dumps(chunk)}"
try:
while True:
try:
line = await lines_gen.__anext__()
except StopAsyncIteration:
break
if not line or line.startswith("event:"):
continue
if not line.startswith("data:"):
continue
data_str = line[len("data:") :].strip()
if not data_str:
continue
if data_str == "[DONE]":
# Flush any held-over partial marker; strip
# private-use bytes so garbled glyphs don't leak.
if pending_marker_tail:
flushed = _flush_pending_marker_tail(
pending_marker_tail
)
pending_marker_tail = ""
if flushed:
if reasoning_open:
yield _chunk_with_text("")
reasoning_open = False
yield _chunk_with_text(flushed)
# Force-drain any segment still awaiting an
# annotation; lingering codepoints are stripped.
tail_flushed = _drain_pending_segments(
force = True,
)
if tail_flushed:
if reasoning_open:
yield _chunk_with_text("")
reasoning_open = False
yield _chunk_with_text(tail_flushed)
if not done_emitted:
yield "data: [DONE]"
done_emitted = True
break
try:
event = _json.loads(data_str)
except _json.JSONDecodeError:
continue
event_type = event.get("type")
_record_openai_response_id(event)
if event_type == "response.output_text.delta":
delta_text = event.get("delta", "")
# Process inline annotations first so source_ids
# referenced by same-delta markers are in the lookup
# before the rewriter runs. Some API versions inline
# url citations on the delta event itself.
for ann in event.get("annotations") or []:
if isinstance(ann, dict):
_record_url_citation(ann)
if delta_text or pending_marker_tail:
# Prepend any held-over tail so a marker
# straddling two SSE events resolves cleanly.
combined = pending_marker_tail + delta_text
head, pending_marker_tail = (
_split_pending_citation_tail(combined)
)
if head:
if reasoning_open:
yield _chunk_with_text("")
reasoning_open = False
# Re-attempt earlier deferred segments first
# so output stays in order; the needed
# annotation may have arrived inline above.
flushed = _drain_pending_segments(
force = False,
)
if flushed:
yield _chunk_with_text(flushed)
head_rewritten, has_unresolved = (
_rewrite_citation_markers_partial(
head,
all_url_citations,
)
)
if has_unresolved or pending_citation_segments:
pending_citation_segments.append(
head_rewritten
)
elif head_rewritten:
yield _chunk_with_text(head_rewritten)
elif event_type == "response.output_text.annotation.added":
ann = event.get("annotation")
if isinstance(ann, dict):
_record_url_citation(ann)
flushed = _drain_pending_segments(
force = False,
)
if flushed:
if reasoning_open:
yield _chunk_with_text("")
reasoning_open = False
yield _chunk_with_text(flushed)
elif event_type == "response.output_item.added":
item = event.get("item", {})
if (
isinstance(item, dict)
and item.get("type") == "web_search_call"
):
item_id = item.get("id", "") or (
f"ws_{len(web_search_calls)}"
)
web_search_calls.setdefault(item_id, {"query": ""})
# Shell-tool: register the call eagerly so
# the matching shell_call_output can link
# back even if `done` arrives out of order.
# Also probe for container_id on the
# environment field — when container_auto
# auto-creates one, this is the first place
# the new id might surface (OpenAI doesn't
# promise this in docs, but the field is
# cheap to scan and lets us emit
# container_ready earlier than
# response.completed).
if (
isinstance(item, dict)
and item.get("type") == "shell_call"
):
item_id = item.get("id", "") or (
f"sc_{len(shell_calls)}"
)
shell_calls.setdefault(
item_id,
{"commands": [], "output": None},
)
env = item.get("environment")
if isinstance(env, dict):
probe = env.get("container_id") or env.get("id")
if (
isinstance(probe, str)
and probe.startswith("cntr_")
and latched_container_id is None
):
latched_container_id = probe
if (
isinstance(item, dict)
and item.get("type") == "image_generation_call"
):
raw_item_id = item.get("id")
if isinstance(raw_item_id, str) and raw_item_id:
arguments = _image_generation_arguments(
"",
raw_item_id,
)
image_generation_calls_started.add(raw_item_id)
yield _emit_tool_event(
{
"type": "tool_start",
"tool_name": "image_generation",
"tool_call_id": raw_item_id,
"arguments": arguments,
}
)
elif event_type == "response.output_item.done":
item = event.get("item", {})
if not isinstance(item, dict):
continue
if item.get("type") == "reasoning":
last_openai_reasoning_replay_item = (
_record_openai_reasoning_replay_item(item)
)
summary_text = _extract_reasoning_text(
item.get("summary")
)
if summary_text and not reasoning_emitted:
if not reasoning_open:
summary_text = f"{summary_text}"
reasoning_open = True
yield _chunk_with_text(summary_text)
reasoning_emitted = True
elif item.get("type") == "web_search_call":
# done is the canonical place to read the
# query, so emit both tool_start and tool_end
# here. Frontend then renders a card per call
# with the proper "Searching: " label.
# Citations are aggregated separately and the
# *last* call's result is overwritten at
# response.completed with the citation list
# (so the source-pill extraction at message
# tail surfaces them once).
item_id = item.get("id", "") or (
f"ws_{len(web_search_calls)}"
)
action = item.get("action")
query = (
action.get("query", "")
if isinstance(action, dict)
else ""
)
web_search_calls[item_id] = {"query": query}
yield _emit_tool_event(
{
"type": "tool_start",
"tool_name": "web_search",
"tool_call_id": item_id,
"arguments": (
{"query": query} if query else {}
),
}
)
# Per-card text; last call gets overwritten
# with citations at response.completed.
per_call_result = (
f"Searching: {query}" if query else ""
)
yield _emit_tool_event(
{
"type": "tool_end",
"tool_call_id": item_id,
"result": per_call_result,
}
)
elif item.get("type") == "shell_call":
# OpenAI ships the commands array on the
# action field. Join them onto one
# command string for the tool card —
# the renderer is shared with Anthropic
# bash, which only carries a single
# `command`. Multiple commands in one
# shell_call get joined with newlines so
# they still render as one card.
item_id = item.get("id", "") or (
f"sc_{len(shell_calls)}"
)
action = item.get("action") or {}
commands = (
action.get("commands")
if isinstance(action, dict)
else None
) or []
joined_command = (
"\n".join(str(c) for c in commands)
if isinstance(commands, list)
else ""
)
shell_calls.setdefault(
item_id,
{
"commands": [],
"output": None,
"tool_end_emitted": False,
},
)
shell_calls[item_id]["commands"] = (
list(commands)
if isinstance(commands, list)
else []
)
yield _emit_tool_event(
{
"type": "tool_start",
"tool_name": "code_execution",
"tool_call_id": item_id,
"arguments": {
"kind": "bash",
"command": joined_command,
},
}
)
# Fallback: output may be bundled on the
# shell_call done event itself.
embedded_output = item.get("output")
if (
isinstance(embedded_output, list)
and embedded_output
):
shell_calls[item_id]["output"] = embedded_output
shell_calls[item_id]["tool_end_emitted"] = True
yield _emit_tool_event(
{
"type": "tool_end",
"tool_call_id": item_id,
"result": _format_shell_output(
embedded_output
),
}
)
elif item.get("type") == "shell_call_output":
# `call_id` links back to the shell_call's
# `id`, which is what we used as the
# tool_call_id on tool_start. Match on
# call_id when present so the matching
# card transitions to complete.
call_id = (
item.get("call_id") or item.get("id") or ""
)
output = item.get("output") or []
# Skip if bundled-output path already
# finalised this card.
if shell_calls.get(call_id, {}).get(
"tool_end_emitted"
):
continue
if call_id in shell_calls:
shell_calls[call_id]["output"] = output
shell_calls[call_id]["tool_end_emitted"] = True
result_text = _format_shell_output(output)
yield _emit_tool_event(
{
"type": "tool_end",
"tool_call_id": call_id,
"result": result_text,
}
)
elif item.get("type") == "image_generation_call":
# OpenAI's image_generation tool returns
# a single output item with the base64
# PNG/WebP/JPEG on `result` (sometimes
# `b64_json` depending on output_format).
# `revised_prompt` is what the gpt-image
# backbone actually used after refinement
# of the assistant's request. Emit
# tool_start + tool_end so the chat card
# renders the prompt + the generated
# image inline. The frontend chat-adapter
# decides how to render the base64 blob
# (likely an
)
# based on the `kind: "image"` hint we
# set on tool_start arguments.
# `time_ns()` (nanoseconds) instead of
# millisecond resolution so synthesised
# ids stay unique even when two image
# generations resolve in the same ms.
raw_item_id = item.get("id")
item_id = raw_item_id or f"img_{time.time_ns()}"
prompt_in = (
item.get("revised_prompt")
or item.get("prompt")
or ""
)
done_arguments = _image_generation_arguments(
prompt_in,
raw_item_id,
)
if item_id not in image_generation_calls_started:
yield _emit_tool_event(
{
"type": "tool_start",
"tool_name": "image_generation",
"tool_call_id": item_id,
"arguments": done_arguments,
}
)
b64 = (
item.get("result") or item.get("b64_json") or ""
)
output_format = item.get("output_format") or "png"
yield _emit_tool_event(
{
"type": "tool_end",
"tool_call_id": item_id,
"result": "",
"arguments": done_arguments,
"image_b64": b64,
"image_mime": (f"image/{output_format}"),
"size": item.get("size"),
"quality": item.get("quality"),
"background": item.get("background"),
"prompt": prompt_in,
}
)
elif (
isinstance(event_type, str)
and "reasoning" in event_type
):
recorded_reasoning = (
_record_openai_reasoning_replay_item(event)
)
if recorded_reasoning:
last_openai_reasoning_replay_item = (
recorded_reasoning
)
reasoning_delta = _extract_reasoning_text(event)
if reasoning_delta:
if not reasoning_open:
reasoning_delta = f"{reasoning_delta}"
reasoning_open = True
yield _chunk_with_text(reasoning_delta)
reasoning_emitted = True
elif event_type == "response.completed":
completed_usage = (event.get("response") or {}).get(
"usage"
)
if isinstance(completed_usage, dict):
last_usage = completed_usage
# Flush any unterminated citation tail
# held over from the last delta. By
# the time we get here every annotation
# has been recorded so a late-arriving
# source_id may resolve cleanly; if it
# still doesn't, the helper strips the
# private-use bytes so no garbled
# glyph reaches the user.
if pending_marker_tail:
flushed = _flush_pending_marker_tail(
pending_marker_tail
)
pending_marker_tail = ""
if flushed:
if reasoning_open:
yield _chunk_with_text("")
reasoning_open = False
yield _chunk_with_text(flushed)
# Force-drain any segment still awaiting an
# annotation; lingering codepoints are stripped.
tail_flushed = _drain_pending_segments(
force = True,
)
if tail_flushed:
if reasoning_open:
yield _chunk_with_text("")
reasoning_open = False
yield _chunk_with_text(tail_flushed)
if reasoning_open:
yield _chunk_with_text("")
reasoning_open = False
# Probe response.container_id (top-level) and
# response.container.id for the shell-tool
# container id. OpenAI's docs don't pin the
# exact field, so we scan both. Emit
# `container_ready` only when the value
# differs from the inbound one — no churn on
# reuse.
response_obj = event.get("response") or {}
if isinstance(response_obj, dict):
probe_id = response_obj.get("container_id")
if not probe_id:
container_field = response_obj.get("container")
if isinstance(container_field, dict):
probe_id = container_field.get("id")
if (
isinstance(probe_id, str)
and probe_id.startswith("cntr_")
and latched_container_id is None
):
latched_container_id = probe_id
if (
latched_container_id
and not container_id_emitted
and latched_container_id
!= openai_code_exec_container_id
):
yield _emit_tool_event(
{
"type": "container_ready",
"container_id": latched_container_id,
}
)
container_id_emitted = True
# Overwrite the last web_search call with the
# citation list; the source-pill extractor
# flatMaps across cards. Earlier cards keep
# their per-call "Searching:" text.
if web_search_calls and all_url_citations:
last_id = list(web_search_calls.keys())[-1]
blocks: list[str] = []
for cit in all_url_citations:
line = (
f"Title: {cit['title']}\nURL: {cit['url']}"
)
if cit.get("snippet"):
line += f"\nSnippet: {cit['snippet']}"
blocks.append(line)
yield _emit_tool_event(
{
"type": "tool_end",
"tool_call_id": last_id,
"result": "\n---\n".join(blocks),
}
)
# Final flush: finalise any orphan shell_call
# so the card stops spinning.
for sc_id, sc_state in shell_calls.items():
if sc_state.get("tool_end_emitted"):
continue
yield _emit_tool_event(
{
"type": "tool_end",
"tool_call_id": sc_id,
"result": _format_shell_output(
sc_state.get("output") or []
),
}
)
sc_state["tool_end_emitted"] = True
chunk = {
"id": completion_id,
"object": "chat.completion.chunk",
"choices": [
{
"index": 0,
"delta": {},
"finish_reason": "stop",
}
],
}
yield f"data: {_json.dumps(chunk)}"
# Emit include_usage-style chunk after the
# finish_reason so callers can surface
# cached_tokens in their UI.
usage_line = _build_usage_chunk(
completion_id,
"openai",
last_usage,
)
if usage_line:
yield usage_line
elif event_type == "response.incomplete":
incomplete_usage = (event.get("response") or {}).get(
"usage"
)
if isinstance(incomplete_usage, dict):
last_usage = incomplete_usage
# Same flush as response.completed --
# truncated streams can leave a half-
# marker in the buffer.
if pending_marker_tail:
flushed = _flush_pending_marker_tail(
pending_marker_tail
)
pending_marker_tail = ""
if flushed:
if reasoning_open:
yield _chunk_with_text("")
reasoning_open = False
yield _chunk_with_text(flushed)
# Force-drain any segment still awaiting an
# annotation; lingering codepoints are stripped.
tail_flushed = _drain_pending_segments(
force = True,
)
if tail_flushed:
if reasoning_open:
yield _chunk_with_text("")
reasoning_open = False
yield _chunk_with_text(tail_flushed)
if reasoning_open:
yield _chunk_with_text("")
reasoning_open = False
# Same backfill as response.completed — apply
# whatever citations we managed to gather
# before truncation onto the last call. All
# earlier tool cards already have their proper
# query + empty placeholder result from the
# output_item.done emissions above.
if web_search_calls and all_url_citations:
last_id = list(web_search_calls.keys())[-1]
blocks = []
for cit in all_url_citations:
line = (
f"Title: {cit['title']}\nURL: {cit['url']}"
)
if cit.get("snippet"):
line += f"\nSnippet: {cit['snippet']}"
blocks.append(line)
yield _emit_tool_event(
{
"type": "tool_end",
"tool_call_id": last_id,
"result": "\n---\n".join(blocks),
}
)
# Mirror the response.completed flush so
# truncated streams also finalise orphan
# shell_calls.
for sc_id, sc_state in shell_calls.items():
if sc_state.get("tool_end_emitted"):
continue
yield _emit_tool_event(
{
"type": "tool_end",
"tool_call_id": sc_id,
"result": _format_shell_output(
sc_state.get("output") or []
),
}
)
sc_state["tool_end_emitted"] = True
chunk = {
"id": completion_id,
"object": "chat.completion.chunk",
"choices": [
{
"index": 0,
"delta": {},
"finish_reason": "length",
}
],
}
yield f"data: {_json.dumps(chunk)}"
# Emit include_usage-style chunk after the
# length-truncated finish_reason too, so
# incomplete responses still report
# cached_tokens.
usage_line = _build_usage_chunk(
completion_id,
"openai",
last_usage,
)
if usage_line:
yield usage_line
elif event_type in ("response.failed", "error"):
# Surface the failure to the client; let the
# outer route emit [DONE] as part of its cleanup.
error_payload = event.get("response", {}).get(
"error", {}
) or {
"message": event.get("message", "Unknown error"),
"code": event.get("code"),
}
yield _error_sse_line(
502,
_json.dumps(error_payload),
self.provider_type,
)
break
except GeneratorExit:
await response.aclose()
await lines_gen.aclose()
raise
finally:
# Summarise what the model actually did this turn so
# support reports of "I clicked Search and got nothing"
# can be triaged at a glance: was the tool requested,
# did OpenAI invoke it, and how many sources came back?
web_search_requested = bool(
enabled_tools and "web_search" in enabled_tools
)
web_search_invocations = len(web_search_calls)
total_citations = len(all_url_citations)
queries = [
sc["query"]
for sc in web_search_calls.values()
if sc.get("query")
]
# cached_input_tokens > 0 on turn N proves
# prompt_cache_retention="24h" is letting the previous
# turn's prefix hit the cache instead of being
# recomputed. On /v1/responses the field is nested as
# usage.input_tokens_details.cached_tokens (not
# prompt_tokens_details, which is the /v1/chat/completions
# shape).
cached_input_tokens = None
if isinstance(last_usage, dict):
details = last_usage.get("input_tokens_details")
if isinstance(details, dict):
cached_input_tokens = details.get("cached_tokens")
code_execution_requested = code_execution_enabled_openai
code_execution_invocations = len(shell_calls)
code_execution_results = sum(
1
for sc in shell_calls.values()
if sc.get("output") is not None
)
logger.info(
"OpenAI Responses stream complete (model=%s, "
"web_search_requested=%s, web_search_invocations=%s, "
"citations=%s, queries=%s, reasoning_emitted=%s, "
"code_execution_requested=%s, "
"code_execution_invocations=%s, "
"code_execution_results=%s, "
"container_id_in=%s, container_id_out=%s, "
"input_tokens=%s, output_tokens=%s, "
"cached_input_tokens=%s)",
model,
web_search_requested,
web_search_invocations,
total_citations,
queries,
reasoning_emitted,
code_execution_requested,
code_execution_invocations,
code_execution_results,
openai_code_exec_container_id,
latched_container_id,
(last_usage or {}).get("input_tokens"),
(last_usage or {}).get("output_tokens"),
cached_input_tokens,
)
await response.aclose()
await lines_gen.aclose()
return
except httpx.ConnectError as exc:
logger.error("Connection error to %s: %s", self.provider_type, exc)
yield _error_sse_line(
502,
f"Failed to connect to {self.provider_type}: {exc}",
self.provider_type,
)
except httpx.ReadTimeout as exc:
logger.error("Read timeout from %s: %s", self.provider_type, exc)
yield _error_sse_line(
504,
f"Timeout waiting for {self.provider_type} response",
self.provider_type,
)
except httpx.HTTPError as exc:
logger.error("HTTP error from %s: %s", self.provider_type, exc)
yield _error_sse_line(
502,
f"Error communicating with {self.provider_type}: {exc}",
self.provider_type,
)
async def chat_completion(
self,
messages: list[dict[str, Any]],
model: str,
temperature: float = 0.7,
top_p: float = 0.95,
max_tokens: Optional[int] = None,
presence_penalty: float = 0.0,
) -> dict[str, Any]:
"""Non-streaming chat completion. Returns the full response dict.
Note: only valid for OpenAI-compatible providers. Anthropic requires its
own Messages API; use stream_chat_completion (with stream=False) instead
if a non-streaming Anthropic path is needed in the future.
"""
body: dict[str, Any] = {
"model": model,
"messages": messages,
"stream": False,
"temperature": temperature,
"top_p": top_p,
"presence_penalty": presence_penalty,
}
if max_tokens is not None:
if self.provider_type == "openai":
body["max_completion_tokens"] = max_tokens
else:
body["max_tokens"] = max_tokens
response = await _http_client.post(
f"{self.base_url}/chat/completions",
json = body,
headers = self._auth_headers(),
timeout = self._timeout,
)
response.raise_for_status()
return response.json()
async def list_models(self) -> list[dict[str, Any]]:
"""
Call GET /models on the provider to discover available models.
Returns a list of model dicts with at least 'id' and optionally
'created', 'owned_by', etc.
All supported providers expose a /models endpoint:
- OpenAI-compatible: standard {"data": [...]} response
- Anthropic: https://api.anthropic.com/v1/models — same {"data": [...]} shape
"""
try:
response = await _http_client.get(
f"{self.base_url}/models",
headers = self._auth_headers(),
timeout = self._timeout,
)
response.raise_for_status()
data = response.json()
# OpenAI format: {"data": [{"id": "...", ...}, ...]}
# Some local servers (Ollama with no models) return data: null.
models: list[dict[str, Any]] = []
if isinstance(data, dict):
raw_models = data.get("data") or []
if isinstance(raw_models, list):
models = [model for model in raw_models if isinstance(model, dict)]
if not models and self.provider_type == "ollama":
models = await self._list_ollama_native_models()
return models
except httpx.HTTPError as exc:
logger.error("Failed to list models from %s: %s", self.provider_type, exc)
raise
async def _list_ollama_native_models(self) -> list[dict[str, Any]]:
"""Fallback when Ollama's /v1/models returns an empty or null catalog."""
root = self.base_url.removesuffix("/v1").rstrip("/")
response = await _http_client.get(
f"{root}/api/tags",
headers = self._auth_headers(),
timeout = self._timeout,
)
response.raise_for_status()
payload = response.json()
if not isinstance(payload, dict):
return []
raw_models = payload.get("models") or []
if not isinstance(raw_models, list):
return []
return [
{"id": entry.get("name", "").strip(), "owned_by": "ollama"}
for entry in raw_models
if isinstance(entry, dict) and entry.get("name", "").strip()
]
async def verify_models_endpoint_lightweight(self) -> None:
"""
Confirm GET /models returns 200 without buffering the full response body.
Used for providers with enormous catalogs (e.g. OpenRouter, Hugging Face router)
where downloading the full JSON would be prohibitive.
"""
url = f"{self.base_url}/models"
try:
async with _http_client.stream(
"GET",
url,
headers = self._auth_headers(),
timeout = self._timeout,
) as response:
if response.status_code != 200:
response.raise_for_status()
async for _chunk in response.aiter_bytes(chunk_size = 2048):
break
except httpx.HTTPError as exc:
logger.error(
"Lightweight /models check failed for %s: %s",
self.provider_type,
exc,
)
raise
def _container_headers(self) -> dict[str, str]:
"""Auth headers plus the OpenAI-Beta opt-in for /v1/containers.
OpenAI's containers API requires ``OpenAI-Beta: containers=v1``.
Without it, DELETE silently no-ops: the API returns 200 with a
``{"deleted": true}`` body but does not actually remove the
container (verified 2026-05-15). The header is required for
list / create / delete to behave consistently.
"""
headers = self._auth_headers()
headers["OpenAI-Beta"] = "containers=v1"
return headers
async def list_openai_containers(self) -> list[dict[str, Any]]:
"""
GET /v1/containers on the user's OpenAI account.
Returns the raw container records (id, name, created_at,
last_active_at, expires_after, status). The route layer
reshapes these into the UI summary shape.
Only valid against api.openai.com — non-cloud OpenAI-compat
servers don't implement /v1/containers and would 404 here.
Caller is responsible for the is_openai_cloud guard.
"""
response = await _http_client.get(
f"{self.base_url}/containers",
headers = self._container_headers(),
timeout = self._timeout,
)
response.raise_for_status()
data = response.json()
containers = data.get("data") if isinstance(data, dict) else None
result = list(containers) if isinstance(containers, list) else []
logger.info(
"openai_container_list.response count=%s items=%s",
len(result),
[
{"id": c.get("id"), "status": c.get("status")}
for c in result
if isinstance(c, dict)
],
)
return result
async def create_openai_container(
self,
name: str,
ttl_minutes: int,
) -> dict[str, Any]:
"""
POST /v1/containers with ``expires_after.anchor="last_active_at"``.
``ttl_minutes`` is the idle timeout — every API call that
touches the container resets the timer.
"""
body = {
"name": name,
"expires_after": {
"anchor": "last_active_at",
"minutes": ttl_minutes,
},
}
response = await _http_client.post(
f"{self.base_url}/containers",
json = body,
headers = self._container_headers(),
timeout = self._timeout,
)
response.raise_for_status()
return response.json()
async def delete_openai_container(self, container_id: str) -> None:
"""DELETE /v1/containers/{id}. 404s are surfaced as HTTPError.
Uses a fresh httpx client (not the shared ``_http_client``) so
connection-pool state from earlier chat requests cannot
interfere — observed in the wild that DELETEs over the shared
pool returned ``deleted: true`` while the container persisted
in subsequent /containers list calls, even though the same
DELETE issued from a fresh client genuinely removed it.
Verifies the response body reports ``deleted: true``. OpenAI
returns a 2xx ``deleted: true`` body even when the request is
silently rejected (e.g. missing OpenAI-Beta header), so a
status-only check is not sufficient.
"""
url = f"{self.base_url}/containers/{container_id}"
headers = self._container_headers()
logger.info(
"openai_container_delete.outbound url=%s has_auth=%s openai_beta=%s",
url,
"Authorization" in headers,
headers.get("OpenAI-Beta"),
)
async with httpx.AsyncClient(timeout = self._timeout) as fresh_client:
response = await fresh_client.delete(url, headers = headers)
logger.info(
"openai_container_delete.response status=%s cf_ray=%s "
"request_id=%s organization=%s project=%s processing_ms=%s body=%s",
response.status_code,
response.headers.get("cf-ray"),
response.headers.get("x-request-id"),
response.headers.get("openai-organization"),
response.headers.get("openai-project"),
response.headers.get("openai-processing-ms"),
response.text[:300],
)
response.raise_for_status()
try:
payload = response.json()
except ValueError:
payload = None
if not (isinstance(payload, dict) and payload.get("deleted") is True):
raise httpx.HTTPError(
f"OpenAI did not confirm container deletion: {response.text[:200]}"
)
async def close(self) -> None:
"""No-op — the underlying client is shared across requests."""
def _provider_display_name(provider_type: str) -> str:
from core.inference.providers import get_provider_info
info = get_provider_info(provider_type) or {}
return str(info.get("display_name") or provider_type)
def _friendly_provider_error_text(
provider_type: str,
status_code: int,
raw_message: str,
*,
model: str | None = None,
) -> str:
"""Rewrite common provider errors into actionable Studio copy."""
if status_code == 404 and model:
lowered = raw_message.lower()
if "not found" in lowered or "not_found" in lowered:
if provider_type == "ollama":
label = _provider_display_name(provider_type)
return (
f"Model '{model}' is not installed in {label}. "
f"Run `ollama pull {model}` in a terminal, then retry."
)
if provider_type in ("vllm", "llama_cpp"):
label = _provider_display_name(provider_type)
return (
f"Model '{model}' is not available on the {label} server. "
"Check that the server is running and the model is loaded, "
"then retry."
)
return raw_message
def _error_sse_line(status_code: int, message: str, provider_type: str) -> str:
"""Format an error as an SSE data line in OpenAI error format."""
import json
error_obj = {
"error": {
"message": message,
"type": "provider_error",
"code": str(status_code),
"provider": provider_type,
}
}
return f"data: {json.dumps(error_obj)}"
def _build_usage_chunk(
completion_id: str,
provider: Literal["anthropic", "openai"],
last_usage: Optional[dict],
) -> Optional[str]:
"""Build an OpenAI ``include_usage``-style SSE chunk that carries the
upstream prompt-cache accounting back to the client.
Until now Studio captured ``cache_creation_input_tokens`` /
``cache_read_input_tokens`` (Anthropic) and
``input_tokens_details.cached_tokens`` (OpenAI Responses) on
``last_usage`` and only wrote them to the structlog stream.
Browser / SDK clients had no way to see how many tokens hit the cache
-- so the "you saved $X" UX in the chat panel was impossible without
scraping the server log.
This helper emits the standard OpenAI chunk shape -- ``choices: []``
with a populated ``usage`` block -- so any client that already
consumes ``stream_options={"include_usage": true}`` keeps working,
and the Anthropic-native counts are surfaced as extra keys on the
same ``usage`` dict:
usage.prompt_tokens_details.cached_tokens
normalised cache-read count, present for both providers.
usage.cache_creation_input_tokens
Anthropic-only; tokens billed at the cache-write premium.
usage.cache_read_input_tokens
Anthropic-only; same value as cached_tokens, kept for
callers that already key off the native Anthropic name.
Anthropic's ``input_tokens`` excludes the cache buckets -- the
real prompt size is ``input_tokens + cache_creation_input_tokens
+ cache_read_input_tokens``. Emitting ``input_tokens`` alone as
``prompt_tokens`` undercounts cache-heavy turns and breaks
downstream context / cost displays, so we add all three input
buckets together. OpenAI Responses already folds cached tokens
into ``input_tokens`` so no extra arithmetic is needed there.
Returns ``None`` when there are no usage numbers to report (e.g. an
upstream error before ``message_start`` / ``response.completed``).
"""
if not isinstance(last_usage, dict):
return None
completion_tokens = last_usage.get("output_tokens") or 0
if provider == "anthropic":
uncached_input = last_usage.get("input_tokens") or 0
cache_creation = last_usage.get("cache_creation_input_tokens") or 0
cache_read = last_usage.get("cache_read_input_tokens") or 0
prompt_tokens = uncached_input + cache_creation + cache_read
if not (prompt_tokens or completion_tokens):
return None
usage_block: dict[str, Any] = {
"prompt_tokens": prompt_tokens,
"completion_tokens": completion_tokens,
"total_tokens": prompt_tokens + completion_tokens,
"prompt_tokens_details": {"cached_tokens": cache_read},
"cache_creation_input_tokens": cache_creation,
"cache_read_input_tokens": cache_read,
}
# Forward 5m/1h cache-write breakdown so cost calc applies the
# 2x 1h premium instead of defaulting to 5m on chat-style.
cc_breakdown = last_usage.get("cache_creation")
if isinstance(cc_breakdown, dict) and cc_breakdown:
usage_block["cache_creation"] = cc_breakdown
# Propagate fast-mode `usage.speed` so the cost ledger can apply
# the 6x multiplier without re-derivation (Anthropic falls back
# to "standard" when fast-mode is unsupported or rate-limited).
speed = last_usage.get("speed")
if speed in ("fast", "standard"):
usage_block["speed"] = speed
else:
prompt_tokens = last_usage.get("input_tokens") or 0
cached = 0
details = last_usage.get("input_tokens_details")
if isinstance(details, dict):
cached = details.get("cached_tokens") or 0
if not (prompt_tokens or completion_tokens or cached):
return None
usage_block = {
"prompt_tokens": prompt_tokens,
"completion_tokens": completion_tokens,
"total_tokens": prompt_tokens + completion_tokens,
"prompt_tokens_details": {"cached_tokens": cached},
}
chunk = {
"id": completion_id,
"object": "chat.completion.chunk",
"choices": [],
"usage": usage_block,
}
return f"data: {_json.dumps(chunk)}"