hermes-agent/tests/gateway/test_abandoned_turn_process...

235 lines
6.7 KiB
Python

"""Regression coverage for abandoned gateway-turn subprocess cleanup (#76115)."""
import threading
from gateway.run import (
_abandon_timed_out_gateway_turn,
_reap_gateway_turn_processes,
_watch_gateway_turn_inactivity,
)
from tools.process_registry import process_registry
class _IdleAgent:
def __init__(self, idle_seconds=60.0):
self.idle_seconds = idle_seconds
self.interrupts = []
def get_activity_summary(self):
return {"seconds_since_activity": self.idle_seconds}
def interrupt(self, reason):
self.interrupts.append(reason)
def _state():
return threading.Event(), threading.Event(), threading.Lock()
def test_thread_watchdog_reaps_only_processes_created_by_timed_out_turn(monkeypatch):
agent = _IdleAgent()
worker_done, timeout_fired, cleanup_lock = _state()
calls = []
monkeypatch.setattr(
process_registry,
"kill_started_since",
lambda task_id, baseline, *, source: calls.append(
(task_id, baseline, source)
)
or 1,
)
watchdog = threading.Thread(
target=_watch_gateway_turn_inactivity,
kwargs={
"agent_holder": [agent],
"task_id": "session-a",
"process_baseline": frozenset({"proc_existing"}),
"timeout": 30.0,
"worker_done": worker_done,
"timeout_fired": timeout_fired,
"cleanup_lock": cleanup_lock,
"poll_interval": 0.01,
},
)
watchdog.start()
watchdog.join(timeout=1)
assert not watchdog.is_alive()
assert timeout_fired.is_set()
assert agent.interrupts == ["Execution timed out (inactivity)"]
assert calls == [
(
"session-a",
frozenset({"proc_existing"}),
"gateway_turn_timeout",
)
]
def test_completed_worker_wins_race_and_preserves_background_process(monkeypatch):
agent = _IdleAgent()
worker_done, timeout_fired, cleanup_lock = _state()
worker_done.set()
monkeypatch.setattr(
process_registry,
"kill_started_since",
lambda *_args, **_kwargs: (_ for _ in ()).throw(
AssertionError("completed turn must not reap background work")
),
)
assert not _abandon_timed_out_gateway_turn(
agent_holder=[agent],
task_id="session-a",
process_baseline=frozenset(),
worker_done=worker_done,
timeout_fired=timeout_fired,
cleanup_lock=cleanup_lock,
)
assert not timeout_fired.is_set()
assert agent.interrupts == []
def test_timeout_cleanup_is_idempotent(monkeypatch):
agent = _IdleAgent()
worker_done, timeout_fired, cleanup_lock = _state()
calls = []
monkeypatch.setattr(
process_registry,
"kill_started_since",
lambda *_args, **_kwargs: calls.append(True) or 0,
)
kwargs = {
"agent_holder": [agent],
"task_id": "session-a",
"process_baseline": frozenset(),
"worker_done": worker_done,
"timeout_fired": timeout_fired,
"cleanup_lock": cleanup_lock,
}
assert _abandon_timed_out_gateway_turn(**kwargs)
assert not _abandon_timed_out_gateway_turn(**kwargs)
assert len(calls) == 1
assert len(agent.interrupts) == 1
# ---------------------------------------------------------------------------
# Cross-turn race guard (#76188 review): task_id is session-scoped, not
# turn-scoped, so a replacement turn on the same session could otherwise
# have its freshly-spawned process killed by a stale reaper. Gated on
# run_generation via an injected `is_still_current` check.
# ---------------------------------------------------------------------------
def test_reap_skips_when_a_newer_turn_has_claimed_the_session(monkeypatch):
calls = []
monkeypatch.setattr(
process_registry,
"kill_started_since",
lambda *_a, **_k: calls.append(True) or 1,
)
killed = _reap_gateway_turn_processes(
"session-a",
frozenset({"proc_old"}),
source="gateway_turn_timeout",
is_still_current=lambda: False,
)
assert killed == 0
assert calls == []
def test_reap_proceeds_when_this_turn_is_still_current(monkeypatch):
calls = []
monkeypatch.setattr(
process_registry,
"kill_started_since",
lambda task_id, baseline, *, source: calls.append(
(task_id, baseline, source)
)
or 1,
)
killed = _reap_gateway_turn_processes(
"session-a",
frozenset({"proc_old"}),
source="gateway_turn_timeout",
is_still_current=lambda: True,
)
assert killed == 1
assert calls == [("session-a", frozenset({"proc_old"}), "gateway_turn_timeout")]
def test_reap_fails_open_when_is_still_current_raises(monkeypatch):
"""A bug in the generation-check closure must not silently disable the
underlying leak fix — it should log and fall through to reaping."""
calls = []
monkeypatch.setattr(
process_registry,
"kill_started_since",
lambda *_a, **_k: calls.append(True) or 1,
)
def _boom():
raise RuntimeError("session state lookup failed")
killed = _reap_gateway_turn_processes(
"session-a",
frozenset(),
source="gateway_turn_timeout",
is_still_current=_boom,
)
assert killed == 1
assert calls == [True]
def test_reap_skips_empty_task_id(monkeypatch):
"""ProcessSession.task_id defaults to "" — a blank turn id must never
fan out into killing unrelated sessionless processes (#76188 review)."""
calls = []
monkeypatch.setattr(
process_registry,
"kill_started_since",
lambda *_a, **_k: calls.append(True) or 1,
)
killed = _reap_gateway_turn_processes(
"",
frozenset(),
source="gateway_turn_timeout",
)
assert killed == 0
assert calls == []
def test_timeout_abandon_propagates_is_still_current_to_the_reap(monkeypatch):
agent = _IdleAgent()
worker_done, timeout_fired, cleanup_lock = _state()
calls = []
monkeypatch.setattr(
process_registry,
"kill_started_since",
lambda *_a, **_k: calls.append(True) or 1,
)
assert _abandon_timed_out_gateway_turn(
agent_holder=[agent],
task_id="session-a",
process_baseline=frozenset(),
worker_done=worker_done,
timeout_fired=timeout_fired,
cleanup_lock=cleanup_lock,
is_still_current=lambda: False,
)
# The turn was still marked abandoned (interrupt fired), but the actual
# reap was skipped because a newer turn already claimed the session.
assert agent.interrupts == ["Execution timed out (inactivity)"]
assert calls == []