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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions agent_core/loop_types.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
48 changes: 43 additions & 5 deletions agent_core/runtime/loop/agent_loop.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
5 changes: 5 additions & 0 deletions changes/62.feature.md
Original file line number Diff line number Diff line change
@@ -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.
34 changes: 34 additions & 0 deletions docs/llm-runtime-boundary.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
137 changes: 137 additions & 0 deletions tests/test_stream_llm_tokens_toggle.py
Original file line number Diff line number Diff line change
@@ -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
Loading