commit = findworkercommit(cfg.repodir, row.taskid, git=git)
The bug: when a parallel tick fires after a sibling session has already dispatched a task and launched the worker, the tick must adopt the sibling's work — not re-do it. A naive tick re-dispatches (duplicate worker), re-runs the judge (wasted, potentially conflicting verdict), or kills the worker it thinks it owns (orphans/duplicates state). The steward pattern fixes this by gating every mutation behind read-only detection and independent verification, then making exactly one idempotent completion commit.
The fix is a gated transaction (~/steward/steward.py, entry point stew()):
def stew(cfg: StewardConfig, verifiers=None, git=run_git) -> StewardOutcome:
with StewardLock(cfg.lock_path, timeout=cfg.lock_timeout): # 1 tick at a time
board_stale = is_board_stale(cfg.board_path, cfg.stale_after) # sibling dead? (mtime frozen >3-4 min)
orphans = find_orphaned_rows(load_board(cfg.board_path), board_stale)
if not board_stale:
return outcome("board fresh: sibling alive; standing down")
for row in orphans:
# GATE 1 — detect worker commit (read-only; never re-dispatch)
commit = find_worker_commit(cfg.repo_dir, row.task_id, git=git)
if commit is None:
return outcome(f"{row.task_id}: no worker commit; row left open")
# GATE 2 — independent verification: shim/parity/importer/trailer, all must pass
verification = verify_commit(commit, verifiers, cfg.verifier_stages)
if not verification.ok:
return outcome(f"verification failed at '{verification.stage}'; row left open")
# GATE 3 — reuse existing judge verdict; NEVER re-run the judge
verdict = load_judge_verdict(cfg.judge_verdict_path, row.task_id)
if verdict is None:
return outcome(f"no reusable verdict; row left open (judge NOT re-run)")
# GATE 4 — exists-checked completion (idempotent, refuses missing artifacts)
exists_checked_complete(cfg.workdir, row, commit, verdict)
# GATE 5 — close orphaned verification tasks (only this task's pending ones)
outcome.tasks_closed += close_orphaned_verification_tasks(cfg.verification_tasks_path, row.task_id)
cfg.board_path.write_text("\n".join(r.to_line() for r in rows) + "\n")
outcome.pushed = commit_and_push(cfg.repo_dir, [cfg.board_path, cfg.verification_tasks_path],
f"steward: adopt orphaned {row.task_id} @ {commit[:12]}", git=git)
return outcome(f"adopted worker commit {commit[:12]} (verdict {verdict.ref} reused)")
Key defensive pieces:
# staleness: a live sibling keeps touching the board JSONL
def is_board_stale(board_path, stale_after=180, now=None):
return board_staleness_seconds(board_path, now) >= stale_after
# worker-commit detection anchored to the Task-ID trailer, so a baseline/
# dispatch commit that merely *mentions* the task id is never adopted
def find_worker_commit(repo_dir, task_id, git=run_git):
out = git(repo_dir, "log", "--all", "--format=%H", "--grep", f"Task-ID: {task_id}")
return out.splitlines()[0] if out.strip() else None
# verification short-circuits on the first failing stage
def verify_commit(commit, verifiers, stages=("shim","parity","importer","trailer")):
for stage in stages:
res = verifiers[stage](commit)
if not res.ok: return res # e.g. parity: input/output hashes differ
return VerificationResult(True, stages[-1])
# completion only after artifacts exist on disk; no-op when already complete
def exists_checked_complete(workdir, row, commit, verdict, extra_markers=()):
for p in [workdir / "result.json", *extra_markers]:
if not p.exists():
raise FileNotFoundError(f"exists-check failed: {p} missing")
row.status, row.worker_commit, row.verdict_ref = "complete", commit, verdict.ref
row.extra |= {"completed_by": "steward", "verdict_reused": True, "completed_at": int(time.time())}
# flock: a second parallel tick waits/aborts instead of racing (fd never leaks)
class StewardLock:
def __enter__(self):
self._fh = open(self.lock_path, "a+")
try:
while True:
try: fcntl.flock(self._fh.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB); return self
except OSError:
if time.monotonic() >= self.deadline: raise TimeoutError(...)
time.sleep(0.25)
except BaseException:
self._fh.close(); self._fh = None; raise
The do-NOT list is enforced structurally, not just by convention: dispatch and judge are never invoked anywhere in stew.py — the verdict is only read (load_judge_verdict returns None instead of running the judge), worker pids are never signaled, and commit_and_push returns False when nothing is staged so a second tick pushes nothing.
Verified in `~/steward/` with a simulation (`test_steward.py`) that reproduces the exact production layout — git repo + bare origin, board JSONL, verification-tasks JSONL, gitreins judge verdict, dispatch/judge audit logs, and a real live worker process (`sleep 300`, pid tracked). Run: `python3 -W error::ResourceWarning test_steward.py` → **22/22 OK**.
| Test | Asserts |
|---|---|
| `test_adopts_orphaned_row_without_redispatch_or_rejudge` | row → `complete` with worker commit + reused verdict; 4 verification tasks closed; pushed to `origin/main`; `dispatch.log`=1 (no re-dispatch), `judge_runs.log`=1 (no re-run), worker pid still alive (not killed) |
| `test_board_fresh_means_sibling_alive_stand_down` | fresh board → stand down, row untouched, no push |
| `test_no_worker_commit_no_completion_no_redispatch` | dead sibling, no commit → row open, no push, no re-dispatch |
| `test_missing_judge_verdict_leaves_row_open_no_rejudge` | verdict file absent → row open, judge log still 1, no push |
| `test_failed_verification_blocks_completion` | parity mismatch (rows_in≠rows_out) → blocked at `parity` |
| `test_importer_failure_blocks_completion` | `importable:false` → blocked at `importer` |
| `test_commit_without_task_id_trailer_not_adopted` | no `Task-ID:` trailer → not detected as worker output |
| `test_exists_check_blocks_completion_when_artifact_missing` | artifact deleted after commit → graceful `shim` failure, row open |
| `test_already_complete_row_is_idempotent_noop` | second tick: no adoption, no push, `origin/main` unchanged |
| `test_double_tick_does_not_double_complete` | two consecutive ticks → exactly one adoption, dispatch/judge each ran once |
| `test_lock_prevents_concurrent_stewards` | held lock → second steward raises `TimeoutError`; succeeds after release |
| `test_orphaned_verification_tasks_only_touched_for_task` | only T-42 pending tasks closed; T-99 and already-closed untouched |
| `test_worker_alive_when_already_exited_naturally` | worker gone on its own → still adopted, no crash |
| Unit tests | staleness boundary (2:59 not stale / 3:01 and 4:00 stale), orphan requires *both* stale board *and* in_progress, verdict task-id mismatch not reused, verification short-circuits after first failure, close-tasks idempotent, exists-check raises on missing marker, trailer gate rejects bad commit |
Live demo (`demo.py`) output of the full transaction:
```
BEFORE: origin/main:board.jsonl = {"status": "in_progress", "task_id": "T-42"} worker alive, dispatch=1, judge=1
AFTER: outcome.adopted=True pushed=True tasks_closed=4
reason: T-42: adopted worker commit bb2fbd813b6a (verdict gitreins/T-42#1 reused, judge not re-run)
origin/main:board.jsonl = {"status": "complete", "task_id": "T-42", "worker_commit": "bb2fbd…",
"verdict_ref": "gitreins/T-42#1", "verdict_reused": true, "completed_by": "steward"}
4 verification tasks → {"status": "closed", "closed_by": "steward"}
worker pid alive (NOT killed), dispatch.log=1 (no re-dispatch), judge_runs.log=1 (no re-run)
```{"model": "deepseek-v4-flash", "problem_class": "foreman-parallel-tick-stewardship", "result": "passed", "tests": 22}The bug: when a parallel tick fires after a sibling session has already dispatched a task and launched the worker, the tick must adopt the sibling's work — not re-do it. A naive tick re-dispatches (duplicate worker), re-runs the judge (wasted, potentially conflicting verdict), or kills the worker it thinks it owns (orphans/duplicates state). The steward pattern fixes this by gating every mutation behind read-only detection and independent verification, then making exactly one idempotent completion commit.
The fix is a gated transaction (~/steward/steward.py, entry point stew()):
def stew(cfg: StewardConfig, verifiers=None, git=run_git) -> StewardOutcome:
with StewardLock(cfg.lock_path, timeout=cfg.lock_timeout): # 1 tick at a time
board_stale = is_board_stale(cfg.board_path, cfg.stale_after) # sibling dead? (mtime frozen >3-4 min)
orphans = find_orphaned_rows(load_board(cfg.board_path), board_stale)
if not board_stale:
return outcome("board fresh: sibling alive; standing down")
for row in orphans:
# GATE 1 — detect worker commit (read-only; never re-dispatch)
commit = find_worker_commit(cfg.repo_dir, row.task_id, git=git)
if commit is None:
return outcome(f"{row.task_id}: no worker commit; row left open")
# GATE 2 — independent verification: shim/parity/importer/trailer, all must pass
verification = verify_commit(commit, verifiers, cfg.verifier_stages)
if not verification.ok:
return outcome(f"verification failed at '{verification.stage}'; row left open")
# GATE 3 — reuse existing judge verdict; NEVER re-run the judge
verdict = load_judge_verdict(cfg.judge_verdict_path, row.task_id)
if verdict is None:
return outcome(f"no reusable verdict; row left open (judge NOT re-run)")
# GATE 4 — exists-checked completion (idempotent, refuses missing artifacts)
exists_checked_complete(cfg.workdir, row, commit, verdict)
# GATE 5 — close orphaned verification tasks (only this task's pending ones)
outcome.tasks_closed += close_orphaned_verification_tasks(cfg.verification_tasks_path, row.task_id)
cfg.board_path.write_text("\n".join(r.to_line() for r in rows) + "\n")
outcome.pushed = commit_and_push(cfg.repo_dir, [cfg.board_path, cfg.verification_tasks_path],
f"steward: adopt orphaned {row.task_id} @ {commit[:12]}", git=git)
return outcome(f"adopted worker commit {commit[:12]} (verdict {verdict.ref} reused)")
Key defensive pieces:
# staleness: a live sibling keeps touching the board JSONL
def is_board_stale(board_path, stale_after=180, now=None):
return board_staleness_seconds(board_path, now) >= stale_after
# worker-commit detection anchored to the Task-ID trailer, so a baseline/
# dispatch commit that merely *mentions* the task id is never adopted
def find_worker_commit(repo_dir, task_id, git=run_git):
out = git(repo_dir, "log", "--all", "--format=%H", "--grep", f"Task-ID: {task_id}")
return out.splitlines()[0] if out.strip() else None
# verification short-circuits on the first failing stage
def verify_commit(commit, verifiers, stages=("shim","parity","importer","trailer")):
for stage in stages:
res = verifiers[stage](commit)
if not res.ok: return res # e.g. parity: input/output hashes differ
return VerificationResult(True, stages[-1])
# completion only after artifacts exist on disk; no-op when already complete
def exists_checked_complete(workdir, row, commit, verdict, extra_markers=()):
for p in [workdir / "result.json", *extra_markers]:
if not p.exists():
raise FileNotFoundError(f"exists-check failed: {p} missing")
row.status, row.worker_commit, row.verdict_ref = "complete", commit, verdict.ref
row.extra |= {"completed_by": "steward", "verdict_reused": True, "completed_at": int(time.time())}
# flock: a second parallel tick waits/aborts instead of racing (fd never leaks)
class StewardLock:
def __enter__(self):
self._fh = open(self.lock_path, "a+")
try:
while True:
try: fcntl.flock(self._fh.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB); return self
except OSError:
if time.monotonic() >= self.deadline: raise TimeoutError(...)
time.sleep(0.25)
except BaseException:
self._fh.close(); self._fh = None; raise
The do-NOT list is enforced structurally, not just by convention: dispatch and judge are never invoked anywhere in stew.py — the verdict is only read (load_judge_verdict returns None instead of running the judge), worker pids are never signaled, and commit_and_push returns False when nothing is staged so a second tick pushes nothing.
Verified in `~/steward/` with a simulation (`test_steward.py`) that reproduces the exact production layout — git repo + bare origin, board JSONL, verification-tasks JSONL, gitreins judge verdict, dispatch/judge audit logs, and a real live worker process (`sleep 300`, pid tracked). Run: `python3 -W error::ResourceWarning test_steward.py` → **22/22 OK**.
| Test | Asserts |
|---|---|
| `test_adopts_orphaned_row_without_redispatch_or_rejudge` | row → `complete` with worker commit + reused verdict; 4 verification tasks closed; pushed to `origin/main`; `dispatch.log`=1 (no re-dispatch), `judge_runs.log`=1 (no re-run), worker pid still alive (not killed) |
| `test_board_fresh_means_sibling_alive_stand_down` | fresh board → stand down, row untouched, no push |
| `test_no_worker_commit_no_completion_no_redispatch` | dead sibling, no commit → row open, no push, no re-dispatch |
| `test_missing_judge_verdict_leaves_row_open_no_rejudge` | verdict file absent → row open, judge log still 1, no push |
| `test_failed_verification_blocks_completion` | parity mismatch (rows_in≠rows_out) → blocked at `parity` |
| `test_importer_failure_blocks_completion` | `importable:false` → blocked at `importer` |
| `test_commit_without_task_id_trailer_not_adopted` | no `Task-ID:` trailer → not detected as worker output |
| `test_exists_check_blocks_completion_when_artifact_missing` | artifact deleted after commit → graceful `shim` failure, row open |
| `test_already_complete_row_is_idempotent_noop` | second tick: no adoption, no push, `origin/main` unchanged |
| `test_double_tick_does_not_double_complete` | two consecutive ticks → exactly one adoption, dispatch/judge each ran once |
| `test_lock_prevents_concurrent_stewards` | held lock → second steward raises `TimeoutError`; succeeds after release |
| `test_orphaned_verification_tasks_only_touched_for_task` | only T-42 pending tasks closed; T-99 and already-closed untouched |
| `test_worker_alive_when_already_exited_naturally` | worker gone on its own → still adopted, no crash |
| Unit tests | staleness boundary (2:59 not stale / 3:01 and 4:00 stale), orphan requires *both* stale board *and* in_progress, verdict task-id mismatch not reused, verification short-circuits after first failure, close-tasks idempotent, exists-check raises on missing marker, trailer gate rejects bad commit |
Live demo (`demo.py`) output of the full transaction:
```
BEFORE: origin/main:board.jsonl = {"status": "in_progress", "task_id": "T-42"} worker alive, dispatch=1, judge=1
AFTER: outcome.adopted=True pushed=True tasks_closed=4
reason: T-42: adopted worker commit bb2fbd813b6a (verdict gitreins/T-42#1 reused, judge not re-run)
origin/main:board.jsonl = {"status": "complete", "task_id": "T-42", "worker_commit": "bb2fbd…",
"verdict_ref": "gitreins/T-42#1", "verdict_reused": true, "completed_by": "steward"}
4 verification tasks → {"status": "closed", "closed_by": "steward"}
worker pid alive (NOT killed), dispatch.log=1 (no re-dispatch), judge_runs.log=1 (no re-run)
```{"model": "deepseek-v4-flash", "problem_class": "foreman-parallel-tick-stewardship", "result": "passed", "tests": 22}