fix(compression): stop the summary stream at the host's own deadline

CompressionCommitFence.set_total_ceiling_seconds documents its deadline as
"shared by the host and worker", but only the host ever read it. The worker's
streamed summary bounds itself with _aux_stream_total_ceiling() instead —
max(600, 4 * aux_timeout) — which is >= the host's total ceiling for every
configured timeout AND starts counting later (after pool admission,
_serialize_for_summary, prompt build and TTFT). A stream that outlives its
abandoned host is therefore not an edge case; it is the guaranteed outcome of
every total-ceiling timeout.

8207862212 closed the first half: a cancelled fence now releases the
compression owner, freeing its pool slot and session lease. Its own comment
leaves the second half open — the isolated provider daemon that holds the
socket keeps streaming "until the auxiliary stream's longer absolute ceiling
expires". With the #99692 reporter's auxiliary.compression.timeout: 600 that
is 2400s of an orphaned ~500K-token summary the fence is already guaranteed to
refuse, and because the session never shrank, every following turn stacks a
fresh orphan on top of the last.

Publish the fence's deadline as an absolute monotonic instant
(CompressionCommitFence.deadline_monotonic) and give the auxiliary layer the
return leg it was missing: aux_stream_deadline() installs it thread-locally,
_ChatStreamAccumulator.feed() stops the stream once it passes, and
_run_protected_sync_provider_call propagates it onto the provider daemon
(thread-locals do not cross that boundary, so an owner-thread-only install
would be inert on exactly the path large-session compression takes).

