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..086b80dc 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: @@ -156,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.""" @@ -253,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.""" @@ -360,6 +418,58 @@ 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 + _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") + _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 + + 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: @@ -400,6 +510,34 @@ 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 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) @@ -424,6 +562,16 @@ 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 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) except BaseException: fcntl.flock(fd, fcntl.LOCK_UN) @@ -431,6 +579,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 +632,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: + 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: self.check() return self._store.read() @@ -562,3 +742,49 @@ 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: + 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" + self._store._save_transfer_fence(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..6e4522a1 --- /dev/null +++ b/docs/SESSION_TRANSFER_FENCE.md @@ -0,0 +1,134 @@ +# 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; 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. +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: + +```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. + +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 +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. + +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 +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..4fdcd1e6 --- /dev/null +++ b/tests/test_session_transfer_fence.py @@ -0,0 +1,446 @@ +"""Durable execution fences across host adapters and process lifetimes.""" + +from __future__ import annotations + +import json +import errno +import os +from pathlib import Path +import subprocess +import stat +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") + + +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[: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() + 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() + + +@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")