test(terminal): cover hung-wait bound, parent-tid interrupt, and cron inactivity watchdog
Pin that execute() returns at the wall-clock deadline when the inner wait never returns, that /stop on the tool-worker tid still kills the subprocess, that the cron inactivity helper fires while the caller thread is blocked, and that ContextVars plus the activity callback reach the deadline worker.
This commit is contained in:
@@ -224,11 +224,54 @@ class TestRunBoundedSync:
|
||||
assert result.timed_out is True
|
||||
release.set()
|
||||
|
||||
def test_keyboard_interrupt_lands_before_full_deadline(self):
|
||||
"""Sliced Event.wait must observe SetAsyncExc within one poll slice."""
|
||||
release = threading.Event()
|
||||
holder: dict = {}
|
||||
|
||||
def _run():
|
||||
try:
|
||||
holder["result"] = run_bounded_sync(
|
||||
lambda: release.wait(30),
|
||||
10.0,
|
||||
label="ki",
|
||||
)
|
||||
except KeyboardInterrupt:
|
||||
holder["exc"] = "KeyboardInterrupt"
|
||||
|
||||
t = threading.Thread(target=_run)
|
||||
t.start()
|
||||
time.sleep(0.15)
|
||||
import ctypes
|
||||
|
||||
assert t.ident is not None
|
||||
ret = ctypes.pythonapi.PyThreadState_SetAsyncExc(
|
||||
ctypes.c_ulong(t.ident), ctypes.py_object(KeyboardInterrupt),
|
||||
)
|
||||
assert ret == 1
|
||||
t.join(timeout=2.0)
|
||||
release.set()
|
||||
assert not t.is_alive()
|
||||
assert holder.get("exc") == "KeyboardInterrupt"
|
||||
|
||||
def test_deadline_expired_is_a_timeout_error(self):
|
||||
# Error-classification contract: our deadline must be catchable as
|
||||
# TimeoutError but distinguishable by type from transport timeouts.
|
||||
assert issubclass(DeadlineExpired, TimeoutError)
|
||||
|
||||
def test_worker_inherits_caller_contextvars(self):
|
||||
"""Profile secret scope / session id must survive the thread hop."""
|
||||
import contextvars
|
||||
|
||||
var = contextvars.ContextVar("deadline_sync_ctx")
|
||||
token = var.set("from-caller")
|
||||
try:
|
||||
result = run_bounded_sync(lambda: var.get(None), 5.0, label="ctx")
|
||||
finally:
|
||||
var.reset(token)
|
||||
assert result.timed_out is False
|
||||
assert result.value == "from-caller"
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# run_bounded_async
|
||||
|
||||
@@ -11,6 +11,7 @@ Tests cover:
|
||||
import concurrent.futures
|
||||
import os
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
from pathlib import Path
|
||||
|
||||
@@ -235,3 +236,73 @@ class TestSysPathOrdering:
|
||||
from cron.scheduler import _hermes_now
|
||||
assert callable(_hermes_now)
|
||||
|
||||
|
||||
class TestInactivityWatchdogLoop:
|
||||
"""The daemon-thread inactivity helper must not depend on the caller thread."""
|
||||
|
||||
def test_fires_when_idle_crosses_limit(self):
|
||||
from cron.scheduler import _inactivity_watchdog_loop
|
||||
|
||||
stop = threading.Event()
|
||||
idle = {"s": 0.0}
|
||||
results: list = []
|
||||
|
||||
def _watch():
|
||||
results.append(
|
||||
_inactivity_watchdog_loop(
|
||||
get_idle_seconds=lambda: idle["s"],
|
||||
limit_s=0.2,
|
||||
poll_s=0.05,
|
||||
stop=stop,
|
||||
future_done=lambda: False,
|
||||
)
|
||||
)
|
||||
|
||||
watcher = threading.Thread(target=_watch, daemon=True)
|
||||
watcher.start()
|
||||
time.sleep(0.12)
|
||||
idle["s"] = 1.0
|
||||
watcher.join(timeout=2.0)
|
||||
stop.set()
|
||||
assert results == [True]
|
||||
assert not watcher.is_alive()
|
||||
|
||||
def test_stops_when_future_completes_before_idle_limit(self):
|
||||
from cron.scheduler import _inactivity_watchdog_loop
|
||||
|
||||
stop = threading.Event()
|
||||
fired = _inactivity_watchdog_loop(
|
||||
get_idle_seconds=lambda: 0.0,
|
||||
limit_s=10.0,
|
||||
poll_s=0.05,
|
||||
stop=stop,
|
||||
future_done=lambda: True,
|
||||
)
|
||||
assert fired is False
|
||||
|
||||
def test_fires_while_caller_thread_is_blocked(self):
|
||||
"""#94285: a blocked run_job thread must not disable the watchdog."""
|
||||
from cron.scheduler import _inactivity_watchdog_loop
|
||||
|
||||
stop = threading.Event()
|
||||
idle = {"s": 1.0}
|
||||
result = {"fired": None}
|
||||
|
||||
def _watch():
|
||||
result["fired"] = _inactivity_watchdog_loop(
|
||||
get_idle_seconds=lambda: idle["s"],
|
||||
limit_s=0.15,
|
||||
poll_s=0.05,
|
||||
stop=stop,
|
||||
future_done=lambda: False,
|
||||
)
|
||||
|
||||
watcher = threading.Thread(target=_watch, daemon=True)
|
||||
watcher.start()
|
||||
# Simulate the family-A stall: this thread cannot poll.
|
||||
time.sleep(0.4)
|
||||
watcher.join(timeout=2.0)
|
||||
stop.set()
|
||||
assert result["fired"] is True
|
||||
assert not watcher.is_alive()
|
||||
|
||||
|
||||
@@ -27,6 +27,20 @@ class TestInterruptModule:
|
||||
set_interrupt(False)
|
||||
assert not is_interrupted()
|
||||
|
||||
def test_is_thread_interrupted_checks_target_tid_not_caller(self):
|
||||
from tools.interrupt import (
|
||||
set_interrupt, is_interrupted, is_thread_interrupted, _interrupted_threads, _lock,
|
||||
)
|
||||
with _lock:
|
||||
_interrupted_threads.clear()
|
||||
other_tid = threading.get_ident() + 1
|
||||
set_interrupt(True, thread_id=other_tid)
|
||||
assert not is_interrupted()
|
||||
assert is_thread_interrupted(other_tid)
|
||||
assert is_thread_interrupted(None) is False
|
||||
set_interrupt(False, thread_id=other_tid)
|
||||
assert not is_thread_interrupted(other_tid)
|
||||
|
||||
|
||||
def test_clear_current_thread_interrupt_leaves_other_threads(self):
|
||||
"""clear_current_thread_interrupt only touches the calling thread."""
|
||||
|
||||
@@ -0,0 +1,115 @@
|
||||
"""Foreground terminal execute must return when the inner wait loop wedges.
|
||||
|
||||
#94285: a hung ``_wait_for_process`` (Windows pipe/poll, blocked loop thread)
|
||||
silently disabled every asyncio timer in the process. ``execute()`` now
|
||||
bounds spawn+wait with ``run_bounded_sync`` so the wall-clock deadline
|
||||
survives a wedged wait, and ``on_timeout`` kills the process tree.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import time
|
||||
from types import SimpleNamespace
|
||||
|
||||
from tools.environments.local import LocalEnvironment
|
||||
import tools.environments.base as base_mod
|
||||
|
||||
|
||||
def test_execute_returns_when_wait_loop_never_returns(monkeypatch):
|
||||
"""A wedged inner wait cannot hold execute() past timeout + grace."""
|
||||
monkeypatch.setattr(base_mod, "_EXECUTE_WAIT_BOUND_GRACE_S", 0.05)
|
||||
|
||||
env = LocalEnvironment()
|
||||
fake_proc = SimpleNamespace(pid=424242)
|
||||
monkeypatch.setattr(env, "_run_bash", lambda *a, **k: fake_proc)
|
||||
|
||||
def _hang(*_a, **_k):
|
||||
time.sleep(30)
|
||||
return {"output": "late", "returncode": 0}
|
||||
|
||||
monkeypatch.setattr(env, "_wait_for_process", _hang)
|
||||
killed: list = []
|
||||
monkeypatch.setattr(env, "_kill_process", lambda proc: killed.append(("kill", proc)))
|
||||
monkeypatch.setattr(
|
||||
"agent.deadline.kill_process_tree",
|
||||
lambda pid, **_k: killed.append(("tree", pid)),
|
||||
)
|
||||
monkeypatch.setattr(env, "_update_cwd", lambda _result: None)
|
||||
|
||||
start = time.monotonic()
|
||||
result = env.execute("sleep 30", timeout=1)
|
||||
elapsed = time.monotonic() - start
|
||||
|
||||
assert elapsed < 4.0, f"execute hung {elapsed:.1f}s past the 1s bound"
|
||||
assert result["returncode"] == 124
|
||||
assert "timed out" in result["output"].lower()
|
||||
assert ("kill", fake_proc) in killed
|
||||
assert ("tree", 424242) in killed
|
||||
|
||||
|
||||
def test_execute_parent_interrupt_still_kills_wait_on_deadline_worker(monkeypatch):
|
||||
"""/stop targets the tool-worker tid; the deadline worker must honor it."""
|
||||
from tools.interrupt import set_interrupt, is_interrupted
|
||||
|
||||
env = LocalEnvironment()
|
||||
fake_proc = SimpleNamespace(pid=None, poll=lambda: None, stdout=None)
|
||||
monkeypatch.setattr(env, "_run_bash", lambda *a, **k: fake_proc)
|
||||
|
||||
seen = {"parent": False}
|
||||
|
||||
def _wait(_proc, timeout=120, *, bounded_capture=False, watch_interrupt_tid=None):
|
||||
deadline = time.monotonic() + 2.0
|
||||
while time.monotonic() < deadline:
|
||||
from tools.interrupt import is_thread_interrupted
|
||||
|
||||
if is_interrupted() or is_thread_interrupted(watch_interrupt_tid):
|
||||
seen["parent"] = True
|
||||
return {"output": "[Command interrupted]", "returncode": 130}
|
||||
time.sleep(0.02)
|
||||
return {"output": "missed interrupt", "returncode": 0}
|
||||
|
||||
monkeypatch.setattr(env, "_wait_for_process", _wait)
|
||||
monkeypatch.setattr(env, "_kill_process", lambda _proc: None)
|
||||
monkeypatch.setattr(env, "_update_cwd", lambda _result: None)
|
||||
|
||||
parent_tid = __import__("threading").get_ident()
|
||||
|
||||
def _interrupt_soon():
|
||||
time.sleep(0.05)
|
||||
set_interrupt(True, thread_id=parent_tid)
|
||||
|
||||
import threading
|
||||
|
||||
threading.Thread(target=_interrupt_soon, daemon=True).start()
|
||||
result = env.execute("sleep 30", timeout=5)
|
||||
set_interrupt(False, thread_id=parent_tid)
|
||||
|
||||
assert seen["parent"] is True
|
||||
assert result["returncode"] == 130
|
||||
|
||||
|
||||
def test_execute_worker_sees_caller_activity_callback(monkeypatch):
|
||||
"""Heartbeats must fire on the deadline worker, not only the tool thread."""
|
||||
env = LocalEnvironment()
|
||||
fake_proc = SimpleNamespace(pid=None)
|
||||
monkeypatch.setattr(env, "_run_bash", lambda *a, **k: fake_proc)
|
||||
seen = {"cb": "unset"}
|
||||
|
||||
def _wait(*_a, **_k):
|
||||
seen["cb"] = base_mod.get_activity_callback()
|
||||
return {"output": "ok", "returncode": 0}
|
||||
|
||||
monkeypatch.setattr(env, "_wait_for_process", _wait)
|
||||
monkeypatch.setattr(env, "_update_cwd", lambda _r: None)
|
||||
|
||||
def _cb(_msg):
|
||||
pass
|
||||
|
||||
base_mod.set_activity_callback(_cb)
|
||||
try:
|
||||
result = env.execute("echo ok", timeout=5)
|
||||
finally:
|
||||
base_mod.set_activity_callback(None)
|
||||
|
||||
assert result["returncode"] == 0
|
||||
assert seen["cb"] is _cb
|
||||
Reference in New Issue
Block a user