Coverage for src / lexigram / admin / services / collaborative.py: 41%

133 statements  

« prev     ^ index     » next       coverage.py v7.13.5, created at 2026-08-13 22:14 +0800

1"""Collaborative editing support for lexigram-admin. 

2 

3Provides: 

4 

5* :class:`EditLock` — per-record optimistic locking (prevents two users from 

6 editing the same record simultaneously) 

7* :class:`PresenceTracker` — tracks who is currently viewing or editing a 

8 resource record 

9* :class:`CollaborativeEditingService` — orchestrates locks + presence, 

10 emitting change notifications when a record is saved by another user 

11 

12This module integrates with the existing 

13:mod:`lexigram.admin.services.realtime` WebSocket layer to push live updates 

14to connected users. 

15 

16The edit lock store is provided via constructor injection as a 

17:class:`~lexigram.contracts.core.stores.LockStoreProtocol`. In production 

18this is backed by a persistent distributed store (e.g. SQL advisory locks or 

19Redis). In tests, pass a :class:`FakeLockStore` conforming to the protocol. 

20 

21Usage:: 

22 

23 from lexigram.admin.services.collaborative import CollaborativeEditingService 

24 

25 svc = CollaborativeEditingService(lock_store=lock_store) 

26 

27 # Acquire edit lock (raises LockConflictError if another user holds it) 

28 lock = await svc.acquire_lock("user", "record-id", user_id="alice") 

29 

30 # Update presence (call on page load / heartbeat) 

31 await svc.update_presence("user", "record-id", user_id="alice", action="editing") 

32 

33 # List who else is viewing 

34 viewers = svc.get_presence("user", "record-id") 

35 

36 # Notify all viewers that the record changed 

37 await svc.notify_change("user", "record-id", changed_by="alice") 

38 

39 # Release lock when done 

40 await svc.release_lock("user", "record-id", user_id="alice") 

41""" 

42 

43from __future__ import annotations 

44 

45from dataclasses import dataclass, field 

46from datetime import UTC, datetime, timedelta 

47from typing import TYPE_CHECKING 

48 

49from lexigram.contracts.core.stores import LockStoreProtocol 

50from lexigram.di.decorators import inject 

51from lexigram.logging import get_logger 

52 

53if TYPE_CHECKING: 

54 from lexigram.admin.services.realtime import RealtimeService 

55 

56logger = get_logger(__name__) 

57 

58_DEFAULT_LOCK_TTL_SECONDS = 300 # 5 minutes 

59 

60 

61@dataclass 

62class EditLock: 

63 """An exclusive edit lock on a resource record. 

64 

65 Attributes: 

66 resource_type: Resource name (e.g. ``"user"``). 

67 record_id: Record identifier. 

68 user_id: Holder of the lock. 

69 acquired_at: UTC timestamp when the lock was acquired. 

70 expires_at: UTC timestamp when the lock auto-expires. 

71 lock_id: Unique identifier for this lock instance. 

72 """ 

73 

74 resource_type: str 

75 record_id: str 

76 user_id: str 

77 acquired_at: datetime = field(default_factory=lambda: datetime.now(UTC)) 

78 expires_at: datetime = field( 

79 default_factory=lambda: ( 

80 datetime.now(UTC) + timedelta(seconds=_DEFAULT_LOCK_TTL_SECONDS) 

81 ) 

82 ) 

83 lock_id: str = field(default_factory=lambda: "") 

84 

85 def __post_init__(self) -> None: 

86 if not self.lock_id: 

87 import uuid 

88 

89 self.lock_id = str(uuid.uuid4())[:8] 

90 

91 @property 

92 def is_expired(self) -> bool: 

93 """Return ``True`` if the lock has passed its TTL.""" 

94 return datetime.now(UTC) >= self.expires_at 

95 

96 def refresh(self, ttl_seconds: int = _DEFAULT_LOCK_TTL_SECONDS) -> None: 

97 """Extend the lock expiry by *ttl_seconds* from now.""" 

98 self.expires_at = datetime.now(UTC) + timedelta(seconds=ttl_seconds) 

99 

100 

101@dataclass 

102class PresenceEntry: 

