refactor: reuse existing utilities in salvaged PR #72424
Three code-reuse fixes applied during salvage: 1. Reuse _relative_time from hermes_cli/main.py instead of duplicating the relative-time formatting logic in hermes_cli/status.py. 2. Extract _stamp_hygiene_compression_provenance helper in gateway/run.py to deduplicate the two nearly-identical try/except blocks that stamp compression timeout/abort provenance in the hygiene path. 3. Add ContextCompressor.record_timeout_failure() method and use it from the in-agent compress_context timeout callback instead of re-implementing the (60, 300, 900) cooldown ladder inline. The existing summary-LLM exception handler already has this ladder — now both paths share one method.
This commit is contained in:
@@ -1932,6 +1932,24 @@ class ContextCompressor(ContextEngine):
|
||||
self._cooldown_persist_failed = True
|
||||
logger.debug("compression failure cooldown persist failed (non-sqlite): %s", exc)
|
||||
|
||||
def record_timeout_failure(self, error: str) -> None:
|
||||
"""Record a consecutive timeout failure using the shared cooldown ladder.
|
||||
|
||||
Used by both the summary-LLM exception handler (inline at line ~3714)
|
||||
and the host-level ``compress_context`` timeout wrapper in
|
||||
``run_compress_context_with_progress_timeout``. Avoids re-implementing
|
||||
the ladder at each call site (#62452).
|
||||
"""
|
||||
_TIMEOUT_COOLDOWN_LADDER = (60, 300, 900)
|
||||
self._consecutive_timeout_failures = (
|
||||
getattr(self, "_consecutive_timeout_failures", 0) + 1
|
||||
)
|
||||
cooldown = _TIMEOUT_COOLDOWN_LADDER[
|
||||
min(self._consecutive_timeout_failures,
|
||||
len(_TIMEOUT_COOLDOWN_LADDER)) - 1
|
||||
]
|
||||
self._record_compression_failure_cooldown(float(cooldown), error)
|
||||
|
||||
def _clear_compression_failure_cooldown(self) -> None:
|
||||
self._summary_failure_cooldown_until = 0.0
|
||||
self._last_summary_error = None
|
||||
|
||||
+33
-34
@@ -889,6 +889,19 @@ def _float_env(name: str, default: float) -> float:
|
||||
return float(default)
|
||||
|
||||
|
||||
def _stamp_hygiene_compression_provenance(
|
||||
agent: Any,
|
||||
desc: str,
|
||||
provenance: "ActivityProvenance",
|
||||
debug_label: str,
|
||||
) -> None:
|
||||
"""Best-effort activity provenance stamp for hygiene compression transitions."""
|
||||
try:
|
||||
agent._touch_activity(desc, provenance=provenance)
|
||||
except Exception:
|
||||
logger.debug(debug_label, exc_info=True)
|
||||
|
||||
|
||||
def _is_fresh_gateway_interruption(
|
||||
value: Any,
|
||||
*,
|
||||
@@ -16682,23 +16695,16 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
|
||||
self, session_entry.session_id,
|
||||
_hyg_failure_cooldown_seconds,
|
||||
)
|
||||
try:
|
||||
from agent.session_activity import (
|
||||
ActivityProvenance,
|
||||
)
|
||||
|
||||
_hyg_agent._touch_activity(
|
||||
"session hygiene compression timed out",
|
||||
provenance=(
|
||||
ActivityProvenance.AGENT_COMPRESSION_TIMEOUT
|
||||
),
|
||||
)
|
||||
except Exception:
|
||||
logger.debug(
|
||||
"hygiene compression timeout "
|
||||
"activity stamp failed",
|
||||
exc_info=True,
|
||||
)
|
||||
from agent.session_activity import (
|
||||
ActivityProvenance,
|
||||
)
|
||||
_stamp_hygiene_compression_provenance(
|
||||
_hyg_agent,
|
||||
"session hygiene compression timed out",
|
||||
ActivityProvenance.AGENT_COMPRESSION_TIMEOUT,
|
||||
"hygiene compression timeout "
|
||||
"activity stamp failed",
|
||||
)
|
||||
logger.warning(
|
||||
"Session hygiene compression for session %s "
|
||||
"made no progress for %.1fs "
|
||||
@@ -16866,23 +16872,16 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
|
||||
self, session_entry.session_id,
|
||||
_hyg_failure_cooldown_seconds,
|
||||
)
|
||||
try:
|
||||
from agent.session_activity import (
|
||||
ActivityProvenance,
|
||||
)
|
||||
|
||||
_hyg_agent._touch_activity(
|
||||
"session hygiene compression aborted",
|
||||
provenance=(
|
||||
ActivityProvenance.AGENT_COMPRESSION_COOLDOWN
|
||||
),
|
||||
)
|
||||
except Exception:
|
||||
logger.debug(
|
||||
"hygiene compression abort "
|
||||
"activity stamp failed",
|
||||
exc_info=True,
|
||||
)
|
||||
from agent.session_activity import (
|
||||
ActivityProvenance,
|
||||
)
|
||||
_stamp_hygiene_compression_provenance(
|
||||
_hyg_agent,
|
||||
"session hygiene compression aborted",
|
||||
ActivityProvenance.AGENT_COMPRESSION_COOLDOWN,
|
||||
"hygiene compression abort "
|
||||
"activity stamp failed",
|
||||
)
|
||||
_err = getattr(_comp, "_last_summary_error", None) or "unknown error"
|
||||
# Force-redact: provider exception text
|
||||
# may contain credentials; this message
|
||||
|
||||
+2
-16
@@ -65,23 +65,9 @@ def _format_iso_timestamp(value) -> str:
|
||||
|
||||
def _format_relative_ts(ts: float) -> str:
|
||||
"""Format an epoch timestamp as a short relative age for status output."""
|
||||
if not ts:
|
||||
return "?"
|
||||
import time as _time
|
||||
from datetime import datetime
|
||||
from hermes_cli.main import _relative_time
|
||||
|
||||
delta = _time.time() - float(ts)
|
||||
if delta < 60:
|
||||
return "just now"
|
||||
if delta < 3600:
|
||||
return f"{int(delta / 60)}m ago"
|
||||
if delta < 86400:
|
||||
return f"{int(delta / 3600)}h ago"
|
||||
if delta < 172800:
|
||||
return "yesterday"
|
||||
if delta < 604800:
|
||||
return f"{int(delta / 86400)}d ago"
|
||||
return datetime.fromtimestamp(ts).strftime("%Y-%m-%d")
|
||||
return _relative_time(ts)
|
||||
|
||||
|
||||
def _configured_model_label(config: dict) -> str:
|
||||
|
||||
+2
-11
@@ -7169,21 +7169,12 @@ class AIAgent:
|
||||
# (#62452): avoid re-burning the full idle budget every turn.
|
||||
compressor = getattr(self, "context_compressor", None)
|
||||
if compressor is not None:
|
||||
streak = (
|
||||
getattr(compressor, "_consecutive_timeout_failures", 0) + 1
|
||||
)
|
||||
compressor._consecutive_timeout_failures = streak
|
||||
ladder = (60, 300, 900)
|
||||
cooldown = ladder[min(streak, len(ladder)) - 1]
|
||||
record = getattr(
|
||||
compressor, "_record_compression_failure_cooldown", None
|
||||
)
|
||||
record = getattr(compressor, "record_timeout_failure", None)
|
||||
if callable(record):
|
||||
try:
|
||||
record(
|
||||
float(cooldown),
|
||||
"host compress_context timeout "
|
||||
"(no summary progress)",
|
||||
"(no summary progress)"
|
||||
)
|
||||
except Exception:
|
||||
logger.debug(
|
||||
|
||||
@@ -243,6 +243,14 @@ class TestCompressContextForwarderOwnsTimeout:
|
||||
agent._conversation_root_id = MagicMock(return_value=None)
|
||||
agent.context_compressor = MagicMock()
|
||||
agent.context_compressor._consecutive_timeout_failures = 0
|
||||
# Use the real record_timeout_failure method so the cooldown ladder
|
||||
# is exercised end-to-end (not auto-mocked by MagicMock).
|
||||
from agent.context_compressor import ContextCompressor
|
||||
agent.context_compressor.record_timeout_failure = (
|
||||
ContextCompressor.record_timeout_failure.__get__(
|
||||
agent.context_compressor, MagicMock
|
||||
)
|
||||
)
|
||||
agent.context_compressor._record_compression_failure_cooldown = MagicMock()
|
||||
|
||||
hang = threading.Event()
|
||||
|
||||
Reference in New Issue
Block a user