Coverage for /usr/local/lib/python3.10/site-packages/opal_server-0.0.0-py3.10.egg/opal_server/scopes/purge.py: 80%

114 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 11:54 +0000

1"""Fleet-wide purge of GitPolicyFetcher caches (PR3 of the leak series). 

2 

3Every worker subscribes ``handle_purge_message`` to SCOPES_PURGE_CHANNEL at 

4startup (see ``server.py``) and drops its in-memory cache entries for the 

5purged source. The leader additionally registers ``LeaderScopePurger.handle`` 

6(at watcher start), which sibling-checks and then authorizes that drop by 

7broadcasting the confirmation. 

8 

9WHO MUTATES THE CLONE TREE. Not "only the leader" — that claim was in this 

10docstring and was never true. The leader is the only mutator on SYNC paths 

11(clone/fetch are leader-only). A delete additionally removes the dir on the 

12worker that SERVED the DELETE, best-effort, exactly as master did 

13(``ScopesService._purge_local_clone_best_effort``). Its guards — ``lock_source`` 

14(an asyncio.Lock) and ``git_op_in_flight`` (a module-global set) — are 

15PROCESS-LOCAL, so they do not serialize that worker against the leader cloning 

16the same source in a sibling process on the same pod. Master had the identical 

17exposure with no in-flight guard at all, so this is not a regression — but note 

18the branch does not reclaim that orphan either: delivering the purge broadcast 

19drains MEMORY on every worker and touches no disk. It is also not an invariant. 

20Closing it cross-process is PER-15612. 

21 

22NOTE: the reconciliation sweep that reclaimed clone dirs referencing no live 

23scope was split out of this PR and is tracked as PER-15612. What remains here is 

24the purge driven by an actual scope delete/repoint. 

25 

26There is therefore NO reconciliation of any kind in this PR, and exactly ONE 

27path RECLAIMS a clone dir — removes it and leaves it removed: the best-effort 

28local floor ``ScopesService.delete_scope`` spawns on the worker serving the 

29DELETE. 

30 

31Two other ``rmtree`` sites exist and are not reclamation: ``git_fetcher``'s 

32invalid-repo recovery and ``_clone``'s partial-dir wipe both delete-and-replace 

33a dir they are about to re-create. Both run under ``lock_source`` and target 

34``self._repo_path``, derived from ``GitPolicyFetcher.source_id(source)`` — a 

35sha256 digest of the URL, never wire input. An audit of "can caller-controlled 

36input reach an rmtree?" has to account for all three. That is 

37master's behaviour (master removed it inline there, with no broadcast involved), 

38kept so a dropped broadcast does not regress against the merge base. 

39 

40What that leaves on disk, stated plainly because nothing else will reclaim it: 

41 

42- a DELETE's dir on every pod EXCEPT the serving one, when the broadcast is 

43 lost. Since 0.9.9-rc.3 every worker keeps a backbone reader for the whole 

44 process when SCOPES is on and the broadcaster is the reconnecting one 

45 (server.py holds the global listening context; the default), so a 

46 client-less non-leader is no longer deaf to the purge; the residual is 

47 BROADCAST_RECONNECT_ENABLED=false (legacy broadcaster, reader only while a 

48 client is connected) and a backbone outage, which still loses the message 

49 for everyone; 

50- a REPOINT's old dir on EVERY pod, always — there is no floor on that path; 

51- a dir whose source_id is unknowable because the prior record would not parse 

52 (see ``scopes/api.py``). 

53 

54All three are PER-15612. The memory purge is unaffected by them: it is 

55authorized by the leader's confirmation and is self-healing in both directions. 

56""" 

57import asyncio 

58import re 

59from pathlib import Path 

60from typing import Any, Optional 

61 

62from opal_common.logger import logger 

63from opal_common.schemas.policy_source import GitPolicyScopeSource 

64from opal_server.config import opal_server_config 

65from opal_server.git_fetcher import GitPolicyFetcher, git_op_in_flight 

66from pydantic import BaseModel, ValidationError 