103 """A single user's presence on a resource record. 

104 

105 Attributes: 

106 user_id: User identifier. 

107 resource_type: Resource name. 

108 record_id: Record identifier. 

109 action: ``"viewing"`` or ``"editing"``. 

110 last_seen: UTC timestamp of last heartbeat. 

111 display_name: Optional human-readable name for the UI. 

112 """ 

113 

114 user_id: str 

115 resource_type: str 

116 record_id: str 

117 action: str = "viewing" 

118 last_seen: datetime = field(default_factory=lambda: datetime.now(UTC)) 

119 display_name: str = "" 

120 

121 @property 

122 def is_stale(self) -> bool: 

123 """Return ``True`` if no heartbeat in the last 60 seconds.""" 

124 return (datetime.now(UTC) - self.last_seen).total_seconds() > 60 

125 

126 

127class LockConflictError(RuntimeError): 

128 """Raised when trying to acquire a lock already held by another user.""" 

129 

130 _code: str = "LEX_ERR_ADMIN_026" 

131 

132 def __init__(self, holder_id: str, expires_at: datetime) -> None: 

133 self.holder_id = holder_id 

134 self.expires_at = expires_at 

135 super().__init__( 

136 f"Record is locked by '{holder_id}' until {expires_at.isoformat()}" 

137 ) 

138 

139 

140@inject 

141class CollaborativeEditingService: 

142 """Manages edit locks, presence, and change notifications. 

143 

144 Edit lock coordination is delegated to an injected 

145 :class:`~lexigram.contracts.core.stores.LockStoreProtocol`, allowing 

146 distributed deployments to use a shared persistent store (Redis, SQL 

147 advisory locks, etc.) without this service knowing the implementation. 

148 Rich lock metadata (:class:`EditLock`) is kept in-process as a local 

149 cache; the store is the authoritative source for *whether* a lock is held. 

150 

151 Args: 

152 lock_store: Injected distributed lock store. 

153 realtime_service: Optional 

154 :class:`~lexigram.admin.services.realtime.RealtimeDataService` 

155 instance. When provided, ``notify_change`` broadcasts over 

156 WebSocket to connected users. 

157 lock_ttl_seconds: How long (in seconds) an edit lock remains valid 

158 without being refreshed. 

159 """ 

160 

161 def __init__( 

162 self, 

163 lock_store: LockStoreProtocol, 

164 realtime_service: RealtimeService | None = None, 

165 lock_ttl_seconds: int = _DEFAULT_LOCK_TTL_SECONDS, 

166 ) -> None: 

167 self._lock_store = lock_store 

168 self._realtime = realtime_service 

169 self._lock_ttl = lock_ttl_seconds 

170 # Local metadata cache: key "type:id" → EditLock. 

171 # The distributed store is authoritative for coordination; this dict 

172 # stores the holder identity and expiry so we can surface rich errors. 

173 self._lock_metadata: dict[str, EditLock] = {} 

174 self._presence: dict[ 

175 str, dict[str, PresenceEntry] 

176 ] = {} # key: "type:id" → {user_id: entry} 

177 

178 # ------------------------------------------------------------------ 

179 # Internal helpers 

180 # ------------------------------------------------------------------ 

181 

182 @staticmethod 

183 def _record_key(resource_type: str, record_id: str) -> str: 

184 return f"{resource_type}:{record_id}" 

185 

186 # ------------------------------------------------------------------ 

187 # Locking 

188 # ------------------------------------------------------------------ 

189 

190 async def acquire_lock( 

191 self, 

192 resource_type: str, 

193 record_id: str, 

194 *, 

195 user_id: str, 

196 ) -> EditLock: 

197 """Acquire an exclusive edit lock on a record. 

198 

199 Args: 

200 resource_type: Resource name. 

201 record_id: Record identifier. 

202 user_id: User attempting to acquire the lock. 

203 

204 Returns: 

205 The acquired :class:`EditLock`. 

206 

207 Raises: 

208 LockConflictError: If another user holds a non-expired lock. 

209 """ 

210 key = self._record_key(resource_type, record_id) 

211 existing = self._lock_metadata.get(key) 

