diff --git a/MagicMock/mock._session_db.db_path/132161382957648 b/MagicMock/mock._session_db.db_path/132161382957648 new file mode 100644 index 0000000000..ccf54b739f Binary files /dev/null and b/MagicMock/mock._session_db.db_path/132161382957648 differ diff --git a/MagicMock/mock._session_db.db_path/132161382957648.fts_rebuild.lock b/MagicMock/mock._session_db.db_path/132161382957648.fts_rebuild.lock new file mode 100644 index 0000000000..e69de29bb2 diff --git a/MagicMock/mock._session_db.db_path/132161382957648.quarantine.lock b/MagicMock/mock._session_db.db_path/132161382957648.quarantine.lock new file mode 100644 index 0000000000..e69de29bb2 diff --git a/MagicMock/mock._session_db.db_path/132161402292368 b/MagicMock/mock._session_db.db_path/132161402292368 new file mode 100644 index 0000000000..fb203e1636 Binary files /dev/null and b/MagicMock/mock._session_db.db_path/132161402292368 differ diff --git a/MagicMock/mock._session_db.db_path/132161402292368.fts_rebuild.lock b/MagicMock/mock._session_db.db_path/132161402292368.fts_rebuild.lock new file mode 100644 index 0000000000..e69de29bb2 diff --git a/MagicMock/mock._session_db.db_path/132161402292368.quarantine.lock b/MagicMock/mock._session_db.db_path/132161402292368.quarantine.lock new file mode 100644 index 0000000000..e69de29bb2 diff --git a/tools/delegate_tool_child_run.py b/tools/delegate_tool_child_run.py index a95c2b10c4..c2293846eb 100644 --- a/tools/delegate_tool_child_run.py +++ b/tools/delegate_tool_child_run.py @@ -404,36 +404,15 @@ class _WorktreeReporter: info = self.info if info is None: return - try: - from tools import subagent_worktree + from tools import subagent_worktree + try: entry_dict["worktree"] = subagent_worktree.finalize_subagent_worktree(info) except Exception as e: # State is unknown: emit the SAME flagged schema the parent expects, # via the shared factory so the two producers never drift. logger.warning("worktree finalize failed: %s", e) - try: - from tools import subagent_worktree as _sw - - entry_dict["worktree"] = _sw.unproven_worktree_payload(info, f"finalize raised: {e}") - except Exception: - # Import itself failed — inline the same shape rather than - # dropping the flag (the parent must still see the warning). - entry_dict["worktree"] = { - "path": info.get("path", ""), - "branch": info.get("branch", ""), - "commits": 0, - "dirty": False, - "pruned": False, - "inspection_failed": True, - "note": ( - f"worktree finalize raised ({e}) and the reporting " - "helper was unavailable: 'commits' and 'dirty' are " - "UNKNOWN, not zero/clean. Inspect " - f"{info.get('path', '')} before assuming " - "no work." - ), - } + entry_dict["worktree"] = subagent_worktree.unproven_worktree_payload(info, f"finalize raised: {e}") @dataclass @@ -445,6 +424,35 @@ class _ChildWorkspace: parent_reads_snapshot: list +def _create_isolated_worktree(parent_agent: Any, parent_task_id: Any, subagent_id: Optional[str]): + """Opt-in worktree isolation: own git worktree off the parent's HEAD (the + child's terminal starts there). Git-only, local-backend-only; failures + degrade silently to the shared workspace. Returns the worktree info or None.""" + from tools.delegate_tool import _get_worktree_isolation, _resolve_workspace_hint + + if not _get_worktree_isolation(): + return None + try: + from tools import subagent_worktree + + if not subagent_worktree.local_backend_active(): + logger.debug("worktree isolation skipped: non-local terminal backend") + return None + _parent_cwd = None + try: + from tools.terminal_tool import get_session_cwd as _gsc + + _parent_cwd = _gsc(parent_task_id) + except Exception: + pass + return subagent_worktree.create_subagent_worktree( + _parent_cwd or _resolve_workspace_hint(parent_agent), subagent_id=subagent_id, + ) + except Exception as e: + logger.debug("worktree isolation setup failed: %s", e) + return None + + def _seed_child_workspace( child: Any, parent_agent: Any, @@ -458,7 +466,6 @@ def _seed_child_workspace( Returns the ids the run needs; ``goal`` comes back extended with the worktree contract note when isolation engaged. """ - from tools.delegate_tool import _get_worktree_isolation, _resolve_workspace_hint import uuid as _uuid child_task_id = subagent_id or f"subagent-{task_index}-{_uuid.uuid4().hex[:8]}" @@ -474,42 +481,19 @@ def _seed_child_workspace( except Exception as e: logger.debug("Child cwd seed failed: %s", e) - # Opt-in worktree isolation: own git worktree off the parent's HEAD, terminal - # started there. Git-only, local-backend-only; failures degrade silently. - _worktree_info = None - if _get_worktree_isolation(): + _worktree_info = _create_isolated_worktree(parent_agent, parent_task_id, subagent_id) + if _worktree_info is not None: try: - from tools import subagent_worktree + from tools.terminal_tool import record_session_cwd as _rsc - if subagent_worktree.local_backend_active(): - _parent_cwd = None - try: - from tools.terminal_tool import get_session_cwd as _gsc - - _parent_cwd = _gsc(parent_task_id) - except Exception: - pass - _worktree_info = subagent_worktree.create_subagent_worktree( - _parent_cwd or _resolve_workspace_hint(parent_agent), - subagent_id=subagent_id, - ) - else: - logger.debug("worktree isolation skipped: non-local terminal backend") + _rsc(child_task_id, _worktree_info["path"]) except Exception as e: - logger.debug("worktree isolation setup failed: %s", e) - if _worktree_info is not None: - try: - from tools.terminal_tool import record_session_cwd as _rsc + logger.debug("worktree cwd seed failed: %s", e) + # The child's context is already built; carry the isolation contract on + # the goal message instead (same turn, no system-prompt mutation). + from tools.subagent_worktree import build_worktree_context_note - _rsc(child_task_id, _worktree_info["path"]) - except Exception as e: - logger.debug("worktree cwd seed failed: %s", e) - # The child's context is already built; carry the isolation - # contract on the goal message instead (same turn, no - # system-prompt mutation). - from tools.subagent_worktree import build_worktree_context_note - - goal = goal + build_worktree_context_note(_worktree_info) + goal = goal + build_worktree_context_note(_worktree_info) worktree.info = _worktree_info parent_reads_snapshot = list(file_state.known_reads(parent_task_id)) if parent_task_id else [] diff --git a/tools/delegate_tool_config.py b/tools/delegate_tool_config.py index c1bb9c0e36..f2429684dc 100644 --- a/tools/delegate_tool_config.py +++ b/tools/delegate_tool_config.py @@ -1,7 +1,6 @@ """Delegation config knobs (delegation.* keys) and child credential/provider resolution. -Split out of ``tools/delegate_tool.py``; every moved name is re-imported there, so -``tools.delegate_tool.`` keeps resolving (and monkeypatching) as before. +Split out of ``tools/delegate_tool.py``, which re-imports every name (patch targets stay valid). """ from __future__ import annotations diff --git a/tools/delegate_tool_dispatch.py b/tools/delegate_tool_dispatch.py index 84d92d5b85..7993c795b6 100644 --- a/tools/delegate_tool_dispatch.py +++ b/tools/delegate_tool_dispatch.py @@ -32,6 +32,29 @@ def _future_entry(future: Any, idx: int, child: Any) -> Dict[str, Any]: return _fabricated_entry(idx, "error", str(exc), child) +def _report_child_done(parent_agent, spinner_ref, entry, tag, task_labels, n_tasks, remaining) -> None: + """Print one completion line for a finished child and refresh the spinner text.""" + idx = entry["task_index"] + label = task_labels[idx] if idx < len(task_labels) else f"Task {idx}" + status = entry.get("status", "?") + _slot = f"{tag} · {idx+1}/{n_tasks}" if tag else f"{idx+1}/{n_tasks}" + completion_line = f"{'✓' if status == 'completed' else '✗'} [{_slot}] {label} ({entry.get('duration_seconds', 0)}s)" + # Failed/errored/timed-out children: say WHY on the same line — a bare ✗ + # reads as "silently dropped". + if status in SUBAGENT_FAILURE_STATUSES: + _err_line = _clean_error_text(entry.get("error"), max_chars=120) + if _err_line: + completion_line += f" — {_err_line}" + _print_completion_line(parent_agent, spinner_ref, completion_line) + if spinner_ref and remaining > 0: + try: + spinner_ref.update_text( + f"🔀 {'[' + tag + '] ' if tag else ''}{remaining} task{'s' if remaining != 1 else ''} remaining" + ) + except Exception as e: + logger.debug("Spinner update_text failed: %s", e) + + def _run_children_parallel( children: List[tuple], results: list, @@ -111,28 +134,7 @@ def _run_children_parallel( results.append(entry) completed_count += 1 - idx = entry["task_index"] - label = task_labels[idx] if idx < len(task_labels) else f"Task {idx}" - status = entry.get("status", "?") - icon = "✓" if status == "completed" else "✗" - remaining = n_tasks - completed_count - _slot = f"{_tag} · {idx+1}/{n_tasks}" if _tag else f"{idx+1}/{n_tasks}" - completion_line = f"{icon} [{_slot}] {label} ({entry.get('duration_seconds', 0)}s)" - # Failed/errored/timed-out children: say WHY on the same line — - # a bare ✗ reads as "silently dropped". - if status in SUBAGENT_FAILURE_STATUSES: - _err_line = _clean_error_text(entry.get("error"), max_chars=120) - if _err_line: - completion_line += f" — {_err_line}" - _print_completion_line(parent_agent, spinner_ref, completion_line) - - if spinner_ref and remaining > 0: - try: - spinner_ref.update_text( - f"🔀 {'[' + _tag + '] ' if _tag else ''}{remaining} task{'s' if remaining != 1 else ''} remaining" - ) - except Exception as e: - logger.debug("Spinner update_text failed: %s", e) + _report_child_done(parent_agent, spinner_ref, entry, _tag, task_labels, n_tasks, n_tasks - completed_count) # Sort by task_index so results match input order results.sort(key=lambda r: r["task_index"]) diff --git a/tools/delegate_tool_progress.py b/tools/delegate_tool_progress.py index 1fbfd41417..27a1b3ba72 100644 --- a/tools/delegate_tool_progress.py +++ b/tools/delegate_tool_progress.py @@ -1,7 +1,6 @@ """Child progress relay, console formatting and child system-prompt construction for delegate_task. -Split out of ``tools/delegate_tool.py``; every moved name is re-imported there, so -``tools.delegate_tool.`` keeps resolving (and monkeypatching) as before. +Split out of ``tools/delegate_tool.py``, which re-imports every name (patch targets stay valid). """ from __future__ import annotations diff --git a/tools/delegate_tool_registry.py b/tools/delegate_tool_registry.py index 21d366eb5d..df6d606bb4 100644 --- a/tools/delegate_tool_registry.py +++ b/tools/delegate_tool_registry.py @@ -1,7 +1,6 @@ """Live-subagent registry + model-facing control plane (list/steer/stop) for delegate_task. -Split out of ``tools/delegate_tool.py``; every moved name is re-imported there, so -``tools.delegate_tool.`` keeps resolving (and monkeypatching) as before. +Split out of ``tools/delegate_tool.py``, which re-imports every name (patch targets stay valid). """ from __future__ import annotations @@ -217,6 +216,10 @@ def _capture_gateway_steer_authority( return None, None +# Registry record fields never exposed to the TUI/RPC snapshot. +_PRIVATE_RECORD_KEYS = frozenset({"agent", "owner_session_id", "owner_transport", "owner_session_record", "accepting_steer"}) + + def list_active_subagents() -> List[Dict[str, Any]]: """Snapshot of the currently running subagent tree. @@ -224,21 +227,7 @@ def list_active_subagents() -> List[Dict[str, Any]]: tool_count, status}. Safe to call from any thread — returns a copy. """ with _active_subagents_lock: - return [ - { - k: v - for k, v in r.items() - if k - not in { - "agent", - "owner_session_id", - "owner_transport", - "owner_session_record", - "accepting_steer", - } - } - for r in _active_subagents.values() - ] + return [{k: v for k, v in r.items() if k not in _PRIVATE_RECORD_KEYS} for r in _active_subagents.values()] def _is_descendant_of(child_agent: Any, parent_agent: Any, max_hops: int = 8) -> bool: diff --git a/tools/delegate_tool_results.py b/tools/delegate_tool_results.py index 3963506314..f2b13ec278 100644 --- a/tools/delegate_tool_results.py +++ b/tools/delegate_tool_results.py @@ -1,7 +1,6 @@ """Subagent result post-processing: summary budget/spill, tool-trace summaries, lifecycle hooks and cost rollup. -Split out of ``tools/delegate_tool.py``; every moved name is re-imported there, so -``tools.delegate_tool.`` keeps resolving (and monkeypatching) as before. +Split out of ``tools/delegate_tool.py``, which re-imports every name (patch targets stay valid). """ from __future__ import annotations @@ -83,17 +82,12 @@ def _stringify_tool_content(content: Any) -> str: if isinstance(content, str): return content if isinstance(content, list): - parts = [] - for item in content: - if isinstance(item, dict): - text = item.get("text") - if isinstance(text, str): - parts.append(text) - else: - parts.append(json.dumps(item, ensure_ascii=False, default=str)) - else: - parts.append(str(item)) - return "\n".join(parts) + return "\n".join( + item["text"] if isinstance(item, dict) and isinstance(item.get("text"), str) + else json.dumps(item, ensure_ascii=False, default=str) if isinstance(item, dict) + else str(item) + for item in content + ) if isinstance(content, dict): return json.dumps(content, ensure_ascii=False, default=str) return str(content) @@ -249,12 +243,7 @@ def _looks_like_error_output(content: Any) -> bool: pass first = content.splitlines()[0].strip().lower() if content.splitlines() else "" - return ( - first.startswith("error:") - or first.startswith("failed:") - or first.startswith("traceback ") - or first.startswith("exception:") - ) + return first.startswith(("error:", "failed:", "traceback ", "exception:")) # Hard per-summary character ceiling layered on top of the dynamic @@ -494,6 +483,76 @@ def _parent_finalization_lock(parent_agent) -> threading.RLock: return lock +def _notify_memory_manager(results, task_list, child_by_index, parent_agent) -> None: + memory = getattr(parent_agent, "_memory_manager", None) if parent_agent else None + if not memory: + return + for entry in results: + try: + task_index = entry.get("task_index", -1) + in_range = isinstance(task_index, int) and 0 <= task_index < len(task_list) + memory.on_delegation( + task=task_list[task_index]["goal"] if in_range else "", + result=entry.get("summary", "") or "", + child_session_id=getattr(child_by_index.get(task_index), "session_id", ""), + ) + except Exception: + pass + + +def _fire_subagent_stop_hooks(results, child_by_index, parent_agent) -> float: + """Pop the model-hidden ``_child_role`` / ``_child_cost_usd`` fields from every + entry, fire ``subagent_stop`` per child, and return the summed child cost.""" + 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 = child_by_index.get(entry.get("task_index", -1)) + invoke_hook( + "subagent_stop", + parent_session_id=getattr(parent_agent, "session_id", None), + 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) + return children_cost_total + + +def _rollup_children_cost(parent_agent, children_cost_total: float) -> None: + """Fold the children's spend into the parent's session cost (source/status + only set when the parent had none of its own).""" + if children_cost_total <= 0.0: + return + 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 _finalize_child_results( results: List[Dict[str, Any]], task_list: List[Dict[str, Any]], @@ -504,82 +563,8 @@ def _finalize_child_results( 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) + _notify_memory_manager(results, task_list, child_by_index, parent_agent) + _rollup_children_cost(parent_agent, _fire_subagent_stop_hooks(results, child_by_index, parent_agent)) def _run_child_lifecycle(