"""Shared flock-based JSON locking helper. Lock ordering invariant: always acquire locks in alphabetical order by filename to avoid deadlock when a caller holds multiple locks simultaneously. Implementation note (2026-04 — Lane C2 fix): We lock on a sidecar `.lock` file rather than the data file itself. That's because `write_locked` uses `os.replace(tmp, target)` for atomic publish — but `replace` swaps the inode behind `target`, which means a flock held on the *old* fd no longer guards the new file. Concurrent writers on a sidecar lockfile (whose inode is stable) get correct serialisation across threads and processes. """ import errno import fcntl import json import logging import os import threading import time from typing import Callable _TIMEOUT_SEC = 5.0 _POLL_INTERVAL = 0.05 _log = logging.getLogger(__name__) class LockTimeoutError(Exception): pass _local = threading.local() def _held_locks() -> dict: """Per-thread map: abspath → (lockfd, refcount). Used for re-entrancy.""" if not hasattr(_local, 'locks'): _local.locks = {} return _local.locks def _try_lock(fd: int, lock_type: int, timeout: float) -> bool: deadline = time.monotonic() + timeout while True: try: fcntl.flock(fd, lock_type | fcntl.LOCK_NB) return True except BlockingIOError: if time.monotonic() >= deadline: return False time.sleep(_POLL_INTERVAL) def _acquire(fd: int, lock_type: int) -> None: if _try_lock(fd, lock_type, _TIMEOUT_SEC): return if _try_lock(fd, lock_type, _TIMEOUT_SEC): return raise LockTimeoutError( f"could not acquire flock within {2 * _TIMEOUT_SEC}s (after retry)" ) def _open_lockfile(abspath: str) -> int: """Open (creating if needed) the sidecar `.lock` file.""" lock_path = abspath + ".lock" # Ensure the parent dir exists — write_locked auto-creates the data file # too, so we should be tolerant of the parent dir not having ever been # touched. parent = os.path.dirname(lock_path) if parent: try: os.makedirs(parent, exist_ok=True) except OSError as exc: if exc.errno != errno.EEXIST: raise return os.open(lock_path, os.O_RDWR | os.O_CREAT, 0o644) def read_locked(path: str) -> dict: abspath = os.path.abspath(path) held = _held_locks() if abspath in held: _log.debug("re-entrant read on %s; skipping flock", abspath) with open(abspath, 'r', encoding='utf-8') as f: return json.load(f) lock_fd = _open_lockfile(abspath) try: _acquire(lock_fd, fcntl.LOCK_SH) held[abspath] = (lock_fd, 1) try: with open(abspath, 'r', encoding='utf-8') as f: return json.load(f) finally: held.pop(abspath, None) try: fcntl.flock(lock_fd, fcntl.LOCK_UN) except OSError: pass finally: os.close(lock_fd) def write_locked(path: str, mutator: Callable[[dict], dict]) -> dict: abspath = os.path.abspath(path) held = _held_locks() reentrant = abspath in held lock_fd = -1 if reentrant else _open_lockfile(abspath) try: if not reentrant: _acquire(lock_fd, fcntl.LOCK_EX) held[abspath] = (lock_fd, 1) else: _log.debug("re-entrant write on %s; skipping flock", abspath) try: # Read current data (file may not exist yet — treat as {}). try: with open(abspath, 'r', encoding='utf-8') as f: text = f.read() data = json.loads(text) if text.strip() else {} except FileNotFoundError: data = {} new_data = mutator(data) # Atomic-rename invariant: tmp file MUST be on the same filesystem # as the target (sibling path guarantees this). tmp_path = abspath + ".tmp" with open(tmp_path, 'w', encoding='utf-8') as tmp: json.dump(new_data, tmp, indent=2) tmp.flush() os.fsync(tmp.fileno()) os.replace(tmp_path, abspath) return new_data finally: if not reentrant: held.pop(abspath, None) try: fcntl.flock(lock_fd, fcntl.LOCK_UN) except OSError: pass finally: if not reentrant and lock_fd >= 0: os.close(lock_fd)