refactor(agent): compact turn-loop comments/docstrings to invariant statements (AST-identical)
This commit is contained in:
+917
-2589
File diff suppressed because it is too large
Load Diff
+148
-486
File diff suppressed because it is too large
Load Diff
+101
-273
@@ -1,24 +1,9 @@
|
|||||||
"""Post-loop turn finalization for ``run_conversation``.
|
"""Post-loop turn finalization for ``run_conversation``.
|
||||||
|
|
||||||
Extracted from ``agent/conversation_loop.py`` as part of the god-file
|
Lifted verbatim: budget summary, trajectory save, persist, diagnostics, response
|
||||||
decomposition campaign (``~/.hermes/plans/god-file-decomposition.md``, Phase 1
|
transforms, result assembly, steer drain, memory/skill review. Synchronous, single
|
||||||
step 4 — the post-loop ``TurnFinalizer`` seam). ``run_conversation``'s tail
|
return. ``logger`` is imported lazily from ``agent.conversation_loop`` (no cycle,
|
||||||
(everything after the main tool-calling ``while`` loop) is lifted here verbatim:
|
same logger name)."""
|
||||||
budget-exhaustion summary, trajectory save, session persist, turn diagnostics,
|
|
||||||
response transforms, result-dict assembly, steer drain, and the memory/skill
|
|
||||||
review trigger.
|
|
||||||
|
|
||||||
Behavior-neutral: the body is moved unchanged. All ``agent.*`` side effects fire
|
|
||||||
exactly as before; only the post-loop *locals* are passed in as keyword args, and
|
|
||||||
the assembled ``result`` dict is returned to ``run_conversation`` which returns it
|
|
||||||
to the caller. The function is synchronous with a single return — mirroring the
|
|
||||||
region it replaces (no awaits, no early returns).
|
|
||||||
|
|
||||||
Module ``logger`` is imported lazily inside the body (``from
|
|
||||||
agent.conversation_loop import logger``) so this module never imports
|
|
||||||
``agent.conversation_loop`` at import time -> no import cycle, and the log records
|
|
||||||
keep the exact logger name (``"agent.conversation_loop"``).
|
|
||||||
"""
|
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
@@ -54,11 +39,9 @@ def _fill_assistant_tail_content(agent, tail: dict, final_response) -> None:
|
|||||||
agent._db_flush_scan_prefix = None
|
agent._db_flush_scan_prefix = None
|
||||||
|
|
||||||
|
|
||||||
# Verification continuation scaffolding flags: verify-on-stop / pre_verify
|
# Verification-continuation nudges (verify-on-stop / pre_verify) must be stripped from
|
||||||
# inject a synthetic user nudge to keep the agent going one more turn.
|
# returned/live history to avoid role-alternation breaks; the assistant response is
|
||||||
# These nudges must be stripped from returned/live history to avoid
|
# real content and is not flagged. (#65919)
|
||||||
# role-alternation breaks and poisoning the resumed transcript. The
|
|
||||||
# assistant response is real content and is not flagged. (#65919 §7)
|
|
||||||
_VERIFICATION_CONTINUATION_FLAGS = (
|
_VERIFICATION_CONTINUATION_FLAGS = (
|
||||||
"_verification_stop_synthetic",
|
"_verification_stop_synthetic",
|
||||||
"_pre_verify_synthetic",
|
"_pre_verify_synthetic",
|
||||||
@@ -71,14 +54,10 @@ def _record_kanban_budget_exhausted(
|
|||||||
max_iterations: int,
|
max_iterations: int,
|
||||||
logger: logging.Logger,
|
logger: logging.Logger,
|
||||||
) -> None:
|
) -> None:
|
||||||
"""Record a terminal ``timed_out`` outcome for a kanban worker that
|
"""Record a terminal ``timed_out`` outcome for a kanban worker out of budget.
|
||||||
exhausted its iteration budget.
|
|
||||||
|
|
||||||
This is a bounded fallback (#87096): the CAS invariant in ``_end_run``
|
Idempotent via the ``_end_run`` CAS (``WHERE ended_at IS NULL``): a no-op if
|
||||||
(``WHERE ended_at IS NULL``) guarantees idempotence — if another path
|
another path already closed the run, so safe from multiple exit paths (#87096)."""
|
||||||
already closed the run this is a no-op — so it is safe to call from
|
|
||||||
multiple exit paths.
|
|
||||||
"""
|
|
||||||
try:
|
try:
|
||||||
from hermes_cli import kanban_db as _kb
|
from hermes_cli import kanban_db as _kb
|
||||||
_conn = _kb.connect()
|
_conn = _kb.connect()
|
||||||
@@ -116,10 +95,8 @@ def _record_kanban_budget_exhausted(
|
|||||||
def _drop_verification_continuation_scaffolding(messages) -> None:
|
def _drop_verification_continuation_scaffolding(messages) -> None:
|
||||||
"""Remove verification-continuation nudge messages from *messages* in place.
|
"""Remove verification-continuation nudge messages from *messages* in place.
|
||||||
|
|
||||||
Only the synthetic nudges carry these flags, so this strips just the
|
Only the synthetic nudges carry these flags, so the real attempted-final-answer
|
||||||
nudges while preserving the real attempted-final-answer that was
|
persisted to state.db survives."""
|
||||||
persisted to state.db.
|
|
||||||
"""
|
|
||||||
messages[:] = [
|
messages[:] = [
|
||||||
m for m in messages
|
m for m in messages
|
||||||
if not (isinstance(m, dict) and any(m.get(f) for f in _VERIFICATION_CONTINUATION_FLAGS))
|
if not (isinstance(m, dict) and any(m.get(f) for f in _VERIFICATION_CONTINUATION_FLAGS))
|
||||||
@@ -153,11 +130,7 @@ def finalize_turn(
|
|||||||
_pending_verification_response=None,
|
_pending_verification_response=None,
|
||||||
_pending_verification_response_previewed=False,
|
_pending_verification_response_previewed=False,
|
||||||
):
|
):
|
||||||
"""Run the post-loop finalization and return the turn ``result`` dict.
|
"""Run the post-loop finalization and return the turn ``result`` dict."""
|
||||||
|
|
||||||
Lifted verbatim from ``run_conversation`` (the region after the main agent
|
|
||||||
loop). See module docstring.
|
|
||||||
"""
|
|
||||||
from agent.conversation_loop import logger
|
from agent.conversation_loop import logger
|
||||||
|
|
||||||
budget_exhausted = (
|
budget_exhausted = (
|
||||||
@@ -179,24 +152,19 @@ def finalize_turn(
|
|||||||
iteration_limit_fallback = False
|
iteration_limit_fallback = False
|
||||||
preserved_verification_fallback = False
|
preserved_verification_fallback = False
|
||||||
if continuation_budget_exhausted:
|
if continuation_budget_exhausted:
|
||||||
# A verification/continuation gate deliberately withheld a composed
|
# A verification gate withheld a composed answer, then the budget ran out:
|
||||||
# answer, then consumed the remaining budget before producing a newer
|
# preserve it rather than make another fallible call. The explicit pending
|
||||||
# one. Preserve that exact answer instead of replacing it with another
|
# value is the provenance guard; unrelated error exits never enter here.
|
||||||
# fallible model call. The explicit pending value is the provenance
|
|
||||||
# guard: unrelated error/recovery exits can never enter this branch.
|
|
||||||
final_response = _pending_verification_response
|
final_response = _pending_verification_response
|
||||||
# Mark the turn as previewed only when the reused candidate was
|
# Previewed only if the reused candidate was actually streamed as interim.
|
||||||
# actually streamed to the user as interim content. (#65919 review:
|
|
||||||
# response-loss blocker)
|
|
||||||
if _pending_verification_response_previewed:
|
if _pending_verification_response_previewed:
|
||||||
agent._response_was_previewed = True
|
agent._response_was_previewed = True
|
||||||
_turn_exit_reason = f"max_iterations_reached({api_call_count}/{agent.max_iterations})"
|
_turn_exit_reason = f"max_iterations_reached({api_call_count}/{agent.max_iterations})"
|
||||||
iteration_limit_fallback = True
|
iteration_limit_fallback = True
|
||||||
preserved_verification_fallback = True
|
preserved_verification_fallback = True
|
||||||
elif final_response is None and budget_fallback_eligible:
|
elif final_response is None and budget_fallback_eligible:
|
||||||
# Budget exhausted — ask the model for a summary via one extra
|
# Budget exhausted: _handle_max_iterations makes one extra toolless request
|
||||||
# API call with tools stripped. _handle_max_iterations injects a
|
# for a summary.
|
||||||
# user message and makes a single toolless request.
|
|
||||||
_turn_exit_reason = f"max_iterations_reached({api_call_count}/{agent.max_iterations})"
|
_turn_exit_reason = f"max_iterations_reached({api_call_count}/{agent.max_iterations})"
|
||||||
agent._emit_status(
|
agent._emit_status(
|
||||||
f"⚠️ Iteration budget exhausted ({api_call_count}/{agent.max_iterations}) "
|
f"⚠️ Iteration budget exhausted ({api_call_count}/{agent.max_iterations}) "
|
||||||
@@ -211,29 +179,18 @@ def finalize_turn(
|
|||||||
iteration_limit_fallback = True
|
iteration_limit_fallback = True
|
||||||
|
|
||||||
if iteration_limit_fallback:
|
if iteration_limit_fallback:
|
||||||
# If running as a kanban worker, signal the dispatcher that the
|
# Kanban worker: signal the dispatcher the worker could not complete. Route
|
||||||
# worker could not complete (rather than treating it as a
|
# via ``_record_task_failure(outcome="timed_out")`` (not ``kanban_block``) so
|
||||||
# protocol violation). This applies whether the user-facing fallback
|
# it counts toward the consecutive-failure circuit breaker (#29747).
|
||||||
# came from the summary call or an explicitly pending continuation;
|
|
||||||
# both exhausted the task budget and must advance the failure circuit.
|
|
||||||
#
|
|
||||||
# We route through ``_record_task_failure(outcome="timed_out")``
|
|
||||||
# rather than ``kanban_block`` so this counts toward the dispatcher's
|
|
||||||
# consecutive-failure circuit breaker (#29747 gap 2).
|
|
||||||
_kanban_task = os.environ.get("HERMES_KANBAN_TASK")
|
_kanban_task = os.environ.get("HERMES_KANBAN_TASK")
|
||||||
if _kanban_task:
|
if _kanban_task:
|
||||||
_record_kanban_budget_exhausted(
|
_record_kanban_budget_exhausted(
|
||||||
_kanban_task, api_call_count, agent.max_iterations, logger,
|
_kanban_task, api_call_count, agent.max_iterations, logger,
|
||||||
)
|
)
|
||||||
elif budget_exhausted:
|
elif budget_exhausted:
|
||||||
# Bounded fallback (#87096): budget was exhausted but none of the
|
# Bounded fallback: budget exhausted with no eligible fallback path. A kanban
|
||||||
# normal fallback paths were eligible (interrupted / failed /
|
# worker must still record a terminal outcome; the ``_end_run`` CAS makes it
|
||||||
# anomalous exit_reason). If running as a kanban worker we must
|
# idempotent if another path already closed the run (#87096).
|
||||||
# still record a terminal outcome so the task does not remain in
|
|
||||||
# an ambiguous lifecycle state. The worker's run is closed via
|
|
||||||
# ``_record_task_failure`` (compare-and-swap receipt path) which
|
|
||||||
# is a no-op if another path closed it — the CAS invariant in
|
|
||||||
# ``_end_run`` (``WHERE ended_at IS NULL``) guarantees idempotence.
|
|
||||||
_kanban_task = os.environ.get("HERMES_KANBAN_TASK")
|
_kanban_task = os.environ.get("HERMES_KANBAN_TASK")
|
||||||
if _kanban_task:
|
if _kanban_task:
|
||||||
_record_kanban_budget_exhausted(
|
_record_kanban_budget_exhausted(
|
||||||
@@ -251,18 +208,9 @@ def finalize_turn(
|
|||||||
)
|
)
|
||||||
)
|
)
|
||||||
|
|
||||||
# Preflight can seed the display count before the provider receives the
|
# Roll back the preflight-seeded display count only when an interrupt wins
|
||||||
# request. Roll that estimate back only when an interrupt wins the race
|
# before any provider response; compaction state (incl. ``-1``) stays with the
|
||||||
# before any successful provider response. Compaction state remains owned
|
# real-usage path. Type-pinned guards keep MagicMock/SimpleNamespace doubles inert.
|
||||||
# by the real-usage/post-compaction path, including its ``-1`` sentinel.
|
|
||||||
# Guard rules (test-double density on this path is high):
|
|
||||||
# - snapshot is type-pinned to a real int — MagicMock agents auto-create
|
|
||||||
# truthy Mock attributes that must never arm the rollback;
|
|
||||||
# - the received-response flag is pinned to ``is not True`` — its real
|
|
||||||
# domain is True/False, and only a literal True means a provider
|
|
||||||
# response completed;
|
|
||||||
# - the compressor method gets a getattr+callable guard — SimpleNamespace
|
|
||||||
# compressor doubles and plugin context engines lack it.
|
|
||||||
_preflight_snapshot = getattr(
|
_preflight_snapshot = getattr(
|
||||||
agent, "_turn_preflight_display_snapshot", None
|
agent, "_turn_preflight_display_snapshot", None
|
||||||
)
|
)
|
||||||
@@ -281,17 +229,9 @@ def finalize_turn(
|
|||||||
if callable(_rollback_fn):
|
if callable(_rollback_fn):
|
||||||
_rollback_fn(_preflight_snapshot)
|
_rollback_fn(_preflight_snapshot)
|
||||||
|
|
||||||
# Post-loop cleanup must never lose the response. Trajectory save,
|
# Post-loop cleanup must never lose the response: trajectory save, teardown,
|
||||||
# resource teardown, and session persistence all touch fallible
|
# and session persist are guarded independently and errors surface via
|
||||||
# surfaces — file I/O / JSON serialization (_save_trajectory), remote
|
# ``cleanup_errors`` rather than killing the turn (#8049).
|
||||||
# VM/browser teardown over the network (_cleanup_task_resources), and
|
|
||||||
# SQLite writes (_persist_session). A raise from any of them used to
|
|
||||||
# propagate straight out of run_conversation, discarding the partial
|
|
||||||
# final_response the caller is waiting for (subprocess wrappers saw an
|
|
||||||
# empty stdout with no traceback — #8049). Each step is now guarded
|
|
||||||
# independently so one failure can't skip the others, and any errors
|
|
||||||
# are surfaced on the result dict via ``cleanup_errors`` rather than
|
|
||||||
# killing the turn.
|
|
||||||
_cleanup_errors = []
|
_cleanup_errors = []
|
||||||
|
|
||||||
# Save trajectory if enabled. ``user_message`` may be a multimodal
|
# Save trajectory if enabled. ``user_message`` may be a multimodal
|
||||||
@@ -309,22 +249,17 @@ def finalize_turn(
|
|||||||
_cleanup_errors.append(f"cleanup_task_resources: {_cleanup_err}")
|
_cleanup_errors.append(f"cleanup_task_resources: {_cleanup_err}")
|
||||||
logger.error("finalize_turn: _cleanup_task_resources failed: %s", _cleanup_err, exc_info=True)
|
logger.error("finalize_turn: _cleanup_task_resources failed: %s", _cleanup_err, exc_info=True)
|
||||||
|
|
||||||
# Persist session to both JSON log and SQLite only after private retry
|
# Persist only after private retry scaffolding is removed, or a later "continue"
|
||||||
# scaffolding has been removed. Otherwise a later user "continue" turn
|
# replays assistant("(empty)") / recovery nudges into the same empty-response loop.
|
||||||
# can replay assistant("(empty)") / recovery nudges and fall into the
|
|
||||||
# same empty-response loop again.
|
|
||||||
try:
|
try:
|
||||||
agent._drop_trailing_empty_response_scaffolding(messages)
|
agent._drop_trailing_empty_response_scaffolding(messages)
|
||||||
|
|
||||||
# Drop verification-continuation nudges (synthetic user messages)
|
# Strip only the synthetic verification nudges before the tail-assistant
|
||||||
# from the live history before the tail-assistant check — only the
|
# check; the assistant candidate persists in state.db. (#65919)
|
||||||
# nudges need stripping; the assistant candidate persists in
|
|
||||||
# state.db. (#65919 §7)
|
|
||||||
_drop_verification_continuation_scaffolding(messages)
|
_drop_verification_continuation_scaffolding(messages)
|
||||||
|
|
||||||
# #95514: an empty terminal completion is not authoritative when the
|
# An empty terminal completion is not authoritative when the stream already
|
||||||
# stream already delivered text. Recover before persist so a blank
|
# delivered text; recover before persist so a blank tail isn't frozen (#95514).
|
||||||
# assistant tail is filled instead of frozen as content=''.
|
|
||||||
_recovered_from_stream = False
|
_recovered_from_stream = False
|
||||||
if not interrupted and not failed:
|
if not interrupted and not failed:
|
||||||
_streamed = getattr(agent, "_current_streamed_assistant_text", "") or ""
|
_streamed = getattr(agent, "_current_streamed_assistant_text", "") or ""
|
||||||
@@ -337,39 +272,16 @@ def finalize_turn(
|
|||||||
final_response = _streamed
|
final_response = _streamed
|
||||||
_recovered_from_stream = True
|
_recovered_from_stream = True
|
||||||
|
|
||||||
# When the turn was interrupted and the last message is a tool
|
# An interrupt can leave a tool result as the tail (no scaffolding flag rewinds
|
||||||
# result, append a synthetic assistant message to close the
|
# it); close the sequence so strict providers don't see ``tool → user``. An
|
||||||
# tool-call sequence. Without this, the session persists a
|
# explicit placeholder is used since final_response is usually empty (#48879).
|
||||||
# ``tool → user`` alternation that strict providers (Gemini,
|
|
||||||
# Claude) reject, causing them to hallucinate a continuation of
|
|
||||||
# the user's message on the next turn (#48879).
|
|
||||||
#
|
|
||||||
# ``_drop_trailing_empty_response_scaffolding`` only rewinds the
|
|
||||||
# tool tail when an empty-response scaffolding flag is present; a
|
|
||||||
# clean ``/stop`` interrupt after a successful tool sets no such
|
|
||||||
# flag, so the tool result survives as the tail and we close it
|
|
||||||
# here instead. On an interrupt ``final_response`` is typically
|
|
||||||
# empty, so fall back to an explicit placeholder rather than
|
|
||||||
# persisting an empty-content assistant turn.
|
|
||||||
if interrupted:
|
if interrupted:
|
||||||
from agent.message_sanitization import close_interrupted_tool_sequence
|
from agent.message_sanitization import close_interrupted_tool_sequence
|
||||||
close_interrupted_tool_sequence(messages, final_response)
|
close_interrupted_tool_sequence(messages, final_response)
|
||||||
|
|
||||||
# Some recovery/fallback paths return a real final_response without
|
# Recovery ``break`` sites can return a final_response with no closing
|
||||||
# adding a closing assistant message to the transcript (e.g. the
|
# assistant row; enforce "delivered final_response ⇒ assistant row" here.
|
||||||
# partial-stream and prior-turn-content recovery ``break`` sites in
|
# Compare content, not role, so a matching verification candidate isn't dup'd.
|
||||||
# ``conversation_loop``). If persisted as-is, the durable session can
|
|
||||||
# end at a tool/user message even though the caller — and the gateway
|
|
||||||
# platform — already saw a completed assistant response. The next turn
|
|
||||||
# then replays a user-only backlog and the model re-answers every
|
|
||||||
# "unanswered" message. Close the durable turn at the source, at the
|
|
||||||
# single chokepoint every recovery ``break`` flows through, so the
|
|
||||||
# invariant "delivered final_response ⇒ assistant row in transcript"
|
|
||||||
# holds regardless of which path produced it. (#43849 / #44100)
|
|
||||||
#
|
|
||||||
# Compare content (not just role) so a verification candidate that
|
|
||||||
# matches the final response is not duplicated at budget
|
|
||||||
# exhaustion. (#65919 §7)
|
|
||||||
if final_response and not interrupted:
|
if final_response and not interrupted:
|
||||||
try:
|
try:
|
||||||
_tail = messages[-1] if messages else None
|
_tail = messages[-1] if messages else None
|
||||||
@@ -394,66 +306,45 @@ def finalize_turn(
|
|||||||
)
|
)
|
||||||
)
|
)
|
||||||
):
|
):
|
||||||
# The tail IS an assistant row, but a *pure tool-call turn* or
|
# Tail is an assistant row (pure tool-call turn or stream-recovered
|
||||||
# a blank assistant tail whose content was recovered from the
|
# blank, #95514): fill its content rather than append a second row.
|
||||||
# stream buffer (#95514). Fill that row's content instead of
|
|
||||||
# appending, so the durable turn ends with the answer without
|
|
||||||
# creating an assistant→assistant pair.
|
|
||||||
_fill_assistant_tail_content(agent, _tail, final_response)
|
_fill_assistant_tail_content(agent, _tail, final_response)
|
||||||
|
|
||||||
# The model has completed its request, so replace API-local
|
# Request is complete, so replace API-local voice/model/skill guidance with
|
||||||
# voice/model/skill guidance with the clean user input before writing the
|
# the clean user input before the durable snapshot; earlier flushes used the
|
||||||
# final durable snapshot and returning the continuation history. Earlier
|
# DB-only override as their messages were still needed (#48677 / #63766).
|
||||||
# turn-start flushes use the DB-only override because their messages are
|
|
||||||
# still needed for the API request; this finalizer runs after that request
|
|
||||||
# is complete (#48677 / #63766).
|
|
||||||
_apply_override = getattr(agent, "_apply_persist_user_message_override", None)
|
_apply_override = getattr(agent, "_apply_persist_user_message_override", None)
|
||||||
if callable(_apply_override):
|
if callable(_apply_override):
|
||||||
_apply_override(messages)
|
_apply_override(messages)
|
||||||
|
|
||||||
# ── Post-turn micro-compaction ────────────────────────────
|
# ── Post-turn micro-compaction ────────────────────────────
|
||||||
# After the assistant response is finalized but before the session is
|
# Absorb the oldest uncompacted exchange into the rolling summary before
|
||||||
# persisted, run micro-compaction to absorb the oldest uncompacted
|
# persist, amortizing compression across turns instead of one big pause.
|
||||||
# exchange into the rolling summary. This amortizes compression
|
|
||||||
# across turns rather than batching it into one big pause.
|
|
||||||
if not interrupted and not failed:
|
if not interrupted and not failed:
|
||||||
try:
|
try:
|
||||||
_compressor = getattr(agent, "context_compressor", None)
|
_compressor = getattr(agent, "context_compressor", None)
|
||||||
# Strict `is True` + isinstance gates: plugin context engines
|
# Strict `is True` + isinstance gates: plugin context engines and
|
||||||
# (and MagicMock compressors in tests) satisfy getattr/duck
|
# MagicMock compressors pass duck checks and would wipe the transcript.
|
||||||
# checks with truthy auto-attributes — a bare truthiness check
|
|
||||||
# here called _micro_compact on a mock and spliced its (empty-
|
|
||||||
# iterating) return value over the transcript, wiping it.
|
|
||||||
if (
|
if (
|
||||||
_compressor
|
_compressor
|
||||||
and getattr(_compressor, '_micro_compact_enabled', False) is True
|
and getattr(_compressor, '_micro_compact_enabled', False) is True
|
||||||
and callable(getattr(_compressor, '_micro_compact', None))
|
and callable(getattr(_compressor, '_micro_compact', None))
|
||||||
and final_response
|
and final_response
|
||||||
# compression.checkpoint_required: agent init already
|
# Micro-compaction has no checkpoint hook, so it must never run
|
||||||
# forces _micro_compact_enabled off, but the compressor
|
# while compression.checkpoint_required is armed.
|
||||||
# attribute is plain state a future path could flip on a
|
|
||||||
# live agent. Micro-compaction has no checkpoint hook in
|
|
||||||
# its path, so it must never run while the gate is armed.
|
|
||||||
and getattr(
|
and getattr(
|
||||||
agent, "compression_checkpoint_required", False
|
agent, "compression_checkpoint_required", False
|
||||||
) is not True
|
) is not True
|
||||||
# Persistence-isolated agents (background review fork)
|
# Persistence-isolated agents (background review fork) must not
|
||||||
# must not micro-compact: the pass burns a real aux-LLM
|
# micro-compact: it burns an aux-LLM call on a throwaway transcript
|
||||||
# call on a throwaway replay transcript, and if the
|
# and could archive_and_compact the CANONICAL session rows.
|
||||||
# compressor ever holds a session_db binding it would
|
|
||||||
# archive_and_compact the CANONICAL session rows — the
|
|
||||||
# exact write class _persist_disabled exists to stop.
|
|
||||||
and not getattr(agent, "_persist_disabled", False)
|
and not getattr(agent, "_persist_disabled", False)
|
||||||
):
|
):
|
||||||
_before = len(messages)
|
_before = len(messages)
|
||||||
_compacted = _compressor._micro_compact(messages)
|
_compacted = _compressor._micro_compact(messages)
|
||||||
# Micro-compaction defrag rewrites the newest MICRO
|
# Defrag rewrites the newest MICRO marker in place and pops
|
||||||
# marker's content and pops _db_persisted from the live
|
# _db_persisted; the compressor flags us to invalidate the flush-
|
||||||
# dict in place — the sibling of the pop site above. The
|
# scan cursor, else the rewritten row is identity-skipped (stale).
|
||||||
# compressor has no agent reference, so it raises a flag
|
|
||||||
# for us to invalidate the bounded flush-scan cursor;
|
|
||||||
# otherwise the rewritten marker row is identity-skipped
|
|
||||||
# and the stale summary persists to state.db.
|
|
||||||
if getattr(
|
if getattr(
|
||||||
_compressor, "_flush_scan_cursor_invalidated", False
|
_compressor, "_flush_scan_cursor_invalidated", False
|
||||||
):
|
):
|
||||||
@@ -475,18 +366,16 @@ def finalize_turn(
|
|||||||
_cleanup_errors.append(f"persist_session: {_persist_err}")
|
_cleanup_errors.append(f"persist_session: {_persist_err}")
|
||||||
logger.error("finalize_turn: _persist_session failed: %s", _persist_err, exc_info=True)
|
logger.error("finalize_turn: _persist_session failed: %s", _persist_err, exc_info=True)
|
||||||
|
|
||||||
# The gateway owns a separate in-memory history snapshot. Keep it current
|
# Keep the gateway's separate in-memory history snapshot current even on
|
||||||
# even when finalization reports a cleanup error: a later prompt must not be
|
# cleanup error, so a later prompt isn't sent with a pre-turn snapshot.
|
||||||
# sent with the pre-turn snapshot while the durable DB already has this turn.
|
|
||||||
try:
|
try:
|
||||||
agent._session_messages = messages
|
agent._session_messages = messages
|
||||||
except Exception:
|
except Exception:
|
||||||
pass
|
pass
|
||||||
|
|
||||||
# ── Turn-exit diagnostic log ─────────────────────────────────────
|
# ── Turn-exit diagnostic log ─────────────────────────────────────
|
||||||
# Always logged at INFO so agent.log captures WHY every turn ended.
|
# Always INFO so agent.log captures WHY every turn ended; WARNING when the last
|
||||||
# When the last message is a tool result (agent was mid-work), log
|
# message is a tool result (the "just stops" scenario).
|
||||||
# at WARNING — this is the "just stops" scenario users report.
|
|
||||||
_last_msg_role = messages[-1].get("role") if messages else None
|
_last_msg_role = messages[-1].get("role") if messages else None
|
||||||
_last_tool_name = None
|
_last_tool_name = None
|
||||||
if _last_msg_role == "tool":
|
if _last_msg_role == "tool":
|
||||||
@@ -527,21 +416,9 @@ def finalize_turn(
|
|||||||
else:
|
else:
|
||||||
logger.info(_diag_msg, *_diag_args)
|
logger.info(_diag_msg, *_diag_args)
|
||||||
|
|
||||||
# File-mutation verifier footer.
|
# File-mutation verifier footer: if ``write_file`` / ``patch`` calls failed and
|
||||||
# If one or more ``write_file`` / ``patch`` calls failed during this
|
# were never superseded by a successful write to the same path, append an
|
||||||
# turn and were never superseded by a successful write to the same
|
# advisory so over-claiming is surfaced. Only on real, uninterrupted responses.
|
||||||
# path, append an advisory footer to the assistant response. This
|
|
||||||
# catches the specific case — reported by Ben Eng (#15524-adjacent)
|
|
||||||
# — where a model issues a batch of parallel patches, half of them
|
|
||||||
# fail with "Could not find old_string", and the model summarises
|
|
||||||
# the turn claiming every file was edited. The user then has to
|
|
||||||
# manually run ``git status`` to catch the lie. With this footer
|
|
||||||
# the truth is surfaced on every turn, so over-claiming is
|
|
||||||
# structurally impossible past the model.
|
|
||||||
#
|
|
||||||
# Gate: only applied when a real text response exists for this
|
|
||||||
# turn and the user didn't interrupt. Empty/interrupted turns
|
|
||||||
# already have other surface text that shouldn't be augmented.
|
|
||||||
if final_response and not interrupted:
|
if final_response and not interrupted:
|
||||||
try:
|
try:
|
||||||
_failed = getattr(agent, "_turn_failed_file_mutations", None) or {}
|
_failed = getattr(agent, "_turn_failed_file_mutations", None) or {}
|
||||||
@@ -552,30 +429,16 @@ def finalize_turn(
|
|||||||
except Exception as _ver_err:
|
except Exception as _ver_err:
|
||||||
logger.debug("file-mutation verifier footer failed: %s", _ver_err)
|
logger.debug("file-mutation verifier footer failed: %s", _ver_err)
|
||||||
|
|
||||||
# Turn-completion explainer.
|
# Turn-completion explainer: on abnormal exits, surface one explanation from
|
||||||
# When a turn ends abnormally after substantive work — empty content
|
# ``_turn_exit_reason``. Only acts when no usable reply exists (empty, "(empty)",
|
||||||
# after retries, a partial/truncated stream, a still-pending tool
|
# or a short unpunctuated fragment); ``text_response(...)`` exits stay silent.
|
||||||
# result, or an iteration/budget limit — the user otherwise gets a
|
|
||||||
# blank or fragmentary response box with no consolidated reason why
|
|
||||||
# the agent stopped (#34452). Surface a single user-visible
|
|
||||||
# explanation derived from ``_turn_exit_reason``, mirroring the
|
|
||||||
# file-mutation verifier footer pattern above.
|
|
||||||
#
|
|
||||||
# Gate carefully so healthy turns stay quiet:
|
|
||||||
# - ``text_response(...)`` exits never produce an explanation
|
|
||||||
# (handled inside the formatter), so a terse ``Done.`` is silent.
|
|
||||||
# - We only ACT when there is no genuinely usable reply this turn:
|
|
||||||
# an empty response, the "(empty)" terminal sentinel, or a
|
|
||||||
# suspiciously short partial fragment with no terminating
|
|
||||||
# punctuation (e.g. "The"). A real short answer keeps its text.
|
|
||||||
if not interrupted:
|
if not interrupted:
|
||||||
try:
|
try:
|
||||||
if agent._turn_completion_explainer_enabled():
|
if agent._turn_completion_explainer_enabled():
|
||||||
_stripped = (final_response or "").strip()
|
_stripped = (final_response or "").strip()
|
||||||
_is_empty_terminal = _stripped == "" or _stripped == "(empty)"
|
_is_empty_terminal = _stripped == "" or _stripped == "(empty)"
|
||||||
# A short fragment that is not a normal text_response exit
|
# A short fragment not from a text_response exit and lacking sentence-
|
||||||
# and lacks sentence-ending punctuation is treated as a
|
# ending punctuation is treated as a truncated partial (#34452).
|
||||||
# truncated partial (the "The" case from #34452).
|
|
||||||
_is_partial_fragment = (
|
_is_partial_fragment = (
|
||||||
not _is_empty_terminal
|
not _is_empty_terminal
|
||||||
and not preserved_verification_fallback
|
and not preserved_verification_fallback
|
||||||
@@ -601,9 +464,7 @@ def finalize_turn(
|
|||||||
# the actionable explanation.
|
# the actionable explanation.
|
||||||
final_response = _explanation
|
final_response = _explanation
|
||||||
else:
|
else:
|
||||||
# Keep the partial fragment, append the reason so
|
# Keep the partial fragment and append why it stopped.
|
||||||
# the user sees both what arrived and why it
|
|
||||||
# stopped.
|
|
||||||
final_response = (
|
final_response = (
|
||||||
_stripped + "\n\n" + _explanation
|
_stripped + "\n\n" + _explanation
|
||||||
)
|
)
|
||||||
@@ -613,10 +474,8 @@ def finalize_turn(
|
|||||||
_response_transformed = False
|
_response_transformed = False
|
||||||
_pre_transform_response = None
|
_pre_transform_response = None
|
||||||
|
|
||||||
# Plugin hook: transform_llm_output
|
# Plugin hook: transform_llm_output — fired once per turn after the tool loop.
|
||||||
# Fired once per turn after the tool-calling loop completes.
|
# First hook to return a string wins; None/empty leaves the text unchanged.
|
||||||
# Plugins can transform the LLM's output text before it's returned.
|
|
||||||
# First hook to return a string wins; None/empty return leaves text unchanged.
|
|
||||||
if final_response and not interrupted:
|
if final_response and not interrupted:
|
||||||
try:
|
try:
|
||||||
from hermes_cli.lifecycle import invoke_hook as _invoke_hook
|
from hermes_cli.lifecycle import invoke_hook as _invoke_hook
|
||||||
@@ -636,10 +495,8 @@ def finalize_turn(
|
|||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
logger.warning("transform_llm_output hook failed: %s", exc)
|
logger.warning("transform_llm_output hook failed: %s", exc)
|
||||||
|
|
||||||
# Plugin hook: post_llm_call
|
# Plugin hook: post_llm_call — fired once per turn after the tool loop (e.g. sync
|
||||||
# Fired once per turn after the tool-calling loop completes.
|
# conversation data to an external memory system).
|
||||||
# Plugins can use this to persist conversation data (e.g. sync
|
|
||||||
# to an external memory system).
|
|
||||||
if final_response and not interrupted:
|
if final_response and not interrupted:
|
||||||
try:
|
try:
|
||||||
from hermes_cli.lifecycle import invoke_hook as _invoke_hook
|
from hermes_cli.lifecycle import invoke_hook as _invoke_hook
|
||||||
@@ -657,18 +514,12 @@ def finalize_turn(
|
|||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
logger.warning("post_llm_call hook failed: %s", exc)
|
logger.warning("post_llm_call hook failed: %s", exc)
|
||||||
|
|
||||||
# Context engine observation hook: notify the active engine that this
|
# Context engine observation hook (complements per-request select_context()):
|
||||||
# turn has finished, with the finalized transcript. Complements the
|
# notify the engine the turn finished with the finalized transcript. Fail-open.
|
||||||
# per-request select_context() hook (selection before the request;
|
|
||||||
# observation after the turn). No-op default, fail-open.
|
|
||||||
try:
|
try:
|
||||||
from agent.conversation_loop import _notify_context_engine_turn_complete
|
from agent.conversation_loop import _notify_context_engine_turn_complete
|
||||||
# Forward the turn's canonical usage when the host has it. The loop
|
# ``_last_turn_usage`` holds the last API response's canonical usage dict, or
|
||||||
# stashes the most recent API response's usage dict (the same
|
# ``None`` on turns that never reached a provider response — by contract.
|
||||||
# canonical buckets fed to ``update_from_response``) on the agent as
|
|
||||||
# ``_last_turn_usage``. It is ``None`` on turns that never reached a
|
|
||||||
# provider response (early failure / interrupt), which is exactly the
|
|
||||||
# contract: real usage when available, ``None`` otherwise.
|
|
||||||
_turn_usage = getattr(agent, "_last_turn_usage", None)
|
_turn_usage = getattr(agent, "_last_turn_usage", None)
|
||||||
_notify_context_engine_turn_complete(
|
_notify_context_engine_turn_complete(
|
||||||
agent,
|
agent,
|
||||||
@@ -685,15 +536,9 @@ def finalize_turn(
|
|||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
logger.warning("on_turn_complete notification failed: %s", exc)
|
logger.warning("on_turn_complete notification failed: %s", exc)
|
||||||
|
|
||||||
# Extract reasoning from the CURRENT turn only. Walk backwards
|
# Reasoning from the CURRENT turn only: stop at this turn's user message
|
||||||
# but stop at the user message that started this turn — anything
|
# (#17055), but take the most recent non-empty reasoning since many providers
|
||||||
# earlier is from a prior turn and must not leak into the reasoning
|
# emit it on the tool-call step and leave the final step with reasoning=None.
|
||||||
# box (confusing stale display; #17055). Within the current turn
|
|
||||||
# we still want the *most recent* non-empty reasoning: many
|
|
||||||
# providers (Claude thinking, DeepSeek v4, Codex Responses) emit
|
|
||||||
# reasoning on the tool-call step and leave the final-answer step
|
|
||||||
# with reasoning=None, so picking only the last assistant would
|
|
||||||
# silently drop legitimate same-turn reasoning.
|
|
||||||
last_reasoning = None
|
last_reasoning = None
|
||||||
for msg in reversed(messages):
|
for msg in reversed(messages):
|
||||||
if msg.get("role") == "user":
|
if msg.get("role") == "user":
|
||||||
@@ -702,15 +547,9 @@ def finalize_turn(
|
|||||||
last_reasoning = msg["reasoning"]
|
last_reasoning = msg["reasoning"]
|
||||||
break
|
break
|
||||||
|
|
||||||
# Class-level surrogate chokepoint (#80366, #55143, #55309, #19819):
|
# Surrogate chokepoint: ``final_response`` may be RAW SDK content, and a lone UTF-16
|
||||||
# ``final_response`` is often the RAW SDK content
|
# surrogate crashes downstream consumers (stdout, Telegram ``utf16_len``, JSON).
|
||||||
# (``assistant_message.content``), not the sanitized copy stored in
|
# Scrub once where model text leaves the loop (#80366).
|
||||||
# history by ``build_assistant_message``. Any lone UTF-16 surrogate
|
|
||||||
# (U+D800–U+DFFF) in it crashes downstream consumers — oneshot stdout
|
|
||||||
# writes, Telegram's ``utf16_len`` length check, Signal formatting,
|
|
||||||
# JSON envelope encodes — on every provider (Ollama, NVIDIA NIM, …).
|
|
||||||
# Scrub once here, where model text leaves the conversation loop, so
|
|
||||||
# every delivery surface receives valid Unicode.
|
|
||||||
if isinstance(final_response, str):
|
if isinstance(final_response, str):
|
||||||
final_response = _sanitize_surrogates(final_response)
|
final_response = _sanitize_surrogates(final_response)
|
||||||
|
|
||||||
@@ -752,30 +591,27 @@ def finalize_turn(
|
|||||||
}
|
}
|
||||||
if agent._tool_guardrail_halt_decision is not None:
|
if agent._tool_guardrail_halt_decision is not None:
|
||||||
result["guardrail"] = agent._tool_guardrail_halt_decision.to_metadata()
|
result["guardrail"] = agent._tool_guardrail_halt_decision.to_metadata()
|
||||||
# Persistence failures already set failed=True + an explanation in
|
# Persistence failures already set failed=True; also stamp `error` so the gateway
|
||||||
# final_response; also stamp `error` so gateway surfaces status="error"
|
# surfaces status="error" (and desktop can toast) instead of a quiet complete frame.
|
||||||
# (and desktop can toast the cause) instead of a quiet complete frame.
|
|
||||||
if failed and str(_turn_exit_reason) == "session_persistence_failed":
|
if failed and str(_turn_exit_reason) == "session_persistence_failed":
|
||||||
result["error"] = final_response or (
|
result["error"] = final_response or (
|
||||||
"session storage could not be written — check the state database "
|
"session storage could not be written — check the state database "
|
||||||
"health (`hermes doctor`), then send your message again"
|
"health (`hermes doctor`), then send your message again"
|
||||||
)
|
)
|
||||||
# Machine-readable cause for the gateway/desktop: exactly
|
# Machine-readable cause for the gateway/desktop, exactly
|
||||||
# 'session_persistence_failed:<locked|compression|turn_lease|corrupt|replaced|disk|unknown>'.
|
# 'session_persistence_failed:<locked|compression|turn_lease|corrupt|...>'.
|
||||||
# Never clobber a failure_reason another path already stamped.
|
# Never clobber a failure_reason another path already stamped.
|
||||||
if "failure_reason" not in result:
|
if "failure_reason" not in result:
|
||||||
_cause = getattr(agent, "_last_persistence_error_cause", None)
|
_cause = getattr(agent, "_last_persistence_error_cause", None)
|
||||||
result["failure_reason"] = (
|
result["failure_reason"] = (
|
||||||
"session_persistence_failed:" + (_cause or "unknown")
|
"session_persistence_failed:" + (_cause or "unknown")
|
||||||
)
|
)
|
||||||
# Surface any post-loop cleanup failures so the caller can distinguish a
|
# Surface post-loop cleanup failures so the caller can tell a clean turn from one
|
||||||
# clean turn from one whose trajectory/session/resource teardown raised
|
# whose teardown raised; the response is returned either way (#8049).
|
||||||
# (the response is still returned either way — #8049).
|
|
||||||
if _cleanup_errors:
|
if _cleanup_errors:
|
||||||
result["cleanup_errors"] = _cleanup_errors
|
result["cleanup_errors"] = _cleanup_errors
|
||||||
# If a /steer landed after the final assistant turn (no more tool
|
# A /steer landing after the final assistant turn has no tool batch to drain into;
|
||||||
# batches to drain into), hand it back to the caller so it can be
|
# hand it back so it becomes the next user turn instead of being lost.
|
||||||
# delivered as the next user turn instead of being silently lost.
|
|
||||||
_leftover_steer = agent._drain_pending_steer()
|
_leftover_steer = agent._drain_pending_steer()
|
||||||
if _leftover_steer:
|
if _leftover_steer:
|
||||||
result["pending_steer"] = _leftover_steer
|
result["pending_steer"] = _leftover_steer
|
||||||
@@ -807,11 +643,9 @@ def finalize_turn(
|
|||||||
messages=messages,
|
messages=messages,
|
||||||
)
|
)
|
||||||
|
|
||||||
# Background memory/skill review — runs AFTER the response is delivered
|
# Background memory/skill review runs AFTER delivery so it never competes with the
|
||||||
# so it never competes with the user's task for model attention.
|
# user's task. Suppressed by skip_background_review (e.g. cron): the fork costs
|
||||||
# Suppressed when skip_background_review=True (e.g. cron) — review forks
|
# ~30K tokens / event with no human-in-the-loop benefit.
|
||||||
# spawn another AIAgent (~30K tokens / event) and cron sessions have no
|
|
||||||
# human-in-the-loop benefit from the review.
|
|
||||||
if (
|
if (
|
||||||
final_response
|
final_response
|
||||||
and not interrupted
|
and not interrupted
|
||||||
@@ -829,16 +663,10 @@ def finalize_turn(
|
|||||||
except Exception:
|
except Exception:
|
||||||
pass # Background review is best-effort
|
pass # Background review is best-effort
|
||||||
|
|
||||||
# Note: Memory provider on_session_end() + shutdown_all() are NOT
|
# Memory provider on_session_end()/shutdown_all() are NOT called here:
|
||||||
# called here — run_conversation() is called once per user message in
|
# run_conversation() runs once per message; CLI/gateway own session-end cleanup.
|
||||||
# multi-turn sessions. Shutting down after every turn would kill the
|
|
||||||
# provider before the second message. Actual session-end cleanup is
|
|
||||||
# handled by the CLI (atexit / /reset) and gateway (session expiry /
|
|
||||||
# _reset_session).
|
|
||||||
|
|
||||||
# Plugin hook: on_session_end
|
# Plugin hook: on_session_end — fired at the end of every run_conversation call.
|
||||||
# Fired at the very end of every run_conversation call.
|
|
||||||
# Plugins can use this for cleanup, flushing buffers, etc.
|
|
||||||
try:
|
try:
|
||||||
from hermes_cli.lifecycle import invoke_hook as _invoke_hook
|
from hermes_cli.lifecycle import invoke_hook as _invoke_hook
|
||||||
_invoke_hook(
|
_invoke_hook(
|
||||||
|
|||||||
+15
-45
@@ -1,28 +1,8 @@
|
|||||||
"""Per-attempt recovery bookkeeping for the conversation turn loop.
|
"""Per-attempt recovery bookkeeping (``TurnRetryState``) for the conversation turn loop.
|
||||||
|
|
||||||
The inner retry loop in ``run_conversation`` (``while retry_count <
|
Each one-shot recovery branch of the inner retry loop is guarded by a flag here so it
|
||||||
max_retries``) makes several distinct recovery attempts on a single model API
|
fires at most once per attempt. Loop-control (``retry_count``, ``max_retries``) stays
|
||||||
call: a credential-pool 429 retry, a per-provider OAuth refresh (codex,
|
as plain locals. Dependency-free so it imports without a cycle."""
|
||||||
anthropic, nous, copilot), a long-context compression restart, a length-
|
|
||||||
continuation restart, and a handful of format-recovery branches (thinking-
|
|
||||||
signature stripping, multimodal-tool-content stripping, llama.cpp grammar
|
|
||||||
fallback, image shrink, invalid-encrypted-content, 1M-beta header).
|
|
||||||
|
|
||||||
Each of those branches is guarded by a one-shot boolean so it fires at most
|
|
||||||
once per attempt. They used to be ~16 bare ``*_attempted`` / ``has_retried_*``
|
|
||||||
/ ``restart_with_*`` locals declared inline before the loop and threaded
|
|
||||||
through its 2,400-line body. ``TurnRetryState`` collapses them into one object
|
|
||||||
the loop mutates in place (``state.codex_auth_retry_attempted = True``), giving
|
|
||||||
the recovery bookkeeping a single named, testable home.
|
|
||||||
|
|
||||||
Loop-control variables (``retry_count``, ``max_retries``,
|
|
||||||
``max_compression_attempts``) intentionally stay as plain locals — they are the
|
|
||||||
``while`` mechanics, not recovery bookkeeping, and putting them on the object
|
|
||||||
would add indirection without clarifying anything.
|
|
||||||
|
|
||||||
This module is dependency-free so it can be unit-tested in isolation and
|
|
||||||
imported by the turn loop without an import cycle.
|
|
||||||
"""
|
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
@@ -33,11 +13,8 @@ from dataclasses import dataclass, fields
|
|||||||
class TurnRetryState:
|
class TurnRetryState:
|
||||||
"""One-shot recovery guards + restart signals for a single API-call attempt.
|
"""One-shot recovery guards + restart signals for a single API-call attempt.
|
||||||
|
|
||||||
A fresh instance is created for each iteration of the outer turn loop
|
A fresh instance is created per ``api_call_count`` iteration; each guard fires at
|
||||||
(once per ``api_call_count``). Each guard fires its recovery branch at most
|
most once, and ``restart_with_*`` signals tell the loop to rebuild and retry."""
|
||||||
once; the ``restart_with_*`` signals are read by the loop after the attempt
|
|
||||||
to decide whether to rebuild the request and retry.
|
|
||||||
"""
|
|
||||||
|
|
||||||
# ── Per-provider OAuth / credential refresh guards ───────────────────
|
# ── Per-provider OAuth / credential refresh guards ───────────────────
|
||||||
codex_auth_retry_attempted: bool = False
|
codex_auth_retry_attempted: bool = False
|
||||||
@@ -46,12 +23,8 @@ class TurnRetryState:
|
|||||||
nous_paid_entitlement_refresh_attempted: bool = False
|
nous_paid_entitlement_refresh_attempted: bool = False
|
||||||
copilot_auth_retry_attempted: bool = False
|
copilot_auth_retry_attempted: bool = False
|
||||||
# Copilot surfaces a stale/degraded credential as a 400
|
# Copilot surfaces a stale/degraded credential as a 400
|
||||||
# ``model_not_available_for_integrator`` / ``model_not_supported`` instead
|
# ``model_not_available_for_integrator`` / ``model_not_supported``, not a 401.
|
||||||
# of a clean 401 (e.g. a raw OAuth token seeded when the token exchange
|
# Single-shot forced re-exchange + rebuild, separate from the 401 guard.
|
||||||
# degraded at startup, routing the request to the restricted
|
|
||||||
# ``copilot-language-server`` integrator). Guard a single-shot forced
|
|
||||||
# re-exchange + client rebuild for that case, separate from the 401 guard
|
|
||||||
# so both can fire within one attempt if needed.
|
|
||||||
copilot_stale_cred_retry_attempted: bool = False
|
copilot_stale_cred_retry_attempted: bool = False
|
||||||
vertex_auth_retry_attempted: bool = False
|
vertex_auth_retry_attempted: bool = False
|
||||||
|
|
||||||
@@ -69,22 +42,19 @@ class TurnRetryState:
|
|||||||
has_retried_429: bool = False
|
has_retried_429: bool = False
|
||||||
|
|
||||||
# ── Auth-failure provider failover ───────────────────────────────────
|
# ── Auth-failure provider failover ───────────────────────────────────
|
||||||
# Set once we've escalated a persistent 401/403 (after the per-provider
|
# Set once a persistent 401/403 has been escalated to the fallback chain, so
|
||||||
# credential-refresh attempt above failed) to the fallback chain, so we
|
# we don't loop on the same auth failover within one attempt.
|
||||||
# don't loop on the same auth failover within one attempt.
|
|
||||||
auth_failover_attempted: bool = False
|
auth_failover_attempted: bool = False
|
||||||
|
|
||||||
# ── Restart signals (read by the outer loop after the attempt) ───────
|
# ── Restart signals (read by the outer loop after the attempt) ───────
|
||||||
restart_with_compressed_messages: bool = False
|
restart_with_compressed_messages: bool = False
|
||||||
restart_with_length_continuation: bool = False
|
restart_with_length_continuation: bool = False
|
||||||
# Set when a content-filter stream stall (e.g. MiniMax "new_sensitive")
|
# Set when a content-filter stream stall (e.g. MiniMax "new_sensitive") was
|
||||||
# has been escalated to the fallback chain: the partial-stream content
|
# escalated to the fallback chain: partial content was rolled back off
|
||||||
# was rolled back off ``messages`` and the loop should re-issue the API
|
# ``messages``; re-issue the call against the new provider (#32421).
|
||||||
# call against the newly-activated provider (#32421).
|
|
||||||
restart_with_rebuilt_messages: bool = False
|
restart_with_rebuilt_messages: bool = False
|
||||||
# A user correction cancelled the in-flight provider request. The outer
|
# A user correction cancelled the in-flight request: append a role-safe checkpoint +
|
||||||
# loop must append a role-safe checkpoint + user message, rebuild the API
|
# user message, rebuild the payload, and retry the same logical iteration.
|
||||||
# payload, and retry the same logical iteration.
|
|
||||||
restart_with_redirected_messages: bool = False
|
restart_with_redirected_messages: bool = False
|
||||||
|
|
||||||
def __iter__(self):
|
def __iter__(self):
|
||||||
|
|||||||
Reference in New Issue
Block a user