refactor(agent/anthropic_*,background_review): final density pass — walrus guards, single-line SDK errors, screenshot placeholder hoist, run-token latch folds

This commit is contained in:
Teknium
2026-09-02 20:45:26 -07:00
parent 1bc8337bd0
commit 96fb75621b
3 changed files with 66 additions and 118 deletions

View File

@@ -63,10 +63,7 @@ def _require_sdk(purpose: str, verb: str = "Install it with"):
"""``_get_anthropic_sdk()`` or ImportError naming the feature that needs it."""
sdk = _get_anthropic_sdk()
if sdk is None:
raise ImportError(
f"The 'anthropic' package is required for {purpose}. "
f"{verb}: pip install 'anthropic>=0.39.0'"
)
raise ImportError(f"The 'anthropic' package is required for {purpose}. {verb}: pip install 'anthropic>=0.39.0'")
return sdk
@@ -333,9 +330,8 @@ def _base_client_kwargs(base_url, timeout) -> tuple[str, Dict[str, Any]]:
SDK appends ``/v1/messages``. Azure's ``api-version`` goes through ``default_query`` so the
base_url is not corrupted into ``/anthropic?api-version=.../v1/messages``."""
kwargs: Dict[str, Any] = {"timeout": _client_timeout(timeout), "max_retries": 0}
normalized = _normalize_base_url_text(base_url)
normalized = re.sub(r"/v1/?$", "", _normalize_base_url_text(base_url).rstrip("/"))
if normalized:
normalized = re.sub(r"/v1/?$", "", normalized.rstrip("/"))
kwargs["base_url"] = normalized
if _is_azure_anthropic_endpoint(normalized) and "api-version" not in normalized:
kwargs["default_query"] = {"api-version": "2025-04-15"}
@@ -431,13 +427,9 @@ def build_anthropic_bedrock_client(region: str):
``context-1m-2025-08-07`` are attached: without the latter Bedrock caps Opus 4.6/4.7 at 200K."""
sdk = _require_sdk("the Bedrock provider")
if not hasattr(sdk, "AnthropicBedrock"):
raise ImportError(
"anthropic.AnthropicBedrock not available. "
"Upgrade with: pip install 'anthropic>=0.39.0'"
)
raise ImportError("anthropic.AnthropicBedrock not available. Upgrade with: pip install 'anthropic>=0.39.0'")
return sdk.AnthropicBedrock(
aws_region=region,
timeout=_client_timeout(None),
aws_region=region, timeout=_client_timeout(None),
max_retries=0, # retry belongs to hermes's outer loop (honors Retry-After)
default_headers=_beta_header([*_COMMON_BETAS, _CONTEXT_1M_BETA]),
)
@@ -451,9 +443,7 @@ def _normalize_to_mcp_wire(name: str) -> str:
land on the double-underscore form. normalize_response reverses both via registry lookup."""
if name.startswith("mcp__"):
return name # already correct, don't double-prefix
if name.startswith("mcp_"):
return "mcp__" + name[len("mcp_"):]
return _MCP_TOOL_PREFIX + name
return _MCP_TOOL_PREFIX + name.removeprefix("mcp_")
def _oauth_wire_namer(anthropic_tools: List[Dict[str, Any]]):
@@ -500,15 +490,12 @@ def _apply_claude_code_identity(system, anthropic_tools, anthropic_messages, to_
for tool in anthropic_tools or []:
if "name" in tool:
tool["name"] = to_wire(tool["name"])
description = tool.get("description")
if isinstance(description, str):
tool["description"] = _apply_oauth_prose_aliases(description) # prose-safe aliases only
if isinstance(tool.get("description"), str):
tool["description"] = _apply_oauth_prose_aliases(tool["description"]) # prose-safe aliases only
for msg in anthropic_messages:
content = msg.get("content")
if isinstance(content, list):
for block in content:
if isinstance(block, dict) and block.get("type") == "tool_use" and "name" in block:
block["name"] = to_wire(block["name"]) # tool_result pairs by id, not name
for block in msg.get("content") if isinstance(msg.get("content"), list) else []:
if isinstance(block, dict) and block.get("type") == "tool_use" and "name" in block:
block["name"] = to_wire(block["name"]) # tool_result pairs by id, not name
return system
@@ -526,15 +513,12 @@ def _thinking_kwargs(reasoning_config: Dict[str, Any], model: str, effective_max
if "haiku" in model.lower():
return {}
effort = str(reasoning_config.get("effort", "medium")).lower()
budget = THINKING_BUDGET.get(effort, 8000)
if _supports_adaptive_thinking(model):
adaptive_effort = ADAPTIVE_EFFORT_MAP.get(effort, "medium")
if adaptive_effort == "xhigh" and not _supports_xhigh_effort(model):
adaptive_effort = "max"
return {
"thinking": {"type": "adaptive", "display": "summarized"},
"output_config": {"effort": adaptive_effort},
}
return {"thinking": {"type": "adaptive", "display": "summarized"}, "output_config": {"effort": adaptive_effort}}
budget = THINKING_BUDGET.get(effort, 8000)
return {
"thinking": {"type": "enabled", "budget_tokens": budget},
"temperature": 1, # required when thinking is enabled on older models
@@ -648,18 +632,17 @@ def _stream_final_message(stream_fn, api_kwargs, log_prefix, on_stream_event, on
on_response(getattr(stream, "response", None))
except Exception:
logger.debug("%son_response callback failed", log_prefix, exc_info=True)
if callable(on_stream_event):
# Consume manually so each event ticks the progress callback; get_final_message then
# returns the accumulated snapshot.
for _event in stream:
try:
on_stream_event(_event)
except TimeoutError:
# The callback is the caller's deadline seam: the host has given up, so abandon
# the stream (``with`` closes it) instead of streaming an answer nobody reads.
raise
except Exception:
logger.debug("%son_stream_event callback failed", log_prefix, exc_info=True)
# Consume manually so each event ticks the progress callback; get_final_message then
# returns the accumulated snapshot. TimeoutError is the caller's deadline seam: the host
# has given up, so abandon the stream (``with`` closes it) instead of streaming an answer
# nobody reads.
for event in stream if callable(on_stream_event) else ():
try:
on_stream_event(event)
except TimeoutError:
raise
except Exception:
logger.debug("%son_stream_event callback failed", log_prefix, exc_info=True)
return stream.get_final_message()

View File

@@ -1,11 +1,8 @@
"""OpenAI-style -> Anthropic Messages API request conversion.
Everything here rewrites *request payloads*: model-id normalization, tool schemas, and the
message list (content blocks, thinking blocks and their signatures, tool_use/tool_result
pairing, cache_control placement, screenshot eviction, blank-block scrubbing). Endpoint
predicates come from ``agent/anthropic_endpoints.py``, so this module never imports the adapter
and there is no cycle. ``agent.anthropic_adapter`` re-exports every name below.
"""
"""OpenAI-style -> Anthropic Messages API request conversion: model-id normalization, tool
schemas, and the message list (content blocks, thinking blocks and their signatures,
tool_use/tool_result pairing, cache_control placement, screenshot eviction, blank-block
scrubbing). Endpoint predicates come from ``agent/anthropic_endpoints.py`` so this module never
imports the adapter (no cycle); ``agent.anthropic_adapter`` re-exports the public names."""
import copy
import json
@@ -24,9 +21,7 @@ _THINKING_TYPES = frozenset(("thinking", "redacted_thinking"))
_CACHEABLE_TYPES = frozenset(("text", "tool_use"))
_EMPTY_TEXT_PLACEHOLDER = "(empty)"
_EMPTY_SCHEMA = {"type": "object", "properties": {}}
_BEDROCK_REGION_PREFIXES = (
"global.", "us.", "eu.", "apac.", "ap.", "au.", "jp.", "ca.", "sa.", "me.", "af.",
)
_BEDROCK_REGION_PREFIXES = ("global.", "us.", "eu.", "apac.", "ap.", "au.", "jp.", "ca.", "sa.", "me.", "af.")
def _block_type(b: Any) -> Any:
@@ -41,10 +36,7 @@ def _has_block_type(blocks: List[Any], types) -> bool:
def _is_blank_text_block(b: Any) -> bool:
"""A text block whose ``text`` is not a non-whitespace string (None/int/blank all count) —
Anthropic 400s on them ("text content blocks must contain non-whitespace text")."""
if _block_type(b) != "text":
return False
text = b.get("text")
return not (isinstance(text, str) and text.strip())
return _block_type(b) == "text" and not (isinstance(b.get("text"), str) and b["text"].strip())
def _cache_control_of(b: Any) -> Optional[Dict[str, Any]]:
@@ -97,17 +89,10 @@ def _carry_cache_control(out: Dict[str, Any], b: Any, *, copy: bool = False) ->
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
dropped = [i for i, blk in enumerate(blocks) if _is_blank_text_block(blk)]
kept = [blk for i, blk in enumerate(blocks) if i not in dropped]
relocated = [cc for i in dropped if (cc := _cache_control_of(blocks[i])) is not None]
return kept, relocated[-1] if relocated else None, dropped
def _is_bedrock_model_id(model: str) -> bool:
@@ -143,10 +128,9 @@ def _normalize_tool_input_schema(schema: Any) -> Dict[str, Any]:
optionality is already expressed by ``required``; ``keep_nullable_hint=False`` because the
OpenAPI ``nullable`` keyword is not recognized. Top-level oneOf/allOf/anyOf are rejected with a
generic 400, so they are dropped in favour of a plain object."""
if not schema:
return dict(_EMPTY_SCHEMA)
from tools.schema_sanitizer import strip_nullable_unions
normalized = strip_nullable_unions(schema, keep_nullable_hint=False)
normalized = strip_nullable_unions(schema, keep_nullable_hint=False) if schema else None
if not isinstance(normalized, dict):
return dict(_EMPTY_SCHEMA)
banned = {"oneOf", "allOf", "anyOf"}
@@ -207,8 +191,7 @@ def _convert_content_part_to_anthropic(part: Any) -> Optional[Dict[str, Any]]:
block = {"type": "image", "source": _image_source_from_openai_url(url)}
else:
block = dict(part)
cache_control = _cache_control_of(part)
if cache_control is not None:
if (cache_control := _cache_control_of(part)) is not None:
block.setdefault("cache_control", dict(cache_control))
return block
@@ -359,15 +342,12 @@ def _replay_ordered_blocks(m: Dict[str, Any], ordered_blocks: List[Any]) -> Opti
for b in ordered_blocks:
clean = _sanitize_replay_block(b)
if clean is None:
if _block_type(b) == "text":
dropped_blank_text = True
if _cache_control_of(b) is not None: # relocate a dropped block's breakpoint
relocated_cc = b["cache_control"]
dropped_blank_text = dropped_blank_text or _block_type(b) == "text"
if (cc := _cache_control_of(b)) is not None: # relocate a dropped block's breakpoint
relocated_cc = cc
continue
if clean.get("type") == "tool_use":
redacted = redacted_input_by_id.get(clean.get("id", ""))
if redacted is not None:
clean["input"] = redacted
if clean.get("type") == "tool_use" and (redacted := redacted_input_by_id.get(clean.get("id", ""))) is not None:
clean["input"] = redacted
replayed.append(clean)
# Nothing cacheable survived (e.g. signed thinking + blank text): emit the placeholder so the
# turn stays schema-valid and a relocated marker has a carrier.
@@ -610,9 +590,8 @@ def _evict_old_screenshots(result: List[Dict[str, Any]]) -> None:
continue
image_count += 1
if image_count > 3:
block["content"] = [
b if b.get("type") != "image" else _text_block("[screenshot removed to save context]") for b in inner
]
placeholder = _text_block("[screenshot removed to save context]")
block["content"] = [placeholder if b.get("type") == "image" else b for b in inner]
def _ensure_leading_user_turn(result: List[Dict[str, Any]]) -> None:
@@ -637,8 +616,7 @@ def _fix_blank_text_blocks_in_list(
"(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))
kept = kept or [_text_block(placeholder_text)]
_apply_assistant_cache_control_to_last_cacheable_block(kept, relocated_cache_control)
return kept

View File

@@ -31,8 +31,7 @@ class _BackgroundReviewRun:
self.request_done = threading.Event()
self._lock = threading.Lock()
self._review_agent = None
self._request_finished = False
self._cancel_dispatched = False
self._request_finished = self._cancel_dispatched = False
def begin_request(self, review_agent: Any) -> bool:
"""Atomically admit the first provider-capable review phase."""
@@ -46,18 +45,17 @@ class _BackgroundReviewRun:
"""Fence startup and return the running fork, if one was admitted."""
with self._lock:
self.cancel_requested.set()
if self._review_agent is not None and not self._cancel_dispatched:
self._cancel_dispatched = True
return self._review_agent
return None
if self._review_agent is None or self._cancel_dispatched:
return None
self._cancel_dispatched = True
return self._review_agent
def mark_request_finished(self) -> bool:
"""Latch request completion once; the caller publishes the event."""
with self._lock:
if self._request_finished:
return False
self._request_finished = True
self._review_agent = None
self._request_finished, self._review_agent = True, None
return True
@@ -251,11 +249,9 @@ def _parent_can_emit_tool_calls(agent: Any) -> bool:
def _msg_text(m: Dict) -> str:
c = m.get("content")
if isinstance(c, str):
return c.strip()
if isinstance(c, list):
return " ".join(b.get("text", "") for b in c if isinstance(b, dict)).strip()
return ""
c = " ".join(b.get("text", "") for b in c if isinstance(b, dict))
return c.strip() if isinstance(c, str) else ""
def _digest_history(messages_snapshot: List[Dict], tail: int = 24) -> List[Dict]:
@@ -274,8 +270,7 @@ def _digest_history(messages_snapshot: List[Dict], tail: int = 24) -> List[Dict]
for m in msgs[:-len(keep)]:
if not isinstance(m, dict):
continue
role = m.get("role")
text = _msg_text(m).replace("\n", " ")
role, text = m.get("role"), _msg_text(m).replace("\n", " ")
if role == "user" and text:
lines.append(f"USER: {text[:300]}")
elif role == "assistant":
@@ -507,13 +502,12 @@ def _verbose_skill_line(data: Dict, detail: Dict, message: str) -> str:
change: dict = change_raw if isinstance(change_raw, dict) else {}
old_string = change.get("old", "") or detail.get("old_string", "")
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, new_preview = (_preview(t, 80).replace("\n", " ") for t in (old_string, new_string))
return f"📝 Skill '{skill_name}' patched: \"{old_preview}\" → \"{new_preview}\""
verb = {"create": "created", "edit": "rewritten"}.get(action)
if verb and description:
return f"📝 Skill '{skill_name}' {verb}: {description}"
if verb and change.get("description"):
return f"📝 Skill '{skill_name}' {verb}: {change['description']}"
return f"📝 {message}" if message else f"Skill {action}"
@@ -581,24 +575,16 @@ def _action_lines(data: Dict, detail: Dict, verbose: bool) -> List[str]:
message = data.get("message", "")
target = data.get("target", "") or detail.get("target", "")
is_skill = detail.get("tool") == "skill_manage"
message_lower = message.lower()
if not verbose and (
"created" in message_lower or "updated" in message_lower or (is_skill and "patched" in message_lower)
):
lower = message.lower()
if not verbose and ("created" in lower or "updated" in lower or (is_skill and "patched" in lower)):
return [message]
if is_skill:
label = "Skill"
elif target:
label = "Memory" if target == "memory" else "User profile" if target == "user" else target
else:
if not is_skill and not target:
return []
label = "Skill" if is_skill else {"memory": "Memory", "user": "User profile"}.get(target, target)
if verbose:
return [_verbose_skill_line(data, detail, message)] if is_skill else _verbose_memory_lines(label, detail)
if any(k in message_lower for k in ("added", "replaced", "removed", "applied")) or (
target and "add" in message_lower
):
return [f"{label} updated"]
return []
hit = any(k in lower for k in ("added", "replaced", "removed", "applied")) or (target and "add" in lower)
return [f"{label} updated"] if hit else []
def summarize_background_review_actions(
@@ -847,6 +833,7 @@ def _bg_review_auto_deny(command, description, **kwargs):
def _set_thread_approval_callback(callback: Any) -> None:
from tools.terminal_tool import set_approval_callback
with suppress(Exception):
set_approval_callback(callback)
@@ -940,6 +927,7 @@ def _run_review_fork(
)
with suppress(Exception):
from tools.skill_manager_tool import _reset_background_review_read_marks
_reset_background_review_read_marks()
try:
if review_run is None or review_run.begin_request(st.review_agent):
@@ -992,7 +980,7 @@ def _run_review_in_thread(
# A client that can't carry Hermes tool calls back would spawn a fork that cannot write
# anything. Checked BEFORE the thread-scoped silence so the warning is not swallowed; cheap
# check first so the normal path never resolves the runtime twice.
if not _parent_can_emit_tool_calls(agent) and not bool(_resolve_review_runtime(agent, task_cfg).get("routed")):
if not _parent_can_emit_tool_calls(agent) and not _resolve_review_runtime(agent, task_cfg).get("routed"):
logger.warning(
"Background review skipped: provider %r cannot emit Hermes tool calls, "
"so the review fork could not write memories or skills. Set "
@@ -1064,8 +1052,7 @@ def spawn_background_review_thread(
# Per-agent overrides (agent._MEMORY_REVIEW_PROMPT etc.) keep working.
name = _PROMPT_NAME_BY_SCOPE[(review_memory, review_skills)]
prompt = getattr(agent, name, globals()[name])
focus = (focus or "").strip()
if focus:
if focus := (focus or "").strip():
prompt = (
f"{prompt}\n\nThe user explicitly requested this review with the following "
f"focus — prioritize it over the general instructions above:\n{focus}"