Files
automaton/tests/test_state_loop_lock.py

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"]