diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 35430b1c79..3256bca9ec 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -825,6 +825,7 @@ def run_compress_context_with_progress_timeout( on_timeout: Optional[Callable[[float, float, float], None]] = None, on_commit_overrun: Optional[Callable[[float, float], None]] = None, fence: Optional[CompressionCommitFence] = None, + telemetry_agent: Any = None, ) -> Tuple[list, str]: """Run ``worker(fence)`` under a sync progress-aware timeout. @@ -891,6 +892,17 @@ def run_compress_context_with_progress_timeout( "provider health.", _COMPRESS_EXECUTOR_MAX_WORKERS, ) + # Round-2 #6: saturation refusals must be visible in the same + # telemetry stream as every other failed attempt, or a wedged pool + # looks like compression simply stopped being attempted. + if telemetry_agent is not None: + _emit_compression_attempt_telemetry( + telemetry_agent, + started_at=time.monotonic(), + commit_status="aborted", + split_status="aborted", + failure_class="pool_saturated", + ) return messages, _resolve_fallback_prompt() def _fence_gated_worker(worker_fence: CompressionCommitFence): diff --git a/run_agent.py b/run_agent.py index eee84344ef..39a1689bad 100644 --- a/run_agent.py +++ b/run_agent.py @@ -7277,6 +7277,7 @@ class AIAgent: on_timeout=_on_timeout, on_commit_overrun=_on_commit_overrun, fence=active_fence, + telemetry_agent=self, ) # compress_context ran on a daemon pool worker thread; the session # id rotation updated hermes_logging._session_context (a diff --git a/tests/agent/test_compression_review_76354.py b/tests/agent/test_compression_review_76354.py index abeb9153cc..11609c4b78 100644 --- a/tests/agent/test_compression_review_76354.py +++ b/tests/agent/test_compression_review_76354.py @@ -318,14 +318,49 @@ class TestF6ExecutorSaturation: return ([], "5th") fifth_msgs = [{"role": "user", "content": "fifth"}] + # Round-2 #6: the fail-fast refusal must emit the standard + # compression-attempt telemetry with failure_class=pool_saturated. + class _TelemetryAgent: + session_id = "SATURATED_SESSION" + _compression_attempt_id = "sat-attempt" + + class context_compressor: # noqa: D106 — minimal stub + _last_compression_telemetry = None + _last_summary_fallback_used = False + _last_aux_model_failure_model = None + + import json as _json + import logging as _logging + + class _CaptureHandler(_logging.Handler): + def __init__(self): + super().__init__() + self.payloads = [] + + def emit(self, record): + msg = record.getMessage() + if "compression attempt telemetry" in msg: + self.payloads.append( + _json.loads(msg.split(": ", 1)[1]) + ) + + capture = _CaptureHandler() + _prev_level = cc.logger.level + cc.logger.addHandler(capture) + cc.logger.setLevel(_logging.DEBUG) t0 = time.monotonic() - msgs, prompt = run_compress_context_with_progress_timeout( - worker=fifth_worker, - messages=fifth_msgs, - system_prompt_fallback="fifth-fallback", - idle_timeout_seconds=5.0, - total_ceiling_seconds=5.0, - ) + try: + msgs, prompt = run_compress_context_with_progress_timeout( + worker=fifth_worker, + messages=fifth_msgs, + system_prompt_fallback="fifth-fallback", + idle_timeout_seconds=5.0, + total_ceiling_seconds=5.0, + telemetry_agent=_TelemetryAgent(), + ) + finally: + cc.logger.removeHandler(capture) + cc.logger.setLevel(_prev_level) elapsed = time.monotonic() - t0 # ── Assert while the 4 workers are STILL wedged ─────────────── assert not release.is_set() @@ -335,6 +370,16 @@ class TestF6ExecutorSaturation: assert msgs is fifth_msgs assert prompt == "fifth-fallback" assert not fifth_ran.is_set() + saturated = [ + p for p in capture.payloads + if p.get("failure_class") == "pool_saturated" + ] + assert saturated, ( + "fail-fast admission refusal must emit compression-attempt " + "telemetry with failure_class='pool_saturated'" + ) + assert saturated[0]["commit_status"] == "aborted" + assert saturated[0]["session_id"] == "SATURATED_SESSION" finally: release.set()