diff --git a/agent/compression_facade.py b/agent/compression_facade.py index c043fd9e5c..fdc3f6847f 100644 --- a/agent/compression_facade.py +++ b/agent/compression_facade.py @@ -66,8 +66,7 @@ def _report_compression_timeout( touch("context compression timed out", provenance=ActivityProvenance.AGENT_COMPRESSION_TIMEOUT) except Exception: logger.debug("compress_context timeout activity touch failed", exc_info=True) - # Same timeout cooldown ladder as summary-LLM timeouts: avoid re-burning the - # full idle budget every turn. + # Same timeout cooldown ladder as summary-LLM timeouts: avoid re-burning the full idle budget every turn. compressor = getattr(agent, "context_compressor", None) record = getattr(compressor, "record_timeout_failure", None) if compressor is not None else None if callable(record): @@ -133,10 +132,9 @@ def _run_under_progress_timeout( ): """Run ``run(fence, target_messages=snapshot)`` on the pool under the progress-aware timeout. - The pooled worker must NEVER share the caller's live transcript — a late engine after a - host timeout could rewrite it. It deep-snapshots on the worker and publishes only via an - ADMITTED commit; a no-op/abort returns the snapshot unchanged, so the ORIGINAL list is - handed back to keep identity semantics. + The pooled worker must NEVER share the caller's live transcript — a late engine after a host timeout could rewrite + it. It deep-snapshots on the worker and publishes only via an ADMITTED commit; a no-op/abort returns the snapshot + unchanged, so the ORIGINAL list is handed back to keep identity semantics. """ from agent.conversation_compression import CompressionCommitFence, run_compress_context_with_progress_timeout def _snapshot_worker(fence=None): @@ -314,8 +312,7 @@ class CompressionFacadeMixin: vars(self).pop("_active_compression_commit_fence", None) else: self._active_compression_commit_fence = previous_fence - # Restore whatever the caller had, so a compaction never leaks its - # tag into the surrounding scope. + # Restore whatever the caller had, so a compaction never leaks its tag into the surrounding scope. if token is not None: reset_conversation_context(token) if affinity_token is not None: diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 91a260f587..3fa927d2ac 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -53,9 +53,8 @@ def _swallow(message: str, *, exc_info: bool = False): logger.debug(message, exc) -# Terminal outcomes from host/hygiene timeout or cooldown writers. Detached -# heartbeat workers must not clobber these (timeout unobservable). Seeing one -# latches the heartbeat silent so a later UNKNOWN rewrite can't re-arm a zombie. +# Terminal outcomes from host/hygiene timeout or cooldown writers. Detached heartbeat workers must not clobber these +# (timeout unobservable). Seeing one latches the heartbeat silent so a later UNKNOWN rewrite can't re-arm a zombie. _TERMINAL_COMPRESSION_PROVENANCES = frozenset( {ActivityProvenance.AGENT_COMPRESSION_TIMEOUT, ActivityProvenance.AGENT_COMPRESSION_COOLDOWN} ) @@ -64,9 +63,8 @@ _TERMINAL_COMPRESSION_PROVENANCES = frozenset( # timeout-ladder rung (60s), not the 600s summary-provider cooldown. _SPLIT_FAILURE_COOLDOWN_SECONDS = 60 -# Marker tui_gateway/server.py::_status_update matches to tag kind="compacting" -# for drivers' "Summarizing…" UI. Keep the phrase intact when rewording. Idle/ -# preflight/retry lines lack it; is_compaction_progress_status covers those. +# Marker tui_gateway/server.py::_status_update matches to tag kind="compacting" for drivers' "Summarizing…" UI. Keep +# the phrase intact when rewording. Idle/preflight/retry lines lack it; is_compaction_progress_status covers those. COMPACTION_STATUS_MARKER = "Compacting context" COMPACTION_STATUS = f"🗜️ {COMPACTION_STATUS_MARKER} — summarizing earlier conversation so I can continue..." @@ -76,9 +74,8 @@ COMPACTION_DONE_STATUS = "✓ Context compaction complete — continuing turn... def _strip_marker_for_comparison(msgs: Any) -> Any: """Copy ``msgs`` with the ``_db_persisted`` marker removed for no-op comparison. - Live dicts carry the marker while ``compress()`` output is swept, so a raw - ``==`` would misclassify an identical no-op copy as progress. Non-list inputs - and non-dict entries pass through unchanged. + Live dicts carry the marker while ``compress()`` output is swept, so a raw ``==`` would misclassify an identical + no-op copy as progress. Non-list inputs and non-dict entries pass through unchanged. """ from agent.context_compressor import _DB_PERSISTED_MARKER if not isinstance(msgs, list): @@ -116,17 +113,15 @@ COMPRESSION_RETRY_CONTEXT_REDUCED_STATUS_TEMPLATE = ( "🗜️ Context reduced to {new_ctx:,} tokens (was {old_ctx:,}), retrying..." ) -# FAILURE-class notice: compression blocked, so the session grows until the -# provider limit kills it. Must stay visible on gateways: never add it to -# ROUTINE_COMPRESSION_STATUS_SAMPLES or _TELEGRAM_NOISY_STATUS_RE. +# FAILURE-class notice: compression blocked, so the session grows until the provider limit kills it. Must stay visible +# on gateways: never add it to ROUTINE_COMPRESSION_STATUS_SAMPLES or _TELEGRAM_NOISY_STATUS_RE. CONTEXT_OVERFLOW_BLOCKED_WARNING_TEMPLATE = ( "⚠ Context is over the compression threshold (~{tokens:,} tokens >= {threshold:,}) " "but compression is currently blocked ({reason}). The model may stop responding. Run /new to start a fresh " "session or /compress to retry immediately." ) -# Formatted from the same constants the emission sites use, so noise-filter -# tests exercise the ACTUAL wording. +# Formatted from the same constants the emission sites use, so noise-filter tests exercise the ACTUAL wording. ROUTINE_COMPRESSION_STATUS_SAMPLES = ( COMPACTION_STATUS, COMPACTION_DONE_STATUS, @@ -143,10 +138,9 @@ ROUTINE_COMPRESSION_STATUS_SAMPLES = ( def is_compaction_progress_status(text: str | None) -> bool: """True for in-progress auto-compaction lifecycle lines (not the done edge). - The gateway re-tags matches as ``kind="compacting"`` for the whole pause; - matching only the marker left idle/preflight/retry lines looking hung. - ``COMPACTION_DONE_STATUS`` is emitted as ``kind="compacted"`` and must not - match here. + The gateway re-tags matches as ``kind="compacting"`` for the whole pause; matching only the marker left + idle/preflight/retry lines looking hung. ``COMPACTION_DONE_STATUS`` is emitted as ``kind="compacted"`` and must + not match here. """ if not isinstance(text, str): return False @@ -170,9 +164,8 @@ def is_compaction_progress_status(text: str | None) -> bool: def _refresh_agent_tool_definitions(agent) -> bool: """Rebuild agent.tools at the compaction commit boundary. - Forever-sessions never restart, so this is the only moment config changes - reach the frozen dynamic tool schemas; the prompt cache is already invalid. - Delegates to refresh_agent_mcp_tools in content_aware mode (swaps on schema + Forever-sessions never restart, so this is the only moment config changes reach the frozen dynamic tool schemas; + the prompt cache is already invalid. Delegates to refresh_agent_mcp_tools in content_aware mode (swaps on schema CONTENT change). Returns True when tools were added. Never raises. """ from tools.mcp_tool import refresh_agent_mcp_tools @@ -213,9 +206,8 @@ def _snapshot_compressor_attempt_state(compressor: Any) -> dict[str, Any]: return copy.deepcopy(selected) -# Attempt ownership: stall-fallback detaches a timed-out worker and reuses the -# compressor, so its late unwind could restore a stale snapshot or clear the -# fallback's cancel check. Generation guards ATTRIBUTE writes; fence, COMMITs. +# Attempt ownership: stall-fallback detaches a timed-out worker and reuses the compressor, so its late unwind could +# restore a stale snapshot or clear the fallback's cancel check. Generation guards ATTRIBUTE writes; fence, COMMITs. _COMPRESSOR_ATTEMPT_LOCK = threading.Lock() @@ -256,8 +248,7 @@ def _install_compression_cancelled_check(compressor: Any, check: Any, generation def _clear_compression_cancelled_check_if_owner(compressor: Any, generation: int) -> bool: """Clear the cancellation consult only when *generation* installed it. - Prevents a detached late primary from tearing down a newer fallback's - callback. Returns True when cleared. + Prevents a detached late primary from tearing down a newer fallback's callback. Returns True when cleared. """ with _COMPRESSOR_ATTEMPT_LOCK: owner = getattr(compressor, "_compression_cancelled_check_owner", None) @@ -375,8 +366,7 @@ def _capture_authoritative_cooldown_under_lease( return None, None if not callable(raw_reader): return False, None - # Read the raw persisted row: the active getter filters expired rows and is not - # a lossless rollback snapshot. + # Read the raw persisted row: the active getter filters expired rows and is not a lossless rollback snapshot. durable_state = raw_reader(session_db, session_id) if not isinstance(durable_state, dict): raise TypeError("raw compression cooldown snapshot must be a mapping") @@ -569,10 +559,9 @@ class CompressionCommitFence: def revoke_commit_admission(self) -> None: """Revoke FUTURE commit admission without blocking on the fence lock. - An in-flight commit is never abandoned, but ``begin_commit`` re-checks the - flag under the lock so no new commit is admitted. The lease release must not - run mid-commit (a second compressor could interleave): released now if the - lock is free, else deferred to ``finish_commit``/refusal (holder-qualified). + An in-flight commit is never abandoned, but ``begin_commit`` re-checks the flag under the lock so no new + commit is admitted. The lease release must not run mid-commit (a second compressor could interleave): released + now if the lock is free, else deferred to ``finish_commit``/refusal (holder-qualified). """ self._admission_revoked = True if self._lock.acquire(blocking=False): @@ -605,8 +594,7 @@ class CompressionCommitFence: def register_cancelled_lock_release(self, release: Callable[[], None]) -> bool: """Publish the timed-out worker's holder-qualified lock release. - Returns whether cleanup was already requested; in that race the release runs - synchronously before returning. + Returns whether cleanup was already requested; in that race the release runs synchronously before returning. """ with self._lock_release_guard: self._cancelled_lock_release = release @@ -663,9 +651,8 @@ _CANCELLED_WORKER_TEARDOWN_GRACE_SECONDS = 5.0 def _join_cancelled_worker(future: Any, grace_seconds: float) -> bool: """Best-effort bounded join of a fence-cancelled compression worker. - Returns True when the future settled within ``grace_seconds`` (thread provably - exited); False for a still-running worker, which the caller must treat as an - orphan behind the poison fence. + Returns True when the future settled within ``grace_seconds`` (thread provably exited); False for a still-running + worker, which the caller must treat as an orphan behind the poison fence. """ try: grace = max(float(grace_seconds), 0.0) @@ -835,10 +822,9 @@ def _record_stall_interrupted_backoff( def resolve_compression_fallback_route() -> Optional[dict]: """Return the first usable ``auxiliary.compression.fallback_chain`` entry. - The aux client applies the chain only from its exception handler, so a silent - stall never reaches it; this pins the route onto one bounded retry instead. - Only the first complete entry: if it errors, the aux client's own exception - path walks the rest. ``None`` when none is usable (skip compression). + The aux client applies the chain only from its exception handler, so a silent stall never reaches it; this pins + the route onto one bounded retry instead. Only the first complete entry: if it errors, the aux client's own + exception path walks the rest. ``None`` when none is usable (skip compression). """ try: from agent.auxiliary_client import _fallback_entry_api_key, _get_auxiliary_task_config @@ -971,8 +957,7 @@ def _await_worker_within_budget( remaining_ceiling = ceiling - waited if remaining_ceiling <= 0: return False, None - # Charge idle budget from LAST PROGRESS, not slice start, or silence could - # approach 2x the budget. + # Charge idle budget from LAST PROGRESS, not slice start, or silence could approach 2x the budget. since_progress = fence.seconds_since_progress() wait_slice = min(max(idle - since_progress, 0.005), remaining_ceiling) try: @@ -1005,8 +990,7 @@ def _await_in_flight_commit( waited = time.monotonic() - wait_started remaining = ceiling - waited if remaining <= 0: - # Bounded increments so each overrun window is visible in logs rather than one - # silent unbounded block. + # Bounded increments so each overrun window is visible in logs rather than one silent unbounded block. remaining = min(_COMMIT_OVERRUN_WAIT_SLICE_SECONDS, max(ceiling, 0.05)) overrun_reports += 1 log = logger.warning if overrun_reports <= 2 else logger.error @@ -1087,10 +1071,9 @@ def run_compress_context_with_progress_timeout( ) -> Tuple[list, str]: """Run ``worker(fence)`` under a sync progress-aware (idle + ceiling) timeout. - Budgets bound the PRE-commit phase only: an admitted commit always completes - (overrun logged, surfaced once via ``on_commit_overrun``). A pre-commit cancel - returns ``(messages, system_prompt_fallback)`` (lazy callable), detaching the - worker; a stall first retries the chain once on ``new_fence``, then on_timeout + Budgets bound the PRE-commit phase only: an admitted commit always completes (overrun logged, surfaced once via + ``on_commit_overrun``). A pre-commit cancel returns ``(messages, system_prompt_fallback)`` (lazy callable), + detaching the worker; a stall first retries the chain once on ``new_fence``, then on_timeout """ if idle_timeout_seconds <= 0: raise ValueError( @@ -1340,9 +1323,8 @@ def _restore_messages_snapshot(messages: list, snapshot: Optional[list]) -> None def _restore_prune_rearm_tokens(compressor: Any, snapshot: dict) -> None: """Restore ONLY the prune runway from the attempt snapshot. - compress() zeroes it in memory while the durable copy only clears on a - successful commit; a kept transcript keeps its cached prefix, and 0 would let - the next prune break that cache. + compress() zeroes it in memory while the durable copy only clears on a successful commit; a kept transcript keeps + its cached prefix, and 0 would let the next prune break that cache. """ if "_proactive_prune_rearm_tokens" in snapshot: compressor._proactive_prune_rearm_tokens = snapshot["_proactive_prune_rearm_tokens"] @@ -1421,9 +1403,8 @@ def context_compression_timed_out(agent: Any) -> bool: def _automatic_compression_gate_blocks(agent: Any, bypass_cooldown: bool, *, include_cooldown: bool = True) -> bool: """Refresh durable guards, then evaluate the compressor's automatic breaker gate. - ``bypass_cooldown`` ignores the cooldown when the gate accepts ``ignore_cooldown`` - (engines predating it get the legacy no-argument call). When blocked, the - transient-block signal is published for automatic-path consumers. + ``bypass_cooldown`` ignores the cooldown when the gate accepts ``ignore_cooldown`` (engines predating it get the + legacy no-argument call). When blocked, the transient-block signal is published for automatic-path consumers. """ compressor = agent.context_compressor _refresh_persisted_compression_guards(compressor, include_cooldown=include_cooldown) @@ -1445,9 +1426,8 @@ def _automatic_compression_gate_blocks(agent: Any, bypass_cooldown: bool, *, inc def compression_blocked_transiently(agent: Any) -> bool: """Type-pinned read of the transient-block signal. - Set when an automatic pass no-ops on a TRANSIENT guard (summary-failure - cooldown or structural backoff). Consumers must defer, not count it toward - ``compression_exhausted``, or an overflow auto-reset wipes a session that was + Set when an automatic pass no-ops on a TRANSIENT guard (summary-failure cooldown or structural backoff). Consumers + must defer, not count it toward ``compression_exhausted``, or an overflow auto-reset wipes a session that was merely cooling down. The permanent ``ineffective`` breaker never sets it. """ _sig = getattr(agent, "_compression_blocked_transient", None) @@ -1495,9 +1475,8 @@ def _adopt_live_compression_child( ) -> Optional[List[Dict[str, Any]]]: """Move a stale compression contender onto the live continuation tip. - Resolve and load first, then mutate the agent, so ambiguous lineage or an - unreadable handoff fails closed. Uses the transitive ``get_compression_tip`` - walk; a tip is adopted only while its row is still live. + Resolve and load first, then mutate the agent, so ambiguous lineage or an unreadable handoff fails closed. Uses + the transitive ``get_compression_tip`` walk; a tip is adopted only while its row is still live. """ resolver = getattr(type(session_db), "get_compression_tip", None) row_getter = getattr(type(session_db), "get_session", None) @@ -1724,8 +1703,7 @@ def _direct_messages_for_pre_compress_memory(messages: Any) -> list[dict[str, An Summaries, tool rows and system messages are omitted; assistant prose is kept with ``tool_calls`` stripped, and pure tool-call wrappers are dropped. """ - # Deferred import: context_compressor → turn_context → this module would form - # an import cycle. + # Deferred import: context_compressor → turn_context → this module would form an import cycle. from agent.context_compressor import COMPRESSED_SUMMARY_METADATA_KEY direct_messages: list[dict[str, Any]] = [] for message in messages or []: @@ -1838,8 +1816,7 @@ def _lower_threshold_to_aux_context( getattr(compressor, "max_tokens", None), ) threshold_suggestion_viable = recomputed_threshold is None or recomputed_threshold <= aux_context - # "model (provider)" labels for both sides; empty/"auto" provider falls back to - # the client's base_url hostname. + # "model (provider)" labels for both sides; empty/"auto" provider falls back to the client's base_url hostname. _main_model = getattr(agent, "model", "") or "?" _main_provider = getattr(agent, "provider", "") or "" _aux_provider_label = aux_provider if aux_provider and aux_provider != "auto" else "" @@ -1989,8 +1966,7 @@ def check_compression_model_feasibility(agent: Any) -> None: aux_base_url=aux_base_url, ) except ValueError: - # Hard rejections (aux below minimum context) must propagate - # so the session refuses to start. + # Hard rejections (aux below minimum context) must propagate so the session refuses to start. raise except Exception as exc: logger.debug("Compression feasibility check failed (non-fatal): %s", exc) @@ -2013,10 +1989,9 @@ def conversation_history_after_compression( ) -> Optional[list]: """Return the correct flush baseline after a compression boundary. - Session rotation returns ``None`` so the child gets the full compacted list. - In-place compaction returns a shallow copy of the already-persisted rows (else - the identity flush re-appends them). Aborted/no-op attempts keep the baseline: - marking all persisted drops unflushed turns; clearing re-appends rows. + Session rotation returns ``None`` so the child gets the full compacted list. In-place compaction returns a shallow + copy of the already-persisted rows (else the identity flush re-appends them). Aborted/no-op attempts keep the + baseline: marking all persisted drops unflushed turns; clearing re-appends rows. """ if bool(getattr(agent, "_last_compression_attempt_recorded", False)): attempt_in_place = getattr(agent, "_last_compression_attempt_in_place", None) @@ -2112,8 +2087,7 @@ def _extract_steer_text_from_message(message: Any) -> Optional[str]: open_marker, close_marker = _steer_markers() start = text.find(open_marker) if start == -1: - # Fallback: marker wording may evolve; look for the stable prefix, then skip to - # the end of the opening line. + # Fallback: marker wording may evolve; look for the stable prefix, then skip to the end of the opening line. start = text.find(_STEER_FALLBACK_OPEN) if start == -1: return None @@ -2132,8 +2106,7 @@ def _extract_steer_text_from_message(message: Any) -> Optional[str]: def _compressed_has_busy_steer(messages: list) -> bool: """Whether *messages* already carries a steer marker in a ``role=tool`` row. - Only tool rows count, so a summary merely quoting the marker text is not - mistaken for live intent. + Only tool rows count, so a summary merely quoting the marker text is not mistaken for live intent. """ for msg in messages: if not isinstance(msg, dict) or msg.get("role") != "tool": @@ -2273,10 +2246,9 @@ def _insert_real_user_anchor(messages: list, anchor: dict) -> CompressedUserTurn for index, message in enumerate(messages): if _role(message) == "assistant" and (index == 0 or _role(messages[index - 1]) != "user"): return _place(index) - # Every assistant is user-preceded (or there are none). Appending is safe whenever the - # transcript does not already end with a user turn. Never merge into a summary either: - # its prefix must stay at message start for summary detection; repair_message_sequence - # merges adjacent user turns summary-first. + # Every assistant is user-preceded (or there are none). Appending is safe whenever the transcript does not already + # end with a user turn. Never merge into a summary either: its prefix must stay at message start for summary + # detection; repair_message_sequence merges adjacent user turns summary-first. if ( not messages or _role(messages[-1]) != "user" @@ -2441,9 +2413,8 @@ class _CompactionLifecycle: class _CompressionLease: """The per-attempt durable compression lock plus its lifecycle plumbing. - ``holder`` is None when no durable lock is owned (legacy DB, no session db); - ``watermark`` is MAX(id) of active rows at lease start (None = archive - everything, no concurrent-tail preservation this cycle). + ``holder`` is None when no durable lock is owned (legacy DB, no session db); ``watermark`` is MAX(id) of active + rows at lease start (None = archive everything, no concurrent-tail preservation this cycle). """ def __init__( @@ -2568,11 +2539,10 @@ def _abort_lease( def _try_acquire_durable_lock(lease: _CompressionLease, try_acquire: Any, commit_fence: Any) -> bool: """Acquire the durable lock for ``lease.holder`` and capture the start watermark. - Watermark = MAX(id) of active rows at START: appends aren't blocked during summary; - later rows are concurrent tail that archive_and_compact re-sequences. Capture is - safety-additive (fallback archives everything), so its failure never aborts. An - acquire that raises is not version skew: fail closed and release holder-qualified - best-effort (safe if never acquired). + Watermark = MAX(id) of active rows at START: appends aren't blocked during summary; later rows are concurrent tail + that archive_and_compact re-sequences. Capture is safety-additive (fallback archives everything), so its failure + never aborts. An acquire that raises is not version skew: fail closed and release holder-qualified best-effort + (safe if never acquired). """ try: acquired = try_acquire(lease.sid, lease.holder, ttl_seconds=lease.ttl) @@ -2658,11 +2628,10 @@ def _acquire_compression_lease( ) -> Tuple[Optional[_CompressionLease], Optional[str]]: """Take the per-session compression lock; ``(None, prompt)`` means sit out. - Two AIAgents sharing a session_id (e.g. background review fork) would both - rotate and orphan a child. Keyed on the OLD id (what rivals read from - SessionEntry). Loser sits out: messages unchanged, caller sees no-op. - Only structural absence of the lock API (version skew) fails open; once - resolved, any exception fails closed since unlocked runs can fork lineage. + Two AIAgents sharing a session_id (e.g. background review fork) would both rotate and orphan a child. Keyed on the + OLD id (what rivals read from SessionEntry). Loser sits out: messages unchanged, caller sees no-op. Only + structural absence of the lock API (version skew) fails open; once resolved, any exception fails closed since + unlocked runs can fork lineage. """ _lock_db = getattr(agent, "_session_db", None) _lock_sid = agent.session_id or "" @@ -2727,8 +2696,7 @@ def _acquire_compression_lease( if lease.holder is not None: agent._active_compression_lock_holder = lease.holder if commit_fence is not None and commit_fence.register_cancelled_lock_release(lease.release_holder_only): - # Cancellation won during lock setup (hook ran synchronously, lease gone): - # abort before any summary work. + # Cancellation won during lock setup (hook ran synchronously, lease gone): abort before any summary work. logger.info( "Compression commit cancelled before summary dispatch (session=%s).", agent.session_id or "none" ) @@ -2831,10 +2799,9 @@ def _adopt_grown_durable_parent(agent: Any, lease: _CompressionLease, messages: def _pre_compress_memory_context(agent: Any, messages: list, checkpoint_required: bool) -> str: """Provider ``on_pre_compress()`` insights to surface in the summary ("" if none). - Raw messages stay the API v1 provider contract; normalized evidence goes only - to API v2+ checkpoint providers inside MemoryManager.on_pre_compress(). - Raises :class:`CompressionCheckpointUnavailable` when a required checkpoint - cannot be taken. + Raw messages stay the API v1 provider contract; normalized evidence goes only to API v2+ checkpoint providers + inside MemoryManager.on_pre_compress(). Raises :class:`CompressionCheckpointUnavailable` when a required + checkpoint cannot be taken. """ memory_context = "" memory_manager = getattr(agent, "_memory_manager", None) @@ -2952,8 +2919,7 @@ def _run_summary_dispatch( aux_interrupt_protection(cancel_check=_compression_cancel_requested), ): compressed = compress_fn(messages, **compress_kwargs) - # Freeze a hard stop that arrived after the last provider attempt but before - # session state rotates. + # Freeze a hard stop that arrived after the last provider attempt but before session state rotates. if hard_cancel_event is not None and hard_cancel_event.is_set(): raise AuxiliaryExplicitCancellation() finally: @@ -3144,10 +3110,9 @@ def _parent_deliberately_ended(session_db: Any, session_id: str) -> bool: def _carry_session_state_to_child(agent: Any, old_session_id: str, old_title: Any) -> None: """Migrate /goal, /heartbeat, /loop state and the title from the parent to the child. - Each lookup is a flat per-session read with no parent walk, so state would silently - die at the boundary. The title is carried unchanged (renumbering per rotation made - one session look like many); its provenance is read BEFORE the transfer clears the - ancestor's row, then restored so an inherited auto-title stays upgradeable. + Each lookup is a flat per-session read with no parent walk, so state would silently die at the boundary. The title + is carried unchanged (renumbering per rotation made one session look like many); its provenance is read BEFORE the + transfer clears the ancestor's row, then restored so an inherited auto-title stays upgradeable. """ with _swallow('Could not migrate goal on compression: %s'): from hermes_cli.goals import migrate_goal_to_session @@ -3186,9 +3151,8 @@ def _publish_rotated_compaction( ) -> None: """Rotate the session: flush the parent, publish the child, re-point the agent. - Flushes current-turn msgs to the OLD session, passing the durable prefix - (messages[:persist idx]) so preflight, which runs before rows are - marker-stamped, can't re-append them. + Flushes current-turn msgs to the OLD session, passing the durable prefix (messages[:persist idx]) so preflight, + which runs before rows are marker-stamped, can't re-append them. """ current_idx = getattr(agent, "_persist_user_message_idx", None) persisted_history = ( @@ -3204,8 +3168,7 @@ def _publish_rotated_compaction( try: _foreign_tail_ceiling = agent._session_db.get_active_message_watermark(agent.session_id) except Exception: - # No trustworthy ceiling: the clone could duplicate the handoff, so skip tail - # preservation this rotation. + # No trustworthy ceiling: the clone could duplicate the handoff, so skip tail preservation this rotation. _foreign_tail_ceiling = None with contextlib.suppress(Exception): # best-effort — don't block compression on a flush error agent._flush_messages_to_session_db(messages, conversation_history=persisted_history) @@ -3253,10 +3216,9 @@ def _publish_rotated_compaction( _compressed_anchor_source[_DB_PERSISTED_MARKER] = True _session_messages = getattr(agent, "_session_messages", None) if isinstance(_session_messages, list) and _session_messages is not messages: - # Adoption may leave _session_messages on the pre-adoption list with an out-of- - # range idx; stamp every scoped twin against the ANCHOR SOURCE, as the wrapper. - # An already-stamped exact twin still suppresses the broad pass here, or a - # content-equal old duplicate would get stamped. + # Adoption may leave _session_messages on the pre-adoption list with an out-of- range idx; stamp every + # scoped twin against the ANCHOR SOURCE, as the wrapper. An already-stamped exact twin still + # suppresses the broad pass here, or a content-equal old duplicate would get stamped. _stamp_scoped_twins(_session_messages, _compressed_anchor_source, exact_counts_stamped=True) for _handoff_message in compressed: if isinstance(_handoff_message, dict): @@ -3295,8 +3257,7 @@ def _warn_summary_or_aux_fallback(agent: Any) -> None: def _reset_read_dedup_caches(task_id: str, *, skills: bool = True) -> None: """Clear the file-read (and skill_view) repeat-read dedup caches after a boundary. - Original read content was summarized away, so a re-read must return full content, - not a "file unchanged" stub. + Original read content was summarized away, so a re-read must return full content, not a "file unchanged" stub. """ with contextlib.suppress(Exception): from tools.file_tools import reset_file_dedup @@ -3436,9 +3397,8 @@ def _candidate_rejected( ) -> bool: """Reject an unusable compression candidate before any session mutation. - Order matters: compressor-reported abort, no progress, empty transcript, - superseded attempt. Each branch surfaces its own warning/telemetry; the - caller releases the lease and returns the input unchanged when True. + Order matters: compressor-reported abort, no progress, empty transcript, superseded attempt. Each branch surfaces + its own warning/telemetry; the caller releases the lease and returns the input unchanged when True. """ # Aborted compression returns input unchanged: surface the error, skip rotation # (no session ended); auto-compress callers detect no-op via equal lengths. @@ -3539,9 +3499,8 @@ def _commit_compaction( """Persist the compacted transcript: memory extraction, anti-growth guard, then the in-place archive or the parent->child rotation. - Failures roll the live list back and arm the split-failure cooldown; a refused - (would-grow) candidate returns ``refused_prompt`` so the caller hands back the - input unchanged. + Failures roll the live list back and arm the split-failure cooldown; a refused (would-grow) candidate returns + ``refused_prompt`` so the caller hands back the input unchanged. """ session_commit_succeeded = False compacted_in_place = False @@ -3619,14 +3578,12 @@ def _commit_compaction( agent._flushed_db_message_session_id = agent.session_id session_commit_succeeded = True except Exception as e: - # Rotation: atomic publication failed (including lease loss) — keep the parent - # live and discard the stale compacted snapshot. In-place: archive_and_compact is - # atomic so old rows stay active, but marker-swept `compressed` would re-INSERT on - # top of them (doubling each try); gate on split_status (set right after commit). - # Either way the deepcopy keeps markers/identity and only the prune runway rolls - # back (the full snapshot restore is for pre-commit cancels; telemetry keeps the - # failed values). _db_flush_scan_prefix is intentionally NOT cleared: the scan is - # identity-based and the deepcopy replaces every row. + # Rotation: atomic publication failed (including lease loss) — keep the parent live and discard the stale + # compacted snapshot. In-place: archive_and_compact is atomic so old rows stay active, but marker-swept + # `compressed` would re-INSERT on top of them (doubling each try); gate on split_status (set right after + # commit). Either way the deepcopy keeps markers/identity and only the prune runway rolls back (the full + # snapshot restore is for pre-commit cancels; telemetry keeps the failed values). _db_flush_scan_prefix is + # intentionally NOT cleared: the scan is identity-based and the deepcopy replaces every row. rotation_rollback = not in_place and old_session_id and agent.session_id == old_session_id if rotation_rollback or ( in_place and split_status != "in_place_committed" and messages_before_compression is not None @@ -3699,9 +3656,8 @@ def _run_summary_phase( ) -> _SummaryPhase: """Adopt a grown durable parent, gather memory context and run the summarizer. - A hard cancel restores the compressor snapshot + live list, records a stall - backoff while the lease is still held, and aborts; any other failure releases - the lease and re-raises. + A hard cancel restores the compressor snapshot + live list, records a stall backoff while the lease is still held, + and aborts; any other failure releases the lease and re-raises. """ pre_msg_count = len(messages) _activity_heartbeat: Optional[_CompressionActivityHeartbeat] = None @@ -3782,8 +3738,7 @@ def _run_summary_phase( ) return _SummaryPhase(messages=messages, abort_prompt=_existing_system_prompt(agent, system_message)) except BaseException as _compress_exc: - # Any failure after lock acquisition must release it or the session is - # permanently blocked from compression. + # Any failure after lock acquisition must release it or the session is permanently blocked from compression. if _activity_heartbeat is not None: _activity_heartbeat.stop("context compression failed") _activity_heartbeat = None @@ -3805,12 +3760,11 @@ def _run_summary_phase( def _begin_compression_attempt(agent: Any, *, force: bool, defer_notification: bool) -> Tuple[dict, int, float]: """Snapshot + claim the compressor, reset per-attempt agent signals, seed telemetry. - Returns ``(attempt_snapshot, attempt_generation, attempt_started_at)``. The claim - stops a late-unwinding sibling (stall-fallback overlap) from restoring its snapshot - over ours or clearing our cancellation consult. Signals are cleared at the VERY TOP, - before codex/breaker early-returns, so a stale value cannot make a later no-op look - like lock contention; ``_last_compression_attempt_in_place=None`` means aborted/no - boundary for ``conversation_history_after_compression()``. + Returns ``(attempt_snapshot, attempt_generation, attempt_started_at)``. The claim stops a late-unwinding sibling + (stall-fallback overlap) from restoring its snapshot over ours or clearing our cancellation consult. Signals are + cleared at the VERY TOP, before codex/breaker early-returns, so a stale value cannot make a later no-op look like + lock contention; ``_last_compression_attempt_in_place=None`` means aborted/no boundary for + ``conversation_history_after_compression()``. """ snapshot = _snapshot_compressor_attempt_state(agent.context_compressor) generation = _claim_compressor_attempt(agent.context_compressor) @@ -3907,10 +3861,9 @@ def compress_context( ) -> Tuple[list, str]: """Compress conversation context and split the session in SQLite. - ``force`` (manual /compress) clears the summary-failure cooldown; - ``bypass_cooldown`` (provider-proven overflow) skips it once, breakers still - apply. ``commit_fence`` stops a timed-out worker mutating session state. - Returns ``(messages, system_prompt)``; on abort input is unchanged, NOT split. + ``force`` (manual /compress) clears the summary-failure cooldown; ``bypass_cooldown`` (provider-proven overflow) + skips it once, breakers still apply. ``commit_fence`` stops a timed-out worker mutating session state. Returns + ``(messages, system_prompt)``; on abort input is unchanged, NOT split. """ _compressor_attempt_snapshot, _attempt_generation, _attempt_started_at = _begin_compression_attempt( agent, force=force, defer_notification=defer_context_engine_notification @@ -4314,11 +4267,10 @@ def _decode_pixels(data_url: str) -> Optional[tuple]: def _shrink_data_url(url: str, *, max_dimension: int, resize_fn: Any) -> tuple: """Return ``(resized_url, unshrinkable)`` for a data URL. - ``resized_url`` is None when no rewrite applied. ``unshrinkable`` is True only - when the image violated a constraint and resizing failed to satisfy that same - constraint, so the caller knows a retry is pointless. The accept gate MUST use - the axis that triggered the shrink: a pixel downscale can re-encode to MORE - bytes (PNG non-monotonic); a byte-only reject wedges. + ``resized_url`` is None when no rewrite applied. ``unshrinkable`` is True only when the image violated a + constraint and resizing failed to satisfy that same constraint, so the caller knows a retry is pointless. The + accept gate MUST use the axis that triggered the shrink: a pixel downscale can re-encode to MORE bytes (PNG + non-monotonic); a byte-only reject wedges. """ target_bytes = _IMAGE_SHRINK_TARGET_BYTES if not isinstance(url, str) or not url.startswith("data:"): @@ -4397,9 +4349,8 @@ def _write_data_url_to_source(source: dict, data_url: str) -> dict: def try_shrink_image_parts_in_messages(api_messages: list, *, max_dimension: int = 8000) -> bool: """Re-encode oversized native image parts to recover from image-too-large errors. - Mutates ``api_messages`` in place. Returns True if any part was replaced, - False if nothing to shrink or Pillow could not help. Targets data-URL parts - over 4 MB or ``max_dimension`` (Anthropic's per-side pixel cap, parsed from + Mutates ``api_messages`` in place. Returns True if any part was replaced, False if nothing to shrink or Pillow + could not help. Targets data-URL parts over 4 MB or ``max_dimension`` (Anthropic's per-side pixel cap, parsed from the rejection by the caller); http(s) image URLs are left untouched. """ if not api_messages: