feat(agent): emit pool_saturated compression-attempt telemetry (re-review #6)
The fail-fast admission path (bounded compress pool, F6) only logged a WARNING; in the compression-attempt telemetry stream a wedged pool looked like compression simply stopped being attempted. Emit the existing attempt telemetry with failure_class='pool_saturated' (commit_status=aborted, split_status=aborted) on refusal, following _emit_compression_attempt_telemetry's existing call shape. Regression extends the F6 saturation test (sabotage-verified).
This commit is contained in:
@@ -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):
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user