refactor(agent/background_review,anthropic_*): inline request-phase finisher, walrus pin/prefill folds, drop intra-function padding
This commit is contained in:
@@ -322,7 +322,6 @@ def _attribution_headers() -> Dict[str, str]:
|
||||
def _client_timeout(timeout):
|
||||
"""httpx.Timeout with the caller's read timeout (default 900s) and a 10s connect."""
|
||||
from httpx import Timeout
|
||||
|
||||
read = timeout if (isinstance(timeout, (int, float)) and timeout > 0) else 900.0
|
||||
return Timeout(timeout=float(read), connect=10.0)
|
||||
|
||||
@@ -353,9 +352,7 @@ def _build_anthropic_client_with_bearer_hook(
|
||||
placeholder ``auth_token`` is still required at construction and makes any leak diagnosable."""
|
||||
sdk = _require_sdk("Azure Foundry Anthropic-style endpoints with Entra ID auth", verb="Install with")
|
||||
normalize_proxy_env_vars()
|
||||
|
||||
from agent.azure_identity_adapter import build_bearer_http_client
|
||||
|
||||
normalized_base_url, kwargs = _base_client_kwargs(base_url, timeout)
|
||||
kwargs["http_client"] = build_bearer_http_client(token_provider, timeout=kwargs["timeout"])
|
||||
kwargs["auth_token"] = "entra-id-bearer-via-http-hook"
|
||||
@@ -406,13 +403,11 @@ def build_anthropic_client(api_key, base_url: str = None, timeout: float = None,
|
||||
return _build_anthropic_client_with_bearer_hook(
|
||||
api_key, base_url, timeout, drop_context_1m_beta=drop_context_1m_beta
|
||||
)
|
||||
|
||||
normalize_proxy_env_vars()
|
||||
normalized_base_url, kwargs = _base_client_kwargs(base_url, timeout)
|
||||
if "default_query" in kwargs: # historical: this path also strips a stray trailing slash on Azure
|
||||
kwargs["base_url"] = normalized_base_url.rstrip("/")
|
||||
common_betas = _common_betas_for_base_url(normalized_base_url, drop_context_1m_beta=drop_context_1m_beta)
|
||||
|
||||
style = _auth_style(api_key, base_url, normalized_base_url)
|
||||
kwargs["auth_token" if style in ("bearer", "oauth") else "api_key"] = api_key
|
||||
headers = _beta_header(common_betas + _OAUTH_ONLY_BETAS if style == "oauth" else common_betas)
|
||||
@@ -421,7 +416,6 @@ def build_anthropic_client(api_key, base_url: str = None, timeout: float = None,
|
||||
elif style == "oauth":
|
||||
headers["user-agent"] = f"claude-code/{_get_claude_code_version()} (external, cli)"
|
||||
headers["x-app"] = "cli"
|
||||
|
||||
if _is_opencode_endpoint(base_url):
|
||||
# OpenCode identifies clients by request headers (like OpenRouter). The OpenAI-wire paths
|
||||
# get these from profile.default_headers, but this route never sees the profile.
|
||||
@@ -503,14 +497,12 @@ def _apply_claude_code_identity(system, anthropic_tools, anthropic_messages, to_
|
||||
for old, new in _OAUTH_SYSTEM_REPLACEMENTS:
|
||||
text = text.replace(old, new)
|
||||
block["text"] = _apply_oauth_prose_aliases(text)
|
||||
|
||||
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
|
||||
|
||||
for msg in anthropic_messages:
|
||||
content = msg.get("content")
|
||||
if isinstance(content, list):
|
||||
@@ -645,7 +637,6 @@ def _is_stream_unavailable_error(exc: Exception) -> bool:
|
||||
if "invokemodelwithresponsestream" not in err_lower:
|
||||
return False
|
||||
from agent.bedrock_adapter import is_streaming_access_denied_error
|
||||
|
||||
return is_streaming_access_denied_error(exc)
|
||||
|
||||
|
||||
|
||||
@@ -89,7 +89,6 @@ def _is_nous_portal_endpoint(base_url: str | None) -> bool:
|
||||
return True
|
||||
try:
|
||||
from hermes_cli.auth import _nous_inference_env_override
|
||||
|
||||
override = _nous_inference_env_override()
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
@@ -145,9 +145,7 @@ def _normalize_tool_input_schema(schema: Any) -> Dict[str, Any]:
|
||||
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)
|
||||
if not isinstance(normalized, dict):
|
||||
return dict(_EMPTY_SCHEMA)
|
||||
@@ -398,7 +396,6 @@ def _convert_assistant_message(m: Dict[str, Any]) -> Dict[str, Any]:
|
||||
replayed = _replay_ordered_blocks(m, ordered_blocks)
|
||||
if replayed:
|
||||
return {"role": "assistant", "content": replayed}
|
||||
|
||||
blocks = _extract_preserved_thinking_blocks(m)
|
||||
# Blank text blocks are dropped; a cache marker riding on one is relocated onto the last
|
||||
# surviving cacheable block (prompt_caching sets cache_control on content[-1], which may be
|
||||
@@ -505,7 +502,6 @@ def _strip_orphaned_tool_blocks(result: List[Dict[str, Any]]) -> None:
|
||||
if len(kept) != len(m["content"]) and _has_block_type(m["content"], _THINKING_TYPES):
|
||||
m["_thinking_signature_invalidated"] = True
|
||||
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: set = set()
|
||||
for _, m in _assistant_block_lists(result):
|
||||
@@ -580,7 +576,6 @@ def _manage_thinking_signatures(result: List[Dict[str, Any]], base_url: str | No
|
||||
is_kimi = _is_kimi_family_endpoint(base_url, model)
|
||||
is_deepseek = _is_deepseek_anthropic_endpoint(base_url)
|
||||
last_assistant_idx = next((i for i in range(len(result) - 1, -1, -1) if result[i].get("role") == "assistant"), None)
|
||||
|
||||
for idx, m in _assistant_block_lists(result):
|
||||
if is_kimi:
|
||||
pass # shared cleanup below still strips cache markers + the flag
|
||||
@@ -596,7 +591,6 @@ def _manage_thinking_signatures(result: List[Dict[str, Any]], base_url: str | No
|
||||
else:
|
||||
new_content = _keep_valid_latest_thinking(m["content"], bool(m.get("_thinking_signature_invalidated")))
|
||||
m["content"] = new_content or [_text_block("(empty)")]
|
||||
|
||||
# cache_control on thinking blocks interferes with signature validation.
|
||||
for b in m["content"]:
|
||||
if _block_type(b) in _THINKING_TYPES:
|
||||
|
||||
+17
-48
@@ -102,7 +102,6 @@ def finish_background_review_run(agent: Any, run: Optional[_BackgroundReviewRun]
|
||||
def _interrupt_background_review(review_agent: Any) -> None:
|
||||
"""Request abort off-thread so a wedged abort hook cannot stall the live turn (the bounded
|
||||
``request_done`` wait in the canceller relies on this returning fast)."""
|
||||
|
||||
def _interrupt() -> None:
|
||||
try:
|
||||
from agent.interrupt_compat import request_hard_interrupt
|
||||
@@ -165,7 +164,6 @@ def _background_review_task_config(task_cfg: Optional[Dict[str, Any]] = None) ->
|
||||
return task_cfg if isinstance(task_cfg, dict) else {}
|
||||
try:
|
||||
from hermes_cli.config import load_config_readonly
|
||||
|
||||
return _task_block(load_config_readonly())
|
||||
except Exception:
|
||||
return {}
|
||||
@@ -187,7 +185,6 @@ def load_background_review_settings() -> tuple[bool, Dict[str, Any]]:
|
||||
try:
|
||||
from hermes_cli.config import load_config_readonly
|
||||
from utils import is_truthy_value
|
||||
|
||||
task = _task_block(load_config_readonly())
|
||||
return is_truthy_value(task.get("enabled"), default=True), task
|
||||
except Exception:
|
||||
@@ -618,7 +615,6 @@ def summarize_background_review_actions(
|
||||
verbose = mode == "verbose"
|
||||
existing_tool_call_ids, existing_tool_contents = _prior_tool_keys(prior_snapshot)
|
||||
all_tool_call_ids, call_details = _collect_review_call_details(review_messages)
|
||||
|
||||
actions: List[str] = []
|
||||
for msg in _tool_messages(review_messages):
|
||||
tcid = msg.get("tool_call_id")
|
||||
@@ -679,11 +675,9 @@ def _record_review_usage_to_parent(parent_agent: Any, usage: Dict[str, Any]) ->
|
||||
try:
|
||||
session_db = getattr(parent_agent, "_session_db", None)
|
||||
session_id = getattr(parent_agent, "session_id", None)
|
||||
if session_db is None or not session_id:
|
||||
return
|
||||
counts = {key: int(usage.get(key) or 0) for key in _USAGE_COUNTERS}
|
||||
if not any(counts.values()):
|
||||
return # fork made no successful API calls (e.g. failed at spawn)
|
||||
if session_db is None or not session_id or not any(counts.values()):
|
||||
return # no DB, or the fork made no successful API calls (e.g. failed at spawn)
|
||||
session_db.record_auxiliary_usage(
|
||||
session_id, task="background_review", model=usage.get("model"),
|
||||
billing_provider=usage.get("provider"), billing_base_url=usage.get("base_url"),
|
||||
@@ -734,17 +728,13 @@ def _same_model_parity_kwargs(agent: Any) -> Dict[str, Any]:
|
||||
# is appended to the cached system prompt at API-call time (without it the prompt diverges).
|
||||
"reasoning_config": getattr(agent, "reasoning_config", None),
|
||||
"ephemeral_system_prompt": getattr(agent, "ephemeral_system_prompt", None),
|
||||
**{attr: val for attr in _PROVIDER_PIN_ATTRS if (val := getattr(agent, attr, None))},
|
||||
}
|
||||
# Prefill sits right after the system message, so a parent with prefill would diverge at
|
||||
# index 1. Deep copy: unicode-error recovery sanitizes prefill entries IN PLACE and must not
|
||||
# rewrite the parent's bytes.
|
||||
parent_prefill = copy.deepcopy(getattr(agent, "prefill_messages", None) or [])
|
||||
if parent_prefill:
|
||||
if parent_prefill := copy.deepcopy(getattr(agent, "prefill_messages", None) or []):
|
||||
kwargs["prefill_messages"] = parent_prefill
|
||||
for attr in _PROVIDER_PIN_ATTRS:
|
||||
val = getattr(agent, attr, None)
|
||||
if val:
|
||||
kwargs[attr] = val
|
||||
return kwargs
|
||||
|
||||
|
||||
@@ -797,8 +787,7 @@ def _fork_init_kwargs(agent: Any, rt: Dict[str, Any], routed: bool, max_iteratio
|
||||
if isinstance(rt.get("max_tokens"), int):
|
||||
kwargs["max_tokens"] = rt["max_tokens"]
|
||||
if isinstance(rt.get("command"), str) and rt["command"]:
|
||||
kwargs["acp_command"] = rt["command"]
|
||||
kwargs["acp_args"] = rt.get("args") or []
|
||||
kwargs.update(acp_command=rt["command"], acp_args=rt.get("args") or [])
|
||||
if not routed:
|
||||
kwargs.update(_same_model_parity_kwargs(agent))
|
||||
return kwargs
|
||||
@@ -816,7 +805,6 @@ def build_cache_parity_fork(
|
||||
replay a digest). The caller owns registration, whitelisting, running, usage attribution and
|
||||
teardown."""
|
||||
from run_agent import AIAgent # local: avoids a circular import at load
|
||||
|
||||
# Inherit the parent's live runtime: AIAgent.__init__'s env auto-resolution fails for
|
||||
# OAuth-only providers, session-scoped creds and credential pools.
|
||||
_rt = _resolve_review_runtime(agent, task_cfg)
|
||||
@@ -859,7 +847,6 @@ 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)
|
||||
|
||||
@@ -891,12 +878,10 @@ def _review_tool_whitelist(review_agent: Any, task_cfg: Optional[Dict[str, Any]]
|
||||
"""``(whitelist, configured_extra_tools)`` for the review fork — DISPATCH-side only, so the
|
||||
advertised ``tools[]`` stays byte-identical to the parent's (prompt-cache parity)."""
|
||||
from model_tools import get_tool_definitions
|
||||
|
||||
# Gate the built-in memory tool on the profile's memory flags so a memory-disabled profile
|
||||
# is never contaminated by the review LLM.
|
||||
review_toolsets = ["skills"]
|
||||
if review_agent._memory_enabled or review_agent._user_profile_enabled:
|
||||
review_toolsets.insert(0, "memory")
|
||||
memory_on = review_agent._memory_enabled or review_agent._user_profile_enabled
|
||||
review_toolsets = ["memory", "skills"] if memory_on else ["skills"]
|
||||
whitelist = {t["function"]["name"] for t in get_tool_definitions(enabled_toolsets=review_toolsets, quiet_mode=True)}
|
||||
# Read-only file tools: denying read_file/search_files caused a per-review denial storm that
|
||||
# starved the loop (read_file also registers the read with the read-before-write guard).
|
||||
@@ -906,9 +891,9 @@ def _review_tool_whitelist(review_agent: Any, task_cfg: Optional[Dict[str, Any]]
|
||||
# can only admit, never advertise: a listed tool must already exist in the inherited schema.
|
||||
configured_extra_tools: set = set()
|
||||
try:
|
||||
_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()}
|
||||
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()}
|
||||
except Exception:
|
||||
logger.debug("background_review extra_tools parse failed", exc_info=True)
|
||||
return whitelist | configured_extra_tools, configured_extra_tools
|
||||
@@ -930,12 +915,6 @@ def _release_fork_clients(review_agent: Any) -> None:
|
||||
review_agent.release_clients()
|
||||
|
||||
|
||||
def _finish_request_phase(agent: Any, review_agent: Any, review_run: Optional[_BackgroundReviewRun]) -> None:
|
||||
"""Unregister the fork and publish request completion (identity-scoped, idempotent)."""
|
||||
_track_review_fork(agent, review_agent, register=False)
|
||||
finish_background_review_run(agent, review_run)
|
||||
|
||||
|
||||
def _run_review_fork(
|
||||
agent: Any, messages_snapshot: List[Dict], prompt: str, task_cfg: Optional[Dict[str, Any]],
|
||||
review_run: Optional[_BackgroundReviewRun], st: _ReviewForkState,
|
||||
@@ -945,9 +924,7 @@ def _run_review_fork(
|
||||
so the caller's error path still sees usage and the fork to clean up."""
|
||||
st.review_agent, _rt, _routed = build_cache_parity_fork(agent, task_cfg, max_iterations=_REVIEW_MAX_ITERATIONS)
|
||||
_track_review_fork(agent, st.review_agent, register=True)
|
||||
|
||||
from hermes_cli.plugins import set_thread_tool_whitelist, clear_thread_tool_whitelist
|
||||
|
||||
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 ""
|
||||
@@ -963,9 +940,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):
|
||||
# Routed -> digest (cache cold anyway); same model -> full snapshot (warm cache reads).
|
||||
@@ -986,9 +961,9 @@ def _run_review_fork(
|
||||
st.review_usage.update(_snapshot_review_usage(st.review_agent))
|
||||
_record_review_usage_to_parent(agent, st.review_usage)
|
||||
# Publish completion as soon as the provider-capable phase has returned or startup
|
||||
# cancellation has fenced it out.
|
||||
_finish_request_phase(agent, st.review_agent, review_run)
|
||||
|
||||
# cancellation has fenced it out (unregister + finish are identity-scoped and idempotent).
|
||||
_track_review_fork(agent, st.review_agent, register=False)
|
||||
finish_background_review_run(agent, review_run)
|
||||
st.review_messages = list(getattr(st.review_agent, "_session_messages", []))
|
||||
_release_fork_clients(st.review_agent)
|
||||
st.review_agent = None
|
||||
@@ -997,10 +972,9 @@ def _run_review_fork(
|
||||
def _publish_review_summary(agent: Any, actions: List[str]) -> None:
|
||||
summary = " · ".join(dict.fromkeys(actions))
|
||||
agent._safe_print(f" 💾 Self-improvement review: {summary}")
|
||||
_bg_cb = agent.background_review_callback
|
||||
if _bg_cb:
|
||||
if agent.background_review_callback:
|
||||
with suppress(Exception):
|
||||
_bg_cb(f"💾 Self-improvement review: {summary}")
|
||||
agent.background_review_callback(f"💾 Self-improvement review: {summary}")
|
||||
|
||||
|
||||
def _run_review_in_thread(
|
||||
@@ -1014,9 +988,7 @@ def _run_review_in_thread(
|
||||
if review_run is not None and review_run.cancel_requested.is_set():
|
||||
finish_background_review_run(agent, review_run)
|
||||
return
|
||||
|
||||
_set_thread_approval_callback(_bg_review_auto_deny)
|
||||
|
||||
# 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.
|
||||
@@ -1029,14 +1001,12 @@ def _run_review_in_thread(
|
||||
)
|
||||
_set_thread_approval_callback(None)
|
||||
return
|
||||
|
||||
st = _ReviewForkState()
|
||||
try:
|
||||
# Silence stdout/stderr for THIS thread only: a process-global redirect would blank every
|
||||
# other thread's console for the whole review.
|
||||
with thread_scoped_silence():
|
||||
_run_review_fork(agent, messages_snapshot, prompt, task_cfg, review_run, st)
|
||||
|
||||
# A buggy/legacy tool response shape must NOT take down the whole review (the outer
|
||||
# except would discard every action the fork DID complete), so coerce to an empty list.
|
||||
try:
|
||||
@@ -1052,11 +1022,9 @@ def _run_review_in_thread(
|
||||
e,
|
||||
)
|
||||
actions = []
|
||||
|
||||
_log_review_completion(st.review_usage, _classify_review_result(actions))
|
||||
if actions:
|
||||
_publish_review_summary(agent, actions)
|
||||
|
||||
except Exception as e:
|
||||
logger.warning("Background memory/skill review failed: %s", e)
|
||||
if st.review_usage:
|
||||
@@ -1066,7 +1034,8 @@ def _run_review_in_thread(
|
||||
# Safety net for the exception path (setup failures before the request-phase finally).
|
||||
# Both cleanups are identity-scoped and idempotent; re-enter thread-scoped silence so
|
||||
# cleanup output stays quiet without blanking other threads.
|
||||
_finish_request_phase(agent, st.review_agent, review_run)
|
||||
_track_review_fork(agent, st.review_agent, register=False)
|
||||
finish_background_review_run(agent, review_run)
|
||||
if st.review_agent is not None:
|
||||
with suppress(Exception), thread_scoped_silence():
|
||||
_release_fork_clients(st.review_agent)
|
||||
|
||||
Reference in New Issue
Block a user