refactor: hoist preflight clear into restart handler, single-source qwen predicate
Simplify-pass follow-ups on the salvage stack (all guard tests re-run, mutation-checked): 1. conversation_loop.py: moved `_preflight_compression_blocked = False` from 9 per-site copies into the restart_with_rebuilt_messages handler (its single consumer). Besides removing the 9 duplicated blocks, this fixes a 10th pre-existing retry-loop site (content-filter stall failover, #32421) that set the flag and broke WITHOUT clearing the preflight block — a content-filter failover previously restarted with preflight compression still blocked against the fallback's smaller window, the same #84733 bug class. The outer-loop empty-response site keeps its own clear (it never passes through the handler). New AST guard test_restart_handler_clears_preflight_block pins the hoisted clear (mutation-checked). 2. agent_runtime_helpers.py: extracted _raw_cache_ttl_from_config() — prompt_caching_disabled_from_config and configured_cache_ttl were verbatim copies of the same config read. Added VALID_CACHE_TTLS. 3. prompt_caching.py: added is_qwen_model() next to ALIBABA_FAMILY_PROVIDERS; effective_cache_ttl and anthropic_prompt_cache_policy now share both the family set and the qwen predicate — neither can desync. 4. Guard-test hardening: assert every _try_activate_fallback reference is a direct `if agent._try_activate_fallback(...):` site, so a future `activated = ...` form can't silently escape the restart-discipline guard.
This commit is contained in:
@@ -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
|
||||
|
||||
+11
-45
@@ -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:
|
||||
|
||||
+11
-1
@@ -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"
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user