From 4dfc8b9def6f2d89738a77de186637f9113c5fdf Mon Sep 17 00:00:00 2001 From: Brian Krabach Date: Tue, 22 Sep 2026 08:19:56 -0700 Subject: [PATCH 1/3] Fence native session execution during durable transfer --- amplifier_foundation/session/__init__.py | 4 + amplifier_foundation/session/shared_state.py | 151 +++++++++++++- docs/API_REFERENCE.md | 1 + docs/SESSION_TRANSFER_FENCE.md | 105 ++++++++++ docs/SHARED_SESSION_STATE.md | 7 +- tests/test_session_init_exports.py | 10 +- tests/test_session_transfer_fence.py | 200 +++++++++++++++++++ 7 files changed, 472 insertions(+), 6 deletions(-) create mode 100644 docs/SESSION_TRANSFER_FENCE.md create mode 100644 tests/test_session_transfer_fence.py diff --git a/amplifier_foundation/session/__init__.py b/amplifier_foundation/session/__init__.py index 4496f5a1..0a21fc37 100644 --- a/amplifier_foundation/session/__init__.py +++ b/amplifier_foundation/session/__init__.py @@ -154,7 +154,9 @@ from .shared_state import ( FileStamp, HeldSession, + HeldTransfer, SessionBusyError, + SessionTransferFencedError, SharedSessionStore, file_stamp, ) @@ -231,6 +233,8 @@ "FileStamp", "file_stamp", "SessionBusyError", + "SessionTransferFencedError", "SharedSessionStore", "HeldSession", + "HeldTransfer", ] diff --git a/amplifier_foundation/session/shared_state.py b/amplifier_foundation/session/shared_state.py index f3d5ea89..8ad49580 100644 --- a/amplifier_foundation/session/shared_state.py +++ b/amplifier_foundation/session/shared_state.py @@ -21,7 +21,7 @@ from dataclasses import dataclass from datetime import datetime, timezone from pathlib import Path -from typing import Any +from typing import Any, cast try: # Fail loudly on platforms without the locking primitive. import fcntl @@ -80,6 +80,14 @@ class SharedStateError(RuntimeError): """Raised when shared-state storage or a checkpoint is invalid.""" +class SessionTransferFencedError(RuntimeError): + """A durable transfer marker prevents ordinary session execution.""" + + def __init__(self, fence: dict[str, Any]) -> None: + self.fence = copy.deepcopy(fence) + super().__init__("session execution is fenced by a durable transfer") + + def _default_root() -> Path: configured = os.environ.get("AMPLIFIER_SESSION_STATE_HOME") if configured: @@ -360,6 +368,53 @@ def __init__(self, workspace: str | os.PathLike[str], session_id: str, *, root: self.session_id = _session_id(session_id) self.root = Path(root).expanduser() if root is not None else _default_root() self._directory = self.root / "v1" / _workspace_key(self.workspace) / self.session_id + home = Path(os.environ.get("AMPLIFIER_HOME") or Path.home() / ".amplifier").expanduser().absolute() + slug = str(self.workspace).replace("/", "-").replace("\\", "-").replace(":", "") + self._native_root = home + self._native_directory = home / "projects" / (slug if slug.startswith("-") else "-" + slug) / "sessions" / self.session_id + + @property + def transfer_fence_path(self) -> Path: + """Native-history marker; reading this property creates no files.""" + return self._native_directory / "transfer-fence.json" + + def _transfer_directory(self, *, create: bool = False) -> None: + # Native histories predate private coordination roots, so do not change + # existing directory modes. Reject links, non-directories and foreign + # owners; the marker itself must always be private regular data. + relative = self._native_directory.relative_to(self._native_root) + directories = [self._native_root] + for part in relative.parts: + directories.append(directories[-1] / part) + for directory in directories: + try: + info = directory.lstat() + except FileNotFoundError: + if not create: + return + directory.mkdir(mode=0o700, parents=True, exist_ok=True) + info = directory.lstat() + if not stat.S_ISDIR(info.st_mode) or info.st_uid != os.getuid(): + raise SharedStateError("native session transfer directory is unsafe") + _validate_private_file(self.transfer_fence_path) + + def transfer_fence(self) -> dict[str, Any] | None: + """Inspect durable transfer state without acquiring or starting work.""" + self._transfer_directory() + value = _read_json(self.transfer_fence_path, missing_ok=True) + if value is None: + return None + required = {"version", "session_id", "transfer_id", "destination_host", "role", "phase"} + if (set(value) != required or type(value["version"]) is not int or value["version"] != 1 + or value["session_id"] != self.session_id + or not isinstance(value["transfer_id"], str) or not _ID_RE.fullmatch(value["transfer_id"]) + or not isinstance(value["destination_host"], str) or not 1 <= len(value["destination_host"]) <= 255 + or any(ord(char) < 32 for char in value["destination_host"]) + or not isinstance(value["role"], str) or value["role"] not in {"source", "destination"} + or not isinstance(value["phase"], str) or value["phase"] not in {"staged", "committed"} + or (value["role"] == "destination" and value["phase"] != "staged")): + raise SharedStateError("invalid session transfer fence") + return value @property def checkpoint_path(self) -> Path: @@ -400,6 +455,19 @@ def read(self) -> dict[str, Any] | None: def acquire(self, *, app: str, **diagnostics: Any) -> "HeldSession": """Acquire the stable OS lock and publish advisory owner diagnostics.""" + return self._acquire(app=app, diagnostics=diagnostics) + + def acquire_transfer(self, transfer_id: str, *, app: str, **diagnostics: Any) -> "HeldTransfer": + """Acquire only a matching staged fence for explicit adapter recovery. + + This does not grant execution or validate remote release evidence. The + application must authenticate the transfer before clearing/committing. + Committed source markers cannot be acquired through this API. + """ + _session_id(transfer_id) + return cast("HeldTransfer", self._acquire(app=app, diagnostics=diagnostics, transfer_id=transfer_id)) + + def _acquire(self, *, app: str, diagnostics: dict[str, Any], transfer_id: str | None = None) -> "HeldSession": _ensure_supported() owner = _owner_details(app, diagnostics, self.workspace, self.session_id, self.root) _ensure_private_state_path(self.root, self._directory) @@ -424,6 +492,12 @@ def acquire(self, *, app: str, **diagnostics: Any) -> "HeldSession": pass raise SharedStateError("cannot acquire stable session lock") from exc try: + fence = self.transfer_fence() + if transfer_id is None: + if fence is not None: + raise SessionTransferFencedError(fence) + elif fence is None or fence["transfer_id"] != transfer_id or fence["phase"] != "staged": + raise SharedStateError("transfer recovery requires the exact staged fence") _atomic_json(self._owner_path, owner) except BaseException: fcntl.flock(fd, fcntl.LOCK_UN) @@ -431,6 +505,8 @@ def acquire(self, *, app: str, **diagnostics: Any) -> "HeldSession": raise _PROCESS_LOCKS.add(lock_name) _LOCK_FDS.add(fd) + if transfer_id is not None: + return HeldTransfer(self, fd, owner, transfer_id) return HeldSession(self, fd, owner) @classmethod @@ -482,11 +558,41 @@ def active(self) -> bool: return self._active and self._pid == os.getpid() def check(self) -> None: + self._check_active() + fence = self._store.transfer_fence() + if fence is not None: + raise SessionTransferFencedError(fence) + + def _check_active(self) -> None: if self._pid != os.getpid(): raise RuntimeError("HeldSession belongs to a different process") if not self._active: raise RuntimeError("HeldSession is no longer active") + def fence_transfer(self, transfer_id: str, destination_host: str, *, role: str = "source") -> dict[str, Any]: + """Fence this saved, quiescent writer before releasing it for transfer. + + The host must settle work and persist native history first. The marker + does not stop a runtime or prove any remote outcome. Once written this + handle cannot authorize further execution or checkpoint writes. + """ + _session_id(transfer_id) + if (not isinstance(destination_host, str) or not 1 <= len(destination_host) <= 255 + or any(ord(char) < 32 for char in destination_host) or role not in {"source", "destination"}): + raise ValueError("invalid transfer destination or role") + record = {"version": 1, "session_id": self._store.session_id, "transfer_id": transfer_id, + "destination_host": destination_host, "role": role, "phase": "staged"} + with self._mutex: + self._check_active() + existing = self._store.transfer_fence() + if existing is not None: + if existing == record: + return existing + raise SharedStateError("a different transfer fence already exists") + self._store._transfer_directory(create=True) + _atomic_json(self._store.transfer_fence_path, record) + return copy.deepcopy(record) + def read(self) -> dict[str, Any] | None: self.check() return self._store.read() @@ -562,3 +668,46 @@ def release(self) -> None: _PROCESS_LOCKS.discard(lock_name) _LOCK_FDS.discard(self._fd) os.close(self._fd) + + +class HeldTransfer(HeldSession): + """A receipt-matched lock for resolving a staged fence, never execution.""" + + def __init__(self, store: SharedSessionStore, fd: int, owner: dict[str, Any], transfer_id: str) -> None: + super().__init__(store, fd, owner) + self._transfer_id = transfer_id + + def check(self) -> None: + self._check_active() + raise SharedStateError("transfer recovery capability cannot authorize execution") + + def fence_transfer(self, transfer_id: str, destination_host: str, *, role: str = "source") -> dict[str, Any]: + raise SharedStateError("transfer recovery capability cannot replace a fence") + + def _staged(self) -> dict[str, Any]: + self._check_active() + record = self._store.transfer_fence() + if record is None or record["transfer_id"] != self._transfer_id or record["phase"] != "staged": + raise SharedStateError("transfer recovery requires the exact staged fence") + return record + + def commit_transfer(self) -> dict[str, Any]: + """Permanently fence the source after the host verifies its transfer.""" + with self._mutex: + record = self._staged() + if record["role"] != "source": + raise SharedStateError("only a source transfer fence can be committed") + record["phase"] = "committed" + _atomic_json(self._store.transfer_fence_path, record) + return record + + def clear_transfer(self) -> None: + """Explicitly cancel a staged source or activate a staged destination. + + Authenticate the corresponding cancellation or release receipt before + calling. Release this capability and acquire normally before execution. + """ + with self._mutex: + self._staged() + self._store.transfer_fence_path.unlink() + _fsync_directory(self._store.transfer_fence_path.parent) diff --git a/docs/API_REFERENCE.md b/docs/API_REFERENCE.md index 8ad64b17..1ee94b80 100644 --- a/docs/API_REFERENCE.md +++ b/docs/API_REFERENCE.md @@ -116,6 +116,7 @@ Portable same-host checkpoints live in `session/shared_state.py`. They are suppo | `SharedSessionStore` | `session/shared_state.py` | Addresses a checkpoint; `acquire(app=...)` returns the exclusive writer capability. | | `HeldSession` | `session/shared_state.py` | Process-bound, non-copyable writer; atomically writes messages, portable bundle reference, and credential-safe metadata. Its live `delete_checkpoint()` safely removes only the authoritative checkpoint. | | `SessionBusyError` | `session/shared_state.py` | Contention error with advisory, bounded owner diagnostics. | +| `SessionTransferFencedError`, `HeldTransfer` | `session/shared_state.py` | Durable native-history admission fence and restricted receipt-matched resolution capability; see [session transfer fences](SESSION_TRANSFER_FENCE.md). | Import these names from `amplifier_foundation.session`. Store construction, `read`, `stamp`, and `list_ids` never create state directories. `acquire` alone creates or validates private state directories and the stable `session.lock`; it rejects symlinked, foreign-owned, or group/world-accessible state paths. `release` never rewrites the authoritative checkpoint. See [SHARED_SESSION_STATE.md](SHARED_SESSION_STATE.md) for the same-host participant contract, warm-reuse guidance, and safe checkpoint deletion. diff --git a/docs/SESSION_TRANSFER_FENCE.md b/docs/SESSION_TRANSFER_FENCE.md new file mode 100644 index 00000000..89c7af7f --- /dev/null +++ b/docs/SESSION_TRANSFER_FENCE.md @@ -0,0 +1,105 @@ +# Durable session transfer fences + +The shared-session lock prevents concurrent local writers. A transfer fence adds +durable admission state for an application that moves a saved session to another +host. It remains effective after process exit or restart. It is not a transport, +distributed lock, remote authentication mechanism, or transfer implementation. + +All execution participants must use a Foundation version with this API and call +`SharedSessionStore.acquire()` before execution and `HeldSession.check()` before +persistence. Older versions and writers that bypass the API do not honor the +fence. Hosts must use the same native history root, canonical workspace and exact +session ID. Altering the native root or manually deleting the marker is outside +this contract; the OS user controlling these private files remains trusted. + +## Source handoff + +The application blocks new work, cooperatively stops its runtime and dependent +work, checkpoints native history, and obtains the shared lock. It then fences +the saved writer before releasing the handle: + +```python +held = store.acquire(app="transfer-adapter") +try: + held.fence_transfer(transfer_id, destination_host, role="source") +finally: + held.release() +``` + +`fence_transfer` does not stop work or save history. The host must establish that +boundary first. It writes a staged marker atomically, fsyncs its directory where +supported, and makes that handle fail execution checks and checkpoint writes. +An identical call is idempotent; a different existing marker is rejected. + +The application exports saved data without replaying inputs or uncertain effects. +After authenticating destination readiness, it permanently commits the source: + +```python +recovery = store.acquire_transfer(transfer_id, app="transfer-adapter") +try: + committed = recovery.commit_transfer() +finally: + recovery.release() +``` + +A committed source cannot be acquired normally, recovered with +`acquire_transfer`, or cleared by this API. A later return transfer therefore +needs a new destination workspace/native history location. Preserve the source +marker when archiving inactive history. A failed or unknown application-level +transfer must retain its marker and expose its evidence for reconciliation. + +## Destination staging and explicit resolution + +Before installing transferred native history, an application acquires its new +destination session and calls `fence_transfer(..., role="destination")`. This +works before transcript or metadata files exist. Stage data without copying +process locks, advisory owner records, credentials or live runtime state. + +After verifying the authenticated source release certificate and its exact +transfer identity, destination activation is explicit: + +```python +recovery = destination.acquire_transfer(transfer_id, app="transfer-adapter") +try: + recovery.clear_transfer() +finally: + recovery.release() +``` + +The same `clear_transfer()` operation permits explicit cancellation of a staged +source before committing its release. The application owns the evidence and +authorization policy; matching an ID alone is not proof of a remote result. +`HeldTransfer` cannot authorize execution, read a checkpoint as a writer, write, +delete checkpoints, or replace a fence. Even after clearing, release it and +acquire a new ordinary handle before accepting later deliberate work. Clearing +does not start a runtime, submit input, or replay an operation. + +## Native storage and failure behavior + +`store.transfer_fence_path` identifies the marker without creating files: + +```text +${AMPLIFIER_HOME:-~/.amplifier}/projects//sessions//transfer-fence.json +``` + +The slug uses the native convention: replace `/` and `\\` with `-`, remove `:`, +and ensure a leading `-`. The private marker contains only `version` (1), +`session_id`, `transfer_id`, `destination_host`, `role` and `phase`. It is separate +from transcript/metadata and does not introduce a second history. Native history +remains readable without acquiring execution ownership. + +`store.transfer_fence()` is a read-only inspection API. Malformed markers, +unsupported phases, wrong session IDs, symlinks and unsafe marker permissions +fail closed. Acquisition checks the marker after obtaining the local OS lock +and before replacing advisory owner metadata. A rejected acquisition releases +the temporary lock and does not change history, marker or owner metadata. +Ordinary `acquire()` raises `SessionTransferFencedError`, whose `fence` attribute +is a detached record. `acquire_transfer()` requires the exact staged transfer; +passing a transfer ID among ordinary acquisition diagnostics never bypasses it. + +The marker is addressed through native history rather than the configurable +coordination root, so choosing another coordination root does not erase this +admission check. Applications must still share the coordination root to exclude +concurrent writers before a marker exists. Cross-host single-owner safety comes +from the application's authenticated staged-transfer protocol; these local +primitives do not provide that protocol or validate its receipts. diff --git a/docs/SHARED_SESSION_STATE.md b/docs/SHARED_SESSION_STATE.md index f6698d2c..4f27032c 100644 --- a/docs/SHARED_SESSION_STATE.md +++ b/docs/SHARED_SESSION_STATE.md @@ -25,6 +25,11 @@ No live attach protocol, TUI/web multiwriter protocol, or remote protocol is provided. There is no force unlock, TTL, lease, takeover, or shared lock-file deletion API. +Applications implementing host migration can use the separate +[durable transfer fence](SESSION_TRANSFER_FENCE.md) to keep ordinary execution +blocked after a saved writer releases its lock. This persists admission state +beside native history and does not implement remote transfer or authentication. + ## Minimal participant pattern ```python @@ -139,4 +144,4 @@ stable lock file or directory to try to take over a session. An installed Digital Twin Universe exercise, if needed, is requested from the manager. This guide's example is illustrative and does not claim that such an -exercise was run. \ No newline at end of file +exercise was run. diff --git a/tests/test_session_init_exports.py b/tests/test_session_init_exports.py index 38a6fc5f..f69cf165 100644 --- a/tests/test_session_init_exports.py +++ b/tests/test_session_init_exports.py @@ -111,14 +111,14 @@ def test_session_info_export(self): class TestDunderAll: - """Verify __all__ contains all 41 expected names.""" + """Verify the public session utility names.""" - def test_all_contains_60_names(self): + def test_all_contains_62_names(self): import amplifier_foundation.session as session assert hasattr(session, "__all__"), "__all__ must be defined" - assert len(session.__all__) == 60, ( - f"Expected 60 names in __all__, got {len(session.__all__)}: " + assert len(session.__all__) == 62, ( + f"Expected 62 names in __all__, got {len(session.__all__)}: " f"{sorted(session.__all__)}" ) @@ -161,6 +161,8 @@ def test_all_includes_new_names(self): "SessionBusyError", "SharedSessionStore", "HeldSession", + "HeldTransfer", + "SessionTransferFencedError", } for name in expected_new: assert name in session.__all__, f"{name!r} missing from __all__" diff --git a/tests/test_session_transfer_fence.py b/tests/test_session_transfer_fence.py new file mode 100644 index 00000000..52fd998e --- /dev/null +++ b/tests/test_session_transfer_fence.py @@ -0,0 +1,200 @@ +"""Durable execution fences across host adapters and process lifetimes.""" + +from __future__ import annotations + +import json +import os +import subprocess +import sys + +import pytest + +from amplifier_foundation.session import ( + HeldTransfer, + SessionTransferFencedError, + SharedSessionStore, +) +from amplifier_foundation.session.shared_state import SharedStateError + + +pytestmark = pytest.mark.skipif(os.name != "posix", reason="requires POSIX locks") + + +@pytest.fixture +def store(tmp_path, monkeypatch): + monkeypatch.setenv("AMPLIFIER_HOME", str(tmp_path / "native")) + workspace = tmp_path / "workspace" + workspace.mkdir() + return SharedSessionStore(workspace, "task-1", root=tmp_path / "locks") + + +def stage(store, *, role="source"): + held = store.acquire(app="exporter") + try: + return held.fence_transfer("transfer-1", "destination", role=role) + finally: + held.release() + + +def test_marker_is_native_private_and_does_not_change_saved_history(store): + assert store.transfer_fence() is None + assert not store.transfer_fence_path.parent.exists() + held = store.acquire(app="source") + try: + held.write([{"role": "user", "content": "keep"}], bundle="portable") + before = store.checkpoint_path.read_bytes() + directory = store.transfer_fence_path.parent + directory.mkdir(mode=0o700, parents=True) + native = {"transcript.jsonl": b'{"role":"user","content":"original"}\n', + "metadata.json": b'{"session_id":"task-1","unknown":"retained"}'} + for name, content in native.items(): + (directory / name).write_bytes(content) + marker = held.fence_transfer("transfer-1", "destination") + assert marker == held.fence_transfer("transfer-1", "destination") + assert store.checkpoint_path.read_bytes() == before + assert store.transfer_fence_path.name == "transfer-fence.json" + assert store.transfer_fence_path.parent.name == "task-1" + assert store.transfer_fence_path.stat().st_mode & 0o777 == 0o600 + assert all((directory / name).read_bytes() == content for name, content in native.items()) + for operation in (held.check, lambda: held.write([], bundle="bad"), held.delete_checkpoint): + with pytest.raises(SessionTransferFencedError): + operation() + finally: + held.release() + assert store.read()["messages"] == [{"role": "user", "content": "keep"}] + + +def test_fence_survives_process_death_and_blocks_other_adapter(store): + script = """ +import os, sys +from amplifier_foundation.session import SharedSessionStore +store = SharedSessionStore(sys.argv[1], 'task-1', root=sys.argv[2]) +held = store.acquire(app='source') +held.write([{'role':'user','content':'saved'}], bundle='portable') +held.fence_transfer('transfer-1', 'destination') +os._exit(0) +""" + result = subprocess.run([sys.executable, "-c", script, str(store.workspace), str(store.root)], capture_output=True, text=True) + assert result.returncode == 0, result.stderr + owner_before = store._owner_path.read_bytes() + for _ in range(2): + with pytest.raises(SessionTransferFencedError) as error: + store.acquire(app="different-adapter") + assert error.value.fence["transfer_id"] == "transfer-1" + assert store._owner_path.read_bytes() == owner_before + assert store.read()["messages"][0]["content"] == "saved" + + +def test_matching_recovery_cannot_execute_and_committed_source_cannot_reopen(store): + stage(store) + with pytest.raises(SharedStateError, match="exact staged"): + store.acquire_transfer("wrong-receipt", app="recovery") + held = store.acquire_transfer("transfer-1", app="recovery") + assert isinstance(held, HeldTransfer) + try: + for operation in (held.check, held.read, lambda: held.write([], bundle="bad"), held.delete_checkpoint): + with pytest.raises(SharedStateError, match="cannot authorize execution"): + operation() + assert held.commit_transfer()["phase"] == "committed" + with pytest.raises(SharedStateError, match="exact staged"): + held.clear_transfer() + finally: + held.release() + fresh = SharedSessionStore(store.workspace, store.session_id, root=store.root) + with pytest.raises(SharedStateError, match="exact staged"): + fresh.acquire_transfer("transfer-1", app="recovery") + with pytest.raises(SessionTransferFencedError): + fresh.acquire(app="ordinary") + # Native storage owns the fence even if an adapter selects another local + # coordination root. The marker is not advisory owner.json metadata. + alternate = SharedSessionStore(store.workspace, store.session_id, root=store.root.parent / "other-locks") + with pytest.raises(SessionTransferFencedError): + alternate.acquire(app="other-adapter") + + +@pytest.mark.parametrize("role", ["source", "destination"]) +def test_explicit_clear_requires_new_normal_acquisition(store, role): + stage(store, role=role) + with pytest.raises(SessionTransferFencedError): + store.acquire(app="ordinary", transfer_id="transfer-1") + held = store.acquire_transfer("transfer-1", app="recovery") + try: + if role == "destination": + with pytest.raises(SharedStateError, match="only a source"): + held.commit_transfer() + held.clear_transfer() + with pytest.raises(SharedStateError, match="cannot authorize execution"): + held.check() + finally: + held.release() + assert store.transfer_fence() is None + ordinary = store.acquire(app="destination") + try: + ordinary.write([], bundle="portable") + finally: + ordinary.release() + + +@pytest.mark.parametrize("mutation", ["missing-role", "unknown-phase", "wrong-session", "non-string-role", "invalid-json"]) +def test_corrupt_marker_fails_closed_without_changing_it(store, mutation): + stage(store) + record = store.transfer_fence() + if mutation == "missing-role": + record.pop("role") + elif mutation == "unknown-phase": + record["phase"] = "released" + elif mutation == "wrong-session": + record["session_id"] = "other" + elif mutation == "non-string-role": + record["role"] = [] + store.transfer_fence_path.write_text("{" if mutation == "invalid-json" else json.dumps(record)) + before = store.transfer_fence_path.read_bytes() + for acquire in (lambda: store.acquire(app="ordinary"), lambda: store.acquire_transfer("transfer-1", app="recovery")): + with pytest.raises(SharedStateError): + acquire() + assert store.transfer_fence_path.read_bytes() == before + + +def test_marker_and_native_directory_symlinks_are_rejected(store, tmp_path): + stage(store) + original = store.transfer_fence_path.read_bytes() + target = tmp_path / "marker.json" + target.write_bytes(original) + target.chmod(0o600) + store.transfer_fence_path.unlink() + store.transfer_fence_path.symlink_to(target) + with pytest.raises(SharedStateError, match="safe regular"): + store.acquire(app="ordinary") + store.transfer_fence_path.unlink() + directory = store.transfer_fence_path.parent + moved = directory.with_name("moved") + directory.rename(moved) + directory.symlink_to(moved, target_is_directory=True) + with pytest.raises(SharedStateError, match="directory is unsafe"): + store.acquire(app="ordinary") + assert target.read_bytes() == original + + +def test_released_source_handle_and_mismatched_fence_cannot_change_state(store): + held = store.acquire(app="source") + first = held.fence_transfer("transfer-1", "destination") + with pytest.raises(SharedStateError, match="different transfer"): + held.fence_transfer("transfer-2", "destination") + held.release() + with pytest.raises(RuntimeError, match="no longer active"): + held.fence_transfer("transfer-1", "destination") + assert store.transfer_fence() == first + + +def test_recovery_requires_an_existing_marker(store): + with pytest.raises(SharedStateError, match="exact staged"): + store.acquire_transfer("transfer-1", app="recovery") + held = store.acquire(app="ordinary") + held.release() + + +def test_nonprivate_marker_fails_closed(store): + stage(store) + store.transfer_fence_path.chmod(0o644) + with pytest.raises(SharedStateError, match="unsafe permissions"): + store.acquire(app="ordinary") From 59b7b17dccba54a08f1653157aea42e00f9b07cb Mon Sep 17 00:00:00 2001 From: Brian Krabach Date: Tue, 22 Sep 2026 08:40:52 -0700 Subject: [PATCH 2/3] Persist new native directory ancestors before transfer fences --- amplifier_foundation/session/shared_state.py | 32 ++++++++- docs/SESSION_TRANSFER_FENCE.md | 4 ++ tests/test_session_transfer_fence.py | 71 ++++++++++++++++++++ 3 files changed, 106 insertions(+), 1 deletion(-) diff --git a/amplifier_foundation/session/shared_state.py b/amplifier_foundation/session/shared_state.py index 8ad49580..c23f665d 100644 --- a/amplifier_foundation/session/shared_state.py +++ b/amplifier_foundation/session/shared_state.py @@ -164,6 +164,36 @@ def _validate_private_file(path: Path) -> None: raise SharedStateError(f"state file {path.name} has unsafe permissions") +def _mkdir_private_durable(path: Path) -> None: + """Create at most 128 private ancestors, persisting each new directory entry. + + A marker's own parent fsync cannot make a newly created ancestor durable. + Build from the existing ancestor outward and sync each containing directory + before publishing any descendant. Existing directory modes are untouched. + """ + missing = [] + current = path + while True: + try: + current.lstat() + except FileNotFoundError: + if len(missing) >= 128 or current == current.parent: + raise SharedStateError("native transfer directory exceeds the creation depth limit") from None + missing.append(current) + current = current.parent + else: + if not current.is_dir(): + raise SharedStateError("native transfer ancestor is not a directory") + break + for directory in reversed(missing): + try: + directory.mkdir(mode=0o700) + except FileExistsError: + pass + _validate_private_directory(directory, create=False) + _fsync_directory(directory.parent) + + def _ensure_private_state_path(root: Path, directory: Path) -> None: """Create the state-root chain, validating each state-owned component.""" @@ -392,7 +422,7 @@ def _transfer_directory(self, *, create: bool = False) -> None: except FileNotFoundError: if not create: return - directory.mkdir(mode=0o700, parents=True, exist_ok=True) + _mkdir_private_durable(directory) info = directory.lstat() if not stat.S_ISDIR(info.st_mode) or info.st_uid != os.getuid(): raise SharedStateError("native session transfer directory is unsafe") diff --git a/docs/SESSION_TRANSFER_FENCE.md b/docs/SESSION_TRANSFER_FENCE.md index 89c7af7f..d2060954 100644 --- a/docs/SESSION_TRANSFER_FENCE.md +++ b/docs/SESSION_TRANSFER_FENCE.md @@ -30,6 +30,10 @@ finally: boundary first. It writes a staged marker atomically, fsyncs its directory where supported, and makes that handle fail execution checks and checkpoint writes. An identical call is idempotent; a different existing marker is rejected. +If native history does not yet exist, creation is bounded to 128 missing +ancestors at a time. Each new directory is private and its containing directory +is fsynced before descendants or the marker are published, where supported. +Existing native directory modes are unchanged. The application exports saved data without replaying inputs or uncertain effects. After authenticating destination readiness, it permanently commits the source: diff --git a/tests/test_session_transfer_fence.py b/tests/test_session_transfer_fence.py index 52fd998e..0611e05e 100644 --- a/tests/test_session_transfer_fence.py +++ b/tests/test_session_transfer_fence.py @@ -4,6 +4,7 @@ import json import os +from pathlib import Path import subprocess import sys @@ -198,3 +199,73 @@ def test_nonprivate_marker_fails_closed(store): store.transfer_fence_path.chmod(0o644) with pytest.raises(SharedStateError, match="unsafe permissions"): store.acquire(app="ordinary") + + +def test_fresh_native_root_persists_all_ancestors_before_marker(tmp_path, monkeypatch): + import amplifier_foundation.session.shared_state as shared_state + + native = tmp_path / "new-parent" / "new-account" / "native" + monkeypatch.setenv("AMPLIFIER_HOME", str(native)) + workspace = tmp_path / "workspace" + workspace.mkdir() + store = SharedSessionStore(workspace, "fresh-task", root=tmp_path / "locks") + held = store.acquire(app="destination") + events = [] + sync = shared_state._fsync_directory + atomic = shared_state._atomic_json + + def recorded_sync(directory): + sync(directory) + events.append(("synced", Path(directory))) + + def recorded_atomic(path, value): + events.append(("published", Path(path))) + return atomic(path, value) + + monkeypatch.setattr(shared_state, "_fsync_directory", recorded_sync) + monkeypatch.setattr(shared_state, "_atomic_json", recorded_atomic) + try: + held.fence_transfer("fresh-transfer", "destination", role="destination") + marker_index = events.index(("published", store.transfer_fence_path)) + created = [] + directory = store.transfer_fence_path.parent + while directory != tmp_path: + created.append(directory) + directory = directory.parent + expected = [("synced", directory.parent) for directory in reversed(created)] + assert events[:marker_index] == expected + assert ("synced", store.transfer_fence_path.parent) in events[marker_index + 1:] + assert all(directory.stat().st_mode & 0o777 == 0o700 for directory in created) + assert not (store.transfer_fence_path.parent / "transcript.jsonl").exists() + finally: + held.release() + with pytest.raises(SessionTransferFencedError): + SharedSessionStore(workspace, "fresh-task", root=store.root).acquire(app="another-host") + + +def test_existing_native_directories_keep_modes_and_skip_creation_sync(store, monkeypatch): + import amplifier_foundation.session.shared_state as shared_state + + directory = store.transfer_fence_path.parent + directory.mkdir(parents=True, mode=0o755) + directory.chmod(0o755) + root = Path(os.environ["AMPLIFIER_HOME"]) + existing = [directory, *directory.parents] + existing = existing[:existing.index(root) + 1] + before = {path: path.stat().st_mode for path in existing} + + def no_creation(path): + pytest.fail("existing native history must not create ancestors") + + monkeypatch.setattr(shared_state, "_mkdir_private_durable", no_creation) + stage(store) + assert {path: path.stat().st_mode for path in existing} == before + + +def test_private_ancestor_creation_is_bounded_and_has_no_partial_tree(tmp_path): + from amplifier_foundation.session.shared_state import _mkdir_private_durable + + requested = tmp_path.joinpath(*["d"] * 129) + with pytest.raises(SharedStateError, match="depth limit"): + _mkdir_private_durable(requested) + assert not (tmp_path / "d").exists() From 907c442b17a83707102c01e80cb5d1a115fdff88 Mon Sep 17 00:00:00 2001 From: Brian Krabach Date: Tue, 22 Sep 2026 08:52:20 -0700 Subject: [PATCH 3/3] Fail transfer acknowledgments on persistence errors and verify retries --- amplifier_foundation/session/shared_state.py | 77 ++++++-- docs/SESSION_TRANSFER_FENCE.md | 27 ++- tests/test_session_transfer_fence.py | 177 ++++++++++++++++++- 3 files changed, 264 insertions(+), 17 deletions(-) diff --git a/amplifier_foundation/session/shared_state.py b/amplifier_foundation/session/shared_state.py index c23f665d..086b80dc 100644 --- a/amplifier_foundation/session/shared_state.py +++ b/amplifier_foundation/session/shared_state.py @@ -291,19 +291,39 @@ def _atomic_json(path: Path, value: dict[str, Any]) -> None: def _fsync_directory(directory: Path) -> None: """Persist a directory entry where the local POSIX filesystem supports it.""" - - try: - directory_fd = os.open(directory, os.O_RDONLY) - except OSError: - return + directory_fd = os.open(directory, os.O_RDONLY | getattr(os, "O_DIRECTORY", 0)) try: os.fsync(directory_fd) - except OSError: - pass # Some POSIX filesystems do not support directory fsync. + except OSError as exc: + # These errors specifically mean directory fsync is unsupported. I/O, + # space, permission, bad-descriptor and open errors must prevent ack. + unsupported = {errno.EINVAL, errno.ENOSYS, errno.ENOTSUP, errno.EOPNOTSUPP} + if exc.errno not in unsupported: + raise finally: os.close(directory_fd) +def _fsync_ancestor_chain(directory: Path) -> None: + """Re-establish ancestor durability, including after a failed earlier mkdir. + + An existing directory may be the residue of a failed unacknowledged write, + so existence cannot stand in for parent durability on a retry. Bound the + complete chain before doing I/O, then sync from filesystem root inward. + """ + ancestors = [] + current = directory.absolute() + while True: + if len(ancestors) >= 128: + raise SharedStateError("native transfer directory exceeds the durability depth limit") + ancestors.append(current) + if current == current.parent: + break + current = current.parent + for ancestor in reversed(ancestors): + _fsync_directory(ancestor) + + def _process_start_identity(path: Path = Path("/proc/self/stat")) -> str | None: """Read Linux procfs field 22, handling process names containing spaces.""" @@ -446,6 +466,11 @@ def transfer_fence(self) -> dict[str, Any] | None: raise SharedStateError("invalid session transfer fence") return value + def _save_transfer_fence(self, record: dict[str, Any]) -> None: + self._transfer_directory(create=True) + _fsync_ancestor_chain(self.transfer_fence_path.parent) + _atomic_json(self.transfer_fence_path, record) + @property def checkpoint_path(self) -> Path: """The authoritative checkpoint path; obtaining it never creates directories.""" @@ -497,7 +522,22 @@ def acquire_transfer(self, transfer_id: str, *, app: str, **diagnostics: Any) -> _session_id(transfer_id) return cast("HeldTransfer", self._acquire(app=app, diagnostics=diagnostics, transfer_id=transfer_id)) - def _acquire(self, *, app: str, diagnostics: dict[str, Any], transfer_id: str | None = None) -> "HeldSession": + def confirm_transfer_commit(self, transfer_id: str, *, app: str, **diagnostics: Any) -> dict[str, Any]: + """Durably confirm an exact committed source after an uncertain ack. + + This returns evidence only, never an execution or clearing capability. + It cannot commit a staged marker or change the recorded destination. + """ + _session_id(transfer_id) + held = cast("HeldTransfer", self._acquire(app=app, diagnostics=diagnostics, + transfer_id=transfer_id, committed_only=True)) + try: + return held.commit_transfer() + finally: + held.release() + + def _acquire(self, *, app: str, diagnostics: dict[str, Any], transfer_id: str | None = None, + committed_only: bool = False) -> "HeldSession": _ensure_supported() owner = _owner_details(app, diagnostics, self.workspace, self.session_id, self.root) _ensure_private_state_path(self.root, self._directory) @@ -526,6 +566,10 @@ def _acquire(self, *, app: str, diagnostics: dict[str, Any], transfer_id: str | if transfer_id is None: if fence is not None: raise SessionTransferFencedError(fence) + elif committed_only: + if (fence is None or fence["transfer_id"] != transfer_id + or fence["phase"] != "committed" or fence["role"] != "source"): + raise SharedStateError("commit confirmation requires the exact committed source fence") elif fence is None or fence["transfer_id"] != transfer_id or fence["phase"] != "staged": raise SharedStateError("transfer recovery requires the exact staged fence") _atomic_json(self._owner_path, owner) @@ -616,11 +660,11 @@ def fence_transfer(self, transfer_id: str, destination_host: str, *, role: str = self._check_active() existing = self._store.transfer_fence() if existing is not None: - if existing == record: - return existing - raise SharedStateError("a different transfer fence already exists") - self._store._transfer_directory(create=True) - _atomic_json(self._store.transfer_fence_path, record) + if existing != record: + raise SharedStateError("a different transfer fence already exists") + # Even an identical marker can be left by a failed fsync. Repeat + # file and ancestor durability before acknowledging that request. + self._store._save_transfer_fence(record) return copy.deepcopy(record) def read(self) -> dict[str, Any] | None: @@ -724,11 +768,14 @@ def _staged(self) -> dict[str, Any]: def commit_transfer(self) -> dict[str, Any]: """Permanently fence the source after the host verifies its transfer.""" with self._mutex: - record = self._staged() + self._check_active() + record = self._store.transfer_fence() + if record is None or record["transfer_id"] != self._transfer_id: + raise SharedStateError("transfer commit requires the exact source fence") if record["role"] != "source": raise SharedStateError("only a source transfer fence can be committed") record["phase"] = "committed" - _atomic_json(self._store.transfer_fence_path, record) + self._store._save_transfer_fence(record) return record def clear_transfer(self) -> None: diff --git a/docs/SESSION_TRANSFER_FENCE.md b/docs/SESSION_TRANSFER_FENCE.md index d2060954..6e4522a1 100644 --- a/docs/SESSION_TRANSFER_FENCE.md +++ b/docs/SESSION_TRANSFER_FENCE.md @@ -29,7 +29,8 @@ finally: `fence_transfer` does not stop work or save history. The host must establish that boundary first. It writes a staged marker atomically, fsyncs its directory where supported, and makes that handle fail execution checks and checkpoint writes. -An identical call is idempotent; a different existing marker is rejected. +An identical call is idempotent; it repeats file and ancestor durability before +acknowledging the marker. A different existing marker is rejected. If native history does not yet exist, creation is bounded to 128 missing ancestors at a time. Each new directory is private and its containing directory is fsynced before descendants or the marker are published, where supported. @@ -52,6 +53,22 @@ needs a new destination workspace/native history location. Preserve the source marker when archiving inactive history. A failed or unknown application-level transfer must retain its marker and expose its evidence for reconciliation. +A marker can become visible even when its final directory sync fails. A raw +`transfer_fence()` read therefore does not prove a durable commit acknowledgment. +The same live `HeldTransfer.commit_transfer()` can retry an exact committed +source, repeating marker and ancestor sync. After releasing that handle or +restarting, an application can confirm the existing committed source without +obtaining execution or clearing authority: + +```python +committed = store.confirm_transfer_commit(transfer_id, app="transfer-adapter") +``` + +This takes the exclusive lock, requires the exact committed source transfer, +re-establishes durability, returns the record and releases the lock. It cannot +commit a staged transfer or alter its destination. Confirm durability before +issuing or reissuing a remote release certificate after an uncertain write. + ## Destination staging and explicit resolution Before installing transferred native history, an application acquires its new @@ -101,6 +118,14 @@ Ordinary `acquire()` raises `SessionTransferFencedError`, whose `fence` attribut is a detached record. `acquire_transfer()` requires the exact staged transfer; passing a transfer ID among ordinary acquisition diagnostics never bypasses it. +Genuine file or directory open, I/O, space and permission failures propagate; +the caller must treat the acknowledgment as failed or uncertain. Directory +`fsync` alone may return `EINVAL`, `ENOSYS`, `ENOTSUP` or `EOPNOTSUPP` on filesystems +without this operation; those cases retain the qualified best-effort directory +durability boundary. File fsync errors are never ignored. Transfer marker saves +sync a bounded complete ancestor chain before file replacement, even for +directories already present after an earlier failed creation attempt. + The marker is addressed through native history rather than the configurable coordination root, so choosing another coordination root does not erase this admission check. Applications must still share the coordination root to exclude diff --git a/tests/test_session_transfer_fence.py b/tests/test_session_transfer_fence.py index 0611e05e..4fdcd1e6 100644 --- a/tests/test_session_transfer_fence.py +++ b/tests/test_session_transfer_fence.py @@ -3,9 +3,11 @@ from __future__ import annotations import json +import errno import os from pathlib import Path import subprocess +import stat import sys import pytest @@ -233,7 +235,9 @@ def recorded_atomic(path, value): created.append(directory) directory = directory.parent expected = [("synced", directory.parent) for directory in reversed(created)] - assert events[:marker_index] == expected + assert events[:len(expected)] == expected + ancestors = [store.transfer_fence_path.parent, *store.transfer_fence_path.parent.parents] + assert events[len(expected):marker_index] == [("synced", path) for path in reversed(ancestors)] assert ("synced", store.transfer_fence_path.parent) in events[marker_index + 1:] assert all(directory.stat().st_mode & 0o777 == 0o700 for directory in created) assert not (store.transfer_fence_path.parent / "transcript.jsonl").exists() @@ -269,3 +273,174 @@ def test_private_ancestor_creation_is_bounded_and_has_no_partial_tree(tmp_path): with pytest.raises(SharedStateError, match="depth limit"): _mkdir_private_durable(requested) assert not (tmp_path / "d").exists() + + +@pytest.mark.parametrize("error", [errno.EIO, errno.ENOSPC, errno.EACCES, errno.EBADF]) +def test_directory_fsync_real_failures_propagate_and_close_descriptor(tmp_path, monkeypatch, error): + import amplifier_foundation.session.shared_state as shared_state + + descriptors = [] + + def fail(fd): + descriptors.append(fd) + raise OSError(error, "injected persistence failure") + + monkeypatch.setattr(shared_state.os, "fsync", fail) + with pytest.raises(OSError) as failure: + shared_state._fsync_directory(tmp_path) + assert failure.value.errno == error + with pytest.raises(OSError) as closed: + os.fstat(descriptors[0]) + assert closed.value.errno == errno.EBADF + + +def test_directory_open_failure_is_not_treated_as_unsupported(tmp_path, monkeypatch): + import amplifier_foundation.session.shared_state as shared_state + + def fail(*args, **kwargs): + raise OSError(errno.EACCES, "injected open failure") + + monkeypatch.setattr(shared_state.os, "open", fail) + with pytest.raises(OSError) as failure: + shared_state._fsync_directory(tmp_path) + assert failure.value.errno == errno.EACCES + + +@pytest.mark.parametrize("error", sorted({errno.EINVAL, errno.ENOSYS, errno.ENOTSUP, errno.EOPNOTSUPP})) +def test_only_explicitly_unsupported_directory_fsync_is_tolerated(tmp_path, monkeypatch, error): + import amplifier_foundation.session.shared_state as shared_state + + def unsupported(fd): + raise OSError(error, "directory fsync unsupported") + + monkeypatch.setattr(shared_state.os, "fsync", unsupported) + shared_state._fsync_directory(tmp_path) + + +def test_failed_ancestor_sync_never_acks_and_retry_resyncs_existing_residue(store, monkeypatch): + import amplifier_foundation.session.shared_state as shared_state + + held = store.acquire(app="source") + sync = shared_state._fsync_directory + home = Path(os.environ["AMPLIFIER_HOME"]) + events = [] + try: + def fail_parent(directory): + if directory == home.parent: + raise OSError(errno.EIO, "injected ancestor failure") + sync(directory) + + monkeypatch.setattr(shared_state, "_fsync_directory", fail_parent) + with pytest.raises(OSError, match="ancestor failure"): + held.fence_transfer("transfer-1", "destination") + assert home.exists() and not store.transfer_fence_path.exists() + + def record(directory): + sync(directory) + events.append(directory) + + monkeypatch.setattr(shared_state, "_fsync_directory", record) + assert held.fence_transfer("transfer-1", "destination")["phase"] == "staged" + # The root existed on retry, but its earlier parent sync never succeeded. + assert home.parent in events + assert store.transfer_fence_path.parent in events + finally: + monkeypatch.setattr(shared_state, "_fsync_directory", sync) + held.release() + + +def test_staged_marker_visible_after_failed_sync_requires_new_durability_ack(store, monkeypatch): + import amplifier_foundation.session.shared_state as shared_state + + held = store.acquire(app="source") + sync = os.fsync + temporary_marker = store.transfer_fence_path + try: + def fail_after_replace(fd): + if stat.S_ISDIR(os.fstat(fd).st_mode) and temporary_marker.exists(): + raise OSError(errno.EIO, "injected marker-directory failure") + sync(fd) + + monkeypatch.setattr(shared_state.os, "fsync", fail_after_replace) + for _ in range(2): + with pytest.raises(OSError, match="marker-directory failure"): + held.fence_transfer("transfer-1", "destination") + assert store.transfer_fence()["phase"] == "staged" + monkeypatch.setattr(shared_state.os, "fsync", sync) + file_syncs = [] + + def record(fd): + file_syncs.append(stat.S_ISREG(os.fstat(fd).st_mode)) + sync(fd) + + monkeypatch.setattr(shared_state.os, "fsync", record) + assert held.fence_transfer("transfer-1", "destination")["phase"] == "staged" + assert True in file_syncs and False in file_syncs + finally: + monkeypatch.setattr(shared_state.os, "fsync", sync) + held.release() + + +def test_commit_sync_failure_never_acks_and_exact_live_retry_can_confirm(store, monkeypatch): + import amplifier_foundation.session.shared_state as shared_state + + stage(store) + held = store.acquire_transfer("transfer-1", app="source") + sync = os.fsync + try: + def fail_committed(fd): + if stat.S_ISDIR(os.fstat(fd).st_mode) and store.transfer_fence()["phase"] == "committed": + raise OSError(errno.EIO, "injected commit failure") + sync(fd) + + monkeypatch.setattr(shared_state.os, "fsync", fail_committed) + for _ in range(2): + with pytest.raises(OSError, match="commit failure"): + held.commit_transfer() + assert store.transfer_fence()["phase"] == "committed" + monkeypatch.setattr(shared_state.os, "fsync", sync) + assert held.commit_transfer()["phase"] == "committed" + with pytest.raises(SharedStateError, match="exact staged"): + held.clear_transfer() + finally: + monkeypatch.setattr(shared_state.os, "fsync", sync) + held.release() + + +def test_restart_commit_confirmation_checks_exact_fence_and_repeats_durability(store, monkeypatch): + import amplifier_foundation.session.shared_state as shared_state + + stage(store) + with pytest.raises(SharedStateError, match="exact committed source"): + store.confirm_transfer_commit("transfer-1", app="confirmer") + held = store.acquire_transfer("transfer-1", app="source") + held.commit_transfer() + held.release() + fresh = SharedSessionStore(store.workspace, store.session_id, root=store.root) + with pytest.raises(SharedStateError, match="exact committed source"): + fresh.confirm_transfer_commit("different-transfer", app="confirmer") + with pytest.raises(SharedStateError, match="exact staged"): + fresh.acquire_transfer("transfer-1", app="recovery") + sync = shared_state._fsync_directory + events = [] + + def fail(directory): + if directory == fresh.transfer_fence_path.parent: + raise OSError(errno.EIO, "injected confirmation failure") + sync(directory) + + monkeypatch.setattr(shared_state, "_fsync_directory", fail) + with pytest.raises(OSError, match="confirmation failure"): + fresh.confirm_transfer_commit("transfer-1", app="confirmer") + + def record(directory): + sync(directory) + events.append(directory) + + monkeypatch.setattr(shared_state, "_fsync_directory", record) + confirmed = fresh.confirm_transfer_commit("transfer-1", app="confirmer") + assert confirmed["phase"] == "committed" + assert fresh.transfer_fence_path.parent in events + assert Path(os.environ["AMPLIFIER_HOME"]).parent in events + with pytest.raises(SessionTransferFencedError): + fresh.acquire(app="ordinary")