diff --git a/MagicMock/mock._session_db.db_path/140588728862416 b/MagicMock/mock._session_db.db_path/140588728862416 new file mode 100644 index 0000000000..960085b500 Binary files /dev/null and b/MagicMock/mock._session_db.db_path/140588728862416 differ diff --git a/MagicMock/mock._session_db.db_path/140588728862416.fts_rebuild.lock b/MagicMock/mock._session_db.db_path/140588728862416.fts_rebuild.lock new file mode 100644 index 0000000000..e69de29bb2 diff --git a/MagicMock/mock._session_db.db_path/140588728862416.quarantine.lock b/MagicMock/mock._session_db.db_path/140588728862416.quarantine.lock new file mode 100644 index 0000000000..e69de29bb2 diff --git a/MagicMock/mock._session_db.db_path/140588740754704 b/MagicMock/mock._session_db.db_path/140588740754704 new file mode 100644 index 0000000000..1b27f54eb2 Binary files /dev/null and b/MagicMock/mock._session_db.db_path/140588740754704 differ diff --git a/MagicMock/mock._session_db.db_path/140588740754704.fts_rebuild.lock b/MagicMock/mock._session_db.db_path/140588740754704.fts_rebuild.lock new file mode 100644 index 0000000000..e69de29bb2 diff --git a/MagicMock/mock._session_db.db_path/140588740754704.quarantine.lock b/MagicMock/mock._session_db.db_path/140588740754704.quarantine.lock new file mode 100644 index 0000000000..e69de29bb2 diff --git a/tools/delegate_tool.py b/tools/delegate_tool.py index 6eeb863982..7d68a0583f 100644 --- a/tools/delegate_tool.py +++ b/tools/delegate_tool.py @@ -71,6 +71,7 @@ from tools.delegate_tool_config import ( # noqa: F401 _load_config, _merge_request_overrides, _resolve_child_credential_pool, + _require_pinned_command, _resolve_delegation_credentials, _subagent_auto_approve, _subagent_auto_deny, @@ -307,28 +308,19 @@ def _resolve_child_toolsets( return child_toolsets, child_disabled_toolsets -@dataclass -class _ChildRuntime: - """Provider/transport/routing settings resolved for one child AIAgent.""" - - model: Any - provider: Any - base_url: Any - api_key: Any - api_mode: Any - capabilities: Optional[dict] - acp_command: Any - acp_args: list - reasoning: Any - fallback: Any - providers_allowed: Any - providers_ignored: Any - providers_order: Any - provider_sort: Any - provider_require_parameters: Any - provider_data_collection: Any - openrouter_min_coding_score: Any - optional_kwargs: Dict[str, Any] +# 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. +# openrouter_min_coding_score stays inherited: model-gated, no-op elsewhere. +_ROUTING_FILTER_DEFAULTS = ( + ("providers_allowed", None), + ("providers_ignored", None), + ("providers_order", None), + ("provider_sort", None), + ("provider_require_parameters", False), + ("provider_data_collection", ""), +) +_NOUS_PROVIDERS = frozenset({"nous", "nous-portal", "nousresearch"}) def _resolve_child_runtime( @@ -344,8 +336,9 @@ def _resolve_child_runtime( override_max_tokens: Optional[int], override_acp_command: Optional[str], override_acp_args: Optional[List[str]], -) -> _ChildRuntime: - """Resolve the child's credentials, transport and routing: config override > parent inherit. +) -> Dict[str, Any]: + """Resolve the child's credentials, transport and routing (config override > + parent inherit) as ``AIAgent`` keyword arguments. Rules that are easy to break: api_mode is re-derived (not inherited) when the child's provider differs from the parent's or is Nous Portal (dual-wire); @@ -355,22 +348,15 @@ def _resolve_child_runtime( """ effective_model = model or parent_agent.model effective_provider = override_provider or getattr(parent_agent, "provider", None) - effective_base_url = override_base_url or parent_agent.base_url - if not override_base_url: - effective_base_url = _inherit_parent_base_url(parent_agent, effective_base_url) - effective_api_key = override_api_key or parent_api_key - child_capabilities = _inherit_parent_capabilities( - parent_agent, override_provider, override_base_url - ) + effective_base_url = override_base_url or _inherit_parent_base_url(parent_agent, parent_agent.base_url) # api_mode: each provider has its own wire, so a different provider re-derives # (None) instead of inheriting (404s otherwise). Nous Portal is dual-wire # within one provider (anthropic/* → Messages, else chat_completions), so # same-provider inheritance would pin the child on the wrong wire — re-derive. _parent_provider = getattr(parent_agent, "provider", None) or "" - _effective_provider_norm = (effective_provider or "").strip().lower() if override_api_mode is not None: effective_api_mode = override_api_mode - elif _effective_provider_norm in {"nous", "nous-portal", "nousresearch"}: + elif (effective_provider or "").strip().lower() in _NOUS_PROVIDERS: from hermes_cli.providers import nous_api_mode effective_api_mode = nous_api_mode(effective_model) @@ -380,34 +366,23 @@ def _resolve_child_runtime( effective_api_mode = getattr(parent_agent, "api_mode", None) # A pinned transport that cannot run must fail the spawn loudly, never fall # back silently (delegate_task pre-validates; this covers direct callers). - if override_acp_command: - import shutil as _shutil - - if not _shutil.which(override_acp_command): - raise ValueError( - f"Pinned delegation command '{override_acp_command}' was not " - f"found on PATH. Install it or remove delegation.command from " - f"config.yaml." - ) - effective_acp_command = override_acp_command or getattr( - parent_agent, "acp_command", None + _require_pinned_command( + override_acp_command, + f"Pinned delegation command '{override_acp_command}' was not " + f"found on PATH. Install it or remove delegation.command from " + f"config.yaml.", ) + effective_acp_command = override_acp_command or getattr(parent_agent, "acp_command", None) effective_acp_args = list( - override_acp_args - if override_acp_args is not None - else (getattr(parent_agent, "acp_args", []) or []) + override_acp_args if override_acp_args is not None else (getattr(parent_agent, "acp_args", []) or []) ) - # A pinned provider must use direct API calls; inheriting the parent's ACP # transport would bypass the override credentials entirely. if override_provider and not override_acp_command: - effective_acp_command = None - effective_acp_args = [] - + effective_acp_command, effective_acp_args = None, [] if override_acp_command: # Forced ACP transport requires provider copilot-acp for run_agent to init the client. - effective_provider = "copilot-acp" - effective_api_mode = "chat_completions" + effective_provider, effective_api_mode = "copilot-acp", "chat_completions" # Reasoning: delegation.reasoning_effort > parent. Keep the raw value — a # YAML ``false`` must disable thinking, not coerce to "" and inherit. @@ -428,66 +403,32 @@ def _resolve_child_runtime( except Exception as exc: logger.debug("Could not load delegation reasoning_effort: %s", exc) - # Inherit the parent's fallback chain EXCEPT under a pinned provider: a - # mid-run 429/auth failure must not silently reroute the quiet child onto - # the parent's fallbacks. Predictability > liveness for explicit pins. - parent_fallback = ( - None - if override_provider - else (getattr(parent_agent, "_fallback_chain", None) or None) - ) - - # OpenRouter routing filters are inherited, but cleared under a pinned - # provider — parent filters (e.g. only=["Anthropic"]) would silently force - # the child back onto the parent's provider. - child_providers_allowed = getattr(parent_agent, "providers_allowed", None) - child_providers_ignored = getattr(parent_agent, "providers_ignored", None) - child_providers_order = getattr(parent_agent, "providers_order", None) - child_provider_sort = getattr(parent_agent, "provider_sort", None) - child_provider_require_parameters = getattr( - parent_agent, "provider_require_parameters", False - ) - child_provider_data_collection = getattr( - parent_agent, "provider_data_collection", None - ) or "" - child_openrouter_min_coding_score = getattr(parent_agent, "openrouter_min_coding_score", None) - if override_provider: - child_providers_allowed = None - child_providers_ignored = None - child_providers_order = None - child_provider_sort = None - child_provider_require_parameters = False - child_provider_data_collection = "" - # openrouter_min_coding_score stays inherited: model-gated, no-op elsewhere. - + kwargs: Dict[str, Any] = { + "base_url": effective_base_url, + "api_key": override_api_key or parent_api_key, + "model": effective_model, + "provider": effective_provider, + "capabilities": _inherit_parent_capabilities(parent_agent, override_provider, override_base_url), + "api_mode": effective_api_mode, + "acp_command": effective_acp_command, + "acp_args": effective_acp_args, + "reasoning_config": child_reasoning, + # Inherit the parent's fallback chain EXCEPT under a pinned provider: a + # mid-run 429/auth failure must not silently reroute the quiet child onto + # the parent's fallbacks. Predictability > liveness for explicit pins. + "fallback_model": None if override_provider else (getattr(parent_agent, "_fallback_chain", None) or None), + "openrouter_min_coding_score": getattr(parent_agent, "openrouter_min_coding_score", None), + } + for attr, pinned_default in _ROUTING_FILTER_DEFAULTS: + kwargs[attr] = pinned_default if override_provider else getattr(parent_agent, attr, pinned_default) + if not override_provider: + kwargs["provider_data_collection"] = kwargs["provider_data_collection"] or "" child_max_tokens = ( - override_max_tokens - if override_max_tokens is not None - else getattr(parent_agent, "max_tokens", None) + override_max_tokens if override_max_tokens is not None else getattr(parent_agent, "max_tokens", None) ) - child_optional_kwargs: Dict[str, Any] = {} if isinstance(child_max_tokens, int): - child_optional_kwargs["max_tokens"] = child_max_tokens - return _ChildRuntime( - model=effective_model, - provider=effective_provider, - base_url=effective_base_url, - api_key=effective_api_key, - api_mode=effective_api_mode, - capabilities=child_capabilities, - acp_command=effective_acp_command, - acp_args=effective_acp_args, - reasoning=child_reasoning, - fallback=parent_fallback, - providers_allowed=child_providers_allowed, - providers_ignored=child_providers_ignored, - providers_order=child_providers_order, - provider_sort=child_provider_sort, - provider_require_parameters=child_provider_require_parameters, - provider_data_collection=child_provider_data_collection, - openrouter_min_coding_score=child_openrouter_min_coding_score, - optional_kwargs=child_optional_kwargs, - ) + kwargs["max_tokens"] = child_max_tokens + return kwargs def _open_child_session_db(parent_agent) -> Any: @@ -515,7 +456,7 @@ def _open_child_session_db(parent_agent) -> Any: def _construct_child_agent( - rt: _ChildRuntime, + rt: Dict[str, Any], *, task_index: int, max_iterations: int, @@ -545,18 +486,9 @@ def _construct_child_agent( with delegated_child_context(): try: return AIAgent( - base_url=rt.base_url, - api_key=rt.api_key, - model=rt.model, - provider=rt.provider, - capabilities=rt.capabilities, - api_mode=rt.api_mode, - acp_command=rt.acp_command, - acp_args=rt.acp_args, + **rt, max_iterations=max_iterations, - reasoning_config=rt.reasoning, prefill_messages=getattr(parent_agent, "prefill_messages", None), - fallback_model=rt.fallback, enabled_toolsets=child_toolsets, disabled_toolsets=child_disabled_toolsets, quiet_mode=True, @@ -569,12 +501,6 @@ def _construct_child_agent( thinking_callback=child_thinking_cb, session_db=child_session_db, parent_session_id=getattr(parent_agent, "session_id", None), - providers_allowed=rt.providers_allowed, - providers_ignored=rt.providers_ignored, - providers_order=rt.providers_order, - provider_sort=rt.provider_sort, - provider_require_parameters=rt.provider_require_parameters, - provider_data_collection=rt.provider_data_collection, request_overrides=( # honored whenever set, incl. the inherit branch where # _resolve_delegation_credentials already merged OVER the parent's @@ -582,10 +508,8 @@ def _construct_child_agent( if override_request_overrides is not None else ({} if override_provider else dict(getattr(parent_agent, "request_overrides", {}) or {})) ), - openrouter_min_coding_score=rt.openrouter_min_coding_score, tool_progress_callback=child_progress_cb, iteration_budget=None, # fresh budget per subagent - **rt.optional_kwargs, ) except BaseException: if child_session_db is not None: @@ -735,7 +659,7 @@ def _build_child_agent( child._session_init_model_config["_delegate_from"] = parent_sid # Shared pool lets children rotate credentials on rate limits. - child_pool = _resolve_child_credential_pool(rt.provider, parent_agent, rt.base_url) + child_pool = _resolve_child_credential_pool(rt["provider"], parent_agent, rt["base_url"]) if child_pool is not None: child._credential_pool = child_pool diff --git a/tools/delegate_tool_child_run.py b/tools/delegate_tool_child_run.py index 5ed5d0864e..a95c2b10c4 100644 --- a/tools/delegate_tool_child_run.py +++ b/tools/delegate_tool_child_run.py @@ -93,10 +93,10 @@ def _detach_child(parent_agent: Any, child: Any) -> None: logger.debug("Could not remove child from active_children: %s", e) -def _signal_child_stop(child: Any) -> None: +def _signal_child_stop(child: Any, *reason: str) -> None: """Cooperative interrupt so the child's worker thread can exit cleanly.""" try: - if child is not None and not request_hard_interrupt(child) and hasattr(child, "_interrupt_requested"): + if child is not None and not request_hard_interrupt(child, *reason) and hasattr(child, "_interrupt_requested"): child._interrupt_requested = True except Exception: pass diff --git a/tools/delegate_tool_dispatch.py b/tools/delegate_tool_dispatch.py index 0aa6b738f7..84d92d5b85 100644 --- a/tools/delegate_tool_dispatch.py +++ b/tools/delegate_tool_dispatch.py @@ -12,8 +12,7 @@ import logging from concurrent.futures import FIRST_COMPLETED, wait as _cf_wait from typing import Any, Dict, List, Optional -from agent.interrupt_compat import request_hard_interrupt -from tools.delegate_tool_child_run import _fabricated_entry +from tools.delegate_tool_child_run import _detach_child, _fabricated_entry, _signal_child_stop from tools.delegate_tool_progress import ( SUBAGENT_FAILURE_STATUSES, _clean_error_text, @@ -228,24 +227,6 @@ def _resolve_async_session_key(parent_agent: Any, origin_ui_session_id: str) -> return session_key, origin_ui_session_id -def _detach_from_parent(parent_agent: Any, child_agents: List[Any]) -> None: - """Drop the children from the parent's interrupt-propagation list — the - batch's lifecycle is owned by the async registry now (_build_child_agent - attached them, correct for sync runs).""" - if not hasattr(parent_agent, "_active_children"): - return - lock = getattr(parent_agent, "_active_children_lock", None) - for c in child_agents: - try: - if lock: - with lock: - parent_agent._active_children.remove(c) - else: - parent_agent._active_children.remove(c) - except ValueError: - pass - - def _batch_progress_token(child_agents: List[Any]) -> tuple: """Progress token for the async registry's stale monitor. @@ -349,16 +330,16 @@ def _dispatch_background( session_key, origin_ui_session_id = _resolve_async_session_key(parent_agent, origin_ui_session_id) child_agents = [c for (_, _, c) in children] - _detach_from_parent(parent_agent, child_agents) + # The batch's lifecycle is owned by the async registry now: drop the children + # from the parent's interrupt-propagation list (_build_child_agent attached + # them, which is correct for sync runs). + for c in child_agents: + _detach_child(parent_agent, c) def _batch_interrupt(): # Cancellation path for the detached batch (owned by the async registry). for c in child_agents: - try: - if not request_hard_interrupt(c, "Async delegation cancelled") and hasattr(c, "_interrupt_requested"): - c._interrupt_requested = True - except Exception: - pass + _signal_child_stop(c, "Async delegation cancelled") goals = [t["goal"] for t in task_list] dispatch = dispatch_async_delegation_batch( diff --git a/tools/delegate_tool_registry.py b/tools/delegate_tool_registry.py index a6e1df794a..21d366eb5d 100644 --- a/tools/delegate_tool_registry.py +++ b/tools/delegate_tool_registry.py @@ -50,21 +50,10 @@ def get_subagent_attribution(task_id: Optional[str]) -> Optional[Dict[str, Any]] if not task_id or not isinstance(task_id, str): return None with _active_subagents_lock: - record = _active_subagents.get(task_id) - if record is not None: - return { - "subagent_id": task_id, - "goal": record.get("goal"), - "delegation_id": record.get("delegation_id"), - } - retained = _recent_subagents.get(task_id) - if retained is not None: - return { - "subagent_id": task_id, - "goal": retained.get("goal"), - "delegation_id": retained.get("delegation_id"), - } - return None + record = _active_subagents.get(task_id) or _recent_subagents.get(task_id) + if record is None: + return None + return {"subagent_id": task_id, "goal": record.get("goal"), "delegation_id": record.get("delegation_id")} def set_spawn_paused(paused: bool) -> bool: @@ -404,54 +393,41 @@ def _handle_control_action( "completion message). Use action='list' to see live children." ) - if action == "stop": - if interrupt_subagent(sid): - return json.dumps( - { - "action": "stop", - "subagent_id": sid, - "status": "interrupt_requested", - "note": ( - "The subagent stops at its next iteration boundary " - "(in-flight tool calls are asked to cancel). Its " - "partial result still re-enters the conversation as a " - "completion message — do not wait or poll." - ), - }, - ensure_ascii=False, - ) + if action == "steer" and not (message or "").strip(): return tool_error( - f"Could not interrupt '{sid}' — it likely finished in the last " - "moment. Its result arrives as a normal completion message." + "action='steer' requires a non-empty 'message' describing the " + "course correction." ) + outcome = _CONTROL_OUTCOMES.get(action) + if outcome is None: + return tool_error(f"Unknown action '{action}'. Use spawn, list, steer, or stop.") + status, note, failure = outcome + ok = interrupt_subagent(sid) if action == "stop" else steer_subagent(sid, message.strip()) + if ok: + return json.dumps({"action": action, "subagent_id": sid, "status": status, "note": note}, ensure_ascii=False) + return tool_error(failure.format(sid=sid)) - if action == "steer": - text = (message or "").strip() - if not text: - return tool_error( - "action='steer' requires a non-empty 'message' describing the " - "course correction." - ) - if steer_subagent(sid, text): - return json.dumps( - { - "action": "steer", - "subagent_id": sid, - "status": "queued", - "note": ( - "Steering text queued. The subagent sees it appended " - "to its next tool result — the current tool call is " - "never cut. If the child finishes before a delivery " - "boundary remains, the text is reported back as " - "missed_steer in its completion entry." - ), - }, - ensure_ascii=False, - ) - return tool_error( - f"Subagent '{sid}' is no longer accepting steering (finishing or " - "already finished). Its result arrives as a normal completion " - "message; re-delegate a follow-up task if more work is needed." - ) - return tool_error(f"Unknown action '{action}'. Use spawn, list, steer, or stop.") +# action -> (success status, success note, failure error template) +_CONTROL_OUTCOMES = { + "stop": ( + "interrupt_requested", + "The subagent stops at its next iteration boundary " + "(in-flight tool calls are asked to cancel). Its " + "partial result still re-enters the conversation as a " + "completion message — do not wait or poll.", + "Could not interrupt '{sid}' — it likely finished in the last " + "moment. Its result arrives as a normal completion message.", + ), + "steer": ( + "queued", + "Steering text queued. The subagent sees it appended " + "to its next tool result — the current tool call is " + "never cut. If the child finishes before a delivery " + "boundary remains, the text is reported back as " + "missed_steer in its completion entry.", + "Subagent '{sid}' is no longer accepting steering (finishing or " + "already finished). Its result arrives as a normal completion " + "message; re-delegate a follow-up task if more work is needed.", + ), +}