diff --git a/agent/agent_runtime_helpers.py b/agent/agent_runtime_helpers.py index dfcf1806ab..a2248e9b63 100644 --- a/agent/agent_runtime_helpers.py +++ b/agent/agent_runtime_helpers.py @@ -1964,6 +1964,20 @@ def cache_ttl_means_disabled(ttl: Any) -> bool: return str(ttl).lower() in ("off", "false", "disabled", "no", "none") +# The two cache_ttl tiers accepted by config (anything else is either a +# disable synonym or ignored). Shared by the config readers below and +# mirrored by agent_init's live-agent snapshot. +VALID_CACHE_TTLS = ("5m", "1h") + + +def _raw_cache_ttl_from_config() -> Any: + """Read the raw ``prompt_caching.cache_ttl`` config value (may raise).""" + from hermes_cli.config import load_config_readonly + + pc_cfg = load_config_readonly().get("prompt_caching", {}) or {} + return pc_cfg.get("cache_ttl", "5m") + + def prompt_caching_disabled_from_config() -> bool: """Return True when ``prompt_caching.cache_ttl`` is configured as off. @@ -1973,10 +1987,7 @@ def prompt_caching_disabled_from_config() -> bool: ``AIAgent`` (#76085 / #33555). """ try: - from hermes_cli.config import load_config_readonly - - pc_cfg = load_config_readonly().get("prompt_caching", {}) or {} - ttl = pc_cfg.get("cache_ttl", "5m") + ttl = _raw_cache_ttl_from_config() except Exception: return False return cache_ttl_means_disabled(ttl) @@ -1992,13 +2003,10 @@ def configured_cache_ttl() -> Optional[str]: ``effective_cache_ttl`` resolves ``None`` to ``5m`` downstream. """ try: - from hermes_cli.config import load_config_readonly - - pc_cfg = load_config_readonly().get("prompt_caching", {}) or {} - ttl = pc_cfg.get("cache_ttl", "5m") + ttl = _raw_cache_ttl_from_config() except Exception: return None - return ttl if ttl in ("5m", "1h") else None + return ttl if ttl in VALID_CACHE_TTLS else None def blank_cache_policy_stub(cache_disabled: Optional[bool] = None): @@ -2292,12 +2300,12 @@ def anthropic_prompt_cache_policy( # OpenCode Zen's relay rejects the Anthropic-style content block # format that cache markers produce (content becomes a block array # instead of a plain string), causing HTTP 400 (#77217). - model_is_qwen = "qwen" in model_lower - # Single source of truth for the family set — shared with the - # effective_cache_ttl clamp so the opt-in and the TTL clamp can - # never desync (#84733). - from agent.prompt_caching import ALIBABA_FAMILY_PROVIDERS + # Single source of truth for the family set and the qwen-model + # predicate — shared with the effective_cache_ttl clamp so the + # opt-in and the TTL clamp can never desync (#84733). + from agent.prompt_caching import ALIBABA_FAMILY_PROVIDERS, is_qwen_model + model_is_qwen = is_qwen_model(model_lower) provider_is_alibaba_family = provider_lower in ALIBABA_FAMILY_PROVIDERS if provider_is_alibaba_family and model_is_qwen: # Envelope layout (native_anthropic=False): markers on inner diff --git a/agent/conversation_loop.py b/agent/conversation_loop.py index 28560ae62d..bc7dad72c3 100644 --- a/agent/conversation_loop.py +++ b/agent/conversation_loop.py @@ -2471,11 +2471,6 @@ def run_conversation( retry_count = 0 compression_attempts = 0 _retry.primary_recovery_attempted = False - # Failover shrank the compressor's context window to - # the fallback's; restart the outer iteration so the - # pre-API preflight re-runs against the new threshold - # before the first fallback call (#84733). - _preflight_compression_blocked = False _retry.restart_with_rebuilt_messages = True break # No fallback available — surface buffered context @@ -2929,11 +2924,6 @@ def run_conversation( retry_count = 0 compression_attempts = 0 _retry.primary_recovery_attempted = False - # Failover shrank the compressor's context window to - # the fallback's; restart the outer iteration so the - # pre-API preflight re-runs against the new threshold - # before the first fallback call (#84733). - _preflight_compression_blocked = False _retry.restart_with_rebuilt_messages = True break @@ -3008,11 +2998,6 @@ def run_conversation( retry_count = 0 compression_attempts = 0 _retry.primary_recovery_attempted = False - # Failover shrank the compressor's context window to - # the fallback's; restart the outer iteration so the - # pre-API preflight re-runs against the new threshold - # before the first fallback call (#84733). - _preflight_compression_blocked = False _retry.restart_with_rebuilt_messages = True break # Terminal — flush buffered retry trace so user sees what happened. @@ -3191,11 +3176,6 @@ def run_conversation( retry_count = 0 compression_attempts = 0 _retry.primary_recovery_attempted = False - # Failover shrank the compressor's context window to - # the fallback's; restart the outer iteration so the - # pre-API preflight re-runs against the new threshold - # before the first fallback call (#84733). - _preflight_compression_blocked = False _retry.restart_with_rebuilt_messages = True break @@ -4898,11 +4878,6 @@ def run_conversation( retry_count = 0 compression_attempts = 0 _retry.primary_recovery_attempted = False - # Failover shrank the compressor's context window to - # the fallback's; restart the outer iteration so the - # pre-API preflight re-runs against the new threshold - # before the first fallback call (#84733). - _preflight_compression_blocked = False _retry.restart_with_rebuilt_messages = True break @@ -4937,11 +4912,6 @@ def run_conversation( retry_count = 0 compression_attempts = 0 _retry.primary_recovery_attempted = False - # Failover shrank the compressor's context window to - # the fallback's; restart the outer iteration so the - # pre-API preflight re-runs against the new threshold - # before the first fallback call (#84733). - _preflight_compression_blocked = False _retry.restart_with_rebuilt_messages = True break @@ -5550,11 +5520,6 @@ def run_conversation( retry_count = 0 compression_attempts = 0 _retry.primary_recovery_attempted = False - # Failover shrank the compressor's context window to - # the fallback's; restart the outer iteration so the - # pre-API preflight re-runs against the new threshold - # before the first fallback call (#84733). - _preflight_compression_blocked = False _retry.restart_with_rebuilt_messages = True break if api_kwargs is not None: @@ -5779,11 +5744,6 @@ def run_conversation( retry_count = 0 compression_attempts = 0 _retry.primary_recovery_attempted = False - # Failover shrank the compressor's context window to - # the fallback's; restart the outer iteration so the - # pre-API preflight re-runs against the new threshold - # before the first fallback call (#84733). - _preflight_compression_blocked = False _retry.restart_with_rebuilt_messages = True break # Terminal — flush buffered retry/fallback trace. @@ -6101,14 +6061,20 @@ def run_conversation( continue if _retry.restart_with_rebuilt_messages: - # A content-filter stream stall (#32421) was escalated to the - # fallback chain and the partial content rolled back. Re-issue - # the API call against the now-active fallback provider. Refund - # the budget/count for the stalled attempt so the fallback gets a - # fair turn. + # A stream stall or provider failure was escalated to the + # fallback chain (10 activation sites in the retry loop set this + # flag and break here). Re-issue the API call against the + # now-active fallback provider. Refund the budget/count for the + # stalled attempt so the fallback gets a fair turn. api_call_count -= 1 agent.iteration_budget.refund() _retry.restart_with_rebuilt_messages = False + # Failover shrank the compressor's context window to the + # fallback's; clear the preflight block so the pre-API preflight + # re-runs against the new threshold before the first fallback + # call (#84733). Hoisted here (the single consumer) so every + # activation site — including ones added later — gets it. + _preflight_compression_blocked = False continue if _retry.restart_with_length_continuation: diff --git a/agent/prompt_caching.py b/agent/prompt_caching.py index 2ef50d5aca..8643ad533e 100644 --- a/agent/prompt_caching.py +++ b/agent/prompt_caching.py @@ -132,6 +132,16 @@ ALIBABA_FAMILY_PROVIDERS = frozenset({ }) +def is_qwen_model(model: str) -> bool: + """True when ``model`` names a Qwen-family model (case-insensitive). + + Shared by the TTL clamp below and + ``agent_runtime_helpers.anthropic_prompt_cache_policy`` so the + cache-policy opt-in and the clamp can never desync (#84733). + """ + return "qwen" in (model or "").lower() + + def effective_cache_ttl( ttl: str | None, *, @@ -150,7 +160,7 @@ def effective_cache_ttl( """ if ttl != "1h": return ttl or "5m" - if "qwen" in (model or "").lower(): + if is_qwen_model(model): return "5m" if (provider or "").lower() in ALIBABA_FAMILY_PROVIDERS: return "5m" diff --git a/tests/agent/test_prompt_cache_ttl_propagation.py b/tests/agent/test_prompt_cache_ttl_propagation.py index 683b7ce712..f2ab8d97f3 100644 --- a/tests/agent/test_prompt_cache_ttl_propagation.py +++ b/tests/agent/test_prompt_cache_ttl_propagation.py @@ -236,6 +236,21 @@ class TestFailoverRestartsPreflight: and node.test.func.attr == "_try_activate_fallback" ] assert fallback_ifs, "expected _try_activate_fallback sites in run_conversation" + # Every reference to _try_activate_fallback must be one of the matched + # `if agent._try_activate_fallback(...):` sites — a site written as + # `activated = agent._try_activate_fallback()` would silently escape + # this guard. + all_refs = [ + node + for node in ast.walk(tree) + if isinstance(node, ast.Attribute) + and node.attr == "_try_activate_fallback" + ] + assert len(all_refs) == len(fallback_ifs), ( + "every _try_activate_fallback reference must be a direct " + "`if agent._try_activate_fallback(...):` site so this guard " + "can bind its restart discipline (#84733)" + ) for node in fallback_ifs: if _inside_retry_loop(node): assert any(isinstance(stmt, ast.Break) for stmt in node.body), ( @@ -259,6 +274,51 @@ class TestFailoverRestartsPreflight: "exits the conversation loop and ends the turn (#84733)" ) + def test_restart_handler_clears_preflight_block(self): + """The single consumer of restart_with_rebuilt_messages must clear + _preflight_compression_blocked, so every retry-loop failover gets a + fresh preflight against the fallback's context window (#84733).""" + from agent import conversation_loop + + tree = ast.parse(inspect.getsource(conversation_loop.run_conversation)) + handlers = [ + node + for node in ast.walk(tree) + if isinstance(node, ast.If) + and isinstance(node.test, ast.Attribute) + and node.test.attr == "restart_with_rebuilt_messages" + ] + assert handlers, "expected the restart_with_rebuilt_messages handler" + consumer = [ + node + for node in handlers + if any( + isinstance(stmt, ast.Assign) + and any( + isinstance(t, ast.Attribute) + and t.attr == "restart_with_rebuilt_messages" + for t in stmt.targets + ) + for stmt in node.body + ) + ] + assert consumer, "expected the flag-consuming handler" + for node in consumer: + assert any( + isinstance(stmt, ast.Assign) + and any( + isinstance(t, ast.Name) + and t.id == "_preflight_compression_blocked" + for t in stmt.targets + ) + and isinstance(stmt.value, ast.Constant) + and stmt.value.value is False + for stmt in node.body + ), ( + "the restart handler must clear _preflight_compression_blocked " + "so the re-run preflight isn't skipped (#84733)" + ) + class TestAuxFallbackReplanThreadsTtl: """#84733 follow-up: the auxiliary fallback replan path threads the