67 

68# \Z (not $) so a trailing newline can't sneak past validation: in Python `$` 

69# also matches just before a final "\n", so "<64hex>-0\n" would wrongly pass. 

70# [0-9], not \d: for a `str` pattern \d is Unicode and also matches e.g. Arabic- 

71# Indic digits. Containment is unaffected either way (the derived path stays 

72# under base_dir and could not exist), but the id this validates is a sha256 

73# hex digest + an ASCII shard index, so say that. 

74_SOURCE_ID_RE = re.compile(r"\A[0-9a-f]{64}-[0-9]+\Z") 

75 

76 

77def confined_clone_path(base_dir, source_id: str): 

78 """Derive the on-disk clone dir for ``source_id``, or ``None`` if the id is 

79 malformed. 

80 

81 SECURITY: a purge command's ``clone_path`` field arrives over pub/sub and 

82 must NEVER reach the filesystem. ``source_id`` is a sha256 hex digest + 

83 shard index (no separators, no traversal), so the derived path is always 

84 confined to ``base_dir/git_sources``. The result is also the exact key used 

85 in ``GitPolicyFetcher.repos``. 

86 

87 What a forged message can still reach, now that the disk reclaim is cut: a 

88 ``free()`` of a cached pygit2 handle under an attacker-chosen key (via 

89 ``purge_local_memory`` -> ``forget_repo``). It can no longer reach an 

90 ``rmtree`` — no wire-driven path in this module removes a directory, and the 

91 one remaining removal (``ScopesService``'s delete floor) derives its target 

92 from the stored scope record, not from a message. Keeping the validation is 

93 still right: it is the only thing standing between a forged ``source_id`` 

94 and an arbitrary ``repos`` key, and it is what makes the confinement 

95 property survive PER-15612 putting an rmtree back on this path. 

96 """ 

97 if not _SOURCE_ID_RE.match(source_id): 97 ↛ 98line 97 didn't jump to line 98 because the condition on line 97 was never true

98 return None 

99 return str(GitPolicyFetcher.base_dir(Path(base_dir)) / source_id) 

100 

101 

102class ScopePurgeCommand(BaseModel): 

103 source_id: str # cache key for repos_last_fetched / repo_locks 

104 # Informational only: every handler re-derives the clone dir from source_id 

105 # (never trusts this path); kept for readable logs and forward-compat. 

106 clone_path: str 

107 scope_id: str # logging / tracing only 

108 # Load-bearing, NOT just logging: the leader's sibling-check fail-open 

109 # branches on reason. A repoint's old source still has a LIVE record (it 

110 # just moved), so an unreadable or unanswered scan can't be told apart from 

111 # "still shared" and the clone is kept; a delete's record is already gone, 

112 # so the same scan failure purges defensively — under-purging there is a 

113 # permanent leak. See LeaderScopePurger.purge_source_if_unshared. 

114 reason: str # "delete" | "repoint" 

115 confirmed: bool = False # set by the leader after the sibling-check; 

116 # memory handlers act only on confirmed commands 

117 

118 

119def purge_local_memory(source_id: str, clone_path: str) -> None: 

120 """Drop this process's in-memory cache entries for a source. 

121 

122 Never pops ``repo_locks``: lock_source's recheck loop protects waiters, 

123 not a current holder — popping is only safe while holding the lock, which 

124 this function does not (its callers that DO hold it pop there). ``forget_repo`` is skipped while a 

125 git op is in flight: freeing a pygit2 handle a pool thread still uses 

126 (e.g. a lingering timed-out fetch) is a crash risk. 

127 

128 A skipped free is NOT self-healing on this worker: ``forget_repo`` is 

129 otherwise only reached from the invalid-repo branch of a SYNC, and a purged 

130 source by construction has no live scope to sync. So a handle pinned by a 

131 lingering timed-out op stays cached for the life of the process. Nothing in 

132 this PR revisits it — that is PER-15612's job, along with the clone dir. 

133 """ 

