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
« 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).
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.
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.
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.
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.
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.
40What that leaves on disk, stated plainly because nothing else will reclaim it:
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``).
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
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
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")
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.
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``.
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)
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
119def purge_local_memory(source_id: str, clone_path: str) -> None:
120 """Drop this process's in-memory cache entries for a source.
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.
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)
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)
192async def subscribe_worker_purge_handler(endpoint) -> None:
193 await endpoint.subscribe(
194 [opal_server_config.SCOPES_PURGE_CHANNEL], handle_purge_message
195 )
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.
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 )
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.
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)
233class LeaderScopePurger:
234 """Leader-only: removes clone dirs for purged sources.
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 """
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
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 )
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))
290 def signal_stop(self) -> None:
291 """Refuse new purge requests from ``handle`` (idempotent)."""
292 self._stopping = True
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).
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.
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)
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
323 async def purge_source_if_unshared(self, cmd: ScopePurgeCommand) -> None:
324 """Leader-only: authorize the fleet-wide MEMORY purge for a source.
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.
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)