refactor(delegate): split _finalize_child_results (memory/hooks/cost rollup), lift _create_isolated_worktree and _report_child_done, collapse dead worktree-attach import fallback
This commit is contained in:
BIN
MagicMock/mock._session_db.db_path/132161382957648
Normal file
BIN
MagicMock/mock._session_db.db_path/132161382957648
Normal file
Binary file not shown.
BIN
MagicMock/mock._session_db.db_path/132161402292368
Normal file
BIN
MagicMock/mock._session_db.db_path/132161402292368
Normal file
Binary file not shown.
@@ -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 []
|
||||
|
||||
@@ -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.<name>`` 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
|
||||
|
||||
@@ -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"])
|
||||
|
||||
@@ -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.<name>`` 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
|
||||
|
||||
@@ -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.<name>`` 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:
|
||||
|
||||
@@ -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.<name>`` 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(
|
||||
|
||||
Reference in New Issue
Block a user