212 

213 if existing and not existing.is_expired: 

214 if existing.user_id == user_id: 

215 # Refresh own lock in the store and update local metadata. 

216 await self._lock_store.extend(key, user_id, self._lock_ttl) 

217 existing.refresh(self._lock_ttl) 

218 return existing 

219 raise LockConflictError(existing.user_id, existing.expires_at) 

220 

221 if existing and existing.is_expired: 

222 # Best-effort release of the stale entry from the store. 

223 await self._lock_store.release(key, existing.user_id) 

224 del self._lock_metadata[key] 

225 

226 acquired = await self._lock_store.acquire(key, user_id, self._lock_ttl) 

227 if not acquired: 

228 # Another process holds the lock; we don't have local metadata for it. 

229 raise LockConflictError( 

230 "unknown", 

231 datetime.now(UTC) + timedelta(seconds=self._lock_ttl), 

232 ) 

233 

234 lock = EditLock( 

235 resource_type=resource_type, 

236 record_id=str(record_id), 

237 user_id=user_id, 

238 expires_at=datetime.now(UTC) + timedelta(seconds=self._lock_ttl), 

239 ) 

240 self._lock_metadata[key] = lock 

241 logger.debug( 

242 "lock_acquired", 

243 user_id=user_id, 

244 resource_type=resource_type, 

245 record_id=record_id, 

246 ) 

247 return lock 

248 

249 async def release_lock( 

250 self, 

251 resource_type: str, 

252 record_id: str, 

253 *, 

254 user_id: str, 

255 ) -> bool: 

256 """Release a lock held by *user_id*. 

257 

258 Args: 

259 resource_type: Resource name. 

260 record_id: Record identifier. 

261 user_id: User releasing the lock. 

262 

263 Returns: 

264 ``True`` if released, ``False`` if lock not found or belongs to 

265 another user. 

266 """ 

267 key = self._record_key(resource_type, record_id) 

268 lock = self._lock_metadata.get(key) 

269 if lock is None: 

270 return False 

271 if lock.user_id != user_id: 

272 logger.warning( 

273 "lock_release_wrong_user", 

274 user_id=user_id, 

275 holder_id=lock.user_id, 

276 resource_type=resource_type, 

277 record_id=record_id, 

278 ) 

279 return False 

280 # Best-effort release from the distributed store; always clean local metadata. 

281 await self._lock_store.release(key, user_id) 

282 del self._lock_metadata[key] 

283 logger.debug( 

284 "lock_released", 

285 user_id=user_id, 

286 resource_type=resource_type, 

287 record_id=record_id, 

288 ) 

289 return True 

290 

291 def get_lock(self, resource_type: str, record_id: str) -> EditLock | None: 

292 """Return the active lock for a record, or ``None``. 

293 

294 Expired locks are evicted from local metadata on access. 

295 

296 Args: 

297 resource_type: Resource name. 

298 record_id: Record identifier. 

299 """ 

300 key = self._record_key(resource_type, record_id) 

301 lock = self._lock_metadata.get(key) 

302 if lock and lock.is_expired: 

303 del self._lock_metadata[key] 

304 return None 

305 return lock 

306 

307 def is_locked(self, resource_type: str, record_id: str) -> bool: 

308 """Return ``True`` if the record has an active lock.""" 

309 return self.get_lock(resource_type, record_id) is not None 

310 

311 def is_locked_by(self, resource_type: str, record_id: str, *, user_id: str) -> bool: 

312 """Return ``True`` if the record is locked by *user_id*.""" 

313 lock = self.get_lock(resource_type, record_id) 

314 return lock is not None and lock.user_id == user_id 

315 

316 # ------------------------------------------------------------------ 

317 # Presence 

318 # ------------------------------------------------------------------ 

319 

320 async def update_presence( 

321 self, 

322 resource_type: str, 

323 record_id: str, 

324 *, 

325 user_id: str, 

326 action: str = "viewing", 

327 display_name: str = "", 

328 ) -> PresenceEntry: 

