diff --git a/agent/client_lifecycle.py b/agent/client_lifecycle.py index fd1ae1f395..168e948074 100644 --- a/agent/client_lifecycle.py +++ b/agent/client_lifecycle.py @@ -755,8 +755,8 @@ class ClientLifecycleMixin: ) -> bool: if self.provider != "nous": return False - # Portal serves anthropic/* on the native Messages route, so either client kind may hold the - # expiring invoke JWT. + # Portal serves anthropic/* on the native Messages route, so either client kind may hold the expiring + # invoke JWT. if self.api_mode not in ("chat_completions", "anthropic_messages"): return False diff --git a/agent/interrupt_control.py b/agent/interrupt_control.py index e6cf998acf..f55cb61290 100644 --- a/agent/interrupt_control.py +++ b/agent/interrupt_control.py @@ -57,8 +57,7 @@ class InterruptControlMixin: return if not getattr(fence, "commit_in_flight", False): # No commit in flight — cancel_before_commit here WOULD cancel the pending commit; leave it to - # the - # destructive half. + # the destructive half. return cancel_before_commit = getattr( type(fence), "cancel_before_commit", None @@ -66,8 +65,7 @@ class InterruptControlMixin: if callable(cancel_before_commit): try: # A commit holds the fence lock through finish_commit: this blocks until it finishes and - # returns - # False WITHOUT setting _cancelled. + # returns False WITHOUT setting _cancelled. cancel_before_commit(fence) except Exception: logger.debug( @@ -90,8 +88,7 @@ class InterruptControlMixin: if callable(cancel_before_commit): try: # Marks the fence cancelled (or waits out a just-started commit) without touching the - # hard-stop - # Event, which was published at the claim edge. + # hard-stop Event, which was published at the claim edge. cancel_before_commit(fence) except Exception: logger.debug( @@ -116,8 +113,7 @@ class InterruptControlMixin: # activity stamp, or the stamp landed first and the abort declines without publishing. if require_generation is None: # No claim to race: publish WITHOUT the liveness lock. Bare AIAgent stand-ins in other suites - # lack - # the liveness seam and would AttributeError. + # lack the liveness seam and would AttributeError. _publish_interrupt_state() return True with self._liveness_activity_lock(): @@ -142,8 +138,7 @@ class InterruptControlMixin: if _redirect_lock is not None: with _redirect_lock: # The blocking in-flight-commit wait runs BEFORE the atomic claim edge (redirect lock still - # held); - # the destructive pending-commit cancel runs AFTER the claim survives (#99758 P1). + # held); the destructive pending-commit cancel runs AFTER the claim survives (#99758 P1). if hard_cancel: _wait_for_compression_commit() if not _consume_claim_and_publish_first_state(): diff --git a/agent/session_persistence.py b/agent/session_persistence.py index ac2808cc86..30eb76635a 100644 --- a/agent/session_persistence.py +++ b/agent/session_persistence.py @@ -53,11 +53,11 @@ def _is_ephemeral_scaffolding(msg: Any) -> bool: ) -# `_DB_PERSISTED_MARKER` (agent.context_compressor) — intrinsic "already written to SQLite" marker. An id(msg) dedup set can alias a freed dict's address onto -# a new message and silently skip persisting it; a marker on the dict cannot. The `_` prefix is mandatory: -# wire sanitizers strip `_`-prefixed keys. CONTRACT (#92231): the marker asserts the dict's CONTENT is -# durable as written — any in-place mutation that must persist MUST pop it (see turn_finalizer, -# context_compressor). +# `_DB_PERSISTED_MARKER` (agent.context_compressor) — intrinsic "already written to SQLite" marker. An id(msg) +# dedup set can alias a freed dict's address onto a new message and silently skip persisting it; a marker on +# the dict cannot. The `_` prefix is mandatory: wire sanitizers strip `_`-prefixed keys. CONTRACT (#92231): +# the marker asserts the dict's CONTENT is durable as written — any in-place mutation that must persist MUST +# pop it (see turn_finalizer, context_compressor). def _safe_session_filename_component(session_id: str) -> str: @@ -99,11 +99,9 @@ class SessionPersistenceMixin: msg = messages[idx] if isinstance(msg, dict) and msg.get("role") == "user": # A plain-text override must not replace native image/audio blocks; a list override is the - # clean - # multimodal payload and does. Preflight compaction may re-anchor this index at a message - # MERGED - # with the compaction summary — overwriting it would drop the summary (see the twin guard in - # _flush_messages_to_session_db_unlocked). + # clean multimodal payload and does. Preflight compaction may re-anchor this index at a + # message MERGED with the compaction summary — overwriting it would drop the summary (see the + # twin guard in _flush_messages_to_session_db_unlocked). if ( override is not None and not msg.get(COMPRESSED_SUMMARY_METADATA_KEY) @@ -249,8 +247,7 @@ class SessionPersistenceMixin: } # Bounded scan: skip the identity-matched prefix of the previous flush's snapshot. Every message - # in - # it already got its final disposition, and no live dict has its marker popped in place. + # in it already got its final disposition, and no live dict has its marker popped in place. _scan_start = 0 _prev_prefix = getattr(self, "_db_flush_scan_prefix", None) if isinstance(_prev_prefix, list): @@ -271,8 +268,8 @@ class SessionPersistenceMixin: if not isinstance(msg, dict): continue # Never write ephemeral scaffolding: the flush is append-only, so a mid-turn persist could - # commit a - # synthetic turn that the end-of-turn drop cannot un-write. Skip regardless of position. + # commit a synthetic turn that the end-of-turn drop cannot un-write. Skip regardless of + # position. if _is_ephemeral_scaffolding(msg): continue if msg.get(_DB_PERSISTED_MARKER): @@ -285,34 +282,30 @@ class SessionPersistenceMixin: role = msg.get("role", "unknown") content = msg.get("content") # api_content sidecar: exact bytes sent to the API when they differ from clean content, so - # replay - # reproduces the sent prefix byte-for-byte. + # replay reproduces the sent prefix byte-for-byte. _row_api_content = msg.get("api_content") if not isinstance(_row_api_content, str): _row_api_content = None _row_timestamp = msg.get("timestamp") # Apply the persist override to THIS row only. A list override replaces a noted payload; a - # text - # override must not erase an image/audio summary. Also match the staged CLI dict by identity — - # the close safety-net may flush a shortened snapshot whose turn index refers to the full - # history. + # text override must not erase an image/audio summary. Also match the staged CLI dict by + # identity — the close safety-net may flush a shortened snapshot whose turn index refers to + # the full history. pending_cli_message = getattr(self, "_pending_cli_user_message", None) is_current_turn_user = ( _ov_idx == _msg_idx or msg is pending_cli_message ) if is_current_turn_user and msg.get("role") == "user": # Preflight compaction may have re-anchored the index at a message MERGED with the - # compaction - # summary; overwriting it with the clean text would drop the summary from the durable - # transcript. + # compaction summary; overwriting it with the clean text would drop the summary from the + # durable transcript. if ( _ov_content is not None and (not isinstance(content, list) or isinstance(_ov_content, list)) and not msg.get(COMPRESSED_SUMMARY_METADATA_KEY) ): # Live content is what the wire sent, the override is the clean transcript; keep the - # sent bytes in - # api_content so replay matches the wire (#48677). + # sent bytes in api_content so replay matches the wire (#48677). if ( _row_api_content is None and isinstance(content, str) @@ -374,8 +367,7 @@ class SessionPersistenceMixin: "timestamp": _row_timestamp, "api_content": _row_api_content, # Standalone reference handoffs are always hidden so they never occupy the active user - # slot in - # retry/undo dispatch (#80622); merge-into-tail carriers keep prior visibility. + # slot in retry/undo dispatch (#80622); merge-into-tail carriers keep prior visibility. "display_kind": ( "hidden" if ( @@ -461,8 +453,7 @@ class SessionPersistenceMixin: if isinstance(e, CompressionSessionClosedError): # Compression race: another path rotated this session mid-write. Adopt the continuation tip # (get_compression_tip) ONLY when it is a different, live row, and retry exactly once; a - # second - # closed-parent write fails closed. tip == session_id means no continuation exists. + # second closed-parent write fails closed. tip == session_id means no continuation exists. if _adoption_budget > 0: old_id = self.session_id tip = None @@ -497,8 +488,7 @@ class SessionPersistenceMixin: _adoption_budget=0, ) # No live tip or budget exhausted: fail closed. The flag lets the turn explanation name - # compression - # rotation instead of misleading full-disk advice. + # compression rotation instead of misleading full-disk advice. self._compression_adoption_failed = True logger.warning("Session DB append_message failed: %s", e) return False diff --git a/run_agent.py b/run_agent.py index 2e1ad348ea..fb16873eb4 100644 --- a/run_agent.py +++ b/run_agent.py @@ -44,9 +44,8 @@ import threading import uuid import warnings from typing import List, Dict, Any, Optional, Callable -# `OpenAI` is a lazy proxy (SDK import costs ~240ms) that keeps the single `OpenAI(**kw)` call site -# and `patch("run_agent.OpenAI")` working. `fire` is imported only in __main__ so library imports never need -# it. +# `OpenAI` is a lazy proxy (SDK import costs ~240ms) that keeps the single `OpenAI(**kw)` call site and +# `patch("run_agent.OpenAI")` working. `fire` is imported only in __main__ so library imports never need it. from datetime import datetime from pathlib import Path @@ -461,8 +460,7 @@ class AIAgent( from hermes_cli.profiles import get_active_profile_name _profile_for_session = get_active_profile_name() # Persist the profile name explicitly, including "default": profile-keyed consumers treat NULL - # as - # unowned (#94724 backfill, #99222). + # as unowned (#94724 backfill, #99222). except Exception: _profile_for_session = None # Carry the live YOLO bypass into model_config: the row is created lazily on the first turn, so @@ -1109,8 +1107,7 @@ class AIAgent( if "glm" not in model_lower and provider_lower != "zai": return False base = self._base_url_lower - # Ollama Cloud (hosted service or :cloud proxy) forwards finish_reason - # faithfully — do not rewrite. + # Ollama Cloud (hosted service or :cloud proxy) forwards finish_reason faithfully — do not rewrite. if "ollama.com" in base or ":cloud" in model_lower: return False if "ollama" in base or ":11434" in base: @@ -1578,8 +1575,8 @@ class AIAgent( except Exception: pass - # 7. Free conversation history proactively (close() is the hard teardown; callers may still hold - # the closed agent). + # 7. Free conversation history proactively (close() is the hard teardown; callers may still hold the + # closed agent). try: self._session_messages = [] except Exception: @@ -2772,23 +2769,20 @@ class AIAgent( else: def _snapshot_worker(fence=None): # #76354 F3: the pooled worker must NEVER share the caller's live transcript — a late - # engine after a - # host timeout could rewrite it. Deep-snapshot on the worker; results publish only via an - # ADMITTED commit. + # engine after a host timeout could rewrite it. Deep-snapshot on the worker; results + # publish only via an ADMITTED commit. snapshot = copy.deepcopy(messages) result_msgs, result_prompt = _run( fence, target_messages=snapshot ) if result_msgs is snapshot: # No-op/abort returned the snapshot unchanged: hand back the ORIGINAL list so - # identity-based - # semantics keep working. + # identity-based semantics keep working. return messages, result_prompt return result_msgs, result_prompt # Resolve the fallback prompt lazily: an eager rebuild would raise before compress_context - # runs - # when _cached_system_prompt is unset and _build_system_prompt fails. + # runs when _cached_system_prompt is unset and _build_system_prompt fails. def _fallback_prompt(): cached = getattr(self, "_cached_system_prompt", None) if cached: @@ -2911,8 +2905,7 @@ class AIAgent( def _publish_new_fence(): # The stall-fallback retry (#78981) needs a fence the aborted attempt cannot veto; publish - # it on - # the slot hard_interrupt() reads. The finally restores the caller's fence. + # it on the slot hard_interrupt() reads. The finally restores the caller's fence. retry_fence = CompressionCommitFence() with fence_registration_lock: self._active_compression_commit_fence = retry_fence @@ -2945,8 +2938,7 @@ class AIAgent( ): return # Stamps land on the worker's snapshot first; mirror them onto the live lists by scoped - # identity. - # Timestamp-less repeated content is ambiguous, so every scoped match is stamped. + # identity. Timestamp-less repeated content is ambiguous, so every scoped match is stamped. for source_message in source_messages: if not ( isinstance(source_message, dict) @@ -3335,8 +3327,7 @@ class AIAgent( _durable_session_exists = _turn_db.get_session(session_id) is not None except Exception: # A locked / non-WAL read is not proof the row is absent; treating probe failure as - # "fresh" ran - # fail-open at the exact contention point (#84234). Acquire, or fail closed. + # "fresh" ran fail-open at the exact contention point (#84234). Acquire, or fail closed. logger.warning( "Could not check durable session before turn lease; " "will acquire rather than run without serialization", @@ -3348,8 +3339,7 @@ class AIAgent( and session_id and not getattr(self, "_persist_disabled", False) # 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. + # memory seed before the row exists — reloading would erase it. and _durable_session_exists # Check the concrete type: MagicMock-style shims accept any attribute without the protocol. and callable( @@ -3458,8 +3448,7 @@ class AIAgent( ) # 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). + # admission, and skip when acquisition was immediate (avoids a needless prompt-cache miss). if _lease_waited: latest_session_id = _turn_db.resolve_resume_session_id(session_id) if latest_session_id: @@ -3472,8 +3461,7 @@ class AIAgent( ) # Long turns outlive a fixed TTL: refresh in a daemon thread; holder-qualified UPDATE/DELETE - # fence - # a late refresher from a successor lease. + # fence a late refresher from a successor lease. durable_turn_lease_stop = threading.Event() _lease_refresh_interval = float( getattr(self, "_session_turn_lease_refresh_interval", 60.0) @@ -3498,8 +3486,8 @@ class AIAgent( def _interrupt_turn(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. + # this process no longer owns the session. Only the watchdog's stalls can be spuriously + # stale. nonlocal durable_turn_lease_interrupt_message with durable_turn_lease_activity_lock: if ( @@ -3553,8 +3541,7 @@ class AIAgent( ) except Exception: # Round-4 (#95663): fail closed — an exceptional path must not turn an unvalidated - # claim into - # unconditional abort authority. + # claim into unconditional abort authority. logger.debug( "Turn liveness abort interrupt raised; " "declining the abort", @@ -3682,9 +3669,8 @@ class AIAgent( with durable_turn_lease_activity_lock: durable_turn_lease_turn_active = True # Stamp the activity clock at turn entry (#95663): `_last_activity_ts` persists across - # turns, so - # without this the watchdog would measure idle from the PREVIOUS turn and abort a - # fresh one. + # turns, so without this the watchdog would measure idle from the PREVIOUS turn and + # abort a fresh one. self._touch_activity("starting new turn") durable_turn_lease_thread.start() if durable_turn_liveness_thread is not None: @@ -3707,8 +3693,7 @@ class AIAgent( # Post-loop relay/task finalization must not receive a late refresh interrupt. _stop_durable_turn_lease_refresher() # Interrupt clear is deferred until after thread join (outer finally) so a refresher - # firing - # between stop and join cannot leave an interrupt behind. + # firing between stop and join cannot leave an interrupt behind. terminal = result if isinstance(result, dict) else {} if terminal.get("interrupted") is True: relay_outcome = "cancelled"