refactor(agent/anthropic_*,background_review): tool_choice table, shared block helpers (_split_blank_text_blocks, _carry_cache_control, _block_ids), predicate collapse in endpoints, suppress() for swallow-all guards
This commit is contained in:
+13
-13
@@ -12,6 +12,7 @@ import logging
|
||||
import math
|
||||
import re
|
||||
import subprocess
|
||||
from contextlib import suppress
|
||||
from pathlib import Path # noqa: F401 (tests patch ``anthropic_adapter.Path.home``)
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
@@ -51,11 +52,9 @@ def _get_anthropic_sdk():
|
||||
"""Return the ``anthropic`` SDK module, importing lazily. None if not installed."""
|
||||
global _anthropic_sdk
|
||||
if _anthropic_sdk is ...:
|
||||
try:
|
||||
with suppress(Exception): # ImportError or FeatureUnavailable — fall through to the import below
|
||||
from tools.lazy_deps import ensure as _lazy_ensure
|
||||
_lazy_ensure("provider.anthropic", prompt=False)
|
||||
except Exception: # ImportError or FeatureUnavailable — fall through to the import below
|
||||
pass
|
||||
try:
|
||||
import anthropic as _sdk
|
||||
_anthropic_sdk = _sdk
|
||||
@@ -603,6 +602,10 @@ def _thinking_kwargs(reasoning_config: Dict[str, Any], model: str, effective_max
|
||||
}
|
||||
|
||||
|
||||
# OpenAI tool_choice -> Anthropic; any other string is a forced tool name.
|
||||
_TOOL_CHOICE_MAP = {None: {"type": "auto"}, "auto": {"type": "auto"}, "required": {"type": "any"}}
|
||||
|
||||
|
||||
def build_anthropic_kwargs(
|
||||
model: str, messages: List[Dict], tools: Optional[List[Dict]], max_tokens: Optional[int],
|
||||
reasoning_config: Optional[Dict[str, Any]], tool_choice: Optional[str] = None,
|
||||
@@ -644,17 +647,14 @@ def build_anthropic_kwargs(
|
||||
|
||||
if anthropic_tools:
|
||||
kwargs["tools"] = anthropic_tools
|
||||
if tool_choice == "auto" or tool_choice is None:
|
||||
kwargs["tool_choice"] = {"type": "auto"}
|
||||
elif tool_choice == "required":
|
||||
kwargs["tool_choice"] = {"type": "any"}
|
||||
elif tool_choice == "none":
|
||||
if tool_choice == "none":
|
||||
kwargs.pop("tools", None) # no Anthropic "none" — omit tools to prevent use
|
||||
elif isinstance(tool_choice, str):
|
||||
# Under OAuth every tools[] entry is mcp__-prefixed/aliased; the forced name must go
|
||||
# through the same normalizer or it (a) leaks the literal trigger string and (b) names
|
||||
# a tool that no longer exists -> 400.
|
||||
kwargs["tool_choice"] = {"type": "tool", "name": to_wire(tool_choice) if to_wire else tool_choice}
|
||||
elif tool_choice is None or isinstance(tool_choice, str):
|
||||
# A forced tool name goes through the OAuth normalizer too: every tools[] entry is
|
||||
# mcp__-prefixed/aliased there, so the literal would leak and name a nonexistent tool.
|
||||
kwargs["tool_choice"] = _TOOL_CHOICE_MAP.get(tool_choice) or {
|
||||
"type": "tool", "name": to_wire(tool_choice) if to_wire else tool_choice
|
||||
}
|
||||
|
||||
if reasoning_config and isinstance(reasoning_config, dict):
|
||||
kwargs.update(_thinking_kwargs(reasoning_config, model, effective_max_tokens))
|
||||
|
||||
@@ -56,12 +56,8 @@ _KIMI_FAMILY_EXACT_SLUGS = frozenset({"k3"})
|
||||
def _model_name_is_kimi_family(model: str | None) -> bool:
|
||||
if not isinstance(model, str):
|
||||
return False
|
||||
m = model.strip().lower()
|
||||
if not m:
|
||||
return False
|
||||
if "/" in m: # ``moonshotai/kimi-k2.5`` -> ``kimi-k2.5``
|
||||
m = m.rsplit("/", 1)[-1]
|
||||
return m in _KIMI_FAMILY_EXACT_SLUGS or m.startswith(_KIMI_FAMILY_MODEL_PREFIXES)
|
||||
m = model.strip().lower().rsplit("/", 1)[-1] # ``moonshotai/kimi-k2.5`` -> ``kimi-k2.5``
|
||||
return bool(m) and (m in _KIMI_FAMILY_EXACT_SLUGS or m.startswith(_KIMI_FAMILY_MODEL_PREFIXES))
|
||||
|
||||
|
||||
def _is_kimi_family_endpoint(base_url: str | None, model: str | None = None) -> bool:
|
||||
@@ -69,11 +65,11 @@ def _is_kimi_family_endpoint(base_url: str | None, model: str | None = None) ->
|
||||
moonshot.ai / moonshot.cn host, or any endpoint (e.g. a private gateway) whose *model* is in
|
||||
the Kimi family — the upstream enforces Kimi's thinking semantics regardless of hostname.
|
||||
Decides whether unsigned reasoning_content-derived thinking blocks are preserved on replay."""
|
||||
if _is_kimi_coding_endpoint(base_url):
|
||||
return True
|
||||
if any(base_url_host_matches(base_url or "", d) for d in ("api.kimi.com", "moonshot.ai", "moonshot.cn")):
|
||||
return True
|
||||
return _model_name_is_kimi_family(model)
|
||||
return (
|
||||
_is_kimi_coding_endpoint(base_url)
|
||||
or any(base_url_host_matches(base_url or "", d) for d in ("api.kimi.com", "moonshot.ai", "moonshot.cn"))
|
||||
or _model_name_is_kimi_family(model)
|
||||
)
|
||||
|
||||
|
||||
def _is_deepseek_anthropic_endpoint(base_url: str | None) -> bool:
|
||||
@@ -81,9 +77,7 @@ def _is_deepseek_anthropic_endpoint(base_url: str | None) -> bool:
|
||||
blocks to round-trip while the generic third-party path strips them; its blocks are unsigned,
|
||||
so it gets the same strip-signed / keep-unsigned policy as Kimi. Pinned to the ``/anthropic``
|
||||
path so the OpenAI-compatible base URL is not misclassified."""
|
||||
if not base_url_host_matches(base_url or "", "api.deepseek.com"):
|
||||
return False
|
||||
return "/anthropic" in _normalized_lower(base_url)
|
||||
return base_url_host_matches(base_url or "", "api.deepseek.com") and "/anthropic" in _normalized_lower(base_url)
|
||||
|
||||
|
||||
def _is_nous_portal_endpoint(base_url: str | None) -> bool:
|
||||
@@ -99,9 +93,7 @@ def _is_nous_portal_endpoint(base_url: str | None) -> bool:
|
||||
override = _nous_inference_env_override()
|
||||
except Exception:
|
||||
return False
|
||||
if not override:
|
||||
return False
|
||||
override_host = base_url_hostname(override)
|
||||
override_host = base_url_hostname(override) if override else ""
|
||||
return bool(override_host) and base_url_hostname(base_url or "") == override_host
|
||||
|
||||
|
||||
@@ -109,13 +101,10 @@ def _requires_bearer_auth(base_url: str | None) -> bool:
|
||||
"""Providers needing ``Authorization: Bearer`` instead of ``x-api-key``: MiniMax, Azure AI
|
||||
Foundry, Palantir Foundry's LLM proxy, CommandCode, Nous Portal. Palantir/CommandCode use
|
||||
hostname matching (not substring) so ``evil.com/palantirfoundry`` paths don't trigger it."""
|
||||
if _is_nous_portal_endpoint(base_url):
|
||||
return True
|
||||
normalized = _normalized_lower(base_url)
|
||||
if not normalized:
|
||||
return False
|
||||
return (
|
||||
normalized.startswith(_MINIMAX_ANTHROPIC_PREFIXES)
|
||||
_is_nous_portal_endpoint(base_url)
|
||||
or normalized.startswith(_MINIMAX_ANTHROPIC_PREFIXES)
|
||||
or "azure.com" in normalized
|
||||
or base_url_host_matches(normalized, "palantirfoundry.com")
|
||||
or base_url_host_matches(normalized, "api.commandcode.ai")
|
||||
|
||||
@@ -29,8 +29,6 @@ _BEDROCK_REGION_PREFIXES = (
|
||||
)
|
||||
|
||||
|
||||
# ----- small shared predicates -----
|
||||
|
||||
|
||||
def _block_type(b: Any) -> Any:
|
||||
"""``type`` of a dict block, None for non-dicts."""
|
||||
@@ -59,6 +57,14 @@ def _text_block(text: str) -> Dict[str, str]:
|
||||
return {"type": "text", "text": text}
|
||||
|
||||
|
||||
def _text_block_with_citations(text: Any, cits: Any) -> Dict[str, Any]:
|
||||
"""Text block carrying ``citations`` only when it is a non-empty list (the only input-valid shape)."""
|
||||
block: Dict[str, Any] = _text_block(text)
|
||||
if isinstance(cits, list) and cits:
|
||||
block["citations"] = cits
|
||||
return block
|
||||
|
||||
|
||||
def _parse_tool_args(raw: Any) -> Any:
|
||||
"""JSON-decode a tool_call ``arguments`` string; non-strings pass through, bad JSON -> {}."""
|
||||
try:
|
||||
@@ -71,7 +77,32 @@ def _strip_thinking(blocks: List[Any]) -> List[Any]:
|
||||
return [b for b in blocks if _block_type(b) not in _THINKING_TYPES]
|
||||
|
||||
|
||||
# ----- model / tool conversion -----
|
||||
def _block_ids(blocks: List[Any], btype: str, key: str) -> set:
|
||||
return {b.get(key) for b in blocks if _block_type(b) == btype}
|
||||
|
||||
|
||||
def _carry_cache_control(out: Dict[str, Any], b: Dict[str, Any]) -> Dict[str, Any]:
|
||||
"""Copy a dict-valued ``cache_control`` marker from ``b`` onto ``out`` (returned)."""
|
||||
if _cache_control_of(b) is not None:
|
||||
out["cache_control"] = b["cache_control"]
|
||||
return out
|
||||
|
||||
|
||||
def _split_blank_text_blocks(blocks: List[Any]) -> Tuple[List[Any], Any, List[int]]:
|
||||
"""``(kept, relocated_cache_control, dropped_indexes)``: drop blank text blocks, remembering
|
||||
the cache_control of the last one dropped so the caller can relocate the breakpoint."""
|
||||
kept: List[Any] = []
|
||||
relocated_cc = None
|
||||
dropped: List[int] = []
|
||||
for i, blk in enumerate(blocks):
|
||||
if _is_blank_text_block(blk):
|
||||
if _cache_control_of(blk) is not None:
|
||||
relocated_cc = blk["cache_control"]
|
||||
dropped.append(i)
|
||||
else:
|
||||
kept.append(blk)
|
||||
return kept, relocated_cc, dropped
|
||||
|
||||
|
||||
|
||||
def _is_bedrock_model_id(model: str) -> bool:
|
||||
@@ -99,6 +130,10 @@ def _sanitize_tool_id(tool_id: str) -> str:
|
||||
return re.sub(r"[^a-zA-Z0-9_-]", "_", tool_id) or "tool_0"
|
||||
|
||||
|
||||
def _tool_use_block(tool_id: Any, name: Any, tool_input: Any) -> Dict[str, Any]:
|
||||
return {"type": "tool_use", "id": _sanitize_tool_id(tool_id), "name": name, "input": tool_input}
|
||||
|
||||
|
||||
def _normalize_tool_input_schema(schema: Any) -> Dict[str, Any]:
|
||||
"""Normalize a tool schema for Anthropic's validator: collapse nullable unions (``anyOf:
|
||||
[{type: string}, {type: null}]`` from Pydantic/MCP optional fields) to the non-null branch —
|
||||
@@ -149,8 +184,6 @@ def convert_tools_to_anthropic(tools: List[Dict]) -> List[Dict]:
|
||||
return result
|
||||
|
||||
|
||||
# ----- content-part conversion -----
|
||||
|
||||
|
||||
def _image_source_from_openai_url(url: str) -> Dict[str, str]:
|
||||
"""OpenAI image URL / data URL -> Anthropic image ``source``."""
|
||||
@@ -176,10 +209,7 @@ def _convert_content_part_to_anthropic(part: Any) -> Optional[Dict[str, Any]]:
|
||||
if ptype in ("input_text", "text"):
|
||||
# Rebuild from whitelisted fields only: stored SDK text blocks carry output-only siblings
|
||||
# (parsed_output, citations=None) that the INPUT schema rejects with 400.
|
||||
block: Dict[str, Any] = _text_block(part.get("text", ""))
|
||||
cits = part.get("citations")
|
||||
if ptype == "text" and isinstance(cits, list) and cits:
|
||||
block["citations"] = cits
|
||||
block = _text_block_with_citations(part.get("text", ""), part.get("citations") if ptype == "text" else None)
|
||||
elif ptype in {"image_url", "input_image"}:
|
||||
image_value = part.get("image_url", {})
|
||||
url = image_value.get("url", "") if isinstance(image_value, dict) else str(image_value or "")
|
||||
@@ -269,15 +299,10 @@ def _safe_text(text: Any) -> str:
|
||||
"""``text`` if non-whitespace, else the placeholder. A blank text block stored in history (e.g.
|
||||
by compression) is replayed on every turn and wedges the session with HTTP 400; the placeholder
|
||||
is self-healing. Mirrors ``bedrock_adapter._safe_text`` (kept separate on purpose)."""
|
||||
if text is None:
|
||||
return _EMPTY_TEXT_PLACEHOLDER
|
||||
if not isinstance(text, str):
|
||||
text = str(text)
|
||||
text = "" if text is None else str(text)
|
||||
return text if text.strip() else _EMPTY_TEXT_PLACEHOLDER
|
||||
|
||||
|
||||
# ----- replay-block sanitizing (per-type whitelist) -----
|
||||
|
||||
|
||||
def _replay_text(b: Dict[str, Any]) -> Optional[Dict[str, Any]]:
|
||||
# Drop blank blocks rather than coerce in place: the caller relocates any cache_control and
|
||||
@@ -285,13 +310,7 @@ def _replay_text(b: Dict[str, Any]) -> Optional[Dict[str, Any]]:
|
||||
# model-visible noise next to real blocks.
|
||||
if _is_blank_text_block(b):
|
||||
return None
|
||||
out: Dict[str, Any] = _text_block(b["text"])
|
||||
cits = b.get("citations") # input-valid ONLY as a non-empty list
|
||||
if isinstance(cits, list) and cits:
|
||||
out["citations"] = cits
|
||||
if _cache_control_of(b) is not None:
|
||||
out["cache_control"] = b["cache_control"]
|
||||
return out
|
||||
return _carry_cache_control(_text_block_with_citations(b["text"], b.get("citations")), b)
|
||||
|
||||
|
||||
def _replay_thinking(b: Dict[str, Any]) -> Dict[str, Any]:
|
||||
@@ -306,13 +325,7 @@ def _replay_redacted_thinking(b: Dict[str, Any]) -> Optional[Dict[str, Any]]:
|
||||
|
||||
|
||||
def _replay_tool_use(b: Dict[str, Any]) -> Dict[str, Any]:
|
||||
out = {
|
||||
"type": "tool_use", "id": _sanitize_tool_id(b.get("id", "")), "name": b.get("name", ""),
|
||||
"input": b.get("input", {}),
|
||||
}
|
||||
if _cache_control_of(b) is not None:
|
||||
out["cache_control"] = b["cache_control"]
|
||||
return out
|
||||
return _carry_cache_control(_tool_use_block(b.get("id", ""), b.get("name", ""), b.get("input", {})), b)
|
||||
|
||||
|
||||
def _replay_image(b: Dict[str, Any]) -> Optional[Dict[str, Any]]:
|
||||
@@ -347,8 +360,6 @@ def _apply_assistant_cache_control_to_last_cacheable_block(blocks: List[Dict[str
|
||||
break
|
||||
|
||||
|
||||
# ----- per-message conversion -----
|
||||
|
||||
|
||||
def _replay_ordered_blocks(m: Dict[str, Any], ordered_blocks: List[Any]) -> Optional[List[Dict[str, Any]]]:
|
||||
"""Interleaved-thinking replay: rebuild the assistant turn from the verbatim block list
|
||||
@@ -411,22 +422,15 @@ def _convert_assistant_message(m: Dict[str, Any]) -> Dict[str, Any]:
|
||||
# exactly the blank block).
|
||||
relocated_cc = None
|
||||
if isinstance(content, list):
|
||||
for blk in _convert_content_to_anthropic(content):
|
||||
if _is_blank_text_block(blk):
|
||||
if _cache_control_of(blk) is not None:
|
||||
relocated_cc = blk["cache_control"]
|
||||
continue
|
||||
blocks.append(blk)
|
||||
kept, relocated_cc, _ = _split_blank_text_blocks(_convert_content_to_anthropic(content))
|
||||
blocks.extend(kept)
|
||||
elif content and str(content).strip():
|
||||
blocks.append(_text_block(str(content)))
|
||||
for tc in m.get("tool_calls", []):
|
||||
if not tc or not isinstance(tc, dict):
|
||||
continue
|
||||
fn = tc.get("function", {})
|
||||
blocks.append({
|
||||
"type": "tool_use", "id": _sanitize_tool_id(tc.get("id", "")), "name": fn.get("name", ""),
|
||||
"input": _parse_tool_args(fn.get("arguments", "{}")),
|
||||
})
|
||||
blocks.append(_tool_use_block(tc.get("id", ""), fn.get("name", ""), _parse_tool_args(fn.get("arguments", "{}"))))
|
||||
# Kimi's /coding endpoint requires reasoning_content on replayed tool-call turns — even ""
|
||||
# (injected as a fallback upstream). Prepend, since thinking must precede text/tool_use. Skip
|
||||
# when reasoning_details already supplied (signed) thinking blocks: a duplicate unsigned one
|
||||
@@ -499,8 +503,6 @@ def _convert_user_message(content: Any) -> Dict[str, Any]:
|
||||
return {"role": "user", "content": content}
|
||||
|
||||
|
||||
# ----- whole-list passes -----
|
||||
|
||||
|
||||
def _strip_orphaned_tool_blocks(result: List[Dict[str, Any]]) -> None:
|
||||
"""Strip tool_use blocks with no matching tool_result, and vice versa. Compression/truncation
|
||||
@@ -511,14 +513,14 @@ def _strip_orphaned_tool_blocks(result: List[Dict[str, Any]]) -> None:
|
||||
for i, m in enumerate(result):
|
||||
if m.get("role") != "assistant" or not isinstance(m.get("content"), list):
|
||||
continue
|
||||
tool_use_ids_in_turn = {b.get("id") for b in m["content"] if _block_type(b) == "tool_use"}
|
||||
tool_use_ids_in_turn = _block_ids(m["content"], "tool_use", "id")
|
||||
if not tool_use_ids_in_turn:
|
||||
continue
|
||||
adjacent_result_ids: set = set()
|
||||
if i + 1 < len(result):
|
||||
nxt = result[i + 1]
|
||||
if nxt.get("role") == "user" and isinstance(nxt.get("content"), list):
|
||||
adjacent_result_ids = {b.get("tool_use_id") for b in nxt["content"] if _block_type(b) == "tool_result"}
|
||||
adjacent_result_ids = _block_ids(nxt["content"], "tool_result", "tool_use_id")
|
||||
orphaned = tool_use_ids_in_turn - adjacent_result_ids
|
||||
if not orphaned:
|
||||
continue
|
||||
@@ -531,13 +533,10 @@ def _strip_orphaned_tool_blocks(result: List[Dict[str, Any]]) -> None:
|
||||
m["content"] = kept if kept else [_text_block("(tool call removed)")]
|
||||
|
||||
# Pass 2: tool_result whose tool_use no longer exists anywhere.
|
||||
surviving_tool_use_ids = {
|
||||
b.get("id")
|
||||
for m in result
|
||||
if m.get("role") == "assistant" and isinstance(m.get("content"), list)
|
||||
for b in m["content"]
|
||||
if _block_type(b) == "tool_use"
|
||||
}
|
||||
surviving_tool_use_ids: set = set()
|
||||
for m in result:
|
||||
if m.get("role") == "assistant" and isinstance(m.get("content"), list):
|
||||
surviving_tool_use_ids |= _block_ids(m["content"], "tool_use", "id")
|
||||
for m in result:
|
||||
if m.get("role") != "user" or not isinstance(m.get("content"), list):
|
||||
continue
|
||||
@@ -674,19 +673,13 @@ def _fix_blank_text_blocks_in_list(
|
||||
"""Drop blank text blocks; relocate any cache_control they carried onto the last surviving
|
||||
cacheable block; if nothing survives, substitute one placeholder block (carrying the relocated
|
||||
marker). Non-text blocks and order are untouched. Returns a new list; logs structure only."""
|
||||
kept: List[Any] = []
|
||||
relocated_cache_control = None
|
||||
for block_index, blk in enumerate(blocks):
|
||||
if _is_blank_text_block(blk):
|
||||
if _cache_control_of(blk) is not None:
|
||||
relocated_cache_control = blk["cache_control"]
|
||||
logger.warning(
|
||||
"Pre-call sanitizer: dropped blank text content block "
|
||||
"(message_index=%d role=%s location=%s block_index=%d block_type=text)",
|
||||
msg_index, role, location, block_index,
|
||||
)
|
||||
continue
|
||||
kept.append(blk)
|
||||
kept, relocated_cache_control, dropped = _split_blank_text_blocks(blocks)
|
||||
for block_index in dropped:
|
||||
logger.warning(
|
||||
"Pre-call sanitizer: dropped blank text content block "
|
||||
"(message_index=%d role=%s location=%s block_index=%d block_type=text)",
|
||||
msg_index, role, location, block_index,
|
||||
)
|
||||
if not kept:
|
||||
kept.append(_text_block(placeholder_text))
|
||||
_apply_assistant_cache_control_to_last_cacheable_block(kept, relocated_cache_control)
|
||||
|
||||
+32
-77
@@ -15,7 +15,7 @@ import json
|
||||
import logging
|
||||
import os
|
||||
import threading
|
||||
from contextlib import contextmanager
|
||||
from contextlib import contextmanager, suppress
|
||||
from dataclasses import dataclass, field
|
||||
from typing import Any, Dict, Iterator, List, Optional, Tuple
|
||||
|
||||
@@ -112,8 +112,7 @@ def _interrupt_background_review(review_agent: Any) -> None:
|
||||
from agent.interrupt_compat import request_hard_interrupt
|
||||
|
||||
request_hard_interrupt(
|
||||
review_agent, "superseded by a new live turn",
|
||||
tool_reason="background review superseded",
|
||||
review_agent, "superseded by a new live turn", tool_reason="background review superseded"
|
||||
)
|
||||
except Exception:
|
||||
logger.debug("Failed to cancel in-flight background review for a new turn", exc_info=True)
|
||||
@@ -255,13 +254,9 @@ def _resolve_review_runtime(agent: Any, task_cfg: Optional[Dict[str, Any]] = Non
|
||||
return {
|
||||
"provider": rp.get("provider") or task_provider,
|
||||
"model": rp.get("model") or task_model,
|
||||
"api_key": rp.get("api_key"),
|
||||
"base_url": rp.get("base_url"),
|
||||
"api_mode": rp.get("api_mode"),
|
||||
"credential_pool": rp.get("credential_pool"),
|
||||
**{key: rp.get(key) for key in ("api_key", "base_url", "api_mode", "credential_pool", "command")},
|
||||
"request_overrides": dict(rp.get("request_overrides") or {}),
|
||||
"max_tokens": rp.get("max_output_tokens"),
|
||||
"command": rp.get("command"),
|
||||
"args": list(rp.get("args") or []),
|
||||
"routed": True,
|
||||
}
|
||||
@@ -296,14 +291,13 @@ def _digest_history(messages_snapshot: List[Dict], tail: int = 24) -> List[Dict]
|
||||
messages verbatim (extended so the kept run never starts on a tool result) and collapse older
|
||||
turns into one synthetic user-role digest, preserving role alternation."""
|
||||
msgs = list(messages_snapshot or [])
|
||||
if len(msgs) <= tail:
|
||||
return msgs
|
||||
keep = msgs[-tail:]
|
||||
while keep and isinstance(keep[0], dict) and keep[0].get("role") == "tool":
|
||||
tail += 1
|
||||
if len(msgs) <= tail:
|
||||
return msgs
|
||||
while len(msgs) > tail:
|
||||
keep = msgs[-tail:]
|
||||
if not (isinstance(keep[0], dict) and keep[0].get("role") == "tool"):
|
||||
break
|
||||
tail += 1
|
||||
else:
|
||||
return msgs
|
||||
lines: List[str] = []
|
||||
for m in msgs[:-len(keep)]:
|
||||
if not isinstance(m, dict):
|
||||
@@ -626,13 +620,11 @@ def _verbose_skill_line(data: Dict, detail: Dict, message: str) -> str:
|
||||
new_string = change.get("new", "") or detail.get("new_string", "")
|
||||
description = change.get("description", "")
|
||||
if action == "patch" and (old_string or new_string):
|
||||
old_preview = _preview(old_string, 80).replace("\n", " ")
|
||||
new_preview = _preview(new_string, 80).replace("\n", " ")
|
||||
old_preview, new_preview = (_preview(t, 80).replace("\n", " ") for t in (old_string, new_string))
|
||||
return f"📝 Skill '{skill_name}' patched: \"{old_preview}\" → \"{new_preview}\""
|
||||
if action == "create" and description:
|
||||
return f"📝 Skill '{skill_name}' created: {description}"
|
||||
if action == "edit" and description:
|
||||
return f"📝 Skill '{skill_name}' rewritten: {description}"
|
||||
verb = {"create": "created", "edit": "rewritten"}.get(action)
|
||||
if verb and description:
|
||||
return f"📝 Skill '{skill_name}' {verb}: {description}"
|
||||
return f"📝 {message}" if message else f"Skill {action}"
|
||||
|
||||
|
||||
@@ -837,16 +829,9 @@ def _classify_review_result(actions: List[str]) -> str:
|
||||
``📝 Skill …``, ``Memory …``, ``User profile …``), so a free-text line like ``Skipped: no
|
||||
skill worth saving`` stays ``none``.
|
||||
"""
|
||||
has_skill = has_memory = False
|
||||
for action in actions or []:
|
||||
text = str(action).lstrip()
|
||||
if text.startswith("📝"):
|
||||
text = text[1:].lstrip()
|
||||
lower = text.lower()
|
||||
if lower.startswith("skill"):
|
||||
has_skill = True
|
||||
elif lower.startswith("memory") or lower.startswith("user profile"):
|
||||
has_memory = True
|
||||
lowers = [str(action).lstrip().removeprefix("📝").lstrip().lower() for action in actions or []]
|
||||
has_skill = any(t.startswith("skill") for t in lowers)
|
||||
has_memory = any(t.startswith(("memory", "user profile")) for t in lowers)
|
||||
return "+".join(kind for kind, hit in (("skill", has_skill), ("memory", has_memory)) if hit) or "none"
|
||||
|
||||
|
||||
@@ -986,15 +971,13 @@ def build_cache_parity_fork(
|
||||
_rt = _resolve_review_runtime(agent, task_cfg)
|
||||
_routed = bool(_rt.get("routed"))
|
||||
review_agent = AIAgent(**_fork_init_kwargs(agent, _rt, _routed, max_iterations))
|
||||
review_agent._memory_write_origin = write_origin
|
||||
review_agent._memory_write_context = write_origin
|
||||
review_agent._memory_write_origin = review_agent._memory_write_context = write_origin
|
||||
# The between-turns MCP refresh would add late-connecting MCP tools and break tools[] parity.
|
||||
review_agent._skip_mcp_refresh = True
|
||||
review_agent._memory_store = agent._memory_store
|
||||
review_agent._memory_enabled = agent._memory_enabled
|
||||
review_agent._user_profile_enabled = agent._user_profile_enabled
|
||||
review_agent._memory_nudge_interval = 0
|
||||
review_agent._skill_nudge_interval = 0
|
||||
review_agent._memory_nudge_interval = review_agent._skill_nudge_interval = 0
|
||||
# PERSISTENCE ISOLATION (curator-takeover root cause): sharing the parent's session_id, the
|
||||
# fork would otherwise write its harness turn into the REAL session, which the next live turn
|
||||
# re-reads as a standing instruction.
|
||||
@@ -1030,10 +1013,8 @@ def _bg_review_auto_deny(command, description, **kwargs):
|
||||
def _set_thread_approval_callback(callback: Any) -> None:
|
||||
from tools.terminal_tool import set_approval_callback
|
||||
|
||||
try:
|
||||
with suppress(Exception):
|
||||
set_approval_callback(callback)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def _track_review_fork(agent: Any, review_agent: Any, *, register: bool) -> None:
|
||||
@@ -1084,10 +1065,9 @@ def _review_tool_whitelist(review_agent: Any, task_cfg: Optional[Dict[str, Any]]
|
||||
_extra_raw = _background_review_task_config(task_cfg).get("extra_tools", [])
|
||||
if isinstance(_extra_raw, list):
|
||||
configured_extra_tools = {name.strip() for name in _extra_raw if isinstance(name, str) and name.strip()}
|
||||
whitelist |= configured_extra_tools
|
||||
except Exception:
|
||||
logger.debug("background_review extra_tools parse failed", exc_info=True)
|
||||
return whitelist, configured_extra_tools
|
||||
return whitelist | configured_extra_tools, configured_extra_tools
|
||||
|
||||
|
||||
@dataclass
|
||||
@@ -1102,10 +1082,8 @@ class _ReviewForkState:
|
||||
def _release_fork_clients(review_agent: Any) -> None:
|
||||
"""The fork shares the foreground session ID: close() / shutdown_memory_provider() are
|
||||
session-bound (close() kills that session's terminal processes), so release only clients."""
|
||||
try:
|
||||
with suppress(Exception):
|
||||
review_agent.release_clients()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def _finish_request_phase(agent: Any, review_agent: Any, review_run: Optional[_BackgroundReviewRun]) -> None:
|
||||
@@ -1128,6 +1106,8 @@ def _run_review_fork(
|
||||
|
||||
review_whitelist, configured_extra_tools = _review_tool_whitelist(st.review_agent, task_cfg)
|
||||
extra_list = ", ".join(sorted(configured_extra_tools))
|
||||
deny_extra = f" Configured extra tools also allowed: {extra_list}." if configured_extra_tools else ""
|
||||
prompt_extra = f" Exception — these configured tools are also allowed: {extra_list}." if configured_extra_tools else ""
|
||||
set_thread_tool_whitelist(
|
||||
review_whitelist,
|
||||
deny_msg_fmt=(
|
||||
@@ -1135,37 +1115,22 @@ def _run_review_fork(
|
||||
"{tool_name}. Allowed here: skill_view/skills_list/"
|
||||
"read_file/search_files to read, "
|
||||
"skill_manage(action='patch'|...) to change skills, and "
|
||||
"memory for notes."
|
||||
+ (
|
||||
" Configured extra tools also allowed: " + extra_list + "."
|
||||
if configured_extra_tools
|
||||
else ""
|
||||
)
|
||||
+ " Do not retry {tool_name}."
|
||||
"memory for notes." + deny_extra + " Do not retry {tool_name}."
|
||||
),
|
||||
)
|
||||
try:
|
||||
with suppress(Exception):
|
||||
from tools.skill_manager_tool import _reset_background_review_read_marks
|
||||
|
||||
_reset_background_review_read_marks()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
try:
|
||||
if review_run is None or review_run.begin_request(st.review_agent):
|
||||
# Routed -> digest (cache cold anyway); same model -> full snapshot (warm cache reads).
|
||||
st.review_agent.run_conversation(
|
||||
user_message=(
|
||||
prompt
|
||||
+ "\n\nYou can only call memory and skill "
|
||||
prompt + "\n\nYou can only call memory and skill "
|
||||
"management tools. Other tools will be denied "
|
||||
"at runtime — do not attempt them."
|
||||
+ (
|
||||
" Exception — these configured tools are "
|
||||
"also allowed: " + extra_list + "."
|
||||
if configured_extra_tools
|
||||
else ""
|
||||
)
|
||||
"at runtime — do not attempt them." + prompt_extra
|
||||
),
|
||||
conversation_history=_digest_history(messages_snapshot) if _routed else messages_snapshot,
|
||||
)
|
||||
@@ -1191,10 +1156,8 @@ def _publish_review_summary(agent: Any, actions: List[str]) -> None:
|
||||
agent._safe_print(f" 💾 Self-improvement review: {summary}")
|
||||
_bg_cb = agent.background_review_callback
|
||||
if _bg_cb:
|
||||
try:
|
||||
with suppress(Exception):
|
||||
_bg_cb(f"💾 Self-improvement review: {summary}")
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def _run_review_in_thread(
|
||||
@@ -1263,11 +1226,8 @@ def _run_review_in_thread(
|
||||
# cleanup output stays quiet without blanking other threads.
|
||||
_finish_request_phase(agent, st.review_agent, review_run)
|
||||
if st.review_agent is not None:
|
||||
try:
|
||||
with thread_scoped_silence():
|
||||
_release_fork_clients(st.review_agent)
|
||||
except Exception:
|
||||
pass
|
||||
with suppress(Exception), thread_scoped_silence():
|
||||
_release_fork_clients(st.review_agent)
|
||||
# Clear the approval callback so a recycled thread-id doesn't inherit it.
|
||||
_set_thread_approval_callback(None)
|
||||
|
||||
@@ -1310,11 +1270,6 @@ def spawn_background_review_thread(
|
||||
|
||||
|
||||
__all__ = [
|
||||
"_MEMORY_REVIEW_PROMPT",
|
||||
"_SKILL_REVIEW_PROMPT",
|
||||
"_COMBINED_REVIEW_PROMPT",
|
||||
"load_background_review_settings",
|
||||
"spawn_background_review_thread",
|
||||
"summarize_background_review_actions",
|
||||
"build_memory_write_metadata",
|
||||
"_MEMORY_REVIEW_PROMPT", "_SKILL_REVIEW_PROMPT", "_COMBINED_REVIEW_PROMPT", "load_background_review_settings",
|
||||
"spawn_background_review_thread", "summarize_background_review_actions", "build_memory_write_metadata",
|
||||
]
|
||||
|
||||
Reference in New Issue
Block a user