diff --git a/agent/turn_facade.py b/agent/turn_facade.py index 17a3d32bbe..7b7d3e3f4e 100644 --- a/agent/turn_facade.py +++ b/agent/turn_facade.py @@ -29,8 +29,8 @@ class TurnFacadeMixin: persist_user_platform_id: Optional[str]=None, moa_config: Optional[dict[str, Any]]=None, ) -> Dict[str, Any]: """Forwarder — see ``agent.conversation_loop.run_conversation``.""" - # A review shares this session_id for cache parity: fence review startup or interrupt an admitted - # request and await its exit before opening live-turn instrumentation. + # A review shares this session_id for cache parity: fence review startup or interrupt + # an admitted request and await its exit before opening live-turn instrumentation. from agent.background_review import cancel_background_review_for_live_turn cancel_background_review_for_live_turn(self) @@ -64,15 +64,14 @@ class TurnFacadeMixin: else "" ) relay_lease = relay_turn = lease = None - # Scope tokens start None: early returns leave the try before the set_*() calls and the finally - # resets each one unconditionally. + # Scope tokens start None: early returns leave the try before the set_*() calls and + # the finally resets each one unconditionally. token = affinity_token = acct_token = None task_started = task_finished = False relay_outcome = "failed" try: - # Turn liveness for the deferred-review idle queue: first statement of the try so the finally's - # note_turn_finished balances every exit. + # First statement of the try so the finally's note_turn_finished balances every exit. _review_queue.note_turn_started() admission = admit_durable_turn_lease( self, session_id=session_id, relay_turn_id=relay_turn_id, task_context=task_context, @@ -95,24 +94,23 @@ class TurnFacadeMixin: relay_turn = relay_runtime.SESSION_COORDINATOR.begin_turn( relay_lease, turn_id=relay_turn_id, task_id=effective_task_id ) - # Minimal relay-runtime shims (tests, external) may lack the opt-out flag: default enabled. + # Minimal relay-runtime shims may lack the opt-out flag: default enabled. if getattr(relay_turn, "relay_enabled", True): start_task_run( **task_context, parent_session_id=getattr(self, "_parent_session_id", None) or "", ) task_started = True - # Publish the conversation id for ambient Nous Portal tagging: every LLM call in this turn - # (loop, compression, vision, MoA, review forks) inherits `conversation=`. + # Ambient Nous Portal tagging: every LLM call in this turn (loop, compression, + # vision, MoA, review forks) inherits `conversation=`; host-declared + # affinity scope falls back to it; accounting handles route aux usage to the session. token = set_conversation_context(self._conversation_root_id()) - # Routing/affinity scope the HOST declared; providers fall back to the id above when unset. affinity_token = set_affinity_scope(declared_conversation_scope_safe(self)) - # Session accounting handles so auxiliary calls record usage into session_model_usage. acct_token = set_accounting_context( getattr(self, "_session_db", None), getattr(self, "session_id", None) ) - # Keep the ContextVar scope local (tokens on the agent may be observed from another thread). + # Keep the ContextVar scope local (agent tokens may be observed from another thread). with bind_subagent_parent(self), scoped_runtime_main({}): try: if lease is not None: @@ -126,8 +124,8 @@ class TurnFacadeMixin: persist_user_platform_id=persist_user_platform_id, moa_config=moa_config, ) finally: - # Post-loop relay/task finalization must not receive a late refresh interrupt. The - # interrupt clear itself waits for the thread join in the outer finally. + # Post-loop relay/task finalization must not receive a late refresh interrupt; + # the interrupt clear itself waits for the thread join in the outer finally. if lease is not None: lease.stop_refresher() terminal = result if isinstance(result, dict) else {} @@ -168,11 +166,10 @@ class TurnFacadeMixin: if lease is not None: lease.stop_refresher() lease.join_threads() - # Clear any refresher interrupt fired between stop and join; must run AFTER join. - lease.clear_interrupt() + lease.clear_interrupt() # refresher interrupt between stop and join; AFTER join lease.release() - # Always clear mid-turn labels when the turn exits — including interrupted early - # returns that skip finalize_turn. Keep ts. + # Always clear mid-turn labels on exit — including interrupted early returns + # that skip finalize_turn. Keep ts. with suppress(Exception): self._reset_activity_labels_after_turn() if getattr(self, "_relay_pending_turn_id", None) == relay_turn_id: @@ -183,8 +180,7 @@ class TurnFacadeMixin: reset_conversation_context(token) if affinity_token is not None: reset_affinity_scope(affinity_token) - # Balance note_turn_started on every exit so the idle queue's live-turn count - # cannot leak. + # Balance note_turn_started so the idle queue's live-turn count cannot leak. with suppress(Exception): _review_queue.note_turn_finished() diff --git a/agent/turn_facade_lease.py b/agent/turn_facade_lease.py index d0e1581c18..7f66863151 100644 --- a/agent/turn_facade_lease.py +++ b/agent/turn_facade_lease.py @@ -44,11 +44,8 @@ class DurableTurnLease: return getattr(self.agent, "session_id", None) or self.session_id def build_threads(self) -> None: - """Create (not start) the refresher thread and, when configured, the liveness watchdog. - - Lease renewal is NOT evidence of progress; a silently stalled turn would renew forever, so the - watchdog (policy in ``agent/turn_liveness.py``) is wired with the commit/deactivate callbacks. - """ + """Create (not start) the refresher thread and, when configured, the liveness watchdog: + lease renewal is NOT evidence of progress; a silently stalled turn would renew forever.""" self.refresh_thread = threading.Thread( target=self.refresh_loop, name="session-turn-lease-refresh", daemon=True ) @@ -73,8 +70,8 @@ class DurableTurnLease: def start(self) -> None: with self._lock: self.turn_active = True - # Stamp the activity clock at turn entry: `_last_activity_ts` persists across turns, so without - # this the watchdog would measure idle from the PREVIOUS turn and abort a fresh one. + # Stamp the activity clock at turn entry: `_last_activity_ts` persists across turns, so + # without this the watchdog would measure idle from the PREVIOUS turn and abort a fresh one. self.agent._touch_activity("starting new turn") self.refresh_thread.start() if self.liveness_thread is not None: @@ -111,28 +108,26 @@ class DurableTurnLease: return self.turn_active def _interrupt_turn(self, message: str) -> None: - # Lease-loss interrupts fire UNCONDITIONALLY (no generation claim): a lost lease means this - # process no longer owns the session. Only the watchdog's stalls can be spuriously stale. - agent = self.agent + """Lease-loss interrupts fire UNCONDITIONALLY (no generation claim): a lost lease means + this process no longer owns the session. Only the watchdog's stalls can be spuriously stale.""" with self._lock: if self.stop.is_set() or not self.turn_active: return self.interrupt_message = message try: - agent.interrupt(message, hard_cancel=True) + self.agent.interrupt(message, hard_cancel=True) except Exception: - agent._interrupt_requested = True - agent._interrupt_message = message + self.agent._interrupt_requested = True + self.agent._interrupt_message = message def commit_liveness_abort(self, snapshot, message: str) -> bool: """Commit point for the watchdog's stall observation. Revalidates the observed ``(generation, timestamp)`` under the SAME lock ``_touch_activity`` - uses, so a turn that resumed while the stall was logged is never hard-cancelled. The - revalidated generation is carried into ``interrupt`` as ``require_generation``, which consumes - it with the first publication in ONE critical section. If ``interrupt`` raises, the abort - declines FAIL-CLOSED. Returns False when stale or already winding down. - """ + uses, so a turn that resumed while the stall was logged is never hard-cancelled; the + revalidated generation is consumed by ``interrupt(require_generation=...)`` with the first + publication in ONE critical section. If ``interrupt`` raises, the abort declines FAIL-CLOSED. + Returns False when stale or already winding down.""" agent = self.agent with agent._liveness_activity_lock(): current_generation = getattr(agent, "_turn_liveness_activity_generation", 0) @@ -148,10 +143,7 @@ class DurableTurnLease: message, hard_cancel=True, require_generation=current_generation ) except Exception: - # Fail closed: an exceptional path must not turn an unvalidated claim into abort authority. - logger.debug( - "Turn liveness abort interrupt raised; declining the abort", exc_info=True - ) + logger.debug("Turn liveness abort interrupt raised; declining the abort", exc_info=True) published = False if published is False: # Claim went stale between revalidation and the hammer: real progress landed. @@ -182,10 +174,9 @@ class DurableTurnLease: def refresh_loop(self) -> None: """Renew the lease every ``refresh_interval``; a miss or error interrupts the turn. - Long turns outlive the fixed TTL; the holder-qualified UPDATE fences a late refresher from a - successor lease. The façade's finally sets ``stop`` before releasing, so a holder-fenced miss - observed after stop is not a loss. - """ + The holder-qualified UPDATE fences a late refresher from a successor lease. The façade's + finally sets ``stop`` before releasing, so a holder-fenced miss observed after stop is not + a loss.""" while not self.stop.wait(self.refresh_interval): try: if self.db.refresh_session_turn_lease( @@ -197,16 +188,12 @@ class DurableTurnLease: logger.error( "Lost session turn lease while turn is active: %s", self._current_session_id() ) - self._interrupt_turn( - "Session turn lease lost; stopping to protect the transcript." - ) + self._interrupt_turn("Session turn lease lost; stopping to protect the transcript.") except Exception: if self.stop.is_set(): return logger.warning( - "Failed to refresh session turn lease: %s", - self._current_session_id(), - exc_info=True, + "Failed to refresh session turn lease: %s", self._current_session_id(), exc_info=True, ) self._interrupt_turn( "Session turn lease could not be refreshed; stopping to protect the transcript." @@ -245,15 +232,14 @@ def admit_durable_turn_lease( Mutates ``task_context["session_id"]`` and ``agent.session_id`` when the wait forced a resume-id reload. Returns an ``early_result`` (interrupted / timed out) instead of a lease when admission - fails; the caller returns it verbatim. - """ + fails; the caller returns it verbatim.""" db = getattr(agent, "_session_db", None) admission = TurnLeaseAdmission(conversation_history=conversation_history) if db is None or not session_id: return admission - # A fresh session id has no durable transcript to race over, and callers may supply an in-memory - # seed before the row exists — reloading would erase it. Check the concrete type: MagicMock-style - # shims accept any attribute without the protocol. + # A fresh session id has no durable transcript to race over, and callers may supply an + # in-memory seed before the row exists — reloading would erase it. Check the concrete type: + # MagicMock-style shims accept any attribute without the protocol. if ( getattr(agent, "_persist_disabled", False) or not _durable_session_exists(db, session_id) @@ -270,16 +256,12 @@ def admit_durable_turn_lease( def _on_wait(elapsed: float) -> None: nonlocal waited waited = True - if elapsed < 1.0: - agent._emit_status( - "⏳ Another Hermes process is using this session; " - "waiting for it to finish before starting your turn..." - ) - else: - agent._emit_status( - "⏳ Still waiting for the other Hermes process on " - f"this session ({int(elapsed)}s)..." - ) + agent._emit_status( + "⏳ Another Hermes process is using this session; " + "waiting for it to finish before starting your turn..." + if elapsed < 1.0 else + f"⏳ Still waiting for the other Hermes process on this session ({int(elapsed)}s)..." + ) if not db.acquire_session_turn_lease( session_id, holder, ttl_seconds=LEASE_TTL_SECONDS, wait_seconds=LEASE_WAIT_SECONDS, @@ -288,16 +270,16 @@ def admit_durable_turn_lease( admission.early_result = _lease_not_acquired_result(agent, session_id, conversation_history) return admission - # Assign only after admission so the finally cannot release a holder that never owned the row; - # persist paths read the agent attr so a late flush is fenced in the same SQLite transaction. + # Assign only after admission so the finally cannot release a holder that never owned the + # row; persist paths read the agent attr so a late flush is fenced in the same transaction. lease = DurableTurnLease(agent, db, session_id, holder) agent._active_session_turn_lease_holder = holder agent._active_session_turn_lease_ttl_seconds = LEASE_TTL_SECONDS try: if waited: agent._emit_status("Session is free; loading the latest transcript...") - # The holder may have compressed/rotated the session while we waited: reload only AFTER - # admission, and skip when acquisition was immediate (avoids a needless prompt-cache miss). + # The holder may have compressed/rotated the session while we waited: reload only + # AFTER admission; an immediate acquisition skips this (needless prompt-cache miss). latest_session_id = db.resolve_resume_session_id(session_id) if latest_session_id: agent.session_id = latest_session_id @@ -326,11 +308,10 @@ def _lease_not_acquired_result(agent, session_id: str, conversation_history) -> **base, "interrupted": True, } - interrupt_message = getattr(agent, "_interrupt_message", None) - if interrupt_message: - result["interrupt_message"] = interrupt_message - # The finalizer never runs on this early return; clear so a cached agent doesn't fail-close - # the next turn. + if getattr(agent, "_interrupt_message", None): + result["interrupt_message"] = agent._interrupt_message + # The finalizer never runs on this early return; clear so a cached agent doesn't + # fail-close the next turn. try: agent.clear_interrupt() except Exception: diff --git a/agent/turn_final_response.py b/agent/turn_final_response.py index 906308a1e5..fca3dbe7f6 100644 --- a/agent/turn_final_response.py +++ b/agent/turn_final_response.py @@ -72,15 +72,12 @@ def finish_text_response( result=result, ) - # No tool calls — final response. (Dropped tool-call recovery lives at - # the finalization chokepoint below so it catches every path.) final_response = assistant_message.content or "" - - # Unmute: _mute_post_response from a housekeeping tool turn must not - # silence empty-response warnings on the final response path. + # Unmute: _mute_post_response from a housekeeping tool turn must not silence + # empty-response warnings on the final response path. agent._mute_post_response = False - # Check if response only has think block with no actual content after it + # Think-block-only / empty content: recovery path. if not agent._has_content_after_think_block(final_response): _ev = recover_empty_response( agent, assistant_message, response, finish_reason, final_response=final_response, @@ -99,11 +96,10 @@ def finish_text_response( return _verdict("break") return _verdict("continue") - # Reset retry counter/signature on successful content agent._empty_content_retries = 0 agent._thinking_prefill_retries = 0 - # Surface the one-shot fallback switch notice before dropping the retry - # buffer so a provider/model switch stays visible on success. + # Surface the one-shot fallback switch notice before dropping the retry buffer so a + # provider/model switch stays visible on success. agent._emit_pending_fallback_notice() agent._clear_status_buffer() @@ -112,16 +108,13 @@ def finish_text_response( ) _ack_mode = intent_ack_continuation_mode(agent) - # Said-continue-but-stopped guard: no tool calls but the short reply - # TAILS with an announced next action. Fires mid-task too; reuses the - # SAME bounded continuation path and counter (max 2 per turn). + # Said-continue-but-stopped guard: no tool calls but the short reply TAILS with an + # announced next action. Reuses the SAME bounded continuation counter (max 2 per turn). _stall_continue_intent = ( bool(getattr(agent, "_stall_guards", True)) and agent.valid_tool_names and codex_ack_continuations < 2 - and trailing_continue_intent( - agent._strip_think_blocks(final_response or "") - ) + and trailing_continue_intent(agent._strip_think_blocks(final_response or "")) ) if _stall_continue_intent or ( _ack_mode != "off" @@ -144,8 +137,8 @@ def finish_text_response( agent._emit_interim_assistant_message(interim_msg) append_message(messages, {"role": "user", "content": _CODEX_ACK_CONTINUATION_NUDGE}) agent._session_messages = messages - # An acknowledgment is non-final: its text must not suppress - # iteration-limit summarization if the continuation exhausts budget. + # An acknowledgment is non-final: its text must not suppress iteration-limit + # summarization if the continuation exhausts budget. final_response = None return _verdict("continue") @@ -165,9 +158,8 @@ def finish_text_response( final_msg = agent._build_assistant_message(assistant_message, finish_reason) - # ── Dropped tool-call recovery (copilot/Claude) ──────── - # finish_reason="tool_calls" with empty tool_calls would end the turn - # unstarted; re-prompt (max 3 CONSECUTIVE stalls, reset per tool round). + # Dropped tool-call recovery (copilot/Claude): finish_reason="tool_calls" with empty + # tool_calls would end the turn unstarted; re-prompt (max 3 CONSECUTIVE stalls). if ( finish_reason == "tool_calls" and not assistant_message.tool_calls @@ -184,9 +176,8 @@ def finish_text_response( "↻ Model signaled a tool call but sent none — " f"re-prompting ({agent._dropped_toolcall_retries}/3)" ) - # Both halves of the re-prompt pair are ephemeral scaffolding; flag - # them so the flush never persists them and the finalization pop - # can strip an unanswered tail pair. + # Both halves of the re-prompt pair are ephemeral scaffolding: never persisted, + # and the finalization pop strips an unanswered tail pair. final_msg["_dropped_toolcall_nudge"] = True append_message(messages, final_msg) append_message(messages, { @@ -223,9 +214,8 @@ def finish_text_response( return _verdict("continue") append_message(messages, final_msg) - # Make the answer durable before leaving the loop; _DB_PERSISTED_MARKER - # keeps _persist_session idempotent. Failure must NOT abort the turn: - # _persist_session retries the write. (#81641) + # Make the answer durable before leaving the loop (_DB_PERSISTED_MARKER keeps + # _persist_session idempotent). Failure must NOT abort the turn: finalize retries. try: agent._flush_messages_to_session_db(messages, conversation_history) except Exception: diff --git a/agent/turn_finalizer.py b/agent/turn_finalizer.py index 9b49898045..5c0973973d 100644 --- a/agent/turn_finalizer.py +++ b/agent/turn_finalizer.py @@ -20,9 +20,7 @@ from agent.message_sanitization import _sanitize_surrogates # Verification-continuation nudges (verify-on-stop / pre_verify) must be stripped from # returned/live history to avoid role-alternation breaks; the assistant response is # real content and is not flagged. (#65919) -_VERIFICATION_CONTINUATION_FLAGS = ( - "_verification_stop_synthetic", "_pre_verify_synthetic" -) +_VERIFICATION_CONTINUATION_FLAGS = ("_verification_stop_synthetic", "_pre_verify_synthetic") _SENTENCE_END = {".", "!", "?", "。", "!", "?", "`", ")"} @@ -50,17 +48,13 @@ def _record_kanban_budget_exhausted( _conn, kanban_task, error=( - f"Iteration budget exhausted " - f"({api_call_count}/{max_iterations}) — " - "task could not complete within the allowed " - "iterations" + f"Iteration budget exhausted ({api_call_count}/{max_iterations}) — " + "task could not complete within the allowed iterations" ), outcome="timed_out", release_claim=True, end_run=True, - event_payload_extra={ - "budget_used": api_call_count, "budget_max": max_iterations - }, + event_payload_extra={"budget_used": api_call_count, "budget_max": max_iterations}, ) finally: with suppress(Exception): @@ -191,9 +185,8 @@ def _close_transcript_tail(agent, messages, final_response, interrupted, failed) final_response = _streamed _recovered_from_stream = True - # An interrupt can leave a tool result as the tail (no scaffolding flag rewinds - # it); close the sequence so strict providers don't see ``tool → user``. An - # explicit placeholder is used since final_response is usually empty (#48879). + # An interrupt can leave a tool result as the tail; close the sequence so strict + # providers don't see ``tool → user`` (placeholder: final_response is usually empty). if interrupted: from agent.message_sanitization import close_interrupted_tool_sequence close_interrupted_tool_sequence(messages, final_response) @@ -204,7 +197,6 @@ def _close_transcript_tail(agent, messages, final_response, interrupted, failed) if final_response and not interrupted: _tail = messages[-1] if messages else None if not isinstance(_tail, dict) or _tail.get("role") != "assistant": - # Append so the durable turn closes with the answer (#43849/#44100). append_message(messages, {"role": "assistant", "content": final_response}) elif ( _tail.get("content") != final_response @@ -219,8 +211,7 @@ def _close_transcript_tail(agent, messages, final_response, interrupted, failed) agent._db_flush_scan_prefix = None # Request is complete, so replace API-local voice/model/skill guidance with the - # clean user input before the durable snapshot; earlier flushes used the DB-only - # override as their messages were still needed (#48677 / #63766). + # clean user input before the durable snapshot (earlier flushes still needed them). _apply_override = getattr(agent, "_apply_persist_user_message_override", None) if callable(_apply_override): _apply_override(messages) @@ -460,12 +451,11 @@ def finalize_turn( turn_id=turn_id, original_user_message=original_user_message, messages=messages, ) - # Context engine observation hook (complements per-request select_context()): - # notify the engine the turn finished with the finalized transcript. Fail-open. + # Context engine observation hook: the turn finished with the finalized transcript. + # Fail-open. ``_last_turn_usage`` is the last response's canonical usage dict, or + # ``None`` on turns that never reached a provider response — by contract. try: from agent.conversation_loop import _notify_context_engine_turn_complete - # ``_last_turn_usage`` holds the last API response's canonical usage dict, or - # ``None`` on turns that never reached a provider response — by contract. _notify_context_engine_turn_complete( agent, messages, usage=getattr(agent, "_last_turn_usage", None), logger=logger, turn_id=turn_id, task_id=effective_task_id, api_call_count=api_call_count, @@ -474,9 +464,8 @@ def finalize_turn( except Exception as exc: logger.warning("on_turn_complete notification failed: %s", exc) - # Surrogate chokepoint: ``final_response`` may be RAW SDK content, and a lone UTF-16 - # surrogate crashes downstream consumers (stdout, Telegram ``utf16_len``, JSON). - # Scrub once where model text leaves the loop (#80366). + # Surrogate chokepoint: RAW SDK text with a lone UTF-16 surrogate crashes downstream + # consumers (stdout, Telegram ``utf16_len``, JSON); scrub once where it leaves the loop. if isinstance(final_response, str): final_response = _sanitize_surrogates(final_response) @@ -508,8 +497,7 @@ def finalize_turn( "estimated_cost_usd": agent.session_estimated_cost_usd, "cost_status": agent.session_cost_status, "cost_source": agent.session_cost_source, - # Requested service tier (from request_overrides.extra_body), for billing - # audits by callers like `hermes -z --usage-file`. + # Requested service tier, for billing audits (`hermes -z --usage-file`). "service_tier": ( (getattr(agent, "request_overrides", {}) or {}).get("extra_body") or {} ).get("service_tier"), @@ -518,9 +506,8 @@ def finalize_turn( if agent._tool_guardrail_halt_decision is not None: result["guardrail"] = agent._tool_guardrail_halt_decision.to_metadata() # Persistence failures already set failed=True; also stamp `error` so the gateway - # surfaces status="error" (and desktop can toast) instead of a quiet complete frame, - # plus the machine-readable cause, exactly - # 'session_persistence_failed:'. + # surfaces status="error" (desktop can toast) instead of a quiet complete frame, plus + # the machine-readable cause 'session_persistence_failed:'. if failed and str(_turn_exit_reason) == "session_persistence_failed": result["error"] = final_response or ( "session storage could not be written — check the state database " @@ -528,8 +515,7 @@ def finalize_turn( ) _cause = getattr(agent, "_last_persistence_error_cause", None) result["failure_reason"] = "session_persistence_failed:" + (_cause or "unknown") - # Surface post-loop cleanup failures so the caller can tell a clean turn from one - # whose teardown raised; the response is returned either way (#8049). + # Cleanup failures are surfaced, but the response is returned either way (#8049). if _cleanup_errors: result["cleanup_errors"] = _cleanup_errors # A /steer landing after the final assistant turn has no tool batch to drain into; @@ -541,8 +527,7 @@ def finalize_turn( if interrupted and agent._interrupt_message: result["interrupt_message"] = agent._interrupt_message agent.clear_interrupt() - # Clear stream callback so it doesn't leak into future calls. - agent._stream_callback = None + agent._stream_callback = None # don't leak into future calls # Skill trigger is checked NOW — based on how many tool iterations THIS turn used. _should_review_skills = ( @@ -561,7 +546,8 @@ def finalize_turn( # Background memory/skill review runs AFTER delivery so it never competes with the # user's task. Suppressed by skip_background_review (e.g. cron): the fork costs - # ~30K tokens / event with no human-in-the-loop benefit. Best-effort. + # ~30K tokens / event with no human-in-the-loop benefit. Best-effort; the review + # clones the snapshot structurally so its sanitizers can't reach the live transcript. if ( final_response and not interrupted @@ -569,8 +555,6 @@ def finalize_turn( and (_should_review_memory or _should_review_skills) ): with suppress(Exception): - # _spawn_background_review clones the snapshot structurally so the fork's - # in-place sanitizers can't reach the live transcript. agent._spawn_background_review( messages_snapshot=list(messages), review_memory=_should_review_memory, review_skills=_should_review_skills, diff --git a/agent/turn_liveness.py b/agent/turn_liveness.py index 7a006938f2..a1f5a552f8 100644 --- a/agent/turn_liveness.py +++ b/agent/turn_liveness.py @@ -109,8 +109,8 @@ class TurnLivenessWatchdog: self._deactivate_turn = deactivate_turn def make_thread(self) -> threading.Thread: - """Build the (not yet started) watcher thread; run_agent starts it at - turn entry, after the turn-active flag and activity clock are stamped.""" + """Build the (not yet started) watcher thread; started at turn entry, after the + turn-active flag and activity clock are stamped.""" return threading.Thread(target=self._watch, name="turn-liveness-watchdog", daemon=True) def _watch(self) -> None: @@ -120,15 +120,13 @@ class TurnLivenessWatchdog: return # turn no longer active if snapshot.idle_seconds < self._timeout_s: continue - # Observational only: the commit below can still veto the abort - # if progress resumed; the definitive settlement is published - # by _surface_committed_abort after commit + deactivate. + # Observational only: the commit below can still veto the abort if progress + # resumed; the definitive settlement is _surface_committed_abort. self._surface_stall(snapshot) message = f"Turn made no progress for {int(snapshot.idle_seconds)}s; aborting to release the session." if not self._commit_abort(snapshot, message): continue - # Stop renewing the lease so a wedge the interrupt cannot unwind - # expires via TTL instead of masking forever. + # Stop renewing the lease so a wedge the interrupt cannot unwind expires via TTL. self._deactivate_turn() self._surface_committed_abort(snapshot) return @@ -152,11 +150,8 @@ class TurnLivenessWatchdog: logger.debug(debug_msg, exc_info=True) def _surface_stall(self, snapshot: ActivitySnapshot) -> None: - """Log + UI-warn that recovery is beginning (not that it committed). - - Rate-limited per activity generation so repeatedly declined aborts - do not re-log every poll; a new generation re-arms the surface. - """ + """Log + UI-warn that recovery is beginning (not that it committed). Rate-limited per + activity generation so repeatedly declined aborts do not re-log every poll.""" generation = snapshot.generation if getattr(self, "_last_surfaced_generation", None) == generation: return diff --git a/agent/turn_summary.py b/agent/turn_summary.py index eada53ff2c..55d27fd6df 100644 --- a/agent/turn_summary.py +++ b/agent/turn_summary.py @@ -1,14 +1,10 @@ """Per-turn accounting for the interactive CLI (display-only, pure). -:class:`TurnSummaryCollector` rides the existing ``tool_progress_callback`` feed -(``tool.completed`` events carry the tool name and raw result) and tallies what a -turn did — no agent-loop state is threaded through. :func:`format_turn_summary` -renders a tally plus wall-clock duration into one dim line, e.g.:: - - ⋯ 12.4s · edited 2 files +18 -3 · read 4 files · ran 3 commands - -(ported from Claude Code's post-turn accounting line). :func:`format_token_flow` -is the spinner-side cumulative token readout (``↓ 1.2k tok``). +:class:`TurnSummaryCollector` rides the ``tool_progress_callback`` feed (``tool.completed`` +events carry the tool name and raw result) and tallies what a turn did — no agent-loop state +is threaded through. :func:`format_turn_summary` renders a tally plus wall-clock duration +into one dim line: ``⋯ 12.4s · edited 2 files +18 -3 · read 4 files · ran 3 commands``. +:func:`format_token_flow` is the spinner-side cumulative token readout (``↓ 1.2k tok``). """ from __future__ import annotations @@ -26,18 +22,14 @@ __all__ = [ ] -# Leading glyph: terminal chrome, deliberately not an emoji. -SUMMARY_PREFIX = "⋯" - +SUMMARY_PREFIX = "⋯" # terminal chrome, deliberately not an emoji # A tool-less turn faster than this is a plain chat reply: formatter returns "". _MIN_TOOLLESS_SECONDS = 2.0 - # Max "verb + count" segments before collapsing the rest into "+N more". _MAX_SEGMENTS = 4 - -# Tool name -> (verb, singular noun, plural noun). Past tense: printed after the -# turn. Unlisted tools (plugin/MCP) fall into a generic "called N tools" bucket. +# Tool name -> (verb, singular noun, plural noun), past tense. Unlisted tools (plugin/MCP) +# fall into a generic "called N tools" bucket. _VERB_GROUPS: dict[str, tuple[str, str, str]] = { "write_file": ("edited", "file", "files"), "patch": ("edited", "file", "files"), @@ -57,12 +49,9 @@ _VERB_GROUPS: dict[str, tuple[str, str, str]] = { "memory": ("updated", "memory", "memories"), } -# Verb group that carries file-edit line deltas (+X -Y) when known. -_EDIT_VERB = "edited" - +_EDIT_VERB = "edited" # verb group that carries file-edit line deltas (+X -Y) when known # Render order: edits first, then reads, then commands; others in first-seen order. _VERB_PRIORITY: tuple[str, ...] = ("edited", "read", "ran") - # Tools whose results may report a unified diff we can count lines from. _DIFF_RESULT_TOOLS = frozenset({"patch"}) diff --git a/agent/turn_tool_round.py b/agent/turn_tool_round.py index 7c507e1911..f2c7a65fb5 100644 --- a/agent/turn_tool_round.py +++ b/agent/turn_tool_round.py @@ -113,12 +113,10 @@ def run_tool_round( tc for tc in assistant_message.tool_calls if tc.function.name in agent.valid_tool_names ] + # Persist the tool-call turn before any tool side effects so resume sees the executed + # block if a destructive tool restarts Hermes. try: - # Persist the tool-call turn before any tool side effects so resume - # sees the executed block if a destructive tool restarts Hermes. - _tool_turn_persisted = agent._flush_messages_to_session_db( - messages, conversation_history - ) + _tool_turn_persisted = agent._flush_messages_to_session_db(messages, conversation_history) except Exception as exc: _tool_turn_persisted = False from hermes_state import classify_persistence_error @@ -131,9 +129,8 @@ def run_tool_round( ) if _tool_turn_persisted is False: - # Canonical append failed: never project the row or run tools from - # process-only state; break rather than retry the unpersisted turn. - # If the flush recorded no cause, the cause is genuinely unknown. + # Canonical append failed: never project the row or run tools from process-only + # state; break rather than retry. No recorded cause means genuinely unknown. if getattr(agent, "_last_persistence_error_cause", None) is None: agent._last_persistence_error_cause = "unknown" _turn_exit_reason = "session_persistence_failed" @@ -141,14 +138,13 @@ def run_tool_round( failed = True return _verdict("break") - # A UI must never observe an assistant/tool-call row that is only an - # in-memory projection: emit interim commentary after the DB append. + # A UI must never observe an assistant/tool-call row that is only an in-memory + # projection: emit interim commentary after the DB append. if not duplicate_previous_interim: agent._emit_interim_assistant_message(assistant_msg) - # Flush open streaming boxes before tools so early content doesn't wrap - # tool feed lines. Display callback only — TTS (_stream_callback) must - # NOT receive None (its end-of-stream marker). + # Flush open streaming boxes before tools so early content doesn't wrap tool feed + # lines. Display callback only — TTS (_stream_callback) must NOT receive None (EOS). if agent.stream_delta_callback: with suppress(Exception): agent.stream_delta_callback(None) @@ -156,8 +152,8 @@ def run_tool_round( agent._execute_tool_calls(assistant_message, messages, effective_task_id, api_call_count) if getattr(agent, "_incremental_persistence_failed", False): - # Tool result could not be made canonical: never send the in-memory - # result to the model or project later events from this turn. + # Tool result could not be made canonical: never send the in-memory result to + # the model or project later events from this turn. _turn_exit_reason = "session_persistence_failed" final_response = "" failed = True @@ -167,12 +163,10 @@ def run_tool_round( decision = agent._tool_guardrail_halt_decision _turn_exit_reason = "guardrail_halt" final_response = agent._toolguard_controlled_halt_response(decision) - agent._emit_status( - f"⚠️ Tool guardrail halted {decision.tool_name}: {decision.code}" - ) + agent._emit_status(f"⚠️ Tool guardrail halted {decision.tool_name}: {decision.code}") append_message(messages, {"role": "assistant", "content": final_response}) - # Emit the halt message so it isn't mistaken for a crash; the stream - # callback is still alive, so SSE/TUI clients see the explanation. + # Emit the halt so it isn't mistaken for a crash; the stream callback is still + # alive, so SSE/TUI clients see the explanation. if final_response: agent._safe_print(f"\n{final_response}\n") if agent.stream_delta_callback: @@ -183,13 +177,11 @@ def run_tool_round( # Reset per-turn retry counters so one truncation can't poison the turn. truncated_tool_call_retries = 0 - - # Defer the paragraph break: _fire_stream_delta() prepends one "\n\n" - # when real text arrives, so tool iterations don't stack blank lines. + # Defer the paragraph break: _fire_stream_delta() prepends one "\n\n" when real + # text arrives, so tool iterations don't stack blank lines. agent._stream_needs_break = True - - # Refund the iteration when the ONLY tool was execute_code (programmatic - # tool calling) — cheap RPC-style calls shouldn't eat the budget. + # Refund the iteration when the ONLY tool was execute_code (programmatic tool + # calling) — cheap RPC-style calls shouldn't eat the budget. if {tc.function.name for tc in assistant_message.tool_calls} == {"execute_code"}: agent.iteration_budget.refund() @@ -211,9 +203,8 @@ def run_tool_round( # Save session log incrementally (so progress is visible even if interrupted) agent._session_messages = messages - - # Touch activity so slow post-tool work plus a slow follow-up API call - # can't exceed the gateway inactivity timeout (HERMES_AGENT_TIMEOUT). + # Touch activity so slow post-tool work plus a slow follow-up API call can't exceed + # the gateway inactivity timeout (HERMES_AGENT_TIMEOUT). agent._touch_activity(f"tool results posted, continuing iteration #{api_call_count}") return _verdict("continue") @@ -235,40 +226,32 @@ def stage_tool_call_message( turn_content = assistant_message.content or "" - # A bare bracketed token (e.g. ``[memory]``) beside a function call is - # protocol scaffolding; persisting it lets the post-tool fallback replay - # it forever (#78148). - if ( - assistant_message.tool_calls and _STALE_MARKER_RE.fullmatch(turn_content.strip()) - ): + # A bare bracketed token (e.g. ``[memory]``) beside a function call is protocol + # scaffolding; persisting it lets the post-tool fallback replay it forever (#78148). + if assistant_message.tool_calls and _STALE_MARKER_RE.fullmatch(turn_content.strip()): logger.warning( "Discarding bare tool-call marker from assistant content: %s", turn_content ) turn_content = "" assistant_msg["content"] = "" - # Classify tools regardless of visible content: a substantive tool-only - # turn must invalidate any older housekeeping fallback. + # Classify tools regardless of visible content: a substantive tool-only turn must + # invalidate any older housekeeping fallback (so a two-turn-old housekeeping + # narration isn't attributed to the preceding tool turn), and clear the mute flag a + # prior housekeeping turn set, else _vprint suppresses this turn's tool progress. _all_housekeeping = all( tc.function.name in _HOUSEKEEPING_TOOLS for tc in assistant_message.tool_calls ) - - # Substantive tools clear any older fallback so a two-turn-old - # housekeeping narration isn't attributed to the preceding tool turn. if assistant_message.tool_calls and not _all_housekeeping: agent._last_content_with_tools = None agent._last_content_tools_all_housekeeping = False - # Also clear the mute flag a prior housekeeping turn may have set, - # else _vprint suppresses this turn's tool progress until the - # no-tool-call branch clears it. agent._mute_post_response = False - # Content + tool_calls in one turn: keep the content as a fallback final - # response in case the follow-up turn after tools is empty. + # Content + tool_calls in one turn: keep the content as a fallback final response in + # case the follow-up turn after tools is empty. Mute only when EVERY tool call is + # post-response housekeeping; substantive tools keep output on. if turn_content and agent._has_content_after_think_block(turn_content): agent._last_content_with_tools = turn_content - # Mute only when EVERY tool call is post-response housekeeping - # (memory, todo, skill_manage); substantive tools keep output on. agent._last_content_tools_all_housekeeping = _all_housekeeping if _all_housekeeping and agent._has_stream_consumers(): agent._mute_post_response = True @@ -281,18 +264,15 @@ def stage_tool_call_message( # final-response path). Tool calls after a prefill recovery reset the prefill # counter, so each tool-call success is a fresh start, not a cumulative burn. _had_prefill = False - while ( - messages and isinstance(messages[-1], dict) and messages[-1].get("_thinking_prefill") - ): + while messages and isinstance(messages[-1], dict) and messages[-1].get("_thinking_prefill"): messages.pop() _had_prefill = True if _had_prefill: agent._thinking_prefill_retries = 0 agent._empty_content_retries = 0 - # Re-arm the post-tool nudge so it can fire on a LATER tool round. + # Re-arm the post-tool nudge so it can fire on a LATER tool round; a landed tool call + # recovers any dropped-tool-call stall, so refresh that budget per stall. agent._post_tool_empty_retried = False - # A landed tool call recovers any dropped-tool-call stall; refresh that - # budget so it guards each stall independently, not the whole run. agent._dropped_toolcall_retries = 0 previous_msg = messages[-1] if messages else None diff --git a/agent/turn_usage.py b/agent/turn_usage.py index d1c158cff5..425c9e1076 100644 --- a/agent/turn_usage.py +++ b/agent/turn_usage.py @@ -79,18 +79,15 @@ def record_response_usage( compressor.update_from_response({}) return ResponseUsageOutcome(compression_attempts=compression_attempts, rearmed=rearmed) - canonical_usage = normalize_usage( - response.usage, provider=agent.provider, api_mode=agent.api_mode - ) - # Aggregator-only usage kept for pricing: advisor tokens are priced - # at each advisor's OWN model rate and added as dollars below. + canonical_usage = normalize_usage(response.usage, provider=agent.provider, api_mode=agent.api_mode) + # Aggregator-only usage kept for pricing: advisor tokens are priced at each advisor's + # OWN model rate and added as dollars below. aggregator_usage = canonical_usage _moa_client, canonical_usage, _moa_ref_cost = _fold_moa_usage(agent, canonical_usage) prompt_tokens = canonical_usage.prompt_tokens completion_tokens = canonical_usage.output_tokens total_tokens = canonical_usage.total_tokens - # Forward canonical token + cache buckets for context engines; - # legacy keys stay for back-compat. + # Canonical token + cache buckets for context engines; legacy keys stay for back-compat. usage_dict = { "prompt_tokens": prompt_tokens, "completion_tokens": completion_tokens, @@ -101,24 +98,21 @@ def record_response_usage( "cache_write_tokens": canonical_usage.cache_write_tokens, "reasoning_tokens": canonical_usage.reasoning_tokens, } - # Capture the boundary latch before update_from_response() consumes - # it: only the real prompt count right after a compaction rearms the - # budget. + # Capture the boundary latch before update_from_response() consumes it: only the real + # prompt count right after a compaction rearms the budget. _completed_compaction_pending = bool( getattr(compressor, "_verify_compaction_cleared_threshold", False) ) compressor.update_from_response(usage_dict) - # Usage-anchored accounting: snapshot exact provider usage against - # the durable transcript; main-loop ONLY. MoA uses pre-fold - # aggregator usage. + # Usage-anchored accounting: snapshot exact provider usage against the durable + # transcript (main-loop ONLY; MoA uses pre-fold aggregator usage). The display meter + # anchors on the turn's FIRST response: later same-turn responses inflate + # prompt_tokens with replayed thinking. Display-only; compression math uses real usage. _new_anchor = capture_usage_anchor( aggregator_usage.prompt_tokens, aggregator_usage.output_tokens, messages ) if _new_anchor is not None: agent._usage_anchor = _new_anchor - # Anchor the display meter on the turn's FIRST response: - # later same-turn responses inflate prompt_tokens with replayed - # thinking. Display-only; compression math uses real usage. if api_call_count == 1: agent._turn_base_usage_anchor = _new_anchor _compression_threshold = int(getattr(compressor, "threshold_tokens", 0) or 0) @@ -135,13 +129,11 @@ def record_response_usage( max_compression_attempts, ) compression_attempts = 0 - # Confirmed recovery also clears the loop's stale insufficient-progress - # verdict (``_preflight_compression_blocked``), else it stays armed all - # turn and a later pressure spike grows unchecked. + # Confirmed recovery also clears the loop's stale insufficient-progress verdict + # (``_preflight_compression_blocked``), else a later pressure spike grows unchecked. rearmed = True - # Stash canonical usage for on_turn_complete() (same shape as - # update_from_response); keep the latest call's — last request. + # Stash canonical usage for on_turn_complete(); keep the latest call's. agent._last_turn_usage = dict(usage_dict) # Persist only provider-confirmed context lengths, not probe tiers. @@ -183,9 +175,8 @@ def record_response_usage( api_duration, _cache_pct, ) - # MoA: agent.model/provider are the virtual preset/"moa" with no - # pricing entry, silently dropping aggregator spend. Price at the - # REAL model/provider from the MoA client's aggregator slot. + # MoA: agent.model/provider are the virtual preset/"moa" with no pricing entry, silently + # dropping aggregator spend. Price at the REAL model/provider from the aggregator slot. _agg_cost_model, _agg_cost_provider, _agg_cost_base_url = agent.model, agent.provider, agent.base_url _agg_slot = getattr(_moa_client, "last_aggregator_slot", None) if _moa_client is not None else None if _agg_slot and _agg_slot.get("model"): @@ -214,18 +205,16 @@ def record_response_usage( agent.session_cost_status = cost_result.status agent.session_cost_source = cost_result.source - # Persist per-call token deltas for any session_id so non-CLI runs - # can't lose accounting; gateway/session-store writes use absolute - # totals and safely overwrite these deltas. + # Persist per-call token deltas for any session_id so non-CLI runs can't lose + # accounting; gateway/session-store writes use absolute totals and safely overwrite + # these deltas. Enqueued, not written (a cold state.db UPDATE here stalled the tool + # loop); drained at finalize via _persist_session. if agent._session_db and agent.session_id: try: - # Ensure the row exists: under concurrent SQLite load the - # initial _ensure_db_session() may fail, and UPDATE on a - # missing row silently affects 0 rows. + # Ensure the row exists: under concurrent SQLite load the initial + # _ensure_db_session() may fail, and UPDATE on a missing row affects 0 rows. if not agent._session_db_created: agent._ensure_db_session() - # Enqueued, not written: a cold state.db UPDATE here stalled - # the tool loop. Drained at finalize via _persist_session. agent._session_db.queue_token_counts( agent.session_id, input_tokens=canonical_usage.input_tokens, @@ -243,8 +232,7 @@ def record_response_usage( model=agent.model, api_call_count=1, ) - except Exception as e: - # Log failures — silent loss here undercounts analytics. + except Exception as e: # silent loss here undercounts analytics logger.debug( "Token persistence failed (session=%s, tokens=%d): %s", agent.session_id, total_tokens, e, @@ -253,9 +241,8 @@ def record_response_usage( if agent.verbose_logging: logging.debug(f"Token usage: prompt={usage_dict['prompt_tokens']:,}, completion={usage_dict['completion_tokens']:,}, total={usage_dict['total_tokens']:,}") - # Report cache stats for any provider that returns - # ``prompt_tokens_details.cached_tokens``, not only when we inject - # cache_control markers. ``canonical_usage`` is already normalised. + # Report cache stats for any provider that returns ``prompt_tokens_details.cached_tokens``, + # not only when we inject cache_control markers. cached = canonical_usage.cache_read_tokens written = canonical_usage.cache_write_tokens prompt = usage_dict["prompt_tokens"]