From 9c4ed242962f5072f3d399437a8439c93182db28 Mon Sep 17 00:00:00 2001 From: Alexandre Teixeira <111787685+alteixeira20@users.noreply.github.com> Date: Fri, 2 Oct 2026 20:02:31 +0100 Subject: [PATCH] feat(runtime): add durable append-only effect log with interrupted replay Claims are fsynced to a per-lineage JSONL log before a caller may invoke a backend; a persistence failure raises EffectPersistenceError (a ResourceIdentityError) so dispatch fails closed. Outcomes and observations are appended; nothing is rewritten. Reload validates every record strictly, ignores only a torn final write, and fails closed on corruption, forgery or hardlink aliasing. recover_interrupted appends INTERRUPTED/possible-impact outcomes for claims that never settled and leaves RUNNING background effects alone. The effect store is added to Wave 3 control-plane paths (prefix check only; the log itself refuses aliased files), so filesystem tools cannot forge it. Tests redirect the store to a session tmp directory. --- src/agent_runtime/effect_log.py | 177 +++++++++++++++++++++++ src/agent_runtime/resources.py | 8 + tests/conftest.py | 22 +++ tests/test_effect_journal_persistence.py | 164 +++++++++++++++++++++ 4 files changed, 371 insertions(+) create mode 100644 src/agent_runtime/effect_log.py create mode 100644 tests/test_effect_journal_persistence.py diff --git a/src/agent_runtime/effect_log.py b/src/agent_runtime/effect_log.py new file mode 100644 index 000000000..34ead8372 --- /dev/null +++ b/src/agent_runtime/effect_log.py @@ -0,0 +1,177 @@ +"""Durable append-only effect log for one root run lineage. + +This is the Wave 4 semantic store: claims, outcomes and observations only. It +is not a resource database, a process/containment store or an authority source. +A claim is fsynced before the backend is invoked; if that fails, the caller must +refuse the invocation. Later records are appended; nothing is rewritten. + +On reload, a claim without a settled outcome becomes an appended INTERRUPTED +outcome with possible impact. Reload never manufactures success and never +upgrades an old report to fresh state. +""" +from __future__ import annotations + +import json +import os +from pathlib import Path +import re +import stat +import threading +from typing import Any + +from src.constants import DATA_DIR +from src.agent_runtime.effects import ( + EffectAssessment, EffectClaim, EffectHistory, EffectOutcome, Observation, assess_all, replay_interrupted, +) +from src.agent_runtime.resources import ResourceIdentityError + + +EFFECTS_DIR = os.path.join(DATA_DIR, "effects") +_RUN_ID = re.compile(r"[a-f0-9]{32}") +_TYPES = {"claim": EffectClaim, "outcome": EffectOutcome, "observation": Observation} +_VERSION = 1 + + +class EffectPersistenceError(ResourceIdentityError): + """A pre-invocation claim could not be made durable; do not invoke.""" + + +def effects_dir() -> Path: + return Path(EFFECTS_DIR) + + +class EffectLog: + def __init__(self, run_id: str, *, durable: bool = True, directory: str | os.PathLike | None = None) -> None: + if not isinstance(run_id, str) or not _RUN_ID.fullmatch(run_id): + raise ValueError("Effect log requires a server-generated run identifier") + self.run_id = run_id + self.path = (Path(directory) if directory is not None else effects_dir()) / f"{run_id}.jsonl" if durable else None + self._claims: list[EffectClaim] = [] + self._outcomes: list[EffectOutcome] = [] + self._observations: list[Observation] = [] + self._sequence = 0 + # A non-claim record failed to persist. In-memory history stays + # truthful for this process; replay may lack the later record. + self.degraded = False + self._lock = threading.RLock() + + # -- persistence ------------------------------------------------------- + + def _write(self, kind: str, record: Any) -> None: + assert self.path is not None + line = json.dumps({"v": _VERSION, "type": kind, "record": record.to_dict()}, + sort_keys=True, separators=(",", ":"), ensure_ascii=False) + "\n" + self.path.parent.mkdir(mode=0o700, parents=True, exist_ok=True) + flags = os.O_WRONLY | os.O_APPEND | os.O_CREAT | getattr(os, "O_NOFOLLOW", 0) | getattr(os, "O_CLOEXEC", 0) + descriptor = os.open(self.path, flags, 0o600) + try: + info = os.fstat(descriptor) + if info.st_nlink != 1 or not stat.S_ISREG(info.st_mode): + raise OSError("Effect log is aliased") + data = line.encode("utf-8") + while data: + written = os.write(descriptor, data) + data = data[written:] + os.fsync(descriptor) + finally: + os.close(descriptor) + + def _append(self, kind: str, build, *, required: bool): + with self._lock: + record = build(self._sequence + 1) + # Read-only runs need no durable file: replay concerns claims, and + # observations matter on disk only alongside them. + if self.path is not None and (kind != "observation" or self._claims): + try: + self._write(kind, record) + except OSError as error: + if required: + raise EffectPersistenceError("Effect claim could not be persisted durably") from error + self.degraded = True + self._sequence = record.sequence + {"claim": self._claims, "outcome": self._outcomes, "observation": self._observations}[kind].append(record) + return record + + # -- records ----------------------------------------------------------- + + def claim(self, **fields: Any) -> EffectClaim: + """Persist a claim before invocation; raises if it is not durable.""" + run_id = fields.pop("run_id", self.run_id) + return self._append("claim", lambda seq: EffectClaim(sequence=seq, run_id=run_id, **fields), required=True) + + def outcome(self, **fields: Any) -> EffectOutcome: + return self._append("outcome", lambda seq: EffectOutcome(sequence=seq, **fields), required=False) + + def observe(self, **fields: Any) -> Observation: + return self._append("observation", lambda seq: Observation(sequence=seq, **fields), required=False) + + def history(self) -> EffectHistory: + with self._lock: + return EffectHistory(tuple(self._claims), tuple(self._outcomes), tuple(self._observations)) + + def assessments(self) -> tuple[EffectAssessment, ...]: + return assess_all(self.history()) + + # -- replay ------------------------------------------------------------ + + @classmethod + def load(cls, run_id: str, *, directory: str | os.PathLike | None = None) -> "EffectLog": + """Reload a persisted log. A malformed record fails closed. + + A torn final line (no newline) is the only tolerated damage: it was a + write interrupted by a crash, so its claim never returned to a caller + and no backend invocation followed it. + """ + log = cls(run_id, directory=directory) + assert log.path is not None + try: + descriptor = os.open(log.path, os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0) | getattr(os, "O_CLOEXEC", 0)) + except FileNotFoundError: + return log + except OSError as error: + raise EffectPersistenceError("Effect log is unreadable") from error + try: + info = os.fstat(descriptor) + if info.st_nlink != 1 or not stat.S_ISREG(info.st_mode): + raise EffectPersistenceError("Effect log is aliased") + with os.fdopen(descriptor, "rb") as stream: + descriptor = None + raw = stream.read() + except OSError as error: + raise EffectPersistenceError("Effect log is unreadable") from error + finally: + if descriptor is not None: + os.close(descriptor) + lines = raw.split(b"\n") + if lines and lines[-1] == b"": + lines.pop() + elif lines: + lines.pop() # torn final write + records: dict[str, list] = {"claim": [], "outcome": [], "observation": []} + for line in lines: + try: + entry = json.loads(line.decode("utf-8")) + if (not isinstance(entry, dict) or set(entry) != {"v", "type", "record"} + or entry["v"] != _VERSION or entry["type"] not in _TYPES): + raise ValueError("unsupported effect record") + records[entry["type"]].append(_TYPES[entry["type"]].from_dict(entry["record"])) + except (ValueError, TypeError, KeyError, UnicodeDecodeError) as error: + raise EffectPersistenceError("Effect log is corrupt") from error + try: + history = EffectHistory(tuple(records["claim"]), tuple(records["outcome"]), tuple(records["observation"])) + except ValueError as error: + raise EffectPersistenceError("Effect log history is inconsistent") from error + log._claims, log._outcomes, log._observations = (list(history.claims), list(history.outcomes), + list(history.observations)) + log._sequence = max((r.sequence for r in (*history.claims, *history.outcomes, *history.observations)), + default=0) + return log + + def recover_interrupted(self) -> tuple[EffectOutcome, ...]: + """Append INTERRUPTED outcomes for claims that never settled.""" + with self._lock: + pending = replay_interrupted(self.history(), self._sequence + 1) + for outcome in pending: + self._append("outcome", lambda seq, o=outcome: EffectOutcome( + o.effect_id, seq, o.execution, o.impact, replayed=True), required=False) + return tuple(self._outcomes[-len(pending):]) if pending else () diff --git a/src/agent_runtime/resources.py b/src/agent_runtime/resources.py index 8513d087f..66461744b 100644 --- a/src/agent_runtime/resources.py +++ b/src/agent_runtime/resources.py @@ -45,6 +45,14 @@ def _control_plane_path(path): processes = sys.modules.get("src.agent_runtime.process_resources") if processes is not None: job_dirs.add(canonical_root(processes._LAUNCH_DIR)) + # Durable effect claims/outcomes/observations are server evidence state. + # The log refuses hardlinked files itself, so a prefix check suffices. + effect_dirs = {canonical_root(os.path.join(constants.DATA_DIR, "effects"))} + effect_log = sys.modules.get("src.agent_runtime.effect_log") + if effect_log is not None: + effect_dirs.add(canonical_root(effect_log.EFFECTS_DIR)) + if any(Path(path).is_relative_to(directory) for directory in effect_dirs): + return True # Producers may have configured paths different from the default constants. # Inspect already-loaded server metadata without initializing a store here. bg = sys.modules.get("src.bg_jobs") diff --git a/tests/conftest.py b/tests/conftest.py index 4f3eb2c07..0915d519c 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -227,6 +227,28 @@ def _serve_test_static(): server.server_close() +@pytest.fixture(scope="session") +def _effects_store_root(tmp_path_factory): + return tmp_path_factory.mktemp("effects") + + +@pytest.fixture(autouse=True) +def _isolated_effects_store(_effects_store_root): + """Keep durable effect claims out of the developer's real data directory. + + Restored manually: requesting the shared ``monkeypatch`` here would move + its teardown after ``_no_leaked_module_stubs`` and misreport test stubs. + """ + from src.agent_runtime import effect_log + + previous = effect_log.EFFECTS_DIR + effect_log.EFFECTS_DIR = str(_effects_store_root) + try: + yield + finally: + effect_log.EFFECTS_DIR = previous + + @pytest.fixture(autouse=True) def _no_leaked_module_stubs(): """Fail the test that leaves a bare ``src.*``/``core.*`` stub behind. diff --git a/tests/test_effect_journal_persistence.py b/tests/test_effect_journal_persistence.py new file mode 100644 index 000000000..dea9125c7 --- /dev/null +++ b/tests/test_effect_journal_persistence.py @@ -0,0 +1,164 @@ +"""Durable Wave 4 effect log: pre-invocation claims, append-only replay.""" +from __future__ import annotations + +import json +import os + +import pytest + +from src.agent_runtime import effects as fx +from src.agent_runtime.effect_log import EffectLog, EffectPersistenceError +from src.agent_runtime.resources import FilesystemResource, FilesystemRoot, OwnedResource + + +RUN = "a" * 32 + + +@pytest.fixture +def target(tmp_path): + workspace = tmp_path / "ws" + workspace.mkdir() + root = FilesystemRoot.seal(str(workspace)) + return fx.resource_ref(FilesystemResource.resolve(root, str(workspace / "a.txt"), allow_missing=True), "destination") + + +def claim(log, target, effect_id="e1", action_id="act-1"): + return log.claim(effect_id=effect_id, action_id=action_id, operation=fx.OperationRef("write_file", "", "0" * 64), + impact_scope=(target,), obligations=(fx.Postcondition(target, fx.Predicate.EXISTS),)) + + +def records(path): + return [json.loads(line) for line in path.read_text().splitlines()] + + +def test_claim_is_fsynced_to_disk_before_returning(tmp_path, target, monkeypatch): + synced = [] + real_fsync = os.fsync + monkeypatch.setattr(os, "fsync", lambda fd: (synced.append(fd), real_fsync(fd))) + log = EffectLog(RUN, directory=tmp_path / "fx") + made = claim(log, target) + assert synced, "claim must be fsynced before the caller can invoke a backend" + on_disk = records(log.path) + assert [r["type"] for r in on_disk] == ["claim"] + assert fx.EffectClaim.from_dict(on_disk[0]["record"]) == made + assert oct(log.path.stat().st_mode & 0o777) == "0o600" + + +def test_claim_persistence_failure_raises_and_records_nothing(tmp_path, target): + blocker = tmp_path / "not-a-directory" + blocker.write_text("x") + log = EffectLog(RUN, directory=blocker) + with pytest.raises(EffectPersistenceError): + claim(log, target) + assert log.history().claims == () + + +def test_non_claim_failure_degrades_without_losing_in_memory_truth(tmp_path, target, monkeypatch): + log = EffectLog(RUN, directory=tmp_path / "fx") + made = claim(log, target) + monkeypatch.setattr(log, "_write", lambda kind, record: (_ for _ in ()).throw(OSError("disk full"))) + log.outcome(effect_id=made.effect_id, execution=fx.ExecutionOutcome.REPORTED_SUCCESS, impact=fx.Impact.POSSIBLE) + assert log.degraded + assert fx.assess(made, log.history()).execution is fx.ExecutionOutcome.REPORTED_SUCCESS + # Replay only sees the durable claim: it stays unknown, never success. + reloaded = EffectLog.load(RUN, directory=tmp_path / "fx") + assert fx.assess(made, reloaded.history()).verdict is fx.EffectVerdict.PENDING + + +def test_replay_after_restart_marks_unsettled_claims_interrupted(tmp_path, target): + directory = tmp_path / "fx" + log = EffectLog(RUN, directory=directory) + settled = claim(log, target, "e1", "a1") + log.outcome(effect_id="e1", execution=fx.ExecutionOutcome.REPORTED_SUCCESS, impact=fx.Impact.POSSIBLE, + execution_id="a1:x") + log.observe(observation_id="o1", resource=target, mechanism=fx.ObservationMechanism.FILESYSTEM_READ, + coverage=fx.Coverage.COMPLETE, source_action_id="r1", exists=True) + pending = claim(log, target, "e2", "a2") + del log # process "crashes" before e2 settles + + reloaded = EffectLog.load(RUN, directory=directory) + assert fx.assess(settled, reloaded.history()).verdict is fx.EffectVerdict.UNVERIFIED # e2 made o1 stale + appended = reloaded.recover_interrupted() + assert [(o.effect_id, o.execution, o.impact, o.replayed) for o in appended] == [ + ("e2", fx.ExecutionOutcome.INTERRUPTED, fx.Impact.POSSIBLE, True)] + assessment = fx.assess(pending, reloaded.history()) + assert assessment.verdict is fx.EffectVerdict.UNVERIFIED and assessment.unresolved_impact + # Recovery is append-only and idempotent across another restart. + again = EffectLog.load(RUN, directory=directory) + assert again.recover_interrupted() == () + assert [r["type"] for r in records(again.path)] == ["claim", "outcome", "observation", "claim", "outcome"] + + +def test_running_background_claim_is_not_converted_by_replay(tmp_path, target): + log = EffectLog(RUN, directory=tmp_path / "fx") + made = claim(log, target) + log.outcome(effect_id=made.effect_id, execution=fx.ExecutionOutcome.RUNNING, impact=fx.Impact.POSSIBLE) + reloaded = EffectLog.load(RUN, directory=tmp_path / "fx") + assert reloaded.recover_interrupted() == () + assert fx.assess(made, reloaded.history()).verdict is fx.EffectVerdict.PENDING + + +def test_torn_final_write_is_ignored_but_corruption_fails_closed(tmp_path, target): + directory = tmp_path / "fx" + log = EffectLog(RUN, directory=directory) + claim(log, target) + with open(log.path, "ab") as stream: + stream.write(b'{"v":1,"type":"outcome","rec') # crash mid-append + assert len(EffectLog.load(RUN, directory=directory).history().claims) == 1 + with open(log.path, "ab") as stream: + stream.write(b'\n{"v":1,"type":"outcome","record":{"forged":true}}\n') + with pytest.raises(EffectPersistenceError): + EffectLog.load(RUN, directory=directory) + + +def test_forged_success_record_cannot_be_replayed_into_verification(tmp_path, target): + directory = tmp_path / "fx" + log = EffectLog(RUN, directory=directory) + made = claim(log, target) + forged = {"v": 1, "type": "outcome", "record": {**fx.EffectOutcome( + made.effect_id, 2, fx.ExecutionOutcome.FAILED, fx.Impact.POSSIBLE).to_dict(), "execution": "verified"}} + with open(log.path, "a") as stream: + stream.write(json.dumps(forged) + "\n") + with pytest.raises(EffectPersistenceError): + EffectLog.load(RUN, directory=directory) + + +def test_hardlinked_log_is_refused(tmp_path, target): + directory = tmp_path / "fx" + log = EffectLog(RUN, directory=directory) + claim(log, target) + os.link(log.path, tmp_path / "alias.jsonl") + with pytest.raises(EffectPersistenceError): + EffectLog.load(RUN, directory=directory) + with pytest.raises(EffectPersistenceError): + claim(log, target, "e2", "a2") + + +def test_read_only_runs_write_no_file(tmp_path, target): + log = EffectLog(RUN, directory=tmp_path / "fx") + log.observe(observation_id="o1", resource=target, mechanism=fx.ObservationMechanism.FILESYSTEM_READ, + coverage=fx.Coverage.PARTIAL, source_action_id="r1", exists=True) + assert not log.path.exists() and len(log.history().observations) == 1 + + +def test_run_identifier_must_be_server_generated(tmp_path): + for forged in ("../escape", "", "A" * 32, "a" * 31): + with pytest.raises(ValueError): + EffectLog(forged, directory=tmp_path) + + +def test_effect_store_is_server_control_state(tmp_path): + from src.agent_runtime import effect_log + store = effect_log.effects_dir() + store.mkdir(parents=True, exist_ok=True) + root = FilesystemRoot.seal(str(store.parent)) + with pytest.raises(ValueError, match="sensitive"): + FilesystemResource.resolve(root, str(store / ("b" * 32 + ".jsonl")), allow_missing=True) + + +def test_owned_revision_scope_round_trips(tmp_path): + record = fx.resource_ref(OwnedResource("notes", "u", "t", "notes", "n1", "rev-1"), "record") + log = EffectLog(RUN, directory=tmp_path / "fx") + log.claim(effect_id="e1", action_id="a1", operation=fx.OperationRef("manage_notes", "", "0" * 64), + impact_scope=(record,)) + assert EffectLog.load(RUN, directory=tmp_path / "fx").history().claims[0].impact_scope == (record,)