refactor(turn): split recover_from_overflow into per-error handlers on a _Recovery state (501-LOC function -> max 76)
This commit is contained in:
@@ -5,7 +5,7 @@ Extracted from ``run_conversation``'s ``except`` branch. Each path either compre
|
||||
signals a restart, defers softly (compression lock / transient block), or ends the turn
|
||||
with a typed result. Nothing here imports ``agent.conversation_loop`` at module level
|
||||
(cycle); loop-internal helpers and the token estimators that tests patch on the loop
|
||||
module are imported lazily inside the handler so they keep resolving through the loop.
|
||||
module are imported lazily inside the handlers so they keep resolving through the loop.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -35,6 +35,16 @@ from utils import base_url_host_matches
|
||||
|
||||
logger = logging.getLogger("agent.conversation_loop")
|
||||
|
||||
_RETRY_HINT = " 💡 Try /new to start a fresh conversation, or /compress to retry compression."
|
||||
|
||||
_GITHUB_MODELS_HINT = (
|
||||
" 💡 GitHub Models free tier (models.inference.ai.azure.com) caps every",
|
||||
" request at ~8K tokens. Hermes' system prompt + tool schemas baseline",
|
||||
" exceeds that floor, so this endpoint cannot run an agentic loop.",
|
||||
" Use the `copilot` provider with a Copilot subscription token (`hermes",
|
||||
" setup` → GitHub Copilot), or pick any other provider.",
|
||||
)
|
||||
|
||||
|
||||
@dataclass
|
||||
class OverflowVerdict:
|
||||
@@ -57,6 +67,376 @@ class OverflowVerdict:
|
||||
is_context_length_error: bool
|
||||
|
||||
|
||||
@dataclass
|
||||
class _Recovery:
|
||||
"""Mutable working state shared by the overflow sub-handlers; the loop locals
|
||||
they may rebind live here and are handed back through ``verdict()``."""
|
||||
|
||||
agent: Any
|
||||
api_messages: Any
|
||||
system_message: Any
|
||||
effective_task_id: Any
|
||||
api_call_count: int
|
||||
max_compression_attempts: int
|
||||
messages: List[Dict[str, Any]]
|
||||
active_system_prompt: Any
|
||||
conversation_history: Any
|
||||
approx_tokens: int
|
||||
compression_attempts: int
|
||||
provider_overflow_recovery_pending: bool = False
|
||||
is_context_length_error: bool = False
|
||||
|
||||
def verdict(self, action: str, result: Optional[Dict[str, Any]] = None) -> OverflowVerdict:
|
||||
return OverflowVerdict(
|
||||
action=action,
|
||||
result=result,
|
||||
messages=self.messages,
|
||||
active_system_prompt=self.active_system_prompt,
|
||||
conversation_history=self.conversation_history,
|
||||
approx_tokens=self.approx_tokens,
|
||||
compression_attempts=self.compression_attempts,
|
||||
provider_overflow_recovery_pending=self.provider_overflow_recovery_pending,
|
||||
is_context_length_error=self.is_context_length_error,
|
||||
)
|
||||
|
||||
def fail_turn(
|
||||
self,
|
||||
final_response: str,
|
||||
*,
|
||||
notices: tuple = (),
|
||||
log: Optional[tuple] = None,
|
||||
compression_exhausted: bool = True,
|
||||
**extra: Any,
|
||||
) -> OverflowVerdict:
|
||||
"""End the turn as failed/partial. ``notices`` flush the buffered retry trace
|
||||
first so the user sees what compression attempts were made."""
|
||||
agent = self.agent
|
||||
if notices:
|
||||
agent._flush_status_buffer()
|
||||
for line in notices:
|
||||
agent._vprint(f"{agent.log_prefix}{line}", force=True)
|
||||
if log:
|
||||
logger.error(*log)
|
||||
agent._persist_session(self.messages, self.conversation_history)
|
||||
result = {
|
||||
"final_response": final_response,
|
||||
"messages": self.messages,
|
||||
"completed": False,
|
||||
"api_calls": self.api_call_count,
|
||||
"error": final_response,
|
||||
"partial": True,
|
||||
"failed": True,
|
||||
}
|
||||
if compression_exhausted:
|
||||
result["compression_exhausted"] = True
|
||||
result.update(extra)
|
||||
return self.verdict("return", result)
|
||||
|
||||
def exhausted(self, *, payload_too_large: bool = False) -> OverflowVerdict:
|
||||
"""Terminal: ``compression_attempts`` exceeded ``max_compression_attempts``."""
|
||||
cap = self.max_compression_attempts
|
||||
if payload_too_large:
|
||||
return self.fail_turn(
|
||||
f"Request payload too large: max compression attempts ({cap}) reached.",
|
||||
notices=(
|
||||
f"❌ Max compression attempts ({cap}) reached for payload-too-large error.",
|
||||
_RETRY_HINT,
|
||||
),
|
||||
log=("%s413 compression failed after %d attempts.", self.agent.log_prefix, cap),
|
||||
)
|
||||
return self.fail_turn(
|
||||
f"Context length exceeded: max compression attempts ({cap}) reached.",
|
||||
notices=(f"❌ Max compression attempts ({cap}) reached.", _RETRY_HINT),
|
||||
log=("%sContext compression failed after %d attempts.", self.agent.log_prefix, cap),
|
||||
)
|
||||
|
||||
def compress(self, request_tokens: int, *, fail_on_timeout: bool = False) -> Optional[OverflowVerdict]:
|
||||
"""One compression pass with the summary-failure cooldown bypassed (the
|
||||
provider proved the request doesn't fit, #100661). Returns ``None`` when
|
||||
history was compressed, or a soft-defer verdict when another path holds the
|
||||
compression lock (#69870) or a timed guard no-oped the pass (#97488): the
|
||||
attempt is refunded and the turn ends WITHOUT ``compression_exhausted`` so the
|
||||
gateway does not auto-reset. With ``fail_on_timeout`` a host timeout (recovery
|
||||
spent its wait budget with no committed summary) ends the turn via the typed
|
||||
contract, since re-sending would hit the same overflow (#98722)."""
|
||||
from agent.conversation_loop import (
|
||||
_COMPRESSION_TIMEOUT_FINAL_RESPONSE,
|
||||
_compression_deferred_result,
|
||||
conversation_history_after_compression,
|
||||
)
|
||||
|
||||
agent = self.agent
|
||||
before = self.messages
|
||||
self.messages, self.active_system_prompt = agent._compress_context(
|
||||
before, self.system_message,
|
||||
approx_tokens=request_tokens,
|
||||
task_id=self.effective_task_id,
|
||||
bypass_cooldown=True,
|
||||
)
|
||||
if self.messages is before:
|
||||
deferred = None
|
||||
if compression_skipped_due_to_lock(agent):
|
||||
deferred = _compression_deferred_result(agent, self.messages, self.api_call_count)
|
||||
elif compression_blocked_transiently(agent):
|
||||
deferred = _compression_deferred_result(
|
||||
agent, self.messages, self.api_call_count, reason="transient_block",
|
||||
)
|
||||
if deferred is not None:
|
||||
self.compression_attempts -= 1
|
||||
agent._persist_session(self.messages, self.conversation_history)
|
||||
return self.verdict("return", deferred)
|
||||
if fail_on_timeout and context_compression_timed_out(agent):
|
||||
return self.fail_turn(
|
||||
_COMPRESSION_TIMEOUT_FINAL_RESPONSE,
|
||||
turn_exit_reason="context_compression_timeout",
|
||||
)
|
||||
self.conversation_history = conversation_history_after_compression(
|
||||
agent, self.messages, self.conversation_history
|
||||
)
|
||||
return None
|
||||
|
||||
def request_tokens(self) -> int:
|
||||
"""Overhead-aware request size (msgs + tools + system) so LCM forced-overflow
|
||||
recovery arms on the TRUE request, not the tool-blind message count."""
|
||||
from agent.conversation_loop import estimate_request_tokens_rough
|
||||
|
||||
return estimate_request_tokens_rough(self.api_messages, tools=self.agent.tools or None)
|
||||
|
||||
|
||||
def _recover_payload_too_large(st: _Recovery, _retry: TurnRetryState) -> OverflowVerdict:
|
||||
"""413: compress and retry. A 413 is a BYTE-size error, so progress is scored in
|
||||
payload bytes — never the token estimate, which is deliberately byte-blind to images
|
||||
and wedged sessions on "no progress" (#88960 / #47339)."""
|
||||
from agent.conversation_loop import estimate_messages_tokens_rough
|
||||
|
||||
agent = st.agent
|
||||
st.compression_attempts += 1
|
||||
if st.compression_attempts > st.max_compression_attempts:
|
||||
return st.exhausted(payload_too_large=True)
|
||||
agent._buffer_status(
|
||||
f"⚠️ Request payload too large (413) — compression attempt "
|
||||
f"{st.compression_attempts}/{st.max_compression_attempts}..."
|
||||
)
|
||||
|
||||
messages = st.messages
|
||||
original_len = len(messages)
|
||||
original_bytes = serialized_messages_bytes(messages)
|
||||
deferred = st.compress(st.request_tokens())
|
||||
if deferred is not None:
|
||||
return deferred
|
||||
|
||||
# Re-measure: same-count compression and media aging can shrink the request
|
||||
# without shrinking the array. Tokens only for status display.
|
||||
messages = st.messages
|
||||
st.approx_tokens = estimate_messages_tokens_rough(messages)
|
||||
new_bytes = serialized_messages_bytes(messages)
|
||||
if len(messages) < original_len or (new_bytes > 0 and new_bytes < original_bytes * 0.95):
|
||||
if len(messages) < original_len:
|
||||
agent._buffer_status(COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE.format(before=original_len, after=len(messages)))
|
||||
else:
|
||||
agent._buffer_status(
|
||||
f"🗜️ Compressed {original_bytes:,} → {new_bytes:,} "
|
||||
f"payload bytes, retrying..."
|
||||
)
|
||||
time.sleep(2) # Brief pause between compression retries
|
||||
_retry.restart_with_compressed_messages = True
|
||||
return st.verdict("break")
|
||||
|
||||
if agent._try_strip_image_parts_from_tool_messages(st.api_messages, remember_model=False):
|
||||
agent._buffer_status(
|
||||
"📐 Compression could not reduce the request further — "
|
||||
"removed retained vision payloads and retrying..."
|
||||
)
|
||||
return st.verdict("continue")
|
||||
|
||||
return st.fail_turn(
|
||||
"Request payload too large (413). Cannot compress further.",
|
||||
notices=("❌ Payload too large and cannot compress further.", _RETRY_HINT),
|
||||
log=("%s413 payload too large. Cannot compress further.", agent.log_prefix),
|
||||
)
|
||||
|
||||
|
||||
def _clamp_output_cap(st: _Recovery, _retry: TurnRetryState, available_out: int, old_ctx: int) -> OverflowVerdict:
|
||||
"""Output-cap error ("max_tokens too large": input fits but input + max_tokens >
|
||||
window). The provider's available_tokens is the authoritative bound; also estimate
|
||||
the real request shape (API-only content) and use the smaller minus a margin."""
|
||||
from agent.conversation_loop import estimate_messages_tokens_rough
|
||||
|
||||
agent = st.agent
|
||||
request_input_estimate = st.request_tokens()
|
||||
local_available_out = old_ctx - request_input_estimate
|
||||
if local_available_out > 0:
|
||||
safe_out = max(1, min(available_out, local_available_out) - 64)
|
||||
else:
|
||||
# Local estimate can overshoot; fall back to the provider-reported budget.
|
||||
safe_out = max(1, available_out - 64)
|
||||
agent._ephemeral_max_output_tokens = safe_out
|
||||
agent._buffer_vprint(
|
||||
f"⚠️ Output cap too large for current prompt — "
|
||||
f"retrying with max_tokens={safe_out:,} "
|
||||
f"(provider_available={available_out:,}, "
|
||||
f"estimated_request_tokens={request_input_estimate:,}; "
|
||||
f"context_length unchanged at {old_ctx:,})"
|
||||
)
|
||||
# Still count against compression_attempts so a recurring error can't loop forever.
|
||||
st.compression_attempts += 1
|
||||
if st.compression_attempts > st.max_compression_attempts:
|
||||
return st.exhausted()
|
||||
# Also compress history so the retry doesn't spin on max_tokens alone; dropping the
|
||||
# middle window makes the total fit (#55546). Compression must never turn an
|
||||
# output-cap error fatal — on error, fall through and retry on max_tokens alone.
|
||||
try:
|
||||
original_len = len(st.messages)
|
||||
original_tokens = estimate_messages_tokens_rough(st.messages)
|
||||
deferred = st.compress(request_input_estimate)
|
||||
if deferred is not None:
|
||||
return deferred
|
||||
messages = st.messages
|
||||
new_tokens = estimate_messages_tokens_rough(messages)
|
||||
if len(messages) < original_len:
|
||||
agent._buffer_status(COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE.format(before=original_len, after=len(messages)))
|
||||
elif new_tokens > 0 and new_tokens < original_tokens * 0.95:
|
||||
agent._buffer_status(COMPRESSION_RETRY_TOKENS_STATUS_TEMPLATE.format(before=original_tokens, after=new_tokens))
|
||||
except Exception:
|
||||
logger.warning(
|
||||
"%sOutput-cap compression hit an error; retrying on max_tokens only.",
|
||||
agent.log_prefix,
|
||||
)
|
||||
_retry.restart_with_compressed_messages = True
|
||||
return st.verdict("break")
|
||||
|
||||
|
||||
def _adopt_provider_context_limit(st: _Recovery, error_msg: str, old_ctx: int) -> Optional[int]:
|
||||
"""Shrink context_length only when the provider reports the real limit; else keep
|
||||
the window and compress. Guessed probe tiers can turn a configured 1M window into
|
||||
256K/128K/64K. Returns the provider-reported limit, or ``None``."""
|
||||
from agent.conversation_loop import save_context_length
|
||||
|
||||
agent = st.agent
|
||||
compressor = agent.context_compressor
|
||||
new_ctx = get_context_length_from_provider_error(error_msg, old_ctx)
|
||||
if new_ctx is not None:
|
||||
agent._buffer_vprint(f"Context limit detected from API: {new_ctx:,} tokens (was {old_ctx:,})")
|
||||
compressor.update_model(
|
||||
model=agent.model,
|
||||
context_length=new_ctx,
|
||||
base_url=agent.base_url,
|
||||
api_key=getattr(agent, "api_key", ""),
|
||||
provider=agent.provider,
|
||||
api_mode=agent.api_mode,
|
||||
)
|
||||
# Persist the provider-reported limit BEFORE compression/retry: rate limit,
|
||||
# missing usage, or restart must not lose confirmed metadata. Probe flags
|
||||
# remain a fallback if this write fails.
|
||||
save_context_length(agent.model, agent.base_url, new_ctx)
|
||||
# Probe flags only on the built-in compressor (plugin engines manage their
|
||||
# own); provider-sourced value, so safe to cache.
|
||||
if hasattr(compressor, "_context_probed"):
|
||||
compressor._context_probed = True
|
||||
compressor._context_probe_persistable = True
|
||||
agent._buffer_vprint(f"⚠️ Context length exceeded — using provider limit: {old_ctx:,} → {new_ctx:,} tokens")
|
||||
return new_ctx
|
||||
|
||||
_provider_lower = (getattr(agent, "provider", "") or "").lower()
|
||||
_base_lower = (getattr(agent, "base_url", "") or "").rstrip("/").lower()
|
||||
is_minimax_provider = (
|
||||
_provider_lower in {"minimax", "minimax-cn"}
|
||||
or _base_lower.startswith((
|
||||
"https://api.minimax.io/anthropic",
|
||||
"https://api.minimaxi.com/anthropic",
|
||||
))
|
||||
)
|
||||
if is_minimax_provider and "context window exceeds limit (" in error_msg:
|
||||
agent._buffer_vprint(
|
||||
f"Provider reported overflow amount only; "
|
||||
f"keeping context_length at {old_ctx:,} tokens and compressing."
|
||||
)
|
||||
else:
|
||||
agent._buffer_vprint(
|
||||
f"⚠️ Context length exceeded, but provider did not report a max context length; "
|
||||
f"keeping context_length at {old_ctx:,} tokens and compressing."
|
||||
)
|
||||
return None
|
||||
|
||||
|
||||
def _recover_context_length(st: _Recovery, _retry: TurnRetryState, error_msg: str) -> OverflowVerdict:
|
||||
"""Context-length error. Two shapes: "prompt too long" = INPUT overflows the window
|
||||
(shrink context_length + compress); "max_tokens too large" = input fits but
|
||||
input + max_tokens > window (shrink the OUTPUT cap only)."""
|
||||
from agent.conversation_loop import estimate_messages_tokens_rough
|
||||
|
||||
agent = st.agent
|
||||
old_ctx = agent.context_compressor.context_length
|
||||
|
||||
available_out = parse_available_output_tokens_from_error(error_msg)
|
||||
if available_out is not None:
|
||||
return _clamp_output_cap(st, _retry, available_out, old_ctx)
|
||||
|
||||
# Output-cap error with unparseable budget: compression can't help (input already
|
||||
# fits) and would death-loop on the same 400. Fail fast (#55546).
|
||||
if is_output_cap_error(error_msg):
|
||||
return st.fail_turn(
|
||||
"max_tokens exceeds the provider's output cap for this model. "
|
||||
"Lower model.max_tokens in config.yaml.",
|
||||
notices=(
|
||||
"❌ The provider rejected the request because "
|
||||
"max_tokens exceeds its output cap for this model.",
|
||||
" 💡 Lower model.max_tokens in your config.yaml to "
|
||||
"at or below the model's max-output limit. "
|
||||
"(This is an output-cap error, not a context overflow — "
|
||||
"compression cannot fix it.)",
|
||||
),
|
||||
log=(
|
||||
f"{agent.log_prefix}Output-cap error not routed into compression "
|
||||
f"(max_tokens over provider cap): {error_msg[:200]}",
|
||||
),
|
||||
compression_exhausted=False,
|
||||
)
|
||||
|
||||
new_ctx = _adopt_provider_context_limit(st, error_msg, old_ctx)
|
||||
|
||||
st.compression_attempts += 1
|
||||
if st.compression_attempts > st.max_compression_attempts:
|
||||
return st.exhausted()
|
||||
agent._buffer_status(COMPRESSION_RETRY_TOO_LARGE_STATUS_TEMPLATE.format(
|
||||
tokens=st.approx_tokens, attempt=st.compression_attempts, cap=st.max_compression_attempts,
|
||||
))
|
||||
|
||||
original_len = len(st.messages)
|
||||
original_tokens = estimate_messages_tokens_rough(st.messages)
|
||||
deferred = st.compress(st.request_tokens(), fail_on_timeout=True)
|
||||
if deferred is not None:
|
||||
return deferred
|
||||
|
||||
# Re-estimate: same-message-count compression (tool-result pruning, in-place
|
||||
# summarization) can shrink the request (#39550).
|
||||
messages = st.messages
|
||||
new_tokens = estimate_messages_tokens_rough(messages)
|
||||
st.approx_tokens = new_tokens
|
||||
shrank_tokens = new_tokens > 0 and new_tokens < original_tokens * 0.95
|
||||
if len(messages) < original_len or shrank_tokens or (new_ctx and new_ctx < old_ctx):
|
||||
if len(messages) < original_len:
|
||||
agent._buffer_status(COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE.format(before=original_len, after=len(messages)))
|
||||
elif shrank_tokens:
|
||||
agent._buffer_status(COMPRESSION_RETRY_TOKENS_STATUS_TEMPLATE.format(before=original_tokens, after=new_tokens))
|
||||
time.sleep(2) # Brief pause between compression retries
|
||||
# Rebuild the full request and force normal preflight to honor it; message
|
||||
# count alone doesn't prove system/tool-inclusive pressure fell.
|
||||
st.provider_overflow_recovery_pending = True
|
||||
_retry.restart_with_compressed_messages = True
|
||||
return st.verdict("break")
|
||||
|
||||
# Can't compress further and already at minimum tier.
|
||||
return st.fail_turn(
|
||||
f"Context length exceeded ({new_tokens:,} tokens). Cannot compress further.",
|
||||
notices=(
|
||||
"❌ Context length exceeded and cannot compress further.",
|
||||
" 💡 The conversation has accumulated too much content. Try /new to start fresh, or /compress to manually trigger compression.",
|
||||
),
|
||||
log=("%sContext length exceeded: %s tokens. Cannot compress further.", agent.log_prefix, f"{new_tokens:,}"),
|
||||
)
|
||||
|
||||
|
||||
def recover_from_overflow(
|
||||
agent: Any,
|
||||
api_error: Exception,
|
||||
@@ -83,478 +463,40 @@ def recover_from_overflow(
|
||||
errors (incl. relay-wrapped output-cap 429s) BEFORE non-retryable client errors.
|
||||
Compression progress is scored in payload BYTES for 413 (never the byte-blind token
|
||||
estimate) and in tokens/message count for context overflow."""
|
||||
# Token estimators + loop-internal helpers resolve through the loop module so
|
||||
# existing ``patch("agent.conversation_loop.X")`` mocks keep intercepting.
|
||||
from agent.conversation_loop import (
|
||||
_COMPRESSION_TIMEOUT_FINAL_RESPONSE,
|
||||
_compression_deferred_result,
|
||||
conversation_history_after_compression,
|
||||
estimate_messages_tokens_rough,
|
||||
estimate_request_tokens_rough,
|
||||
save_context_length,
|
||||
st = _Recovery(
|
||||
agent=agent,
|
||||
api_messages=api_messages,
|
||||
system_message=system_message,
|
||||
effective_task_id=effective_task_id,
|
||||
api_call_count=api_call_count,
|
||||
max_compression_attempts=max_compression_attempts,
|
||||
messages=messages,
|
||||
active_system_prompt=active_system_prompt,
|
||||
conversation_history=conversation_history,
|
||||
approx_tokens=approx_tokens,
|
||||
compression_attempts=compression_attempts,
|
||||
)
|
||||
|
||||
_provider_overflow_recovery_pending = False
|
||||
is_context_length_error = False
|
||||
_wrapped_output_cap_budget = wrapped_output_cap_budget
|
||||
|
||||
def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> OverflowVerdict:
|
||||
return OverflowVerdict(
|
||||
action=action,
|
||||
result=result,
|
||||
messages=messages,
|
||||
active_system_prompt=active_system_prompt,
|
||||
conversation_history=conversation_history,
|
||||
approx_tokens=approx_tokens,
|
||||
compression_attempts=compression_attempts,
|
||||
provider_overflow_recovery_pending=_provider_overflow_recovery_pending,
|
||||
is_context_length_error=is_context_length_error,
|
||||
)
|
||||
|
||||
is_payload_too_large = (
|
||||
classified.reason == FailoverReason.payload_too_large
|
||||
)
|
||||
|
||||
# GitHub Models free tier caps requests at 8K tokens, under the system
|
||||
# prompt + tool schema floor; compression can't help, so say so.
|
||||
# GitHub Models free tier caps requests at 8K tokens, under the system prompt +
|
||||
# tool schema floor; compression can't help, so say so.
|
||||
if (
|
||||
status_code == 413
|
||||
and isinstance(agent.base_url, str)
|
||||
and base_url_host_matches(agent.base_url, "models.inference.ai.azure.com")
|
||||
):
|
||||
agent._vprint(
|
||||
f"{agent.log_prefix} 💡 GitHub Models free tier (models.inference.ai.azure.com) caps every",
|
||||
force=True,
|
||||
)
|
||||
agent._vprint(
|
||||
f"{agent.log_prefix} request at ~8K tokens. Hermes' system prompt + tool schemas baseline",
|
||||
force=True,
|
||||
)
|
||||
agent._vprint(
|
||||
f"{agent.log_prefix} exceeds that floor, so this endpoint cannot run an agentic loop.",
|
||||
force=True,
|
||||
)
|
||||
agent._vprint(
|
||||
f"{agent.log_prefix} Use the `copilot` provider with a Copilot subscription token (`hermes",
|
||||
force=True,
|
||||
)
|
||||
agent._vprint(
|
||||
f"{agent.log_prefix} setup` → GitHub Copilot), or pick any other provider.",
|
||||
force=True,
|
||||
)
|
||||
for line in _GITHUB_MODELS_HINT:
|
||||
agent._vprint(f"{agent.log_prefix}{line}", force=True)
|
||||
|
||||
if is_payload_too_large:
|
||||
compression_attempts += 1
|
||||
if compression_attempts > max_compression_attempts:
|
||||
# Terminal — surface the buffered retry trace.
|
||||
agent._flush_status_buffer()
|
||||
agent._vprint(f"{agent.log_prefix}❌ Max compression attempts ({max_compression_attempts}) reached for payload-too-large error.", force=True)
|
||||
agent._vprint(f"{agent.log_prefix} 💡 Try /new to start a fresh conversation, or /compress to retry compression.", force=True)
|
||||
logger.error("%s413 compression failed after %d attempts.", agent.log_prefix, max_compression_attempts)
|
||||
agent._persist_session(messages, conversation_history)
|
||||
_final_response = f"Request payload too large: max compression attempts ({max_compression_attempts}) reached."
|
||||
return _verdict("return", {
|
||||
"final_response": _final_response,
|
||||
"messages": messages,
|
||||
"completed": False,
|
||||
"api_calls": api_call_count,
|
||||
"error": _final_response,
|
||||
"partial": True,
|
||||
"failed": True,
|
||||
"compression_exhausted": True,
|
||||
})
|
||||
agent._buffer_status(f"⚠️ Request payload too large (413) — compression attempt {compression_attempts}/{max_compression_attempts}...")
|
||||
if classified.reason == FailoverReason.payload_too_large:
|
||||
return _recover_payload_too_large(st, _retry)
|
||||
|
||||
original_len = len(messages)
|
||||
# A 413 is a BYTE-size error: score progress in payload bytes,
|
||||
# never the token estimate, which is deliberately byte-blind to
|
||||
# images and wedged sessions on "no progress" (#88960 / #47339).
|
||||
original_bytes = serialized_messages_bytes(messages)
|
||||
_overflow_input = messages
|
||||
# Option A (LCM issue 441): overhead-aware request size so recovery arms on the
|
||||
# true request (msgs + tools + system), not the tool-blind message count.
|
||||
messages, active_system_prompt = agent._compress_context(
|
||||
messages, system_message,
|
||||
approx_tokens=estimate_request_tokens_rough(api_messages, tools=agent.tools or None),
|
||||
task_id=effective_task_id,
|
||||
# Provider proved the request doesn't fit: ignore the
|
||||
# summary-failure cooldown for this ONE attempt (#100661).
|
||||
bypass_cooldown=True,
|
||||
)
|
||||
if messages is _overflow_input and compression_skipped_due_to_lock(agent):
|
||||
# Lock-skip: another path holds the compression lock. A
|
||||
# temporary defer, not exhaustion — refund the attempt and
|
||||
# end softly so the gateway does NOT auto-reset (#69870).
|
||||
compression_attempts -= 1
|
||||
agent._persist_session(messages, conversation_history)
|
||||
return _verdict("return", _compression_deferred_result(
|
||||
agent, messages, api_call_count
|
||||
))
|
||||
if messages is _overflow_input and compression_blocked_transiently(agent):
|
||||
# Transient-block: a timed guard no-oped compression. A
|
||||
# defer, never compression_exhausted (auto-reset) (#97488).
|
||||
compression_attempts -= 1
|
||||
agent._persist_session(messages, conversation_history)
|
||||
return _verdict("return", _compression_deferred_result(
|
||||
agent, messages, api_call_count,
|
||||
reason="transient_block",
|
||||
))
|
||||
conversation_history = conversation_history_after_compression(
|
||||
agent, messages, conversation_history
|
||||
)
|
||||
|
||||
# Re-measure: same-count compression and media aging can shrink
|
||||
# the request without shrinking the array. Bytes are the yardstick
|
||||
# for a 413; tokens only for status display.
|
||||
new_tokens = estimate_messages_tokens_rough(messages)
|
||||
approx_tokens = new_tokens # update for downstream logging
|
||||
new_bytes = serialized_messages_bytes(messages)
|
||||
|
||||
made_progress = (
|
||||
len(messages) < original_len
|
||||
or (new_bytes > 0 and new_bytes < original_bytes * 0.95)
|
||||
)
|
||||
if made_progress:
|
||||
if len(messages) < original_len:
|
||||
agent._buffer_status(COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE.format(before=original_len, after=len(messages)))
|
||||
else:
|
||||
agent._buffer_status(
|
||||
f"🗜️ Compressed {original_bytes:,} → {new_bytes:,} "
|
||||
f"payload bytes, retrying..."
|
||||
)
|
||||
time.sleep(2) # Brief pause between compression retries
|
||||
_retry.restart_with_compressed_messages = True
|
||||
return _verdict("break")
|
||||
else:
|
||||
if agent._try_strip_image_parts_from_tool_messages(
|
||||
api_messages,
|
||||
remember_model=False,
|
||||
):
|
||||
agent._buffer_status(
|
||||
"📐 Compression could not reduce the request further — "
|
||||
"removed retained vision payloads and retrying..."
|
||||
)
|
||||
return _verdict("continue")
|
||||
|
||||
# Terminal — surface buffered context so the user
|
||||
# sees what compression attempts were made.
|
||||
agent._flush_status_buffer()
|
||||
agent._vprint(f"{agent.log_prefix}❌ Payload too large and cannot compress further.", force=True)
|
||||
agent._vprint(f"{agent.log_prefix} 💡 Try /new to start a fresh conversation, or /compress to retry compression.", force=True)
|
||||
logger.error("%s413 payload too large. Cannot compress further.", agent.log_prefix)
|
||||
agent._persist_session(messages, conversation_history)
|
||||
_final_response = "Request payload too large (413). Cannot compress further."
|
||||
return _verdict("return", {
|
||||
"final_response": _final_response,
|
||||
"messages": messages,
|
||||
"completed": False,
|
||||
"api_calls": api_call_count,
|
||||
"error": _final_response,
|
||||
"partial": True,
|
||||
"failed": True,
|
||||
"compression_exhausted": True,
|
||||
})
|
||||
|
||||
# Check context-length errors BEFORE the generic 4xx handler; the
|
||||
# classifier also covers 400/disconnect + large-session heuristics.
|
||||
is_context_length_error = (
|
||||
# Relay-wrapped output-cap 429s (parsed by the caller) go to the clamp, not
|
||||
# failover or generic retries (#72281). The classifier also covers
|
||||
# 400/disconnect + large-session heuristics.
|
||||
st.is_context_length_error = (
|
||||
classified.reason == FailoverReason.context_overflow
|
||||
# Relay-wrapped output-cap 429s (parsed above) go to the clamp
|
||||
# below, not failover or generic retries (#72281).
|
||||
or _wrapped_output_cap_budget is not None
|
||||
or wrapped_output_cap_budget is not None
|
||||
)
|
||||
|
||||
if is_context_length_error:
|
||||
compressor = agent.context_compressor
|
||||
old_ctx = compressor.context_length
|
||||
|
||||
# Two errors: "prompt too long" = INPUT overflows the window (shrink
|
||||
# context_length + compress); "max_tokens too large" = input fits
|
||||
# but input + max_tokens > window (shrink OUTPUT cap only).
|
||||
available_out = parse_available_output_tokens_from_error(error_msg)
|
||||
if available_out is not None:
|
||||
# Output-cap error: provider available_tokens is the
|
||||
# authoritative bound; also estimate the real request shape
|
||||
# (API-only content), use the smaller minus a margin.
|
||||
request_input_estimate = estimate_request_tokens_rough(
|
||||
api_messages, tools=agent.tools or None,
|
||||
)
|
||||
local_available_out = old_ctx - request_input_estimate
|
||||
if local_available_out > 0:
|
||||
safe_out = max(1, min(available_out, local_available_out) - 64)
|
||||
else:
|
||||
# Local estimate can overshoot; fall back to the
|
||||
# authoritative provider-reported budget.
|
||||
safe_out = max(1, available_out - 64)
|
||||
agent._ephemeral_max_output_tokens = safe_out
|
||||
agent._buffer_vprint(
|
||||
f"⚠️ Output cap too large for current prompt — "
|
||||
f"retrying with max_tokens={safe_out:,} "
|
||||
f"(provider_available={available_out:,}, "
|
||||
f"estimated_request_tokens={request_input_estimate:,}; "
|
||||
f"context_length unchanged at {old_ctx:,})"
|
||||
)
|
||||
# Still count against compression_attempts so we don't
|
||||
# loop forever if the error keeps recurring.
|
||||
compression_attempts += 1
|
||||
if compression_attempts > max_compression_attempts:
|
||||
agent._flush_status_buffer()
|
||||
agent._vprint(f"{agent.log_prefix}❌ Max compression attempts ({max_compression_attempts}) reached.", force=True)
|
||||
agent._vprint(f"{agent.log_prefix} 💡 Try /new to start a fresh conversation, or /compress to retry compression.", force=True)
|
||||
logger.error("%sContext compression failed after %d attempts.", agent.log_prefix, max_compression_attempts)
|
||||
agent._persist_session(messages, conversation_history)
|
||||
_final_response = f"Context length exceeded: max compression attempts ({max_compression_attempts}) reached."
|
||||
return _verdict("return", {
|
||||
"final_response": _final_response,
|
||||
"messages": messages,
|
||||
"completed": False,
|
||||
"api_calls": api_call_count,
|
||||
"error": _final_response,
|
||||
"partial": True,
|
||||
"failed": True,
|
||||
"compression_exhausted": True,
|
||||
})
|
||||
# Also compress history so the output-cap retry doesn't spin on
|
||||
# max_tokens alone; dropping the middle window makes the total
|
||||
# fit. (#55546)
|
||||
try:
|
||||
original_len = len(messages)
|
||||
original_tokens = estimate_messages_tokens_rough(messages)
|
||||
_overflow_input = messages
|
||||
messages, active_system_prompt = agent._compress_context(
|
||||
messages, system_message,
|
||||
approx_tokens=request_input_estimate,
|
||||
task_id=effective_task_id,
|
||||
bypass_cooldown=True, # #100661 provider-proven overflow
|
||||
)
|
||||
if messages is _overflow_input and compression_skipped_due_to_lock(agent):
|
||||
compression_attempts -= 1
|
||||
agent._persist_session(messages, conversation_history)
|
||||
return _verdict("return", _compression_deferred_result(
|
||||
agent, messages, api_call_count
|
||||
))
|
||||
if messages is _overflow_input and compression_blocked_transiently(agent):
|
||||
# #97488: timed transient guard — defer, never
|
||||
# exhaustion (gateway auto-reset).
|
||||
compression_attempts -= 1
|
||||
agent._persist_session(messages, conversation_history)
|
||||
return _verdict("return", _compression_deferred_result(
|
||||
agent, messages, api_call_count,
|
||||
reason="transient_block",
|
||||
))
|
||||
conversation_history = conversation_history_after_compression(
|
||||
agent, messages, conversation_history
|
||||
)
|
||||
new_tokens = estimate_messages_tokens_rough(messages)
|
||||
if len(messages) < original_len:
|
||||
agent._buffer_status(COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE.format(before=original_len, after=len(messages)))
|
||||
elif new_tokens > 0 and new_tokens < original_tokens * 0.95:
|
||||
agent._buffer_status(COMPRESSION_RETRY_TOKENS_STATUS_TEMPLATE.format(before=original_tokens, after=new_tokens))
|
||||
except Exception:
|
||||
# Compression must never turn an output-cap error
|
||||
# fatal — fall through and retry on max_tokens alone.
|
||||
logger.warning(
|
||||
"%sOutput-cap compression hit an error; retrying on max_tokens only.",
|
||||
agent.log_prefix,
|
||||
)
|
||||
_retry.restart_with_compressed_messages = True
|
||||
return _verdict("break")
|
||||
|
||||
# Output-cap error with unparseable budget: compression can't help
|
||||
# (input already fits) and would death-loop on the same 400. Fail
|
||||
# fast. (#55546)
|
||||
if is_output_cap_error(error_msg):
|
||||
agent._flush_status_buffer()
|
||||
agent._vprint(
|
||||
f"{agent.log_prefix}❌ The provider rejected the request because "
|
||||
f"max_tokens exceeds its output cap for this model.",
|
||||
force=True,
|
||||
)
|
||||
agent._vprint(
|
||||
f"{agent.log_prefix} 💡 Lower model.max_tokens in your config.yaml to "
|
||||
f"at or below the model's max-output limit. "
|
||||
f"(This is an output-cap error, not a context overflow — "
|
||||
f"compression cannot fix it.)",
|
||||
force=True,
|
||||
)
|
||||
logger.error(
|
||||
f"{agent.log_prefix}Output-cap error not routed into compression "
|
||||
f"(max_tokens over provider cap): {error_msg[:200]}"
|
||||
)
|
||||
agent._persist_session(messages, conversation_history)
|
||||
_final_response = (
|
||||
"max_tokens exceeds the provider's output cap for this model. "
|
||||
"Lower model.max_tokens in config.yaml."
|
||||
)
|
||||
return _verdict("return", {
|
||||
"final_response": _final_response,
|
||||
"messages": messages,
|
||||
"completed": False,
|
||||
"api_calls": api_call_count,
|
||||
"error": _final_response,
|
||||
"partial": True,
|
||||
"failed": True,
|
||||
})
|
||||
|
||||
# Input too large: shrink context_length only when the provider
|
||||
# reports the real limit; else keep the window and compress. Guessed
|
||||
# probe tiers can turn a configured 1M window into 256K/128K/64K.
|
||||
new_ctx = get_context_length_from_provider_error(error_msg, old_ctx)
|
||||
_provider_lower = (getattr(agent, "provider", "") or "").lower()
|
||||
_base_lower = (getattr(agent, "base_url", "") or "").rstrip("/").lower()
|
||||
is_minimax_provider = (
|
||||
_provider_lower in {"minimax", "minimax-cn"}
|
||||
or _base_lower.startswith((
|
||||
"https://api.minimax.io/anthropic",
|
||||
"https://api.minimaxi.com/anthropic",
|
||||
))
|
||||
)
|
||||
minimax_delta_only_overflow = (
|
||||
is_minimax_provider
|
||||
and new_ctx is None
|
||||
and "context window exceeds limit (" in error_msg
|
||||
)
|
||||
|
||||
if new_ctx is not None:
|
||||
agent._buffer_vprint(f"Context limit detected from API: {new_ctx:,} tokens (was {old_ctx:,})")
|
||||
compressor.update_model(
|
||||
model=agent.model,
|
||||
context_length=new_ctx,
|
||||
base_url=agent.base_url,
|
||||
api_key=getattr(agent, "api_key", ""),
|
||||
provider=agent.provider,
|
||||
api_mode=agent.api_mode,
|
||||
)
|
||||
# Persist the provider-reported limit before compression/retry:
|
||||
# rate limit, missing usage, or restart must not lose confirmed
|
||||
# metadata. Probe flags remain a fallback if this write fails.
|
||||
save_context_length(agent.model, agent.base_url, new_ctx)
|
||||
# Probe flags only on the built-in compressor (plugin engines
|
||||
# manage their own); provider-sourced value, so safe to cache.
|
||||
if hasattr(compressor, "_context_probed"):
|
||||
compressor._context_probed = True
|
||||
compressor._context_probe_persistable = True
|
||||
agent._buffer_vprint(f"⚠️ Context length exceeded — using provider limit: {old_ctx:,} → {new_ctx:,} tokens")
|
||||
elif minimax_delta_only_overflow:
|
||||
agent._buffer_vprint(
|
||||
f"Provider reported overflow amount only; "
|
||||
f"keeping context_length at {old_ctx:,} tokens and compressing."
|
||||
)
|
||||
else:
|
||||
agent._buffer_vprint(
|
||||
f"⚠️ Context length exceeded, but provider did not report a max context length; "
|
||||
f"keeping context_length at {old_ctx:,} tokens and compressing."
|
||||
)
|
||||
|
||||
compression_attempts += 1
|
||||
if compression_attempts > max_compression_attempts:
|
||||
agent._flush_status_buffer()
|
||||
agent._vprint(f"{agent.log_prefix}❌ Max compression attempts ({max_compression_attempts}) reached.", force=True)
|
||||
agent._vprint(f"{agent.log_prefix} 💡 Try /new to start a fresh conversation, or /compress to retry compression.", force=True)
|
||||
logger.error("%sContext compression failed after %d attempts.", agent.log_prefix, max_compression_attempts)
|
||||
agent._persist_session(messages, conversation_history)
|
||||
_final_response = f"Context length exceeded: max compression attempts ({max_compression_attempts}) reached."
|
||||
return _verdict("return", {
|
||||
"final_response": _final_response,
|
||||
"messages": messages,
|
||||
"completed": False,
|
||||
"api_calls": api_call_count,
|
||||
"error": _final_response,
|
||||
"partial": True,
|
||||
"failed": True,
|
||||
"compression_exhausted": True,
|
||||
})
|
||||
agent._buffer_status(COMPRESSION_RETRY_TOO_LARGE_STATUS_TEMPLATE.format(tokens=approx_tokens, attempt=compression_attempts, cap=max_compression_attempts))
|
||||
|
||||
original_len = len(messages)
|
||||
original_tokens = estimate_messages_tokens_rough(messages)
|
||||
_overflow_input = messages
|
||||
# Pass the OVERHEAD-AWARE size (msgs + tool schemas + system) so LCM
|
||||
# forced-overflow recovery arms on the TRUE request; approx_tokens
|
||||
# stays for status. See hermes-lcm _should_force_overflow_recovery.
|
||||
messages, active_system_prompt = agent._compress_context(
|
||||
messages, system_message,
|
||||
approx_tokens=estimate_request_tokens_rough(api_messages, tools=agent.tools or None),
|
||||
task_id=effective_task_id,
|
||||
# Provider proved the request doesn't fit: ignore the
|
||||
# summary-failure cooldown for this ONE attempt (bounded by
|
||||
# max_compression_attempts). (#100661)
|
||||
bypass_cooldown=True,
|
||||
)
|
||||
if messages is _overflow_input and compression_skipped_due_to_lock(agent):
|
||||
# Lock-skip: another path holds the compression lock, so this
|
||||
# pass no-oped. Temporary defer, not exhaustion — refund the
|
||||
# attempt, end the turn softly, no auto-reset. (#69870)
|
||||
compression_attempts -= 1
|
||||
agent._persist_session(messages, conversation_history)
|
||||
return _verdict("return", _compression_deferred_result(
|
||||
agent, messages, api_call_count
|
||||
))
|
||||
if messages is _overflow_input and compression_blocked_transiently(agent):
|
||||
# Transient block: a timed guard (host-timeout cooldown /
|
||||
# structural backoff) no-oped this pass — defer softly, never
|
||||
# compression_exhausted (auto-reset). (#97488)
|
||||
compression_attempts -= 1
|
||||
agent._persist_session(messages, conversation_history)
|
||||
return _verdict("return", _compression_deferred_result(
|
||||
agent, messages, api_call_count,
|
||||
reason="transient_block",
|
||||
))
|
||||
if context_compression_timed_out(agent):
|
||||
# Host timeout: recovery spent its wait budget with no committed
|
||||
# summary. Re-sending would hit the same overflow; end the turn
|
||||
# via the typed recovery contract. (#98722)
|
||||
agent._persist_session(messages, conversation_history)
|
||||
_final_response = _COMPRESSION_TIMEOUT_FINAL_RESPONSE
|
||||
return _verdict("return", {
|
||||
"final_response": _final_response,
|
||||
"messages": messages,
|
||||
"completed": False,
|
||||
"api_calls": api_call_count,
|
||||
"error": _final_response,
|
||||
"partial": True,
|
||||
"failed": True,
|
||||
"compression_exhausted": True,
|
||||
"turn_exit_reason": "context_compression_timeout",
|
||||
})
|
||||
conversation_history = conversation_history_after_compression(
|
||||
agent, messages, conversation_history
|
||||
)
|
||||
|
||||
# Re-estimate after compression: same-message-count compression
|
||||
# (tool-result pruning, in-place summarization) can shrink the
|
||||
# request. (#39550)
|
||||
new_tokens = estimate_messages_tokens_rough(messages)
|
||||
approx_tokens = new_tokens # update for downstream logging
|
||||
|
||||
if len(messages) < original_len or (new_tokens > 0 and new_tokens < original_tokens * 0.95) or (new_ctx and new_ctx < old_ctx):
|
||||
if len(messages) < original_len:
|
||||
agent._buffer_status(COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE.format(before=original_len, after=len(messages)))
|
||||
elif new_tokens > 0 and new_tokens < original_tokens * 0.95:
|
||||
agent._buffer_status(COMPRESSION_RETRY_TOKENS_STATUS_TEMPLATE.format(before=original_tokens, after=new_tokens))
|
||||
time.sleep(2) # Brief pause between compression retries
|
||||
# Rebuild the full request and force normal preflight to honor
|
||||
# it; message count alone doesn't prove system/tool-inclusive
|
||||
# pressure fell.
|
||||
_provider_overflow_recovery_pending = True
|
||||
_retry.restart_with_compressed_messages = True
|
||||
return _verdict("break")
|
||||
else:
|
||||
# Can't compress further and already at minimum tier
|
||||
agent._flush_status_buffer()
|
||||
agent._vprint(f"{agent.log_prefix}❌ Context length exceeded and cannot compress further.", force=True)
|
||||
agent._vprint(f"{agent.log_prefix} 💡 The conversation has accumulated too much content. Try /new to start fresh, or /compress to manually trigger compression.", force=True)
|
||||
logger.error("%sContext length exceeded: %s tokens. Cannot compress further.", agent.log_prefix, f"{new_tokens:,}")
|
||||
agent._persist_session(messages, conversation_history)
|
||||
_final_response = f"Context length exceeded ({new_tokens:,} tokens). Cannot compress further."
|
||||
return _verdict("return", {
|
||||
"final_response": _final_response,
|
||||
"messages": messages,
|
||||
"completed": False,
|
||||
"api_calls": api_call_count,
|
||||
"error": _final_response,
|
||||
"partial": True,
|
||||
"failed": True,
|
||||
"compression_exhausted": True,
|
||||
})
|
||||
return _verdict("fallthrough")
|
||||
if st.is_context_length_error:
|
||||
return _recover_context_length(st, _retry, error_msg)
|
||||
return st.verdict("fallthrough")
|
||||
|
||||
Reference in New Issue
Block a user