diff --git a/MagicMock/mock._session_db.db_path/126670844021456 b/MagicMock/mock._session_db.db_path/126670844021456 new file mode 100644 index 0000000000..d0a247e7aa Binary files /dev/null and b/MagicMock/mock._session_db.db_path/126670844021456 differ diff --git a/MagicMock/mock._session_db.db_path/126670844021456.fts_rebuild.lock b/MagicMock/mock._session_db.db_path/126670844021456.fts_rebuild.lock new file mode 100644 index 0000000000..e69de29bb2 diff --git a/MagicMock/mock._session_db.db_path/126670844021456.quarantine.lock b/MagicMock/mock._session_db.db_path/126670844021456.quarantine.lock new file mode 100644 index 0000000000..e69de29bb2 diff --git a/MagicMock/mock._session_db.db_path/126670871778000 b/MagicMock/mock._session_db.db_path/126670871778000 new file mode 100644 index 0000000000..4e547c7618 Binary files /dev/null and b/MagicMock/mock._session_db.db_path/126670871778000 differ diff --git a/MagicMock/mock._session_db.db_path/126670871778000.fts_rebuild.lock b/MagicMock/mock._session_db.db_path/126670871778000.fts_rebuild.lock new file mode 100644 index 0000000000..e69de29bb2 diff --git a/MagicMock/mock._session_db.db_path/126670871778000.quarantine.lock b/MagicMock/mock._session_db.db_path/126670871778000.quarantine.lock new file mode 100644 index 0000000000..e69de29bb2 diff --git a/tools/delegate_tool.py b/tools/delegate_tool.py index e24f6a3c88..8dbac85c3b 100644 --- a/tools/delegate_tool.py +++ b/tools/delegate_tool.py @@ -186,9 +186,7 @@ def _expand_parent_toolsets(parent_toolsets: set) -> set: return expanded -def _preserve_parent_mcp_toolsets( - child_toolsets: List[str], parent_toolsets: set[str] -) -> List[str]: +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) for toolset_name in sorted(parent_toolsets): @@ -276,9 +274,7 @@ def _resolve_child_toolsets( expanded_parent = _expand_parent_toolsets(parent_toolsets) child_toolsets = [t for t in toolsets if t in expanded_parent] if _get_inherit_mcp_toolsets(): - child_toolsets = _preserve_parent_mcp_toolsets( - child_toolsets, parent_toolsets - ) + child_toolsets = _preserve_parent_mcp_toolsets(child_toolsets, parent_toolsets) child_toolsets = _strip_blocked_tools(child_toolsets) elif parent_agent and parent_enabled is not None: child_toolsets = _strip_blocked_tools(parent_enabled) @@ -293,9 +289,7 @@ def _resolve_child_toolsets( else: inherited_disabled = [] if effective_role == "orchestrator": - inherited_disabled = [ - name for name in inherited_disabled if name != "delegation" - ] + inherited_disabled = [name for name in inherited_disabled if name != "delegation"] child_disabled_toolsets = list( dict.fromkeys( inherited_disabled + _blocked_toolsets_for_role(effective_role) + ["kanban"] @@ -395,10 +389,7 @@ def _resolve_child_runtime( if parsed is not None: child_reasoning = parsed else: - logger.warning( - "Unknown delegation.reasoning_effort '%s', inheriting parent level", - delegation_effort, - ) + logger.warning("Unknown delegation.reasoning_effort '%s', inheriting parent level", delegation_effort) except Exception as exc: logger.debug("Could not load delegation reasoning_effort: %s", exc) @@ -447,10 +438,7 @@ def _open_child_session_db(parent_agent) -> Any: _parent_db_path = getattr(parent_session_db, "db_path", None) return get_shared_session_db(_parent_db_path) if _parent_db_path is not None else get_shared_session_db() except Exception: - logger.debug( - "subagent: failed to open dedicated SessionDB; child persistence disabled", - exc_info=True, - ) + logger.debug("subagent: failed to open dedicated SessionDB; child persistence disabled", exc_info=True) return None @@ -737,9 +725,7 @@ def _run_single_child( _child_close_deferred = failure.close_deferred return failure.entry - schema = _validate_child_output_schema( - child, result, task_index, ws.child_task_id, _relay_child_text - ) + schema = _validate_child_output_schema(child, result, task_index, ws.child_task_id, _relay_child_text) _merge_late_steer(result, _subagent_id, child) # Flush any remaining batched progress to gateway @@ -757,9 +743,7 @@ def _run_single_child( return entry except Exception as exc: - _late_pending_steer = ( - _close_subagent_steering(_subagent_id, child) if _subagent_id else None - ) + _late_pending_steer = (_close_subagent_steering(_subagent_id, child) if _subagent_id else None) duration = round(time.monotonic() - child_start, 2) logging.exception(f"[subagent-{task_index}] failed") _safe_progress( @@ -787,9 +771,7 @@ def _run_single_child( ) -def _recover_tasks_from_json_string( - tasks: Any, -) -> tuple[Optional[List[Dict[str, Any]]], Optional[str]]: +def _recover_tasks_from_json_string(tasks: Any) -> tuple[Optional[List[Dict[str, Any]]], Optional[str]]: if not isinstance(tasks, str): return None, None raw = tasks.strip() @@ -803,10 +785,7 @@ def _recover_tasks_from_json_string( f"that could not be parsed as JSON ({exc.msg})." ) if not isinstance(parsed, list): - return None, ( - f"tasks must be a JSON array of task objects; parsed " - f"{type(parsed).__name__} instead." - ) + return None, (f"tasks must be a JSON array of task objects; parsed " f"{type(parsed).__name__} instead.") return parsed, None @@ -1280,10 +1259,7 @@ def _build_top_level_description() -> str: "derives this from depth automatically.\n" ) else: - restrictions_rule = ( - "- Children cannot call delegate_task, clarify, memory, or " - "cronjob.\n" - ) + restrictions_rule = ("- Children cannot call delegate_task, clarify, memory, or " "cronjob.\n") return ( "Spawn subagents in isolated contexts; each gets its own conversation, " @@ -1341,19 +1317,14 @@ def _build_dynamic_schema_overrides() -> dict: get_definitions() pass rewrites the description fields to the user's actual limits. """ - overrides_params = { - **DELEGATE_TASK_SCHEMA["parameters"], - } + overrides_params = {**DELEGATE_TASK_SCHEMA["parameters"]} # Deep-copy properties so we don't mutate the static schema dict. overrides_params["properties"] = { k: dict(v) for k, v in DELEGATE_TASK_SCHEMA["parameters"]["properties"].items() } overrides_params["properties"]["tasks"]["description"] = _build_tasks_param_description() - return { - "description": _build_top_level_description(), - "parameters": overrides_params, - } + return {"description": _build_top_level_description(), "parameters": overrides_params} DELEGATE_TASK_SCHEMA = { @@ -1496,11 +1467,7 @@ def _strip_model_hidden_task_fields(tasks: Any) -> Any: if not isinstance(task, dict): stripped_tasks.append(task) continue - stripped = { - key: value - for key, value in task.items() - if key not in _MODEL_HIDDEN_TASK_FIELDS - } + stripped = {key: value for key, value in task.items() if key not in _MODEL_HIDDEN_TASK_FIELDS} changed = changed or len(stripped) != len(task) stripped_tasks.append(stripped) return stripped_tasks if changed else tasks diff --git a/tools/delegate_tool_child_run.py b/tools/delegate_tool_child_run.py index d2b4b3a5d0..adc58821b7 100644 --- a/tools/delegate_tool_child_run.py +++ b/tools/delegate_tool_child_run.py @@ -50,10 +50,7 @@ def _append_missed_steer(entry: Dict[str, Any], late_steer: Optional[str]) -> No """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}]" - ) + entry["error"] += (" [steer did not land before the subagent stopped: " f"{late_steer}]") def _close_child(child: Any, log_message: str) -> None: @@ -263,11 +260,7 @@ def _start_heartbeat(child: Any, parent_agent: Any, task_index: int) -> tuple: ``try`` so a failed ``start()`` (OS thread exhaustion) leaves ``ident`` None and the finally-path join can be skipped safely. """ - from tools.delegate_tool import ( - _HEARTBEAT_INTERVAL, - _HEARTBEAT_STALE_CYCLES_IDLE, - _HEARTBEAT_STALE_CYCLES_IN_TOOL, - ) + from tools.delegate_tool import (_HEARTBEAT_INTERVAL, _HEARTBEAT_STALE_CYCLES_IDLE, _HEARTBEAT_STALE_CYCLES_IN_TOOL) _heartbeat_stop = threading.Event() # Stale detection: a cycle counts as stale when (tool, iteration, @@ -593,9 +586,7 @@ def _await_child( back to ``input()`` and deadlock the parent TUI (deny vs approve follows delegation.subagent_auto_approve). """ - from tools.delegate_tool import ( - _get_child_timeout, _get_subagent_approval_callback, _set_subagent_approval_cb, - ) + from tools.delegate_tool import (_get_child_timeout, _get_subagent_approval_callback, _set_subagent_approval_cb) from tools.daemon_pool import DaemonThreadPoolExecutor child_timeout = _get_child_timeout() @@ -610,9 +601,7 @@ def _await_child( from agent.delegation_context import delegated_child_context with delegated_child_context(str(getattr(child, "session_id", "") or "")): - return child.run_conversation( - user_message=goal, task_id=ws.child_task_id, stream_callback=relay_child_text, - ) + return child.run_conversation(user_message=goal, task_id=ws.child_task_id, stream_callback=relay_child_text) future = executor.submit(contextvars.copy_context().run, _run_with_thread_capture) try: @@ -702,11 +691,7 @@ def _handle_child_wait_failure( goal=goal, ) if diagnostic_path: - logger.warning( - "Subagent %d 0-API-call timeout — diagnostic written to %s", - task_index, - diagnostic_path, - ) + logger.warning("Subagent %d 0-API-call timeout — diagnostic written to %s", task_index, diagnostic_path) status = "timeout" if is_timeout else "error" _safe_progress( @@ -949,10 +934,7 @@ def _build_result_entry( _missed_steer = result.get("pending_steer") if isinstance(_missed_steer, str) and _missed_steer.strip(): entry["missed_steer"] = _missed_steer - _miss_note = ( - "[steer did not land — the subagent finished before it could " - f"be delivered: {_missed_steer}]" - ) + _miss_note = ("[steer did not land — the subagent finished before it could " f"be delivered: {_missed_steer}]") entry["summary"] = f"{summary}\n\n{_miss_note}" if summary else _miss_note return entry diff --git a/tools/delegate_tool_config.py b/tools/delegate_tool_config.py index 9a27b31d12..935a1a4c3c 100644 --- a/tools/delegate_tool_config.py +++ b/tools/delegate_tool_config.py @@ -57,10 +57,7 @@ def _subagent_auto_deny(command: str, description: str, **kwargs) -> str: 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, - ) + logger.warning("Subagent auto-approved dangerous command: %s (%s)", command, description) return "once" @@ -180,20 +177,11 @@ def _get_max_spawn_depth() -> int: try: ival = int(val) except (TypeError, ValueError): - logger.warning( - "delegation.max_spawn_depth=%r is not a valid integer; " "using default %d", - val, - MAX_DEPTH, - ) + logger.warning("delegation.max_spawn_depth=%r is not a valid integer; " "using default %d", val, MAX_DEPTH) return MAX_DEPTH floored = max(_MIN_SPAWN_DEPTH, ival) if floored != ival: - logger.warning( - "delegation.max_spawn_depth=%d below floor %d; using %d", - ival, - _MIN_SPAWN_DEPTH, - floored, - ) + logger.warning("delegation.max_spawn_depth=%d below floor %d; using %d", ival, _MIN_SPAWN_DEPTH, floored) return floored @@ -217,9 +205,7 @@ 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]: +def _inherit_parent_capabilities(parent_agent, override_provider, override_base_url) -> Optional[dict]: """Parent's endpoint-trust capability map for a child, or None. ``agent.capabilities`` is a trust decision scoped to one provider+endpoint: @@ -231,11 +217,7 @@ def _inherit_parent_capabilities( parent_caps = getattr(parent_agent, "capabilities", None) if not isinstance(parent_caps, dict): return None - return { - key: value - for key, value in parent_caps.items() - if isinstance(key, str) and isinstance(value, bool) - } + 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]: @@ -294,9 +276,7 @@ def _resolve_child_credential_pool( # Unregistered endpoint (no custom_providers entry): keep the # child's fixed credential rather than inherit the parent's. return None - parent_key = get_custom_provider_pool_key( - getattr(parent_agent, "base_url", None) - ) + parent_key = get_custom_provider_pool_key(getattr(parent_agent, "base_url", None)) if ( parent_pool is not None and parent_provider == "custom" @@ -318,11 +298,7 @@ def _resolve_child_credential_pool( try: return _loaded_pool(effective_provider) except Exception as exc: - logger.debug( - "Could not load credential pool for child provider '%s': %s", - effective_provider, - exc, - ) + logger.debug("Could not load credential pool for child provider '%s': %s", effective_provider, exc) return None @@ -406,9 +382,7 @@ def _direct_endpoint_credentials(cfg_values: dict, explicit_request_overrides) - try: from hermes_cli.runtime_provider import resolve_runtime_provider - runtime = resolve_runtime_provider( - requested=configured_provider, target_model=configured_model - ) + runtime = resolve_runtime_provider(requested=configured_provider, target_model=configured_model) request_overrides = dict(runtime.get("request_overrides") or {}) or None max_output_tokens = runtime.get("max_output_tokens") except Exception as exc: diff --git a/tools/delegate_tool_dispatch.py b/tools/delegate_tool_dispatch.py index 7993c795b6..b535b65b63 100644 --- a/tools/delegate_tool_dispatch.py +++ b/tools/delegate_tool_dispatch.py @@ -107,10 +107,7 @@ def _run_children_parallel( pending = set(futures.keys()) while pending: - if ( - honor_parent_interrupt - and getattr(parent_agent, "_interrupt_requested", False) is True - ): + if (honor_parent_interrupt and getattr(parent_agent, "_interrupt_requested", False) is True): # Parent interrupted — collect whatever finished and abandon the # rest (children already got the interrupt signal). for f in pending: diff --git a/tools/delegate_tool_progress.py b/tools/delegate_tool_progress.py index 27a1b3ba72..1e640990fe 100644 --- a/tools/delegate_tool_progress.py +++ b/tools/delegate_tool_progress.py @@ -111,9 +111,7 @@ _EVENT_HANDLERS: Dict[Any, Optional[str]] = { def _normalize_event(event_type: Any) -> Any: """Lifecycle string / DelegateEvent / legacy string / ``delegate.*`` string → dispatch key; None for unknown events.""" - if isinstance(event_type, DelegateEvent) or ( - isinstance(event_type, str) and event_type in _LIFECYCLE_EVENTS - ): + if isinstance(event_type, DelegateEvent) or (isinstance(event_type, str) and event_type in _LIFECYCLE_EVENTS): return event_type event = _LEGACY_EVENT_MAP.get(event_type) if event is not None: @@ -139,11 +137,7 @@ def _build_child_system_prompt( OpenClaw's buildSubagentSystemPrompt); its depth note is literal truth grounded in the passed config so the LLM can't confabulate nesting. """ - parts = [ - "You are a focused subagent working on a specific delegated task.", - "", - f"YOUR TASK:\n{goal}", - ] + parts = ["You are a focused subagent working on a specific delegated task.", "", f"YOUR TASK:\n{goal}"] if context and context.strip(): parts.append(f"\nCONTEXT:\n{context}") if workspace_path and str(workspace_path).strip(): @@ -162,13 +156,9 @@ def _build_child_system_prompt( try: from agent.prompt_builder import build_context_files_prompt - _ctx_files = build_context_files_prompt( - cwd=str(workspace_path), skip_soul=True - ) + _ctx_files = build_context_files_prompt(cwd=str(workspace_path), skip_soul=True) except Exception: - logger.debug( - "subagent: workspace context-files load failed", exc_info=True - ) + logger.debug("subagent: workspace context-files load failed", exc_info=True) _ctx_files = "" if _ctx_files.strip(): parts.append( @@ -338,9 +328,7 @@ class _ChildProgressRelay: return _batch_prefix(deleg, self.task_index, self.task_count) def _identity_kwargs(self) -> Dict[str, Any]: - kw: Dict[str, Any] = { - "task_index": self.task_index, "task_count": self.task_count, "goal": self.goal_label, - } + kw: Dict[str, Any] = {"task_index": self.task_index, "task_count": self.task_count, "goal": self.goal_label} for key in ("subagent_id", "parent_id", "depth", "model"): if getattr(self, key) is not None: kw[key] = getattr(self, key) diff --git a/tools/delegate_tool_registry.py b/tools/delegate_tool_registry.py index df6d606bb4..d066f6d131 100644 --- a/tools/delegate_tool_registry.py +++ b/tools/delegate_tool_registry.py @@ -198,9 +198,7 @@ def steer_subagent( return False -def _capture_gateway_steer_authority( - owner_session_id: Optional[str], -) -> tuple[Any, Any]: +def _capture_gateway_steer_authority(owner_session_id: Optional[str]) -> tuple[Any, Any]: """Capture exact request transport + live session generation, if any. This is intentionally an in-process bridge, not a serializable capability. @@ -316,12 +314,7 @@ def _owns_subagent_record(record: Dict[str, Any], parent_agent: Any) -> bool: } -def _handle_control_action( - action: str, - subagent_id: Optional[str], - message: Optional[str], - parent_agent: Any, -) -> str: +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. Runs in-turn (never backgrounded) and only over subagents descended from @@ -353,11 +346,7 @@ def _handle_control_action( "live_transcript": getattr(agent, "_live_transcript_path", None), } ) - payload: Dict[str, Any] = { - "action": "list", - "count": len(entries), - "subagents": entries, - } + payload: Dict[str, Any] = {"action": "list", "count": len(entries), "subagents": entries} if not entries: payload["note"] = ( "No live subagents right now. Children that already finished " @@ -383,10 +372,7 @@ def _handle_control_action( ) if action == "steer" and not (message or "").strip(): - return tool_error( - "action='steer' requires a non-empty 'message' describing the " - "course correction." - ) + return tool_error("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.") diff --git a/tools/delegate_tool_results.py b/tools/delegate_tool_results.py index f2b13ec278..b1bcc891ce 100644 --- a/tools/delegate_tool_results.py +++ b/tools/delegate_tool_results.py @@ -116,10 +116,7 @@ _TOOL_INPUT_URL_KEYS = frozenset({"endpoint", "url", "urls"}) def _sanitize_tool_target(key: str, value: Any) -> Any: """Keep bounded side-effect targets while dropping URL secrets.""" if isinstance(value, list): - cleaned = [ - item for item in (_sanitize_tool_target(key, item) for item in value[:16]) - if item is not None - ] + cleaned = [item for item in (_sanitize_tool_target(key, item) for item in value[:16]) if item is not None] return cleaned or None if not isinstance(value, str) or not value: return None @@ -292,9 +289,7 @@ def _spill_summary_to_file(task_index: int, summary: str) -> Optional[str]: return None -def _trim_summary_with_footer( - summary: str, cap: int, task_index: int -) -> tuple[str, Optional[str]]: +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. Mirrors web_extract's ``_truncate_with_footer``: keep a head+tail window @@ -402,9 +397,7 @@ def _apply_summary_budget(results: List[Dict[str, Any]], parent_agent) -> None: compression/429 death spiral. """ from tools.delegate_tool import _load_config - summaries = [ - r for r in results if isinstance(r, dict) and isinstance(r.get("summary"), str) and r["summary"] - ] + summaries = [r for r in results if isinstance(r, dict) and isinstance(r.get("summary"), str) and r["summary"]] if not summaries: return @@ -427,9 +420,7 @@ def _apply_summary_budget(results: List[Dict[str, Any]], parent_agent) -> None: if len(summary) <= cap: continue original_len = len(summary) - model_text, spill_path = _trim_summary_with_footer( - summary, cap, entry.get("task_index", -1) - ) + model_text, spill_path = _trim_summary_with_footer(summary, cap, entry.get("task_index", -1)) entry["summary"] = model_text entry["summary_truncated"] = True if spill_path: @@ -567,12 +558,7 @@ def _finalize_child_results( _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]: +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 result = _run_single_child(task_index, goal, child, parent_agent)