134 if not git_op_in_flight(source_id): 134 ↛ 136line 134 didn't jump to line 136 because the condition on line 134 was always true

135 GitPolicyFetcher.forget_repo(clone_path) 

136 GitPolicyFetcher.repos_last_fetched.pop(source_id, None) 

137 # Unconditionally, unlike forget_repo: this holds no handle a lingering 

138 # pool thread could still be reading, so the in-flight guard does not apply. 

139 # Dropped for the same reason as repos_last_fetched — the source has no live 

140 # scope on this worker any more, so a kept entry is counted in the 

141 # sources_in_backoff gauge for the life of the process, and would suppress 

142 # the first sync of a scope later re-created against the same URL. 

143 GitPolicyFetcher.forget_source_backoff(source_id) 

144 

145 

146async def handle_purge_message(subscription, data: Any) -> None: 

147 """Every-worker subscriber for SCOPES_PURGE_CHANNEL.""" 

148 try: 

149 cmd = ScopePurgeCommand(**data) 

150 except (ValidationError, TypeError): 

151 logger.warning("Ignoring malformed scope purge message: {data}", data=data) 

152 return 

153 if not cmd.confirmed: 

154 # A request — only the leader acts on those (sibling-check first). 

155 return 

156 safe_path = confined_clone_path(opal_server_config.BASE_DIR, cmd.source_id) 

157 if safe_path is None: 157 ↛ 158line 157 didn't jump to line 158 because the condition on line 157 was never true

158 logger.warning( 

159 "Ignoring scope purge with malformed source_id: {sid}", 

160 sid=cmd.source_id, 

161 ) 

162 return 

163 logger.info( 

164 "Purging local caches for source {source_id} (scope {scope_id}, {reason})", 

165 source_id=cmd.source_id, 

166 scope_id=cmd.scope_id, 

167 reason=cmd.reason, 

168 ) 

169 # Under lock_source — a FRESH lock, not the one the leader holds while it 

170 # publishes the confirmation: the leader pops the repo_locks entry before 

171 # publishing precisely so this handler's setdefault mints a new one instead 

172 # of deadlocking on the held one (see purge_source_if_unshared). 

173 # forget_repo -> Repository.free() must not run concurrently with a 

174 # re-created scope's sync on THIS process: fetch_and_notify_on_changes holds 

175 # its handle across an await and then set_target()s it, all under 

176 # lock_source. The git_op_in_flight guard inside purge_local_memory only 

177 # covers pool-thread ops, not that event-loop handle-holding — so without 

178 # the lock this every-worker handler is the use-after-free the leader path 

179 # is careful to avoid, on every process except the publisher. 

180 async with GitPolicyFetcher.lock_source(cmd.source_id): 

181 purge_local_memory(cmd.source_id, safe_path) 

182 # Pop the repo_locks entry lock_source just minted (via setdefault), 

183 # under the lock — the same lock-identity rule the leader follows in 

184 # purge_source_if_unshared. purge_local_memory deliberately never pops it 

185 # (it held no lock); now that this handler does, popping here is what 

186 # keeps a purged source from leaving a stray repo_locks key (invariant 

187 # I4). lock_source waiters re-check the dict and re-mint a fresh lock, so 

188 # a concurrently re-created scope is unaffected. 

189 GitPolicyFetcher.repo_locks.pop(cmd.source_id, None) 

190 

191 

192async def subscribe_worker_purge_handler(endpoint) -> None: 

193 await endpoint.subscribe( 

194 [opal_server_config.SCOPES_PURGE_CHANNEL], handle_purge_message 

195 ) 

196 

197 

198def _scope_sharing_source( 

199 scopes_snapshot, source_id: str, excluded_scope_id: Optional[str] = None 

200) -> Optional[str]: 

201 """First live scope in a pre-fetched list mapping to ``source_id``, else 

202 None. 

203 

204 Pure (no I/O): RAISES if a scope's ``source_id()`` derivation raises — 

205 the caller owns the fail-open/fail-closed policy for that. 

206 """ 

