Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,9 @@ jobs:
# a product worded from it vanishes on adoption. See the script's docstring.
- run: python3 scripts/check_unconsumed_fields.py
- run: uv run pytest -q
# Keep current-model request/response handling usable at our SDK floor,
# including SDK 0.x (httpx) versus 1.x (httpx2) parsing and transport.
- run: uv run --isolated --with 'anthropic[bedrock]==0.69.0' --extra dev pytest -q tests/test_anthropic_latest_models.py tests/test_provider_native_clients.py tests/test_provider_protocol_reasoning.py
# The gate pins tiktoken's own cache key and content hash, and only a real
# tiktoken can falsify them. Scope that optional dependency to this contract
# test; the isolated run uses the tokenizer version pinned in uv.lock.
Expand Down
18 changes: 18 additions & 0 deletions agent_core/components/observers/trajectory.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,9 @@

_FORMATS: tuple[str, ...] = ("json", "jsonl")
_DEFAULT_FORMATS: tuple[str, ...] = _FORMATS
# Normal completion markers (Anthropic / OpenAI-compatible). They carry no
# information and appear on almost every turn, so trajectories omit them.
_ORDINARY_FINISH_REASONS = frozenset({"end_turn", "stop"})
_STREAM_ENCODER = json.JSONEncoder(ensure_ascii=False, separators=(",", ":"))
_DEFAULT_FORMAT_ENV_VARS = (
"AGENT_CORE_TRAJECTORY_FORMATS",
Expand Down Expand Up @@ -504,6 +507,17 @@ async def on_llm_response(
# encrypted_content / text) so the (sub-)agent trajectory stays
# replay-able with signatures / encrypted reasoning intact.
record["thinking_blocks"] = ctx.thinking_blocks
# Why the turn ended. Without it a trajectory cannot distinguish a
# finished answer from one the output cap truncated, or from a request
# a safety classifier declined — all three look like a turn whose text
# stops, and a refusal looks like a turn that said nothing at all.
# ``end_turn`` / ``stop`` are the uninformative common case and are
# omitted to keep the line small; everything else (``length``, ``tool_use``,
# ``refusal``, …) is recorded.
if ctx.finish_reason and ctx.finish_reason not in _ORDINARY_FINISH_REASONS:
record["finish_reason"] = ctx.finish_reason
if ctx.stop_details:
record["stop_details"] = ctx.stop_details
self._write_jsonl(record)

# JSON envelope: emit prior context once, then this turn.
Expand Down Expand Up @@ -531,6 +545,10 @@ async def on_llm_response(
# Verbatim thinking / reasoning blocks (signatures /
# encrypted_content) so the JSON envelope stays replay-able.
msg["thinking_blocks"] = ctx.thinking_blocks
if ctx.finish_reason and ctx.finish_reason not in _ORDINARY_FINISH_REASONS:
msg["finish_reason"] = ctx.finish_reason
if ctx.stop_details:
msg["stop_details"] = ctx.stop_details
if ctx.tool_calls:
ids: list[str] = []
tcs: list[dict] = []
Expand Down
14 changes: 14 additions & 0 deletions agent_core/llm.py
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,20 @@ class StreamDelta:
reasoning_blocks: list[dict[str, Any]] = field(
default_factory=list[dict[str, Any]],
)
# Structured refusal detail, carried only by a provider that reports one
# (Anthropic populates ``stop_details`` solely when a safety classifier
# declined the request). A refusal arrives as a successful stream that
# simply produces no text and no tool call, so without this channel the
# streamed turn is indistinguishable from the model choosing to say
# nothing — the same blind spot the non-streaming path closes by putting
# it in ``response_metadata``. The assembler folds it into exactly that,
# so both paths expose a refusal identically.
stop_details: dict[str, Any] = field(default_factory=dict[str, Any])
# Provider's original marker, before finish_reason normalization.
stop_reason: str = ""
# A successful provider recovery discarded historical signed reasoning.
# The loop must reset it before appending the newly generated blocks.
thinking_history_reset: bool = False


@runtime_checkable
Expand Down
16 changes: 16 additions & 0 deletions agent_core/loop_types.py
Original file line number Diff line number Diff line change
Expand Up @@ -222,6 +222,22 @@ class TurnContext:
# turn — the condition under which a leaked call can appear in
# ``blocked_tool_calls``.
tool_schemas_stripped: bool = False
# Why the provider stopped generating, normalised (``length`` for every
# transport's output-cap marker; other values pass through — ``end_turn``,
# ``tool_use``, ``refusal``, …).
#
# An observer could not previously tell a finished answer from a truncated
# one or a declined one: all three arrive as a turn whose text simply ends,
# and a refusal arrives with no text at all. Trajectory writers therefore
# recorded runs whose failure mode could not be reconstructed afterwards.
# Defaulted so every existing TurnContext construction stays valid; a
# producer that does not set it leaves consumers exactly as blind as before,
# never wrong.
finish_reason: str = ""
# Provider-reported detail for a declined request (Anthropic's
# ``stop_details``: ``type`` / ``category`` / ``explanation``). Empty for
# every other stop reason, and for providers that report nothing.
stop_details: dict[str, Any] = field(default_factory=dict[str, Any])


@dataclass
Expand Down
180 changes: 164 additions & 16 deletions agent_core/providers/anthropic.py
Original file line number Diff line number Diff line change
Expand Up @@ -119,11 +119,11 @@ def _build_kwargs(
if system:
kwargs["system"] = system
if self._thinking:
# ``effort`` rides on ``extra_body.output_config`` so any value (incl. ``xhigh``)
# reaches ``messages.create`` without the SDK's stricter validation.
kwargs["thinking"] = self._thinking
if self._effort:
kwargs["extra_body"] = {"output_config": {"effort": self._effort}}
# Current Claude models think adaptively even when thinking is omitted.
# Keep effort independent of that optional display/configuration field.
if self._effort:
kwargs["extra_body"] = {"output_config": {"effort": self._effort}}
if tools:
kwargs["tools"] = [_to_anthropic_tool(t) for t in tools]
if extra_headers:
Expand All @@ -138,6 +138,37 @@ def _build_kwargs(
kwargs["messages"] = _fold_transient_tail(kwargs["messages"], transient_tail)
return kwargs

async def _create_message(self, kwargs: dict[str, Any]) -> tuple[Any, bool]:
"""Retry once when the API rejects a historical thinking signature.

Fable 5.1 / Opus 5.5 bind thinking to its conversation prefix. Runtime
compaction, tool filtering, or a changed system prompt can invalidate
that prefix. Remove all thinking for this request only, leaving durable
history intact; unrelated 400s and a second rejection propagate.

Matching is on "signature" + "thinking" rather than one exact wording:
the phrasing is not a documented contract, and omitting historical
thinking is always a valid request, so any signature rejection is
recoverable the same way. A narrower match fails silently on rewording.
"""
from anthropic import BadRequestError

try:
return await self._client.messages.create(**kwargs), False
except BadRequestError as exc:
error = str(exc).lower()
if not ("signature" in error and "thinking" in error):
raise
messages, stripped = _without_thinking_blocks(kwargs["messages"])
if not stripped:
raise
logger.warning(
"Anthropic thinking signature no longer matches the conversation; "
"retrying once without historical thinking blocks",
)
raw = await self._client.messages.create(**{**kwargs, "messages": messages})
return raw, True

async def chat(
self,
messages: list[Message],
Expand All @@ -152,8 +183,11 @@ async def chat(
messages, tools=tools, temperature=temperature,
max_tokens=max_tokens, extra_headers=extra_headers, timeout=timeout,
)
raw = await self._client.messages.create(**kwargs)
return _to_llm_response(raw)
raw, reset = await self._create_message(kwargs)
response = _to_llm_response(raw)
if reset:
response.response_metadata["thinking_history_reset"] = True
return response

async def stream(
self,
Expand Down Expand Up @@ -182,6 +216,7 @@ async def stream(
reasoning_tokens: int | None = None
model = ""
stop_reason = ""
stop_details: dict[str, Any] = {}
# Verbatim block list, in the provider's own emission order, rebuilt
# from the event stream so the streamed turn replays exactly like the
# non-streaming one (``_to_llm_response``). Keyed by the stream's block
Expand All @@ -195,7 +230,7 @@ async def stream(
# string, so a streamed thinking turn used to yield reasoning that
# ``thinking_format="content_block"`` could not replay.
blocks: dict[int, dict[str, Any]] = {}
stream = await self._client.messages.create(**kwargs)
stream, reset = await self._create_message(kwargs)
async for event in stream:
etype = getattr(event, "type", "")
if etype == "message_start":
Expand All @@ -217,10 +252,13 @@ async def stream(
if cbtype == "tool_use":
# Open a tool-call slot: id + name set once; arguments
# arrive as ``input_json_delta`` partial-JSON fragments.
# Deliberately NOT recorded in ``blocks``: tool calls ride
# the ``tool_call_deltas`` channel and are re-emitted from
# ``Message.tool_calls`` on replay, exactly as the
# non-streaming ``_to_llm_response`` does.
# Also keep its position among signed thinking blocks.
# Reordering tool calls changes the prefix of later thinking.
blocks[idx] = {
"type": "tool_use", "id": getattr(cb, "id", "") or "",
"name": getattr(cb, "name", "") or "",
"input": getattr(cb, "input", {}) or {},
}
yield StreamDelta(tool_call_deltas=[{
"index": idx,
"id": getattr(cb, "id", "") or "",
Expand Down Expand Up @@ -269,6 +307,12 @@ async def stream(
if blk is not None and blk.get("type") == "thinking":
blk["signature"] += getattr(d, "signature", "") or ""
elif dtype == "input_json_delta":
blk = blocks.get(idx)
if blk is not None and blk.get("type") == "tool_use":
blk["_partial_json"] = (
blk.get("_partial_json", "")
+ (getattr(d, "partial_json", "") or "")
)
yield StreamDelta(tool_call_deltas=[{
"index": idx,
"id": None,
Expand All @@ -278,6 +322,12 @@ async def stream(
elif etype == "message_delta":
d = getattr(event, "delta", None)
stop_reason = getattr(d, "stop_reason", "") or stop_reason
# A classifier refusal ends the stream here, with the detail on
# the same delta that carries the stop reason. Captured so the
# streamed turn reports it exactly like the non-streaming one.
details = _anthropic_stop_details(d)
if details is not None:
stop_details = details
u = getattr(event, "usage", None)
if u is not None:
ot = getattr(u, "output_tokens", None)
Expand All @@ -297,6 +347,15 @@ async def stream(
# text turn the flattened ``content`` string is the faithful shape and
# the assembler should keep using it.
ordered = _ordered_blocks(blocks)
for block in ordered:
partial = block.pop("_partial_json", None)
if partial is not None:
try:
block["input"] = json.loads(partial)
except (ValueError, TypeError):
# The tool-call channel retains truncated arguments for
# the runtime's repair/replay logic.
block["input"] = {}
yield StreamDelta(
usage=_anthropic_usage_dict(
input_tokens,
Expand All @@ -310,6 +369,9 @@ async def stream(
finish_reason=normalize_finish_reason(stop_reason),
model=model,
reasoning_blocks=ordered if _has_thinking(ordered) else [],
stop_details=stop_details,
stop_reason=stop_reason,
thinking_history_reset=reset,
)


Expand Down Expand Up @@ -379,6 +441,26 @@ async def _prepare_request(self, request: httpx.Request) -> None:
# ── Conversion helpers ───────────────────────────────────────────────────


def _without_thinking_blocks(
messages: list[dict[str, Any]],
) -> tuple[list[dict[str, Any]], bool]:
out: list[dict[str, Any]] = []
stripped = False
for message in messages:
content = message.get("content")
if not isinstance(content, list):
out.append(message)
continue
kept = [b for b in content if b.get("type") not in ("thinking", "redacted_thinking")]
if len(kept) == len(content):
out.append(message)
continue
stripped = True
if kept:
out.append({**message, "content": kept})
return out, stripped


def _merge_tool_results(
pairs: list[tuple[dict[str, Any], bool]],
) -> list[tuple[dict[str, Any], bool]]:
Expand Down Expand Up @@ -509,6 +591,11 @@ def _to_anthropic_msg(m: Message) -> dict[str, Any] | None:
}
if role == "assistant":
blocks: list[dict[str, Any]] = []
calls = [
converted for tc in (m.get("tool_calls") or [])
if (converted := _to_anthropic_tool_use(tc)) is not None
]
remaining = {call["id"]: call for call in calls}
raw = m.get("content")
if isinstance(raw, list):
# Extended-thinking continuation: history kept the VERBATIM block
Expand Down Expand Up @@ -538,14 +625,17 @@ def _to_anthropic_msg(m: Message) -> dict[str, Any] | None:
txt = block.get("text", "") or ""
if txt:
blocks.append({"type": "text", "text": txt})
elif bt == "tool_use":
# The canonical tool-call channel remains authoritative
# after runtime repair/filtering; native blocks supply order.
call = remaining.pop(block.get("id"), None)
if call is not None:
blocks.append(call)
else:
body = text_of(raw or "")
if body:
blocks.append({"type": "text", "text": body})
for tc in m.get("tool_calls", []) or []:
block = _to_anthropic_tool_use(tc)
if block is not None:
blocks.append(block)
blocks.extend(remaining.values())
if not blocks:
# Nothing to say and nothing to call. An empty ``text`` block is
# NOT a usable placeholder — Anthropic rejects zero-length text
Expand Down Expand Up @@ -683,6 +773,46 @@ def _anthropic_cache_write_tokens(usage: Any) -> int | None:
return max(0, int(raw or 0)) + max(0, int(extension or 0))


def _anthropic_stop_details(raw: Any) -> dict[str, Any] | None:
"""Structured refusal detail off a response, or ``None`` when absent.

Anthropic populates ``stop_details`` ONLY when ``stop_reason ==
"refusal"`` — a safety classifier declined the request. That arrives as a
normal HTTP 200 with empty ``content``, so without this the turn is
indistinguishable from "the model chose to say nothing": the loop sees no
text and no tool call, takes its no-tool exit, and the run ends looking
clean. The 2026-10-05 gdpval triage had to rule a refusal in or out by
hand for exactly this reason.

``category`` is an open set (``cyber``, ``bio``, ``reasoning_extraction``,
``frontier_llm``, ``general_harms``, ``None``, …) and models keep adding
to it, so every field is passed through as-is rather than validated
against a list this module would have to chase. Returns ``None`` when the
payload carries nothing, which keeps the key out of ``response_metadata``
for the overwhelmingly common non-refusal turn.
"""
details = getattr(raw, "stop_details", None)
if details is None and isinstance(raw, dict):
details = raw.get("stop_details")
if details is None:
return None
if isinstance(details, dict):
out = {k: v for k, v in details.items() if v is not None}
return out or None
model_dump = getattr(details, "model_dump", None)
if callable(model_dump):
dumped = model_dump(mode="json", exclude_none=True)
if isinstance(dumped, dict):
return dumped or None
# SDK model object: read the documented fields off it.
out = {}
for field in ("type", "category", "explanation"):
value = getattr(details, field, None)
if value is not None:
out[field] = value
return out or None


def _anthropic_reasoning_tokens(usage: Any) -> int | None:
"""Best-effort extended-thinking token count off an Anthropic usage object.

Expand Down Expand Up @@ -760,6 +890,11 @@ def _to_llm_response(raw: Any) -> LLMResponse:
"data": getattr(block, "data", "") or "",
})
elif btype == "tool_use":
blocks_out.append({
"type": "tool_use", "id": getattr(block, "id", ""),
"name": getattr(block, "name", ""),
"input": getattr(block, "input", {}) or {},
})
tool_calls.append({
"id": getattr(block, "id", ""),
"type": "function",
Expand Down Expand Up @@ -787,14 +922,27 @@ def _to_llm_response(raw: Any) -> LLMResponse:
_anthropic_reasoning_tokens(usage),
)

# ``stop_reason`` is normalised for ``finish_reason`` (``max_tokens`` →
# ``length``), which is what the loop's truncation checks need. The RAW
# value is kept alongside it: ``refusal`` survives normalisation today,
# but a consumer asking "did the provider decline?" should not have to
# know which markers this function rewrites.
metadata: dict[str, Any] = {"id": getattr(raw, "id", "")}
stop_reason_raw = str(getattr(raw, "stop_reason", "") or "")
if stop_reason_raw:
metadata["stop_reason"] = stop_reason_raw
stop_details = _anthropic_stop_details(raw)
if stop_details is not None:
metadata["stop_details"] = stop_details

return LLMResponse(
content=content,
tool_calls=tool_calls,
reasoning_content="\n".join(thinking_parts),
finish_reason=normalize_finish_reason(getattr(raw, "stop_reason", "")),
model=getattr(raw, "model", "") or "",
usage=usage_dict,
response_metadata={"id": getattr(raw, "id", "")},
response_metadata=metadata,
)


Expand Down
Loading
Loading