diff --git a/agent_core/loop_types.py b/agent_core/loop_types.py index c0ced30..cb6ba95 100644 --- a/agent_core/loop_types.py +++ b/agent_core/loop_types.py @@ -150,6 +150,14 @@ class LoopConfig: # Abort reasoning-only streams after either enabled bound. reasoning_only_timeout_s: float | None = None reasoning_only_max_tokens: int | None = None + # Stream LLM tokens, overriding what the run would decide on its own. + # ``None`` (default) keeps the existing behaviour: stream when something + # needs the deltas and the protocol's streamed replay is verified. ``True`` + # streams regardless — the transport, not the observers, is then the reason + # (a gateway that times out waiting for a non-streaming response's headers + # leaves no other option). ``False`` never streams. An explicit value makes + # the host responsible for its protocol's streamed-replay fidelity. + stream_llm_tokens: bool | None = None # Total budget across admission, attempts, backoff, and recovery. logical_call_timeout_s: float | None = None context_token_limit: int = 120_000 diff --git a/agent_core/runtime/loop/agent_loop.py b/agent_core/runtime/loop/agent_loop.py index 3f4af92..d8a2414 100644 --- a/agent_core/runtime/loop/agent_loop.py +++ b/agent_core/runtime/loop/agent_loop.py @@ -93,6 +93,27 @@ # loops (e.g. a flaky LLM that keeps emitting refusals/duplicates). EXTRA_ATTEMPTS_BUFFER = 200 +#: Protocols whose STREAMED turn is not proven here to replay like its +#: non-streaming twin, so ``stream_llm_tokens`` stays off for them unless the +#: host states otherwise. +#: +#: Every one of these clients implements ``stream``; the question is fidelity of +#: the replayed turn, and the bar is a test in this repository. ``anthropic`` +#: clears it — ``stream`` rebuilds the provider's verbatim block list, including +#: the ``signature_delta`` that arrives after a thinking block's text and the +#: payload a ``redacted_thinking`` block carries with no deltas at all, and +#: ``test_anthropic_latest_models.py`` runs the signature round-trip, the +#: prefix-bound retry and the reset commit parametrized over streaming. +#: ``responses`` has one streaming test and it only covers error routing; +#: ``bedrock`` has none, and its transport is a different SDK client. +#: +#: The three were suppressed together and unconditionally when the loop engine +#: was extracted (2026-09-01), one day BEFORE the provider substrate landed +#: that block-list fidelity (2026-09-02) — so this was once correct for all +#: three and nobody revisited it. Removing a protocol from this set is a matter +#: of writing the test that proves its streamed replay. +UNVERIFIED_STREAM_PROTOCOLS = frozenset({"responses", "bedrock"}) + # Signature: (turn_index, messages_snapshot, metadata) -> awaitable None. # Fires once per completed turn, after observer `on_turn_end` and any @@ -423,16 +444,33 @@ async def _run_loop_inner( last_input_tokens = 0 last_output_tokens = 0 - stream_llm_tokens = bool( + # Something in the run needs the token deltas: the reasoning-only watchdog + # reads the stream, and so does any observer that declares it. + stream_requested = bool( cfg.reasoning_only_timeout_s or cfg.reasoning_only_max_tokens ) or any( bool(getattr(observer, "wants_llm_delta", False)) for observer in obs ) - if getattr(profile, "protocol", "chat_completions") in ( - "anthropic", "responses", "bedrock", - ): - stream_llm_tokens = False + protocol = getattr(profile, "protocol", "chat_completions") + if cfg.stream_llm_tokens is not None: + # The host stated it; the transport may be the reason (see LoopConfig). + stream_llm_tokens = cfg.stream_llm_tokens + else: + stream_llm_tokens = ( + stream_requested and protocol not in UNVERIFIED_STREAM_PROTOCOLS + ) + if stream_requested and not stream_llm_tokens: + # Previously silent, and silence here is expensive: a profile that + # configures the reasoning-only watchdog on one of these protocols + # gets no watchdog at all, and nothing says so. + logger.warning( + "agent_loop: deltas were requested but protocol %r streams " + "unverified here, so the run is non-streaming and any " + "reasoning-only watchdog is inert; set " + "LoopConfig.stream_llm_tokens=True to stream anyway", + protocol, + ) max_attempts = cfg.max_turns + EXTRA_ATTEMPTS_BUFFER turn = 0 diff --git a/changes/62.feature.md b/changes/62.feature.md new file mode 100644 index 0000000..48b477d --- /dev/null +++ b/changes/62.feature.md @@ -0,0 +1,5 @@ +`LoopConfig.stream_llm_tokens` (`None` by default) overrides the loop's own choice of transport: `True` streams the turn, `False` does not, `None` keeps the existing automatic selection (the reasoning-only watchdog or an observer's `wants_llm_delta`). It exists because the transport can be the reason to stream rather than any observer — a gateway that abandons a non-streaming request while waiting for response headers leaves no alternative, and for native Anthropic those headers arrive only when generation finishes. Measured on llm-hub's "Claude Code" channel on 2026-10-05, an identical `claude-opus-5-5` request at `effort=max` returned `502 upstream_unreachable — first byte timeout` after 76s non-streaming and HTTP 200 when streamed (first byte 3.9s, generation 242s); concurrency was ruled out. + +Behaviour change in automatic mode: `anthropic` now streams when deltas are requested, where all three native protocols were previously suppressed unconditionally. That suppression predated, by one day, the provider substrate work that made `AnthropicClient.stream` rebuild the verbatim block list (trailing `signature_delta`, `redacted_thinking`) so a streamed turn replays like its non-streaming twin — `test_anthropic_latest_models.py` has covered it parametrized over streaming ever since. `responses` and `bedrock` keep the old behaviour as `UNVERIFIED_STREAM_PROTOCOLS` (exported) until a test proves their streamed replay; removing one from that set is that test's job. Consumers whose Anthropic-protocol profiles set `reasoning_only_*` or register a `wants_llm_delta` observer will now actually stream — their client must implement `stream`, since there is no fallback to `chat`. + +The suppression is also no longer silent: when automatic mode drops deltas something asked for, the loop warns. ApodexHarness's `tests/contract/test_streaming_deltas.py` documents the old selection expression verbatim in its docstring and should be updated when the pin moves. diff --git a/docs/llm-runtime-boundary.md b/docs/llm-runtime-boundary.md index 2fd0438..e4d2922 100644 --- a/docs/llm-runtime-boundary.md +++ b/docs/llm-runtime-boundary.md @@ -72,3 +72,37 @@ hosts are known to add, for anyone who wants to describe such a mapping as a `tests/test_loop_types.py` asserts these fields stay an open mapping, so a later well-meant tightening fails loudly. + +## Who decides whether a turn streams + +Streaming is opt-in and, by default, the RUN decides: it is selected when the +reasoning-only watchdog is configured (`reasoning_only_timeout_s` / +`reasoning_only_max_tokens` — the guard reads the stream) or when an observer +declares `wants_llm_delta = True`. Implementing `on_llm_delta` without that +attribute gets silence. + +`LoopConfig.stream_llm_tokens` overrides that decision: `True` streams, `False` +does not, `None` (the default) keeps the automatic choice. An explicit value +exists because the TRANSPORT can be the reason rather than any observer — a +gateway that abandons a non-streaming request while waiting for its response +headers leaves no other option, and for a non-streaming Anthropic request those +headers arrive only once generation is complete. Measured on llm-hub's "Claude +Code" channel, 2026-10-05: an identical `claude-opus-5-5` request at +`effort=max` returned HTTP 502 `upstream_unreachable — first byte timeout` after +76s non-streaming, and HTTP 200 streamed, first byte at 3.9s and the generation +taking 242s. Concurrency was ruled out (five light requests in flight all +succeeded; five heavy ones failed at the same 76s mark). + +In automatic mode the choice is also gated by protocol: +`UNVERIFIED_STREAM_PROTOCOLS` (`responses`, `bedrock`) stays non-streaming +because no test here proves their streamed turn replays like its non-streaming +twin. `anthropic` is NOT in that set: its `stream` rebuilds the provider's +verbatim block list — including a thinking block's trailing `signature_delta` +and a `redacted_thinking` payload that arrives with no deltas — and +`test_anthropic_latest_models.py` exercises the signature round-trip, the +prefix-bound retry and the reset commit parametrized over streaming. An +explicit `stream_llm_tokens` ignores the gate; the host then owns that fidelity. + +When automatic mode suppresses deltas something asked for, the loop logs a +warning. It used to be silent, which meant a profile configuring the +reasoning-only watchdog on a gated protocol got no watchdog and no sign of it. diff --git a/tests/test_stream_llm_tokens_toggle.py b/tests/test_stream_llm_tokens_toggle.py new file mode 100644 index 0000000..03c48a8 --- /dev/null +++ b/tests/test_stream_llm_tokens_toggle.py @@ -0,0 +1,137 @@ +"""Who decides whether a turn streams, and what a protocol gate may not hide. + +Streaming used to be chosen entirely by the run: the reasoning-only watchdog or +an observer declaring ``wants_llm_delta`` asked for deltas, and three native +protocols were then suppressed unconditionally. That suppression was silent, so +a profile configuring the watchdog on one of them got no watchdog and no +warning, and a host whose GATEWAY requires streaming had no way to say so -- +llm-hub's "Claude Code" channel abandons a non-streaming request after ~76s of +waiting for response headers, which for a slow model is most real turns. +""" +from __future__ import annotations + +from typing import Any + +import pytest + +from agent_core.llm import LLMResponse, StreamDelta +from agent_core.loop_types import LoopConfig, LoopPolicy +from agent_core.runtime.loop.agent_loop import ( + UNVERIFIED_STREAM_PROTOCOLS, + run_agent_loop, +) +from agent_core.runtime.loop.model_profile import ModelProfile + + +class RecordingLLM: + """Answers either way and records which surface the loop reached for.""" + + def __init__(self) -> None: + self.chat_calls = 0 + self.stream_calls = 0 + + async def chat(self, messages, **_kwargs) -> LLMResponse: + self.chat_calls += 1 + return LLMResponse(content="finished") + + async def stream(self, messages, **_kwargs): + self.stream_calls += 1 + yield StreamDelta(content="finished") + yield StreamDelta(content="", finish_reason="stop") + + +class DeltaObserver: + wants_llm_delta = True + + def __init__(self) -> None: + self.deltas: list[str] = [] + + async def on_llm_delta(self, ctx: Any) -> None: + # ``LLMDeltaContext``, not the provider's raw chunk. + self.deltas.append(getattr(ctx, "delta", "") or "") + + +def _config(**kwargs: Any) -> LoopConfig: + return LoopConfig( + max_turns=2, max_llm_retries=1, + loop_policy=LoopPolicy(no_tool_behavior="stop"), **kwargs, + ) + + +async def _run(llm: RecordingLLM, *, protocol: str, observers: list[Any] | None = None, + **cfg: Any) -> None: + await run_agent_loop( + system_prompt="system", user_message="start", llm=llm, tools=[], + config=_config(**cfg), + model_profile=ModelProfile(model_id="m", provider="p", protocol=protocol), + observers=observers or [], + ) + + +@pytest.mark.asyncio +async def test_anthropic_streams_when_an_observer_wants_deltas() -> None: + """The regression this file exists for. + + ``anthropic`` was suppressed one day before the provider substrate taught + ``AnthropicClient.stream`` to rebuild the verbatim block list (signature + deltas, redacted thinking) that makes a streamed turn replay like its + non-streaming twin. It stayed suppressed for a month. + """ + llm, observer = RecordingLLM(), DeltaObserver() + await _run(llm, protocol="anthropic", observers=[observer]) + assert (llm.stream_calls, llm.chat_calls) == (1, 0) + assert "finished" in "".join(observer.deltas) + + +@pytest.mark.asyncio +@pytest.mark.parametrize("protocol", sorted(UNVERIFIED_STREAM_PROTOCOLS)) +async def test_unverified_protocols_stay_non_streaming_but_say_so( + protocol, caplog, +) -> None: + """Silence was the bug, not the suppression: these two have no test proving + their streamed replay, so the default holds — and now announces itself.""" + llm = RecordingLLM() + with caplog.at_level("WARNING"): + await _run(llm, protocol=protocol, observers=[DeltaObserver()]) + assert (llm.stream_calls, llm.chat_calls) == (0, 1) + assert any("watchdog is inert" in r.getMessage() for r in caplog.records) + + +@pytest.mark.asyncio +@pytest.mark.parametrize("protocol", ["anthropic", "bedrock", "chat_completions"]) +async def test_explicit_true_streams_on_any_protocol(protocol) -> None: + """The transport, not the observers, can be the reason to stream: a host + that states it takes responsibility for its protocol's replay fidelity.""" + llm = RecordingLLM() + await _run(llm, protocol=protocol, stream_llm_tokens=True) + assert (llm.stream_calls, llm.chat_calls) == (1, 0) + + +@pytest.mark.asyncio +async def test_explicit_false_outranks_a_delta_hungry_observer() -> None: + llm, observer = RecordingLLM(), DeltaObserver() + await _run(llm, protocol="chat_completions", observers=[observer], + stream_llm_tokens=False) + assert (llm.stream_calls, llm.chat_calls) == (0, 1) + assert observer.deltas == [] + + +@pytest.mark.asyncio +async def test_nothing_asking_keeps_the_cheaper_non_streaming_call() -> None: + llm = RecordingLLM() + await _run(llm, protocol="chat_completions") + assert (llm.stream_calls, llm.chat_calls) == (0, 1) + + +@pytest.mark.asyncio +async def test_reasoning_only_watchdog_also_selects_streaming() -> None: + """The watchdog reads the stream, so configuring it is itself a request.""" + llm = RecordingLLM() + await _run(llm, protocol="anthropic", reasoning_only_timeout_s=120) + assert (llm.stream_calls, llm.chat_calls) == (1, 0) + + +def test_the_verified_set_is_stated_not_guessed() -> None: + """``anthropic`` leaving this set is the behaviour change; pin it so a + future edit has to be deliberate.""" + assert frozenset({"responses", "bedrock"}) == UNVERIFIED_STREAM_PROTOCOLS