Coverage for src/lilbee/runtime/lock.py: 100%
125 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-09-28 17:20 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-09-28 17:20 +0000
1"""Cross-process locks: LanceDB write locking, the sync mark, and the server singleton.
3Write locking combines an in-process mutex with a cross-process file lock
4(filelock) so separate processes also coordinate writes. Read consistency is
5handled by LanceDB's built-in MVCC via ``read_consistency_interval`` in
6``lilbee.data.store``. The server lock makes ``lilbee serve`` a singleton per
7data dir.
8"""
10import asyncio
11import json
12import logging
13import sqlite3
14import threading
15import time
16from collections.abc import AsyncGenerator, Generator
17from contextlib import asynccontextmanager, contextmanager
18from dataclasses import dataclass
19from pathlib import Path
21from filelock import FileLock, ReadWriteLock
22from filelock import Timeout as FileLockTimeout
24from lilbee.core.config import cfg
26log = logging.getLogger(__name__)
28# Default timeout (seconds) for acquiring the write lock
29LOCK_TIMEOUT = 30.0
30# Grace (seconds) for a dying predecessor to release the server lock during a
31# restart handoff before a new `lilbee serve` gives up. This budgets the PROMPT
32# path only: a predecessor whose llama-swaps honor SIGTERM releases well inside
33# it. It deliberately does NOT cover SIGKILL escalation -- the teardown that
34# actually holds the lock (_release_engines -> stop_engine -> _stop_stale_swap in
35# lilbee.providers.fleet.swap_manager) can spend _ORPHAN_STOP_TIMEOUT_S plus the
36# kill and reap waits per group, i.e. tens of seconds across four groups on a
37# wedged machine. Sizing this to that worst case would make every ordinary
38# restart wait on a pathological one; instead the successor exits with
39# LOCK_REFUSAL_EXIT_CODE and the operator retries.
40SERVER_LOCK_TIMEOUT = 15.0
41_SERVER_LOCK_NAME = "server.lock"
42_SCOPE_LOCK_NAME = "server.scope.lock"
43_SCOPE_OWNER_NAME = "server.scope.owner.json"
44_SYNC_LOCK_NAME = "sync.lock"
45# What the filesystem raises when the sync lock cannot be created or taken: OSError
46# from making the data root, sqlite3.Error from SQLite (a file that is not a
47# database, or a mount that refuses its locks). A busy lock arrives as filelock's
48# Timeout instead, which is itself an OSError, so it is caught first.
49_SYNC_LOCK_ERRORS = (OSError, sqlite3.Error)
50_SYNC_LOCK_REFUSED = "Cannot lock %s (%s); a reset refuses until it can."
51_SYNC_RUNNING = "A sync or import is running on this library. Reset again when it finishes."
52_SYNC_LOCK_UNKNOWN = (
53 "Cannot tell whether a sync or import is running: {path} cannot be locked ({error}). "
54 "Stop every lilbee process, delete {path}, and reset again."
55)
56# Minimum blocking wait granted to the in-process mutex even when the file lock
57# consumed the whole budget, so a deadline-edge acquire still gets a real attempt.
58_MUTEX_MIN_WAIT = 0.1
61class LockTimeoutError(TimeoutError):
62 """Raised when a lock cannot be acquired within the timeout."""
65class ResetRefusedError(RuntimeError):
66 """Raised when a reset cannot prove that no sync or import runs on the data root."""
69# In-process write mutex: serializes writers within the same process
70_write_mutex = threading.Lock()
73def _lock_path(lancedb_dir: Path | None) -> Path:
74 return (lancedb_dir if lancedb_dir is not None else cfg.lancedb_dir) / ".lock"
77def server_lock_path(data_dir: Path) -> Path:
78 """Path of the one-server-per-data-dir lock file."""
79 return data_dir / _SERVER_LOCK_NAME
82def acquire_server_lock(data_dir: Path, timeout: float = SERVER_LOCK_TIMEOUT) -> FileLock | None:
83 """Hold the one-server-per-data-dir lock, or None when a live server owns it.
85 The lock is an OS file lock, so the kernel releases it the moment its holder
86 exits, however it died; a crashed or killed server leaves no stale state.
87 """
88 data_dir.mkdir(parents=True, exist_ok=True)
89 lock = FileLock(server_lock_path(data_dir))
90 try:
91 lock.acquire(timeout=timeout)
92 except FileLockTimeout:
93 return None
94 return lock
97@dataclass(frozen=True)
98class ScopeOwner:
99 """The data dir the server holding a scope lock is serving, for the refusal message."""
101 data_dir: str
104@dataclass(frozen=True)
105class ScopeHold:
106 """A held scope lock plus its owner sidecar; release removes both."""
108 lock: FileLock
109 owner_path: Path
111 def release(self) -> None:
112 """Remove the owner sidecar, then free the scope for the next server."""
113 self.owner_path.unlink(missing_ok=True)
114 self.lock.release()
117def acquire_scope_lock(
118 scope_dir: Path, data_dir: Path, timeout: float = SERVER_LOCK_TIMEOUT
119) -> ScopeHold | None:
120 """Hold the one-server-per-scope lock, or None when a live server owns the scope.
122 The scope is a directory shared by several would-be servers (the Obsidian
123 plugin's shared root). Like the data-dir lock, the OS releases it the moment
124 the holder exits. The owner sidecar records which data dir the holder is
125 serving so a refused starter can name it in its message.
126 """
127 scope_dir.mkdir(parents=True, exist_ok=True)
128 lock = FileLock(scope_dir / _SCOPE_LOCK_NAME)
129 try:
130 lock.acquire(timeout=timeout)
131 except FileLockTimeout:
132 return None
133 owner_path = scope_dir / _SCOPE_OWNER_NAME
134 owner_path.write_text(json.dumps({"data_dir": str(data_dir)}), encoding="utf-8")
135 return ScopeHold(lock, owner_path)
138def read_scope_owner(scope_dir: Path) -> ScopeOwner | None:
139 """The scope's recorded owner, or None when absent or unreadable."""
140 try:
141 payload = json.loads((scope_dir / _SCOPE_OWNER_NAME).read_text(encoding="utf-8"))
142 return ScopeOwner(data_dir=str(payload["data_dir"]))
143 except (OSError, ValueError, KeyError, TypeError):
144 return None
147def _acquire_sync_lock(data_root: Path, *, write: bool) -> ReadWriteLock | None:
148 """Hold the data root's sync lock; None lets a sync run unmarked when it cannot lock."""
149 path = data_root / _SYNC_LOCK_NAME
150 try:
151 data_root.mkdir(parents=True, exist_ok=True)
152 lock = ReadWriteLock(path, is_singleton=False)
153 except _SYNC_LOCK_ERRORS as exc:
154 _sync_lock_unavailable(path, exc, write=write)
155 return None
156 try:
157 if write:
158 lock.acquire_write(blocking=False)
159 else:
160 lock.acquire_read()
161 except FileLockTimeout:
162 lock.close()
163 raise ResetRefusedError(_SYNC_RUNNING) from None
164 except _SYNC_LOCK_ERRORS as exc:
165 lock.close()
166 _sync_lock_unavailable(path, exc, write=write)
167 return None
168 return lock
171def _sync_lock_unavailable(path: Path, error: Exception, *, write: bool) -> None:
172 """Refuse a reset that cannot take the lock; warn once and let a sync run unmarked."""
173 if write:
174 raise ResetRefusedError(_SYNC_LOCK_UNKNOWN.format(path=path, error=error)) from error
175 log.warning(_SYNC_LOCK_REFUSED, path, error)
178def _release_sync_lock(lock: ReadWriteLock | None) -> None:
179 if lock is not None:
180 lock.release()
181 lock.close()
184@asynccontextmanager
185async def sync_running(data_root: Path) -> AsyncGenerator[None, None]:
186 """Mark a sync or import running on *data_root*, across processes; they share the mark.
188 The wait for a reset to finish runs in a worker thread, off the event loop.
189 """
190 lock = await asyncio.to_thread(_acquire_sync_lock, data_root, write=False)
191 try:
192 yield
193 finally:
194 _release_sync_lock(lock)
197@contextmanager
198def no_sync_running(data_root: Path) -> Generator[None, None, None]:
199 """Keep syncs off *data_root* for the block; raise ``ResetRefusedError`` unless it can."""
200 lock = _acquire_sync_lock(data_root, write=True)
201 try:
202 yield
203 finally:
204 _release_sync_lock(lock)
207@contextmanager
208def write_lock(
209 lancedb_dir: Path | None = None, timeout: float = LOCK_TIMEOUT
210) -> Generator[None, None, None]:
211 """Acquire the cross-process file lock then the in-process mutex.
213 The file lock lives next to the store's data, so cross-process writers
214 coordinate only when they lock the *same* directory: callers pass their
215 store's ``lancedb_dir`` (a per-instance ``Lilbee`` uses its own dir).
216 ``None`` falls back to the global ``cfg.lancedb_dir``.
218 The two stages share one budget: the time spent waiting on the file lock is
219 deducted before waiting on the mutex (plus a small ``_MUTEX_MIN_WAIT`` floor),
220 so a 30s request cannot stall for roughly twice that.
221 """
222 deadline = time.monotonic() + timeout
223 lock_path = _lock_path(lancedb_dir)
224 # The first write to a per-instance store can run before its data dir exists;
225 # the file lock cannot be created in a missing directory.
226 lock_path.parent.mkdir(parents=True, exist_ok=True)
227 flock = FileLock(lock_path)
228 try:
229 flock.acquire(timeout=timeout)
230 except FileLockTimeout:
231 raise LockTimeoutError("Timed out waiting for exclusive file lock") from None
232 try:
233 # Floor the mutex budget so a file lock that wins right at the deadline
234 # still gets a brief blocking attempt instead of a zero-timeout poll that
235 # spuriously fails when another thread holds the mutex for an instant.
236 remaining = max(_MUTEX_MIN_WAIT, deadline - time.monotonic())
237 acquired = _write_mutex.acquire(timeout=remaining)
238 if not acquired:
239 raise LockTimeoutError("Timed out waiting for write lock")
240 try:
241 yield
242 finally:
243 _write_mutex.release()
244 finally:
245 flock.release()