Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
44 changes: 39 additions & 5 deletions debate_ops_windows.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -296,21 +306,45 @@ 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:
"""Graceful first: stop event → wait for exit; /End only as fallback.

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:
Comment on lines +336 to +338
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))
Expand All @@ -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(
Comment on lines 357 to 360
json.dumps(
{"task": TASK_NAME, "stopped": "graceful", "pump_state": state},
Expand Down
11 changes: 11 additions & 0 deletions hooks/debate_pump.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
117 changes: 117 additions & 0 deletions tests/test_debate_windows_adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand Down