207 return next( 

208 ( 

209 s.scope_id 

210 for s in scopes_snapshot 

211 if s.scope_id != excluded_scope_id 

212 and isinstance(s.policy, GitPolicyScopeSource) 

213 and GitPolicyFetcher.source_id(s.policy) == source_id 

214 ), 

215 None, 

216 ) 

217 

218 

219async def find_scope_sharing_source( 

220 scopes, source_id: str, excluded_scope_id: Optional[str] = None 

221) -> Optional[str]: 

222 """Return the id of a live scope mapping to ``source_id``, or None. 

223 

224 RAISES on a store/scan error (was: swallowed and returned None). The 

225 caller decides the fail-open policy by ``reason``: a repoint's old 

226 source still has a live record (just moved elsewhere), so a raising 

227 scan must NOT read as "unshared"; a delete's record is already gone, 

228 so under-purging there is a permanent leak. 

229 """ 

230 return _scope_sharing_source(await scopes.all(), source_id, excluded_scope_id) 

231 

232 

233class LeaderScopePurger: 

234 """Leader-only: removes clone dirs for purged sources. 

235 

236 Registered on SCOPES_PURGE_CHANNEL when leadership is acquired (the 

237 watcher task's start). It performs the sibling check no other worker can 

238 do, and authorizes the fleet's memory purge; it does not touch the clone 

239 tree (see the module docstring on who does). 

240 """ 

241 

242 def __init__(self, base_dir: Path, scopes, pubsub_endpoint): 

243 self._base_dir = base_dir 

244 self._scopes = scopes 

245 self._pubsub_endpoint = pubsub_endpoint 

246 # Strong refs to in-flight background purges (create_task results are 

247 # otherwise GC-able); discarded on completion. 

248 self._pending_purges = set() 

249 # Set by signal_stop(): no new purge is queued once shutdown started. 

250 self._stopping = False 

251 

252 async def _purge_and_log(self, cmd: ScopePurgeCommand) -> None: 

253 try: 

254 await self.purge_source_if_unshared(cmd) 

255 except Exception: 

256 # Detached background task: without this, an unexpected failure 

257 # (e.g. the confirmation publish hitting a broadcaster error) 

258 # surfaces only as asyncio's unretrieved-exception noise. 

259 logger.exception( 

260 f"Background purge of source {cmd.source_id} " f"({cmd.reason}) failed" 

261 ) 

262 

263 async def handle(self, subscription, data: Any): 

264 try: 

265 cmd = ScopePurgeCommand(**data) 

266 except (ValidationError, TypeError): 

267 # The worker-level handler already logged the malformed payload. 

268 return None 

269 if cmd.confirmed: 

270 return None # our own confirmation broadcast, addressed to workers 

271 if self._stopping: 271 ↛ 280line 271 didn't jump to line 280 because the condition on line 271 was never true

272 # Shutdown started: the watcher has already unsubscribed us and is 

273 # about to drain what is in flight. Queuing more work here would 

274 # either be abandoned by that bounded drain or publish a 

275 # confirmation nothing waits on. NOTHING recovers this: no periodic 

276 # reconciliation exists in this PR, and no later purge will name the 

277 # source (the record is already gone), so the fleet keeps its cache 

278 # entries for it until PER-15612 lands. The serving worker's local 

279 # floor has still reclaimed its own pod's clone dir. 

280 logger.info( 

281 f"Ignoring purge request for {cmd.source_id} ({cmd.reason}): " 

282 "purger is stopping" 

283 ) 

284 return None 

285 # publish() awaits subscriber callbacks inline — never do lock-waiting 

286 # disk work on the publisher's request path (DELETE/PUT latency is 

287 # bounded by contract). The purge proceeds in the background. 

288 return self._schedule(self._purge_and_log(cmd)) 

289 

290 def signal_stop(self) -> None: 

291 """Refuse new purge requests from ``handle`` (idempotent).""" 

292 self._stopping = True 

293 

294 async def stop(self) -> None: 

