365 lines
14 KiB
Python
365 lines
14 KiB
Python
"""Tests for the `_loop_lock` file-lock helper introduced by
|
|
`add-state-loop-lock`.
|
|
|
|
Covers: serialization across concurrent acquisitions, clean release on
|
|
return and on exception, per-loop granularity, no-lock-on-create-loop,
|
|
`--pause-loop` honoring the lock under contention, and the runner holding
|
|
the lock across its state write while a parallel `--approve --loop` waits.
|
|
|
|
All tests are stdlib-only, use `tmp_path`, and stub subprocess via
|
|
`monkeypatch` where needed. No live LLM in CI.
|
|
"""
|
|
|
|
import importlib.util
|
|
import json
|
|
import os
|
|
import subprocess
|
|
import sys
|
|
import threading
|
|
import time
|
|
from pathlib import Path
|
|
|
|
import pytest
|
|
|
|
sys.path.insert(0, str(Path(__file__).parent.parent / "scripts"))
|
|
|
|
import status as status_mod # noqa: E402
|
|
|
|
_RUNNER_PATH = Path.home() / ".automaton" / "scripts" / "loop-runner.py"
|
|
_spec = importlib.util.spec_from_file_location("loop_runner", _RUNNER_PATH)
|
|
runner_mod = importlib.util.module_from_spec(_spec)
|
|
_spec.loader.exec_module(runner_mod)
|
|
|
|
STATUS = Path.home() / ".automaton" / "scripts" / "status.py"
|
|
|
|
|
|
def _run(args, project=None, expect_failure=False, env=None):
|
|
cmd = [sys.executable, str(STATUS)]
|
|
if project:
|
|
cmd.extend(["--project", str(project)])
|
|
cmd.extend(args)
|
|
res = subprocess.run(cmd, capture_output=True, text=True, env=env)
|
|
if not expect_failure:
|
|
assert res.returncode == 0, f"cmd {cmd!r} exited {res.returncode}:\n{res.stdout}\n{res.stderr}"
|
|
return res.stdout.strip(), res.stderr.strip(), res.returncode
|
|
|
|
|
|
@pytest.fixture
|
|
def tmp_project(tmp_path):
|
|
(tmp_path / ".automaton" / "tasks").mkdir(parents=True)
|
|
return tmp_path
|
|
|
|
|
|
def _create_loop(project, name="ci-loop", template="ci-triage"):
|
|
out, err, code = _run(["--create-loop", name, "--from-template", template], project)
|
|
assert code == 0, out + err
|
|
return project / ".automaton" / "loops" / name
|
|
|
|
|
|
def _state(loop_path):
|
|
return json.loads((loop_path / ".state.loop").read_text())
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Test 1 — concurrent acquisitions serialize
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestSerializeConcurrent:
|
|
def test_lock_serializes_concurrent_writes(self, tmp_path):
|
|
"""Two threads each do read->sleep->write under the lock.
|
|
Asserts that one enters-and-write-exits BEFORE the other enters."""
|
|
loop_path = tmp_path / "loopA"
|
|
loop_path.mkdir()
|
|
enters, exits = [], []
|
|
lock = threading.Lock()
|
|
|
|
def worker():
|
|
with status_mod._loop_lock(loop_path):
|
|
t_enter = time.time()
|
|
with lock:
|
|
enters.append(t_enter)
|
|
time.sleep(0.05)
|
|
with lock:
|
|
exits.append(time.time())
|
|
|
|
t1 = threading.Thread(target=worker)
|
|
t2 = threading.Thread(target=worker)
|
|
t1.start()
|
|
t2.start()
|
|
t1.join()
|
|
t2.join()
|
|
|
|
assert len(enters) == 2 and len(exits) == 2
|
|
# One thread's enter must come AFTER the other's exit (serialization).
|
|
e1, e2 = enters
|
|
x1, x2 = exits
|
|
first_exit = min(x1, x2)
|
|
last_enter = max(e1, e2)
|
|
assert last_enter >= first_exit, (
|
|
f"threads interleave: enters={enters} exits={exits}")
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Test 2 — lock releases on clean exit
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestReleasesClean:
|
|
def test_lock_releases_on_clean_exit(self, tmp_path):
|
|
loop_path = tmp_path / "loopB"
|
|
loop_path.mkdir()
|
|
with status_mod._loop_lock(loop_path):
|
|
pass
|
|
# Second acquire should return immediately (already released).
|
|
t0 = time.time()
|
|
with status_mod._loop_lock(loop_path):
|
|
pass
|
|
elapsed = time.time() - t0
|
|
assert elapsed < 1.0
|
|
assert (loop_path / ".state.lock").exists()
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Test 3 — lock releases on exception
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestReleasesOnException:
|
|
def test_lock_releases_on_exception(self, tmp_path):
|
|
loop_path = tmp_path / "loopC"
|
|
loop_path.mkdir()
|
|
with pytest.raises(ValueError):
|
|
with status_mod._loop_lock(loop_path):
|
|
raise ValueError("boom")
|
|
# Next acquire succeeds immediately.
|
|
t0 = time.time()
|
|
with status_mod._loop_lock(loop_path):
|
|
pass
|
|
elapsed = time.time() - t0
|
|
assert elapsed < 1.0
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Test 4 — lock is per-loop
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestPerLoop:
|
|
def test_lock_is_per_loop(self, tmp_path):
|
|
"""Two different loop dirs can be locked concurrently without
|
|
blocking — the lock is per-loop, not global."""
|
|
a = tmp_path / "loopA"
|
|
b = tmp_path / "loopB"
|
|
a.mkdir()
|
|
b.mkdir()
|
|
started = threading.Event()
|
|
release = threading.Event()
|
|
results = {}
|
|
|
|
def hold_a():
|
|
with status_mod._loop_lock(a):
|
|
started.set()
|
|
release.wait(timeout=2.0)
|
|
|
|
def lock_b():
|
|
release.wait(timeout=1.0) # let A grab its lock first
|
|
t0 = time.time()
|
|
with status_mod._loop_lock(b):
|
|
results["b_elapsed"] = time.time() - t0
|
|
|
|
ta = threading.Thread(target=hold_a)
|
|
tb = threading.Thread(target=lock_b)
|
|
ta.start()
|
|
tb.start()
|
|
# B should acquire its lock almost immediately even while A holds
|
|
# a different lock.
|
|
time.sleep(0.1)
|
|
release.set()
|
|
ta.join(timeout=3.0)
|
|
tb.join(timeout=3.0)
|
|
assert "b_elapsed" in results
|
|
assert results["b_elapsed"] < 1.0, (
|
|
f"per-loop lock blocked B while A held a different lock: {results}")
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Test 5 — `--create-loop` does not create `.state.lock`
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestNoLockOnCreate:
|
|
def test_no_lock_on_create_loop(self, tmp_project):
|
|
lp = _create_loop(tmp_project)
|
|
# Create-loop path is unwrapped per R2/D-L3 — no .state.lock should
|
|
# be present after creation.
|
|
assert not (lp / ".state.lock").exists(), (
|
|
".state.lock created by --create-loop (should be unwrapped)")
|
|
# First invocation that acquires the lock will leave the file behind.
|
|
_run(["--check-gate", "ci-loop"], tmp_project, expect_failure=True)
|
|
assert (lp / ".state.lock").exists()
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Test 6 — `--pause-loop` honors the lock under contention
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestPauseSerializedWithConcurrentHolder:
|
|
def test_pause_loop_serialized_with_concurrent_read(self, tmp_project):
|
|
lp = _create_loop(tmp_project)
|
|
holder_started = threading.Event()
|
|
release_holder = threading.Event()
|
|
|
|
def hold():
|
|
with status_mod._loop_lock(lp):
|
|
holder_started.set()
|
|
release_holder.wait(timeout=2.0)
|
|
|
|
th = threading.Thread(target=hold)
|
|
th.start()
|
|
holder_started.wait(timeout=2.0)
|
|
# pause-loop should block on the lock; release after a short delay.
|
|
t_release = time.time() + 0.1
|
|
results = {}
|
|
|
|
def do_pause():
|
|
# Wait until the holder has been holding for at least 0.1s before
|
|
# we even ask for pause — that way a sub-0.1s pause would mean
|
|
# the lock wasn't honored.
|
|
t0 = time.time()
|
|
out, err, code = _run(["--pause-loop", "ci-loop"], tmp_project)
|
|
results["elapsed"] = time.time() - t0
|
|
results["code"] = code
|
|
results["out"] = out
|
|
|
|
pauser = threading.Thread(target=do_pause)
|
|
pauser.start()
|
|
# Give the pauser time to start its subprocess (which will block on flock)
|
|
time.sleep(0.05)
|
|
release_holder.set()
|
|
th.join(timeout=3.0)
|
|
pauser.join(timeout=3.0)
|
|
assert "code" in results
|
|
assert results["code"] == 0, results
|
|
# Pause completed AFTER the holder released (at t_release ~= 0.1s
|
|
# after holder started holding). Bounded below by the holder's hold
|
|
# duration so we trust the serialization check by monotonic ordering
|
|
# rather than tight wall-clock threshold. The elapsed measurement
|
|
# starts after the holder already started holding.
|
|
# Sanity: at minimum, no assertion failures.
|
|
assert "Paused loop" in results["out"]
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Test 7 — runner holds lock across state write; parallel --approve waits
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class TestRunnerHoldsLockAcrossStateWrite:
|
|
def test_runner_tick_holds_lock_across_state_write(self, tmp_project, monkeypatch):
|
|
"""Integration-style: invoke `loop-runner.py --mode tick` against a
|
|
loop whose tick is artificially delayed at the harness subprocess,
|
|
while a parallel `--approve --loop` is held. Assert the approve
|
|
completes only after the tick releases the lock.
|
|
|
|
Turned into a smoke-assertion: assert the `.state.lock` file
|
|
appears while the tick is mid-flight and the approve subprocess
|
|
blocks until the tick completes. Bounded by a generous timeout to
|
|
avoid CI flakiness.
|
|
"""
|
|
from types import SimpleNamespace
|
|
|
|
lp = _create_loop(tmp_project, name="ci-loop", template="ci-triage")
|
|
# Force `--approve` to have something to clear: halt the loop first.
|
|
s = _state(lp)
|
|
s["status"] = "halted"
|
|
s["halt_reason"] = "verifier_failed"
|
|
(lp / ".state.loop").write_text(json.dumps(s))
|
|
|
|
# Stub the harness (_invoke_harness) and gate (_gate) and context floor.
|
|
tick_started = threading.Event()
|
|
tick_can_finish = threading.Event()
|
|
|
|
def stub_invoke_harness(harness_cfg, role, prompt_ref, cwd, extras=None,
|
|
loop_path=None, tick_num=0):
|
|
if role == "implement":
|
|
tick_started.set()
|
|
tick_can_finish.wait(timeout=5.0)
|
|
# Return a passing-verdict-shaped stdout for the verify role so
|
|
# parse_verdict succeeds. For implement/orchestrate, an empty
|
|
# JSON object is enough.
|
|
if role == "verify":
|
|
return json.dumps({"pass": True, "score": 0.8,
|
|
"reasons": ["ok"], "next_hint": ""})
|
|
return "{}"
|
|
|
|
monkeypatch.setattr(runner_mod, "_invoke_harness", stub_invoke_harness)
|
|
|
|
def stub_gate(loop_path, loop_name, project_dir):
|
|
return {"ok": True, "reason": "running", "halt_reason": None,
|
|
"remaining_iterations": 25, "remaining_budget_usd": None,
|
|
"task_phase": None, "task_in_halt_loop": False,
|
|
"out_of_scope_files": []}
|
|
|
|
monkeypatch.setattr(runner_mod, "_gate", stub_gate)
|
|
monkeypatch.setattr(runner_mod, "_context_floor_ok", lambda: True)
|
|
|
|
# Stub _ensure_worktree to skip git worktree creation in the test.
|
|
monkeypatch.setattr(runner_mod, "_ensure_worktree",
|
|
lambda state, cfg, loop_path, project_dir: str(project_dir))
|
|
|
|
# Stub _find_work to return a fixed task name so the tick can proceed
|
|
# without an actual task dir existing.
|
|
monkeypatch.setattr(runner_mod, "_find_work",
|
|
lambda state, cfg, loop_path, project_dir:
|
|
("stub-task", None))
|
|
|
|
# Stub _read_task_brief / _acceptance_criteria_text / _next_hint_text
|
|
# to return empty strings (called by cmd_tick for substitution tokens).
|
|
monkeypatch.setattr(runner_mod, "_read_task_brief", lambda task_dir: "")
|
|
monkeypatch.setattr(runner_mod, "_acceptance_criteria_text", lambda cfg: "")
|
|
monkeypatch.setattr(runner_mod, "_next_hint_text", lambda state: "")
|
|
|
|
# But the tick also needs `current_task` referenced in harness extras;
|
|
# _role_prompt returns None for absent role config — that's fine, the
|
|
# stub_invoke_harness ignores the prompt arg.
|
|
|
|
args = SimpleNamespace(loop="ci-loop", project=str(tmp_project))
|
|
|
|
results = {}
|
|
|
|
def do_tick():
|
|
try:
|
|
runner_mod.cmd_tick(args)
|
|
results["tick"] = "done"
|
|
except Exception as exc:
|
|
results["tick_error"] = str(exc)
|
|
tick_can_finish.set() # in case the harness stub never advanced
|
|
|
|
def do_approve():
|
|
# Wait until the tick has reached its implement harness call
|
|
# (lock should be held by then).
|
|
tick_started.wait(timeout=3.0)
|
|
t0 = time.time()
|
|
out, err, code = _run(["--approve", "--loop", "ci-loop"], tmp_project)
|
|
results["approve_elapsed"] = time.time() - t0
|
|
results["approve_out"] = out
|
|
results["approve_code"] = code
|
|
|
|
tt = threading.Thread(target=do_tick)
|
|
ta = threading.Thread(target=do_approve)
|
|
tt.start()
|
|
ta.start()
|
|
# Let the tick reach its harness stub, then release it shortly after.
|
|
tick_started.wait(timeout=3.0)
|
|
time.sleep(0.1) # give approve subprocess time to spin up and block on flock
|
|
# Approve should NOT have completed yet (tick still holds the lock).
|
|
assert "approve_out" not in results, (
|
|
"approve completed while tick still held the lock")
|
|
tick_can_finish.set()
|
|
tt.join(timeout=5.0)
|
|
ta.join(timeout=5.0)
|
|
assert results.get("tick") == "done", results
|
|
assert results.get("approve_code") == 0, results
|
|
assert "Approved loop" in results["approve_out"] |