refactor(agent/context_compressor): reflow docstrings/comments to width (every word preserved)
This commit is contained in:
+149
-235
@@ -1,9 +1,8 @@
|
||||
"""Automatic context window compression for long conversations.
|
||||
|
||||
Uses a cheap auxiliary model to summarize middle turns while protecting head
|
||||
and tail context: structured iterative summaries, token-budget tail protection,
|
||||
tool-output pruning before summarization, and scaled summary budgets.
|
||||
"""
|
||||
Uses a cheap auxiliary model to summarize middle turns while protecting head and tail context:
|
||||
structured iterative summaries, token-budget tail protection, tool-output pruning before summarization,
|
||||
and scaled summary budgets."""
|
||||
|
||||
import contextlib
|
||||
import contextvars
|
||||
@@ -66,8 +65,7 @@ _PINNED_ROUTE_FIELDS: tuple[str, ...] = ("provider", "model", "base_url", "api_k
|
||||
def pin_summary_route(route: Optional[Dict[str, Any]]):
|
||||
"""Pin the next summary LLM call in this context to an explicit route.
|
||||
|
||||
``None`` is a no-op passthrough. Re-entrant-safe: restores the previous pin on exit.
|
||||
"""
|
||||
``None`` is a no-op passthrough. Re-entrant-safe: restores the previous pin on exit."""
|
||||
token = _SUMMARY_ROUTE_PIN.set(route if isinstance(route, dict) else None)
|
||||
try:
|
||||
yield
|
||||
@@ -78,8 +76,7 @@ def pin_summary_route(route: Optional[Dict[str, Any]]):
|
||||
def take_pinned_summary_route() -> Optional[Dict[str, Any]]:
|
||||
"""Read and consume the pinned summary route, if one is installed.
|
||||
|
||||
Single use by design: the main-model retry must not re-issue the pin.
|
||||
"""
|
||||
Single use by design: the main-model retry must not re-issue the pin."""
|
||||
route = _SUMMARY_ROUTE_PIN.get()
|
||||
if route is None:
|
||||
return None
|
||||
@@ -110,9 +107,8 @@ _HYGIENE_PREAGENT_ONLY_COOLDOWN_MARKERS: tuple[str, ...] = (
|
||||
def _is_hygiene_preagent_only_cooldown(error: object) -> bool:
|
||||
"""Return True for a cooldown that belongs only to pre-agent hygiene.
|
||||
|
||||
Hygiene watchdog timeouts / turn-hold deferrals are not evidence of an auxiliary-model
|
||||
failure and must never block the in-agent compressor.
|
||||
"""
|
||||
Hygiene watchdog timeouts / turn-hold deferrals are not evidence of an auxiliary-model failure and
|
||||
must never block the in-agent compressor."""
|
||||
text = str(error or "").strip().casefold()
|
||||
return any(marker in text for marker in _HYGIENE_PREAGENT_ONLY_COOLDOWN_MARKERS)
|
||||
|
||||
@@ -120,8 +116,7 @@ def _is_hygiene_preagent_only_cooldown(error: object) -> bool:
|
||||
def _response_finish_reason(response: Any) -> str:
|
||||
"""Return lowercased ``choices[0].finish_reason`` from a dict- or object-shaped response.
|
||||
|
||||
Returns ``""`` when absent/unreadable.
|
||||
"""
|
||||
Returns ``""`` when absent/unreadable."""
|
||||
try:
|
||||
if isinstance(response, dict):
|
||||
choices = response.get("choices") or [{}]
|
||||
@@ -248,8 +243,7 @@ _BACKGROUND_PROCESS_NOTIFICATION_PREFIX = "[IMPORTANT: Background process "
|
||||
def _fresh_compaction_message_copy(msg: Dict[str, Any]) -> Dict[str, Any]:
|
||||
"""Copy a message for compaction assembly without persistence markers.
|
||||
|
||||
The authoritative guarantee is the terminal sweep ``_strip_persistence_markers``.
|
||||
"""
|
||||
The authoritative guarantee is the terminal sweep ``_strip_persistence_markers``."""
|
||||
fresh = msg.copy()
|
||||
fresh.pop(_DB_PERSISTED_MARKER, None)
|
||||
return fresh
|
||||
@@ -258,9 +252,8 @@ def _fresh_compaction_message_copy(msg: Dict[str, Any]) -> Dict[str, Any]:
|
||||
def _template_visible_role(message: Any) -> Optional[str]:
|
||||
"""Role as counted by strict chat-template alternation checks.
|
||||
|
||||
Mistral-family templates exempt ``tool`` rows and assistant rows with ``tool_calls``
|
||||
from alternation. Returns ``None`` for messages the check skips.
|
||||
"""
|
||||
Mistral-family templates exempt ``tool`` rows and assistant rows with ``tool_calls`` from
|
||||
alternation. Returns ``None`` for messages the check skips."""
|
||||
if not isinstance(message, dict):
|
||||
return None
|
||||
role = message.get("role")
|
||||
@@ -274,11 +267,10 @@ def _template_visible_role(message: Any) -> Optional[str]:
|
||||
def _strip_persistence_markers(messages: List[Dict[str, Any]]) -> None:
|
||||
"""Enforce the invariant: no assembled message carries a persistence marker.
|
||||
|
||||
A leaked ``_db_persisted`` makes the child-session rotation flush skip the row, losing it
|
||||
from state.db. Per-copy-site strips are positional and re-leak when a copy site is added;
|
||||
this terminal sweep makes the guarantee structural. Run once on the fully assembled list;
|
||||
mutates in place (compaction-local copies).
|
||||
"""
|
||||
A leaked ``_db_persisted`` makes the child-session rotation flush skip the row, losing it from
|
||||
state.db. Per-copy-site strips are positional and re-leak when a copy site is added; this terminal
|
||||
sweep makes the guarantee structural. Run once on the fully assembled list; mutates in place
|
||||
(compaction-local copies)."""
|
||||
for msg in messages:
|
||||
if isinstance(msg, dict):
|
||||
msg.pop(_DB_PERSISTED_MARKER, None)
|
||||
@@ -287,11 +279,10 @@ def _strip_persistence_markers(messages: List[Dict[str, Any]]) -> None:
|
||||
def stamp_db_persisted_markers(messages: List[Dict[str, Any]]) -> None:
|
||||
"""Fulfil the post-commit contract of ``SessionDB.archive_and_compact()``.
|
||||
|
||||
Single stamp site for all callers. Call ONLY after the commit succeeded, on the
|
||||
dict instances the caller keeps live. Needed because compress() output is marker-swept
|
||||
for the ROTATION flush; an in-place commit returned unstamped is re-INSERTed as new by
|
||||
the next persist walk and the transcript doubles on every compaction.
|
||||
"""
|
||||
Single stamp site for all callers. Call ONLY after the commit succeeded, on the dict instances the
|
||||
caller keeps live. Needed because compress() output is marker-swept for the ROTATION flush; an
|
||||
in-place commit returned unstamped is re-INSERTed as new by the next persist walk and the transcript
|
||||
doubles on every compaction."""
|
||||
for msg in messages:
|
||||
if isinstance(msg, dict):
|
||||
msg[_DB_PERSISTED_MARKER] = True
|
||||
@@ -300,12 +291,10 @@ def stamp_db_persisted_markers(messages: List[Dict[str, Any]]) -> None:
|
||||
def _prune_stale_reasoning_replay(messages: List[Dict[str, Any]]) -> int:
|
||||
"""Strip stale ``codex_reasoning_items`` from assistant turns older than the active one.
|
||||
|
||||
Boundary is the last USER message (a turn spans several assistant rows): the Responses
|
||||
API replays a turn's bridging reasoning items together, so cutting at the last ASSISTANT
|
||||
would strip mid-chain. ``type: "compaction"`` items are cumulative context carriers that
|
||||
must survive on every retained message — filter items, never pop the key. In place;
|
||||
returns pruned message count.
|
||||
"""
|
||||
Boundary is the last USER message (a turn spans several assistant rows): the Responses API replays a
|
||||
turn's bridging reasoning items together, so cutting at the last ASSISTANT would strip mid-chain.
|
||||
``type: "compaction"`` items are cumulative context carriers that must survive on every retained
|
||||
message — filter items, never pop the key. In place; returns pruned message count."""
|
||||
# Active turn = everything after the last real user message; synthetic
|
||||
# continuation rows and tool results never mark a turn boundary.
|
||||
last_user_idx = -1
|
||||
@@ -377,8 +366,7 @@ def _looks_like_compaction_summary(msg: Dict[str, Any], content: str) -> bool:
|
||||
def _salvage_reduce_todo_snapshot(out: List[Dict[str, Any]]) -> None:
|
||||
"""Last-resort shrink: reduce or drop the synthetic todo snapshot.
|
||||
|
||||
A snapshot carrying a pruned-skill reload notice keeps just the notice.
|
||||
"""
|
||||
A snapshot carrying a pruned-skill reload notice keeps just the notice."""
|
||||
from agent.conversation_compression import _PRUNED_SKILL_RELOAD_NOTICE_HEADER
|
||||
|
||||
for i in range(len(out) - 1, -1, -1):
|
||||
@@ -404,8 +392,7 @@ def salvage_grown_transcript(
|
||||
) -> Optional[List[Dict[str, Any]]]:
|
||||
"""Mechanically shrink a compression candidate, or return ``None``.
|
||||
|
||||
Works on copies; cheapest-information-loss first; admitted only when strictly smaller.
|
||||
"""
|
||||
Works on copies; cheapest-information-loss first; admitted only when strictly smaller."""
|
||||
if not candidate or not original:
|
||||
return None
|
||||
if budget is None:
|
||||
@@ -715,9 +702,8 @@ _TIMEOUT_COOLDOWN_LADDER = (60, 300, 900)
|
||||
def _next_timeout_cooldown(compressor: Any) -> int:
|
||||
"""Bump ``compressor._consecutive_timeout_failures`` and return the ladder rung for it.
|
||||
|
||||
Module-level (not a method) so callers that bind a single real method onto a stub still
|
||||
exercise the ladder.
|
||||
"""
|
||||
Module-level (not a method) so callers that bind a single real method onto a stub still exercise the
|
||||
ladder."""
|
||||
compressor._consecutive_timeout_failures = getattr(compressor, "_consecutive_timeout_failures", 0) + 1
|
||||
return _TIMEOUT_COOLDOWN_LADDER[
|
||||
min(compressor._consecutive_timeout_failures, len(_TIMEOUT_COOLDOWN_LADDER)) - 1
|
||||
@@ -753,10 +739,9 @@ _CLARIFY_NON_RESPONSE_PREFIXES = (
|
||||
def _is_clarify_non_response_sentinel(response: Any) -> bool:
|
||||
"""Return True when a clarify ``user_response`` is runtime sentinel prose, not an answer.
|
||||
|
||||
For lists, ANY sentinel item poisons the whole response: real producers only emit scalar
|
||||
sentinels, so a mixed list is forged/corrupt content — fall back to the generic path (may
|
||||
lose info, never misattributes a user answer).
|
||||
"""
|
||||
For lists, ANY sentinel item poisons the whole response: real producers only emit scalar sentinels,
|
||||
so a mixed list is forged/corrupt content — fall back to the generic path (may lose info, never
|
||||
misattributes a user answer)."""
|
||||
if isinstance(response, str):
|
||||
return response.lstrip().startswith(_CLARIFY_NON_RESPONSE_PREFIXES)
|
||||
if isinstance(response, list):
|
||||
@@ -804,8 +789,7 @@ def _extract_pruned_skill_names(text: str) -> list[str]:
|
||||
def _collect_ghosted_skill_names(turns: List[Dict[str, Any]]) -> list[str]:
|
||||
"""Skill names whose instructions are about to be lost in compaction.
|
||||
|
||||
Covers both already-demoted ``skill_view`` rows and raw, never-demoted bodies.
|
||||
"""
|
||||
Covers both already-demoted ``skill_view`` rows and raw, never-demoted bodies."""
|
||||
names: list[str] = []
|
||||
|
||||
def _add(name: str) -> None:
|
||||
@@ -845,9 +829,8 @@ _PRUNED_SKILLS_SECTION_HEADING = "## Pruned Skills"
|
||||
def _reinject_pruned_skill_markers(summary: str, skill_names: list[str]) -> str:
|
||||
"""Deterministically restore prune markers the summarizer dropped.
|
||||
|
||||
Presence is checked against the canonical marker string; the appended block is
|
||||
plain body text (no handoff prefix/scaffolding) and is redacted like all others.
|
||||
"""
|
||||
Presence is checked against the canonical marker string; the appended block is plain body text (no
|
||||
handoff prefix/scaffolding) and is redacted like all others."""
|
||||
if not skill_names:
|
||||
return summary
|
||||
missing = [name for name in skill_names if _skill_pruned_marker(name) not in summary]
|
||||
@@ -907,8 +890,7 @@ def _synthetic_user_row(content: str) -> bool:
|
||||
def _build_verbatim_user_section(turns: List[Dict[str, Any]]) -> str:
|
||||
"""Embed the compacted region's REAL user messages verbatim in the summary.
|
||||
|
||||
Newest-first under a char budget, straddler truncated. Returns "" when none.
|
||||
"""
|
||||
Newest-first under a char budget, straddler truncated. Returns "" when none."""
|
||||
collected: list[str] = []
|
||||
used = 0
|
||||
for msg in reversed(turns):
|
||||
@@ -943,9 +925,8 @@ def _build_verbatim_user_section(turns: List[Dict[str, Any]]) -> str:
|
||||
def _build_recovery_footer(session_id: str, region_len: int) -> str:
|
||||
"""Deterministic pointer to the compacted region in session history.
|
||||
|
||||
state.db keeps every pre-compaction message; naming the session_search re-access path
|
||||
lets the model treat compaction as deferred retrieval, not loss.
|
||||
"""
|
||||
state.db keeps every pre-compaction message; naming the session_search re-access path lets the model
|
||||
treat compaction as deferred retrieval, not loss."""
|
||||
if not session_id:
|
||||
return ""
|
||||
return (
|
||||
@@ -986,8 +967,7 @@ _ANCHOR_NOISE = frozenset({
|
||||
def _build_anchor_index(turns: List[Dict[str, Any]]) -> str:
|
||||
"""Regex-harvest exact identifiers from the compacted region (LLM-free).
|
||||
|
||||
Per-category caps; most-frequent first, ties by last-seen order.
|
||||
"""
|
||||
Per-category caps; most-frequent first, ties by last-seen order."""
|
||||
text_parts: list[str] = []
|
||||
for msg in turns:
|
||||
c = msg.get("content")
|
||||
@@ -1061,9 +1041,8 @@ def _skill_view_call_sites(messages: List[Dict[str, Any]]) -> list[tuple[int, st
|
||||
def _collect_protected_skill_names(messages: List[Dict[str, Any]], prune_boundary: int) -> set[str]:
|
||||
"""Skill names (lower-cased) whose skill_view bodies must survive Phase-1 demotion.
|
||||
|
||||
Recently loaded, loaded inside the protected tail, or named by a tail user message.
|
||||
Applies to Phase-1/2 only; the Pass-4 pressure demotion ignores it.
|
||||
"""
|
||||
Recently loaded, loaded inside the protected tail, or named by a tail user message. Applies to
|
||||
Phase-1/2 only; the Pass-4 pressure demotion ignores it."""
|
||||
total = len(messages)
|
||||
if not total:
|
||||
return set()
|
||||
@@ -1129,9 +1108,8 @@ _HISTORICAL_TASK_SECTION_RE = re.compile(rf"(?ms)^{re.escape(HISTORICAL_TASK_HEA
|
||||
def _redact_compaction_text(text: Any) -> str:
|
||||
"""Redact text that crosses a compaction summary boundary (strict mode).
|
||||
|
||||
``force=True`` overrides ``security.redact_secrets: false``; URL credentials are
|
||||
redacted too, since summaries persist and re-enter every later prompt.
|
||||
"""
|
||||
``force=True`` overrides ``security.redact_secrets: false``; URL credentials are redacted too, since
|
||||
summaries persist and re-enter every later prompt."""
|
||||
return redact_sensitive_text(text or "", force=True, redact_url_credentials=True)
|
||||
|
||||
|
||||
@@ -1222,8 +1200,7 @@ def _bullets(items: list[str], limit: int = 8) -> str:
|
||||
def _content_length_for_budget(raw_content: Any) -> int:
|
||||
"""Return the effective char-length of a message's content for token budgeting.
|
||||
|
||||
Text parts by length plus a flat ``_IMAGE_CHAR_EQUIVALENT`` per image part.
|
||||
"""
|
||||
Text parts by length plus a flat ``_IMAGE_CHAR_EQUIVALENT`` per image part."""
|
||||
if isinstance(raw_content, str):
|
||||
return len(raw_content)
|
||||
if not isinstance(raw_content, list):
|
||||
@@ -1269,8 +1246,7 @@ _STALE_REPLAY_PRUNE_KEYS = "codex_reasoning_items",
|
||||
def _reasoning_details_text_chars(value: Any) -> int:
|
||||
"""Textual thinking chars inside a ``reasoning_details`` envelope.
|
||||
|
||||
Counts only thinking text, never signed/base64 envelope blobs.
|
||||
"""
|
||||
Counts only thinking text, never signed/base64 envelope blobs."""
|
||||
if not value:
|
||||
return 0
|
||||
if isinstance(value, str):
|
||||
@@ -1291,12 +1267,11 @@ def _reasoning_details_text_chars(value: Any) -> int:
|
||||
def _estimate_msg_budget_tokens(msg: dict, charge_stale_thinking: bool = True) -> int:
|
||||
"""Token estimate for one message in the tail-protection budget walks.
|
||||
|
||||
Counts content, the full ``tool_call`` envelope (arguments-only undercounted parallel-call
|
||||
turns by 2-15x), and always-replayed provider fields. Always-replayed fields are charged
|
||||
because the preflight estimator sees the full shape; a mismatched size class protects
|
||||
blob-heavy rows as "small" and compaction re-fires. ``charge_stale_thinking=False`` skips
|
||||
newest-turn-only thinking keys. Accounting only; never mutates.
|
||||
"""
|
||||
Counts content, the full ``tool_call`` envelope (arguments-only undercounted parallel-call turns by
|
||||
2-15x), and always-replayed provider fields. Always-replayed fields are charged because the
|
||||
preflight estimator sees the full shape; a mismatched size class protects blob-heavy rows as "small"
|
||||
and compaction re-fires. ``charge_stale_thinking=False`` skips newest-turn-only thinking keys.
|
||||
Accounting only; never mutates."""
|
||||
content = msg.get("content") or ""
|
||||
if isinstance(content, str):
|
||||
tokens = estimate_tokens_rough(content) + 10 # +10 for role/key overhead
|
||||
@@ -1328,8 +1303,7 @@ def _estimate_msg_budget_tokens(msg: dict, charge_stale_thinking: bool = True) -
|
||||
def _last_assistant_index(messages: "List[Dict[str, Any]]") -> int:
|
||||
"""Index of the newest assistant message, or -1 (the one turn whose thinking may replay).
|
||||
|
||||
See ``_NEWEST_TURN_ONLY_BUDGET_KEYS``.
|
||||
"""
|
||||
See ``_NEWEST_TURN_ONLY_BUDGET_KEYS``."""
|
||||
for i in range(len(messages) - 1, -1, -1):
|
||||
msg = messages[i]
|
||||
if isinstance(msg, dict) and msg.get("role") == "assistant":
|
||||
@@ -1388,9 +1362,8 @@ def _tool_content_has_images(content: Any) -> bool:
|
||||
def _strip_images_from_tool_msg(msg: Dict[str, Any]) -> Optional[Dict[str, Any]]:
|
||||
"""Return a copy of a tool message with its image payloads replaced.
|
||||
|
||||
Returns ``None`` when nothing is strippable. Drops the stale ``api_content``
|
||||
sidecar on the copy; never mutates the input.
|
||||
"""
|
||||
Returns ``None`` when nothing is strippable. Drops the stale ``api_content`` sidecar on the copy;
|
||||
never mutates the input."""
|
||||
content = msg.get("content")
|
||||
if isinstance(content, dict) and content.get("_multimodal"):
|
||||
summary = content.get("text_summary") or "[screenshot removed to save context]"
|
||||
@@ -1410,9 +1383,8 @@ def _retire_stale_tool_result_images(
|
||||
) -> int:
|
||||
"""Replace image payloads on older tool results with text placeholders.
|
||||
|
||||
Keeps the newest ``keep_newest`` image-bearing tool messages; user uploads untouched.
|
||||
Mutates ``result`` in place; returns the number of messages rewritten.
|
||||
"""
|
||||
Keeps the newest ``keep_newest`` image-bearing tool messages; user uploads untouched. Mutates
|
||||
``result`` in place; returns the number of messages rewritten."""
|
||||
if keep_newest < 0:
|
||||
keep_newest = 0
|
||||
seen = 0
|
||||
@@ -1437,8 +1409,7 @@ def _retire_stale_tool_result_images(
|
||||
def _truncate_tool_call_args_json(args: str, head_chars: int = 200) -> str:
|
||||
"""Shrink long string leaves inside a tool-call arguments JSON blob, keeping JSON valid.
|
||||
|
||||
Providers 400 on malformed arguments. Non-JSON input is returned unchanged.
|
||||
"""
|
||||
Providers 400 on malformed arguments. Non-JSON input is returned unchanged."""
|
||||
try:
|
||||
parsed = json.loads(args)
|
||||
except (ValueError, TypeError):
|
||||
@@ -1486,10 +1457,9 @@ def _strip_images_from_content(content: Any) -> Any:
|
||||
def _strip_historical_media(messages: List[Dict[str, Any]]) -> List[Dict[str, Any]]:
|
||||
"""Replace image parts in older messages with placeholder text.
|
||||
|
||||
Rule 1: strip everything before the newest image-bearing user message. Rule 1b: the
|
||||
opening attachment ages out once a newer tool image exists. Rule 2: keep only the
|
||||
newest tool-result image. Unchanged list when nothing applies; input never mutated.
|
||||
"""
|
||||
Rule 1: strip everything before the newest image-bearing user message. Rule 1b: the opening
|
||||
attachment ages out once a newer tool image exists. Rule 2: keep only the newest tool-result image.
|
||||
Unchanged list when nothing applies; input never mutated."""
|
||||
if not messages:
|
||||
return messages
|
||||
|
||||
@@ -1557,8 +1527,7 @@ def _summary_part_text(part: Any) -> str:
|
||||
def _image_part_label(part: Dict[str, Any]) -> str:
|
||||
"""Render a multimodal image part as a short text label for the summarizer.
|
||||
|
||||
http(s) URLs are kept as a reusable handle; ``data:`` URLs collapse to ``[image]``.
|
||||
"""
|
||||
http(s) URLs are kept as a reusable handle; ``data:`` URLs collapse to ``[image]``."""
|
||||
url = ""
|
||||
if isinstance(part.get("image_url"), dict):
|
||||
url = str(part["image_url"].get("url") or "")
|
||||
@@ -1582,8 +1551,7 @@ def _str_arg(args: dict, key: str, default: str = "") -> str:
|
||||
def _summarize_tool_result(tool_name: str, tool_args: str, tool_content: str) -> str:
|
||||
"""Create an informative 1-line summary of a tool call + result.
|
||||
|
||||
Never raises: a malformed historical call must not crash-loop compression.
|
||||
"""
|
||||
Never raises: a malformed historical call must not crash-loop compression."""
|
||||
try:
|
||||
return _summarize_tool_result_unguarded(tool_name, tool_args, tool_content)
|
||||
except Exception as exc: # noqa: BLE001 — a summary must never crash compression
|
||||
@@ -1763,9 +1731,8 @@ def _summarize_tool_result_unguarded(tool_name: str, tool_args: str, tool_conten
|
||||
def resolve_model_threshold(model: str, model_thresholds: dict[str, float] | None, default: float) -> float:
|
||||
"""Resolve the effective compression threshold for a given model.
|
||||
|
||||
Longest matching ``model_thresholds`` substring key wins; otherwise ``default``.
|
||||
Module-level so plugin context engines can reuse it.
|
||||
"""
|
||||
Longest matching ``model_thresholds`` substring key wins; otherwise ``default``. Module-level so
|
||||
plugin context engines can reuse it."""
|
||||
if not model_thresholds or not model:
|
||||
return default
|
||||
best_key = ""
|
||||
@@ -1801,8 +1768,7 @@ def _memory_provider_section(memory_context: str) -> str:
|
||||
def _today_for_prompt() -> str:
|
||||
"""Date-only (user tz) for temporal anchoring; "" when the clock fails (never blocks compaction).
|
||||
|
||||
The summary sits outside the cached prefix, so a date in it is cache-safe.
|
||||
"""
|
||||
The summary sits outside the cached prefix, so a date in it is cache-safe."""
|
||||
try:
|
||||
from hermes_time import now as _hermes_now
|
||||
|
||||
@@ -2067,10 +2033,9 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
def on_session_end(self, session_id: str, messages: List[Dict[str, Any]]) -> None:
|
||||
"""Clear all per-session compaction state at a real session boundary.
|
||||
|
||||
Session end (CLI exit, gateway expiry, id rotation) — NOT /new or /reset. Every
|
||||
per-session flag/counter can contaminate the next live session (suppressed compression,
|
||||
stale cooldowns, misleading warnings), so the whole surface is reset here.
|
||||
"""
|
||||
Session end (CLI exit, gateway expiry, id rotation) — NOT /new or /reset. Every per-session
|
||||
flag/counter can contaminate the next live session (suppressed compression, stale cooldowns,
|
||||
misleading warnings), so the whole surface is reset here."""
|
||||
self._reset_session_compaction_state()
|
||||
|
||||
def _reset_real_usage_pairing(self) -> None:
|
||||
@@ -2179,10 +2144,9 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
):
|
||||
"""Best-effort read of a durable per-session value; ``default`` when unbound/unsupported/failed.
|
||||
|
||||
Returns ``(found, value)``: ``found`` is False when no read happened; ``value`` is None
|
||||
when the row held a non-numeric value. Defaults to the bound session row; pass
|
||||
``session_db``/``session_id`` to read another row (parent lineage).
|
||||
"""
|
||||
Returns ``(found, value)``: ``found`` is False when no read happened; ``value`` is None when the
|
||||
row held a non-numeric value. Defaults to the bound session row; pass
|
||||
``session_db``/``session_id`` to read another row (parent lineage)."""
|
||||
if session_db is None:
|
||||
session_db = getattr(self, "_session_db", None)
|
||||
if session_id is None:
|
||||
@@ -2284,9 +2248,8 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
def _record_structural_no_op(self, reason: str) -> None:
|
||||
"""Defer retries after a structural no-op WITHOUT striking the anti-thrash breaker.
|
||||
|
||||
Nothing eligible existed, so nothing was "ineffective"; striking would permanently
|
||||
disarm auto-compaction on short sessions. The backoff still stops per-turn re-scans.
|
||||
"""
|
||||
Nothing eligible existed, so nothing was "ineffective"; striking would permanently disarm
|
||||
auto-compaction on short sessions. The backoff still stops per-turn re-scans."""
|
||||
self._structural_no_op_backoff_until = time.monotonic() + self._STRUCTURAL_NO_OP_BACKOFF_SECONDS
|
||||
if not self.quiet_mode:
|
||||
logger.warning(
|
||||
@@ -2297,8 +2260,7 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
def record_rejected_compaction(self) -> None:
|
||||
"""Record a compaction rejected before commit as one ineffective strike.
|
||||
|
||||
Does not arm real-usage verification or touch the fallback streak (nothing was committed).
|
||||
"""
|
||||
Does not arm real-usage verification or touch the fallback streak (nothing was committed)."""
|
||||
self._record_ineffective_compression_verdict(self._ineffective_compression_count + 1)
|
||||
if not self.quiet_mode:
|
||||
logger.warning(
|
||||
@@ -2312,8 +2274,7 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
) -> None:
|
||||
"""Record one completed boundary and its summary quality.
|
||||
|
||||
``feasibility_skip=True`` is streak-neutral but still arms the real-usage effectiveness verdict.
|
||||
"""
|
||||
``feasibility_skip=True`` is streak-neutral but still arms the real-usage effectiveness verdict."""
|
||||
# A completed boundary proves compressibility: lift any structural no-op backoff.
|
||||
self._structural_no_op_backoff_until = 0.0
|
||||
self._verify_compaction_cleared_threshold = True
|
||||
@@ -2545,8 +2506,7 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
"""Compute the compaction trigger in tokens from the effective input budget.
|
||||
|
||||
Base is ``(context_length - max_tokens) * threshold_percent`` floored at MINIMUM_CONTEXT_LENGTH;
|
||||
when the floor binds it is capped at 85% of the budget so small windows can still fire.
|
||||
"""
|
||||
when the floor binds it is capped at 85% of the budget so small windows can still fire."""
|
||||
effective_window = context_length - (max_tokens or 0)
|
||||
if effective_window <= 0:
|
||||
effective_window = context_length
|
||||
@@ -2701,8 +2661,7 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
def maybe_seed_preflight_display_tokens(self, preflight_tokens: int) -> None:
|
||||
"""Seed ``last_prompt_tokens`` from a rough preflight estimate, display-only.
|
||||
|
||||
Seeds ONLY from the 0 state; the -1 sentinel and any real provider reading are preserved.
|
||||
"""
|
||||
Seeds ONLY from the 0 state; the -1 sentinel and any real provider reading are preserved."""
|
||||
_last = self.last_prompt_tokens
|
||||
if _last == 0 and preflight_tokens > _last:
|
||||
self.last_prompt_tokens = preflight_tokens
|
||||
@@ -2732,8 +2691,7 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
chars/4-underestimated scripts (Cyrillic, Thai, Arabic); bounded by two backstops: a real
|
||||
reading at/over threshold clears the baseline, and the overflow handler compacts reactively.
|
||||
Callers with a smaller (raw-messages) basis can only over-defer; the pre-API pressure check
|
||||
re-runs with the aligned basis.
|
||||
"""
|
||||
re-runs with the aligned basis."""
|
||||
if rough_tokens < self.threshold_tokens:
|
||||
return False
|
||||
# After compaction last_real_prompt_tokens is STALE (above threshold); defer one turn until real usage arrives.
|
||||
@@ -2756,8 +2714,7 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
def should_compress(self, prompt_tokens: int = None) -> bool:
|
||||
"""Return True when compression should run now.
|
||||
|
||||
Includes anti-thrash protection; see :meth:`should_compress_info` for the reason.
|
||||
"""
|
||||
Includes anti-thrash protection; see :meth:`should_compress_info` for the reason."""
|
||||
decision, _reason = self.should_compress_info(prompt_tokens)
|
||||
return decision
|
||||
|
||||
@@ -2765,8 +2722,7 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
"""Return ``(should_compress, reason)``.
|
||||
|
||||
``reason`` is None unless compression is needed but blocked: ``"cooldown:<seconds>"`` or
|
||||
``"ineffective"``. Callers should surface a warning when it is non-None.
|
||||
"""
|
||||
``"ineffective"``. Callers should surface a warning when it is non-None."""
|
||||
tokens = prompt_tokens if prompt_tokens is not None else self.last_prompt_tokens
|
||||
if tokens < self.threshold_tokens:
|
||||
return False, None
|
||||
@@ -2777,8 +2733,7 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
def _compression_block_reason(self) -> "str | None":
|
||||
"""Return the current automatic-compaction block reason, or None.
|
||||
|
||||
One of ``"cooldown:<seconds>"``, ``"structural_backoff:<seconds>"``, ``"ineffective"``.
|
||||
"""
|
||||
One of ``"cooldown:<seconds>"``, ``"structural_backoff:<seconds>"``, ``"ineffective"``."""
|
||||
_cooldown_remaining = self._summary_failure_cooldown_until - time.monotonic()
|
||||
if _cooldown_remaining > 0:
|
||||
return f"cooldown:{_cooldown_remaining:.0f}"
|
||||
@@ -2804,8 +2759,7 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
def _automatic_compression_blocked(self, *, ignore_cooldown: bool = False) -> bool:
|
||||
"""Return whether automatic compaction is in cooldown or tripped.
|
||||
|
||||
``ignore_cooldown=True`` skips only the summary-failure cooldown (overflow recovery path).
|
||||
"""
|
||||
``ignore_cooldown=True`` skips only the summary-failure cooldown (overflow recovery path)."""
|
||||
if not self._automatic_compression_blocked_locally(ignore_cooldown=ignore_cooldown):
|
||||
return False
|
||||
# Blocked locally: durable rows may have been cleared by another agent, so refresh before honouring.
|
||||
@@ -2879,9 +2833,8 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
|
||||
Returns ``(cut_idx, accumulated)``; ``cut_idx`` is the first protected index. On the budget
|
||||
break the cut stays at the last accepted row, or moves onto the breaking row when
|
||||
``cut_at_break``. Only the newest assistant turn's thinking is charged (#73624) unless the
|
||||
route echoes stale thinking every turn — must agree with the preflight estimate (#84371).
|
||||
"""
|
||||
``cut_at_break``. Only the newest assistant turn's thinking is charged (#73624) unless the route
|
||||
echoes stale thinking every turn — must agree with the preflight estimate (#84371)."""
|
||||
n = len(messages)
|
||||
newest_asst_idx = _last_assistant_index(messages)
|
||||
charge_all_thinking = self._stale_thinking_on_wire()
|
||||
@@ -2956,9 +2909,8 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
) -> bool:
|
||||
"""Replace the tool result at ``idx`` with a 1-line summary; True if modified.
|
||||
|
||||
``protected_skills`` (lower-cased) spares matching skill_view bodies; pass None for the
|
||||
pressure pass, which overrides the guard.
|
||||
"""
|
||||
``protected_skills`` (lower-cased) spares matching skill_view bodies; pass None for the pressure
|
||||
pass, which overrides the guard."""
|
||||
msg = result[idx]
|
||||
if msg.get("role") != "tool":
|
||||
return False
|
||||
@@ -2997,8 +2949,7 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
"""Pass 4: demote inside the protected tail when it alone exceeds the soft budget (#61932).
|
||||
|
||||
Keeps a short recent floor verbatim; overrides the skill guard (else the dead-end recurs).
|
||||
Returns the number of tool results demoted (arg truncations are logged but not counted).
|
||||
"""
|
||||
Returns the number of tool results demoted (arg truncations are logged but not counted)."""
|
||||
soft_ceiling = int(protect_tail_tokens * 1.5)
|
||||
demote_end = len(result) - min(_PRESSURE_KEEP_RECENT_MESSAGES, len(result))
|
||||
start = max(0, prune_boundary)
|
||||
@@ -3053,8 +3004,7 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
"""Replace old tool results with 1-line summaries; dedup, arg truncation, pressure demotion.
|
||||
|
||||
Token budget (when given) takes priority over the message-count floor. Returns
|
||||
``(pruned_messages, pruned_count)``.
|
||||
"""
|
||||
``(pruned_messages, pruned_count)``."""
|
||||
if not messages:
|
||||
return messages, 0
|
||||
result = [m.copy() for m in messages]
|
||||
@@ -3085,9 +3035,8 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
"""Deterministic, no-LLM tool-result prune gated on ``proactive_prune_tokens``.
|
||||
|
||||
Protects the tail by message COUNT only. A commit breaks the prompt cache, so it requires
|
||||
``proactive_prune_min_reclaim_tokens`` and a full regrowth runway; otherwise returns the
|
||||
INPUT object as ``(messages, 0)``.
|
||||
"""
|
||||
``proactive_prune_min_reclaim_tokens`` and a full regrowth runway; otherwise returns the INPUT
|
||||
object as ``(messages, 0)``."""
|
||||
if (
|
||||
self.proactive_prune_tokens <= 0
|
||||
or (current_tokens is not None and current_tokens < self.proactive_prune_tokens)
|
||||
@@ -3265,9 +3214,8 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine):
|
||||
) -> str:
|
||||
"""Build a deterministic handoff when the LLM summarizer is unavailable.
|
||||
|
||||
Keeps locally extractable anchors (recent user asks, actions, files/commands, errors)
|
||||
in the normal summary structure so downstream prompts recover gracefully.
|
||||
"""
|
||||
Keeps locally extractable anchors (recent user asks, actions, files/commands, errors) in the
|
||||
normal summary structure so downstream prompts recover gracefully."""
|
||||
anchors = self._fallback_anchors(turns_to_summarize)
|
||||
user_asks = anchors["user_asks"]
|
||||
completed = anchors["completed"]
|
||||
@@ -3337,9 +3285,8 @@ Summary generation was unavailable, so this is a best-effort deterministic fallb
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""Demote old tail tool results to recovery stubs (lean mode).
|
||||
|
||||
Keeps the newest ``_LEAN_TAIL_KEEP_TOOL_ROUNDS`` rounds verbatim; skill-marker rows
|
||||
are never touched. Returns a new list (untouched messages shared, demoted copied).
|
||||
"""
|
||||
Keeps the newest ``_LEAN_TAIL_KEEP_TOOL_ROUNDS`` rounds verbatim; skill-marker rows are never
|
||||
touched. Returns a new list (untouched messages shared, demoted copied)."""
|
||||
session_id = getattr(self, "_session_id", "") or ""
|
||||
tool_indices = [
|
||||
i for i in range(len(messages) - 1, tail_start - 1, -1)
|
||||
@@ -3424,9 +3371,8 @@ Summary generation was unavailable, so this is a best-effort deterministic fallb
|
||||
def _sample_summary_input(cls, content: str) -> str:
|
||||
"""Cap summarizer input by EVEN SAMPLING across the whole region (lean mode).
|
||||
|
||||
The single request also produces the session log, so coverage must be uniform:
|
||||
head+tail truncation would hide the entire middle from it.
|
||||
"""
|
||||
The single request also produces the session log, so coverage must be uniform: head+tail
|
||||
truncation would hide the entire middle from it."""
|
||||
if len(content) <= cls._SUMMARY_INPUT_MAX_CHARS:
|
||||
return content
|
||||
n = max(2, cls._SAMPLED_INPUT_SLICES)
|
||||
@@ -3453,8 +3399,7 @@ Summary generation was unavailable, so this is a best-effort deterministic fallb
|
||||
def _fallback_to_main_for_compression(self, e: Exception, reason: str) -> None:
|
||||
"""Switch from a separate ``summary_model`` back to the main model.
|
||||
|
||||
Records the aux failure, clears the summary model and the cooldown so the retry can run.
|
||||
"""
|
||||
Records the aux failure, clears the summary model and the cooldown so the retry can run."""
|
||||
self._summary_model_fallen_back = True
|
||||
logger.warning(
|
||||
"Summary model '%s' %s (%s). Falling back to main model '%s' for compression.",
|
||||
@@ -3472,9 +3417,8 @@ Summary generation was unavailable, so this is a best-effort deterministic fallb
|
||||
def _call_summary_llm(self, prompt: str, prompt_started_at: float) -> str:
|
||||
"""Issue the single aux summary call; return validated content text.
|
||||
|
||||
Raises RuntimeError for empty content or a length-truncated (PARTIAL) summary so the
|
||||
failure routes through main-model fallback + cooldown instead of wiping the compacted turns.
|
||||
"""
|
||||
Raises RuntimeError for empty content or a length-truncated (PARTIAL) summary so the failure
|
||||
routes through main-model fallback + cooldown instead of wiping the compacted turns."""
|
||||
call_kwargs: Dict[str, Any] = {
|
||||
"task": "compression",
|
||||
"main_runtime": {
|
||||
@@ -3543,8 +3487,7 @@ Summary generation was unavailable, so this is a best-effort deterministic fallb
|
||||
) -> Optional[str]:
|
||||
"""Generate a structured summary of conversation turns.
|
||||
|
||||
Iterative update when a previous summary exists. Returns None if all attempts fail.
|
||||
"""
|
||||
Iterative update when a previous summary exists. Returns None if all attempts fail."""
|
||||
prompt_started_at = time.monotonic()
|
||||
if self._compression_cancelled():
|
||||
raise AuxiliaryExplicitCancellation()
|
||||
@@ -3609,8 +3552,7 @@ Summary generation was unavailable, so this is a best-effort deterministic fallb
|
||||
) -> str:
|
||||
"""Assemble the summarizer prompt (fresh or iterative-update form).
|
||||
|
||||
Focus guidance is appended last so it takes precedence.
|
||||
"""
|
||||
Focus guidance is appended last so it takes precedence."""
|
||||
_memory_section = _memory_provider_section(memory_context)
|
||||
_today_str = _today_for_prompt()
|
||||
_section = _SECTION_INSTRUCTIONS[bool(has_user_turn)]
|
||||
@@ -3759,8 +3701,7 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
) -> Optional[str]:
|
||||
"""Classify a summary-call failure; retry on the main model once or arm a cooldown.
|
||||
|
||||
Returns the retry result, or None when the attempt is given up.
|
||||
"""
|
||||
Returns the retry result, or None when the attempt is given up."""
|
||||
# Only a genuine no-provider RuntimeError gets the long cooldown; empty/invalid-response
|
||||
# RuntimeErrors are transient and must get the main-model retry below first.
|
||||
if isinstance(e, RuntimeError) and "no llm provider configured" in str(e).lower():
|
||||
@@ -3870,9 +3811,8 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
def classify_summary_content(cls, content: Any) -> Optional[str]:
|
||||
"""Classify how *content* relates to a compaction summary.
|
||||
|
||||
Returns ``"standalone"`` (whole message is a handoff), ``"merged"`` (preserved
|
||||
content + delimiter + summary body), or None.
|
||||
"""
|
||||
Returns ``"standalone"`` (whole message is a handoff), ``"merged"`` (preserved content +
|
||||
delimiter + summary body), or None."""
|
||||
text = _content_text_for_contains(content).lstrip()
|
||||
# Merged summaries carry the handoff prefix after the delimiter; detect it there too.
|
||||
if _MERGED_SUMMARY_DELIMITER in text:
|
||||
@@ -4052,9 +3992,8 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
def _latest_user_task_snapshot(cls, messages: List[Dict[str, Any]]) -> Optional[str]:
|
||||
"""Return a deterministic task-snapshot line from the newest real user turn.
|
||||
|
||||
The summarizer must not invent the active-task anchor from a prompt example or a
|
||||
stale prior summary; this grounds it in the exact compacted turns.
|
||||
"""
|
||||
The summarizer must not invent the active-task anchor from a prompt example or a stale prior
|
||||
summary; this grounds it in the exact compacted turns."""
|
||||
# Reuse the runtime's real-user predicate so scaffolding rows can never anchor.
|
||||
from agent.conversation_compression import _is_real_user_message
|
||||
|
||||
@@ -4122,9 +4061,8 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
def _strip_context_summary_handoff_message(cls, message: Dict[str, Any]) -> Optional[Dict[str, Any]]:
|
||||
"""Drop stale handoff data while preserving merged prior-tail content.
|
||||
|
||||
Returns a copy for non-handoff rows, the unwrapped prior-tail content for merged
|
||||
handoffs (delimiter form, or legacy end-marker form), and ``None`` for standalone ones.
|
||||
"""
|
||||
Returns a copy for non-handoff rows, the unwrapped prior-tail content for merged handoffs
|
||||
(delimiter form, or legacy end-marker form), and ``None`` for standalone ones."""
|
||||
if not isinstance(message, dict):
|
||||
return message
|
||||
|
||||
@@ -4208,9 +4146,8 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
def _sanitize_tool_pairs(self, messages: List[Dict[str, Any]]) -> List[Dict[str, Any]]:
|
||||
"""Fix orphaned tool_call / tool_result pairs after compression.
|
||||
|
||||
Removes orphaned results and strips orphaned tool_calls (stubs would be dropped
|
||||
by repair_message_sequence when call_id != id).
|
||||
"""
|
||||
Removes orphaned results and strips orphaned tool_calls (stubs would be dropped by
|
||||
repair_message_sequence when call_id != id)."""
|
||||
from agent.agent_runtime_helpers import _classify_tool_call_orphans
|
||||
|
||||
_, result_call_ids, orphaned_result_msgs, missing_tool_calls = _classify_tool_call_orphans(messages)
|
||||
@@ -4275,9 +4212,8 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
def _effective_protect_first_n(self, messages: Optional[List[Dict[str, Any]]] = None) -> int:
|
||||
"""``protect_first_n`` decayed to 0 once the session has been compressed.
|
||||
|
||||
Otherwise early turns fossilize across compactions. After a restart the decayed
|
||||
state is inferred from handoff summaries in the resumed head.
|
||||
"""
|
||||
Otherwise early turns fossilize across compactions. After a restart the decayed state is
|
||||
inferred from handoff summaries in the resumed head."""
|
||||
if self.compression_count >= 1 or self._previous_summary:
|
||||
return 0
|
||||
if messages and self.protect_first_n > 0:
|
||||
@@ -4293,10 +4229,9 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
def _protect_head_size(self, messages: List[Dict[str, Any]]) -> int:
|
||||
"""Total head messages to protect.
|
||||
|
||||
The system prompt (index 0, if present) is always protected; ``protect_first_n``
|
||||
counts ADDITIONAL messages and decays after the first compaction so early turns
|
||||
don't fossilize (see _effective_protect_first_n).
|
||||
"""
|
||||
The system prompt (index 0, if present) is always protected; ``protect_first_n`` counts
|
||||
ADDITIONAL messages and decays after the first compaction so early turns don't fossilize (see
|
||||
_effective_protect_first_n)."""
|
||||
head = 0
|
||||
if messages and messages[0].get("role") == "system":
|
||||
head = 1
|
||||
@@ -4305,8 +4240,7 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
def _align_boundary_backward(self, messages: List[Dict[str, Any]], idx: int) -> int:
|
||||
"""Pull a compress-end boundary backward so a tool_call/result group is not split.
|
||||
|
||||
A split group leaves orphaned tail results that ``_sanitize_tool_pairs`` then drops.
|
||||
"""
|
||||
A split group leaves orphaned tail results that ``_sanitize_tool_pairs`` then drops."""
|
||||
if idx <= 0 or idx >= len(messages):
|
||||
return idx
|
||||
check = idx - 1
|
||||
@@ -4321,8 +4255,7 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
def _real_user_indices_desc(cls, messages: List[Dict[str, Any]], head_end: int) -> list[int]:
|
||||
"""Indices (newest first) of actionable, non-synthetic user turns at or after *head_end*.
|
||||
|
||||
Compaction handoffs and blank platform echoes never count as the real request.
|
||||
"""
|
||||
Compaction handoffs and blank platform echoes never count as the real request."""
|
||||
return [
|
||||
i for i in range(len(messages) - 1, head_end - 1, -1)
|
||||
if cls._is_actionable_user_turn(messages[i])
|
||||
@@ -4336,8 +4269,7 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
def _find_last_assistant_message_idx(self, messages: List[Dict[str, Any]], head_end: int) -> int:
|
||||
"""Return the last text-bearing, non-summary assistant reply at or after *head_end*, or -1.
|
||||
|
||||
Falls back to the last non-summary assistant of any kind when none has text.
|
||||
"""
|
||||
Falls back to the last non-summary assistant of any kind when none has text."""
|
||||
last_any = -1
|
||||
for i in range(len(messages) - 1, head_end - 1, -1):
|
||||
msg = messages[i]
|
||||
@@ -4361,8 +4293,7 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
) -> int:
|
||||
"""Guarantee the most recent assistant reply is in the protected tail (#29824).
|
||||
|
||||
Re-aligns backward so a preceding tool group is not split.
|
||||
"""
|
||||
Re-aligns backward so a preceding tool group is not split."""
|
||||
last_asst_idx = self._find_last_assistant_message_idx(messages, head_end)
|
||||
if last_asst_idx < 0 or last_asst_idx >= cut_idx:
|
||||
return cut_idx
|
||||
@@ -4381,11 +4312,10 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
) -> int:
|
||||
"""Guarantee the most recent user message is in the protected tail.
|
||||
|
||||
Tool-group alignment can pull the cut past the last user message; once summarized, the
|
||||
prefix tells the model to answer only messages AFTER the summary, so the active ask
|
||||
silently vanishes. If the head_end clamp would strand the user without its reply, the
|
||||
cut is pushed forward past the whole turn-pair instead so it is summarised as completed.
|
||||
"""
|
||||
Tool-group alignment can pull the cut past the last user message; once summarized, the prefix
|
||||
tells the model to answer only messages AFTER the summary, so the active ask silently vanishes.
|
||||
If the head_end clamp would strand the user without its reply, the cut is pushed forward past
|
||||
the whole turn-pair instead so it is summarised as completed."""
|
||||
last_user_idx = self._find_last_user_message_idx(messages, head_end)
|
||||
if last_user_idx < 0 or last_user_idx >= cut_idx:
|
||||
return cut_idx
|
||||
@@ -4417,8 +4347,7 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
) -> int:
|
||||
"""Guarantee the last N actionable user messages are in the protected tail.
|
||||
|
||||
n <= 1 delegates to the single-message method. Only real user turns count.
|
||||
"""
|
||||
n <= 1 delegates to the single-message method. Only real user turns count."""
|
||||
if n <= 1:
|
||||
return self._ensure_last_user_message_in_tail(messages, cut_idx, head_end)
|
||||
|
||||
@@ -4435,8 +4364,7 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
def _find_turn_pair_end(self, messages: List[Dict[str, Any]], user_idx: int) -> int:
|
||||
"""Return the index after the turn-pair (user -> assistant -> tools) at *user_idx*.
|
||||
|
||||
Returns ``user_idx + 1`` when there is no reply yet.
|
||||
"""
|
||||
Returns ``user_idx + 1`` when there is no reply yet."""
|
||||
n = len(messages)
|
||||
idx = user_idx + 1
|
||||
if idx >= n:
|
||||
@@ -4451,9 +4379,8 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
def _stale_thinking_on_wire(self) -> bool:
|
||||
"""Whether the active route replays stale thinking text on every turn.
|
||||
|
||||
Tail-budget walks and the preflight trigger must charge the SAME policy or
|
||||
compaction loops forever.
|
||||
"""
|
||||
Tail-budget walks and the preflight trigger must charge the SAME policy or compaction loops
|
||||
forever."""
|
||||
try:
|
||||
from agent.message_sanitization import stale_thinking_reaches_wire
|
||||
|
||||
@@ -4469,9 +4396,8 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
) -> int:
|
||||
"""Walk backward accumulating tokens until the budget; return the tail start index.
|
||||
|
||||
May exceed the budget by up to 1.5x to avoid cutting inside an oversized message;
|
||||
never splits a tool group; keeps the last user message in the tail.
|
||||
"""
|
||||
May exceed the budget by up to 1.5x to avoid cutting inside an oversized message; never splits a
|
||||
tool group; keeps the last user message in the tail."""
|
||||
if token_budget is None:
|
||||
token_budget = self.tail_token_budget
|
||||
n = len(messages)
|
||||
@@ -4520,8 +4446,7 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
def has_content_to_compress(self, messages: List[Dict[str, Any]]) -> bool:
|
||||
"""Return True if there is a non-empty middle region to compact.
|
||||
|
||||
Lets the gateway ``/compress`` guard skip the LLM call when everything is protected.
|
||||
"""
|
||||
Lets the gateway ``/compress`` guard skip the LLM call when everything is protected."""
|
||||
compress_start = self._align_boundary_forward(messages, self._protect_head_size(messages))
|
||||
compress_end = self._find_tail_cut_by_tokens(messages, compress_start)
|
||||
return compress_start < compress_end
|
||||
@@ -4532,10 +4457,9 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
) -> "_HandoffScan":
|
||||
"""Rehydrate ``_previous_summary`` / user-turn provenance from in-transcript handoffs.
|
||||
|
||||
Handoff rows are removed from the summarizer window (merged handoffs unwrap to their
|
||||
prior-tail content) and ``tail_start`` advances past a handoff beyond the window. The
|
||||
pre-scan state is captured so an aborted attempt can roll the mutation back (#57835).
|
||||
"""
|
||||
Handoff rows are removed from the summarizer window (merged handoffs unwrap to their prior-tail
|
||||
content) and ``tail_start`` advances past a handoff beyond the window. The pre-scan state is
|
||||
captured so an aborted attempt can roll the mutation back (#57835)."""
|
||||
scan = _HandoffScan(
|
||||
turns_to_summarize=turns_to_summarize, summary_indices=set(), tail_start=compress_end,
|
||||
# Snapshot so an aborted attempt can roll back the self-heal mutation (#57835).
|
||||
@@ -4695,9 +4619,8 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
) -> bool:
|
||||
"""Abort (messages unchanged) on a terminal failure or when configured to; True when aborted.
|
||||
|
||||
Access/quota, network, truncated and empty-content failures ALWAYS abort (#29559);
|
||||
otherwise ``abort_on_summary_failure`` decides between abort and the static fallback.
|
||||
"""
|
||||
Access/quota, network, truncated and empty-content failures ALWAYS abort (#29559); otherwise
|
||||
``abort_on_summary_failure`` decides between abort and the static fallback."""
|
||||
terminal_failure = next(
|
||||
(
|
||||
(failure_class, message)
|
||||
@@ -4789,10 +4712,9 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
) -> tuple[str, bool, bool, Optional[int]]:
|
||||
"""Pick the summary row's role so template-visible alternation holds.
|
||||
|
||||
Returns ``(summary_role, merge_into_tail, force_user_leading, first_tail_visible_idx)``.
|
||||
Roles read the assembled (post-strip) head/tail and are TEMPLATE-VISIBLE: Mistral-strict
|
||||
templates skip tool rows for alternation, so alternate against what the template counts.
|
||||
"""
|
||||
Returns ``(summary_role, merge_into_tail, force_user_leading, first_tail_visible_idx)``. Roles
|
||||
read the assembled (post-strip) head/tail and are TEMPLATE-VISIBLE: Mistral-strict templates
|
||||
skip tool rows for alternation, so alternate against what the template counts."""
|
||||
last_head_role: Optional[str] = "user"
|
||||
if compressed:
|
||||
# None = all-exempt head: the summary opens the visible sequence and must be "user".
|
||||
@@ -4811,10 +4733,9 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
),
|
||||
(None, None),
|
||||
)
|
||||
# System-only head: the summary is the first visible message and Anthropic requires
|
||||
# role=user (#52160). Zero-user-turn guard (#58753): if no user row with non-empty TEXT
|
||||
# survives, the summary must be role="user" or OpenAI-compatible backends reject.
|
||||
# Image-only rows don't count.
|
||||
# System-only head: the summary is the first visible message and Anthropic requires role=user
|
||||
# (#52160). Zero-user-turn guard (#58753): if no user row with non-empty TEXT survives, the
|
||||
# summary must be role="user" or OpenAI-compatible backends reject. Image-only rows don't count.
|
||||
force_user_leading = compress_start == 0 or last_head_role == "system" or not any(
|
||||
m.get("role") == "user" and bool(_content_text_for_contains(m.get("content")).strip())
|
||||
for m in (*compressed, *tail_messages)
|
||||
@@ -4922,11 +4843,10 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""Compress conversation messages by summarizing middle turns.
|
||||
|
||||
Prunes tool results and blank echo rows (survives an abort), protects head and a
|
||||
token-budget tail, summarizes the middle, then cleans orphaned tool pairs.
|
||||
``force`` clears the failure cooldown and bypasses the feasibility skip;
|
||||
``bypass_cooldown`` runs the summary LLM without clearing the cooldown (#100661).
|
||||
"""
|
||||
Prunes tool results and blank echo rows (survives an abort), protects head and a token-budget
|
||||
tail, summarizes the middle, then cleans orphaned tool pairs. ``force`` clears the failure
|
||||
cooldown and bypasses the feasibility skip; ``bypass_cooldown`` runs the summary LLM without
|
||||
clearing the cooldown (#100661)."""
|
||||
telemetry = self._begin_compress_attempt(current_tokens, force)
|
||||
n_messages = len(messages)
|
||||
# Only need head + 3 tail messages minimum (token budget decides the real tail size)
|
||||
@@ -5050,9 +4970,8 @@ This compaction should PRIORITISE preserving all information related to the focu
|
||||
def is_compaction_summary_message(message: Any) -> bool:
|
||||
"""Return True when *message* is a context-compaction handoff summary.
|
||||
|
||||
Public API. Uses the metadata key, falling back to content heuristics because the key
|
||||
is stripped by wire sanitizers and some session-store round-trips.
|
||||
"""
|
||||
Public API. Uses the metadata key, falling back to content heuristics because the key is stripped by
|
||||
wire sanitizers and some session-store round-trips."""
|
||||
if isinstance(message, dict):
|
||||
return ContextCompressor._is_context_summary_message(message)
|
||||
return ContextCompressor._is_context_summary_content(message)
|
||||
@@ -5167,8 +5086,7 @@ def history_before_user_originated_turn(
|
||||
) -> tuple[List[Dict[str, Any]], Dict[str, Any]]:
|
||||
"""Return a rewind prefix and canonical live view for ``index``.
|
||||
|
||||
A composite carrier keeps its handoff scaffold at the new history head.
|
||||
"""
|
||||
A composite carrier keeps its handoff scaffold at the new history head."""
|
||||
if index < 0 or index >= len(messages):
|
||||
raise IndexError("user turn index is outside the transcript")
|
||||
handoff, live_view = split_user_originated_turn(messages[index])
|
||||
@@ -5183,8 +5101,7 @@ def history_before_user_originated_turn(
|
||||
def retryable_user_text(content: Any) -> str:
|
||||
"""Return lossless retry text or raise before destructive mutation.
|
||||
|
||||
Media and unknown structured parts fail closed (no attachment replay protocol).
|
||||
"""
|
||||
Media and unknown structured parts fail closed (no attachment replay protocol)."""
|
||||
if isinstance(content, str):
|
||||
text = content
|
||||
elif isinstance(content, list):
|
||||
@@ -5215,8 +5132,7 @@ def retryable_user_text(content: Any) -> str:
|
||||
def _handoff_carries_live_user_content(message: Any) -> bool:
|
||||
"""Return True when a summary-bearing row still carries a live user ask.
|
||||
|
||||
Callers must pre-filter with ``is_compaction_summary_message``.
|
||||
"""
|
||||
Callers must pre-filter with ``is_compaction_summary_message``."""
|
||||
if not isinstance(message, dict):
|
||||
return False
|
||||
return ContextCompressor._strip_context_summary_handoff_message(message) is not None
|
||||
@@ -5225,8 +5141,7 @@ def _handoff_carries_live_user_content(message: Any) -> bool:
|
||||
def reference_handoff_would_drive_next_model_call(messages: Optional[List[Dict[str, Any]]]) -> bool:
|
||||
"""Return True when the next model call would be driven only by a handoff (#80622).
|
||||
|
||||
Mid tool-loop compression is allowed: trailing tool rows mean an in-flight exchange.
|
||||
"""
|
||||
Mid tool-loop compression is allowed: trailing tool rows mean an in-flight exchange."""
|
||||
if not messages:
|
||||
return False
|
||||
|
||||
@@ -5276,6 +5191,5 @@ def reference_handoff_would_drive_next_model_call(messages: Optional[List[Dict[s
|
||||
def is_user_originated_turn(message: Any) -> bool:
|
||||
"""Return True for human-authored user turns (not compaction scaffolding).
|
||||
|
||||
Dispatchers must use this instead of a bare role check (#80622).
|
||||
"""
|
||||
Dispatchers must use this instead of a bare role check (#80622)."""
|
||||
return user_originated_turn_view(message) is not None
|
||||
|
||||
Reference in New Issue
Block a user