diff --git a/agent/anthropic_adapter.py b/agent/anthropic_adapter.py index 49f97f5f29..50ec70c268 100644 --- a/agent/anthropic_adapter.py +++ b/agent/anthropic_adapter.py @@ -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() diff --git a/agent/anthropic_message_convert.py b/agent/anthropic_message_convert.py index 80482c66d4..124574584b 100644 --- a/agent/anthropic_message_convert.py +++ b/agent/anthropic_message_convert.py @@ -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 diff --git a/agent/background_review.py b/agent/background_review.py index 73044a0656..3959ca0d1a 100644 --- a/agent/background_review.py +++ b/agent/background_review.py @@ -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}"