refactor(delegate): join wrapped logical lines that fit in 120 cols (AST-identical)

This commit is contained in:
Teknium
2026-09-02 16:59:26 -07:00
parent e204008f10
commit 016f31375d
13 changed files with 42 additions and 162 deletions
Binary file not shown.
Binary file not shown.
+13 -46
View File
@@ -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
+6 -24
View File
@@ -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
+8 -34
View File
@@ -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:
+1 -4
View File
@@ -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:
+5 -17
View File
@@ -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)
+4 -18
View File
@@ -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.")
+5 -19
View File
@@ -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)