From facab085e1159c38db3a4ffb02dae2fd155b1914 Mon Sep 17 00:00:00 2001 From: zhanghanduo Date: Mon, 5 Oct 2026 15:00:57 +0800 Subject: [PATCH 1/3] =?UTF-8?q?feat(observability):=20=E6=8A=8A=20refusal?= =?UTF-8?q?=20/=20finish=5Freason=20=E6=9A=B4=E9=9C=B2=E5=88=B0=20LLMRespo?= =?UTF-8?q?nse=20=E4=B8=8E=20trace?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 安全分类器拒答是 HTTP 200 + 空 content + stop_reason="refusal",引擎看到的 形状与"模型选择不说话"完全一致:无文本、无 tool call → 走 no_tool 退出 → 整跑结束且 trace 上看起来干净。Anthropic 的 stop_details(type / category / explanation)此前从未被读取,response_metadata 只有 {"id": ...}。 2026-10-05 gdpval-220 取证正是卡在这里:20 个 run 以 llm_error 收场、另有 25 个 以 no_tool + 空交付物收场,而 trace / job.log / trial.log / breaker last_error 四处都没有记录哪些是拒答,只能逐个 case 手工解剖。 本次补齐三段: 1. providers/anthropic.py —— 新增 _anthropic_stop_details(),同时支持 SDK model 对象与 dict;_to_llm_response 把 RAW stop_reason 与 stop_details 放进 response_metadata。raw 值与归一化后的 finish_reason 并存:max_tokens 会被 归一成 length,消费者不该为了判断"provider 是否拒答"而去记住哪些标记被重写。 category 是开放集合(cyber / bio / reasoning_extraction / frontier_llm / general_harms / …,新模型持续新增),原样透传不做白名单。 2. 流式路径 —— StreamDelta 新增 stop_details 字段(沿用 provider 字段"终端 元数据折进 response_metadata"的既有模式),anthropic 的 message_delta 捕获 它,_streaming 的组装器折进 response_metadata。否则一次流式拒答仍然隐形。 3. TurnContext 新增 finish_reason / stop_details 字段,agent_loop 填充, TrajectoryFileObserver 写进 JSONL。finish_reason 此前只在 metadata 里, 观察者拿不到一等字段。end_turn 不写(无信息量且几乎每行都有),其余 (length / tool_use / refusal / …)记录。 两个字段都有默认值,未迁移的 TurnContext 构造点保持有效 —— 不设置的生产者 与改动前一样看不见,但不会出错。 测试:新增 test_trajectory_stop_reason.py(5 条)+ test_provider_native_clients.py 的 refusal 组(7 条,含一条走真实 SDK 解析的 MockTransport 端到端)+ test_llm_runtime_stream_metadata.py 的流式组(2 条)。全量 1779 passed。 ruff 通过。 Co-Authored-By: Claude Opus 5 --- agent_core/components/observers/trajectory.py | 11 ++ agent_core/llm.py | 9 + agent_core/loop_types.py | 16 ++ agent_core/providers/anthropic.py | 58 +++++- agent_core/runtime/loop/_streaming.py | 8 + agent_core/runtime/loop/agent_loop.py | 3 + tests/test_llm_runtime_stream_metadata.py | 76 ++++++++ tests/test_provider_native_clients.py | 168 ++++++++++++++++++ tests/test_trajectory_stop_reason.py | 106 +++++++++++ 9 files changed, 454 insertions(+), 1 deletion(-) create mode 100644 tests/test_trajectory_stop_reason.py diff --git a/agent_core/components/observers/trajectory.py b/agent_core/components/observers/trajectory.py index 8df03d2..92ffeff 100644 --- a/agent_core/components/observers/trajectory.py +++ b/agent_core/components/observers/trajectory.py @@ -504,6 +504,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`` is the uninformative common case and is omitted to keep + # the line small; everything else (``length``, ``tool_use``, + # ``refusal``, …) is recorded. + if ctx.finish_reason and ctx.finish_reason != "end_turn": + 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. diff --git a/agent_core/llm.py b/agent_core/llm.py index e9250d0..4f5549f 100644 --- a/agent_core/llm.py +++ b/agent_core/llm.py @@ -61,6 +61,15 @@ 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]) @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..69b1070 100644 --- a/agent_core/providers/anthropic.py +++ b/agent_core/providers/anthropic.py @@ -182,6 +182,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 @@ -278,6 +279,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) @@ -310,6 +317,7 @@ async def stream( finish_reason=normalize_finish_reason(stop_reason), model=model, reasoning_blocks=ordered if _has_thinking(ordered) else [], + stop_details=stop_details, ) @@ -683,6 +691,41 @@ 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 + # 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. @@ -787,6 +830,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 +850,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..ed2b776 100644 --- a/agent_core/runtime/loop/_streaming.py +++ b/agent_core/runtime/loop/_streaming.py @@ -258,6 +258,10 @@ 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] = {} 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 +290,8 @@ 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 visible_content = accumulated if visible_content: visible_content = ( @@ -387,6 +393,8 @@ 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 # 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..313efc2 100644 --- a/agent_core/runtime/loop/agent_loop.py +++ b/agent_core/runtime/loop/agent_loop.py @@ -1261,6 +1261,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 +1269,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/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..ec7c17d 100644 --- a/tests/test_provider_native_clients.py +++ b/tests/test_provider_native_clients.py @@ -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..9e258de --- /dev/null +++ b/tests/test_trajectory_stop_reason.py @@ -0,0 +1,106 @@ +"""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 +async def test_an_ordinary_turn_stays_lean(tmp_path): + """No empty keys on the overwhelmingly common path.""" + record = await _record(tmp_path, _ctx( + ai_text="hello", finish_reason="end_turn", + )) + + 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" From 20b288f7b0e30aab49f4d2d6bd32bcd649daf1e8 Mon Sep 17 00:00:00 2001 From: zhanghanduo Date: Mon, 5 Oct 2026 15:41:33 +0800 Subject: [PATCH 2/3] fix(anthropic): complete stop telemetry and current Claude replay compatibility --- .github/workflows/ci.yml | 3 + agent_core/components/observers/trajectory.py | 4 + agent_core/llm.py | 5 + agent_core/providers/anthropic.py | 120 ++++++- agent_core/runtime/loop/_streaming.py | 10 + agent_core/runtime/loop/agent_loop.py | 26 ++ changes/58.feature.md | 1 + docs/provider-substrate-boundary.md | 37 +++ tests/test_anthropic_latest_models.py | 308 ++++++++++++++++++ tests/test_provider_native_clients.py | 2 +- tests/test_trajectory_stop_reason.py | 12 + uv.lock | 6 +- 12 files changed, 515 insertions(+), 19 deletions(-) create mode 100644 changes/58.feature.md create mode 100644 tests/test_anthropic_latest_models.py 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 92ffeff..8250a69 100644 --- a/agent_core/components/observers/trajectory.py +++ b/agent_core/components/observers/trajectory.py @@ -542,6 +542,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 != "end_turn": + 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 4f5549f..c6bfb20 100644 --- a/agent_core/llm.py +++ b/agent_core/llm.py @@ -70,6 +70,11 @@ class StreamDelta: # 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/providers/anthropic.py b/agent_core/providers/anthropic.py index 69b1070..a6a713c 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,35 @@ 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 a history edit invalidated signed thinking. + + 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. + """ + 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 + and "bound to a different conversation" 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 +181,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, @@ -196,7 +228,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": @@ -218,10 +250,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 "", @@ -270,6 +305,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, @@ -304,6 +345,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, @@ -318,6 +368,8 @@ async def stream( model=model, reasoning_blocks=ordered if _has_thinking(ordered) else [], stop_details=stop_details, + stop_reason=stop_reason, + thinking_history_reset=reset, ) @@ -387,6 +439,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]]: @@ -517,6 +589,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 @@ -546,14 +623,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 @@ -717,6 +797,11 @@ def _anthropic_stop_details(raw: Any) -> dict[str, Any] | 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"): @@ -803,6 +888,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", diff --git a/agent_core/runtime/loop/_streaming.py b/agent_core/runtime/loop/_streaming.py index ed2b776..8dbf277 100644 --- a/agent_core/runtime/loop/_streaming.py +++ b/agent_core/runtime/loop/_streaming.py @@ -262,6 +262,8 @@ async def _stream_llm_response( # 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) @@ -292,6 +294,10 @@ def _assembled_response() -> LLMResponse: ) 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 = ( @@ -395,6 +401,10 @@ async def _close_chunk_stream() -> 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 313efc2..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) diff --git a/changes/58.feature.md b/changes/58.feature.md new file mode 100644 index 0000000..a74dc58 --- /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, and recover once from thinking signatures invalidated by conversation edits. diff --git a/docs/provider-substrate-boundary.md b/docs/provider-substrate-boundary.md index 105b37f..f2ced52 100644 --- a/docs/provider-substrate-boundary.md +++ b/docs/provider-substrate-boundary.md @@ -67,3 +67,40 @@ 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. 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 explicitly reports a thinking signature bound to +a different conversation (for example 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. 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..61d6dcf --- /dev/null +++ b/tests/test_anthropic_latest_models.py @@ -0,0 +1,308 @@ +"""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), + ("Invalid signature in thinking block", 1), + ("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_provider_native_clients.py b/tests/test_provider_native_clients.py index ec7c17d..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 diff --git a/tests/test_trajectory_stop_reason.py b/tests/test_trajectory_stop_reason.py index 9e258de..cf9ffb9 100644 --- a/tests/test_trajectory_stop_reason.py +++ b/tests/test_trajectory_stop_reason.py @@ -104,3 +104,15 @@ async def test_tool_use_turns_record_their_stop_reason(tmp_path): )) 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] From 2cd1fb5416081aec84b822ef43aba683b73e3816 Mon Sep 17 00:00:00 2001 From: zhanghanduo Date: Mon, 5 Oct 2026 16:11:47 +0800 Subject: [PATCH 3/3] fix(anthropic): broaden signature-retry match and tidy stop telemetry Follow-up to 20b288f (whose body was lost to truncation). Summary of that commit: forward effort without thinking; keep tool_use blocks in provider order among signed thinking (canonical tool_calls stay authoritative); retry once without historical thinking when a thinking signature is rejected, and report thinking_history_reset so the loop drops stale signatures before storing the new turn; record finish_reason/stop_details in JSON trajectories. This commit: - Signature retry now matches "signature" + "thinking" instead of the exact "bound to a different conversation" phrase. The wording is not a documented contract, and omitting historical thinking is always a valid request, so a narrower match would silently disable recovery on any rewording. - Trajectories also omit finish_reason "stop" (OpenAI-compatible normal end), matching the existing end_turn omission. - Docs/changelog: note the per-request double cost for direct client callers that ignore thinking_history_reset, and the effort-without-thinking change. Co-Authored-By: Claude Opus 5.5 (1M context) --- agent_core/components/observers/trajectory.py | 11 +++++++---- agent_core/providers/anthropic.py | 12 +++++++----- changes/58.feature.md | 2 +- docs/provider-substrate-boundary.md | 15 +++++++++------ tests/test_anthropic_latest_models.py | 3 ++- tests/test_trajectory_stop_reason.py | 7 ++++--- 6 files changed, 30 insertions(+), 20 deletions(-) diff --git a/agent_core/components/observers/trajectory.py b/agent_core/components/observers/trajectory.py index 8250a69..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", @@ -508,10 +511,10 @@ async def on_llm_response( # 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`` is the uninformative common case and is omitted to keep - # the line small; everything else (``length``, ``tool_use``, + # ``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 != "end_turn": + 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 @@ -542,7 +545,7 @@ 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 != "end_turn": + 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 diff --git a/agent_core/providers/anthropic.py b/agent_core/providers/anthropic.py index a6a713c..af48724 100644 --- a/agent_core/providers/anthropic.py +++ b/agent_core/providers/anthropic.py @@ -139,12 +139,17 @@ def _build_kwargs( return kwargs async def _create_message(self, kwargs: dict[str, Any]) -> tuple[Any, bool]: - """Retry once when a history edit invalidated signed thinking. + """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 @@ -152,10 +157,7 @@ async def _create_message(self, kwargs: dict[str, Any]) -> tuple[Any, bool]: 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 - and "bound to a different conversation" in error - ): + if not ("signature" in error and "thinking" in error): raise messages, stripped = _without_thinking_blocks(kwargs["messages"]) if not stripped: diff --git a/changes/58.feature.md b/changes/58.feature.md index a74dc58..dba3c0e 100644 --- a/changes/58.feature.md +++ b/changes/58.feature.md @@ -1 +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, and recover once from thinking signatures invalidated by conversation edits. +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 f2ced52..91d315b 100644 --- a/docs/provider-substrate-boundary.md +++ b/docs/provider-substrate-boundary.md @@ -73,7 +73,9 @@ this normalization silently disables truncation recovery for its protocol. 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. Manual budgets and disabled thinking +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. @@ -83,15 +85,16 @@ 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 explicitly reports a thinking signature bound to -a different conversation (for example after client-side compaction), the -adapter retries once without historical thinking or redacted-thinking blocks. +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. Other 400s and a failed retry -propagate. Hosts should keep conversation prefixes stable to preserve reasoning. +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), diff --git a/tests/test_anthropic_latest_models.py b/tests/test_anthropic_latest_models.py index 61d6dcf..b5ac6a2 100644 --- a/tests/test_anthropic_latest_models.py +++ b/tests/test_anthropic_latest_models.py @@ -202,7 +202,8 @@ async def test_prefix_bound_signature_retry_is_once_and_does_not_edit_history(mo @pytest.mark.parametrize("streaming", [False, True]) @pytest.mark.parametrize("message,expected_calls", [ ("max_tokens is too large", 1), - ("Invalid signature in thinking block", 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): diff --git a/tests/test_trajectory_stop_reason.py b/tests/test_trajectory_stop_reason.py index cf9ffb9..2e31e13 100644 --- a/tests/test_trajectory_stop_reason.py +++ b/tests/test_trajectory_stop_reason.py @@ -76,10 +76,11 @@ async def test_a_truncated_turn_is_distinguishable_from_a_finished_one(tmp_path) @pytest.mark.asyncio -async def test_an_ordinary_turn_stays_lean(tmp_path): - """No empty keys on the overwhelmingly common path.""" +@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="end_turn", + ai_text="hello", finish_reason=finish_reason, )) assert "finish_reason" not in record