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

1"""Cross-process locks: LanceDB write locking, the sync mark, and the server singleton. 

2 

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""" 

9 

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 

20 

21from filelock import FileLock, ReadWriteLock 

22from filelock import Timeout as FileLockTimeout 

23 

24from lilbee.core.config import cfg 

25 

26log = logging.getLogger(__name__) 

27 

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 

59 

60 

61class LockTimeoutError(TimeoutError): 

62 """Raised when a lock cannot be acquired within the timeout.""" 

63 

64 

65class ResetRefusedError(RuntimeError): 

66 """Raised when a reset cannot prove that no sync or import runs on the data root.""" 

67 

68 

69# In-process write mutex: serializes writers within the same process 

70_write_mutex = threading.Lock() 

71 

72 

73def _lock_path(lancedb_dir: Path | None) -> Path: 

74 return (lancedb_dir if lancedb_dir is not None else cfg.lancedb_dir) / ".lock" 

75 

76 

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 

80 

81 

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. 

84 

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 

95 

96 

97@dataclass(frozen=True) 

98class ScopeOwner: 

99 """The data dir the server holding a scope lock is serving, for the refusal message.""" 

100 

101 data_dir: str 

102 

103 

104@dataclass(frozen=True) 

105class ScopeHold: 

106 """A held scope lock plus its owner sidecar; release removes both.""" 

107 

108 lock: FileLock 

109 owner_path: Path 

110 

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() 

115 

116 

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. 

121 

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) 

136 

137 

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 

145 

146 

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 

169 

170 

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) 

176 

177 

178def _release_sync_lock(lock: ReadWriteLock | None) -> None: 

179 if lock is not None: 

180 lock.release() 

181 lock.close() 

182 

183 

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. 

187 

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) 

195 

196 

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) 

205 

206 

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. 

212 

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``. 

217 

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()