Files
hermes-agent/tests/agent/test_periodic_scheduler.py
T
Teknium 561b053f79 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".
2026-09-03 02:44:24 -07:00

93 lines
2.9 KiB
Python

"""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()