perf(agents): run per-child timers on one shared scheduler thread
A fan-out of N in-process subagents used to add one sleeping daemon thread per delegated child (delegate heartbeat, 30s) and one or two per active turn (durable turn-lease refresher; turn-liveness watchdog). A profiled session with ~130 children was carrying ~1000 threads. All of these timers now run on a single process-wide daemon thread. - agent/periodic_scheduler.py (new): heap-ordered periodic scheduler on one Condition-driven daemon thread. schedule(fn, interval) -> handle; handle.cancel(wait=) blocks for an in-flight run like the old join. A callback returning False stops itself; a raising callback is logged at debug and rescheduled, so one bad timer cannot kill the rest. - tools/delegate_tool.py: _heartbeat_loop body -> _heartbeat_tick, scheduled at _HEARTBEAT_INTERVAL; stale-cycle closure state and idle/in-tool thresholds unchanged; cancel(wait=5) in finally where the stop-event + join(5) lived. - run_agent.py: _refresh_durable_turn_lease body scheduled at _lease_refresh_interval; lease-lost / refresh-error interrupt paths and the stop-event fencing are unchanged; the join(timeout=1.0) is now cancel(wait=1.0) so the interrupt clear still runs after any in-flight tick. - agent/turn_liveness.py: TurnLivenessWatchdog.make_thread/start -> schedule(); the poll body is _tick(), same sampling state machine. Bench (evals/fanout_resource_bench.py, 30 children / 10 worktrees, ok=30/30 both): peak threads 168 -> 132. At peak the old tree held 30 "Thread-N (_heartbeat_loop)" threads; the new one holds zero plus one "hermes-periodic-scheduler".
This commit is contained in:
@@ -0,0 +1,118 @@
|
||||
"""One process-wide timer thread for periodic maintenance callbacks.
|
||||
|
||||
Replaces the per-child ``while not stop.wait(interval): body()`` daemon
|
||||
threads (delegate heartbeat, durable turn-lease refresher, turn-liveness
|
||||
watchdog). With ~130 in-process subagents those added 2-3 sleeping OS
|
||||
threads per child; this module runs every periodic body on ONE daemon
|
||||
thread ordered by a heap of due times.
|
||||
|
||||
Semantics match the loop they replace: the first call happens ``interval``
|
||||
seconds after :func:`schedule`, and each following call ``interval`` seconds
|
||||
after the previous body *returned* (drift-free wrt. body duration was never
|
||||
a property of the old loops either). A body that returns ``False`` stops
|
||||
itself; a body that raises is logged at debug and rescheduled — one bad
|
||||
callback must never kill the shared thread.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import heapq
|
||||
import itertools
|
||||
import logging
|
||||
import threading
|
||||
import time
|
||||
from typing import Callable, Optional
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
_THREAD_NAME = "hermes-periodic-scheduler"
|
||||
|
||||
|
||||
class ScheduledHandle:
|
||||
"""Cancel token for one scheduled periodic callback."""
|
||||
|
||||
__slots__ = ("_fn", "_interval", "_cancelled", "_scheduler")
|
||||
|
||||
def __init__(self, scheduler: "PeriodicScheduler", fn: Callable[[], object], interval: float):
|
||||
self._scheduler = scheduler
|
||||
self._fn = fn
|
||||
self._interval = interval
|
||||
self._cancelled = False
|
||||
|
||||
@property
|
||||
def cancelled(self) -> bool:
|
||||
return self._cancelled
|
||||
|
||||
def cancel(self, wait: Optional[float] = None) -> None:
|
||||
"""Stop future runs. ``wait`` (seconds) additionally blocks until an
|
||||
in-flight run of this callback finishes — the analogue of
|
||||
``thread.join(timeout=wait)`` on the old per-child thread."""
|
||||
self._scheduler._cancel(self, wait)
|
||||
|
||||
|
||||
class PeriodicScheduler:
|
||||
def __init__(self) -> None:
|
||||
self._cond = threading.Condition()
|
||||
self._heap: list = [] # (due, seq, handle)
|
||||
self._seq = itertools.count()
|
||||
self._thread: Optional[threading.Thread] = None
|
||||
self._running: Optional[ScheduledHandle] = None
|
||||
|
||||
def schedule(self, fn: Callable[[], object], interval: float) -> ScheduledHandle:
|
||||
handle = ScheduledHandle(self, fn, float(interval))
|
||||
with self._cond:
|
||||
heapq.heappush(self._heap, (time.monotonic() + handle._interval, next(self._seq), handle))
|
||||
if self._thread is None or not self._thread.is_alive():
|
||||
self._thread = threading.Thread(target=self._run, name=_THREAD_NAME, daemon=True)
|
||||
self._thread.start()
|
||||
self._cond.notify()
|
||||
return handle
|
||||
|
||||
def _cancel(self, handle: ScheduledHandle, wait: Optional[float]) -> None:
|
||||
with self._cond:
|
||||
handle._cancelled = True
|
||||
self._cond.notify()
|
||||
if wait and self._running is handle and threading.current_thread() is not self._thread:
|
||||
self._cond.wait_for(lambda: self._running is not handle, timeout=wait)
|
||||
|
||||
def _run(self) -> None:
|
||||
while True:
|
||||
with self._cond:
|
||||
while True:
|
||||
if not self._heap:
|
||||
self._cond.wait()
|
||||
continue
|
||||
due, _, handle = self._heap[0]
|
||||
if handle._cancelled:
|
||||
heapq.heappop(self._heap)
|
||||
continue
|
||||
delay = due - time.monotonic()
|
||||
if delay > 0:
|
||||
self._cond.wait(delay)
|
||||
continue
|
||||
heapq.heappop(self._heap)
|
||||
self._running = handle
|
||||
break
|
||||
stop = False
|
||||
try:
|
||||
stop = handle._fn() is False
|
||||
except Exception:
|
||||
logger.debug("periodic callback %r raised", handle._fn, exc_info=True)
|
||||
with self._cond:
|
||||
self._running = None
|
||||
if stop:
|
||||
handle._cancelled = True
|
||||
elif not handle._cancelled:
|
||||
heapq.heappush(
|
||||
self._heap,
|
||||
(time.monotonic() + handle._interval, next(self._seq), handle),
|
||||
)
|
||||
self._cond.notify_all()
|
||||
|
||||
|
||||
_DEFAULT = PeriodicScheduler()
|
||||
|
||||
|
||||
def schedule(fn: Callable[[], object], interval: float) -> ScheduledHandle:
|
||||
"""Run ``fn()`` every ``interval`` seconds on the shared scheduler thread."""
|
||||
return _DEFAULT.schedule(fn, interval)
|
||||
+44
-49
@@ -157,7 +157,8 @@ def resolve_turn_liveness_settings(
|
||||
|
||||
|
||||
class TurnLivenessWatchdog:
|
||||
"""Sampled-idle watchdog thread bound to one conversation turn.
|
||||
"""Sampled-idle watchdog bound to one conversation turn (polls on the
|
||||
shared periodic scheduler thread).
|
||||
|
||||
``run_agent.py`` owns the turn-lease state (stop event, turn-active
|
||||
flag, interrupt plumbing); this class only reads the activity clock
|
||||
@@ -189,59 +190,53 @@ class TurnLivenessWatchdog:
|
||||
self._commit_abort = commit_abort
|
||||
self._deactivate_turn = deactivate_turn
|
||||
|
||||
def make_thread(self) -> threading.Thread:
|
||||
"""Build the (not yet started) watcher thread.
|
||||
def schedule(self):
|
||||
"""Start polling on the shared periodic scheduler thread.
|
||||
|
||||
``run_agent.py`` creates the watchdog before the turn begins but
|
||||
starts the thread at turn entry, right after the turn-active flag
|
||||
and the activity clock are stamped.
|
||||
schedules it at turn entry, right after the turn-active flag and
|
||||
the activity clock are stamped. Returns the cancel handle.
|
||||
"""
|
||||
return threading.Thread(
|
||||
target=self._watch,
|
||||
name="turn-liveness-watchdog",
|
||||
daemon=True,
|
||||
)
|
||||
from agent.periodic_scheduler import schedule
|
||||
|
||||
def start(self) -> threading.Thread:
|
||||
"""Spawn the watcher thread and return it (already running)."""
|
||||
thread = self.make_thread()
|
||||
thread.start()
|
||||
return thread
|
||||
return schedule(self._tick, self._poll_s)
|
||||
|
||||
def _watch(self) -> None:
|
||||
while not self._stop_event.wait(self._poll_s):
|
||||
snapshot = self._sample()
|
||||
if snapshot is None:
|
||||
# Turn is no longer active; nothing to watch.
|
||||
return
|
||||
if snapshot.idle_seconds < self._timeout_s:
|
||||
continue
|
||||
# Pre-commit surface is OBSERVATIONAL only: it reports the
|
||||
# stall and that a recovery attempt is beginning. It must not
|
||||
# claim the abort or the lease withdrawal has committed — the
|
||||
# next operation can still veto the outcome. The definitive
|
||||
# aborted/lease-stopped settlement is published by
|
||||
# _surface_committed_abort only after _commit_abort succeeds
|
||||
# and the turn is deactivated (#95663 review).
|
||||
self._surface_stall(snapshot)
|
||||
# Commit point: bind the abort to the sampled generation/ts
|
||||
# and revalidate under the lock shared with `_touch_activity`.
|
||||
# If progress resumed while the stall was being surfaced, the
|
||||
# turn continues and this loop resumes sampling — the lease
|
||||
# keeps renewing. The commit also carries the revalidated
|
||||
# generation into the interrupt path, which reserves it as a
|
||||
# claim, survives every blocking boundary (compression
|
||||
# fence), and consumes it at the final mutation edge — progress
|
||||
# landing anywhere in that window declines the abort.
|
||||
if not self._commit_abort(snapshot, self._abort_message(snapshot)):
|
||||
continue
|
||||
# Stop renewing the durable lease: a wedge the hard interrupt
|
||||
# cannot unwind must not keep the lease alive forever (the
|
||||
# issue's "lease keeps renewing" masking). The TTL expiry then
|
||||
# lets stale-turn cleanup reclaim the row.
|
||||
self._deactivate_turn()
|
||||
self._surface_committed_abort(snapshot)
|
||||
return
|
||||
def _tick(self):
|
||||
"""One poll. Returns False when the watchdog is finished."""
|
||||
if self._stop_event.is_set():
|
||||
return False
|
||||
snapshot = self._sample()
|
||||
if snapshot is None:
|
||||
# Turn is no longer active; nothing to watch.
|
||||
return False
|
||||
if snapshot.idle_seconds < self._timeout_s:
|
||||
return None
|
||||
# Pre-commit surface is OBSERVATIONAL only: it reports the
|
||||
# stall and that a recovery attempt is beginning. It must not
|
||||
# claim the abort or the lease withdrawal has committed — the
|
||||
# next operation can still veto the outcome. The definitive
|
||||
# aborted/lease-stopped settlement is published by
|
||||
# _surface_committed_abort only after _commit_abort succeeds
|
||||
# and the turn is deactivated (#95663 review).
|
||||
self._surface_stall(snapshot)
|
||||
# Commit point: bind the abort to the sampled generation/ts
|
||||
# and revalidate under the lock shared with `_touch_activity`.
|
||||
# If progress resumed while the stall was being surfaced, the
|
||||
# turn continues and this loop resumes sampling — the lease
|
||||
# keeps renewing. The commit also carries the revalidated
|
||||
# generation into the interrupt path, which reserves it as a
|
||||
# claim, survives every blocking boundary (compression
|
||||
# fence), and consumes it at the final mutation edge — progress
|
||||
# landing anywhere in that window declines the abort.
|
||||
if not self._commit_abort(snapshot, self._abort_message(snapshot)):
|
||||
return None
|
||||
# Stop renewing the durable lease: a wedge the hard interrupt
|
||||
# cannot unwind must not keep the lease alive forever (the
|
||||
# issue's "lease keeps renewing" masking). The TTL expiry then
|
||||
# lets stale-turn cleanup reclaim the row.
|
||||
self._deactivate_turn()
|
||||
self._surface_committed_abort(snapshot)
|
||||
return False
|
||||
|
||||
def _sample(self) -> Optional[ActivitySnapshot]:
|
||||
with self._activity_lock:
|
||||
|
||||
+69
-66
@@ -9324,8 +9324,11 @@ class AIAgent:
|
||||
relay_turn = None
|
||||
durable_turn_lease = None
|
||||
durable_turn_lease_stop = None
|
||||
durable_turn_lease_thread = None
|
||||
durable_turn_liveness_thread = None
|
||||
durable_turn_lease_refresh = None
|
||||
durable_turn_liveness_watchdog = None
|
||||
# Handles on the shared periodic scheduler thread (one per process,
|
||||
# agent/periodic_scheduler.py) instead of 1-2 daemon threads per turn.
|
||||
durable_turn_timer_handles = []
|
||||
durable_turn_lease_activity_lock = threading.Lock()
|
||||
durable_turn_lease_turn_active = False
|
||||
durable_turn_lease_interrupt_message = None
|
||||
@@ -9536,7 +9539,7 @@ class AIAgent:
|
||||
)
|
||||
|
||||
# Long model/tool/compression turns outlive a fixed TTL. Refresh
|
||||
# in a daemon thread; holder-qualified UPDATE and DELETE fence a
|
||||
# on the shared periodic scheduler; holder-qualified UPDATE and DELETE fence a
|
||||
# late refresher/release from a successor lease.
|
||||
durable_turn_lease_stop = threading.Event()
|
||||
_lease_refresh_interval = float(
|
||||
@@ -9553,10 +9556,10 @@ class AIAgent:
|
||||
# "active", and never be force-aborted.
|
||||
#
|
||||
# The watchdog policy (config resolution, sampling state
|
||||
# machine, thread mechanics) lives in agent/turn_liveness.py;
|
||||
# machine, polling mechanics) lives in agent/turn_liveness.py;
|
||||
# this block is only the integration seam: resolve the
|
||||
# config.yaml settings, wire the commit/deactivate callbacks
|
||||
# that own turn-lease state, and start the thread.
|
||||
# that own turn-lease state, and schedule the poll.
|
||||
try:
|
||||
from hermes_cli.config import (
|
||||
load_config_readonly as _liveness_load_config,
|
||||
@@ -9684,60 +9687,57 @@ class AIAgent:
|
||||
with durable_turn_lease_activity_lock:
|
||||
return durable_turn_lease_turn_active
|
||||
|
||||
def _refresh_durable_turn_lease() -> None:
|
||||
while not durable_turn_lease_stop.wait(_lease_refresh_interval):
|
||||
try:
|
||||
if not _turn_db.refresh_session_turn_lease(
|
||||
getattr(self, "session_id", None) or session_id,
|
||||
durable_turn_lease,
|
||||
ttl_seconds=_lease_ttl,
|
||||
):
|
||||
# finally sets the stop event then releases.
|
||||
# A late holder-fenced miss after that join
|
||||
# timeout must not hard-interrupt the next turn.
|
||||
if durable_turn_lease_stop.is_set():
|
||||
return
|
||||
logger.error(
|
||||
"Lost session turn lease while turn is active: %s",
|
||||
getattr(self, "session_id", None) or session_id,
|
||||
)
|
||||
_interrupt_turn(
|
||||
"Session turn lease lost; stopping to protect "
|
||||
"the transcript."
|
||||
)
|
||||
return
|
||||
except Exception:
|
||||
def _refresh_durable_turn_lease():
|
||||
# One periodic tick on the shared scheduler thread every
|
||||
# _lease_refresh_interval; returning False stops it.
|
||||
if durable_turn_lease_stop.is_set():
|
||||
return False
|
||||
try:
|
||||
if not _turn_db.refresh_session_turn_lease(
|
||||
getattr(self, "session_id", None) or session_id,
|
||||
durable_turn_lease,
|
||||
ttl_seconds=_lease_ttl,
|
||||
):
|
||||
# finally sets the stop event then releases.
|
||||
# A late holder-fenced miss after that cancel
|
||||
# wait must not hard-interrupt the next turn.
|
||||
if durable_turn_lease_stop.is_set():
|
||||
return
|
||||
logger.warning(
|
||||
"Failed to refresh session turn lease: %s",
|
||||
return False
|
||||
logger.error(
|
||||
"Lost session turn lease while turn is active: %s",
|
||||
getattr(self, "session_id", None) or session_id,
|
||||
exc_info=True,
|
||||
)
|
||||
_interrupt_turn(
|
||||
"Session turn lease could not be refreshed; "
|
||||
"stopping to protect the transcript."
|
||||
"Session turn lease lost; stopping to protect "
|
||||
"the transcript."
|
||||
)
|
||||
return
|
||||
return False
|
||||
except Exception:
|
||||
if durable_turn_lease_stop.is_set():
|
||||
return False
|
||||
logger.warning(
|
||||
"Failed to refresh session turn lease: %s",
|
||||
getattr(self, "session_id", None) or session_id,
|
||||
exc_info=True,
|
||||
)
|
||||
_interrupt_turn(
|
||||
"Session turn lease could not be refreshed; "
|
||||
"stopping to protect the transcript."
|
||||
)
|
||||
return False
|
||||
|
||||
durable_turn_lease_thread = threading.Thread(
|
||||
target=_refresh_durable_turn_lease,
|
||||
name="session-turn-lease-refresh",
|
||||
daemon=True,
|
||||
)
|
||||
durable_turn_lease_refresh = _refresh_durable_turn_lease
|
||||
if _liveness_timeout is not None:
|
||||
durable_turn_liveness_thread = (
|
||||
turn_liveness.TurnLivenessWatchdog(
|
||||
self,
|
||||
session_id=getattr(self, "session_id", None) or session_id,
|
||||
timeout_s=_liveness_timeout,
|
||||
poll_s=_liveness_poll,
|
||||
stop_event=durable_turn_lease_stop,
|
||||
activity_lock=self._liveness_activity_lock(),
|
||||
is_turn_active=_turn_is_active,
|
||||
commit_abort=_commit_turn_liveness_abort,
|
||||
deactivate_turn=_deactivate_turn_after_liveness_abort,
|
||||
).make_thread()
|
||||
durable_turn_liveness_watchdog = turn_liveness.TurnLivenessWatchdog(
|
||||
self,
|
||||
session_id=getattr(self, "session_id", None) or session_id,
|
||||
timeout_s=_liveness_timeout,
|
||||
poll_s=_liveness_poll,
|
||||
stop_event=durable_turn_lease_stop,
|
||||
activity_lock=self._liveness_activity_lock(),
|
||||
is_turn_active=_turn_is_active,
|
||||
commit_abort=_commit_turn_liveness_abort,
|
||||
deactivate_turn=_deactivate_turn_after_liveness_abort,
|
||||
)
|
||||
|
||||
|
||||
@@ -9790,7 +9790,7 @@ class AIAgent:
|
||||
# which may be observed from another thread.
|
||||
with bind_subagent_parent(self), scoped_runtime_main({}):
|
||||
try:
|
||||
if durable_turn_lease_thread is not None:
|
||||
if durable_turn_lease_refresh is not None:
|
||||
with durable_turn_lease_activity_lock:
|
||||
durable_turn_lease_turn_active = True
|
||||
# Stamp the activity clock at turn entry (#95663
|
||||
@@ -9803,9 +9803,17 @@ class AIAgent:
|
||||
# first poll whenever the agent had been idle longer
|
||||
# than the watchdog bound.
|
||||
self._touch_activity("starting new turn")
|
||||
durable_turn_lease_thread.start()
|
||||
if durable_turn_liveness_thread is not None:
|
||||
durable_turn_liveness_thread.start()
|
||||
from agent.periodic_scheduler import schedule as _schedule_periodic
|
||||
|
||||
durable_turn_timer_handles.append(
|
||||
_schedule_periodic(
|
||||
durable_turn_lease_refresh, _lease_refresh_interval
|
||||
)
|
||||
)
|
||||
if durable_turn_liveness_watchdog is not None:
|
||||
durable_turn_timer_handles.append(
|
||||
durable_turn_liveness_watchdog.schedule()
|
||||
)
|
||||
result = run_conversation(
|
||||
self,
|
||||
user_message,
|
||||
@@ -9874,17 +9882,12 @@ class AIAgent:
|
||||
)
|
||||
finally:
|
||||
_stop_durable_turn_lease_refresher()
|
||||
for _durable_thread in (
|
||||
durable_turn_lease_thread,
|
||||
durable_turn_liveness_thread,
|
||||
):
|
||||
if (
|
||||
_durable_thread is not None
|
||||
and _durable_thread.is_alive()
|
||||
):
|
||||
_durable_thread.join(timeout=1.0)
|
||||
# wait=1.0 mirrors the old thread join(timeout=1.0): an
|
||||
# in-flight tick on the scheduler thread finishes first.
|
||||
for _durable_handle in durable_turn_timer_handles:
|
||||
_durable_handle.cancel(wait=1.0)
|
||||
# Clear any interrupt the refresher may have fired between
|
||||
# the inner stop and this join. Must run AFTER join so a
|
||||
# the inner stop and this cancel. Must run AFTER it so a
|
||||
# late interrupt does not survive into the next turn.
|
||||
_clear_durable_turn_lease_interrupt()
|
||||
if durable_turn_lease is not None:
|
||||
|
||||
@@ -0,0 +1,92 @@
|
||||
"""agent/periodic_scheduler: one shared thread runs every periodic timer."""
|
||||
|
||||
import threading
|
||||
import time
|
||||
|
||||
from agent import periodic_scheduler
|
||||
from agent.periodic_scheduler import PeriodicScheduler, schedule
|
||||
|
||||
|
||||
def _wait_until(pred, timeout=3.0):
|
||||
deadline = time.monotonic() + timeout
|
||||
while time.monotonic() < deadline:
|
||||
if pred():
|
||||
return True
|
||||
time.sleep(0.005)
|
||||
return pred()
|
||||
|
||||
|
||||
def test_two_intervals_fire_proportionally_and_cancel_stops_one():
|
||||
sched = PeriodicScheduler()
|
||||
fast, slow = [], []
|
||||
h_fast = sched.schedule(lambda: fast.append(time.monotonic()), 0.01)
|
||||
h_slow = sched.schedule(lambda: slow.append(time.monotonic()), 0.05)
|
||||
|
||||
assert _wait_until(lambda: len(slow) >= 3)
|
||||
assert len(fast) > len(slow) # 5x interval ratio -> clearly more fast ticks
|
||||
# Both ran on this scheduler's single thread, not on new threads.
|
||||
before = threading.active_count()
|
||||
sched.schedule(lambda: None, 0.01).cancel()
|
||||
assert threading.active_count() == before
|
||||
assert sched._thread is not None and sched._thread.is_alive()
|
||||
|
||||
h_fast.cancel()
|
||||
n_fast = len(fast)
|
||||
time.sleep(0.1)
|
||||
assert len(fast) == n_fast, "cancelled callback kept firing"
|
||||
assert len(slow) > 3, "sibling callback stopped when another was cancelled"
|
||||
h_slow.cancel()
|
||||
|
||||
|
||||
def test_raising_callback_is_rescheduled_and_does_not_kill_sibling():
|
||||
sched = PeriodicScheduler()
|
||||
boom, ok = [], []
|
||||
|
||||
def raises():
|
||||
boom.append(1)
|
||||
raise RuntimeError("bad callback")
|
||||
|
||||
h1 = sched.schedule(raises, 0.01)
|
||||
h2 = sched.schedule(lambda: ok.append(1), 0.01)
|
||||
assert _wait_until(lambda: len(boom) >= 3 and len(ok) >= 3)
|
||||
h1.cancel()
|
||||
h2.cancel()
|
||||
|
||||
|
||||
def test_returning_false_stops_callback_and_cancel_wait_joins_inflight():
|
||||
sched = PeriodicScheduler()
|
||||
calls = []
|
||||
sched.schedule(lambda: (calls.append(1), False)[1], 0.01)
|
||||
assert _wait_until(lambda: len(calls) == 1)
|
||||
time.sleep(0.05)
|
||||
assert calls == [1]
|
||||
|
||||
entered = threading.Event()
|
||||
release = threading.Event()
|
||||
|
||||
def blocking():
|
||||
entered.set()
|
||||
release.wait(2.0)
|
||||
|
||||
h = sched.schedule(blocking, 0.01)
|
||||
assert entered.wait(2.0)
|
||||
threading.Timer(0.05, release.set).start()
|
||||
t0 = time.monotonic()
|
||||
h.cancel(wait=2.0) # returns once the in-flight run finished
|
||||
assert release.is_set()
|
||||
assert time.monotonic() - t0 < 1.5
|
||||
|
||||
|
||||
def test_module_level_schedule_uses_shared_default():
|
||||
hits = []
|
||||
h = schedule(lambda: hits.append(1), 0.01)
|
||||
assert _wait_until(lambda: hits)
|
||||
h.cancel()
|
||||
thread = periodic_scheduler._DEFAULT._thread
|
||||
assert thread is not None and thread.name == "hermes-periodic-scheduler"
|
||||
# Scheduling more timers on the shared default adds no OS threads.
|
||||
before = threading.active_count()
|
||||
handles = [schedule(lambda: None, 0.01) for _ in range(20)]
|
||||
assert threading.active_count() == before
|
||||
for handle in handles:
|
||||
handle.cancel()
|
||||
+81
-82
@@ -2683,7 +2683,8 @@ def _run_single_child(
|
||||
# gateway inactivity timeout doesn't fire while the subagent is working.
|
||||
# Without this, the parent's _last_activity_ts freezes when delegate_task
|
||||
# starts and the gateway eventually kills the agent for "no activity".
|
||||
_heartbeat_stop = threading.Event()
|
||||
# Runs on the shared periodic scheduler thread (agent/periodic_scheduler)
|
||||
# rather than one daemon thread per child; returning False stops it.
|
||||
# Stale detection: track the child's (tool, iteration, activity_ts) across
|
||||
# heartbeat cycles. If none advances, count the cycle as stale.
|
||||
# Different thresholds for idle vs in-tool (see _HEARTBEAT_STALE_CYCLES_*).
|
||||
@@ -2693,87 +2694,85 @@ def _run_single_child(
|
||||
_last_seen_tool = [None] # type: list
|
||||
_last_seen_activity_ts = [None] # type: list
|
||||
_stale_count = [0]
|
||||
_heartbeat_handle = [None] # type: list
|
||||
|
||||
def _heartbeat_loop():
|
||||
while not _heartbeat_stop.wait(_HEARTBEAT_INTERVAL):
|
||||
if parent_agent is None:
|
||||
continue
|
||||
touch = getattr(parent_agent, "_touch_activity", None)
|
||||
if not touch:
|
||||
continue
|
||||
# Pull detail from the child's own activity tracker
|
||||
desc = f"delegate_task: subagent {task_index} working"
|
||||
try:
|
||||
child_summary = child.get_activity_summary()
|
||||
child_tool = child_summary.get("current_tool")
|
||||
child_iter = child_summary.get("api_call_count", 0)
|
||||
child_max = child_summary.get("max_iterations", 0)
|
||||
child_activity_ts = child_summary.get("last_activity_ts")
|
||||
def _heartbeat_tick():
|
||||
if parent_agent is None:
|
||||
return
|
||||
touch = getattr(parent_agent, "_touch_activity", None)
|
||||
if not touch:
|
||||
return
|
||||
# Pull detail from the child's own activity tracker
|
||||
desc = f"delegate_task: subagent {task_index} working"
|
||||
try:
|
||||
child_summary = child.get_activity_summary()
|
||||
child_tool = child_summary.get("current_tool")
|
||||
child_iter = child_summary.get("api_call_count", 0)
|
||||
child_max = child_summary.get("max_iterations", 0)
|
||||
child_activity_ts = child_summary.get("last_activity_ts")
|
||||
|
||||
# Stale detection: count cycles where iteration, current_tool,
|
||||
# AND last_activity_ts are all frozen. A child running a
|
||||
# legitimately long-running tool keeps current_tool set; a
|
||||
# child waiting on a slow model refreshes last_activity_ts
|
||||
# via direct_api_call's activity heartbeat — neither should
|
||||
# look stale at the idle threshold.
|
||||
iter_advanced = child_iter > _last_seen_iter[0]
|
||||
tool_changed = child_tool != _last_seen_tool[0]
|
||||
activity_advanced = (
|
||||
child_activity_ts is not None
|
||||
and (
|
||||
_last_seen_activity_ts[0] is None
|
||||
or child_activity_ts > _last_seen_activity_ts[0]
|
||||
)
|
||||
# Stale detection: count cycles where iteration, current_tool,
|
||||
# AND last_activity_ts are all frozen. A child running a
|
||||
# legitimately long-running tool keeps current_tool set; a
|
||||
# child waiting on a slow model refreshes last_activity_ts
|
||||
# via direct_api_call's activity heartbeat — neither should
|
||||
# look stale at the idle threshold.
|
||||
iter_advanced = child_iter > _last_seen_iter[0]
|
||||
tool_changed = child_tool != _last_seen_tool[0]
|
||||
activity_advanced = (
|
||||
child_activity_ts is not None
|
||||
and (
|
||||
_last_seen_activity_ts[0] is None
|
||||
or child_activity_ts > _last_seen_activity_ts[0]
|
||||
)
|
||||
if iter_advanced or tool_changed or activity_advanced:
|
||||
_last_seen_iter[0] = child_iter
|
||||
_last_seen_tool[0] = child_tool
|
||||
if child_activity_ts is not None:
|
||||
_last_seen_activity_ts[0] = child_activity_ts
|
||||
_stale_count[0] = 0
|
||||
else:
|
||||
_stale_count[0] += 1
|
||||
)
|
||||
if iter_advanced or tool_changed or activity_advanced:
|
||||
_last_seen_iter[0] = child_iter
|
||||
_last_seen_tool[0] = child_tool
|
||||
if child_activity_ts is not None:
|
||||
_last_seen_activity_ts[0] = child_activity_ts
|
||||
_stale_count[0] = 0
|
||||
else:
|
||||
_stale_count[0] += 1
|
||||
|
||||
# Pick threshold based on whether the child is currently
|
||||
# inside a tool call. In-tool threshold is high enough to
|
||||
# cover legitimately slow tools; idle threshold stays
|
||||
# tight so the gateway timeout can fire on a truly wedged
|
||||
# child.
|
||||
stale_limit = (
|
||||
_HEARTBEAT_STALE_CYCLES_IN_TOOL
|
||||
if child_tool
|
||||
else _HEARTBEAT_STALE_CYCLES_IDLE
|
||||
# Pick threshold based on whether the child is currently
|
||||
# inside a tool call. In-tool threshold is high enough to
|
||||
# cover legitimately slow tools; idle threshold stays
|
||||
# tight so the gateway timeout can fire on a truly wedged
|
||||
# child.
|
||||
stale_limit = (
|
||||
_HEARTBEAT_STALE_CYCLES_IN_TOOL
|
||||
if child_tool
|
||||
else _HEARTBEAT_STALE_CYCLES_IDLE
|
||||
)
|
||||
if _stale_count[0] >= stale_limit:
|
||||
logger.warning(
|
||||
"Subagent %d appears stale (no progress for %d "
|
||||
"heartbeat cycles, tool=%s) — stopping heartbeat",
|
||||
task_index,
|
||||
_stale_count[0],
|
||||
child_tool or "<none>",
|
||||
)
|
||||
if _stale_count[0] >= stale_limit:
|
||||
logger.warning(
|
||||
"Subagent %d appears stale (no progress for %d "
|
||||
"heartbeat cycles, tool=%s) — stopping heartbeat",
|
||||
task_index,
|
||||
_stale_count[0],
|
||||
child_tool or "<none>",
|
||||
)
|
||||
break # stop touching parent, let gateway timeout fire
|
||||
return False # stop touching parent, let gateway timeout fire
|
||||
|
||||
if child_tool:
|
||||
if child_tool:
|
||||
desc = (
|
||||
f"delegate_task: subagent running {child_tool} "
|
||||
f"(iteration {child_iter}/{child_max})"
|
||||
)
|
||||
else:
|
||||
child_desc = child_summary.get("last_activity_desc", "")
|
||||
if child_desc:
|
||||
desc = (
|
||||
f"delegate_task: subagent running {child_tool} "
|
||||
f"delegate_task: subagent {child_desc} "
|
||||
f"(iteration {child_iter}/{child_max})"
|
||||
)
|
||||
else:
|
||||
child_desc = child_summary.get("last_activity_desc", "")
|
||||
if child_desc:
|
||||
desc = (
|
||||
f"delegate_task: subagent {child_desc} "
|
||||
f"(iteration {child_iter}/{child_max})"
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
try:
|
||||
touch(desc)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
_heartbeat_thread = threading.Thread(target=_heartbeat_loop, daemon=True)
|
||||
except Exception:
|
||||
pass
|
||||
try:
|
||||
touch(desc)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# Register the live agent in the module-level registry so the TUI can
|
||||
# target it by subagent_id (kill, pause, status queries). Unregistered
|
||||
@@ -2884,7 +2883,9 @@ def _run_single_child(
|
||||
}
|
||||
|
||||
try:
|
||||
_heartbeat_thread.start()
|
||||
from agent.periodic_scheduler import schedule as _schedule_periodic
|
||||
|
||||
_heartbeat_handle[0] = _schedule_periodic(_heartbeat_tick, _HEARTBEAT_INTERVAL)
|
||||
if child_progress_cb:
|
||||
try:
|
||||
child_progress_cb("subagent.start", preview=goal)
|
||||
@@ -3618,14 +3619,12 @@ def _run_single_child(
|
||||
return _error_entry
|
||||
|
||||
finally:
|
||||
# Stop the heartbeat thread so it doesn't keep touching parent activity
|
||||
# after the child has finished (or failed). Guard the join: .start()
|
||||
# now lives inside the try block, so if it raised (OS thread
|
||||
# exhaustion) the thread was never started and Thread.join() would
|
||||
# raise RuntimeError. ident is None until start() succeeds.
|
||||
_heartbeat_stop.set()
|
||||
if _heartbeat_thread.ident is not None:
|
||||
_heartbeat_thread.join(timeout=5)
|
||||
# Stop the heartbeat so it doesn't keep touching parent activity
|
||||
# after the child has finished (or failed). The handle is None if
|
||||
# scheduling itself raised (OS thread exhaustion on first use).
|
||||
# wait=5 mirrors the old thread join: an in-flight tick finishes.
|
||||
if _heartbeat_handle[0] is not None:
|
||||
_heartbeat_handle[0].cancel(wait=5)
|
||||
|
||||
# Drop the TUI-facing registry entry. Safe to call even if the
|
||||
# child was never registered (e.g. ID missing on test doubles).
|
||||
|
||||
Reference in New Issue
Block a user