fix(effects): make journal sequencing crash and concurrency safe

Sequence positions were allocated from each EffectLog object's in-memory
counter, so two objects, threads or processes could reuse a position or
settle one effect twice; replay then failed closed for the whole log. Every
append now takes an exclusive flock, merges the durable records other writers
appended (truncating a torn tail a crashed writer left), allocates from that
merged tail, rejects records the merged history makes invalid (a second
settlement, recovery of a claim another writer settled or marked running),
then appends, fsyncs and releases. history() merges others' records under a
shared lock. An incremental consistency index keeps appends O(1).

The first append of each log object fsyncs the log's directory, and every
directory created for it is fsynced in its parent, all under the lock before
the claim returns. A failed write or directory fsync truncates the record
back, so dispatch is refused and nothing unacknowledged is later merged. The
launch index writes and fsyncs a temp file, replaces it, then fsyncs the
directory. flock and directory fsync are POSIX-only and not claimed
elsewhere.
This commit is contained in:
Alexandre Teixeira
2026-10-03 00:58:32 +01:00
parent f79c2aba7e
commit 04da2f04e8
+345 -72
View File
@@ -8,9 +8,19 @@ 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.
Several writers may append to one log (another ``EffectLog`` object, thread or
process settling a background launch). Each append takes an exclusive advisory
lock on the file, merges every durable record other writers appended, allocates
the next position from that merged tail, checks the record against the merged
history, then appends and fsyncs before releasing the lock. Positions therefore
stay unique and a settled outcome is never appended twice. Cross-process
exclusion relies on POSIX ``flock``; directory fsync relies on POSIX directory
semantics. Neither is claimed where the platform does not provide it.
"""
from __future__ import annotations
from contextlib import contextmanager
import json
import os
from pathlib import Path
@@ -18,11 +28,17 @@ import re
import stat
import threading
import weakref
from typing import Any
from typing import Any, Callable
try: # POSIX only; elsewhere exclusion is per process.
import fcntl
except ImportError: # pragma: no cover - non-POSIX hosts
fcntl = None
from src.constants import DATA_DIR
from src.agent_runtime.effects import (
EffectAssessment, EffectClaim, EffectHistory, EffectOutcome, Observation, assess_all, replay_interrupted,
EffectAssessment, EffectClaim, EffectHistory, EffectOutcome, ExecutionOutcome, Observation, assess_all,
replay_interrupted,
)
from src.agent_runtime.resources import ResourceIdentityError
@@ -33,6 +49,31 @@ _TYPES = {"claim": EffectClaim, "outcome": EffectOutcome, "observation": Observa
_VERSION = 1
def _fsync_directory(directory: str | os.PathLike) -> None:
"""Make a directory's entries durable. POSIX only; a no-op elsewhere."""
if os.name != "posix":
return
descriptor = os.open(os.fspath(directory), os.O_RDONLY | getattr(os, "O_DIRECTORY", 0))
try:
os.fsync(descriptor)
finally:
os.close(descriptor)
def _ensure_directory(directory: Path) -> None:
"""Create ``directory`` and make every newly created entry durable."""
missing = []
current = directory
while not current.exists():
missing.append(current)
if current.parent == current:
break
current = current.parent
directory.mkdir(mode=0o700, parents=True, exist_ok=True)
for created in reversed(missing):
_fsync_directory(created.parent)
class EffectPersistenceError(ResourceIdentityError):
"""A pre-invocation claim could not be made durable; do not invoke."""
@@ -41,10 +82,58 @@ def effects_dir() -> Path:
return Path(EFFECTS_DIR)
def _parse(line: bytes) -> tuple[str, Any]:
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")
return entry["type"], _TYPES[entry["type"]].from_dict(entry["record"])
class _Index:
"""Incremental consistency of an append-ordered record stream.
At least as strict as ``EffectHistory`` validation for records appended in
position order, so a record it accepts never makes the history invalid.
"""
def __init__(self) -> None:
self.positions: set[int] = set()
self.claims: dict[str, int] = {}
self.settled: set[str] = set()
self.last: dict[str, int] = {}
def accepts(self, kind: str, record: Any) -> bool:
if record.sequence in self.positions:
return False
if kind == "claim":
return record.effect_id not in self.claims
if kind == "outcome":
effect = record.effect_id
return (effect in self.claims and record.sequence > self.claims[effect]
and effect not in self.settled and record.sequence > self.last.get(effect, -1))
return True
def add(self, kind: str, record: Any) -> None:
self.positions.add(record.sequence)
if kind == "claim":
self.claims[record.effect_id] = record.sequence
elif kind == "outcome":
self.last[record.effect_id] = record.sequence
if record.execution is not ExecutionOutcome.RUNNING:
self.settled.add(record.effect_id)
def _records() -> dict[str, list]:
return {"claim": [], "outcome": [], "observation": []}
class EffectLog:
# Logs still owned by a live run in this process. A later turn appends to
# the same object rather than a second copy with its own sequence counter.
# the same object rather than a second copy of the same history.
_LIVE: "weakref.WeakValueDictionary[tuple[str, str], EffectLog]" = weakref.WeakValueDictionary()
# Serializes the _LIVE check-and-load in ``open``.
_OPEN_LOCK = threading.Lock()
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):
@@ -53,68 +142,241 @@ class EffectLog:
self.path = (Path(directory) if directory is not None else effects_dir()) / f"{run_id}.jsonl" if durable else None
if self.path is not None:
self._LIVE[(str(self.path.parent), run_id)] = self
self._claims: list[EffectClaim] = []
self._outcomes: list[EffectOutcome] = []
self._observations: list[Observation] = []
self._sequence = 0
# Records known to be on disk, in file order, and the bytes they span.
self._durable, self._durable_index = _records(), _Index()
self._offset = 0
# Records this process holds that are not on disk: a failed non-claim
# write, or an observation of a run that has no durable file yet.
self._volatile: list[tuple[str, Any]] = []
# Durable records plus every volatile record still consistent with
# them. A volatile record that collides with a durable one (another
# writer took its position or settled the same effect) is hidden:
# losing an unpersisted outcome or observation is conservative.
self._view, self._view_index = _records(), _Index()
self._max_sequence = 0
self._descriptor: int | None = None
self._directory_synced = False
# 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()
# -- merged view ---------------------------------------------------------
def _rebuild_view(self) -> None:
self._view, self._view_index = _records(), _Index()
durable = sorted(((kind, r) for kind, records in self._durable.items() for r in records),
key=lambda pair: pair[1].sequence)
for kind, record in durable:
self._view[kind].append(record)
self._view_index.add(kind, record)
for kind, record in self._volatile:
if self._view_index.accepts(kind, record):
self._view_index.add(kind, record)
self._view[kind].append(record)
def _add_durable(self, kind: str, record: Any) -> None:
self._durable[kind].append(record)
self._durable_index.add(kind, record)
self._view[kind].append(record)
self._view_index.add(kind, record)
self._max_sequence = max(self._max_sequence, record.sequence)
def _add_volatile(self, kind: str, record: Any) -> None:
self._volatile.append((kind, record))
self._view[kind].append(record)
self._view_index.add(kind, record)
self._max_sequence = max(self._max_sequence, record.sequence)
def _merged_history(self) -> EffectHistory:
return EffectHistory(tuple(self._view["claim"]), tuple(self._view["outcome"]),
tuple(self._view["observation"]))
# -- persistence -------------------------------------------------------
def _write(self, kind: str, record: Any) -> None:
def _read_tail(self, descriptor: int, *, repair: bool) -> None:
"""Merge complete records other writers appended after our offset.
A trailing partial line is a write that never returned to its caller
(a crash or a failed write), so no backend invocation followed it.
With ``repair`` (exclusive lock held) it is truncated so the next
append starts on a record boundary.
"""
info = os.fstat(descriptor)
if info.st_nlink != 1 or not stat.S_ISREG(info.st_mode):
raise OSError("Effect log is aliased")
if info.st_size < self._offset:
raise OSError("Effect log shrank under its writer")
if info.st_size == self._offset:
return
os.lseek(descriptor, self._offset, os.SEEK_SET)
remaining, chunks = info.st_size - self._offset, []
while remaining:
chunk = os.read(descriptor, remaining)
if not chunk:
break
chunks.append(chunk)
remaining -= len(chunk)
data = b"".join(chunks)
complete = data[:data.rfind(b"\n") + 1]
if repair and len(complete) != len(data):
os.ftruncate(descriptor, self._offset + len(complete))
os.fsync(descriptor)
added: list[tuple[str, Any]] = []
for line in complete.splitlines():
try:
kind, record = _parse(line)
except (ValueError, TypeError, KeyError, UnicodeDecodeError) as error:
raise OSError("Effect log tail is corrupt") from error
if not self._durable_index.accepts(kind, record):
raise OSError("Effect log tail is inconsistent")
self._durable[kind].append(record)
self._durable_index.add(kind, record)
self._max_sequence = max(self._max_sequence, record.sequence)
added.append((kind, record))
self._offset += len(complete)
if added:
self._rebuild_view()
@contextmanager
def _locked(self):
"""Hold the file exclusively with every durable record merged."""
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)
_ensure_directory(self.path.parent)
flags = os.O_RDWR | 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)
if fcntl is not None:
fcntl.flock(descriptor, fcntl.LOCK_EX)
self._read_tail(descriptor, repair=True)
self._descriptor = descriptor
yield
finally:
os.close(descriptor)
self._descriptor = None
os.close(descriptor) # releases the lock
def _append(self, kind: str, build, *, required: bool):
def _write(self, kind: str, record: Any) -> None:
"""Append one record under the held lock and make it durable."""
descriptor = self._descriptor
assert descriptor is not None
line = json.dumps({"v": _VERSION, "type": kind, "record": record.to_dict()},
sort_keys=True, separators=(",", ":"), ensure_ascii=False) + "\n"
data = line.encode("utf-8")
start = os.fstat(descriptor).st_size
try:
view = memoryview(data)
while view:
view = view[os.write(descriptor, view):]
os.fsync(descriptor)
if not self._directory_synced:
# The file's own fsync does not make its directory entry
# durable. Sync it before the first claim returns, still under
# the lock, so no writer can invoke a backend against a log a
# crash could lose.
_fsync_directory(self.path.parent)
self._directory_synced = True
except OSError:
# Unacknowledged: take the record back so the file ends on a
# record boundary and no writer later merges it as durable.
try:
os.ftruncate(descriptor, start)
os.fsync(descriptor)
except OSError:
pass
raise
self._offset = start + len(data)
def _append(self, kind: str, build: Callable[[int, "EffectLog"], Any], *, required: bool):
"""Allocate, validate and persist one record; ``None`` if it no longer applies.
``build(sequence, view)`` may return ``None`` when its precondition no
longer holds against the merged view (``self``).
"""
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
durable = self.path is not None and (kind != "observation" or self._view["claim"]
or self.path.exists())
if not durable:
return self._append_volatile(kind, build, required=required)
try:
with self._locked():
record = build(self._max_sequence + 1, self)
if record is None:
return None
if not (self._view_index.accepts(kind, record) and self._durable_index.accepts(kind, record)):
# Another writer already settled it, or it collides.
if required:
raise ValueError("Effect record conflicts with the durable history")
return None
try:
self._write(kind, record)
except OSError:
if required:
raise
self.degraded = True
self._add_volatile(kind, record)
return record
self._add_durable(kind, record)
return record
except (OSError, ValueError) as error:
if required:
raise EffectPersistenceError("Effect claim could not be persisted durably") from error
self.degraded = True
# The lock or tail could not be taken: keep this process
# truthful without touching the file.
return self._append_volatile(kind, build, required=False)
def _append_volatile(self, kind: str, build, *, required: bool):
record = build(self._max_sequence + 1, self)
if record is None:
return None
if not self._view_index.accepts(kind, record):
if required:
raise ValueError("Effect record conflicts with the history")
return None
self._add_volatile(kind, 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)
return self._append("claim", lambda seq, _h: 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 outcome(self, **fields: Any) -> EffectOutcome | None:
"""Append an outcome; ``None`` if the effect was already settled."""
return self._append("outcome", lambda seq, _h: 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 observe(self, **fields: Any) -> Observation | None:
return self._append("observation", lambda seq, _h: Observation(sequence=seq, **fields), required=False)
def refresh(self) -> None:
"""Merge records other writers appended (shared lock, no repair)."""
if self.path is None:
return
with self._lock:
try:
descriptor = os.open(self.path, os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0)
| getattr(os, "O_CLOEXEC", 0))
except FileNotFoundError:
return
except OSError:
self.degraded = True
return
try:
if fcntl is not None:
fcntl.flock(descriptor, fcntl.LOCK_SH)
self._read_tail(descriptor, repair=False)
except OSError:
self.degraded = True
finally:
os.close(descriptor)
def history(self) -> EffectHistory:
self.refresh()
with self._lock:
return EffectHistory(tuple(self._claims), tuple(self._outcomes), tuple(self._observations))
return self._merged_history()
def assessments(self) -> tuple[EffectAssessment, ...]:
return assess_all(self.history())
@@ -138,6 +400,8 @@ class EffectLog:
except OSError as error:
raise EffectPersistenceError("Effect log is unreadable") from error
try:
if fcntl is not None:
fcntl.flock(descriptor, fcntl.LOCK_SH)
info = os.fstat(descriptor)
if info.st_nlink != 1 or not stat.S_ISREG(info.st_mode):
raise EffectPersistenceError("Effect log is aliased")
@@ -149,29 +413,23 @@ class EffectLog:
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:
complete = raw[:raw.rfind(b"\n") + 1] # drop a torn final write
for line in complete.splitlines():
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"]))
kind, record = _parse(line)
except (ValueError, TypeError, KeyError, UnicodeDecodeError) as error:
raise EffectPersistenceError("Effect log is corrupt") from error
if not log._durable_index.accepts(kind, record):
raise EffectPersistenceError("Effect log history is inconsistent")
log._durable[kind].append(record)
log._durable_index.add(kind, record)
log._max_sequence = max(log._max_sequence, record.sequence)
try:
history = EffectHistory(tuple(records["claim"]), tuple(records["outcome"]), tuple(records["observation"]))
log._rebuild_view()
log._merged_history()
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)
log._offset = len(complete)
return log
@classmethod
@@ -182,12 +440,13 @@ class EffectLog:
unsettled claims are recovered as interrupted before any append.
"""
base = Path(directory) if directory is not None else effects_dir()
live = cls._LIVE.get((str(base), run_id))
if live is not None:
return live
log = cls.load(run_id, directory=base)
log.recover_interrupted()
return log
with cls._OPEN_LOCK:
live = cls._LIVE.get((str(base), run_id))
if live is not None:
return live
log = cls.load(run_id, directory=base)
log.recover_interrupted()
return log
# -- background launch lineage ------------------------------------------
@@ -200,14 +459,18 @@ class EffectLog:
target = self.path.parent / f"launch-{generation}.json"
temporary = target.with_suffix(".tmp")
data = json.dumps({"run_id": self.run_id, "effect_id": effect_id}, sort_keys=True).encode()
flags = os.O_WRONLY | os.O_CREAT | os.O_TRUNC | getattr(os, "O_NOFOLLOW", 0)
_ensure_directory(target.parent)
flags = os.O_WRONLY | os.O_CREAT | os.O_TRUNC | getattr(os, "O_NOFOLLOW", 0) | getattr(os, "O_CLOEXEC", 0)
descriptor = os.open(temporary, flags, 0o600)
try:
os.write(descriptor, data)
view = memoryview(data)
while view:
view = view[os.write(descriptor, view):]
os.fsync(descriptor)
finally:
os.close(descriptor)
os.replace(temporary, target)
_fsync_directory(target.parent)
@staticmethod
def launch_owner(generation: str, *, directory: str | os.PathLike | None = None) -> tuple[str, str] | None:
@@ -229,10 +492,20 @@ class EffectLog:
return value["run_id"], value["effect_id"]
def recover_interrupted(self) -> tuple[EffectOutcome, ...]:
"""Append INTERRUPTED outcomes for claims that never settled."""
"""Append INTERRUPTED outcomes for claims that never settled.
Each claim is rechecked against the merged history under the lock, so
a claim another writer settled (or marked running) meanwhile is left
alone.
"""
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 ()
appended = []
for pending in replay_interrupted(self.history(), self._max_sequence + 1):
def build(seq: int, log: "EffectLog", o: EffectOutcome = pending) -> EffectOutcome | None:
if o.effect_id in log._view_index.last or o.effect_id in log._durable_index.last:
return None
return EffectOutcome(o.effect_id, seq, o.execution, o.impact, replayed=True)
record = self._append("outcome", build, required=False)
if record is not None:
appended.append(record)
return tuple(appended)