diff --git a/src/wf_scheduling/ownership.py b/src/wf_scheduling/ownership.py new file mode 100644 index 00000000..38dfb976 --- /dev/null +++ b/src/wf_scheduling/ownership.py @@ -0,0 +1,123 @@ +"""Exclusive file-store ownership via a held cross-process lock (T11). + +The scheduler process owns the store exclusively before recovery. Ownership +is a held lock on a dedicated lock file (not a stale PID file or an +expiring lease): the handle stays open for the owner's lifetime and the OS +releases it on process death. A second owner attempting acquisition while +the lock is held is rejected. Where locking is unsupported, scheduler +startup is rejected instead of running unprotected. + +Locking design (Windows-tested, per the store transaction boundary): one +new module holding ``msvcrt.locking`` (Windows) / ``fcntl.flock`` (POSIX) +on ``/scheduler.lock``. No ad-hoc per-run lock files. +""" + +from __future__ import annotations + +from pathlib import Path +from typing import BinaryIO + + +class SecondOwnerError(Exception): + """Another process already owns the schedule store.""" + + +class StartupRejected(Exception): + """Scheduler startup rejected where locking is unsupported.""" + + +class SchedulerOwnership: + """Held exclusive ownership of a schedule store root.""" + + def __init__(self, root: Path, *, owner: str) -> None: + self.root = root + self.owner = owner + self._handle: BinaryIO | None = None + self._locked = False + + @property + def lock_path(self) -> Path: + return self.root / "scheduler.lock" + + def acquire(self) -> SchedulerOwnership: + """Acquire the held lock non-blockingly or raise SecondOwnerError.""" + self.lock_path.parent.mkdir(parents=True, exist_ok=True) + try: + handle = open(self.lock_path, "a+b") + except OSError as exc: + raise StartupRejected(f"cannot open scheduler lock: {exc}") from exc + try: + _lock_nonblocking(handle) + except SecondOwnerError: + handle.close() + raise + except StartupRejected: + handle.close() + raise + except OSError as exc: + handle.close() + raise StartupRejected(f"unsupported scheduler locking: {exc}") from exc + self._handle = handle + self._locked = True + return self + + def release(self) -> None: + """Release the held lock (idempotent).""" + if not self._locked or self._handle is None: + return + try: + _unlock(self._handle) + finally: + try: + self._handle.close() + finally: + self._handle = None + self._locked = False + + def __enter__(self) -> SchedulerOwnership: + return self.acquire() + + def __exit__(self, *exc: object) -> None: + self.release() + + +def _lock_nonblocking(handle: BinaryIO) -> None: + import os + + if os.name == "nt": + import msvcrt + + handle.seek(0) + try: + msvcrt.locking(handle.fileno(), msvcrt.LK_NBLCK, 1) + except OSError as exc: + raise SecondOwnerError("schedule store owned by another process") from exc + return + try: + import fcntl + except ImportError as exc: + raise StartupRejected("file locking unsupported on this platform") from exc + try: + fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) + except OSError as exc: + raise SecondOwnerError("schedule store owned by another process") from exc + + +def _unlock(handle: BinaryIO) -> None: + import os + + if os.name == "nt": + import msvcrt + + try: + handle.seek(0) + msvcrt.locking(handle.fileno(), msvcrt.LK_UNLCK, 1) + except OSError: + pass + return + try: + import fcntl + + fcntl.flock(handle.fileno(), fcntl.LOCK_UN) + except OSError, ImportError: + pass diff --git a/tests/scheduling/test_ownership.py b/tests/scheduling/test_ownership.py new file mode 100644 index 00000000..6297679d --- /dev/null +++ b/tests/scheduling/test_ownership.py @@ -0,0 +1,98 @@ +"""Exclusive ownership with real competing processes (T11).""" + +from __future__ import annotations + +import subprocess +import sys +import time +from pathlib import Path + +import pytest + +from wf_scheduling.ownership import ( + SchedulerOwnership, + SecondOwnerError, +) + + +def test_second_owner_rejected_while_held(tmp_path: Path) -> None: + first = SchedulerOwnership(tmp_path, owner="proc-A").acquire() + try: + with pytest.raises(SecondOwnerError): + SchedulerOwnership(tmp_path, owner="proc-B").acquire() + finally: + first.release() + # After release, a new owner acquires. + second = SchedulerOwnership(tmp_path, owner="proc-B").acquire() + second.release() + + +def test_competing_process_cannot_acquire_while_held(tmp_path: Path) -> None: + holder = SchedulerOwnership(tmp_path, owner="parent").acquire() + try: + script_path = tmp_path / "compete.py" + script_path.write_text( + "import sys\n" + "sys.path.insert(0, '.')\n" + "from pathlib import Path\n" + "from wf_scheduling.ownership import SchedulerOwnership, SecondOwnerError\n" + f"root = Path({str(tmp_path)!r})\n" + "try:\n" + " SchedulerOwnership(root, owner='child').acquire()\n" + "except SecondOwnerError:\n" + " print('second-owner-rejected')\n" + " raise SystemExit(0)\n" + "print('unexpectedly-acquired')\n" + "raise SystemExit(1)\n", + encoding="utf-8", + ) + proc = subprocess.run( + [sys.executable, str(script_path)], + capture_output=True, + text=True, + timeout=30, + ) + assert proc.returncode == 0, proc.stderr + assert "second-owner-rejected" in proc.stdout + finally: + holder.release() + + +def test_lock_released_on_process_death(tmp_path: Path) -> None: + script_path = tmp_path / "hold.py" + script_path.write_text( + "import sys, time\n" + "sys.path.insert(0, '.')\n" + "from pathlib import Path\n" + "from wf_scheduling.ownership import SchedulerOwnership\n" + f"root = Path({str(tmp_path)!r})\n" + "own = SchedulerOwnership(root, owner='dying').acquire()\n" + "print('acquired', flush=True)\n" + "time.sleep(30)\n", + encoding="utf-8", + ) + proc = subprocess.Popen( + [sys.executable, str(script_path)], + stdout=subprocess.PIPE, + text=True, + ) + try: + assert proc.stdout is not None + line = proc.stdout.readline() + assert "acquired" in line + # Competing acquisition fails while the child lives. + with pytest.raises(SecondOwnerError): + SchedulerOwnership(tmp_path, owner="parent").acquire() + finally: + proc.kill() + proc.wait(timeout=30) + # After process death the OS releases the held lock; allow a beat. + acquired = None + for _ in range(50): + try: + acquired = SchedulerOwnership(tmp_path, owner="parent").acquire() + break + except SecondOwnerError: + time.sleep(0.1) + assert acquired is not None + acquired.release()