Absolute, not relative: the deadline is unaffected by however long dispatch and
TTFT took before the accumulator was constructed. Checked as well as — not
instead of — the existing ceiling, so every caller without a host deadline is
byte-for-byte unchanged, and the "timed out" phrasing keeps _is_timeout_error
classification identical to a request timeout.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EKrRS7LVgyHf2WQkEahSwu
This commit is contained in:
joaomarcos
2026-08-31 18:27:07 -03:00
committed by Teknium
parent 9514d354ca
commit 904e5bb572
3 changed files with 409 additions and 4 deletions
+80 -3
View File
@@ -453,6 +453,16 @@ _aux_progress = threading.local()
_aux_dispatch = threading.local()
_aux_provider_response = threading.local()
# Absolute wall-clock deadline (time.monotonic) of the HOST waiting for this
# auxiliary call, when it has one (#99692). Liveness alone is not enough: a
# host also stops waiting at its own total ceiling, and the streamed consumer
# below bounds itself only by _aux_stream_total_ceiling() — a budget derived
# from the aux request timeout, which is >= the host ceiling for every
# configured value AND starts counting later. So the stream that outlives its
# abandoned host is not an edge case; it is the guaranteed outcome of every
# total-ceiling timeout.
_aux_stream_deadline = threading.local()
def _notify_aux_progress() -> None:
"""Tick the installed forward-progress hook, if any. Never raises."""
@@ -587,6 +597,38 @@ def aux_progress_hook(hook):
yield
def _current_aux_stream_deadline() -> Optional[float]:
"""The waiting host's absolute monotonic deadline, if one is installed."""
return getattr(_aux_stream_deadline, "value", None)
@contextlib.contextmanager
def aux_stream_deadline(deadline: Optional[float]):
"""Publish the waiting host's absolute deadline to the stream consumer.
*deadline* is a ``time.monotonic()`` timestamp — the same instant the host
itself stops waiting — or ``None`` for callers with no host deadline (a
no-op passthrough, so callers can wire it unconditionally). Re-entrant-safe.
#99692: the progress hook is a one-way channel (worker -> host). This is the
return leg. ``8207862212`` releases the compression OWNER when the fence is
cancelled, but the isolated provider daemon
(:func:`_run_protected_sync_provider_call`) that holds the socket keeps
streaming to its own ``_aux_stream_total_ceiling`` budget — >= the host's
ceiling by construction — billing an abandoned summary the commit fence is
already guaranteed to refuse, and stacking one fresh orphan per turn on a
session that compression never managed to shrink.
"""
previous = getattr(_aux_stream_deadline, "value", None)
_aux_stream_deadline.value = (
deadline if isinstance(deadline, (int, float)) else previous
)
try:
yield
finally:
_aux_stream_deadline.value = previous
# Back-compat alias — the timing hooks were introduced with this name.
_aux_timing_hook = _aux_thread_local_hook
@@ -629,6 +671,11 @@ def _run_protected_sync_provider_call(
# the protected daemon path is taken.
dispatch_hook = getattr(_aux_dispatch, "hook", None)
provider_response_hook = getattr(_aux_provider_response, "hook", None)
# #99692: the stream is consumed on the daemon below, and thread-locals do
# not cross that boundary — an owner-thread-only deadline would leave the
# fix inert on exactly the path large-session compression takes (protected
# call + hard-cancel source installed).
host_deadline = _current_aux_stream_deadline()
provider_context = contextvars.copy_context()
done = threading.Event()
outcome: dict[str, Any] = {}
@@ -639,6 +686,7 @@ def _run_protected_sync_provider_call(
aux_progress_hook(progress_hook),
_aux_thread_local_hook(_aux_dispatch, dispatch_hook),
_aux_thread_local_hook(_aux_provider_response, provider_response_hook),
aux_stream_deadline(host_deadline),
aux_interrupt_protection(cancel_check=cancel_check),
):
outcome["result"] = callback(kwargs)
@@ -9872,7 +9920,11 @@ def _aggregate_chat_stream(
Accumulation is shared with the async mirror via
:class:`_ChatStreamAccumulator`.
"""
acc = _ChatStreamAccumulator(model=model, total_ceiling=total_ceiling)
acc = _ChatStreamAccumulator(
model=model,
total_ceiling=total_ceiling,
host_deadline=_current_aux_stream_deadline(),
)
try:
for chunk in chunks:
acc.feed(chunk)
@@ -9894,9 +9946,20 @@ class _ChatStreamAccumulator:
tool-call delta reassembly, same "timed out" ceiling phrasing).
"""
def __init__(self, model: str = "", total_ceiling: Optional[float] = None):
def __init__(
self,
model: str = "",
total_ceiling: Optional[float] = None,
host_deadline: Optional[float] = None,
):
self._started = time.monotonic()
self._total_ceiling = total_ceiling
# #99692: absolute instant the WAITING HOST gives up. Checked as well
# as (not instead of) the ceiling above: the ceiling still bounds
# callers with no host deadline, and the host deadline is absolute, so
# it is unaffected by however long dispatch and TTFT took before this
# accumulator was constructed.
self._host_deadline = host_deadline
self.content_parts: List[str] = []
self.reasoning_parts: List[str] = []
self.reasoning_details: List[Any] = []
@@ -9920,6 +9983,16 @@ class _ChatStreamAccumulator:
f"Auxiliary streamed call timed out after {self._total_ceiling:.0f}s "
"total ceiling (stream still open but over budget)"
)
if (
self._host_deadline is not None
and time.monotonic() >= self._host_deadline
):
raise TimeoutError(
"Auxiliary streamed call timed out at the host compression "
f"deadline after {time.monotonic() - self._started:.0f}s "
"(the caller already stopped waiting; streaming on would only "
"pin its session lease)"
)
self.resp_id = getattr(chunk, "id", None) or self.resp_id
self.resp_model = getattr(chunk, "model", None) or self.resp_model
chunk_usage = getattr(chunk, "usage", None)
@@ -10031,7 +10104,11 @@ async def _aggregate_chat_stream_async(
the sync helper raises. Same accumulation and ceiling semantics via
:class:`_ChatStreamAccumulator`.
"""
acc = _ChatStreamAccumulator(model=model, total_ceiling=total_ceiling)
acc = _ChatStreamAccumulator(
model=model,
total_ceiling=total_ceiling,
host_deadline=_current_aux_stream_deadline(),
)
try:
async for chunk in chunks:
acc.feed(chunk)
+34 -1
View File
@@ -748,6 +748,20 @@ class CompressionCommitFence:
deadline = self._deadline
return deadline is not None and time.monotonic() >= deadline
@property
def deadline_monotonic(self) -> float | None:
"""The armed deadline as an absolute ``time.monotonic()`` instant.
:meth:`set_total_ceiling_seconds` documents this deadline as "shared by
the host and worker", but until #99692 only the host could read it —
``deadline_exceeded`` answers "is it past?" for a caller that is already
polling, which is useless to a worker blocked inside a provider stream.
Publishing the instant itself lets the worker's stream consumer stop at
exactly the moment the host stops waiting (see
``auxiliary_client.aux_stream_deadline``).
"""
return self._deadline
def seconds_since_progress(self) -> float:
"""Seconds since the worker last reported forward progress."""
return max(0.0, time.monotonic() - self._last_progress)
@@ -3992,11 +4006,28 @@ def compress_context(
from agent.auxiliary_client import (
aux_interrupt_protection,
aux_progress_hook,
aux_stream_deadline,
)
_progress_hook = (
commit_fence.touch_progress if commit_fence is not None
else (lambda: None)
)
# #99692: the progress hook above is the worker -> host leg; this is the
# return leg. _compression_cancel_requested (below) releases the compression
# OWNER when the host gives up, but the isolated provider daemon that
# actually holds the socket keeps streaming to its own budget —
# ``_aux_stream_total_ceiling`` = max(600, 4 * aux_timeout), which is >=
# the host's total ceiling for every configured timeout and starts
# counting later (after admission, serialization, prompt build and TTFT).
# With ``auxiliary.compression.timeout: 600`` that is 2400s of an
# orphaned 500K-token summary the commit fence is already guaranteed to
# refuse: paid tokens, a pinned HTTP connection, and — since every new
# turn re-triggers compression on a session that never shrank — a fresh
# orphan stacked on top of the last one. Sharing the host's absolute
# deadline makes the stream stop when the host it serves stops waiting.
_host_stream_deadline = (
commit_fence.deadline_monotonic if commit_fence is not None else None
)
# F4 state-ordering (#76354): a LATE successful summary must not undo
# the timeout cooldown the host recorded. Install a cancellation
# check the compressor consults BEFORE clearing the failure cooldown;
@@ -4042,7 +4073,9 @@ def compress_context(
)
compressed = messages
else:
with aux_progress_hook(_progress_hook), aux_interrupt_protection(
with aux_progress_hook(_progress_hook), aux_stream_deadline(
_host_stream_deadline
), aux_interrupt_protection(
cancel_check=_compression_cancel_requested
):
compressed = compress_fn(messages, **compress_kwargs)
@@ -0,0 +1,295 @@
"""#99692 — the streamed auxiliary summary must not outlive its compression host.
Background
----------
``run_compress_context_with_progress_timeout`` arms a wall-clock deadline on the
``CompressionCommitFence`` (``set_total_ceiling_seconds``), whose docstring calls
it "the wall-clock deadline **shared by the host and worker**". Only the host
ever read it.
``8207862212`` (fix(compression): stop timeout paths from blocking retries)
closed the first half: a cancelled fence now releases the compression OWNER,
which frees the pool slot and the session lease. It left the second half open
by design — its own comment says the isolated provider daemon runs on "until
the auxiliary stream's longer absolute ceiling expires".
That ceiling is ``_aux_stream_total_ceiling`` = ``max(600, 4 * aux_timeout)``:
>= the default host ceiling (600s) for every configured timeout, and it starts
counting later (after pool admission, serialization, prompt build and TTFT).
So the daemon holding the socket is *always* still streaming when its host gives
up — 2400s with the reporter's ``auxiliary.compression.timeout: 600`` — billing
every token of a summary the fence is already guaranteed to refuse, and stacking
one fresh orphan per turn because the session never shrank.
These tests pin the missing half of that shared deadline: the stream consumer
must stop at the host's deadline, including on the isolated provider daemon
that ``_run_protected_sync_provider_call`` spawns.
"""
from __future__ import annotations
import ast
import asyncio
import inspect
import threading
import time
from pathlib import Path
from types import SimpleNamespace
import pytest
from agent import auxiliary_client as aux
from agent.conversation_compression import (
DEFAULT_CONTEXT_TOTAL_CEILING_SECONDS,
CompressionCommitFence,
)
def _chunk(text: str) -> SimpleNamespace:
return SimpleNamespace(
id="resp-1",
model="test-model",
usage=None,
choices=[
SimpleNamespace(
index=0,
finish_reason=None,
delta=SimpleNamespace(content=text, tool_calls=None),
)
],
)
class _Stream:
"""Chunk iterator that records how far the consumer drained it."""
def __init__(self, count: int = 50) -> None:
self._count = count
self.yielded = 0
self.closed = False
def __iter__(self):
for _ in range(self._count):
self.yielded += 1
yield _chunk("x")
def close(self) -> None:
self.closed = True
class _AsyncStream(_Stream):
async def __aiter__(self): # pragma: no cover - exercised via asyncio.run
for _ in range(self._count):
self.yielded += 1
yield _chunk("x")
# ── The structural gap the bug lives in ──────────────────────────────────
def test_stream_ceiling_structurally_outlives_the_default_host_ceiling():
"""The worker's own budget is >= the host's for every configured timeout.
This is the arithmetic that guarantees the orphan: there is no aux timeout
for which ``_aux_stream_total_ceiling`` lands below the 600s default host
ceiling, and the reporter's ``auxiliary.compression.timeout: 600`` puts it
at 2400s — a 30-minute window in which an abandoned provider daemon keeps
streaming a summary nobody can commit.
"""
for aux_timeout in (None, 0, 30.0, 120.0, 300.0):
assert (
aux._aux_stream_total_ceiling(aux_timeout)
>= DEFAULT_CONTEXT_TOTAL_CEILING_SECONDS
)
assert aux._aux_stream_total_ceiling(600.0) == 2400.0
assert (
aux._aux_stream_total_ceiling(600.0)
- DEFAULT_CONTEXT_TOTAL_CEILING_SECONDS
== 1800.0
)
# ── The fence must publish the deadline it already owns ──────────────────
def test_commit_fence_publishes_its_shared_deadline():
fence = CompressionCommitFence()
assert fence.deadline_monotonic is None
fence.set_total_ceiling_seconds(600.0)
published = fence.deadline_monotonic
assert published is not None
assert 590.0 < published - time.monotonic() <= 600.0
assert not fence.deadline_exceeded
fence.set_total_ceiling_seconds(0.001)
time.sleep(0.01)
assert fence.deadline_exceeded
assert fence.deadline_monotonic <= time.monotonic()
# ── The stream consumer must honour it ───────────────────────────────────
def test_streamed_summary_stops_at_an_elapsed_host_deadline():
"""A host that already gave up must not leave the worker streaming on."""
stream = _Stream(count=50)
with aux.aux_stream_deadline(time.monotonic() - 1.0):
with pytest.raises(TimeoutError) as excinfo:
aux._aggregate_chat_stream(stream, model="m", total_ceiling=2400.0)
# "timed out" keeps _is_timeout_error classification identical to a
# request timeout, so the existing recovery chains are unchanged.
assert "timed out" in str(excinfo.value)
assert "host compression deadline" in str(excinfo.value)
# Stopped on the first frame instead of draining the whole stream, and the
# HTTP response was closed rather than left dangling.
assert stream.yielded == 1
assert stream.closed is True
def test_streamed_summary_runs_to_completion_under_a_live_host_deadline():
stream = _Stream(count=5)
with aux.aux_stream_deadline(time.monotonic() + 600.0):
response = aux._aggregate_chat_stream(
stream, model="m", total_ceiling=2400.0
)
assert response.choices[0].message.content == "xxxxx"
assert stream.yielded == 5
def test_no_host_deadline_keeps_the_historical_ceiling_behaviour():
"""Every non-compression aux caller must be byte-for-byte unchanged."""
stream = _Stream(count=5)
response = aux._aggregate_chat_stream(stream, model="m", total_ceiling=2400.0)
assert response.choices[0].message.content == "xxxxx"
assert stream.yielded == 5
# An installed-then-exited scope must not leak into the next call.
with aux.aux_stream_deadline(time.monotonic() - 1.0):
pass
stream2 = _Stream(count=3)
assert (
aux._aggregate_chat_stream(
stream2, model="m", total_ceiling=2400.0
).choices[0].message.content
== "xxx"
)
def test_none_deadline_is_a_no_op_passthrough():
"""Callers wire the scope unconditionally; a fenceless call must not break."""
stream = _Stream(count=3)
with aux.aux_stream_deadline(None):
response = aux._aggregate_chat_stream(
stream, model="m", total_ceiling=2400.0
)
assert response.choices[0].message.content == "xxx"
def test_nested_none_inherits_rather_than_escaping_the_host_deadline():
"""A fenceless aux call nested inside a fenced one stays bounded.
``None`` means "I have no deadline of my own", not "clear the one in
force" — mirroring ``_aux_thread_local_hook``'s passthrough contract. If it
cleared, any nested auxiliary call made during compression would escape the
host ceiling that the whole attempt is supposed to live inside.
"""
outer = time.monotonic() - 1.0
stream = _Stream(count=50)
with aux.aux_stream_deadline(outer):
with aux.aux_stream_deadline(None):
assert aux._current_aux_stream_deadline() == outer
with pytest.raises(TimeoutError):
aux._aggregate_chat_stream(stream, model="m", total_ceiling=2400.0)
assert stream.yielded == 1
def test_deadline_scope_restores_the_previous_value():
outer = time.monotonic() + 900.0
with aux.aux_stream_deadline(outer):
assert aux._current_aux_stream_deadline() == outer
with aux.aux_stream_deadline(time.monotonic() + 10.0):
assert aux._current_aux_stream_deadline() != outer
assert aux._current_aux_stream_deadline() == outer
assert aux._current_aux_stream_deadline() is None
def test_async_stream_mirror_honours_the_host_deadline():
"""The async consumer must not drift from the sync one."""
stream = _AsyncStream(count=50)
async def _run():
with aux.aux_stream_deadline(time.monotonic() - 1.0):
return await aux._aggregate_chat_stream_async(
stream, model="m", total_ceiling=2400.0
)
with pytest.raises(TimeoutError):
asyncio.run(_run())
assert stream.yielded == 1
# ── The isolated provider daemon must inherit it ─────────────────────────
def test_protected_provider_daemon_inherits_the_host_deadline():
"""``_run_protected_sync_provider_call`` runs the stream on ANOTHER thread.
Thread-locals do not cross that boundary, so without explicit propagation
the fix would be inert on exactly the path large-session compression takes
(protected + hard-cancel source installed).
"""
seen: dict[str, object] = {}
def _callback(_kwargs):
seen["deadline"] = aux._current_aux_stream_deadline()
seen["thread"] = threading.current_thread().name
return "ok"
deadline = time.monotonic() + 42.0
cancel_event = threading.Event()
with aux.aux_progress_hook(lambda: None), aux.aux_interrupt_protection(
cancel_event=cancel_event
), aux.aux_stream_deadline(deadline):
assert aux._run_protected_sync_provider_call(_callback, {}) == "ok"
assert seen["thread"] == "hermes-protected-aux-provider"
assert seen["deadline"] == deadline
# ── The compression worker must actually install it ──────────────────────
def _summary_dispatch_source() -> str:
from agent import conversation_compression
path = Path(inspect.getsourcefile(conversation_compression))
return path.read_text(encoding="utf-8")
def test_compression_summary_dispatch_installs_the_fence_deadline():
"""Source guard: the wiring is one line and trivially droppable.
A behavioural test would have to drive the whole ``compress_context`` body
(durable lock, watermark, telemetry, commit). This asserts the seam itself:
the same ``with`` statement that installs the progress hook must also
install the stream deadline.
"""
tree = ast.parse(_summary_dispatch_source())
wired = False
for node in ast.walk(tree):
if not isinstance(node, ast.With):
continue
names = set()
for item in node.items:
call = item.context_expr
if isinstance(call, ast.Call) and isinstance(call.func, ast.Name):
names.add(call.func.id)
if "aux_progress_hook" in names:
assert "aux_stream_deadline" in names, (
"the summary dispatch scope installs the progress hook but not "
"the host stream deadline — #99692 would regress"
)
wired = True
assert wired, "summary dispatch scope not found"