295 """Await in-flight background purges so a shutdown can't abandon a 

296 sibling check mid-flight (or start a fresh one after the watcher 

297 stopped). 

298 

299 These tasks are spawned detached in ``handle`` and are NOT in the 

300 watcher's ``self._tasks``, so ``BasePolicyWatcherTask.stop`` never waits 

301 on them. 

302 

303 This CAN block for a long time and the caller must bound it: a purge's 

304 first act is to take ``lock_source``, held by a sync across a whole 

305 clone/fetch (unbounded when ``SCOPES_GIT_FETCH_TIMEOUT`` is 0), and it 

306 then does a ``scopes.all()`` and a confirmation ``publish()`` against a 

307 Redis/broadcaster client with no socket timeout. The watcher calls this 

308 after cancelling its tasks (so the lock holders are gone) and under an 

309 ``asyncio.wait_for``. 

310 """ 

311 self.signal_stop() 

312 if self._pending_purges: 312 ↛ 313line 312 didn't jump to line 313 because the condition on line 312 was never true

313 await asyncio.gather(*list(self._pending_purges), return_exceptions=True) 

314 

315 def _schedule(self, coro) -> asyncio.Task: 

316 """Run ``coro`` detached, holding a strong ref so it can't be GC'd, and 

317 make ``stop()``'s drain wait for it.""" 

318 task = asyncio.create_task(coro) 

319 self._pending_purges.add(task) 

320 task.add_done_callback(self._pending_purges.discard) 

321 return task 

322 

323 async def purge_source_if_unshared(self, cmd: ScopePurgeCommand) -> None: 

324 """Leader-only: authorize the fleet-wide MEMORY purge for a source. 

325 

326 Only the leader can sibling-check (it reads the scope store), so only 

327 the leader may authorize a memory purge: acting on the raw request 

328 would drop cache entries a surviving sibling scope still uses. It 

329 publishes the confirmation, and every worker's ``handle_purge_message`` 

330 acts on that. 

331 

332 This deliberately does NOT touch the clone tree. Disk reclaim on delete 

333 happens on the DELETE-serving worker at master's semantics 

334 (``ScopesService._purge_local_clone_best_effort``); distributed disk 

335 reclaim — a leader-side rmtree, a retry for a dir a lingering git op 

336 pins, and the reconciliation sweep — is PER-15612, where the reclaim 

337 policy is agreed before implementation. Under-purging disk leaks 

338 forever and over-purging forces a re-clone, so it has no self-healing 

339 direction; a memory purge self-heals both ways (a wrongly-dropped 

340 handle just re-opens on next use), which is why the two split here. 

341 """ 

342 if confined_clone_path(self._base_dir, cmd.source_id) is None: 342 ↛ 343line 342 didn't jump to line 343 because the condition on line 342 was never true

343 logger.warning( 

344 f"Ignoring leader purge with malformed source_id: {cmd.source_id}" 

345 ) 

346 return 

347 confirm = False 

348 async with GitPolicyFetcher.lock_source(cmd.source_id): 

349 # I4 is "no stray repo_locks key for a source NOBODY holds", so the 

350 # drain below covers the exits where we abandon the source — the two 

351 # fail-open returns, and the confirmed purge — but NOT the path where 

352 # a live sibling still shares it. Popping there is safe (lock_source 

353 # re-mints) but wrong: it churns a lock the sibling is using, and the 

354 # bed asserts the entry survives a sibling delete 

355 # (test_shared_repo_survives_sibling_scope_delete). 

356 minted = GitPolicyFetcher.repo_locks.get(cmd.source_id) 

357 try: 

358 try: 

359 timeout = opal_server_config.SCOPES_STORE_READ_TIMEOUT 

360 check = find_scope_sharing_source(self._scopes, cmd.source_id) 

361 sharer = await ( 

362 asyncio.wait_for(check, timeout=timeout) 

363 if timeout > 0 

364 else check 

365 ) 

366 except asyncio.TimeoutError: 

367 # Bounded because this read is held under lock_source and the 

