From 3da2ff88e33438acc61964b515685777f880fb12 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 17:03:55 -0700 Subject: [PATCH] refactor(delegate): pack re-import block; single blank line between top-level defs (AST-identical) --- tools/delegate_tool.py | 123 +++++-------------------------- tools/delegate_tool_child_run.py | 28 ------- tools/delegate_tool_config.py | 24 ------ tools/delegate_tool_dispatch.py | 10 --- tools/delegate_tool_progress.py | 14 ---- tools/delegate_tool_registry.py | 18 ----- tools/delegate_tool_results.py | 23 ------ 7 files changed, 20 insertions(+), 220 deletions(-) diff --git a/tools/delegate_tool.py b/tools/delegate_tool.py index 7c178d40eb..0decc82898 100644 --- a/tools/delegate_tool.py +++ b/tools/delegate_tool.py @@ -30,88 +30,39 @@ logger = logging.getLogger(__name__) # moved name is re-imported so ``tools.delegate_tool.`` keeps resolving for # callers and patching tests. Mutable flag globals live only in their owning module. from tools.delegate_tool_child_run import ( # noqa: F401 - _WorktreeReporter, - _append_missed_steer, - _append_sibling_write_reminder, - _attach_child, - _await_child, - _build_result_entry, - _cleanup_child_run, - _dump_subagent_timeout_diagnostic, - _emit_child_complete, - _fabricated_entry, - _lease_child_credential, - _make_text_relay, - _merge_late_steer, - _register_child, - _seed_child_workspace, - _start_heartbeat, + _WorktreeReporter, _append_missed_steer, _append_sibling_write_reminder, _attach_child, + _await_child, _build_result_entry, _cleanup_child_run, _dump_subagent_timeout_diagnostic, + _emit_child_complete, _fabricated_entry, _lease_child_credential, _make_text_relay, + _merge_late_steer, _register_child, _seed_child_workspace, _start_heartbeat, _validate_child_output_schema, ) from tools.delegate_tool_config import ( # noqa: F401 - _DEFAULT_MAX_CONCURRENT_CHILDREN, - _get_child_timeout, - _get_inherit_mcp_toolsets, - _get_max_async_children, - _get_max_concurrent_children, - _get_max_spawn_depth, - _get_orchestrator_enabled, - _get_subagent_approval_callback, - _get_worktree_isolation, - _inherit_parent_base_url, - _inherit_parent_capabilities, - _load_config, - _merge_request_overrides, - _resolve_child_credential_pool, - _require_pinned_command, - _resolve_delegation_credentials, - _subagent_auto_approve, - _subagent_auto_deny, + _DEFAULT_MAX_CONCURRENT_CHILDREN, _get_child_timeout, _get_inherit_mcp_toolsets, + _get_max_async_children, _get_max_concurrent_children, _get_max_spawn_depth, + _get_orchestrator_enabled, _get_subagent_approval_callback, _get_worktree_isolation, + _inherit_parent_base_url, _inherit_parent_capabilities, _load_config, _merge_request_overrides, + _resolve_child_credential_pool, _require_pinned_command, _resolve_delegation_credentials, + _subagent_auto_approve, _subagent_auto_deny, ) from tools.delegate_tool_dispatch import ( # noqa: F401 - _dispatch_background, - _run_children_parallel, + _dispatch_background, _run_children_parallel, ) from tools.delegate_tool_progress import ( # noqa: F401 - DelegateEvent, - SUBAGENT_FAILURE_STATUSES, - _batch_prefix, - _build_child_progress_callback, - _build_child_system_prompt, - _clean_error_text, - _emit_parent_console, - _resolve_workspace_hint, - _safe_progress, - format_batch_tag, - format_subagent_failure_line, + DelegateEvent, SUBAGENT_FAILURE_STATUSES, _batch_prefix, _build_child_progress_callback, + _build_child_system_prompt, _clean_error_text, _emit_parent_console, _resolve_workspace_hint, + _safe_progress, format_batch_tag, format_subagent_failure_line, ) from tools.delegate_tool_registry import ( # noqa: F401 - _CONTROL_ACTIONS, - _active_subagents, - _active_subagents_lock, - _capture_gateway_steer_authority, - _close_subagent_steering, - _handle_control_action, - _is_descendant_of, - _owns_subagent_record, - _register_subagent, - _unregister_subagent, - get_subagent_attribution, - interrupt_subagent, - is_spawn_paused, - list_active_subagents, - set_spawn_paused, - steer_subagent, + _CONTROL_ACTIONS, _active_subagents, _active_subagents_lock, _capture_gateway_steer_authority, + _close_subagent_steering, _handle_control_action, _is_descendant_of, _owns_subagent_record, + _register_subagent, _unregister_subagent, get_subagent_attribution, interrupt_subagent, + is_spawn_paused, list_active_subagents, set_spawn_paused, steer_subagent, ) from tools.delegate_tool_results import ( # noqa: F401 - _apply_summary_budget, - _build_child_preserving_parent_tools, - _finalize_child_results, - _run_child_lifecycle, - _summarize_tool_arguments, + _apply_summary_budget, _build_child_preserving_parent_tools, _finalize_child_results, + _run_child_lifecycle, _summarize_tool_arguments, ) - # Tools that children must never have access to DELEGATE_BLOCKED_TOOLS = frozenset( [ @@ -123,7 +74,6 @@ DELEGATE_BLOCKED_TOOLS = frozenset( ] ) - # Nested delegation is granted by depth/role in _build_child_agent, never by the # model naming toolsets (there is no model-facing toolsets argument). def _normalize_role(r: Optional[str]) -> str: @@ -136,7 +86,6 @@ def _normalize_role(r: Optional[str]) -> str: logger.warning("Unknown delegate_task role=%r, coercing to 'leaf'", r) return "leaf" - def _is_mcp_toolset_name(name: str) -> bool: """Return True for canonical MCP toolsets and their registered aliases.""" if not name: @@ -151,7 +100,6 @@ def _is_mcp_toolset_name(name: str) -> bool: target = None return bool(target and str(target).startswith("mcp-")) - def _expand_parent_toolsets(parent_toolsets: set) -> set: """Add every toolset whose tools are a subset of the parent's tools. @@ -176,7 +124,6 @@ def _expand_parent_toolsets(parent_toolsets: set) -> set: expanded.add(ts_name) return expanded - def _preserve_parent_mcp_toolsets(child_toolsets: List[str], parent_toolsets: set[str]) -> List[str]: """Append any parent MCP toolsets that are missing from a narrowed child.""" preserved = list(child_toolsets) @@ -185,7 +132,6 @@ def _preserve_parent_mcp_toolsets(child_toolsets: List[str], parent_toolsets: se preserved.append(toolset_name) return preserved - DEFAULT_MAX_ITERATIONS = 250 _HEARTBEAT_INTERVAL = 30 # seconds between parent activity heartbeats during delegation # Stale-heartbeat thresholds (cycles of _HEARTBEAT_INTERVAL with no progress). @@ -197,12 +143,10 @@ _HEARTBEAT_STALE_CYCLES_IDLE = 15 # 450s idle between turns → stale _HEARTBEAT_STALE_CYCLES_IN_TOOL = 40 # 1200s stuck on same tool → stale DEFAULT_TOOLSETS = ["terminal", "file", "web"] - def check_delegate_requirements() -> bool: """Delegation has no external requirements -- always available.""" return True - def _strip_blocked_tools(toolsets: List[str]) -> List[str]: """Remove toolsets whose tools are ALL blocked (derived from DELEGATE_BLOCKED_TOOLS so the two can't drift) plus composite toolsets children must never get.""" @@ -216,7 +160,6 @@ def _strip_blocked_tools(toolsets: List[str]) -> List[str]: blocked_toolset_names.add("kanban") return [t for t in toolsets if t not in blocked_toolset_names] - def _blocked_toolsets_for_role(role: str) -> List[str]: """One-tool deny toolsets for the role; passed as ``disabled_toolsets`` so blocked names inside mixed bundles are subtracted AFTER composite expansion.""" @@ -230,7 +173,6 @@ def _blocked_toolsets_for_role(role: str) -> List[str]: and set(defn.get("tools", ())).issubset(blocked_names) ) - def _resolve_child_toolsets( parent_agent, toolsets: Optional[List[str]], effective_role: str ) -> tuple[List[str], List[str]]: @@ -285,7 +227,6 @@ def _resolve_child_toolsets( ) return child_toolsets, child_disabled_toolsets - # OpenRouter routing filters: inherited from the parent, but reset to these # defaults under a pinned provider — parent filters (e.g. only=["Anthropic"]) # would silently force the child back onto the parent's provider. @@ -300,7 +241,6 @@ _ROUTING_FILTER_DEFAULTS = ( ) _NOUS_PROVIDERS = frozenset({"nous", "nous-portal", "nousresearch"}) - def _resolve_child_runtime( parent_agent, delegation_cfg: dict, @@ -405,7 +345,6 @@ def _resolve_child_runtime( kwargs["max_tokens"] = child_max_tokens return kwargs - def _open_child_session_db(parent_agent) -> Any: """DEDICATED SessionDB handle for the child, or None. @@ -426,7 +365,6 @@ def _open_child_session_db(parent_agent) -> Any: logger.debug("subagent: failed to open dedicated SessionDB; child persistence disabled", exc_info=True) return None - def _construct_child_agent( rt: Dict[str, Any], *, @@ -492,7 +430,6 @@ def _construct_child_agent( pass raise - def _announce_child_spawn(child, parent_agent, child_progress_cb, *, goal, subagent_id, parent_subagent_id, role) -> None: """spawn_requested event (now — the child may queue for seconds when the pool is saturated) plus the subagent_start lifecycle hook.""" @@ -512,7 +449,6 @@ def _announce_child_spawn(child, parent_agent, child_progress_cb, *, goal, subag except Exception: logger.debug("subagent_start hook invocation failed", exc_info=True) - def _build_child_agent( task_index: int, goal: str, @@ -642,7 +578,6 @@ def _build_child_agent( ) return child - def _run_single_child( task_index: int, goal: str, @@ -755,7 +690,6 @@ def _run_single_child( close_deferred=_child_close_deferred, ) - def _recover_tasks_from_json_string(tasks: Any) -> tuple[Optional[List[Dict[str, Any]]], Optional[str]]: if not isinstance(tasks, str): return None, None @@ -773,7 +707,6 @@ def _recover_tasks_from_json_string(tasks: Any) -> tuple[Optional[List[Dict[str, return None, (f"tasks must be a JSON array of task objects; parsed " f"{type(parsed).__name__} instead.") return parsed, None - # Placeholder shapes for batch goal validation: bare 'TODO' / 'task N' labels, # or unexpanded template markers. The marker regex is deliberately NARROW — # only snake_case / space-separated placeholder identifiers (``, @@ -788,7 +721,6 @@ _TEMPLATE_MARKER_RE = re.compile( ) _MIN_BATCH_GOAL_LEN = 10 - def _validate_batch_tasks(task_list: List[Dict[str, Any]]) -> Optional[str]: """Validate a tasks=[...] batch beyond per-task goal presence; actionable error string or None. @@ -848,7 +780,6 @@ class _Batch: origin_owner_session_record: Any overall_start: float - def _normalize_task_list( goal, context, tasks, output_schema, top_role: str, max_children: int ) -> tuple[Optional[List[Dict[str, Any]]], Optional[str]]: @@ -898,7 +829,6 @@ def _normalize_task_list( return None, batch_error return task_list, None - def _coerce_task_schemas( task_list: List[Dict[str, Any]], output_schema: Optional[Dict[str, Any]] ) -> tuple[List[Optional[Dict[str, Any]]], Optional[str]]: @@ -918,7 +848,6 @@ def _coerce_task_schemas( task_schemas.append(coerced_schema) return task_schemas, None - def _announce_batch(parent_agent, n_tasks: int, live_deleg_id: Optional[str]) -> None: """Announce the batch tag once so interleaved ``[tag n/N]`` lines are attributable.""" if n_tasks <= 1 or not live_deleg_id: @@ -933,7 +862,6 @@ def _announce_batch(parent_agent, n_tasks: int, live_deleg_id: Optional[str]) -> pass _emit_parent_console(parent_agent, _hdr) - def _capture_origin() -> tuple[str, str, Any, Any]: """``(wake_sid, ui_session_id, owner_transport, owner_session_record)`` of the ORIGINATING session, captured BEFORE building any child: AIAgent construction @@ -950,7 +878,6 @@ def _capture_origin() -> tuple[str, str, Any, Any]: transport, record = _capture_gateway_steer_authority(_origin_ui_session_id) return _origin_wake_sid, _origin_ui_session_id, transport, record - def _build_children( task_list: List[Dict[str, Any]], task_schemas: List[Optional[Dict[str, Any]]], @@ -1018,7 +945,6 @@ def _build_children( children.append((i, t, child)) return children, None - def _finalize_live_transcripts(results: list, live_writers: list, live_paths: list) -> None: """Close out live transcripts (files are retained as the full-fidelity record; retention pruning happens on future dispatches).""" @@ -1033,7 +959,6 @@ def _finalize_live_transcripts(results: list, live_writers: list, live_paths: li if _idx < len(live_paths): entry["live_transcript"] = live_paths[_idx] - def _execute_and_aggregate(batch: _Batch, *, honor_parent_interrupt: bool = True) -> dict: """Run all built children, join, finalize (hooks + cost rollup), return the combined dict. @@ -1081,7 +1006,6 @@ def _execute_and_aggregate(batch: _Batch, *, honor_parent_interrupt: bool = True combined["live_transcripts"] = list(batch.live_paths) return combined - def delegate_task( goal: Optional[str] = None, context: Optional[str] = None, @@ -1212,7 +1136,6 @@ def delegate_task( # ── OpenAI function-calling schema ────────────────────────────────────────── - def _build_top_level_description() -> str: """delegate_task description: ONLY guidance stated nowhere else in the schema (limits live in the 'tasks' parameter description, rebuilt per get_definitions()).""" @@ -1268,7 +1191,6 @@ def _build_top_level_description() -> str: "delegation.provider / delegation.model in config.yaml." ) - def _build_tasks_param_description() -> str: """Compose the 'tasks' parameter description with current concurrency limit.""" try: @@ -1282,7 +1204,6 @@ def _build_tasks_param_description() -> str: "is a one-entry array. Required when spawning." ) - def _build_dynamic_schema_overrides() -> dict: """Per-call schema overrides (ToolEntry.dynamic_schema_overrides): every get_definitions() pass rewrites the descriptions to the user's actual limits.""" @@ -1293,7 +1214,6 @@ def _build_dynamic_schema_overrides() -> dict: return {"description": _build_top_level_description(), "parameters": overrides_params} - DELEGATE_TASK_SCHEMA = { "name": "delegate_task", # description / tasks.description are placeholders: the real text is built per @@ -1395,7 +1315,6 @@ DELEGATE_TASK_SCHEMA = { # --- Registry --- from tools.registry import registry, tool_error - def _model_background_value(args: dict, parent_agent=None) -> bool: """Background flag for the MODEL-facing dispatch path (registry fallback). @@ -1408,10 +1327,8 @@ def _model_background_value(args: dict, parent_agent=None) -> bool: """ return not getattr(parent_agent, "_delegate_depth", 0) > 0 - _MODEL_HIDDEN_TASK_FIELDS = {"acp_command", "acp_args"} - def _strip_model_hidden_task_fields(tasks: Any) -> Any: if not isinstance(tasks, list): return tasks diff --git a/tools/delegate_tool_child_run.py b/tools/delegate_tool_child_run.py index a929fb2033..aebace61ab 100644 --- a/tools/delegate_tool_child_run.py +++ b/tools/delegate_tool_child_run.py @@ -27,12 +27,10 @@ from tools.delegate_tool_results import ( # Log-record parity with the origin module. logger = logging.getLogger("tools.delegate_tool") - def _num(value: Any, default: int = 0) -> int: """int() for counters that may be mocks/None on test doubles.""" return int(value) if isinstance(value, (int, float)) else default - def _fabricated_entry(idx: int, status: str, error: str, child: Any, duration: float = 0) -> Dict[str, Any]: """Result entry for a child that raised, never finished, or was abandoned.""" return { @@ -45,14 +43,12 @@ def _fabricated_entry(idx: int, status: str, error: str, child: Any, duration: f "_child_role": getattr(child, "_delegate_role", None), } - def _append_missed_steer(entry: Dict[str, Any], late_steer: Optional[str]) -> None: """Record steer text that won the race with the child's failure/timeout.""" if late_steer: entry["missed_steer"] = late_steer entry["error"] += (" [steer did not land before the subagent stopped: " f"{late_steer}]") - def _close_child(child: Any, log_message: str) -> None: """Best-effort ``child.close()`` (tool sandboxes, browser daemons, httpx clients).""" try: @@ -62,7 +58,6 @@ def _close_child(child: Any, log_message: str) -> None: except Exception: logger.debug(log_message, exc_info=True) - def _attach_child(parent_agent: Any, child: Any) -> None: """Register the child for parent interrupt propagation.""" if not hasattr(parent_agent, "_active_children"): @@ -74,7 +69,6 @@ def _attach_child(parent_agent: Any, child: Any) -> None: else: parent_agent._active_children.append(child) - def _detach_child(parent_agent: Any, child: Any) -> None: """Remove the child from parent interrupt propagation (no-op if absent).""" if not hasattr(parent_agent, "_active_children"): @@ -89,7 +83,6 @@ def _detach_child(parent_agent: Any, child: Any) -> None: except (ValueError, UnboundLocalError) as e: logger.debug("Could not remove child from active_children: %s", e) - def _signal_child_stop(child: Any, *reason: str) -> None: """Cooperative interrupt so the child's worker thread can exit cleanly.""" try: @@ -98,7 +91,6 @@ def _signal_child_stop(child: Any, *reason: str) -> None: except Exception: pass - def _format_thread_stack(frame: Any, indent: str) -> List[str]: import traceback as _traceback @@ -108,14 +100,12 @@ def _format_thread_stack(frame: Any, indent: str) -> List[str]: for sub in frame_line.rstrip().split("\n") ] - _DIAG_CHILD_ATTRS = ( "model", "provider", "api_mode", "base_url", "max_iterations", "quiet_mode", "skip_memory", "skip_context_files", "platform", "_delegate_role", "_delegate_depth", ) - def _diag_sizes(child: Any) -> List[str]: lines: List[str] = ["## Prompt / schema sizes"] try: @@ -133,7 +123,6 @@ def _diag_sizes(child: Any) -> List[str]: lines.append(f" tool_schema: ") return lines - def _diag_threads(worker_thread: Optional[threading.Thread]) -> List[str]: """Worker stack plus all other live threads (bounded to 40): the worker is often parked on a helper thread, so a pre-HTTP wedge is indistinguishable @@ -174,7 +163,6 @@ def _diag_threads(worker_thread: Optional[threading.Thread]) -> List[str]: lines.append(f" ") return lines - def _dump_subagent_timeout_diagnostic( *, child: Any, @@ -252,7 +240,6 @@ def _dump_subagent_timeout_diagnostic( logger.warning("Subagent timeout diagnostic dump failed: %s", exc) return None - def _start_heartbeat(child: Any, parent_agent: Any, task_index: int) -> tuple: """Build the parent-activity heartbeat thread for one child (not started). @@ -319,7 +306,6 @@ def _start_heartbeat(child: Any, parent_agent: Any, task_index: int) -> tuple: return _heartbeat_stop, threading.Thread(target=_heartbeat_loop, daemon=True) - def _register_child( child: Any, parent_agent: Any, @@ -415,7 +401,6 @@ class _ChildWorkspace: wall_start: float 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 @@ -444,7 +429,6 @@ def _create_isolated_worktree(parent_agent: Any, parent_task_id: Any, subagent_i logger.debug("worktree isolation setup failed: %s", e) return None - def _seed_child_workspace( child: Any, parent_agent: Any, @@ -497,7 +481,6 @@ class _ChildFailure: entry: Dict[str, Any] close_deferred: bool - def _defer_close_after_timeout(child: Any, child_future: Any) -> None: """Hand ``child.close()`` to a Future done-callback and drain its transports. @@ -533,7 +516,6 @@ def _defer_close_after_timeout(child: Any, child_future: Any) -> None: _resweep_timer.daemon = True _resweep_timer.start() - def _lease_child_credential(child: Any) -> tuple[Any, Optional[str]]: """Lease a credential from the child's pool (if any) and bind it; ``(pool, lease_id)``.""" child_pool = getattr(child, "_credential_pool", None) @@ -549,7 +531,6 @@ def _lease_child_credential(child: Any) -> tuple[Any, Optional[str]]: logger.debug("Failed to bind child to leased credential: %s", exc) return child_pool, leased_cred_id - def _make_text_relay(child_progress_cb: Any): """Stream callback forwarding the child's reply text up the progress relay so gateway watch windows mirror it live (subagent.text → message.delta). Inert @@ -561,7 +542,6 @@ def _make_text_relay(child_progress_cb: Any): return _relay_child_text - def _await_child( child: Any, goal: str, @@ -622,7 +602,6 @@ def _await_child( # Shut down without waiting — a child stuck on blocking I/O would hang wait=True forever. executor.shutdown(wait=False) - def _merge_late_steer(result: Dict[str, Any], subagent_id: Optional[str], child: Any) -> None: """Linearization boundary for registry steering: from here the child cannot consume another steer. Closing under the registry lock either rejects a @@ -633,7 +612,6 @@ def _merge_late_steer(result: Dict[str, Any], subagent_id: Optional[str], child: existing = result.get("pending_steer") result["pending_steer"] = f"{existing}\n{late}" if isinstance(existing, str) and existing else late - def _handle_child_wait_failure( exc: BaseException, *, @@ -753,7 +731,6 @@ class _SchemaOutcome: errors: List[str] retries: int - def _validate_child_output_schema( child: Any, result: Dict[str, Any], task_index: int, child_task_id: str, relay_child_text: Any ) -> _SchemaOutcome: @@ -797,7 +774,6 @@ def _validate_child_output_schema( _schema_valid, _schema_errors = validate_output(_retry_text, _output_schema) return _SchemaOutcome(_output_schema, _schema_valid, _schema_errors, 1) - def _build_tool_trace(messages: Any) -> list[Dict[str, Any]]: """Tool trace from the child's conversation messages, pairing parallel tool calls with their results by tool_call_id.""" @@ -834,7 +810,6 @@ def _build_tool_trace(messages: Any) -> list[Dict[str, Any]]: tool_trace[-1].update(result_meta) # no tool_call_id: pair with the latest call return tool_trace - def _build_result_entry( child: Any, result: Dict[str, Any], @@ -936,7 +911,6 @@ def _build_result_entry( entry["summary"] = f"{summary}\n\n{_miss_note}" if summary else _miss_note return entry - def _append_sibling_write_reminder(entry: Dict[str, Any], ws: _ChildWorkspace) -> None: """Warn the parent when this child wrote files the parent had already read. @@ -964,7 +938,6 @@ def _append_sibling_write_reminder(entry: Dict[str, Any], ws: _ChildWorkspace) - except Exception: logger.debug("file_state sibling-write check failed", exc_info=True) - def _emit_child_complete( child: Any, result: Dict[str, Any], @@ -1014,7 +987,6 @@ def _emit_child_complete( pass _safe_progress(child_progress_cb, "subagent.complete", **complete_kwargs) - def _cleanup_child_run( child: Any, parent_agent: Any, diff --git a/tools/delegate_tool_config.py b/tools/delegate_tool_config.py index 935a1a4c3c..0cdde7c160 100644 --- a/tools/delegate_tool_config.py +++ b/tools/delegate_tool_config.py @@ -30,7 +30,6 @@ _LEGACY_MAX_ASYNC_WARNED = False # opts back in. DEFAULT_CHILD_TIMEOUT: Optional[float] = None - def _cfg() -> dict: """The ``delegation`` section, read through the origin so tests can patch it.""" from tools.delegate_tool import _load_config @@ -54,20 +53,17 @@ def _subagent_auto_deny(command: str, description: str, **kwargs) -> str: ) return "deny" - def _subagent_auto_approve(command: str, description: str, **kwargs) -> str: """Auto-approve (opt-in YOLO via delegation.subagent_auto_approve): returns 'once'.""" logger.warning("Subagent auto-approved dangerous command: %s (%s)", command, description) return "once" - def _get_subagent_approval_callback(): """Callback for subagent worker threads per delegation.subagent_auto_approve (default False).""" if is_truthy_value(_cfg().get("subagent_auto_approve", False)): return _subagent_auto_approve return _subagent_auto_deny - def _knob(key: str, env_var: str, parse, default, warn_invalid): """delegation. > > default. A config value that fails ``parse`` calls ``warn_invalid(value)`` and yields the default; an env value that fails @@ -87,7 +83,6 @@ def _knob(key: str, env_var: str, parse, default, warn_invalid): pass return default - def _get_max_concurrent_children() -> int: """delegation.max_concurrent_children > DELEGATION_MAX_CONCURRENT_CHILDREN env > 10. @@ -112,7 +107,6 @@ def _get_max_concurrent_children() -> int: ) return result - def _get_worktree_isolation() -> bool: """delegation.worktree_isolation (bool, default False). @@ -122,7 +116,6 @@ def _get_worktree_isolation() -> bool: """ return bool(_cfg().get("worktree_isolation", False)) - def _get_max_async_children() -> int: """Concurrency cap for background delegations == delegation.max_concurrent_children. @@ -142,13 +135,11 @@ def _get_max_async_children() -> int: ) return _get_max_concurrent_children() - def _parse_timeout(raw: Any) -> Optional[float]: """Seconds → None (<= 0 disables) or max(30, value). Raises on non-numeric.""" parsed = float(raw) return None if parsed <= 0 else max(30.0, parsed) - def _get_child_timeout() -> Optional[float]: """Hard wall-clock cap for one child, or None (default: no timeout). @@ -164,7 +155,6 @@ def _get_child_timeout() -> Optional[float]: ), ) - def _get_max_spawn_depth() -> int: """delegation.max_spawn_depth floored at 1 (no ceiling). @@ -184,7 +174,6 @@ def _get_max_spawn_depth() -> int: logger.warning("delegation.max_spawn_depth=%d below floor %d; using %d", ival, _MIN_SPAWN_DEPTH, floored) return floored - def _get_orchestrator_enabled() -> bool: """delegation.orchestrator_enabled kill switch (default True): False forces every child to leaf.""" val = _cfg().get("orchestrator_enabled", True) @@ -195,16 +184,13 @@ def _get_orchestrator_enabled() -> bool: return val.strip().lower() in {"true", "1", "yes", "on"} return True - def _get_inherit_mcp_toolsets() -> bool: """Whether narrowed child toolsets should keep the parent's MCP toolsets.""" return is_truthy_value(_cfg().get("inherit_mcp_toolsets"), default=True) - def _normalized_runtime_url(value: Any) -> str: return str(value or "").strip().rstrip("/") - def _inherit_parent_capabilities(parent_agent, override_provider, override_base_url) -> Optional[dict]: """Parent's endpoint-trust capability map for a child, or None. @@ -219,7 +205,6 @@ def _inherit_parent_capabilities(parent_agent, override_provider, override_base_ return None return {key: value for key, value in parent_caps.items() if isinstance(key, str) and isinstance(value, bool)} - def _inherit_parent_base_url(parent_agent, fallback_base_url: Optional[str]) -> Optional[str]: """Base URL the parent is actually calling (live client), not a stale attribute. @@ -240,7 +225,6 @@ def _inherit_parent_base_url(parent_agent, fallback_base_url: Optional[str]) -> return url return fallback_base_url or None - def _loaded_pool(key: Any): """``load_pool(key)`` when it holds credentials, else None.""" from agent.credential_pool import load_pool @@ -248,7 +232,6 @@ def _loaded_pool(key: Any): pool = load_pool(key) return pool if pool is not None and pool.has_credentials() else None - def _resolve_child_credential_pool( effective_provider: Optional[str], parent_agent, @@ -301,7 +284,6 @@ def _resolve_child_credential_pool( logger.debug("Could not load credential pool for child provider '%s': %s", effective_provider, exc) return None - def _merge_request_overrides(runtime_overrides, explicit_overrides): """Merge explicit ``delegation.request_overrides`` OVER runtime-derived ones. @@ -329,14 +311,12 @@ def _merge_request_overrides(runtime_overrides, explicit_overrides): merged["extra_body"] = explicit_extra return merged or None - # Native-SDK providers speak their own wire protocol and can't be reached via # chat_completions against a base_url: always take the runtime-provider path # (a configured base_url still flows through it, e.g. a Bedrock region). _NATIVE_SDK_PROVIDERS = frozenset({"bedrock", "vertex", "google", "google-genai"}) _EXPLICIT_API_MODES = frozenset({"chat_completions", "codex_responses", "anthropic_messages"}) - def _require_pinned_command(command: Optional[str], message: str) -> None: """A pinned ACP transport command must exist on PATH — refuse loudly rather than let the child silently fall back to another transport.""" @@ -347,7 +327,6 @@ def _require_pinned_command(command: Optional[str], message: str) -> None: if not _shutil.which(command): raise ValueError(message) - def _direct_endpoint_credentials(cfg_values: dict, explicit_request_overrides) -> dict: """``delegation.base_url`` branch: provider/api_mode from URL heuristics.""" configured_model, configured_provider, configured_base_url, configured_api_key, configured_api_mode = ( @@ -402,7 +381,6 @@ def _direct_endpoint_credentials(cfg_values: dict, explicit_request_overrides) - "max_output_tokens": max_output_tokens, } - def _runtime_provider_credentials(cfg_values: dict, explicit_request_overrides) -> dict: """``delegation.provider`` branch: full bundle via the runtime provider system.""" configured_model, configured_provider = cfg_values["model"], cfg_values["provider"] @@ -446,7 +424,6 @@ def _runtime_provider_credentials(cfg_values: dict, explicit_request_overrides) "args": list(runtime.get("args") or []), } - def _resolve_delegation_credentials(cfg: dict, parent_agent) -> dict: """Resolve the child credential bundle from the ``delegation`` config section. @@ -484,7 +461,6 @@ def _resolve_delegation_credentials(cfg: dict, parent_agent) -> dict: } return _runtime_provider_credentials(values, explicit_request_overrides) - def _load_config() -> dict: """Return the ``delegation`` config section (read-only — do NOT mutate). diff --git a/tools/delegate_tool_dispatch.py b/tools/delegate_tool_dispatch.py index 0f593f7b82..2053b66e78 100644 --- a/tools/delegate_tool_dispatch.py +++ b/tools/delegate_tool_dispatch.py @@ -23,7 +23,6 @@ from tools.delegate_tool_progress import ( # Log-record parity with the origin module. logger = logging.getLogger("tools.delegate_tool") - def _future_entry(future: Any, idx: int, child: Any) -> Dict[str, Any]: """The finished Future's entry, or a fabricated error entry if it raised.""" try: @@ -31,7 +30,6 @@ def _future_entry(future: Any, idx: int, child: Any) -> Dict[str, Any]: except Exception as exc: 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"] @@ -54,7 +52,6 @@ def _report_child_done(parent_agent, spinner_ref, entry, tag, task_labels, n_tas except Exception as e: logger.debug("Spinner update_text failed: %s", e) - def _run_children_parallel( children: List[tuple], results: list, @@ -136,7 +133,6 @@ def _run_children_parallel( # Sort by task_index so results match input order results.sort(key=lambda r: r["task_index"]) - _SYNC_FALLBACK_NOTES = { "no_async": ( "background=true is not available in this session — it cannot " @@ -154,7 +150,6 @@ _SYNC_FALLBACK_NOTES = { ), } - def _run_sync_with_note(execute_and_aggregate: Any, reason: str) -> str: """Inline fallback: run the batch now and explain why it was not detached.""" result = execute_and_aggregate() @@ -162,7 +157,6 @@ def _run_sync_with_note(execute_and_aggregate: Any, reason: str) -> str: result["note"] = _SYNC_FALLBACK_NOTES[reason] return json.dumps(result, ensure_ascii=False) - def _resolve_async_wake_sid(origin_wake_sid: str) -> Optional[str]: """Wake target for a detached batch, or None to force synchronous execution. @@ -192,7 +186,6 @@ def _resolve_async_wake_sid(origin_wake_sid: str) -> Optional[str]: return origin_wake_sid return None - def _resolve_async_session_key(parent_agent: Any, origin_ui_session_id: str) -> tuple[str, str]: """``(session_key, origin_ui_session_id)`` the async registry routes completions by. @@ -223,7 +216,6 @@ def _resolve_async_session_key(parent_agent: Any, origin_ui_session_id: str) -> session_key = agent_session_id return session_key, origin_ui_session_id - def _batch_progress_token(child_agents: List[Any]) -> tuple: """Progress token for the async registry's stale monitor: every child's (api_call_count, current_tool, last_activity_ts). last_activity_ts ticks on @@ -243,7 +235,6 @@ def _batch_progress_token(child_agents: List[Any]) -> tuple: parts.append(None) return tuple(parts), in_tool - def _dispatched_payload(dispatch: dict, goals: List[str], child_agents: List[Any], live_paths: List[str]) -> dict: """Model-facing handle for an accepted background batch.""" n = len(goals) @@ -286,7 +277,6 @@ def _dispatched_payload(dispatch: dict, goals: List[str], child_agents: List[Any ) return payload - def _dispatch_background( *, parent_agent: Any, diff --git a/tools/delegate_tool_progress.py b/tools/delegate_tool_progress.py index 1e640990fe..41feee84c6 100644 --- a/tools/delegate_tool_progress.py +++ b/tools/delegate_tool_progress.py @@ -20,7 +20,6 @@ logger = logging.getLogger("tools.delegate_tool") # the parent-facing failure summary so every surface agrees. SUBAGENT_FAILURE_STATUSES = frozenset({"failed", "error", "timeout"}) - def _safe_progress(cb: Any, event_type: Any, *args: Any, **kwargs: Any) -> None: """Invoke a child progress callback; relay failures never reach the run.""" if not cb: @@ -30,7 +29,6 @@ def _safe_progress(cb: Any, event_type: Any, *args: Any, **kwargs: Any) -> None: except Exception as e: logger.debug("Progress callback %s failed: %s", event_type, e) - def _clean_error_text(error: Any, max_chars: int = 200) -> str: """Reduce an error payload (traceback / JSON wall) to one clean line: the exception message (last line of a traceback) or the first non-empty line, @@ -41,7 +39,6 @@ def _clean_error_text(error: Any, max_chars: int = 200) -> str: line = lines[-1] if lines[0].startswith("Traceback") else lines[0] return line[: max_chars - 3] + "..." if len(line) > max_chars else line - def format_subagent_failure_line( goal: Optional[str], status: Optional[str], @@ -83,7 +80,6 @@ class DelegateEvent(str, enum.Enum): TASK_TOOL_STARTED = "delegate.tool_started" TASK_TOOL_COMPLETED = "delegate.tool_completed" - # Legacy child-agent event strings → DelegateEvent. _LEGACY_EVENT_MAP: Dict[str, DelegateEvent] = { "_thinking": DelegateEvent.TASK_THINKING, @@ -107,7 +103,6 @@ _EVENT_HANDLERS: Dict[Any, Optional[str]] = { DelegateEvent.TASK_TOOL_COMPLETED: None, } - def _normalize_event(event_type: Any) -> Any: """Lifecycle string / DelegateEvent / legacy string / ``delegate.*`` string → dispatch key; None for unknown events.""" @@ -121,7 +116,6 @@ def _normalize_event(event_type: Any) -> Any: except (ValueError, TypeError): return None - def _build_child_system_prompt( goal: str, context: Optional[str] = None, @@ -213,7 +207,6 @@ def _build_child_system_prompt( ) return "\n".join(parts) - def _resolve_workspace_hint(parent_agent) -> Optional[str]: """Best-effort local workspace hint for child prompts: only a concrete absolute directory is ever injected (never a fake container path).""" @@ -236,11 +229,9 @@ def _resolve_workspace_hint(parent_agent) -> Optional[str]: return text return None - _BATCH_ORDINALS: Dict[str, int] = {} _BATCH_ORDINALS_LOCK = threading.Lock() - def format_batch_tag(delegation_id: Optional[str]) -> str: """Short human tag for a delegation batch: ``deleg_6a664903`` → ``set 1`` (first batch seen in this process), the next distinct id → ``set 2``. @@ -257,7 +248,6 @@ def format_batch_tag(delegation_id: Optional[str]) -> str: n = _BATCH_ORDINALS.setdefault(delegation_id, len(_BATCH_ORDINALS) + 1) return f"set {n}" - def _batch_prefix(delegation_id: Optional[str], task_index: int, task_count: int) -> str: """``[set 2 · 3/9] `` for batch children, ``[set 2] `` for a lone child, ``[3/9] `` / ``""`` when the batch id is unknown.""" @@ -267,7 +257,6 @@ def _batch_prefix(delegation_id: Optional[str], task_index: int, task_count: int return f"[{inner}] " return f"[{tag}] " if tag else "" - def _emit_parent_console(parent_agent, line: str) -> None: """Emit a progress line through ``parent_agent._safe_print`` when available so headless stdio hosts (ACP, gateway API) can redirect it to stderr; a @@ -281,7 +270,6 @@ def _emit_parent_console(parent_agent, line: str) -> None: pass print(line) - def _print_completion_line(parent_agent: Any, spinner_ref: Any, line: str) -> None: """Above-spinner line when a spinner exists (console fallback if it raises), else console.""" if spinner_ref: @@ -292,7 +280,6 @@ def _print_completion_line(parent_agent: Any, spinner_ref: Any, line: str) -> No pass _emit_parent_console(parent_agent, f" {line}") - def _short(text: str, n: int) -> str: return (text[:n] + "...") if len(text) > n else text @@ -436,7 +423,6 @@ class _ChildProgressRelay: if method is not None: getattr(self, method)(tool_name, preview, args, kwargs) - def _build_child_progress_callback( task_index: int, goal: str, diff --git a/tools/delegate_tool_registry.py b/tools/delegate_tool_registry.py index 50b43b4d93..22561c22ff 100644 --- a/tools/delegate_tool_registry.py +++ b/tools/delegate_tool_registry.py @@ -36,7 +36,6 @@ _RECENT_SUBAGENTS_CAP = 200 _recent_subagents: Dict[str, Dict[str, Any]] = {} - def get_subagent_attribution(task_id: Optional[str]) -> Optional[Dict[str, Any]]: """Resolve a process task_id to its originating delegation, if any. @@ -54,7 +53,6 @@ def get_subagent_attribution(task_id: Optional[str]) -> Optional[Dict[str, Any]] return None return {"subagent_id": task_id, "goal": record.get("goal"), "delegation_id": record.get("delegation_id")} - def set_spawn_paused(paused: bool) -> bool: """Globally block/unblock new delegate_task spawns. @@ -66,12 +64,10 @@ def set_spawn_paused(paused: bool) -> bool: _spawn_paused = bool(paused) return _spawn_paused - def is_spawn_paused() -> bool: with _spawn_pause_lock: return _spawn_paused - def _register_subagent(record: Dict[str, Any]) -> None: sid = record.get("subagent_id") if not sid: @@ -80,7 +76,6 @@ def _register_subagent(record: Dict[str, Any]) -> None: with _active_subagents_lock: _active_subagents[sid] = record - def _retain_recent_subagent(record: Dict[str, Any]) -> None: """Keep a bounded attribution stub after a child finishes (lock held).""" sid = record.get("subagent_id") @@ -94,7 +89,6 @@ def _retain_recent_subagent(record: Dict[str, Any]) -> None: while len(_recent_subagents) > _RECENT_SUBAGENTS_CAP: _recent_subagents.pop(next(iter(_recent_subagents)), None) - def _unregister_subagent(subagent_id: str, *, agent: Any = None) -> None: with _active_subagents_lock: record = _active_subagents.get(subagent_id) @@ -102,7 +96,6 @@ def _unregister_subagent(subagent_id: str, *, agent: Any = None) -> None: _active_subagents.pop(subagent_id, None) _retain_recent_subagent(record) - def _close_subagent_steering(subagent_id: str, agent: Any) -> Optional[str]: """Atomically close steer acceptance and drain its final durable artifact. @@ -126,7 +119,6 @@ def _close_subagent_steering(subagent_id: str, agent: Any) -> Optional[str]: return None return pending if isinstance(pending, str) and pending.strip() else None - def interrupt_subagent(subagent_id: str) -> bool: """Request that a single running subagent stop at its next iteration boundary. @@ -150,7 +142,6 @@ def interrupt_subagent(subagent_id: str) -> bool: return False return True - def steer_subagent( subagent_id: str, text: str, @@ -193,7 +184,6 @@ def steer_subagent( logger.debug("steer_subagent(%s) failed: %s", subagent_id, exc) return False - def _capture_gateway_steer_authority(owner_session_id: Optional[str]) -> tuple[Any, Any]: """Capture exact request transport + live session generation, if any. @@ -209,11 +199,9 @@ def _capture_gateway_steer_authority(owner_session_id: Optional[str]) -> tuple[A except Exception: 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. @@ -223,7 +211,6 @@ def list_active_subagents() -> List[Dict[str, Any]]: with _active_subagents_lock: 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: """True when *child_agent* sits below *parent_agent* in the spawn tree. @@ -244,12 +231,10 @@ def _is_descendant_of(child_agent: Any, parent_agent: Any, max_hops: int = 8) -> cur = ancestor return False - # Model-facing control actions accepted by delegate_task(action=...). # "spawn" (or omitted) keeps the historical spawn semantics. _CONTROL_ACTIONS = frozenset({"list", "steer", "stop"}) - def _resolve_session_lineage(session_id: Optional[str], parent_agent: Any) -> str: """Resolve a session id to the tip of its compression lineage. @@ -269,7 +254,6 @@ def _resolve_session_lineage(session_id: Optional[str], parent_agent: Any) -> st except Exception: return sid - def _owns_subagent_record(record: Dict[str, Any], parent_agent: Any) -> bool: """True when *parent_agent*'s conversation owns this live-child record. @@ -301,7 +285,6 @@ def _owns_subagent_record(record: Dict[str, Any], parent_agent: Any) -> bool: _resolve_session_lineage(parent_sid, parent_agent), } - def _handle_control_action(action: str, subagent_id: Optional[str], message: Optional[str], parent_agent: Any) -> str: """Synchronous control plane for delegate_task: list/steer/stop. @@ -370,7 +353,6 @@ def _handle_control_action(action: str, subagent_id: Optional[str], message: Opt return json.dumps({"action": action, "subagent_id": sid, "status": status, "note": note}, ensure_ascii=False) return tool_error(failure.format(sid=sid)) - # action -> (success status, success note, failure error template) _CONTROL_OUTCOMES = { "stop": ( diff --git a/tools/delegate_tool_results.py b/tools/delegate_tool_results.py index 7c99b982d0..6477ed3457 100644 --- a/tools/delegate_tool_results.py +++ b/tools/delegate_tool_results.py @@ -14,7 +14,6 @@ from urllib.parse import urlsplit, urlunsplit # Log-record parity with the origin module. logger = logging.getLogger("tools.delegate_tool") - def _extract_output_tail( result: Dict[str, Any], *, @@ -69,7 +68,6 @@ def _extract_output_tail( tail.reverse() # restore chronological order for display return tail - def _stringify_tool_content(content: Any) -> str: """Return a stable text representation for tool-result content. @@ -92,7 +90,6 @@ def _stringify_tool_content(content: Any) -> str: return json.dumps(content, ensure_ascii=False, default=str) return str(content) - _TOOL_INPUT_TARGET_KEYS = frozenset({ "cwd", "destination_path", @@ -112,7 +109,6 @@ _TOOL_INPUT_TARGET_KEYS = frozenset({ _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): @@ -140,11 +136,9 @@ def _sanitize_tool_target(key: str, value: Any) -> Any: return None return bounded - def _empty_input_summary() -> Dict[str, Any]: return {"argument_keys": [], "targets": {}} - def _sanitize_targets(mapping: Dict[str, Any]) -> Dict[str, Any]: """Keep only known side-effect target keys, each sanitized (URL secrets dropped).""" targets: Dict[str, Any] = {} @@ -156,7 +150,6 @@ def _sanitize_targets(mapping: Dict[str, Any]) -> Dict[str, Any]: targets[key] = cleaned return targets - def _summarize_tool_arguments(arguments: Any) -> Dict[str, Any]: """Summarize argument names and side-effect targets without raw payloads.""" if not isinstance(arguments, str): @@ -170,7 +163,6 @@ def _summarize_tool_arguments(arguments: Any) -> Dict[str, Any]: keys = sorted(str(key)[:128] for key in parsed)[:64] return {"argument_keys": keys, "targets": _sanitize_targets(parsed)} - def _sanitize_tool_input_summary(summary: Any) -> Dict[str, Any]: """Re-sanitize a stored input summary before handing it to lifecycle hooks.""" if not isinstance(summary, dict): @@ -181,7 +173,6 @@ def _sanitize_tool_input_summary(summary: Any) -> Dict[str, Any]: safe_targets = _sanitize_targets(targets) if isinstance(targets, dict) else {} 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): @@ -211,7 +202,6 @@ def _subagent_stop_tool_call_history(tool_trace: Any) -> List[Dict[str, Any]]: }) return history - def _looks_like_error_output(content: Any) -> bool: """Conservative stderr/error detector for tool-result previews. @@ -242,7 +232,6 @@ def _looks_like_error_output(content: Any) -> bool: first = content.splitlines()[0].strip().lower() if content.splitlines() else "" return first.startswith(("error:", "failed:", "traceback ", "exception:")) - # Hard per-summary character ceiling layered on top of the dynamic # headroom budget (see _apply_summary_budget). Belt-and-suspenders for # models that ignore the "be concise" instruction. 0 disables the ceiling. @@ -258,7 +247,6 @@ _SUMMARY_HEADROOM_FRACTION = 0.5 # already nearly full — below this we'd be truncating to noise. _MIN_SUMMARY_CHARS = 2000 - def _spill_summary_to_file(task_index: int, summary: str) -> Optional[str]: """Write a subagent's full summary to the delegation cache and return path. @@ -288,7 +276,6 @@ def _spill_summary_to_file(task_index: int, summary: str) -> Optional[str]: logger.debug("Failed to spill subagent summary to file: %s", exc) return None - def _trim_summary_with_footer(summary: str, cap: int, task_index: int) -> tuple[str, Optional[str]]: """Return (model_text, spill_path) for one over-budget summary. @@ -341,7 +328,6 @@ def _trim_summary_with_footer(summary: str, cap: int, task_index: int) -> tuple[ model_text = head + "\n\n[... middle omitted — see footer ...]\n\n" + tail + "\n".join(footer_lines) return model_text, spill_path - def _parent_summary_char_budget(parent_agent, n_summaries: int) -> Optional[int]: """Per-summary char budget from the parent's *remaining* context headroom (context length minus prompt tokens minus the compressor's output reserve), @@ -374,7 +360,6 @@ def _parent_summary_char_budget(parent_agent, n_summaries: int) -> Optional[int] logger.debug("Summary budget computation failed", exc_info=True) return None - def _apply_summary_budget(results: List[Dict[str, Any]], parent_agent) -> None: """Trim subagent summaries in-place so a batch can't overflow the parent's context window; full text is spilled to disk so nothing is lost. @@ -422,14 +407,12 @@ def _apply_summary_budget(results: List[Dict[str, Any]], parent_agent) -> None: spill_path or "none", ) - _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.""" from tools.delegate_tool import _build_child_agent @@ -444,7 +427,6 @@ def _build_child_preserving_parent_tools(**kwargs): 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: @@ -462,7 +444,6 @@ def _parent_finalization_lock(parent_agent) -> threading.RLock: return _PARENT_FINALIZATION_FALLBACK_LOCK 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: @@ -479,7 +460,6 @@ def _notify_memory_manager(results, task_list, child_by_index, parent_agent) -> 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.""" @@ -516,7 +496,6 @@ def _fire_subagent_stop_hooks(results, child_by_index, parent_agent) -> float: 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).""" @@ -532,7 +511,6 @@ def _rollup_children_cost(parent_agent, children_cost_total: float) -> None: 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]], @@ -546,7 +524,6 @@ def _finalize_child_results( _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(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.""" from tools.delegate_tool import _run_single_child