hermes: close terminal replay edge cases

This commit is contained in:
jenkins 2026-08-16 20:55:46 -03:00
parent 18be595700
commit 8e0a13e9db
4 changed files with 669 additions and 71 deletions

View File

@ -146,6 +146,68 @@ class AutomaticDecompositionSafetyTests(unittest.TestCase):
self.assertEqual(completions.calls, 1)
self.assertEqual(self._status(task_id), "ready")
def test_latest_ended_run_can_replay_a_durable_completion(self) -> None:
task_id = kanban_db.create_task(self.connection, title="journaled result")
self.assertIsNotNone(kanban_db.claim_task(self.connection, task_id))
run_id = kanban_db.get_task(self.connection, task_id).current_run_id
self.assertTrue(
kanban_db.block_task(
self.connection,
task_id,
reason="legacy post-journal failure",
kind="capability",
expected_run_id=run_id,
)
)
self.assertTrue(
kanban_db.complete_task(
self.connection,
task_id,
result="durable terminal result",
summary="durable terminal result",
replay_ended_run_id=run_id,
)
)
self.assertEqual(self._status(task_id), "done")
def test_older_ended_run_cannot_complete_over_a_replacement(self) -> None:
task_id = kanban_db.create_task(self.connection, title="replacement guard")
self.assertIsNotNone(kanban_db.claim_task(self.connection, task_id))
old_run = kanban_db.get_task(self.connection, task_id).current_run_id
self.assertTrue(
kanban_db.block_task(
self.connection,
task_id,
reason="first run",
kind="capability",
expected_run_id=old_run,
)
)
self.assertTrue(kanban_db.unblock_task(self.connection, task_id))
self.assertIsNotNone(kanban_db.claim_task(self.connection, task_id))
replacement_run = kanban_db.get_task(self.connection, task_id).current_run_id
self.assertNotEqual(old_run, replacement_run)
self.assertTrue(
kanban_db.block_task(
self.connection,
task_id,
reason="replacement run",
kind="capability",
expected_run_id=replacement_run,
)
)
self.assertFalse(
kanban_db.complete_task(
self.connection,
task_id,
result="stale result",
replay_ended_run_id=old_run,
)
)
self.assertIn(self._status(task_id), {"blocked", "triage"})
def test_fresh_triage_task_still_auto_promotes(self) -> None:
task_id = kanban_db.create_task(
self.connection,

View File

@ -1,4 +1,4 @@
"""Patch upstream Hermes auto-decomposition to respect execution history."""
"""Patch upstream Hermes exact-run completion and decomposition safety."""
import os
from pathlib import Path
@ -18,6 +18,112 @@ source_root = Path(os.environ.get("HERMES_SOURCE_ROOT", "/opt/hermes"))
db_path = source_root / "hermes_cli/kanban_db.py"
db = db_path.read_text(encoding="utf-8")
complete_signature_before = '''def complete_task(
conn: sqlite3.Connection,
task_id: str,
*,
result: Optional[str] = None,
summary: Optional[str] = None,
metadata: Optional[dict] = None,
created_cards: Optional[Iterable[str]] = None,
expected_run_id: Optional[int] = None,
) -> bool:
'''
complete_signature_after = '''def complete_task(
conn: sqlite3.Connection,
task_id: str,
*,
result: Optional[str] = None,
summary: Optional[str] = None,
metadata: Optional[dict] = None,
created_cards: Optional[Iterable[str]] = None,
expected_run_id: Optional[int] = None,
replay_ended_run_id: Optional[int] = None,
) -> bool:
'''
db = replace_once(
db,
complete_signature_before,
complete_signature_after,
"exact ended-run completion signature",
)
complete_branch_before = ''' with write_txn(conn):
if expected_run_id is None:
cur = conn.execute(
"""
UPDATE tasks
SET status = 'done',
result = ?,
completed_at = ?,
claim_lock = NULL,
claim_expires= NULL,
worker_pid = NULL,
block_kind = NULL,
block_recurrences = 0
WHERE id = ?
AND status IN ('running', 'ready', 'blocked', 'scheduled')
""",
(result, now, task_id),
)
else:
'''
complete_branch_after = ''' # Journal recovery may arrive after a legacy error path ended the
# accepted run. Verify that it is still the latest run in this same
# transaction; a replacement attempt must make the replay fail closed.
with write_txn(conn):
if expected_run_id is not None and replay_ended_run_id is not None:
raise ValueError("expected_run_id and replay_ended_run_id are mutually exclusive")
if replay_ended_run_id is not None:
cur = conn.execute(
"""
UPDATE tasks
SET status = 'done',
result = ?,
completed_at = ?,
claim_lock = NULL,
claim_expires= NULL,
worker_pid = NULL,
block_kind = NULL,
block_recurrences = 0
WHERE id = ?
AND status IN ('ready', 'blocked')
AND current_run_id IS NULL
AND ? = (
SELECT id FROM task_runs
WHERE task_id = ?
ORDER BY id DESC
LIMIT 1
)
""",
(result, now, task_id, int(replay_ended_run_id), task_id),
)
elif expected_run_id is None:
cur = conn.execute(
"""
UPDATE tasks
SET status = 'done',
result = ?,
completed_at = ?,
claim_lock = NULL,
claim_expires= NULL,
worker_pid = NULL,
block_kind = NULL,
block_recurrences = 0
WHERE id = ?
AND status IN ('running', 'ready', 'blocked', 'scheduled')
""",
(result, now, task_id),
)
else:
'''
db = replace_once(
db,
complete_branch_before,
complete_branch_after,
"transactional exact ended-run completion",
)
specify_signature_before = '''def specify_triage_task(
conn: sqlite3.Connection,
task_id: str,

View File

@ -522,10 +522,27 @@ def _terminal_record_valid(
structured = json.loads(record["result"])
except (TypeError, json.JSONDecodeError):
return False
if not isinstance(structured, dict):
return False
required = set(RESULT_SCHEMA["required"])
if set(structured) != required:
return False
if (
type(structured["status"]) is not str
or structured["status"] not in cli_lane_goal.RESULT_STATUSES
or type(structured["summary"]) is not str
or not structured["summary"].strip()
):
return False
for key in ("changed_files", "tests_run", "artifacts", "findings", "blockers"):
value = structured[key]
if type(value) is not list or any(type(item) is not str for item in value):
return False
valid = (
isinstance(structured, dict)
and structured.get("status") == "completed"
and not structured.get("blockers")
structured["status"] == "completed"
and structured["blockers"] == []
and cli_lane_goal.unfinished_result_reason(structured) is None
and record["summary"] == structured["summary"]
and record.get("kanban_state") in {"pending", "committed"}
)
if not valid or identity is None:
@ -1077,7 +1094,7 @@ def _board_call(
def _quarantine_terminal(path: Path, identity: TerminalIdentity | None, reason: str) -> Path:
"""Move an unusable journal aside with a deterministic, private name."""
"""Copy an unusable journal to a private file, then retire its exact name."""
board_dir = path.parent
try:
board_stat = board_dir.stat(follow_symlinks=False)
@ -1090,44 +1107,135 @@ def _quarantine_terminal(path: Path, identity: TerminalIdentity | None, reason:
flush=True,
)
return path
quarantine = board_dir / "quarantine"
quarantine.mkdir(parents=True, exist_ok=True, mode=0o700)
os.chmod(quarantine, 0o700, follow_symlinks=False)
board_flags = os.O_RDONLY | getattr(os, "O_DIRECTORY", 0)
board_flags |= getattr(os, "O_NOFOLLOW", 0)
file_flags = os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0)
create_flags = os.O_WRONLY | os.O_CREAT | os.O_EXCL
create_flags |= getattr(os, "O_NOFOLLOW", 0)
board_descriptor = None
quarantine_descriptor = None
source_stat = None
destination = path
retired = False
try:
source_stat = path.stat(follow_symlinks=False)
regular = stat.S_ISREG(source_stat.st_mode)
except OSError:
regular = False
try:
if not regular:
raise OSError("journal is not a regular file")
flags = os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0)
descriptor = os.open(path, flags)
with os.fdopen(descriptor, "rb") as stream:
fingerprint = hashlib.sha256(stream.read()).hexdigest()[:16]
except OSError:
fingerprint = hashlib.sha256(path.name.encode("utf-8")).hexdigest()[:16]
safe_reason = re.sub(r"[^a-zA-Z0-9_.-]+", "-", reason).strip("-") or "invalid"
destination = quarantine / f"{path.name}.{safe_reason}.{fingerprint}.quarantine"
source_exists = path.exists() or path.is_symlink()
if source_exists and regular and not destination.exists():
os.chmod(path, 0o600, follow_symlinks=False)
os.replace(path, destination)
os.chmod(destination, 0o600, follow_symlinks=False)
_fsync_directory(board_dir)
_fsync_directory(quarantine)
elif source_exists:
path.unlink()
_fsync_directory(board_dir)
if not destination.exists():
atomic_json(
destination,
{
"original_name": path.name,
"reason": reason,
"recorded_at": utc_now(),
},
board_descriptor = os.open(board_dir, board_flags)
source_stat = os.stat(
path.name,
dir_fd=board_descriptor,
follow_symlinks=False,
)
safe_source = stat.S_ISREG(source_stat.st_mode) and source_stat.st_nlink == 1
payload: bytes
if safe_source:
source_descriptor = os.open(
path.name,
file_flags,
dir_fd=board_descriptor,
)
with os.fdopen(source_descriptor, "rb") as source:
opened_stat = os.fstat(source.fileno())
if (
opened_stat.st_dev != source_stat.st_dev
or opened_stat.st_ino != source_stat.st_ino
or opened_stat.st_nlink != 1
):
raise OSError("terminal journal identity changed while opening")
payload = source.read()
else:
payload = (
json.dumps(
{
"original_name": path.name,
"reason": reason,
"source_kind": "hardlink"
if stat.S_ISREG(source_stat.st_mode) and source_stat.st_nlink != 1
else "non-regular",
},
sort_keys=True,
)
+ "\n"
).encode("utf-8")
fingerprint = hashlib.sha256(payload).hexdigest()[:16]
safe_reason = (
re.sub(r"[^a-zA-Z0-9_.-]+", "-", reason).strip("-") or "invalid"
)
try:
os.mkdir("quarantine", 0o700, dir_fd=board_descriptor)
except FileExistsError:
pass
quarantine_descriptor = os.open(
"quarantine",
board_flags,
dir_fd=board_descriptor,
)
for sequence in range(32):
name = f"{path.name}.{safe_reason}.{fingerprint}.{sequence}.quarantine"
try:
destination_descriptor = os.open(
name,
create_flags,
0o600,
dir_fd=quarantine_descriptor,
)
except FileExistsError:
continue
try:
with os.fdopen(destination_descriptor, "wb") as target:
target.write(payload)
target.flush()
os.fsync(target.fileno())
target_stat = os.fstat(target.fileno())
if (
stat.S_IMODE(target_stat.st_mode) != 0o600
or target_stat.st_nlink != 1
):
raise OSError("quarantine destination is not private and singly linked")
except Exception:
os.unlink(name, dir_fd=quarantine_descriptor)
raise
destination = board_dir / "quarantine" / name
break
else:
raise OSError("could not reserve a unique quarantine destination")
os.fsync(quarantine_descriptor)
current = os.stat(
path.name,
dir_fd=board_descriptor,
follow_symlinks=False,
)
if (
current.st_dev != source_stat.st_dev
or current.st_ino != source_stat.st_ino
):
print(
f"terminal journal changed before retirement; removing new entry: {path}",
file=sys.stderr,
flush=True,
)
os.unlink(path.name, dir_fd=board_descriptor)
retired = True
os.fsync(board_descriptor)
except OSError as error:
# A bad quarantine target must not make an attacker-controlled pending
# path eligible for replay forever. Retire only the directory entry;
# never chmod or write through the source inode.
if board_descriptor is not None and not retired:
try:
os.unlink(path.name, dir_fd=board_descriptor)
retired = True
os.fsync(board_descriptor)
except OSError:
pass
print(
f"terminal journal quarantine degraded safely: {type(error).__name__}: {error}",
file=sys.stderr,
flush=True,
)
finally:
if quarantine_descriptor is not None:
os.close(quarantine_descriptor)
if board_descriptor is not None:
os.close(board_descriptor)
print(
"quarantined terminal journal "
f"{path.name}: {reason}; identity={identity or 'unparseable'}",
@ -1137,12 +1245,12 @@ def _quarantine_terminal(path: Path, identity: TerminalIdentity | None, reason:
return destination
def _recover_quarantined_run(
def _recover_exact_run(
kanban_db: Any,
identity: TerminalIdentity | None,
reason: str,
) -> bool:
"""Make an exact external run retryable after its journal is quarantined."""
"""Make an exact external run retryable after journal recovery fails."""
if identity is None:
return False
@ -1160,7 +1268,7 @@ def _recover_quarantined_run(
kanban_db.reclaim_task(
conn,
identity.task_id,
reason=f"terminal journal quarantined ({reason}); exact run may retry",
reason=f"terminal journal recovery failed ({reason}); exact run may retry",
)
)
@ -1172,7 +1280,7 @@ def _recover_quarantined_run(
if recovered:
print(
f"reclaimed {identity.board}/{identity.task_id} run {identity.run_id} "
f"after terminal journal quarantine: {reason}",
f"after terminal journal recovery: {reason}",
file=sys.stderr,
flush=True,
)
@ -1221,10 +1329,16 @@ def _finalize_terminal_record(
if str(_task_value(task, "result", "") or "") == document["result"]
else "stale"
)
if (
status != "running"
or _task_value(task, "current_run_id", None) != identity.run_id
):
current_run_id = _task_value(task, "current_run_id", None)
completion_guard: dict[str, int]
if status == "running" and current_run_id == identity.run_id:
completion_guard = {"expected_run_id": identity.run_id}
elif status in {"ready", "blocked"} and current_run_id is None:
# The journal may have survived an older post-persistence error
# path that ended its run. The patched DB verifies atomically that
# this is still the latest ended run before allowing completion.
completion_guard = {"replay_ended_run_id": identity.run_id}
else:
return "stale"
completed = kanban_db.complete_task(
conn,
@ -1232,7 +1346,7 @@ def _finalize_terminal_record(
result=document["result"],
summary=document["summary"],
metadata=document["metadata"],
expected_run_id=identity.run_id,
**completion_guard,
)
return "committed" if completed else "pending"
@ -1255,11 +1369,11 @@ def recover_pending_finalizations() -> int:
continue
if not _terminal_record_valid(record):
_quarantine_terminal(path, identity, "malformed-payload")
_recover_quarantined_run(kanban_db, identity, "malformed-payload")
_recover_exact_run(kanban_db, identity, "malformed-payload")
continue
if not _terminal_record_valid(record, identity):
_quarantine_terminal(path, identity, "foreign-identity")
_recover_quarantined_run(kanban_db, identity, "foreign-identity")
_recover_exact_run(kanban_db, identity, "foreign-identity")
continue
try:
outcome = _finalize_terminal_record(kanban_db, path, record)
@ -1270,7 +1384,7 @@ def recover_pending_finalizations() -> int:
recovered += 1
elif outcome in {"invalid", "foreign", "stale"}:
_quarantine_terminal(path, identity, outcome)
_recover_quarantined_run(kanban_db, identity, outcome)
_recover_exact_run(kanban_db, identity, outcome)
return recovered
@ -1329,6 +1443,9 @@ def gc_lane_artifacts(
except OSError:
continue
if not stat.S_ISREG(file_stat.st_mode):
path.unlink()
removed += 1
board_removed += 1
continue
if max_age_seconds >= 0 and current - file_stat.st_mtime > max_age_seconds:
path.unlink()
@ -1631,28 +1748,44 @@ def execute_claim(board: str, task_id: str) -> None:
and result.returncode == 0
and completion_problem is None
):
terminal_file = _terminal_path(state_file, run_id, "pending")
metadata["terminal_record"] = str(
_terminal_path(state_file, run_id, "committed")
)
terminal_file, terminal_record = _write_terminal_record(
state_file,
board=board,
task_id=task_id,
run_id=run_id,
structured=structured,
summary=str(structured.get("summary") or "Completed"),
metadata=metadata,
)
try:
terminal_file, terminal_record = _write_terminal_record(
state_file,
board=board,
task_id=task_id,
run_id=run_id,
structured=structured,
summary=str(structured.get("summary") or "Completed"),
metadata=metadata,
)
outcome = _finalize_terminal_record(
kanban_db,
terminal_file,
terminal_record,
)
except Exception as error:
identity = _terminal_identity(terminal_file)
replayable = bool(
identity is not None
and _terminal_record_valid(
_load_terminal_json(terminal_file, identity),
identity,
)
)
if not replayable:
_recover_exact_run(
kanban_db,
identity,
f"persistence raised {type(error).__name__}",
)
raise TerminalFinalizationPending(
"accepted worker result is durably journaled; "
f"Kanban finalization raised {type(error).__name__}"
"accepted worker result remains outside the generic block path; "
f"terminal replayable={replayable}; persistence/finalization "
f"raised {type(error).__name__}"
) from error
if outcome != "committed":
raise TerminalFinalizationPending(

View File

@ -2,7 +2,9 @@
from __future__ import annotations
import errno
import importlib.util
import hashlib
import json
import os
import signal
@ -864,6 +866,55 @@ def _completed_result(summary: str = "done") -> dict:
}
def test_terminal_recovery_can_complete_the_exact_latest_ended_run(
tmp_path: Path,
monkeypatch,
):
monkeypatch.setattr(lanes, "STATE_ROOT", tmp_path / "cli-lanes")
state_file = lanes.state_path("cassandra", "t_ended")
pending, _record = lanes._write_terminal_record(
state_file,
board="cassandra",
task_id="t_ended",
run_id=17,
structured=_completed_result("accepted before legacy block"),
summary="accepted before legacy block",
metadata={},
)
task = SimpleNamespace(
id="t_ended",
status="blocked",
result=None,
current_run_id=None,
assignee="cli-auto",
)
guards = []
class Connection:
def close(self):
return None
def complete_task(_conn, _task_id, **kwargs):
guards.append(kwargs)
task.status = "done"
task.result = kwargs["result"]
return True
fake_db = SimpleNamespace(
scoped_current_board=lambda _board: nullcontext(),
connect=lambda board: Connection(),
get_task=lambda _conn, _task_id: task,
complete_task=complete_task,
)
monkeypatch.setitem(sys.modules, "hermes_cli", SimpleNamespace(kanban_db=fake_db))
assert lanes.recover_pending_finalizations() == 1
assert guards[0]["replay_ended_run_id"] == 17
assert "expected_run_id" not in guards[0]
assert task.status == "done"
assert not pending.exists()
def test_complete_exception_after_journal_is_replayable_and_never_blocks(
tmp_path: Path,
monkeypatch,
@ -937,6 +988,103 @@ def test_complete_exception_after_journal_is_replayable_and_never_blocks(
assert not pending[0].exists()
def test_terminal_replace_then_directory_fsync_enospc_replays_after_restart(
tmp_path: Path,
monkeypatch,
):
state_root = tmp_path / "cli-lanes"
monkeypatch.setattr(lanes, "STATE_ROOT", state_root)
task = SimpleNamespace(
id="t_enospc",
status="running",
result=None,
current_run_id=52,
assignee="cli-auto",
max_runtime_seconds=60,
)
blocks = []
comments = []
completions = []
class Connection:
def close(self):
return None
def complete_task(_conn, _task_id, **kwargs):
completions.append(kwargs)
task.status = "done"
task.result = kwargs["result"]
task.current_run_id = None
return True
fake_db = SimpleNamespace(
scoped_current_board=lambda _board: nullcontext(),
connect=lambda board: Connection(),
get_task=lambda _conn, _task_id: task,
worker_log_path=lambda _task_id, board: tmp_path / "worker.log",
_resolve_worktree_workspace=lambda _task, board: (tmp_path, "wt/t_enospc"),
set_branch_name=lambda *_args: None,
set_workspace_path=lambda *_args: None,
build_worker_context=lambda *_args: "Finish without losing the result.",
heartbeat_worker=lambda *_args, **_kwargs: True,
add_comment=lambda _conn, _task_id, _author, body: comments.append(body),
complete_task=complete_task,
block_task=lambda *_args, **kwargs: blocks.append(kwargs),
)
monkeypatch.setitem(sys.modules, "hermes_cli", SimpleNamespace(kanban_db=fake_db))
monkeypatch.setattr(
lanes,
"select_route",
lambda *_args, **_kwargs: lanes.Route(
"codex", "gpt-5.6-sol", "high", "codex-high", "test", "test", 1, ()
),
)
monkeypatch.setattr(
lanes,
"run_provider",
lambda *_args, **_kwargs: lanes.ProcessResult(
0, "", _completed_result("persisted before ENOSPC"), False
),
)
real_replace = lanes.os.replace
real_fsync_directory = lanes._fsync_directory
terminal_replaced = {"value": False}
fail_once = {"value": True}
def replace_then_mark(source, destination):
real_replace(source, destination)
if str(destination).endswith(".terminal.pending.json"):
terminal_replaced["value"] = True
def fail_after_terminal_replace(directory):
if terminal_replaced["value"] and fail_once["value"]:
fail_once["value"] = False
raise OSError(errno.ENOSPC, "no space after terminal rename")
real_fsync_directory(directory)
monkeypatch.setattr(lanes.os, "replace", replace_then_mark)
monkeypatch.setattr(lanes, "_fsync_directory", fail_after_terminal_replace)
lanes.execute_claim("cassandra", "t_enospc")
pending = list((state_root / "cassandra").glob("*.terminal.pending.json"))
assert terminal_replaced["value"] is True
assert fail_once["value"] is False
assert len(pending) == 1
assert blocks == []
assert completions == []
assert task.status == "running"
assert any("terminal replayable=True" in body for body in comments)
# A restarted runner sees the exact pending journal and completes the run.
assert lanes.recover_pending_finalizations() == 1
assert len(completions) == 1
assert task.status == "done"
assert not pending[0].exists()
committed = list((state_root / "cassandra").glob("*.terminal.committed.json"))
assert len(committed) == 1
def test_terminal_payload_cannot_select_a_different_board(
tmp_path: Path,
monkeypatch,
@ -1050,6 +1198,147 @@ def test_malformed_exact_run_journal_is_quarantined_and_reclaimed(
assert quarantined[0].stat().st_mode & 0o777 == 0o600
@pytest.mark.parametrize(
"structured",
[
{"status": "completed", "blockers": {}},
{**_completed_result(), "changed_files": "src/a.py"},
{**_completed_result(), "unexpected": True},
{**_completed_result(), "blockers": ["work remains"]},
{**_completed_result("tests are still running")},
],
)
def test_terminal_record_requires_exact_completed_result_contract(structured):
record = {
"board": "cassandra",
"task_id": "t_schema",
"expected_run_id": 1,
"result": json.dumps(structured),
"summary": str(structured.get("summary") or "done"),
"metadata": {},
"kanban_state": "pending",
}
assert lanes._terminal_record_valid(record) is False
def test_minimal_completed_result_with_mapping_blockers_quarantines_without_completion(
tmp_path: Path,
monkeypatch,
):
monkeypatch.setattr(lanes, "STATE_ROOT", tmp_path / "cli-lanes")
path = lanes._terminal_path(lanes.state_path("cassandra", "t_schema"), 6)
lanes.atomic_json(
path,
{
"board": "cassandra",
"task_id": "t_schema",
"expected_run_id": 6,
"result": json.dumps({"status": "completed", "blockers": {}}),
"summary": "done",
"metadata": {},
"kanban_state": "pending",
},
)
task = SimpleNamespace(
id="t_schema", status="running", current_run_id=6, assignee="cli-auto"
)
completed = []
reclaimed = []
class Connection:
def close(self):
return None
fake_db = SimpleNamespace(
scoped_current_board=lambda _board: nullcontext(),
connect=lambda board: Connection(),
get_task=lambda _conn, _task_id: task,
complete_task=lambda *_args, **_kwargs: completed.append(True),
reclaim_task=lambda *_args, **_kwargs: (reclaimed.append(True) or True),
)
monkeypatch.setitem(sys.modules, "hermes_cli", SimpleNamespace(kanban_db=fake_db))
assert lanes.recover_pending_finalizations() == 0
assert completed == []
assert reclaimed == [True]
assert not path.exists()
quarantined = list((path.parent / "quarantine").glob("*.quarantine"))
assert len(quarantined) == 1
assert quarantined[0].stat().st_mode & 0o777 == 0o600
def test_quarantine_avoids_symlink_and_mode_collision_destinations(
tmp_path: Path,
monkeypatch,
):
monkeypatch.setattr(lanes, "STATE_ROOT", tmp_path / "cli-lanes")
path = lanes._terminal_path(lanes.state_path("cassandra", "t_collision"), 2)
path.parent.mkdir(parents=True)
path.write_bytes(b"invalid")
quarantine = path.parent / "quarantine"
quarantine.mkdir()
fingerprint = hashlib.sha256(b"invalid").hexdigest()[:16]
base = f"{path.name}.malformed-payload.{fingerprint}"
victim = tmp_path / "victim"
victim.write_text("unchanged", encoding="utf-8")
(quarantine / f"{base}.0.quarantine").symlink_to(victim)
collision = quarantine / f"{base}.1.quarantine"
collision.write_text("attacker collision", encoding="utf-8")
collision.chmod(0o644)
monkeypatch.setattr(
lanes.os,
"chmod",
lambda *_args, **_kwargs: pytest.fail("quarantine must not chmod foreign inodes"),
)
destination = lanes._quarantine_terminal(
path,
lanes._terminal_identity(path),
"malformed-payload",
)
assert not path.exists()
assert destination.name == f"{base}.2.quarantine"
assert destination.read_bytes() == b"invalid"
assert destination.stat().st_mode & 0o777 == 0o600
assert victim.read_text(encoding="utf-8") == "unchanged"
assert collision.read_text(encoding="utf-8") == "attacker collision"
assert collision.stat().st_mode & 0o777 == 0o644
def test_hardlinked_terminal_source_is_never_chmodded_or_copied_as_data(
tmp_path: Path,
monkeypatch,
):
monkeypatch.setattr(lanes, "STATE_ROOT", tmp_path / "cli-lanes")
path = lanes._terminal_path(lanes.state_path("cassandra", "t_hardlink"), 5)
path.parent.mkdir(parents=True)
foreign = tmp_path / "foreign"
foreign.write_text("foreign inode contents", encoding="utf-8")
foreign.chmod(0o644)
os.link(foreign, path)
monkeypatch.setattr(
lanes.os,
"chmod",
lambda *_args, **_kwargs: pytest.fail("hardlinked source must not be chmodded"),
)
destination = lanes._quarantine_terminal(
path,
lanes._terminal_identity(path),
"malformed-payload",
)
assert not path.exists()
assert foreign.read_text(encoding="utf-8") == "foreign inode contents"
assert foreign.stat().st_mode & 0o777 == 0o644
assert destination.stat().st_mode & 0o777 == 0o600
diagnostic = json.loads(destination.read_text(encoding="utf-8"))
assert diagnostic["source_kind"] == "hardlink"
assert "foreign inode contents" not in destination.read_text(encoding="utf-8")
def test_recovery_does_not_parse_committed_journals_every_tick(
tmp_path: Path,
monkeypatch,
@ -1127,6 +1416,8 @@ def test_artifact_gc_prunes_by_age_without_touching_pending_journals(
]
for artifact in old_artifacts:
artifact.write_text("{}", encoding="utf-8")
quarantine_symlink = quarantine / "attacker.quarantine"
quarantine_symlink.symlink_to(tmp_path / "missing-target")
pending = board / "t.run-1.terminal.pending.json"
pending.write_text("{}", encoding="utf-8")
old_time = 100.0
@ -1136,8 +1427,9 @@ def test_artifact_gc_prunes_by_age_without_touching_pending_journals(
assert lanes.gc_lane_artifacts(
now=1000.0, max_age_seconds=10, max_count=100, max_bytes=10000
) == len(old_artifacts)
) == len(old_artifacts) + 1
assert not any(artifact.exists() for artifact in old_artifacts)
assert not quarantine_symlink.is_symlink()
assert pending.exists()
@ -1189,9 +1481,10 @@ def test_artifact_gc_is_interval_gated(
"status": "completed",
"summary": "done",
"changed_files": ["src/a.py"],
"tests_run": ["pytest -q"],
"artifacts": ["reports/result.json"],
"blockers": [],
"tests_run": ["pytest -q"],
"artifacts": ["reports/result.json"],
"findings": [],
"blockers": [],
},
False,
),
@ -1205,9 +1498,10 @@ def test_artifact_gc_is_interval_gated(
"status": "completed",
"summary": "The full test suite is still running.",
"changed_files": ["src/a.py"],
"tests_run": ["pytest -q — in progress"],
"artifacts": [],
"blockers": [],
"tests_run": ["pytest -q — in progress"],
"artifacts": [],
"findings": [],
"blockers": [],
},
False,
),
@ -1369,6 +1663,7 @@ def test_goal_card_continues_after_local_judge_rejects_progress(
"changed_files": ["src/a.py"],
"tests_run": ["pytest focused: passed"],
"artifacts": [],
"findings": [],
"blockers": [],
},
False,
@ -1382,6 +1677,7 @@ def test_goal_card_continues_after_local_judge_rejects_progress(
"changed_files": ["src/a.py"],
"tests_run": ["pytest full: passed"],
"artifacts": [],
"findings": [],
"blockers": [],
},
False,
@ -1614,6 +1910,7 @@ def test_restart_provider_change_includes_explicit_workspace_handoff(tmp_path: P
"changed_files": [],
"tests_run": [],
"artifacts": [],
"findings": [],
"blockers": [],
},
False,