From 7512f8265a4e90ca508c807ac2f94e485fe16f84 Mon Sep 17 00:00:00 2001 From: RMANOV <96174405+RMANOV@users.noreply.github.com> Date: Wed, 22 Jul 2026 15:19:54 +0300 Subject: [PATCH] fix(debate): make pump restart handoff deterministic --- debate_ops_windows.py | 44 ++++++++-- hooks/debate_pump.py | 11 +++ tests/test_debate_windows_adapter.py | 117 +++++++++++++++++++++++++++ 3 files changed, 167 insertions(+), 5 deletions(-) diff --git a/debate_ops_windows.py b/debate_ops_windows.py index c2c2a18..03911a5 100644 --- a/debate_ops_windows.py +++ b/debate_ops_windows.py @@ -147,6 +147,16 @@ def _spawn_pump_now() -> dict[str, object]: return {"spawned_pid": proc.pid} +def _wait_for_pump_running(timeout_seconds: float = 10.0) -> dict[str, object]: + """Wait for a real heartbeat instead of treating process creation as ready.""" + deadline = time.monotonic() + max(0.0, timeout_seconds) + state = pump_state() + while state.get("state") != "running" and time.monotonic() < deadline: + time.sleep(0.25) + state = pump_state() + return state + + def _schtasks(args: list[str]) -> dict[str, object]: cmd = ["schtasks", *args] result = subprocess.run(cmd, capture_output=True, text=True, check=False) @@ -296,14 +306,22 @@ def cmd_install(*, start: bool = True) -> int: def cmd_start() -> int: + current = pump_state() + if current.get("state") == "running": + print( + json.dumps( + {"task": TASK_NAME, "start": "already_running", "pump_state": current}, + indent=2, + ) + ) + return 0 out = _schtasks(["/Run", "/TN", TASK_NAME]) payload: dict[str, object] = {"task": TASK_NAME, "run": out} if out.get("returncode"): payload["fallback"] = _spawn_pump_now() - time.sleep(2.0) - payload["pump_state"] = pump_state() + payload["pump_state"] = _wait_for_pump_running() print(json.dumps(payload, indent=2)) - return 0 if not out.get("returncode") or "fallback" in payload else 1 + return 0 if payload["pump_state"].get("state") == "running" else 1 def cmd_stop(*, timeout_seconds: float = 20.0) -> int: @@ -311,6 +329,22 @@ def cmd_stop(*, timeout_seconds: float = 20.0) -> int: In-flight workers survive either path (own process group, no console). """ + # Capture the exact process identity before signaling. A clean pump removes + # its heartbeat just before interpreter exit; using only the current file + # state can therefore report "stopped" while that process still owns the + # singleton mutex, making an immediate restart lose a timing race. + target_heartbeat = _read_heartbeat() + target_live = bool(target_heartbeat and _heartbeat_pid_live(target_heartbeat)) + if not target_live: + state = pump_state() + print( + json.dumps( + {"task": TASK_NAME, "stopped": "already_stopped", "pump_state": state}, + indent=2, + ) + ) + return 0 + graceful = False try: sys.path.insert(0, str(ROOT)) @@ -321,8 +355,8 @@ def cmd_stop(*, timeout_seconds: float = 20.0) -> int: graceful = False deadline = time.monotonic() + max(1.0, timeout_seconds) while time.monotonic() < deadline: - state = pump_state() - if not state.get("pid_live"): + if not _heartbeat_pid_live(target_heartbeat): + state = pump_state() print( json.dumps( {"task": TASK_NAME, "stopped": "graceful", "pump_state": state}, diff --git a/hooks/debate_pump.py b/hooks/debate_pump.py index 25c7d37..d3e833f 100644 --- a/hooks/debate_pump.py +++ b/hooks/debate_pump.py @@ -1686,6 +1686,17 @@ def main() -> int: _reap_children() _log("pump_stop", pid=os.getpid(), last_ts=last_ts, last_msg_id=last_msg_id) + if IS_WINDOWS and not args.once: + try: + # Release the singleton before publishing the stopped state. The + # lifecycle command may start the replacement as soon as the + # heartbeat disappears; leaving release to interpreter teardown + # creates a small but real stop -> start race. + from debate_wake_signal import release_pump_singleton + + release_pump_singleton() + except Exception as exc: + _log("pump_singleton_release_failed", error=repr(exc)) try: # Clean exit removes the heartbeat so status reads "stopped", not # "stale"; a crashed pump leaves it behind — which is the signal diff --git a/tests/test_debate_windows_adapter.py b/tests/test_debate_windows_adapter.py index 44b3678..3749321 100644 --- a/tests/test_debate_windows_adapter.py +++ b/tests/test_debate_windows_adapter.py @@ -423,6 +423,49 @@ def test_pump_singleton_claims_existing_unowned_mutex(): observer.wait(timeout=15) +@windows_only +def test_pump_singleton_explicit_release_precedes_process_exit(): + """A clean pump can hand off its lease before interpreter teardown.""" + import subprocess + + env = dict(os.environ) + env["DEBATE_PUMP_SINGLETON_MUTEX"] = rf"Local\DebateReleaseTest{os.getpid()}" + repo_bootstrap = "import sys; sys.path.insert(0, r'" + str(REPO) + "'); " + holder = subprocess.Popen( + [ + sys.executable, + "-c", + repo_bootstrap + + "from debate_wake_signal import acquire_pump_singleton, " + + "release_pump_singleton; import time; " + + "print(acquire_pump_singleton(), flush=True); " + + "release_pump_singleton(); print('released', flush=True); time.sleep(4)", + ], + stdout=subprocess.PIPE, + text=True, + env=env, + ) + try: + assert holder.stdout.readline().strip() == "True" + assert holder.stdout.readline().strip() == "released" + contender = subprocess.run( + [ + sys.executable, + "-c", + repo_bootstrap + + "from debate_wake_signal import acquire_pump_singleton; " + + "print(acquire_pump_singleton())", + ], + capture_output=True, + text=True, + timeout=15, + env=env, + ) + assert contender.stdout.strip() == "True" + finally: + holder.wait(timeout=15) + + def test_agent_log_dir_stays_bounded(tmp_path, monkeypatch): import debate_wake @@ -481,6 +524,80 @@ def test_pump_state_running_stale_stopped(tmp_path, monkeypatch): assert state["pid_live"] is False +def test_stop_waits_for_captured_process_after_heartbeat_disappears( + monkeypatch, capsys +): + import debate_ops_windows as dow + import debate_wake_signal + + heartbeat = {"pid": 1234, "create_time": 5678.0} + liveness = iter([True, True, False]) + monkeypatch.setattr(dow, "_read_heartbeat", lambda: heartbeat) + monkeypatch.setattr(dow, "_heartbeat_pid_live", lambda _hb: next(liveness)) + monkeypatch.setattr( + dow, "pump_state", lambda: {"state": "stopped", "heartbeat": None} + ) + monkeypatch.setattr(debate_wake_signal, "signal_stop", lambda: True) + monkeypatch.setattr(dow.time, "sleep", lambda _seconds: None) + + assert dow.cmd_stop(timeout_seconds=5.0) == 0 + payload = json.loads(capsys.readouterr().out) + assert payload["stopped"] == "graceful" + + +def test_start_requires_observed_running_heartbeat(monkeypatch, capsys): + import debate_ops_windows as dow + + states = iter( + [ + {"state": "stopped", "heartbeat": None}, + {"state": "stopped", "heartbeat": None}, + ] + ) + monkeypatch.setattr(dow, "pump_state", lambda: next(states)) + monkeypatch.setattr( + dow, + "_schtasks", + lambda _args: {"returncode": 1, "stdout": "", "stderr": "missing"}, + ) + monkeypatch.setattr(dow, "_spawn_pump_now", lambda: {"spawned_pid": 9876}) + monkeypatch.setattr( + dow, + "_wait_for_pump_running", + lambda: {"state": "stopped", "heartbeat": None}, + ) + + assert dow.cmd_start() == 1 + payload = json.loads(capsys.readouterr().out) + assert payload["fallback"]["spawned_pid"] == 9876 + assert payload["pump_state"]["state"] == "stopped" + + +def test_stop_when_already_stopped_does_not_leave_stop_event_signaled( + monkeypatch, capsys +): + import debate_ops_windows as dow + import debate_wake_signal + + signaled = False + + def signal_stop(): + nonlocal signaled + signaled = True + return True + + monkeypatch.setattr(dow, "_read_heartbeat", lambda: None) + monkeypatch.setattr( + dow, "pump_state", lambda: {"state": "stopped", "heartbeat": None} + ) + monkeypatch.setattr(debate_wake_signal, "signal_stop", signal_stop) + + assert dow.cmd_stop() == 0 + assert signaled is False + payload = json.loads(capsys.readouterr().out) + assert payload["stopped"] == "already_stopped" + + @windows_only def test_task_xml_encodes_spec_constraints(): import debate_ops_windows as dow