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