diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 99149ba..aae2f4c 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -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. diff --git a/agent_core/components/observers/trajectory.py b/agent_core/components/observers/trajectory.py index 8df03d2..7ab42c2 100644 --- a/agent_core/components/observers/trajectory.py +++ b/agent_core/components/observers/trajectory.py @@ -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", @@ -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. @@ -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] = [] diff --git a/agent_core/llm.py b/agent_core/llm.py index e9250d0..c6bfb20 100644 --- a/agent_core/llm.py +++ b/agent_core/llm.py @@ -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 diff --git a/agent_core/loop_types.py b/agent_core/loop_types.py index 33a9abc..c0ced30 100644 --- a/agent_core/loop_types.py +++ b/agent_core/loop_types.py @@ -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 diff --git a/agent_core/providers/anthropic.py b/agent_core/providers/anthropic.py index 4bf49e8..af48724 100644 --- a/agent_core/providers/anthropic.py +++ b/agent_core/providers/anthropic.py @@ -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: @@ -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], @@ -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, @@ -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 @@ -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": @@ -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 "", @@ -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, @@ -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) @@ -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, @@ -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, ) @@ -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]]: @@ -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 @@ -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 @@ -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. @@ -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", @@ -787,6 +922,19 @@ 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, @@ -794,7 +942,7 @@ def _to_llm_response(raw: Any) -> LLMResponse: 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, ) diff --git a/agent_core/runtime/loop/_streaming.py b/agent_core/runtime/loop/_streaming.py index 53a1aaa..8dbf277 100644 --- a/agent_core/runtime/loop/_streaming.py +++ b/agent_core/runtime/loop/_streaming.py @@ -258,6 +258,12 @@ async def _stream_llm_response( # non-streaming path returns and what ``thinking_format="content_block"`` # needs to replay the signed reasoning state on the next turn. final_reasoning_blocks: list[dict[str, Any]] = [] + # Structured refusal detail, when the provider reports one. A declined + # request streams to a clean close with no text and no tool call, so this + # is the only signal separating it from a turn the model chose to end. + final_stop_details: dict[str, Any] = {} + final_stop_reason = "" + thinking_history_reset = False think_splitter = _ThinkTagSplitter() accepts_tool_call_chunks = _accepts_tool_call_arg_chunks(on_delta) reasoning_timeout_s = max(float(reasoning_only_timeout_s or 0), 0.0) @@ -286,6 +292,12 @@ def _assembled_response() -> LLMResponse: response_metadata = ( {"provider_actually_used": final_provider} if final_provider else {} ) + if final_stop_details: + response_metadata["stop_details"] = final_stop_details + if final_stop_reason: + response_metadata["stop_reason"] = final_stop_reason + if thinking_history_reset: + response_metadata["thinking_history_reset"] = True visible_content = accumulated if visible_content: visible_content = ( @@ -387,6 +399,12 @@ async def _close_chunk_stream() -> None: final_provider = delta.provider if getattr(delta, "reasoning_blocks", None): final_reasoning_blocks = delta.reasoning_blocks + if getattr(delta, "stop_details", None): + final_stop_details = delta.stop_details + if getattr(delta, "stop_reason", ""): + final_stop_reason = delta.stop_reason + if getattr(delta, "thinking_history_reset", False): + thinking_history_reset = True # Inline ``...`` tags (Qwen-style) are # split out so ``delta`` carries answer-only text and # ``thinking_delta`` collects both inline + typed diff --git a/agent_core/runtime/loop/agent_loop.py b/agent_core/runtime/loop/agent_loop.py index 295f588..3f4af92 100644 --- a/agent_core/runtime/loop/agent_loop.py +++ b/agent_core/runtime/loop/agent_loop.py @@ -1097,6 +1097,26 @@ async def _notify_context_compacted( ) +def _reset_signed_thinking_history(messages: list[Message]) -> None: + """Commit a provider's successful signature recovery to loop-owned history.""" + retained: list[Message] = [] + for message in messages: + content = message.get("content") + if message.get("role") != "assistant" or not isinstance(content, list): + retained.append(message) + continue + kept = [ + block for block in content + if not isinstance(block, dict) + or block.get("type") not in ("thinking", "redacted_thinking") + ] + if len(kept) == len(content): + retained.append(message) + elif kept or message.get("tool_calls"): + retained.append({**message, "content": kept}) + messages[:] = retained + + async def _process_llm_response( cfg: LoopConfig, obs: list, tc_parser: Any, profile: Any, policy: Any, thinking_parser: Any, normalizer: Any, tool_names: set[str], messages: list[Message], metadata: dict[str, Any], turn: int, @@ -1107,6 +1127,12 @@ async def _process_llm_response( metadata["llm_duration_ms"] = int((llm_call_finished - llm_call_started) * 1000) metadata["llm_ttft_ms"] = int(((first_delta_at or llm_call_finished) - llm_call_started) * 1000) + if (getattr(response, "response_metadata", None) or {}).get("thinking_history_reset"): + # Only after recovery succeeds, and before storing its new signed blocks. + # Otherwise the next call replays the known-invalid signatures, fails + # again, and discards the newly generated valid reasoning too. + _reset_signed_thinking_history(messages) + tr = thinking_parser.extract(response, profile) history_msg = normalizer.to_history(response, tr, policy, profile.thinking_format) messages.append(history_msg) @@ -1261,6 +1287,7 @@ async def _process_llm_response( post_content = getattr(response, "content", None) ai_text = post_content if isinstance(post_content, str) else tr.visible_content + stop_details = rmd.get("stop_details") ctx = TurnContext( turn=turn, max_turns=cfg.max_turns, task_id=cfg.task_id, role_id=cfg.role_id, ai_text=ai_text, thinking=tr.thinking, tool_calls=parsed_calls, messages=messages, @@ -1268,6 +1295,8 @@ async def _process_llm_response( thinking_blocks=tr.raw_content_blocks or [], blocked_tool_calls=blocked_landing_calls, tool_schemas_stripped=bool(strip_tools), + finish_reason=metadata["finish_reason"], + stop_details=dict(stop_details) if isinstance(stop_details, dict) else {}, ) llm_interventions = await notify_observers(obs, "on_llm_response", ctx) diff --git a/changes/58.feature.md b/changes/58.feature.md new file mode 100644 index 0000000..dba3c0e --- /dev/null +++ b/changes/58.feature.md @@ -0,0 +1 @@ +Expose Anthropic refusal details and original stop reasons consistently in ordinary and streamed responses, and record finish reasons and refusal details in JSON and JSONL trajectories. Preserve signed thinking and tool-call ordering for Claude Fable 5.1 and Opus 5.5, forward effort independently of thinking configuration (direct clients that set effort without thinking now send it), and recover once from thinking signatures the API rejects after conversation edits. diff --git a/docs/provider-substrate-boundary.md b/docs/provider-substrate-boundary.md index 105b37f..91d315b 100644 --- a/docs/provider-substrate-boundary.md +++ b/docs/provider-substrate-boundary.md @@ -67,3 +67,43 @@ names it differently, so every client routes its stop signal through Only truncation markers are rewritten. `tool_use`, `end_turn`, and `stop` pass through unchanged because hosts read them directly. A new transport that skips this normalization silently disables truncation recovery for its protocol. + +## Current Claude models and refusal telemetry + +Native profiles for `claude-fable-5-1` and `claude-opus-5-5` use +`thinking_type: adaptive` and default to `thinking_display: summarized`. +These models also think adaptively when `thinking` is omitted; the direct +client forwards `effort` in that case. A direct `AnthropicClient` built with +`effort` but no `thinking` now sends `output_config.effort`; models without +effort support reject it, so leave `effort` empty for them. Manual budgets and disabled thinking +are unsupported by these models. Sampling parameters and forced tool choice +are omitted. Model selection and output-token budgets remain host-owned. + +An empty thinking string can still carry a signature. The adapter preserves +thinking, redacted thinking, and intervening tool calls in provider order; +otherwise later signatures can become invalid when replayed. Canonical tool +calls remain authoritative after runtime filtering or repair. + +These models bind thinking signatures to the preceding system prompt, tools, +and conversation. If a 400 reports an invalid thinking signature (for example +a signature bound to a different conversation after client-side compaction), +the adapter retries once without historical thinking or redacted-thinking blocks. +It logs the recovery and reports `thinking_history_reset` in response metadata +(carried by `StreamDelta.thinking_history_reset` when streaming). After success, +the agent loop removes invalid historical thinking before storing the new +response. Direct client consumers must apply that reset to their own history; +the client does not mutate caller-owned messages. A consumer that ignores it +keeps replaying the stale signatures, so every later request pays a rejected +call plus the retry. Other 400s and a failed retry propagate. Hosts should keep conversation prefixes stable to preserve reasoning. + +See Anthropic's [Fable 5.1 migration guide](https://platform.claude.com/docs/en/models/fable-5-1/migration-guide), +[Opus 5.5 migration guide](https://platform.claude.com/docs/en/models/opus-5-5/migration-guide), +and [preserved-thinking contract](https://platform.claude.com/docs/en/build-with-claude/preserved-thinking). + +Refusals can contain partial text; they are successful HTTP responses, not +transport exceptions. Both ordinary and streamed responses retain the raw +`stop_reason` and structured `stop_details` in `LLMResponse.response_metadata`. +Unknown detail fields/categories pass through, including with older SDKs. +Observers receive `TurnContext.finish_reason` and `TurnContext.stop_details`; +the trajectory observer writes them to JSON snapshots and JSONL events. +This telemetry does not change loop termination or choose a fallback model. diff --git a/tests/test_anthropic_latest_models.py b/tests/test_anthropic_latest_models.py new file mode 100644 index 0000000..b5ac6a2 --- /dev/null +++ b/tests/test_anthropic_latest_models.py @@ -0,0 +1,309 @@ +"""Current Claude request, SDK parsing, streaming, and signed-history contracts. + +Model IDs and response shapes follow the Fable 5.1 / Opus 5.5 migration guides. +MockTransport exercises real SDK serialization and SSE parsing without live API +credentials; it does not claim to evaluate the models themselves. +""" + +from __future__ import annotations + +import json +from contextlib import asynccontextmanager + +import httpx +import pytest +from anthropic import AsyncAnthropic, BadRequestError, DefaultAsyncHttpxClient + +from agent_core.components.observers.trajectory import TrajectoryFileObserver +from agent_core.loop_types import LoopConfig, LoopPolicy +from agent_core.messages import assistant_msg, tool_msg, user_msg +from agent_core.providers.anthropic import AnthropicClient, _to_anthropic_msg +from agent_core.providers.protocol_client import build_protocol_client +from agent_core.runtime.loop._streaming import _stream_llm_response +from agent_core.runtime.loop.agent_loop import run_agent_loop +from agent_core.runtime.loop.model_profile import ( + DefaultThinkingParser, + HistoryPolicy, + ModelProfile, + NativeMessageNormalizer, +) + +MODELS = ["claude-fable-5-1", "claude-opus-5-5"] +if issubclass(DefaultAsyncHttpxClient, httpx.AsyncClient): + sdk_httpx = httpx +else: + import httpx2 as sdk_httpx + + +def _payload(model, *, content=None, stop_reason="end_turn", **extra): + return { + "id": "msg_latest", "type": "message", "role": "assistant", "model": model, + "content": content or [], "stop_reason": stop_reason, "stop_sequence": None, + "usage": {"input_tokens": 10, "output_tokens": 20}, **extra, + } + + +def _sse(payload): + events = [{ + "type": "message_start", + "message": {**payload, "content": [], "stop_reason": None, "stop_details": None}, + }] + for index, block in enumerate(payload["content"]): + opening = dict(block) + deltas = [] + if block["type"] == "thinking": + opening.update(thinking="", signature="") + deltas = [ + {"type": "thinking_delta", "thinking": block["thinking"]}, + {"type": "signature_delta", "signature": block["signature"]}, + ] + elif block["type"] == "text": + opening["text"] = "" + deltas = [{"type": "text_delta", "text": block["text"]}] + elif block["type"] == "tool_use": + opening["input"] = {} + args = json.dumps(block["input"]) + deltas = [ + {"type": "input_json_delta", "partial_json": args[:3]}, + {"type": "input_json_delta", "partial_json": args[3:]}, + ] + events.append({"type": "content_block_start", "index": index, "content_block": opening}) + events.extend({"type": "content_block_delta", "index": index, "delta": d} for d in deltas) + events.append({"type": "content_block_stop", "index": index}) + events.extend([ + {"type": "message_delta", "delta": { + "stop_reason": payload["stop_reason"], "stop_sequence": None, + "stop_details": payload.get("stop_details"), + }, "usage": payload["usage"]}, + {"type": "message_stop"}, + ]) + return "".join(f"event: {e['type']}\ndata: {json.dumps(e)}\n\n" for e in events) + + +@asynccontextmanager +async def _client(model, payload, *, error=None, **kwargs): + requests = [] + + def respond(request): + body = json.loads(request.content) + requests.append(body) + if error and (len(requests) == 1 or error.get("always")): + return sdk_httpx.Response(400, json={ + "type": "error", "error": { + "type": "invalid_request_error", "message": error["message"], + }, + }) + response_payload = payload(body) if callable(payload) else payload + if body.get("stream"): + return sdk_httpx.Response(200, text=_sse(response_payload), headers={"content-type": "text/event-stream"}) + return sdk_httpx.Response(200, json=response_payload) + + adapter = AnthropicClient(model, api_key="test", **kwargs) + await adapter._client.close() + async with AsyncAnthropic( + api_key="test", base_url="https://anthropic.invalid", max_retries=0, + http_client=DefaultAsyncHttpxClient(transport=sdk_httpx.MockTransport(respond)), + ) as sdk: + adapter._client = sdk + yield adapter, requests + + +async def _call(client, streaming, messages): + if not streaming: + return await client.chat(messages, temperature=0.2) + + async def on_delta(*_): + pass + + return await _stream_llm_response(client, messages, 10, on_delta) + + +@pytest.mark.parametrize("model", MODELS) +@pytest.mark.parametrize("streaming", [False, True]) +@pytest.mark.parametrize("stop_reason", ["refusal", "max_tokens", "end_turn"]) +async def test_latest_stop_metadata_and_effort_through_real_sdk(model, streaming, stop_reason): + extra = {} + if stop_reason == "refusal": + extra["stop_details"] = { + "type": "refusal", "category": "future_category", "explanation": None, + "future_field": "kept", + } + payload = _payload(model, stop_reason=stop_reason, **extra) + async with _client(model, payload, effort="xhigh") as (client, requests): + response = await _call(client, streaming, [user_msg("hi")]) + assert response.finish_reason == ("length" if stop_reason == "max_tokens" else stop_reason) + assert response.response_metadata["stop_reason"] == stop_reason + if stop_reason == "refusal": + assert response.response_metadata["stop_details"] == { + "type": "refusal", "category": "future_category", "future_field": "kept", + } + else: + assert "stop_details" not in response.response_metadata + assert requests[0]["model"] == model + assert requests[0]["output_config"] == {"effort": "xhigh"} + assert "thinking" not in requests[0] + assert {"temperature", "top_p", "top_k"}.isdisjoint(requests[0]) + + +@pytest.mark.parametrize("model", MODELS) +@pytest.mark.parametrize("streaming", [False, True]) +async def test_empty_signed_thinking_and_interleaved_tool_calls_roundtrip(model, streaming): + blocks = [ + {"type": "thinking", "thinking": "", "signature": "sig_first"}, + {"type": "tool_use", "id": "t1", "name": "search", "input": {"q": "first"}}, + {"type": "thinking", "thinking": "progress update", "signature": "sig_second"}, + {"type": "tool_use", "id": "t2", "name": "search", "input": {"q": "second"}}, + {"type": "redacted_thinking", "data": "opaque"}, + ] + payload = _payload(model, content=blocks, stop_reason="tool_use") + async with _client(model, payload, thinking={"type": "adaptive", "display": "summarized"}) as (client, requests): + response = await _call(client, streaming, [user_msg("hi")]) + parsed = DefaultThinkingParser().extract(response, ModelProfile( + model_id=model, provider="anthropic", thinking_format="content_block", + )) + history = NativeMessageNormalizer().to_history(response, parsed, HistoryPolicy(), "content_block") + replay = _to_anthropic_msg(history) + assert response.content == blocks + assert replay["content"] == blocks + assert parsed.thinking == "\nprogress update" + assert len(response.tool_calls) == 2 + assert requests[0]["thinking"] == {"type": "adaptive", "display": "summarized"} + # A filtered canonical call must not be reintroduced from the native blocks. + filtered = {**history, "tool_calls": history["tool_calls"][:1]} + assert [b["id"] for b in _to_anthropic_msg(filtered)["content"] if b["type"] == "tool_use"] == ["t1"] + + +@pytest.mark.parametrize("model", MODELS) +@pytest.mark.parametrize("streaming", [False, True]) +async def test_prefix_bound_signature_retry_is_once_and_does_not_edit_history(model, streaming, caplog): + messages = [ + user_msg("hi"), + assistant_msg([ + {"type": "thinking", "thinking": "", "signature": "old_sig"}, + {"type": "redacted_thinking", "data": "old_opaque"}, + {"type": "text", "text": "checking"}, + ], tool_calls=[{"id": "t1", "type": "function", "function": { + "name": "search", "arguments": '{"q":"test"}', + }}]), + tool_msg("found", "t1"), + ] + snapshot = json.dumps(messages) + error = {"message": "Invalid `signature` in `thinking` block. The block is bound to a different conversation."} + async with _client(model, _payload(model), error=error) as (client, requests): + response = await _call(client, streaming, messages) + assert len(requests) == 2 + assert [b["type"] for b in requests[1]["messages"][1]["content"]] == ["text", "tool_use"] + assert requests[1]["messages"][2] == requests[0]["messages"][2] + assert json.dumps(messages) == snapshot + assert response.response_metadata["thinking_history_reset"] is True + assert "retrying once" in caplog.text + + +@pytest.mark.parametrize("streaming", [False, True]) +@pytest.mark.parametrize("message,expected_calls", [ + ("max_tokens is too large", 1), + ("thinking.budget_tokens must be less than max_tokens", 1), + ("Invalid signature in thinking block", 2), + ("Invalid `signature` in `thinking` block. The block is bound to a different conversation.", 2), +]) +async def test_unrelated_or_repeated_bad_requests_propagate(streaming, message, expected_calls): + messages = [user_msg("hi"), assistant_msg([ + {"type": "thinking", "thinking": "", "signature": "old"}, + {"type": "text", "text": "answer"}, + ]), user_msg("continue")] + async with _client(MODELS[0], _payload(MODELS[0]), error={"message": message, "always": True}) as (client, requests): + with pytest.raises(BadRequestError): + await _call(client, streaming, messages) + assert len(requests) == expected_calls + + +@pytest.mark.parametrize("model", MODELS) +async def test_native_profile_defaults_support_current_models(model): + client = build_protocol_client({"protocol": "anthropic", "model": model, "effort": "medium"}, title="test") + try: + kwargs = client._build_kwargs([user_msg("hi")], tools=None, temperature=0.2, max_tokens=None, extra_headers=None, timeout=None) + assert kwargs["thinking"] == {"type": "adaptive", "display": "summarized"} + assert kwargs["extra_body"]["output_config"]["effort"] == "medium" + assert {"temperature", "top_p", "top_k"}.isdisjoint(kwargs) + finally: + await client._client.close() + + +@pytest.mark.parametrize("model", MODELS) +@pytest.mark.parametrize("streaming", [False, True]) +async def test_loop_commits_signature_reset_and_preserves_new_reasoning(model, streaming): + initial = [user_msg("hi"), assistant_msg([ + {"type": "thinking", "thinking": "", "signature": "invalid_old"}, + {"type": "text", "text": "previous answer"}, + ]), user_msg("continue")] + original = json.dumps(initial) + + def response_payload(request): + signatures = [ + block.get("signature") + for message in request["messages"] if isinstance(message["content"], list) + for block in message["content"] if block["type"] == "thinking" + ] + assert "invalid_old" not in signatures + if "valid_new" in signatures: + return _payload(model, content=[{"type": "text", "text": "done"}]) + return _payload(model, stop_reason="tool_use", content=[ + {"type": "thinking", "thinking": "", "signature": "valid_new"}, + {"type": "tool_use", "id": "t1", "name": "search", "input": {"q": "test"}}, + ]) + + class Search: + name = "search" + + async def ainvoke(self, args): + return "found" + + def to_openai_schema(self): + return {"type": "function", "function": { + "name": self.name, "parameters": {"type": "object"}, + }} + + class DeltaObserver: + wants_llm_delta = streaming + + error = {"message": "Invalid `signature` in `thinking` block. The block is bound to a different conversation."} + async with _client(model, response_payload, error=error) as (client, requests): + result = await run_agent_loop( + system_prompt="s", user_message="hi", initial_messages=initial, + llm=client, tools=[Search()], + observers=[DeltaObserver()], + config=LoopConfig(max_turns=3, max_llm_retries=1, loop_policy=LoopPolicy(no_tool_behavior="stop")), + model_profile=ModelProfile(model_id=model, provider="anthropic", thinking_format="content_block"), + ) + assert result.final_content == "done" + assert len(requests) == 3 # one failed request, recovery, then a clean next turn + assert json.dumps(initial) == original + assert "invalid_old" not in json.dumps(result.messages) + assert "valid_new" in json.dumps(result.messages) + + +@pytest.mark.parametrize("model", MODELS) +@pytest.mark.parametrize("streaming", [False, True]) +async def test_partial_refusal_reaches_both_trajectory_formats(model, streaming, tmp_path): + details = {"type": "refusal", "category": "bio", "explanation": "declined"} + payload = _payload(model, stop_reason="refusal", stop_details=details, + content=[{"type": "text", "text": "Partial output"}]) + + class DeltaObserver: + wants_llm_delta = streaming + + trajectory = TrajectoryFileObserver(tmp_path, filename="refusal", formats=["json", "jsonl"]) + async with _client(model, payload) as (client, _): + await run_agent_loop( + system_prompt="s", user_message="hi", llm=client, tools=[], + observers=[trajectory, DeltaObserver()], + config=LoopConfig(max_turns=3, max_llm_retries=1, loop_policy=LoopPolicy(no_tool_behavior="stop")), + model_profile=ModelProfile(model_id=model, provider="anthropic", thinking_format="content_block"), + ) + events = [json.loads(line) for line in (tmp_path / "refusal.jsonl").read_text().splitlines()] + record = next(event for event in events if event["t"] == "llm") + assistant = next(msg for msg in json.loads((tmp_path / "refusal.json").read_text())["messages"] if msg["role"] == "assistant") + for item in [record, assistant]: + assert item["content"] == "Partial output" + assert item["finish_reason"] == "refusal" + assert item["stop_details"] == details diff --git a/tests/test_llm_runtime_stream_metadata.py b/tests/test_llm_runtime_stream_metadata.py index c6816fd..358436b 100644 --- a/tests/test_llm_runtime_stream_metadata.py +++ b/tests/test_llm_runtime_stream_metadata.py @@ -385,3 +385,79 @@ async def _on_delta(*_): ) assert result is not None assert result.content == "hello world" + + +@pytest.mark.asyncio +async def test_streaming_folds_stop_details_for_a_refusal(monkeypatch): + """A streamed classifier decline must surface the same detail as a + non-streamed one. + + A refusal streams to a clean close carrying no content and no tool call, + so without this the assembled response is identical to a turn the model + chose to end — the loop takes its no-tool exit and the run looks finished. + """ + real_sleep = asyncio.sleep + + async def _noop(_): + await real_sleep(0) + + monkeypatch.setattr("asyncio.sleep", _noop) + + deltas = [ + StreamDelta( + usage={"prompt_tokens": 10, "completion_tokens": 0}, + finish_reason="refusal", + model="claude-x", + stop_details={"type": "refusal", "category": "bio"}, + ), + ] + llm = _StubStreamingLLM(deltas=deltas) + + async def _on_delta(*_): + pass + + result = await call_llm( + llm, [user_msg("hi")], + timeout=10, max_retries=1, turn=0, + on_delta=_on_delta, + ) + + assert result is not None + assert result.content == "" + assert result.finish_reason == "refusal" + assert result.response_metadata["stop_details"] == { + "type": "refusal", "category": "bio", + } + + +@pytest.mark.asyncio +async def test_streaming_without_stop_details_adds_no_key(monkeypatch): + """The ordinary path keeps its metadata exactly as it was.""" + real_sleep = asyncio.sleep + + async def _noop(_): + await real_sleep(0) + + monkeypatch.setattr("asyncio.sleep", _noop) + + deltas = [ + StreamDelta(content="hi"), + StreamDelta( + usage={"prompt_tokens": 7, "completion_tokens": 2}, + finish_reason="stop", + model="claude-x", + ), + ] + llm = _StubStreamingLLM(deltas=deltas) + + async def _on_delta(*_): + pass + + result = await call_llm( + llm, [user_msg("hi")], + timeout=10, max_retries=1, turn=0, + on_delta=_on_delta, + ) + + assert result is not None + assert "stop_details" not in result.response_metadata diff --git a/tests/test_provider_native_clients.py b/tests/test_provider_native_clients.py index 13b85bf..d55fc58 100644 --- a/tests/test_provider_native_clients.py +++ b/tests/test_provider_native_clients.py @@ -706,7 +706,7 @@ def test_anthropic_build_kwargs_matches_installed_sdk_signature(thinking): assert {"temperature", "top_p", "top_k"}.isdisjoint(kwargs.get("extra_body", {})) if thinking: assert kwargs["thinking"] == thinking - assert kwargs["extra_body"] == {"output_config": {"effort": "high"}} + assert kwargs["extra_body"] == {"output_config": {"effort": "high"}} @pytest.mark.asyncio @@ -1332,3 +1332,171 @@ def test_anthropic_folds_every_trailing_per_call_message(monkeypatch): content = msgs[2]["content"] assert [b["type"] for b in content] == ["tool_result", "text", "text"] assert "cache_control" in content[0] and all("cache_control" not in b for b in content[1:]) + + +# ── Refusals: stop_reason / stop_details surfacing ─────────────────────── +# +# A safety classifier decline is an HTTP 200 whose ``content`` is empty, so to +# the loop it is indistinguishable from the model choosing to say nothing: no +# text, no tool call, the no-tool exit, a run that looks clean. These pin the +# detail onto ``response_metadata`` so a trajectory can say WHY. + + +def _refusal_payload(**overrides): + payload = { + "id": "msg_r", "type": "message", "role": "assistant", "model": "claude-x", + "content": [], + "stop_reason": "refusal", "stop_sequence": None, + "stop_details": { + "type": "refusal", "category": "bio", "explanation": "declined", + }, + "usage": {"input_tokens": 10, "output_tokens": 0}, + } + payload.update(overrides) + return payload + + +def test_anthropic_refusal_surfaces_stop_reason_and_details(): + raw = SimpleNamespace( + id="msg_r", model="claude-x", content=[], stop_reason="refusal", + stop_details=SimpleNamespace( + type="refusal", category="bio", explanation="declined", + ), + usage=SimpleNamespace(input_tokens=10, output_tokens=0), + ) + response = ac._to_llm_response(raw) + + assert response.finish_reason == "refusal" + assert response.response_metadata["stop_reason"] == "refusal" + assert response.response_metadata["stop_details"] == { + "type": "refusal", "category": "bio", "explanation": "declined", + } + + +def test_anthropic_keeps_the_raw_stop_reason_next_to_the_normalised_one(): + """``max_tokens`` normalises to ``length``; the provider's word survives.""" + raw = SimpleNamespace( + id="msg_t", model="claude-x", + content=[SimpleNamespace(type="text", text="cut off")], + stop_reason="max_tokens", + usage=SimpleNamespace(input_tokens=1, output_tokens=1), + ) + response = ac._to_llm_response(raw) + + assert response.finish_reason == "length" + assert response.response_metadata["stop_reason"] == "max_tokens" + assert "stop_details" not in response.response_metadata + + +def test_anthropic_ordinary_turn_carries_no_stop_details(): + """The common case stays lean — no empty key on every successful turn.""" + raw = SimpleNamespace( + id="msg_1", model="claude-x", + content=[SimpleNamespace(type="text", text="hello")], + stop_reason="end_turn", + usage=SimpleNamespace(input_tokens=1, output_tokens=1), + ) + response = ac._to_llm_response(raw) + + assert response.response_metadata["stop_reason"] == "end_turn" + assert "stop_details" not in response.response_metadata + + +@pytest.mark.parametrize("details", [ + None, + {}, + SimpleNamespace(type=None, category=None, explanation=None), +]) +def test_anthropic_empty_stop_details_is_omitted_not_recorded_blank(details): + raw = SimpleNamespace( + id="msg_1", model="claude-x", + content=[SimpleNamespace(type="text", text="hi")], + stop_reason="end_turn", stop_details=details, + usage=SimpleNamespace(input_tokens=1, output_tokens=1), + ) + assert "stop_details" not in ac._to_llm_response(raw).response_metadata + + +def test_anthropic_stop_details_passes_through_an_unknown_category(): + """``category`` is an open set; new models add to it. Never filter it.""" + raw = SimpleNamespace( + id="msg_r", model="claude-x", content=[], stop_reason="refusal", + stop_details={"type": "refusal", "category": "some_future_category"}, + usage=SimpleNamespace(input_tokens=1, output_tokens=0), + ) + details = ac._to_llm_response(raw).response_metadata["stop_details"] + assert details["category"] == "some_future_category" + + +@pytest.mark.asyncio +async def test_anthropic_refusal_through_the_real_sdk(): + """End-to-end over the SDK's own parsing, not a hand-built namespace.""" + from anthropic import AsyncAnthropic, DefaultAsyncHttpxClient + + if issubclass(DefaultAsyncHttpxClient, httpx.AsyncClient): + sdk_httpx = httpx + else: + import httpx2 as sdk_httpx + + def respond(request): + return sdk_httpx.Response(200, json=_refusal_payload()) + + c = ac.AnthropicClient("claude-x", api_key="test") + await c._client.close() + async with AsyncAnthropic( + api_key="test", base_url="https://anthropic.invalid", max_retries=0, + http_client=DefaultAsyncHttpxClient(transport=sdk_httpx.MockTransport(respond)), + ) as sdk: + c._client = sdk + response = await c.chat([user_msg("hi")]) + + # The shape that used to be indistinguishable from "said nothing". + assert response.content == "" + assert response.tool_calls == [] + # …and the evidence that it was a decline. + assert response.finish_reason == "refusal" + assert response.response_metadata["stop_reason"] == "refusal" + assert response.response_metadata["stop_details"]["category"] == "bio" + + +@pytest.mark.asyncio +async def test_anthropic_streamed_refusal_reports_the_same_detail(): + """A streamed decline must not be quieter than a non-streamed one.""" + events = [ + SimpleNamespace( + type="message_start", + message=SimpleNamespace( + model="claude-x", + usage=SimpleNamespace(input_tokens=10, cache_read_input_tokens=0), + ), + ), + SimpleNamespace( + type="message_delta", + delta=SimpleNamespace( + stop_reason="refusal", + stop_details=SimpleNamespace( + type="refusal", category="cyber", explanation="declined", + ), + ), + usage=SimpleNamespace(output_tokens=0), + ), + ] + + class _Stream: + def __aiter__(self): + async def gen(): + for event in events: + yield event + return gen() + + c = ac.AnthropicClient("claude-x", api_key="x") + c._client = MagicMock() + c._client.messages.create = AsyncMock(return_value=_Stream()) + + deltas = [d async for d in c.stream([user_msg("hi")])] + + terminal = deltas[-1] + assert terminal.finish_reason == "refusal" + assert terminal.stop_details == { + "type": "refusal", "category": "cyber", "explanation": "declined", + } diff --git a/tests/test_trajectory_stop_reason.py b/tests/test_trajectory_stop_reason.py new file mode 100644 index 0000000..2e31e13 --- /dev/null +++ b/tests/test_trajectory_stop_reason.py @@ -0,0 +1,119 @@ +"""A trajectory must record WHY a turn ended, not only what it produced. + +Three turn endings used to be indistinguishable in a JSONL trajectory: an +answer the model finished, an answer the output cap cut off, and a request a +safety classifier declined. The first two differ only in ``finish_reason``; +the third arrives with empty ``content`` and no tool calls, so it reads as a +turn where the model simply chose to say nothing, and the loop's no-tool exit +ends the run looking clean. + +The 2026-10-05 gdpval triage hit exactly this: 20 runs ended in ``llm_error`` +and 25 more in ``no_tool`` with an empty deliverable, and nothing on disk — +trajectory, job log, or breaker state — recorded which of those were declines. +These pin the fields that answer the question. +""" + +from __future__ import annotations + +import json + +import pytest + +from agent_core.components.observers.trajectory import TrajectoryFileObserver +from agent_core.loop_types import TurnContext + + +def _ctx(**overrides) -> TurnContext: + base = dict( + turn=1, max_turns=10, task_id="t", role_id="r", + ai_text="", thinking="", tool_calls=[], messages=[], + usage=None, metadata={}, + ) + base.update(overrides) + return TurnContext(**base) + + +def _llm_records(tmp_path) -> list[dict]: + lines = (tmp_path / "t.jsonl").read_text().splitlines() + return [r for r in (json.loads(x) for x in lines) if r.get("t") == "llm"] + + +async def _record(tmp_path, ctx: TurnContext) -> dict: + obs = TrajectoryFileObserver(tmp_path, filename="t", formats=["jsonl"]) + await obs.on_llm_response(ctx) + records = _llm_records(tmp_path) + assert len(records) == 1 + return records[0] + + +@pytest.mark.asyncio +async def test_a_refusal_is_recorded_as_a_refusal(tmp_path): + """The shape that used to look like "the model said nothing".""" + record = await _record(tmp_path, _ctx( + finish_reason="refusal", + stop_details={"type": "refusal", "category": "bio", + "explanation": "declined"}, + )) + + assert record["content"] == "" + assert record["tool_calls"] == [] + assert record["finish_reason"] == "refusal" + assert record["stop_details"]["category"] == "bio" + + +@pytest.mark.asyncio +async def test_a_truncated_turn_is_distinguishable_from_a_finished_one(tmp_path): + truncated = await _record(tmp_path / "a", _ctx( + ai_text="half a sent", finish_reason="length", + )) + finished = await _record(tmp_path / "b", _ctx( + ai_text="a whole answer", finish_reason="end_turn", + )) + + assert truncated["finish_reason"] == "length" + # ``end_turn`` carries no information and would be on nearly every line. + assert "finish_reason" not in finished + + +@pytest.mark.asyncio +@pytest.mark.parametrize("finish_reason", ["end_turn", "stop"]) +async def test_an_ordinary_turn_stays_lean(tmp_path, finish_reason): + """No empty keys on the overwhelmingly common path, for either protocol.""" + record = await _record(tmp_path, _ctx( + ai_text="hello", finish_reason=finish_reason, + )) + + assert "finish_reason" not in record + assert "stop_details" not in record + + +@pytest.mark.asyncio +async def test_a_producer_that_sets_nothing_is_no_worse_than_before(tmp_path): + """The fields default, so an un-migrated caller still writes a valid line.""" + record = await _record(tmp_path, _ctx(ai_text="hello")) + + assert record["content"] == "hello" + assert "finish_reason" not in record + assert "stop_details" not in record + + +@pytest.mark.asyncio +async def test_tool_use_turns_record_their_stop_reason(tmp_path): + """Useful for telling a tool-call turn from a text turn when replaying.""" + record = await _record(tmp_path, _ctx( + tool_calls=[{"name": "bash", "args": {}}], finish_reason="tool_use", + )) + + assert record["finish_reason"] == "tool_use" + + +@pytest.mark.parametrize("formats", [["json"], ["json", "jsonl"]]) +async def test_json_snapshot_records_refusal_details(tmp_path, formats): + obs = TrajectoryFileObserver(tmp_path, filename="t", formats=formats) + await obs.on_llm_response(_ctx( + finish_reason="refusal", stop_details={"type": "refusal", "category": "bio"}, + )) + msg = json.loads((tmp_path / "t.json").read_text())["messages"][0] + assert msg["finish_reason"] == "refusal" + assert msg["stop_details"] == {"type": "refusal", "category": "bio"} + assert (tmp_path / "t.jsonl").exists() == ("jsonl" in formats) diff --git a/uv.lock b/uv.lock index dd0800b..d1a4c3c 100644 --- a/uv.lock +++ b/uv.lock @@ -13,7 +13,7 @@ wheels = [ [[package]] name = "anthropic" -version = "1.3.0" +version = "1.11.0" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "anyio" }, @@ -24,9 +24,9 @@ dependencies = [ { name = "sniffio" }, { name = "typing-extensions" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/b4/50/463166f02179ab279edb61de1589a6f69cb3838d6a2fb6f2c92a3f8042f1/anthropic-1.3.0.tar.gz", hash = "sha256:6873492a77ede8849a161ab1bc78bc9a1e492a006d0b5bb4c57ac77845df838a", size = 1148177, upload-time = "2026-09-01T17:37:10.392Z" } +sdist = { url = "https://files.pythonhosted.org/packages/ad/12/9a6ffa397b172adb040008d934a1dfd85d0e4cefc77489416db294ffc880/anthropic-1.11.0.tar.gz", hash = "sha256:3906fabac7ad7b5b46c6186040398fc7826885c77ce34e4dd7849de16fc8d0f8", size = 1390390, upload-time = "2026-09-30T22:58:21.098Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/2c/5d/7863a9961d320c23787c7b594956afe4e878f9c0ae2376b11a20e416791d/anthropic-1.3.0-py3-none-any.whl", hash = "sha256:e7e7dbebf9f3c84a23954ab989378af6ae10a4d1804c81e9fea4b5ced695ce75", size = 1296959, upload-time = "2026-09-01T17:37:08.525Z" }, + { url = "https://files.pythonhosted.org/packages/ab/6f/20edcd2956cac68042b1bcdd13ab3311447519f86b26929ff268dc74c650/anthropic-1.11.0-py3-none-any.whl", hash = "sha256:52f97b2c485cca7ac66058374f5073a3febb7e6849b16989c145d602a3efee21", size = 1681654, upload-time = "2026-09-30T22:58:22.925Z" }, ] [package.optional-dependencies]