Merge origin/main into feat/hermes-relay-shared-metrics

Signed-off-by: Alex Fournier <afournier@nvidia.com>
This commit is contained in:
Alex Fournier 2026-07-27 08:12:37 -07:00
commit 4fe4b0dca7
192 changed files with 10055 additions and 594 deletions

View file

@ -31,6 +31,7 @@ from concurrent.futures import (
TimeoutError as FuturesTimeoutError,
)
from typing import Any, Dict, List, Optional
from urllib.parse import urlsplit, urlunsplit
from toolsets import TOOLSETS
@ -299,6 +300,131 @@ def _stringify_tool_content(content: Any) -> str:
return str(content)
_TOOL_INPUT_TARGET_KEYS = frozenset({
"cwd",
"destination_path",
"directory",
"dst",
"endpoint",
"file_path",
"new_path",
"old_path",
"path",
"source_path",
"src",
"target_path",
"url",
"urls",
})
_TOOL_INPUT_URL_KEYS = frozenset({"endpoint", "url", "urls"})
def _sanitize_tool_target(key: str, value: Any) -> Any:
"""Keep bounded side-effect targets while dropping URL secrets."""
if isinstance(value, list):
cleaned = [
item for item in (_sanitize_tool_target(key, item) for item in value[:16])
if item is not None
]
return cleaned or None
if not isinstance(value, str) or not value:
return None
bounded = value[:1024]
if key in _TOOL_INPUT_URL_KEYS:
try:
parsed = urlsplit(bounded)
if parsed.scheme and parsed.netloc:
hostname = parsed.hostname
if not hostname:
return None
# ``SplitResult.netloc`` includes ``user:password@``. Rebuild
# the authority from parsed host/port so hook-visible history
# cannot carry URL credentials. Bracket IPv6 literals before
# appending a validated port.
host = f"[{hostname}]" if ":" in hostname else hostname
port = parsed.port
netloc = f"{host}:{port}" if port is not None else host
return urlunsplit((parsed.scheme, netloc, parsed.path, "", ""))
except ValueError:
return None
return bounded
def _summarize_tool_arguments(arguments: Any) -> Dict[str, Any]:
"""Summarize argument names and side-effect targets without raw payloads."""
if not isinstance(arguments, str):
return {"argument_keys": [], "targets": {}}
try:
parsed = json.loads(arguments)
except (TypeError, ValueError):
return {"argument_keys": [], "targets": {}}
if not isinstance(parsed, dict):
return {"argument_keys": [], "targets": {}}
keys = sorted(str(key)[:128] for key in parsed)[:64]
targets: Dict[str, Any] = {}
for raw_key, value in parsed.items():
key = str(raw_key).lower()
if key not in _TOOL_INPUT_TARGET_KEYS:
continue
cleaned = _sanitize_tool_target(key, value)
if cleaned is not None:
targets[key] = cleaned
return {"argument_keys": keys, "targets": targets}
def _sanitize_tool_input_summary(summary: Any) -> Dict[str, Any]:
if not isinstance(summary, dict):
return {"argument_keys": [], "targets": {}}
keys = summary.get("argument_keys")
safe_keys = (
[str(key)[:128] for key in keys[:64]]
if isinstance(keys, list)
else []
)
targets = summary.get("targets")
safe_targets: Dict[str, Any] = {}
if isinstance(targets, dict):
for raw_key, value in targets.items():
key = str(raw_key).lower()
if key not in _TOOL_INPUT_TARGET_KEYS:
continue
cleaned = _sanitize_tool_target(key, value)
if cleaned is not None:
safe_targets[key] = cleaned
return {"argument_keys": safe_keys, "targets": safe_targets}
def _subagent_stop_tool_call_history(tool_trace: Any) -> List[Dict[str, Any]]:
"""Build a detached, metadata-only tool history for lifecycle hooks."""
if not isinstance(tool_trace, list):
return []
history: List[Dict[str, Any]] = []
for item in tool_trace:
if not isinstance(item, dict):
continue
tool_name = str(item.get("tool") or "unknown")[:256]
status = str(item.get("status") or "unknown").lower()
if status not in {"ok", "error"}:
status = "unknown"
def _byte_count(key: str) -> int:
value = item.get(key, 0)
if not isinstance(value, (int, float)) or isinstance(value, bool):
return 0
return max(0, int(value))
history.append({
"tool_name": tool_name,
"tool_input": _sanitize_tool_input_summary(item.get("input_summary")),
"input_bytes": _byte_count("args_bytes"),
"output_bytes": _byte_count("result_bytes"),
"status": status,
})
return history
def _looks_like_error_output(content: Any) -> bool:
"""Conservative stderr/error detector for tool-result previews.
@ -1498,6 +1624,7 @@ def _dump_subagent_timeout_diagnostic(
import datetime as _dt
import sys as _sys
import traceback as _traceback
import threading as _threading
hermes_home = get_hermes_home()
logs_dir = hermes_home / "logs"
@ -1604,6 +1731,37 @@ def _dump_subagent_timeout_diagnostic(
_w(" <worker thread already exited>")
_w("")
# All other live threads. The conversation worker's own stack often
# shows it parked waiting on a nested helper thread (interrupt worker,
# daemon-pool sibling) — without the full picture, a pre-HTTP wedge
# (#60203/#62151) is indistinguishable from a slow provider. Best
# effort and bounded: names + stacks for up to 40 threads.
_w("## All thread stacks at timeout")
try:
frames = _sys._current_frames()
by_ident = {
th.ident: th for th in _threading.enumerate() if th.ident
}
worker_ident = worker_thread.ident if worker_thread else None
dumped = 0
for ident, frame in frames.items():
if ident == worker_ident:
continue # already dumped above
if dumped >= 40:
_w(f" <{len(frames) - dumped - 1} more threads omitted>")
break
th = by_ident.get(ident)
name = th.name if th else f"ident={ident}"
daemon = " daemon" if (th and th.daemon) else ""
_w(f" --- {name}{daemon} ---")
for frame_line in _traceback.format_stack(frame):
for sub in frame_line.rstrip().split("\n"):
_w(f" {sub}")
dumped += 1
except Exception as exc:
_w(f" <all-thread dump failed: {exc}>")
_w("")
_w("## Notes")
_w(" This file is written ONLY when a subagent times out with 0 API calls.")
_w(" 0-API-call timeouts mean the child never reached its first LLM request.")
@ -2173,9 +2331,11 @@ def _run_single_child(
if msg.get("role") == "assistant":
for tc in msg.get("tool_calls") or []:
fn = tc.get("function", {})
arguments = fn.get("arguments", "")
entry_t = {
"tool": fn.get("name", "unknown"),
"args_bytes": len(fn.get("arguments", "")),
"args_bytes": len(arguments),
"input_summary": _summarize_tool_arguments(arguments),
}
tool_trace.append(entry_t)
tc_id = tc.get("id")
@ -2438,6 +2598,150 @@ def _run_single_child(
logger.debug("Failed to close child Relay session after delegation")
_PARENT_FINALIZATION_LOCK_GUARD = threading.Lock()
_PARENT_FINALIZATION_FALLBACK_LOCK = threading.RLock()
_CHILD_CONSTRUCTION_LOCK = threading.RLock()
def _build_child_preserving_parent_tools(**kwargs):
"""Build a child without leaking its resolved toolset into the parent."""
import model_tools
with _CHILD_CONSTRUCTION_LOCK:
parent_tool_names = list(model_tools._last_resolved_tool_names)
try:
child = _build_child_agent(**kwargs)
finally:
model_tools._last_resolved_tool_names = parent_tool_names
child._delegate_saved_tool_names = parent_tool_names
return child
def _parent_finalization_lock(parent_agent) -> threading.RLock:
"""Return the per-parent lock that serializes lifecycle side effects."""
if parent_agent is None:
return _PARENT_FINALIZATION_FALLBACK_LOCK
lock = getattr(parent_agent, "_subagent_finalization_lock", None)
if lock is not None:
return lock
with _PARENT_FINALIZATION_LOCK_GUARD:
lock = getattr(parent_agent, "_subagent_finalization_lock", None)
if lock is None:
lock = threading.RLock()
try:
setattr(parent_agent, "_subagent_finalization_lock", lock)
except Exception:
return _PARENT_FINALIZATION_FALLBACK_LOCK
return lock
def _finalize_child_results(
results: List[Dict[str, Any]],
task_list: List[Dict[str, Any]],
children: List[tuple[int, Dict[str, Any], Any]],
parent_agent,
) -> None:
"""Apply host-owned summary, memory, hook, and cost contracts once."""
with _parent_finalization_lock(parent_agent):
_apply_summary_budget(results, parent_agent)
child_by_index = {index: child for index, _task, child in children}
if parent_agent and getattr(parent_agent, "_memory_manager", None):
for entry in results:
try:
task_index = entry.get("task_index", -1)
task_goal = (
task_list[task_index]["goal"]
if isinstance(task_index, int)
and 0 <= task_index < len(task_list)
else ""
)
child = child_by_index.get(task_index)
parent_agent._memory_manager.on_delegation(
task=task_goal,
result=entry.get("summary", "") or "",
child_session_id=getattr(child, "session_id", ""),
)
except Exception:
pass
parent_session_id = getattr(parent_agent, "session_id", None)
try:
from hermes_cli.plugins import invoke_hook as invoke_hook
except Exception:
invoke_hook = None
children_cost_total = 0.0
for entry in results:
child_role = entry.pop("_child_role", None)
child_cost = entry.pop("_child_cost_usd", 0.0)
try:
if child_cost:
children_cost_total += float(child_cost)
except (TypeError, ValueError):
pass
if invoke_hook is None:
continue
try:
child_index = entry.get("task_index", -1)
child = child_by_index.get(child_index)
invoke_hook(
"subagent_stop",
parent_session_id=parent_session_id,
parent_turn_id=getattr(parent_agent, "_current_turn_id", "") or "",
child_session_id=getattr(child, "session_id", None),
child_role=child_role,
child_summary=entry.get("summary"),
child_status=entry.get("status"),
tool_call_history=_subagent_stop_tool_call_history(
entry.get("tool_trace")
),
duration_ms=int((entry.get("duration_seconds") or 0) * 1000),
)
except Exception:
logger.debug("subagent_stop hook invocation failed", exc_info=True)
if children_cost_total > 0.0:
try:
current = float(
getattr(parent_agent, "session_estimated_cost_usd", 0.0) or 0.0
)
parent_agent.session_estimated_cost_usd = current + children_cost_total
if getattr(parent_agent, "session_cost_source", "none") in {
None,
"",
"none",
}:
parent_agent.session_cost_source = "subagent"
if getattr(parent_agent, "session_cost_status", "unknown") in {
None,
"",
"unknown",
}:
parent_agent.session_cost_status = "estimated"
except Exception:
logger.debug("Subagent cost rollup failed", exc_info=True)
def _run_child_lifecycle(
task_index: int,
goal: str,
child=None,
parent_agent=None,
) -> Dict[str, Any]:
"""Run one child and apply the same host lifecycle used by delegate_task."""
result = _run_single_child(task_index, goal, child, parent_agent)
result.setdefault("task_index", task_index)
task = {"goal": goal}
_finalize_child_results(
[result],
[{"goal": ""} for _ in range(task_index)] + [task],
[(task_index, task, child)],
parent_agent,
)
return result
def _recover_tasks_from_json_string(
tasks: Any,
) -> tuple[Optional[List[Dict[str, Any]]], Optional[str]]:
@ -2608,13 +2912,6 @@ def delegate_task(
task_list, context
)
# Save parent tool names BEFORE any child construction mutates the global.
# _build_child_agent() calls AIAgent() which calls get_tool_definitions(),
# which overwrites model_tools._last_resolved_tool_names with child's toolset.
import model_tools as _model_tools
_parent_tool_names = list(_model_tools._last_resolved_tool_names)
# Capture the ORIGINATING session's wake target BEFORE any child agent is
# constructed: _build_child_agent() -> AIAgent() -> agent_init calls
# set_current_session_id(child.session_id), which clobbers the
@ -2627,53 +2924,49 @@ def delegate_task(
_origin_wake_sid = _current_origin_session_id()
# Build all child agents on the main thread (thread-safe construction)
# Wrapped in try/finally so the global is always restored even if a
# child build raises (otherwise _last_resolved_tool_names stays corrupted).
# Build all child agents on the main thread (thread-safe construction).
# _build_child_preserving_parent_tools saves/restores the parent's
# resolved tool names around each construction under a lock, so child
# toolset resolution never leaks into the parent (shared with the plugin
# subagent-lifecycle API).
children = []
try:
for i, t in enumerate(task_list):
# Per-task role beats top-level; normalise again so unknown
# per-task values warn and degrade to leaf uniformly.
effective_role = _normalize_role(t.get("role") or top_role)
child = _build_child_agent(
task_index=i,
goal=t["goal"],
context=t.get("context"),
# Subagents always inherit the parent's toolsets; the model
# cannot choose or narrow them (no model-facing toolsets arg).
toolsets=None,
model=creds["model"],
max_iterations=effective_max_iter,
task_count=n_tasks,
parent_agent=parent_agent,
override_provider=creds["provider"],
override_base_url=creds["base_url"],
override_api_key=creds["api_key"],
override_api_mode=creds["api_mode"],
override_request_overrides=creds.get("request_overrides"),
override_max_tokens=creds.get("max_output_tokens"),
override_acp_command=creds.get("command"),
override_acp_args=creds.get("args"),
role=effective_role,
for i, t in enumerate(task_list):
# Per-task role beats top-level; normalise again so unknown
# per-task values warn and degrade to leaf uniformly.
effective_role = _normalize_role(t.get("role") or top_role)
child = _build_child_preserving_parent_tools(
task_index=i,
goal=t["goal"],
context=t.get("context"),
# Subagents always inherit the parent's toolsets; the model
# cannot choose or narrow them (no model-facing toolsets arg).
toolsets=None,
model=creds["model"],
max_iterations=effective_max_iter,
task_count=n_tasks,
parent_agent=parent_agent,
override_provider=creds["provider"],
override_base_url=creds["base_url"],
override_api_key=creds["api_key"],
override_api_mode=creds["api_mode"],
override_request_overrides=creds.get("request_overrides"),
override_max_tokens=creds.get("max_output_tokens"),
override_acp_command=creds.get("command"),
override_acp_args=creds.get("args"),
role=effective_role,
)
# Tee the child's progress events into its live transcript log.
# wrap_progress_callback preserves the inner callback contract
# (including the _flush attribute) and never lets writer failures
# reach the agent loop. When no parent display exists the inner
# callback is None and the wrapper still records events.
_writer = live_writers[i] if i < len(live_writers) else None
if _writer is not None:
child.tool_progress_callback = wrap_progress_callback(
getattr(child, "tool_progress_callback", None), _writer
)
# Override with correct parent tool names (before child construction mutated global)
child._delegate_saved_tool_names = _parent_tool_names
# Tee the child's progress events into its live transcript log.
# wrap_progress_callback preserves the inner callback contract
# (including the _flush attribute) and never lets writer failures
# reach the agent loop. When no parent display exists the inner
# callback is None and the wrapper still records events.
_writer = live_writers[i] if i < len(live_writers) else None
if _writer is not None:
child.tool_progress_callback = wrap_progress_callback(
getattr(child, "tool_progress_callback", None), _writer
)
child._live_transcript_path = str(_writer.path)
children.append((i, t, child))
finally:
# Authoritative restore: reset global to parent's tool names after all children built
_model_tools._last_resolved_tool_names = _parent_tool_names
child._live_transcript_path = str(_writer.path)
children.append((i, t, child))
def _execute_and_aggregate() -> dict:
"""Run all built children (1 or N), join on them, aggregate results,
@ -2819,101 +3112,7 @@ def delegate_task(
# headroom (split across the batch) before they enter the parent's
# conversation. Full text is spilled to disk so nothing is lost.
# Covers both the single-task and batch paths. See PR #9126.
_apply_summary_budget(results, parent_agent)
# Notify parent's memory provider of delegation outcomes
if (
parent_agent
and hasattr(parent_agent, "_memory_manager")
and parent_agent._memory_manager
):
for entry in results:
try:
_task_goal = (
task_list[entry["task_index"]]["goal"]
if entry["task_index"] < len(task_list)
else ""
)
parent_agent._memory_manager.on_delegation(
task=_task_goal,
result=entry.get("summary", "") or "",
child_session_id=(
getattr(children[entry["task_index"]][2], "session_id", "")
if entry["task_index"] < len(children)
else ""
),
)
except Exception:
pass
# Fire subagent_stop hooks once per child, serialised on the parent thread.
# This keeps Python-plugin and shell-hook callbacks off of the worker threads
# that ran the children, so hook authors don't need to reason about
# concurrent invocation. Role was captured into the entry dict in
# _run_single_child (or the fabricated-entry branches above) before the
# child was closed.
_parent_session_id = getattr(parent_agent, "session_id", None)
try:
from hermes_cli.lifecycle import invoke_hook as _invoke_hook
except Exception:
_invoke_hook = None
# Aggregate child spend here so the parent's footer/UI reflect the true
# cost of a subagent-heavy turn. Port of Kilo-Org/kilocode#9448. Each
# child's cost was captured in _run_single_child before its AIAgent was
# closed; we fold them into the parent in one pass alongside the
# subagent_stop hook loop so we don't walk `results` twice.
_children_cost_total = 0.0
for entry in results:
child_role = entry.pop("_child_role", None)
child_cost = entry.pop("_child_cost_usd", 0.0)
try:
if child_cost:
_children_cost_total += float(child_cost)
except (TypeError, ValueError):
pass
if _invoke_hook is None:
continue
try:
_child_index = entry.get("task_index", -1)
_child_agent = (
children[_child_index][2]
if isinstance(_child_index, int) and 0 <= _child_index < len(children)
else None
)
_invoke_hook(
"subagent_stop",
parent_session_id=_parent_session_id,
parent_turn_id=getattr(parent_agent, "_current_turn_id", "") or "",
child_session_id=getattr(_child_agent, "session_id", None),
child_role=child_role,
child_summary=entry.get("summary"),
child_status=entry.get("status"),
duration_ms=int((entry.get("duration_seconds") or 0) * 1000),
)
except Exception:
logger.debug("subagent_stop hook invocation failed", exc_info=True)
# Fold the aggregated child cost into the parent's session total. This is
# additive — each delegate_task call contributes its own children — so
# nested orchestrator→worker trees roll up naturally: each layer's own
# delegate_task() folds its direct children in, and when the orchestrator
# itself finishes, its parent folds the orchestrator's now-inflated total
# on top. Degrades silently if the parent lacks the counter (older test
# fixtures, etc.).
if _children_cost_total > 0.0:
try:
current = float(getattr(parent_agent, "session_estimated_cost_usd", 0.0) or 0.0)
parent_agent.session_estimated_cost_usd = current + _children_cost_total
# Upgrade the cost_source so the UI doesn't label a partially-real
# total as "none" when the parent itself hadn't billed any calls
# yet (rare but possible when the parent's only action this turn
# was delegate_task).
if getattr(parent_agent, "session_cost_source", "none") in {None, "", "none"}:
parent_agent.session_cost_source = "subagent"
if getattr(parent_agent, "session_cost_status", "unknown") in {None, "", "unknown"}:
parent_agent.session_cost_status = "estimated"
except Exception:
logger.debug("Subagent cost rollup failed", exc_info=True)
_finalize_child_results(results, task_list, children, parent_agent)
total_duration = round(time.monotonic() - overall_start, 2)