368 # Redis client has no socket timeout — an unreachable store 

369 # would otherwise wedge this source's lock for the life of 

370 # the process. 

371 # 

372 # Direction follows master's rule: over-purging self-heals, 

373 # under-purging is a permanent leak. A repoint's old source 

374 # may still be referenced, so an unanswered read must not be 

375 # read as "unshared"; a delete's record is already gone. 

376 if cmd.reason == "repoint": 

377 logger.warning( 

378 "Sibling check for {sid} timed out after {t}s on a " 

379 "repoint; not confirming — the fleet keeps its cache " 

380 "entries for this source until something names it " 

381 "again", 

382 sid=cmd.source_id, 

383 t=timeout, 

384 ) 

385 return 

386 logger.warning( 

387 "Sibling check for {sid} timed out after {t}s on {reason}; " 

388 "confirming defensively — its record is already gone, so " 

389 "withholding the purge would leak the fleet's cache " 

390 "entries permanently (a surviving sibling re-opens its " 

391 "handle on the next sync)", 

392 sid=cmd.source_id, 

393 t=timeout, 

394 reason=cmd.reason, 

395 ) 

396 sharer = None 

397 except Exception as e: 

398 if cmd.reason == "repoint": 

399 logger.warning( 

400 f"Sibling check for {cmd.source_id} failed on repoint; " 

401 f"not confirming: {e!r}" 

402 ) 

403 return 

404 logger.warning( 

405 f"Sibling check for {cmd.source_id} failed on {cmd.reason}; " 

406 f"confirming defensively: {e!r}" 

407 ) 

408 sharer = None 

409 if sharer is not None: 

410 logger.info( 

411 f"Scope {sharer} still shares source {cmd.source_id}, " 

412 "keeping the fleet's cache entries" 

413 ) 

414 minted = None # live source — leave its lock alone 

415 else: 

416 confirm = True 

417 # Published under the lock: publish() runs local subscribers 

418 # inline, so the confirmation frees this process's cached pygit2 

419 # handle. Releasing the lock first would let a re-created scope's 

420 # sync acquire it, cache a fresh handle, and enter 

421 # _notify_on_changes — which holds the handle across an await and 

422 # then calls set_target() on it — while this stale confirmation 

423 # frees it underneath (use-after-free). 

424 if confirm and self._pubsub_endpoint is not None: 

425 # LOAD-BEARING ORDER: pop BEFORE the publish. publish() runs 

426 # local subscribers inline on this task, and 

427 # handle_purge_message re-enters lock_source(source_id) — the 

428 # same non-reentrant Lock still held here. Popping first makes 

429 # its setdefault mint a FRESH lock instead of waiting on ours; 

430 # popping after would wedge this source's lock permanently. 

431 GitPolicyFetcher.repo_locks.pop(cmd.source_id, None) 

432 await self._pubsub_endpoint.publish( 

433 [opal_server_config.SCOPES_PURGE_CHANNEL], 

434 cmd.copy(update={"confirmed": True}).dict(), 

435 ) 

436 finally: 

437 # Lock-identity guarded: the pop-before-publish above hands the 

438 # dict entry off, and the awaited publish lets another coroutine 

439 # mint a SUCCESSOR lock. Popping unconditionally here would 

440 # discard that successor while its holder still runs, putting two 

441 # coroutines inside lock_source for the same source at once. 

442 # 

443 # The identity check is what prevents that, on its own: after the 

444 # hand-off `minted` is no longer the mapped lock, so `is minted` 

445 # is False whether a successor appeared or the key is simply 

446 # absent. An explicit `minted = None` used to sit on that path as 

447 # well; it was removed because it is unreachable-as-true and so 

448 # could not be tested — mutating it changed nothing, while 

449 # mutating this check fails three tests. 

450 if ( 

451 minted is not None 

452 and GitPolicyFetcher.repo_locks.get(cmd.source_id) is minted 

453 ): 

454 GitPolicyFetcher.repo_locks.pop(cmd.source_id, None)