329 """Record that *user_id* is currently on *resource_type/record_id*. 

330 

331 Call this on page load and periodically as a heartbeat (≤30s). 

332 

333 Args: 

334 resource_type: Resource name. 

335 record_id: Record identifier. 

336 user_id: User identifier. 

337 action: ``"viewing"`` or ``"editing"``. 

338 display_name: Optional human-readable name. 

339 

340 Returns: 

341 Updated :class:`PresenceEntry`. 

342 """ 

343 key = self._record_key(resource_type, record_id) 

344 user_map = self._presence.setdefault(key, {}) 

345 entry = user_map.get(user_id) 

346 if entry: 

347 entry.action = action 

348 entry.last_seen = datetime.now(UTC) 

349 if display_name: 

350 entry.display_name = display_name 

351 else: 

352 entry = PresenceEntry( 

353 user_id=user_id, 

354 resource_type=resource_type, 

355 record_id=str(record_id), 

356 action=action, 

357 display_name=display_name, 

358 ) 

359 user_map[user_id] = entry 

360 return entry 

361 

362 async def leave(self, resource_type: str, record_id: str, *, user_id: str) -> None: 

363 """Remove a user from the presence tracking for a record. 

364 

365 Also releases any lock held by the user. 

366 

367 Args: 

368 resource_type: Resource name. 

369 record_id: Record identifier. 

370 user_id: User who left. 

371 """ 

372 key = self._record_key(resource_type, record_id) 

373 user_map = self._presence.get(key, {}) 

374 user_map.pop(user_id, None) 

375 # Also release any lock 

376 await self.release_lock(resource_type, record_id, user_id=user_id) 

377 

378 def get_presence( 

379 self, 

380 resource_type: str, 

381 record_id: str, 

382 *, 

383 exclude_user: str | None = None, 

384 ) -> list[PresenceEntry]: 

385 """Return active (non-stale) presence entries for a record. 

386 

387 Args: 

388 resource_type: Resource name. 

389 record_id: Record identifier. 

390 exclude_user: Optionally exclude one user (e.g. self). 

391 

392 Returns: 

393 List of :class:`PresenceEntry` objects sorted by ``last_seen`` desc. 

394 """ 

395 key = self._record_key(resource_type, record_id) 

396 entries = [ 

397 e 

398 for e in self._presence.get(key, {}).values() 

399 if not e.is_stale and (exclude_user is None or e.user_id != exclude_user) 

400 ] 

401 return sorted(entries, key=lambda e: e.last_seen, reverse=True) 

402 

403 # ------------------------------------------------------------------ 

404 # Change notifications 

405 # ------------------------------------------------------------------ 

406 

407 async def notify_change( 

408 self, 

409 resource_type: str, 

410 record_id: str, 

411 *, 

412 changed_by: str, 

413 change_type: str = "update", 

414 ) -> int: 

415 """Notify all presence-tracked viewers that a record was changed. 

416 

417 Sends a JSON WebSocket message via the realtime service (when wired). 

418 Viewers can use this to trigger a soft refresh of the record. 

419 

420 Args: 

421 resource_type: Resource name. 

422 record_id: Record identifier. 

423 changed_by: User who made the change. 

424 change_type: ``"update"``, ``"delete"``, or ``"restore"``. 

425 

426 Returns: 

427 Number of users notified. 

428 """ 

429 viewers = self.get_presence(resource_type, record_id, exclude_user=changed_by) 

430 if not viewers or self._realtime is None: 

431 return len(viewers) 

432 

433 import lexigram.serialization as json 

434 

435 message = json.dumps( 

436 { 

437 "type": "record_changed", 

438 "resource_type": resource_type, 

439 "record_id": str(record_id), 

440 "changed_by": changed_by, 

441 "change_type": change_type, 

442 } 

443 ) 

444 

445 notified = 0 

446 for entry in viewers: 

447 try: 

448 await self._realtime.broadcast_to_user(entry.user_id, message) # type: ignore[attr-defined] 

449 notified += 1 

450 except Exception: # noqa: BLE001 

451 pass 

452 

453 return notified 

454 

455 

456__all__ = [ 

457 "CollaborativeEditingService", 

458 "EditLock", 

459 "LockConflictError", 

460 "PresenceEntry", 

461]