refactor(delegate): _knob() unifies config>env>default int/float knobs; timeout diagnostic sections -> _diag_sizes/_diag_threads + attr table; drop dead _handle_child_wait_failure re-import

This commit is contained in:
Teknium
2026-09-02 16:57:22 -07:00
parent 884c697291
commit e204008f10
9 changed files with 125 additions and 129 deletions
Binary file not shown.
Binary file not shown.
-1
View File
@@ -47,7 +47,6 @@ from tools.delegate_tool_child_run import ( # noqa: F401
_dump_subagent_timeout_diagnostic,
_emit_child_complete,
_fabricated_entry,
_handle_child_wait_failure,
_lease_child_credential,
_make_text_relay,
_merge_late_steer,
+83 -84
View File
@@ -112,6 +112,72 @@ def _format_thread_stack(frame: Any, indent: str) -> List[str]:
]
_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:
sys_prompt = getattr(child, "ephemeral_system_prompt", None) or getattr(child, "system_prompt", None) or ""
lines.append(f" system_prompt_bytes: {len(sys_prompt.encode('utf-8')) if isinstance(sys_prompt, str) else 'n/a'}")
lines.append(f" system_prompt_chars: {len(sys_prompt) if isinstance(sys_prompt, str) else 'n/a'}")
except Exception as exc:
lines.append(f" system_prompt: <error: {exc}>")
try:
tools_schema = getattr(child, "tools", None)
if tools_schema is not None:
lines.append(f" tool_schema_count: {len(tools_schema)}")
lines.append(f" tool_schema_bytes: {len(json.dumps(tools_schema, default=str).encode('utf-8'))}")
except Exception as exc:
lines.append(f" tool_schema: <error: {exc}>")
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
from a slow provider without the full picture."""
import sys as _sys
import threading as _threading
lines = ["## Worker thread stack at timeout"]
frames = _sys._current_frames()
if worker_thread is not None and worker_thread.is_alive():
worker_frame = frames.get(worker_thread.ident)
lines.extend(
_format_thread_stack(worker_frame, " ") if worker_frame is not None
else [" <worker frame not available>"]
)
elif worker_thread is None:
lines.append(" <no worker thread handle>")
else:
lines.append(" <worker thread already exited>")
lines += ["", "## All thread stacks at timeout"]
try:
frames = _sys._current_frames()
by_ident = {th.ident: th for th in _threading.enumerate() if th.ident}
worker_ident = worker_thread.ident if worker_thread else None
dumped = 0
for ident, frame in frames.items():
if ident == worker_ident:
continue # already dumped above
if dumped >= 40:
lines.append(f" <{len(frames) - dumped - 1} more threads omitted>")
break
th = by_ident.get(ident)
name = th.name if th else f"ident={ident}"
lines.append(f" --- {name}{' daemon' if (th and th.daemon) else ''} ---")
lines.extend(_format_thread_stack(frame, " "))
dumped += 1
except Exception as exc:
lines.append(f" <all-thread dump failed: {exc}>")
return lines
def _dump_subagent_timeout_diagnostic(
*,
child: Any,
@@ -130,8 +196,6 @@ def _dump_subagent_timeout_diagnostic(
try:
from hermes_constants import get_hermes_home
import datetime as _dt
import sys as _sys
import threading as _threading
logs_dir = get_hermes_home() / "logs"
try:
@@ -159,97 +223,32 @@ def _dump_subagent_timeout_diagnostic(
"",
"## Child config",
]
_w = lines.append
for attr in (
"model", "provider", "api_mode", "base_url", "max_iterations",
"quiet_mode", "skip_memory", "skip_context_files", "platform",
"_delegate_role", "_delegate_depth",
):
for attr in _DIAG_CHILD_ATTRS:
try:
_w(f" {attr}: {getattr(child, attr, None)!r}")
lines.append(f" {attr}: {getattr(child, attr, None)!r}")
except Exception:
_w(f" {attr}: <unreadable>")
_w("")
_w("## Toolsets")
_w(f" enabled_toolsets: {getattr(child, 'enabled_toolsets', None)!r}")
lines.append(f" {attr}: <unreadable>")
lines += ["", "## Toolsets", f" enabled_toolsets: {getattr(child, 'enabled_toolsets', None)!r}"]
tool_names = getattr(child, "valid_tool_names", None)
if tool_names:
_w(f" loaded tool count: {len(tool_names)}")
lines.append(f" loaded tool count: {len(tool_names)}")
try:
_w(f" loaded tools: {sorted(tool_names)}")
lines.append(f" loaded tools: {sorted(tool_names)}")
except Exception:
pass
_w("")
_w("## Prompt / schema sizes")
lines += [""] + _diag_sizes(child) + ["", "## Activity summary"]
try:
sys_prompt = getattr(child, "ephemeral_system_prompt", None) or getattr(child, "system_prompt", None) or ""
_w(f" system_prompt_bytes: {len(sys_prompt.encode('utf-8')) if isinstance(sys_prompt, str) else 'n/a'}")
_w(f" system_prompt_chars: {len(sys_prompt) if isinstance(sys_prompt, str) else 'n/a'}")
lines += [f" {k}: {v!r}" for k, v in child.get_activity_summary().items()]
except Exception as exc:
_w(f" system_prompt: <error: {exc}>")
try:
tools_schema = getattr(child, "tools", None)
if tools_schema is not None:
_w(f" tool_schema_count: {len(tools_schema)}")
_w(f" tool_schema_bytes: {len(json.dumps(tools_schema, default=str).encode('utf-8'))}")
except Exception as exc:
_w(f" tool_schema: <error: {exc}>")
_w("")
_w("## Activity summary")
try:
for k, v in child.get_activity_summary().items():
_w(f" {k}: {v!r}")
except Exception as exc:
_w(f" <get_activity_summary failed: {exc}>")
_w("")
_w("## Worker thread stack at timeout")
frames = _sys._current_frames()
if worker_thread is not None and worker_thread.is_alive():
worker_frame = frames.get(worker_thread.ident)
lines.extend(
_format_thread_stack(worker_frame, " ") if worker_frame is not None
else [" <worker frame not available>"]
)
elif worker_thread is None:
_w(" <no worker thread handle>")
else:
_w(" <worker thread already exited>")
_w("")
# All other live threads (bounded to 40): the worker is often parked on
# a helper thread, so a pre-HTTP wedge is indistinguishable from a slow
# provider without the full picture.
_w("## All thread stacks at timeout")
try:
frames = _sys._current_frames()
by_ident = {th.ident: th for th in _threading.enumerate() if th.ident}
worker_ident = worker_thread.ident if worker_thread else None
dumped = 0
for ident, frame in frames.items():
if ident == worker_ident:
continue # already dumped above
if dumped >= 40:
_w(f" <{len(frames) - dumped - 1} more threads omitted>")
break
th = by_ident.get(ident)
name = th.name if th else f"ident={ident}"
_w(f" --- {name}{' daemon' if (th and th.daemon) else ''} ---")
lines.extend(_format_thread_stack(frame, " "))
dumped += 1
except Exception as exc:
_w(f" <all-thread dump failed: {exc}>")
_w("")
_w("## Notes")
_w(" This file is written ONLY when a subagent times out with 0 API calls.")
_w(" 0-API-call timeouts mean the child never reached its first LLM request.")
_w(" Common causes: oversized prompt rejected by provider, transport hang,")
_w(" credential resolution stuck. See issue #14726 for context.")
lines.append(f" <get_activity_summary failed: {exc}>")
lines += [""] + _diag_threads(worker_thread) + [
"",
"## Notes",
" This file is written ONLY when a subagent times out with 0 API calls.",
" 0-API-call timeouts mean the child never reached its first LLM request.",
" Common causes: oversized prompt rejected by provider, transport hang,",
" credential resolution stuck. See issue #14726 for context.",
]
dump_path.write_text("\n".join(lines), encoding="utf-8")
return str(dump_path)
except Exception as exc:
+42 -44
View File
@@ -71,40 +71,49 @@ def _get_subagent_approval_callback():
return _subagent_auto_deny
def _knob(key: str, env_var: str, parse, default, warn_invalid):
"""delegation.<key> > <env_var> > default. A config value that fails ``parse``
calls ``warn_invalid(value)`` and yields the default; an env value that fails
is silently ignored."""
val = _cfg().get(key)
if val is not None:
try:
return parse(val)
except (TypeError, ValueError):
warn_invalid(val)
return default
env_val = os.getenv(env_var)
if env_val:
try:
return parse(env_val)
except (TypeError, ValueError):
pass
return default
def _get_max_concurrent_children() -> int:
"""delegation.max_concurrent_children > DELEGATION_MAX_CONCURRENT_CHILDREN env > 10.
Floor of 1 is the only bound enforced; there is no ceiling.
"""
val = _cfg().get("max_concurrent_children")
if val is not None:
try:
result = max(1, int(val))
except (TypeError, ValueError):
result = _knob(
"max_concurrent_children", "DELEGATION_MAX_CONCURRENT_CHILDREN", lambda v: max(1, int(v)),
_DEFAULT_MAX_CONCURRENT_CHILDREN,
lambda val: logger.warning(
"delegation.max_concurrent_children=%r is not a valid integer; using default %d",
val, _DEFAULT_MAX_CONCURRENT_CHILDREN,
),
)
if result > 10 and _cfg().get("max_concurrent_children") is not None:
global _HIGH_CONCURRENCY_WARNED
if not _HIGH_CONCURRENCY_WARNED:
_HIGH_CONCURRENCY_WARNED = True
logger.warning(
"delegation.max_concurrent_children=%r is not a valid integer; "
"using default %d",
val,
_DEFAULT_MAX_CONCURRENT_CHILDREN,
"delegation.max_concurrent_children=%d: each child consumes API tokens "
"independently. High values multiply cost linearly.",
result,
)
return _DEFAULT_MAX_CONCURRENT_CHILDREN
if result > 10:
global _HIGH_CONCURRENCY_WARNED
if not _HIGH_CONCURRENCY_WARNED:
_HIGH_CONCURRENCY_WARNED = True
logger.warning(
"delegation.max_concurrent_children=%d: each child consumes API tokens "
"independently. High values multiply cost linearly.",
result,
)
return result
env_val = os.getenv("DELEGATION_MAX_CONCURRENT_CHILDREN")
if env_val:
try:
return max(1, int(env_val))
except (TypeError, ValueError):
return _DEFAULT_MAX_CONCURRENT_CHILDREN
return _DEFAULT_MAX_CONCURRENT_CHILDREN
return result
def _get_worktree_isolation() -> bool:
@@ -151,23 +160,12 @@ def _get_child_timeout() -> Optional[float]:
staleness monitor. delegation.child_timeout_seconds > 0 opts in (floor 30 s);
0 or negative disables. Env fallback: DELEGATION_CHILD_TIMEOUT_SECONDS.
"""
val = _cfg().get("child_timeout_seconds")
if val is not None:
try:
return _parse_timeout(val)
except (TypeError, ValueError):
logger.warning(
"delegation.child_timeout_seconds=%r is not a valid number; "
"using default (no timeout)",
val,
)
env_val = os.getenv("DELEGATION_CHILD_TIMEOUT_SECONDS")
if env_val:
try:
return _parse_timeout(env_val)
except (TypeError, ValueError):
pass
return DEFAULT_CHILD_TIMEOUT
return _knob(
"child_timeout_seconds", "DELEGATION_CHILD_TIMEOUT_SECONDS", _parse_timeout, DEFAULT_CHILD_TIMEOUT,
lambda val: logger.warning(
"delegation.child_timeout_seconds=%r is not a valid number; using default (no timeout)", val,
),
)
def _get_max_spawn_depth() -> int: