diff --git a/scripts/check_journal_io_access.py b/scripts/check_journal_io_access.py index fd1461225..574554668 100644 --- a/scripts/check_journal_io_access.py +++ b/scripts/check_journal_io_access.py @@ -59,12 +59,16 @@ GATED_PRIMITIVES: frozenset[str] = frozenset( "append_text", "atomic_replace", "acquire_file_lease", + "adopt_inherited_file_lease_fd", "assert_file_lease_owned", "hold_lock", "install_file", "probe_file_lease_free", "probe_file_lease_held", + "read_file_lease_fd", + "read_file_lease_offset_token", "save_npz", + "set_file_lease_offset_token", "update_npz", "write_bytes_exclusive", "write_json", @@ -91,9 +95,13 @@ MODULE_PRIMITIVES: dict[str, frozenset[str]] = { "solstone.think.journal_io.lease": frozenset( { "acquire_file_lease", + "adopt_inherited_file_lease_fd", "assert_file_lease_owned", "probe_file_lease_free", "probe_file_lease_held", + "read_file_lease_fd", + "read_file_lease_offset_token", + "set_file_lease_offset_token", } ), "solstone.think.journal_io.npz": frozenset({"save_npz", "update_npz", "write_npz"}), diff --git a/solstone/think/journal_io/lease.py b/solstone/think/journal_io/lease.py index 720f50f09..3214a3dcf 100644 --- a/solstone/think/journal_io/lease.py +++ b/solstone/think/journal_io/lease.py @@ -53,6 +53,27 @@ class FileLease: self.release() +@dataclass +class BorrowedFileLease: + """Borrowed file lease backed by a duplicated flock handle.""" + + path: Path + _fd: int | None + + def release(self) -> None: + if self._fd is None: + return + fd = self._fd + self._fd = None + os.close(fd) + + def __enter__(self) -> BorrowedFileLease: + return self + + def __exit__(self, exc_type: Any, exc: Any, tb: Any) -> None: + self.release() + + def _retry_sleep_seconds(deadline: float) -> float: remaining = deadline - time.monotonic() if remaining <= 0: @@ -114,6 +135,31 @@ def assert_file_lease_owned( return lease +def read_file_lease_fd(lease: FileLease, path: Path | None = None) -> int: + """Return the owned lease fd without transferring ownership.""" + + owned = assert_file_lease_owned(lease, path) + assert owned._fd is not None + return owned._fd + + +def set_file_lease_offset_token( + lease: FileLease, token: int, path: Path | None = None +) -> None: + """Set a nonzero offset token on an owned lease fd.""" + + if token <= 0: + raise ValueError("file lease offset token must be nonzero") + fd = read_file_lease_fd(lease, path) + os.lseek(fd, token, os.SEEK_SET) + + +def read_file_lease_offset_token(fd: int) -> int: + """Read an fd's current offset token without moving it.""" + + return os.lseek(fd, 0, os.SEEK_CUR) + + def probe_file_lease_held(path: Path) -> bool: """Return whether an existing lease file is currently held by another process.""" @@ -140,10 +186,46 @@ def probe_file_lease_free(path: Path) -> bool: return not probe_file_lease_held(path) +def adopt_inherited_file_lease_fd( + path: Path, fd: int, token: int +) -> BorrowedFileLease | None: + """Adopt an inherited duplicate only when it proves the held lease.""" + + try: + candidate_stat = os.fstat(fd) + path_stat = os.stat(path) + except OSError: + return None + if ( + candidate_stat.st_dev != path_stat.st_dev + or candidate_stat.st_ino != path_stat.st_ino + ): + return None + try: + if token == 0 or read_file_lease_offset_token(fd) != token: + return None + if not probe_file_lease_held(path): + return None + try: + fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB) + except OSError as exc: + if exc.errno in (errno.EACCES, errno.EAGAIN, errno.EWOULDBLOCK): + return None + raise + return BorrowedFileLease(path=path, _fd=os.dup(fd)) + except OSError: + return None + + __all__ = [ + "BorrowedFileLease", "FileLease", "acquire_file_lease", + "adopt_inherited_file_lease_fd", "assert_file_lease_owned", "probe_file_lease_free", "probe_file_lease_held", + "read_file_lease_fd", + "read_file_lease_offset_token", + "set_file_lease_offset_token", ] diff --git a/tests/test_journal_io.py b/tests/test_journal_io.py index 27ac44ddd..8e36140ad 100644 --- a/tests/test_journal_io.py +++ b/tests/test_journal_io.py @@ -9,6 +9,7 @@ import json import logging import os import stat +import subprocess import sys from pathlib import Path @@ -28,10 +29,15 @@ from solstone.think.journal_io.errors import ( PathEscapeError, ) from solstone.think.journal_io.lease import ( + BorrowedFileLease, acquire_file_lease, + adopt_inherited_file_lease_fd, assert_file_lease_owned, probe_file_lease_free, probe_file_lease_held, + read_file_lease_fd, + read_file_lease_offset_token, + set_file_lease_offset_token, ) from solstone.think.journal_io.locking import hold_lock from solstone.think.journal_io.paths import contained_path @@ -88,6 +94,107 @@ def test_file_lease_holder_and_probe(tmp_path) -> None: assert probe_file_lease_free(path) is True +def test_borrowed_file_lease_adoption_matrix(tmp_path) -> None: + path = tmp_path / "health" / "speakers-analyze.lock" + lease = acquire_file_lease(path, attempts=1) + assert lease is not None + token = 37 + set_file_lease_offset_token(lease, token, path) + owner_fd = read_file_lease_fd(lease, path) + + duplicate_fd = os.dup(owner_fd) + try: + borrowed = adopt_inherited_file_lease_fd(path, duplicate_fd, token) + assert isinstance(borrowed, BorrowedFileLease) + borrowed.release() + assert probe_file_lease_held(path) is True + assert acquire_file_lease(path, attempts=1) is None + + assert adopt_inherited_file_lease_fd(path, duplicate_fd, token + 1) is None + assert adopt_inherited_file_lease_fd(path, duplicate_fd, 0) is None + + wrong_path = tmp_path / "health" / "other.lock" + wrong_path.write_text("", encoding="utf-8") + assert adopt_inherited_file_lease_fd(wrong_path, duplicate_fd, token) is None + + separate_fd = os.open(path, os.O_RDWR) + try: + assert read_file_lease_offset_token(separate_fd) == 0 + assert adopt_inherited_file_lease_fd(path, separate_fd, token) is None + assert read_file_lease_offset_token(separate_fd) == 0 + finally: + os.close(separate_fd) + + closed_fd = os.dup(owner_fd) + os.close(closed_fd) + assert adopt_inherited_file_lease_fd(path, closed_fd, token) is None + + reused_fd = os.dup(owner_fd) + os.close(reused_fd) + reused_path = tmp_path / "health" / "reused.lock" + actual_reused_fd = os.open(reused_path, os.O_RDWR | os.O_CREAT, 0o600) + try: + assert actual_reused_fd == reused_fd + assert adopt_inherited_file_lease_fd(path, actual_reused_fd, token) is None + finally: + os.close(actual_reused_fd) + finally: + os.close(duplicate_fd) + lease.release() + + assert probe_file_lease_free(path) is True + + +def test_inherited_duplicate_preserves_token_and_lock_after_owner_close( + tmp_path, +) -> None: + path = tmp_path / "health" / "speakers-analyze.lock" + lease = acquire_file_lease(path, attempts=1) + assert lease is not None + token = 91 + set_file_lease_offset_token(lease, token, path) + owner_fd = read_file_lease_fd(lease, path) + + script = ( + "import fcntl,json,os,sys\n" + "fd=int(sys.argv[1])\n" + "path=sys.argv[2]\n" + "fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB)\n" + "print(json.dumps({'fd': fd, 'offset': os.lseek(fd, 0, os.SEEK_CUR)}), flush=True)\n" + "sys.stdin.read()\n" + ) + child = subprocess.Popen( + [sys.executable, "-c", script, str(owner_fd), str(path)], + stdin=subprocess.PIPE, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + pass_fds=(owner_fd,), + ) + try: + assert child.stdout is not None + payload = json.loads(child.stdout.readline()) + assert payload == {"fd": owner_fd, "offset": token} + + os.close(owner_fd) + lease._fd = None + assert acquire_file_lease(path, attempts=1) is None + + assert child.stdin is not None + child.stdin.close() + assert child.wait(timeout=5) == 0 + finally: + if child.poll() is None: + child.kill() + child.wait(timeout=5) + if lease._fd is not None: + lease.release() + + reacquired = acquire_file_lease(path, attempts=1) + assert reacquired is not None + reacquired.release() + + def test_install_file_crash_safe(tmp_path, monkeypatch) -> None: dest = tmp_path / "audio.opus" dest.write_bytes(b"OLD")