Coverage for /usr/local/lib/python3.10/site-packages/opal_server-0.0.0-py3.10.egg/opal_server/git_fetcher.py: 69%
517 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
1import asyncio
2import codecs
3import datetime
4import hashlib
5import inspect
6import math
7import os
8import shutil
9import threading
10import time
11import weakref
12from concurrent.futures import ThreadPoolExecutor
13from concurrent.futures import thread as cf_thread
14from contextlib import asynccontextmanager
15from dataclasses import dataclass
16from pathlib import Path
17from typing import Awaitable, Callable, Dict, Optional, cast
19import aiofiles.os
20import pygit2
21from ddtrace import tracer
22from git import Repo
23from opal_common.async_utils import run_sync
24from opal_common.git_utils.bundle_maker import BundleMaker
25from opal_common.http_utils import redact_url
26from opal_common.logger import logger
27from opal_common.monitoring import metrics
28from opal_common.schemas.policy import PolicyBundle
29from opal_common.schemas.policy_source import (
30 GitHubTokenAuthData,
31 GitPolicyScopeSource,
32 SSHAuthData,
33)
34from opal_common.synchronization.named_lock import NamedLock
35from opal_server.config import opal_server_config
36from pygit2 import (
37 KeypairFromMemory,
38 RemoteCallbacks,
39 Repository,
40 Username,
41 UserPass,
42 clone_repository,
43 discover_repository,
44 reference_is_valid_name,
45)
47# Source ids whose scope git op (clone/fetch) is still running on a pool thread
48# — including one that already exceeded its timeout but whose blocking pygit2
49# call has not yet returned. Guarded by a lock because it is cleared from the
50# pool thread (see ``run_in_git_executor``) and read/written from the event
51# loop. Used to guarantee at most one live git op per repository, since pygit2
52# ``Repository`` objects are not thread-safe.
53_git_busy: set = set()
54_git_busy_lock = threading.Lock()
57class GitConcurrencyLimitExceeded(RuntimeError):
58 """Raised when in-flight (live + zombie) git ops reach
59 SCOPES_GIT_MAX_ZOMBIES."""
62class CloneNotPopulatedError(ValueError):
63 """The remote-tracking namespace is EMPTY: no refs/remotes/<remote>/* at
64 all, so this clone has not been populated yet.
66 Distinct from BranchHeadNotFoundError, which means the namespace has refs
67 but not the configured one — a real misconfiguration. This is transient by
68 construction: _clone() rmtree's the destination and clones INTO the final
69 path, so for the whole duration of a recovery re-clone the dir exists with
70 no remote refs.
72 Deliberately derived from DISK, not from the in-flight marker: that marker
73 is a per-process module global, written only by the leader's sync, while
74 GET /scopes/{id}/policy is served by any worker. Keying the 503/409 split
75 on it made every NON-leader worker answer 409 "not retryable" throughout a
76 recovery — the exact inversion the split exists to prevent, on N-1 of N
77 workers. Disk truth is identical on every worker.
79 Subclasses ValueError so broad handlers still catch it.
81 ``waited_seconds`` is how long a request was held waiting for this clone
82 before the error was surfaced, and ``client_disconnected`` says the caller
83 hung up while it was held, so the answer about to be shaped goes nowhere.
84 Both are declared here with defaults rather than read with a getattr
85 default at each handler: a raise path that forgets to set one then shows
86 up as an explicit default in one place, instead of being
87 indistinguishable from "did not wait" at every reader.
88 """
90 waited_seconds: float = 0.0
91 client_disconnected: bool = False
94class BranchHeadNotFoundError(ValueError):
95 """Configured branch has no resolvable HEAD (permanent misconfig), NOT a
96 transient clone gap.
98 Subclasses ValueError so broad handlers still catch it.
99 """
102_zombie_cap_logged = False
105class _DaemonThreadPoolExecutor(ThreadPoolExecutor):
106 """A ``ThreadPoolExecutor`` whose worker threads are daemon threads.
108 A scope git op can stay blocked in a libgit2 network call well past our
109 soft timeout. With the stdlib's non-daemon workers, ``concurrent.futures``'
110 atexit handler would ``join()`` such a thread and hang interpreter shutdown
111 (the pinned libgit2 enforces no network read timeout, so a black-holed
112 remote never unblocks it). Daemon workers let the process exit
113 promptly; abandoning an in-flight fetch at exit is safe (libgit2 stages
114 objects in a temp pack and swaps refs atomically under lockfiles, and a
115 half-written clone dir is detected as invalid and re-cloned on next boot).
117 Only thread creation is customised, mirroring CPython's
118 ``_adjust_thread_count``. If a future CPython changes the internals we rely
119 on, we fall back to the stdlib (non-daemon) behaviour.
120 """
122 def _adjust_thread_count(self) -> None: # pragma: no cover - thread mgmt
123 worker = getattr(cf_thread, "_worker", None)
124 # Fall back to the stdlib if any internal we mirror has moved or changed
125 # shape: _worker must exist and take exactly the 4 positional args we pass,
126 # _threads_queues must exist, and this executor must expose _initializer.
127 if ( 127 ↛ 157line 127 didn't jump to line 157 because the condition on line 127 was never true
128 worker is None
129 or not hasattr(cf_thread, "_threads_queues")
130 or not hasattr(self, "_initializer")
131 # _idle_semaphore is used below and is just as private as the rest.
132 # A CPython that drops it would rewrite its own _adjust_thread_count
133 # accordingly, so the stdlib fallback would still work — but ours
134 # would raise AttributeError out of submit(), failing every scope
135 # git op. Guarding it is what routes that case to the fallback.
136 # (Not demonstrable by deleting the attribute at runtime: that
137 # leaves the stdlib's method still using it, a state CPython can
138 # never actually be in.)
139 # Every private this method goes on to touch, not just the ones
140 # whose absence seemed likely. The argument above for
141 # _idle_semaphore applies verbatim to each: a CPython that drops
142 # one would rewrite its own _adjust_thread_count accordingly, so
143 # the stdlib fallback still works while ours raises AttributeError
144 # out of submit() and fails every scope git op.
145 or not all(
146 hasattr(self, name)
147 for name in (
148 "_idle_semaphore",
149 "_initargs",
150 "_max_workers",
151 "_thread_name_prefix",
152 "_threads",
153 "_work_queue",
154 )
155 )
156 ):
157 return super()._adjust_thread_count()
158 try:
159 if len(inspect.signature(worker).parameters) != 4: 159 ↛ 160line 159 didn't jump to line 160 because the condition on line 159 was never true
160 return super()._adjust_thread_count()
161 except (TypeError, ValueError):
162 return super()._adjust_thread_count()
163 # If idle threads are available, don't spin up new ones.
164 if self._idle_semaphore.acquire(timeout=0): 164 ↛ 165line 164 didn't jump to line 165 because the condition on line 164 was never true
165 return
167 def weakref_cb(_, q=self._work_queue):
168 q.put(None)
170 num_threads = len(self._threads)
171 if num_threads < self._max_workers: 171 ↛ exitline 171 didn't return from function '_adjust_thread_count' because the condition on line 171 was always true
172 thread_name = "%s_%d" % (self._thread_name_prefix or self, num_threads)
173 t = threading.Thread(
174 name=thread_name,
175 target=cf_thread._worker,
176 args=(
177 weakref.ref(self, weakref_cb),
178 self._work_queue,
179 self._initializer,
180 self._initargs,
181 ),
182 daemon=True,
183 )
184 t.start()
185 self._threads.add(t)
186 # Deliberately NOT registered in ``cf_thread._threads_queues``:
187 # the stdlib's ``_python_exit`` atexit handler iterates that global
188 # and ``join()``s every thread in it regardless of ``daemon=True``,
189 # which would block interpreter shutdown on a lingering (timed-out)
190 # git call — the exact "stuck on an offline repo" hang this class
191 # exists to avoid, relocated to shutdown/restart. Normal shutdown
192 # uses ``self._threads`` + queue sentinels and is unaffected.
195def shutdown_git_executor() -> None:
196 """Drop loop-bound live-op accounting before fork.
198 Only the loop-bound semaphores are cleared (meaningless post-fork).
199 _git_busy markers are LEFT in place so reset_caches (next) can skip
200 freeing a handle a lingering git op still holds; the forked child
201 clears the stale markers in _reset_git_executor_after_fork.
202 """
203 _live_ops_semaphores.clear()
206def _reset_git_executor_after_fork() -> None:
207 """after_in_child fork handler: _git_busy_lock is held on entry (the paired
208 'before' handler acquired it and the child inherits it LOCKED). Reinit it in
209 place FIRST (dropping it without a matching acquire — re-acquiring would
210 deadlock), then mutate _git_busy directly (child is single-threaded here)."""
211 global _git_busy_lock
212 reinit = getattr(_git_busy_lock, "_at_fork_reinit", None)
213 if callable(reinit):
214 reinit()
215 else: # pragma: no cover
216 _git_busy_lock = threading.Lock()
217 _live_ops_semaphores.clear()
218 _git_busy.clear()
221if hasattr(os, "register_at_fork"): 221 ↛ 229line 221 didn't jump to line 229 because the condition on line 221 was always true
222 os.register_at_fork(
223 before=_git_busy_lock.acquire,
224 after_in_parent=_git_busy_lock.release,
225 after_in_child=_reset_git_executor_after_fork,
226 )
229def _emit_git_ops_in_flight(count: int) -> None:
230 # Continuous gauge so Datadog can watch the in-flight (incl. timed-out
231 # zombie) git-op count rise and fall — not just the one-shot error log at
232 # the SCOPES_GIT_MAX_ZOMBIES cap. datadog.statsd is fail-silent and
233 # thread-safe, so this is safe to call from the git-op daemon threads even
234 # when metrics are unconfigured.
235 # Tagged by pid: every worker in the gunicorn pool emits this same series,
236 # so untagged it is last-write-wins per flush and reads as one arbitrary
237 # worker's count rather than anything about the pod.
238 metrics.gauge(
239 "opal_server.scopes.git_ops_in_flight",
240 count,
241 tags={"pid": str(os.getpid())},
242 )
245def _mark_git_op_started(key: str) -> None:
246 with _git_busy_lock:
247 _git_busy.add(key)
248 count = len(_git_busy)
249 _emit_git_ops_in_flight(count)
252def _mark_git_op_done(key: str) -> None:
253 global _zombie_cap_logged
254 with _git_busy_lock:
255 _git_busy.discard(key)
256 count = len(_git_busy)
257 if ( 257 ↛ 261line 257 didn't jump to line 261 because the condition on line 257 was never true
258 _zombie_cap_logged
259 and len(_git_busy) < opal_server_config.SCOPES_GIT_MAX_ZOMBIES
260 ):
261 _zombie_cap_logged = False
262 _emit_git_ops_in_flight(count)
265def git_op_in_flight(key: str) -> bool:
266 """True while a git op for ``key`` is still running on a pool thread.
268 Stays True during the "lingering" window after a timeout, until the
269 blocking pygit2 call actually returns.
270 """
271 with _git_busy_lock:
272 return key in _git_busy
275def drain_git_ops(timeout: float) -> bool:
276 """Block up to `timeout`s for all in-flight git ops to finish.
278 Returns True if drained, False if the timeout elapsed with ops still
279 lingering. Lets a clone/fetch that finished at the end of the sync
280 pass clear its marker before reset_caches runs; ops still lingering
281 (unreachable remote) are left running and protected by
282 reset_caches's in-flight guard. timeout<=0 = don't wait.
283 """
284 deadline = time.monotonic() + max(0.0, timeout)
285 while True:
286 with _git_busy_lock:
287 if not _git_busy:
288 return True
289 if time.monotonic() >= deadline:
290 with _git_busy_lock:
291 return not _git_busy
292 time.sleep(0.05)
295def git_busy_count() -> int:
296 """Number of scope git ops holding a pool thread (incl.
298 timed-out zombies).
299 """
300 with _git_busy_lock:
301 return len(_git_busy)
304@dataclass
305class SourceBackoff:
306 """How long a repeatedly-failing source is skipped by the periodic pass.
308 ``next_attempt_at`` is a ``time.monotonic()`` reading, not wall clock: the
309 schedule must survive an NTP step or a container clock jump, and it is only
310 ever compared against another monotonic reading in this process.
311 """
313 consecutive_failures: int
314 next_attempt_at: float
315 last_error: str
318# The exponent is clamped here rather than left to grow with the failure count.
319# `2.0 ** (n-1)` raises OverflowError once the exponent passes ~1023 — from
320# inside the except clause that is handling the git failure. With a 10s base,
321# 2**64 * 10s is ~5.8e12 years: the clamp changes no reachable outcome, it
322# only keeps a very old dead source from raising instead of being skipped.
323_MAX_BACKOFF_DOUBLINGS = 64
325# Past this delay a source is, for practical purposes, abandoned until an
326# explicit refresh/PUT or a process restart — worth one WARNING when crossed.
327_BACKOFF_ABANDONED_SECONDS = 24 * 3600.0
330def _finite_positive_or_zero(value) -> float:
331 """Read a config number as a positive finite float, else 0.0.
333 Confi parses the environment at import, so a non-numeric value fails the
334 process at startup and never reaches this; what this covers is a value
335 assigned to the config object at runtime, plus `nan`/`inf`, which parse
336 cleanly, start the process, and are not durations anyone meant to set.
337 """
338 try:
339 f = float(value)
340 except (TypeError, ValueError):
341 return 0.0
342 if not math.isfinite(f) or f <= 0:
343 return 0.0
344 return f
347def _backoff_base_seconds() -> float:
348 """SCOPES_GIT_BACKOFF_BASE_SECONDS, validated. 0.0 means "disabled".
350 The first delay after a source's first failure; each further
351 consecutive failure doubles it. Deliberately short (10s by default):
352 a delay shorter than the gap to the next pass simply does not skip
353 that pass, so the first few doublings cost one attempt per pass
354 exactly as before, and the schedule bites from roughly the fourth
355 consecutive failure — minutes, then hours, then days. Duplicates of
356 a source in the SAME pass are collapsed regardless of the delay by
357 the re-check under lock_source.
358 """
359 return _finite_positive_or_zero(opal_server_config.SCOPES_GIT_BACKOFF_BASE_SECONDS)
362def _backoff_max_seconds() -> float:
363 """SCOPES_GIT_BACKOFF_MAX_SECONDS, validated. 0.0 means "no cap".
365 Uncapped by default on purpose: a repository that has been unreachable
366 for a day is, in all likelihood, dead — check it again in two days, then
367 four, and before long "at the next restart or explicit refresh". A cap
368 is available for operators who would rather bound the staleness of a
369 repository that comes back on its own without anyone touching the scope.
370 """
371 return _finite_positive_or_zero(opal_server_config.SCOPES_GIT_BACKOFF_MAX_SECONDS)
374def _backoff_delay(n: int) -> float:
375 """The delay armed after the n-th consecutive failure (n >= 1)."""
376 base = _backoff_base_seconds()
377 raw = base * 2.0 ** min(n - 1, _MAX_BACKOFF_DOUBLINGS)
378 cap = _backoff_max_seconds()
379 if cap > 0: 379 ↛ 383line 379 didn't jump to line 383 because the condition on line 379 was never true
380 # A cap below the base would make the feature silently inert (every
381 # delay shorter than one pass, nothing ever skipped): the base is the
382 # floor, so a low cap means "one pass at a time", never "off".
383 raw = min(raw, max(cap, base))
384 return raw
387def _emit_sources_in_backoff() -> None:
388 # Gauge of how many sources the periodic pass is currently
389 # skipping — the one number that says "this pod is not syncing N of your
390 # repos" without reading logs. Tagged by pid for the same reason as
391 # _emit_git_ops_in_flight: every worker emits this series, so untagged it
392 # is last-write-wins per flush. Never tagged by scope or source — that
393 # would make the cardinality proportional to the customer count.
394 # LIVE entries only: an entry whose delay has expired is kept so the
395 # consecutive-failure count survives until the next attempt, but the pass
396 # is not skipping it any more, and the gauge answers "how many sources is
397 # this pod not syncing right now". Emitted on every transition AND once
398 # per pass (sync_scopes), because DogStatsD gauges report nothing between
399 # sends and the steady state this feature creates has few transitions.
400 now = time.monotonic()
401 # With the kill switch on nothing is skipped regardless of the entries
402 # still recorded, so the gauge must read 0 — otherwise the dashboard says
403 # "N sources in backoff" every pass while the feature is off.
404 if _backoff_base_seconds() <= 0: 404 ↛ 405line 404 didn't jump to line 405 because the condition on line 404 was never true
405 live = 0
406 else:
407 live = sum(
408 1
409 for e in GitPolicyFetcher.source_backoff.values()
410 if e.next_attempt_at > now
411 )
412 metrics.gauge(
413 "opal_server.scopes.sources_in_backoff",
414 live,
415 tags={"pid": str(os.getpid())},
416 )
419# Public name for callers outside this module (the per-pass emission in
420# scopes/service.py); the underscore-prefixed one stays for in-module use.
421emit_sources_in_backoff = _emit_sources_in_backoff
424def _consume_future_result(fut) -> None:
425 # A future left running after its awaiter timed out is never awaited again;
426 # retrieve its outcome so asyncio doesn't log "exception never retrieved".
427 if not fut.cancelled(): 427 ↛ exitline 427 didn't return from function '_consume_future_result' because the condition on line 427 was always true
428 try:
429 fut.exception()
430 except Exception:
431 pass
434# Bounds LIVE (non-timed-out) git ops. asyncio primitives are loop-bound, so
435# the semaphore is minted per running loop (WeakKeyDictionary: a dead loop's
436# entry vanishes with it). A timed-out op releases its slot while its zombie
437# thread lingers — capacity is never consumed by zombies (a fixed pool
438# starves once zombies exceed its size; see the offline-repo bed gate).
439_live_ops_semaphores: "weakref.WeakKeyDictionary" = weakref.WeakKeyDictionary()
442def _get_live_ops_semaphore() -> asyncio.Semaphore:
443 loop = asyncio.get_running_loop()
444 sem = _live_ops_semaphores.get(loop)
445 if sem is None:
446 sem = asyncio.Semaphore(max(1, opal_server_config.SCOPES_GIT_MAX_WORKERS))
447 _live_ops_semaphores[loop] = sem
448 return sem
451async def run_in_git_executor(func, *args, timeout: float, busy_key=None, **kwargs):
452 """Run a blocking git call on its own daemon thread with a hard timeout.
454 ``SCOPES_GIT_MAX_WORKERS`` bounds LIVE (non-timed-out) ops via an asyncio
455 semaphore; each op still gets its own single-use daemon-thread executor,
456 so a lingering ("zombie") op after a timeout never occupies a shared pool
457 slot — it keeps running on its private thread but no longer counts
458 against the concurrency bound.
460 Raises the builtin ``TimeoutError`` when the call exceeds ``timeout``
461 seconds (``timeout <= 0`` means no limit). NOTE: the timeout unblocks the
462 event loop and the awaiting coroutine, but the underlying pygit2 call keeps
463 running on its own daemon thread. Nothing forces it to stop — the pinned
464 libgit2 sets no socket/server read timeout — so against a black-holed remote
465 that thread (and this key's in-flight marker) can stay alive for the life of
466 the process. ``SCOPES_GIT_MAX_ZOMBIES`` is the real bound on how many such
467 threads accumulate.
469 When ``busy_key`` is given it is marked in-flight for the *entire real
470 duration* of the call — including any lingering time after a timeout — and
471 cleared only when the blocking call actually returns (on its own thread).
472 Callers use ``git_op_in_flight`` to avoid starting a second git op against
473 the same repository while a timed-out one is still running.
474 """
475 global _zombie_cap_logged
476 # Clamped like every sibling knob in this subsystem: unclamped, a negative
477 # value is truthy AND `count >= cap` holds with nothing in flight, so the
478 # very first git op would be refused and no scope would ever sync. Negative
479 # reads as "no cap" (0), matching the intent of anyone typing -1 to disable.
480 max_zombies = max(0, opal_server_config.SCOPES_GIT_MAX_ZOMBIES)
481 if max_zombies and git_busy_count() >= max_zombies: 481 ↛ 482line 481 didn't jump to line 482 because the condition on line 481 was never true
482 if not _zombie_cap_logged:
483 _zombie_cap_logged = True
484 logger.error(
485 "Refusing new scope git op: {count} in-flight at/over "
486 "SCOPES_GIT_MAX_ZOMBIES={cap}; remotes appear stuck.",
487 count=git_busy_count(),
488 cap=max_zombies,
489 )
490 # Counted on EVERY refusal, unlike the log above which latches once per
491 # episode: the log answers "did we hit the cap", the counter answers
492 # "how hard and for how long" — the part an operator needs mid-outage.
493 metrics.increment(
494 "opal_server.scopes.git_ops_refused", tags={"pid": str(os.getpid())}
495 )
496 raise GitConcurrencyLimitExceeded(
497 f"in-flight git ops ({git_busy_count()}) reached "
498 f"SCOPES_GIT_MAX_ZOMBIES ({max_zombies})"
499 )
501 loop = asyncio.get_running_loop()
503 def _runner():
504 try:
505 return func(*args, **kwargs)
506 finally:
507 if busy_key is not None:
508 _mark_git_op_done(busy_key)
510 sem = _get_live_ops_semaphore()
511 await sem.acquire()
512 released = False
514 def _release_once():
515 nonlocal released
516 if not released: 516 ↛ exitline 516 didn't return from function '_release_once' because the condition on line 516 was always true
517 released = True
518 sem.release()
520 try:
521 # Single-use executor: the op gets a private daemon thread, so a zombie
522 # never blocks the next op the way a fixed shared pool does. shutdown
523 # with wait=False just drops bookkeeping; the daemon thread dies with
524 # the pygit2 call (or the process).
525 executor = _DaemonThreadPoolExecutor(
526 max_workers=1, thread_name_prefix="opal-git"
527 )
528 if busy_key is not None: 528 ↛ 530line 528 didn't jump to line 530 because the condition on line 528 was always true
529 _mark_git_op_started(busy_key)
530 try:
531 fut = loop.run_in_executor(executor, _runner)
532 except BaseException:
533 if busy_key is not None:
534 _mark_git_op_done(busy_key)
535 executor.shutdown(wait=False)
536 raise
537 fut.add_done_callback(lambda f: executor.shutdown(wait=False))
539 if not (timeout and timeout > 0): 539 ↛ 540line 539 didn't jump to line 540 because the condition on line 539 was never true
540 return await fut
542 # asyncio.wait (not wait_for) so a timeout does NOT cancel the future:
543 # the thread runs to completion and clears busy_key; the done-callback
544 # retrieves the eventual result to avoid "exception never retrieved".
545 fut.add_done_callback(_consume_future_result)
546 done, _pending = await asyncio.wait({fut}, timeout=timeout)
547 if not done: 547 ↛ 550line 547 didn't jump to line 550 because the condition on line 547 was never true
548 # Zombie: free the capacity slot; the private daemon thread lingers
549 # until the OS gives up, tracked only by busy_key.
550 raise TimeoutError(f"git operation exceeded {timeout}s")
551 return fut.result()
552 finally:
553 _release_once()
556class PolicyFetcherCallbacks:
557 async def on_update(self, old_head: Optional[str], head: str):
558 pass
561class PolicyFetcher:
562 def __init__(self, callbacks):
563 self.callbacks = callbacks
565 def fetch(self, hinted_hash: Optional[str] = None):
566 raise NotImplementedError()
569class RepoInterface:
570 """Manages a git repo with pygit2."""
572 @staticmethod
573 def create_local_branch_ref(
574 repo: Repository,
575 branch_name: str,
576 remote_name: str,
577 base_branch: str,
578 ) -> pygit2.Reference:
579 if branch_name not in repo.branches.local: 579 ↛ 590line 579 didn't jump to line 590 because the condition on line 579 was always true
580 base_remote_branch = f"{remote_name}/{base_branch}"
581 if repo.branches.remote.get(base_remote_branch) is not None: 581 ↛ 584line 581 didn't jump to line 584 because the condition on line 581 was always true
582 (commit, _) = repo.resolve_refish(base_remote_branch)
583 else:
584 raise RuntimeError("Base branch was not found on remote")
585 logger.debug(
586 f"Created local branch '{branch_name}', pointing to: {commit.hex}"
587 )
588 return repo.create_reference(f"refs/heads/{branch_name}", commit.hex)
589 else:
590 logger.debug(
591 f"No need to create local branch '{branch_name}': already exists!"
592 )
593 return repo.references[f"refs/heads/{branch_name}"]
595 @staticmethod
596 def has_remote_branch(repo: Repository, branch: str, remote: str) -> bool:
597 try:
598 repo.lookup_reference(f"refs/remotes/{remote}/{branch}")
599 return True
600 except KeyError:
601 return False
603 @staticmethod
604 def get_local_branch(repo: Repository, branch: str) -> Optional[pygit2.Reference]:
605 try:
606 return repo.lookup_reference(f"refs/heads/{branch}")
607 except KeyError:
608 return None
610 @staticmethod
611 def get_commit_hash(repo: Repository, branch: str, remote: str) -> Optional[str]:
612 try:
613 (commit, _) = repo.resolve_refish(f"{remote}/{branch}")
614 return commit.hex
615 except (pygit2.GitError, KeyError):
616 return None
618 @staticmethod
619 def verify_found_repo_matches_remote(
620 repo: Repository,
621 expected_remote_url: str,
622 ) -> Repository:
623 """Verifies that the repo we found in the directory matches the repo we
624 are wishing to clone."""
625 for remote in repo.remotes: 625 ↛ 631line 625 didn't jump to line 631 because the loop on line 625 didn't complete
626 if remote.url == expected_remote_url: 626 ↛ 625line 626 didn't jump to line 625 because the condition on line 626 was always true
627 logger.debug(
628 f"found target repo url is referred by remote: {remote.name}, url={redact_url(remote.url)}"
629 )
630 return
631 error: str = f"Repo mismatch! No remote matches target url: {redact_url(expected_remote_url)}, found urls: {[redact_url(remote.url) for remote in repo.remotes]}"
632 logger.error(error)
633 raise ValueError(error)
636class GitPolicyFetcher(PolicyFetcher):
637 repo_locks = {}
638 repos = {}
639 repos_last_fetched = {}
640 # source_id -> how long the periodic pass keeps skipping this source after
641 # consecutive clone/fetch failures. Per process and in memory only.
642 #
643 # Mutated ONLY on the event loop: the awaited outcome of a git op is what
644 # counts, so a daemon thread that finally returns long after its awaiter
645 # timed out never touches this (that late result is unobserved by
646 # construction — see run_in_git_executor). No lock is therefore needed, and
647 # the read in fetch_and_notify_on_changes deliberately happens before
648 # lock_source so a skipped source costs nothing.
649 source_backoff: Dict[str, SourceBackoff] = {}
651 def __init__(
652 self,
653 base_dir: Path,
654 scope_id: str,
655 source: GitPolicyScopeSource,
656 callbacks=PolicyFetcherCallbacks(),
657 remote_name: str = "origin",
658 liveness_probe: Optional[Callable[[], Awaitable[bool]]] = None,
659 ):
660 super().__init__(callbacks)
661 self._base_dir = GitPolicyFetcher.base_dir(base_dir)
662 self._source = source
663 self._source_id = GitPolicyFetcher.source_id(self._source)
664 self._auth_callbacks = GitCallback(self._source)
665 self._repo_path = self._base_dir / self._source_id
666 self._remote = remote_name
667 self._scope_id = scope_id
668 self._liveness_probe = liveness_probe
669 logger.debug(
670 f"Initializing git fetcher: scope_id={scope_id}, url={redact_url(source.url)}, branch={self._source.branch}, source_id={self._source_id}"
671 )
673 @staticmethod
674 @asynccontextmanager
675 async def lock_source(source_id: str):
676 """Serialize all mutation of a source's clone dir and cached handles.
678 Locks are minted on demand into ``repo_locks`` (asyncio.Lock: process-
679 local but fair, unlike the previous file-based lock). A scope delete
680 pops the dict entry while holding the lock, so after acquiring we must
681 re-check that ``repo_locks`` still maps ``source_id`` to the lock we
682 acquired — a waiter woken after a delete would otherwise proceed under
683 the stale lock, unserialized against holders of the freshly-minted one.
684 """
685 while True:
686 lock = GitPolicyFetcher.repo_locks.setdefault(source_id, asyncio.Lock())
687 async with lock:
688 if GitPolicyFetcher.repo_locks.get(source_id) is lock:
689 yield
690 return
692 def _backoff_entry(self) -> Optional[SourceBackoff]:
693 """This source's live backoff entry, or None if it may be attempted.
695 Returns None while the feature is disabled even when an entry exists:
696 an operator who sets SCOPES_GIT_BACKOFF_BASE_SECONDS=0 during an
697 incident must get the old behaviour back on the next pass, not have to
698 wait out the delays already recorded.
699 """
700 if _backoff_base_seconds() <= 0: 700 ↛ 701line 700 didn't jump to line 701 because the condition on line 700 was never true
701 return None
702 entry = GitPolicyFetcher.source_backoff.get(self._source_id)
703 if entry is None or time.monotonic() >= entry.next_attempt_at: 703 ↛ 705line 703 didn't jump to line 705 because the condition on line 703 was always true
704 return None
705 return entry
707 def _record_source_failure(self, err: BaseException) -> None:
708 """Count one failed clone/fetch against this source and re-arm the
709 delay.
711 Called only for failures that say something about the REMOTE (a
712 GitError or a timeout). Backpressure from the global zombie cap is
713 deliberately not recorded: at that ceiling every scope is refused,
714 healthy ones included, so recording it would put the whole fleet into
715 backoff because of one bad repo.
716 """
717 if _backoff_base_seconds() <= 0: 717 ↛ 718line 717 didn't jump to line 718 because the condition on line 717 was never true
718 return # kill switch: record nothing, so nothing is ever skipped
719 previous = GitPolicyFetcher.source_backoff.get(self._source_id)
720 n = (previous.consecutive_failures if previous is not None else 0) + 1
721 delay = _backoff_delay(n)
722 GitPolicyFetcher.source_backoff[self._source_id] = SourceBackoff(
723 consecutive_failures=n,
724 next_attempt_at=time.monotonic() + delay,
725 last_error=repr(err),
726 )
727 # WARNING only when something changes for the operator: the source
728 # ENTERS backoff, its delay first exceeds a day (from here on it is
729 # effectively abandoned until a restart or an explicit refresh), or —
730 # if a cap is configured — its delay first reaches the cap. Every other
731 # recorded failure is DEBUG: the timer's own attempts already get rarer
732 # as the delay grows, but an explicit refresh that keeps failing
733 # (policy-sync re-issues them constantly for a broken repo) bypasses
734 # the backoff and would otherwise WARN on every call, on top of the
735 # ERROR the failing op already logged.
736 previous_delay = (
737 _backoff_delay(previous.consecutive_failures)
738 if previous is not None
739 else None
740 )
741 entering = previous is None
742 crossed_abandoned = delay >= _BACKOFF_ABANDONED_SECONDS and (
743 previous_delay is None or previous_delay < _BACKOFF_ABANDONED_SECONDS
744 )
745 cap = _backoff_max_seconds()
746 reached_cap = (
747 cap > 0
748 and delay >= max(cap, _backoff_base_seconds())
749 and (
750 previous_delay is None
751 or previous_delay < max(cap, _backoff_base_seconds())
752 )
753 )
754 log = (
755 logger.warning
756 if (entering or crossed_abandoned or reached_cap)
757 else logger.debug
758 )
759 log(
760 "Backing off {url} for {delay:.0f}s after {n} consecutive "
761 "failures: {err}",
762 url=redact_url(self._source.url),
763 delay=delay,
764 n=n,
765 err=repr(err),
766 )
767 _emit_sources_in_backoff()
769 def _clear_source_backoff(self) -> None:
770 """A git op against this source succeeded — drop its failure
771 history."""
772 previous = GitPolicyFetcher.source_backoff.pop(self._source_id, None)
773 if previous is None: 773 ↛ 775line 773 didn't jump to line 775 because the condition on line 773 was always true
774 return
775 logger.info(
776 "Source {url} recovered after {n} failures",
777 url=redact_url(self._source.url),
778 n=previous.consecutive_failures,
779 )
780 _emit_sources_in_backoff()
782 @staticmethod
783 def forget_source_backoff(source_id: str) -> None:
784 """Drop a source's backoff entry when the source itself goes away.
786 Called from the purge paths (delete/repoint), never from ``forget_repo``
787 — that one is keyed by clone PATH and is also reached mid-sync from the
788 invalid-repo recovery branch, where the source is very much still ours
789 and its failure history must survive.
790 """
791 if GitPolicyFetcher.source_backoff.pop(source_id, None) is not None:
792 _emit_sources_in_backoff()
794 async def _was_fetched_after(self, t: datetime.datetime):
795 last_fetched = GitPolicyFetcher.repos_last_fetched.get(self._source_id, None)
796 if last_fetched is None:
797 return False
798 return last_fetched > t
800 async def fetch_and_notify_on_changes(
801 self,
802 hinted_hash: Optional[str] = None,
803 force_fetch: bool = False,
804 req_time: datetime.datetime = None,
805 *,
806 honor_backoff: bool = False,
807 ):
808 """Makes sure the repo is already fetched and is up to date.
810 - if no repo is found, the repo will be cloned.
811 - if the repo is found and it is deemed out-of-date, the configured remote will be fetched.
812 - if after a fetch new commits are detected, a callback will be triggered.
813 - if the hinted commit hash is provided and is already found in the local clone
814 we use this hint to avoid an necessary fetch.
816 ``honor_backoff`` says this call is pass-originated (the periodic sync
817 and the boot preload) and may be skipped while the source is serving
818 out a failure backoff. It defaults to False so every explicit path —
819 POST /scopes/{id}/refresh, POST /scopes/refresh, PUT /scopes — attempts
820 the source immediately: those are someone asking for this repo, now,
821 and the most likely reason they are asking is that they just fixed it.
822 """
823 # Checked before lock_source on purpose (and again under it, below).
824 # A hung source holds that lock for the whole clone, so a check ONLY
825 # inside it would make every skipped duplicate queue behind the very
826 # operation the skip exists to avoid; and a skipped source must
827 # consume no git-executor slot, so it can never be refused by (or
828 # contribute to) the SCOPES_GIT_MAX_ZOMBIES cap.
829 if honor_backoff:
830 entry = self._backoff_entry()
831 if entry is not None: 831 ↛ 832line 831 didn't jump to line 832 because the condition on line 831 was never true
832 metrics.increment(
833 "opal_server.scopes.git_op_skipped",
834 tags={"reason": "backoff"},
835 )
836 # DEBUG, not INFO: this fires once per backed-off source per
837 # pass, on every pass, for as long as the repo stays broken.
838 logger.debug(
839 "Skipping sync for {url}: in backoff for another {left:.0f}s "
840 "after {n} consecutive failures ({err})",
841 url=redact_url(self._source.url),
842 left=max(0.0, entry.next_attempt_at - time.monotonic()),
843 n=entry.consecutive_failures,
844 err=entry.last_error,
845 )
846 return
847 async with GitPolicyFetcher.lock_source(self._source_id):
848 # Re-checked under the lock: N pass-originated syncs of one source
849 # that arrive together — phase 2 runs the duplicates concurrently,
850 # and phase 1 may have recorded nothing for it (refused at the
851 # zombie cap, scope gone, no fetch needed) — all pass the cheap
852 # pre-lock check before the first has failed and recorded, then
853 # serialise here; without this second look each would perform its
854 # own full clone attempt against the dead remote.
855 if honor_backoff:
856 entry = self._backoff_entry()
857 if entry is not None: 857 ↛ 858line 857 didn't jump to line 858 because the condition on line 857 was never true
858 metrics.increment(
859 "opal_server.scopes.git_op_skipped",
860 tags={"reason": "backoff"},
861 )
862 logger.debug(
863 "Skipping sync for {url}: in backoff for another "
864 "{left:.0f}s after {n} consecutive failures ({err})",
865 url=redact_url(self._source.url),
866 left=max(0.0, entry.next_attempt_at - time.monotonic()),
867 n=entry.consecutive_failures,
868 err=entry.last_error,
869 )
870 return
871 if git_op_in_flight(self._source_id): 871 ↛ 876line 871 didn't jump to line 876 because the condition on line 871 was never true
872 # A previous git op for this repo exceeded its timeout and is
873 # still running on a pool thread. pygit2 Repository objects are
874 # not thread-safe, so skip this cycle rather than touch the same
875 # repo concurrently; the next cycle retries once it finishes.
876 logger.warning(
877 "Skipping sync for {url}: a previous git operation is still "
878 "running after its timeout.",
879 url=redact_url(self._source.url),
880 )
881 return
882 with tracer.trace(
883 "git_policy_fetcher.fetch_and_notify_on_changes",
884 resource=self._scope_id,
885 ):
886 if self._discover_repository(self._repo_path):
887 logger.debug("Repo found at {path}", path=self._repo_path)
888 # The probe opens/parses a fresh Repository handle from
889 # disk — off the event loop so a slow disk can't stall
890 # every other request being served on this worker.
891 repo = await run_sync(self._get_valid_repo)
892 if repo is not None: 892 ↛ 968line 892 didn't jump to line 968 because the condition on line 892 was always true
893 should_fetch = await self._should_fetch(
894 repo,
895 hinted_hash=hinted_hash,
896 force_fetch=force_fetch,
897 req_time=req_time,
898 )
899 if should_fetch:
900 logger.debug(
901 f"Fetching remote (force_fetch={force_fetch}): {self._remote} ({redact_url(self._source.url)})"
902 )
903 # Record the START time but write it only on
904 # success: a failed fetch must not look "fresh"
905 # to _was_fetched_after(), or it suppresses the
906 # forced refresh a webhook just asked for. The
907 # start time (not completion) is what req_time
908 # comparisons need: a fetch that STARTED after
909 # the request already satisfies it.
910 fetch_started = datetime.datetime.now()
911 try:
912 await run_in_git_executor(
913 repo.remotes[self._remote].fetch,
914 callbacks=self._auth_callbacks,
915 timeout=opal_server_config.SCOPES_GIT_FETCH_TIMEOUT,
916 busy_key=self._source_id,
917 )
918 except TimeoutError as exc:
919 # Expected when a repo is unreachable: log cleanly
920 # (no traceback) and skip, matching the clone path.
921 # repos_last_fetched stays stale so the next cycle
922 # retries and force_fetch is not wrongly suppressed.
923 metrics.increment(
924 "opal_server.scopes.git_op_failures",
925 tags={"op": "fetch", "reason": "timeout"},
926 )
927 self._record_source_failure(exc)
928 logger.error(
929 "Timed out fetching {url}, skipping: {err}",
930 url=redact_url(self._source.url),
931 err=repr(exc),
932 )
933 return
934 except pygit2.GitError as exc:
935 # The fast-fail half of the same problem, on a
936 # source that already has a local copy: revoked
937 # credentials or a deleted remote fail here in
938 # ~1s, every pass, forever. Counted as well as
939 # backed off — git_op_failures previously
940 # covered only the CLONE side of git_error, so
941 # a fleet whose fetches were all failing read
942 # as zero failures on the dashboard.
943 metrics.increment(
944 "opal_server.scopes.git_op_failures",
945 tags={"op": "fetch", "reason": "git_error"},
946 )
947 self._record_source_failure(exc)
948 # Re-raised, not swallowed: sync_scope's
949 # per-scope handler logs it with a traceback
950 # today, and that is left exactly as it was.
951 raise
952 GitPolicyFetcher.repos_last_fetched[
953 self._source_id
954 ] = fetch_started
955 self._clear_source_backoff()
956 logger.debug(
957 f"Fetch completed: {redact_url(self._source.url)}"
958 )
960 # New commits might be present because of a previous fetch made by another scope
961 await self._notify_on_changes(repo)
962 return
963 else:
964 # repo dir exists but invalid -> drop the cached handle
965 # FIRST (it is the thing judging the dir invalid; kept,
966 # it would re-invalidate the fresh clone on every sync
967 # -> infinite re-clone loop), then delete the directory.
968 logger.warning(
969 "Deleting invalid repo: {path}", path=self._repo_path
970 )
971 GitPolicyFetcher.forget_repo(str(self._repo_path))
972 try:
973 await run_sync(shutil.rmtree, str(self._repo_path))
974 except FileNotFoundError:
975 pass # already gone — the intended end state
976 except OSError as e:
977 # A partial dir left by an abandoned (timed-out)
978 # clone may still be written to; a failed delete
979 # self-heals via the clone below (or next cycle).
980 logger.warning(
981 f"Failed to remove clone dir "
982 f"{self._repo_path}: {e!r}"
983 )
984 else:
985 logger.info("Repo not found at {path}", path=self._repo_path)
987 # fallthrough to clean clone
988 # Liveness check before clone (the resurrection point): a
989 # DELETE that landed during this sync already broadcast its
990 # purge; cloning now would resurrect the dead scope's repo
991 # and re-populate the caches. Runs under lock_source, so it
992 # is serialized against any purge on THIS process. Fails open:
993 # a store hiccup must not block the sync.
994 if self._liveness_probe is not None: 994 ↛ 1011line 994 didn't jump to line 1011 because the condition on line 994 was always true
995 try:
996 alive = await self._liveness_probe()
997 except Exception as e:
998 logger.warning(
999 "Liveness probe for scope {scope} failed, "
1000 "proceeding with clone: {err}",
1001 scope=self._scope_id,
1002 err=repr(e),
1003 )
1004 alive = True
1005 if not alive: 1005 ↛ 1006line 1005 didn't jump to line 1006 because the condition on line 1005 was never true
1006 logger.info(
1007 "Scope {scope} was deleted mid-sync, skipping clone",
1008 scope=self._scope_id,
1009 )
1010 return
1011 await self._clone()
1013 def _discover_repository(self, path: Path) -> bool:
1014 git_path: Path = path / ".git"
1015 return discover_repository(str(path)) and git_path.exists()
1017 async def _clone(self):
1018 if self._repo_path.exists(): 1018 ↛ 1022line 1018 didn't jump to line 1022 because the condition on line 1018 was never true
1019 # A failed/interrupted clone leaves a partial dir;
1020 # clone_repository refuses a non-empty destination, which would
1021 # wedge every retry for this source.
1022 try:
1023 await run_sync(shutil.rmtree, str(self._repo_path))
1024 except FileNotFoundError:
1025 pass # already gone — the intended end state
1026 except OSError as e:
1027 logger.warning(f"Failed to remove clone dir {self._repo_path}: {e!r}")
1028 logger.info(
1029 "Cloning repo at '{url}' to '{path}'",
1030 url=redact_url(self._source.url),
1031 path=self._repo_path,
1032 )
1033 # Same start-time rule as the fetch path above: the clone's
1034 # negotiation reflects remote state at clone START, so that is the
1035 # timestamp req_time comparisons need.
1036 clone_started = datetime.datetime.now()
1037 try:
1038 repo: Repository = await run_in_git_executor(
1039 clone_repository,
1040 self._source.url,
1041 str(self._repo_path),
1042 callbacks=self._auth_callbacks,
1043 timeout=opal_server_config.SCOPES_GIT_FETCH_TIMEOUT,
1044 busy_key=self._source_id,
1045 )
1046 except (pygit2.GitError, TimeoutError) as exc:
1047 metrics.increment(
1048 "opal_server.scopes.git_op_failures",
1049 tags={
1050 "op": "clone",
1051 # The distinction the log line cannot carry: a steady rate of
1052 # timeouts against known-unreachable repos is expected after
1053 # SCOPES_GIT_FETCH_TIMEOUT landed; a git_error is not.
1054 "reason": "timeout"
1055 if isinstance(exc, TimeoutError)
1056 else "git_error",
1057 },
1058 )
1059 self._record_source_failure(exc)
1060 logger.error(
1061 "Could not clone repo at {url}: {err}",
1062 url=redact_url(self._source.url),
1063 err=repr(exc),
1064 )
1065 else:
1066 logger.info(f"Clone completed: {redact_url(self._source.url)}")
1067 # Cleared on the awaited SUCCESS of the git op itself, before the
1068 # local bookkeeping below: the remote is demonstrably reachable, so
1069 # a later failure inside _notify_on_changes (a corrupt object
1070 # store, say) must not leave the source marked as unreachable.
1071 self._clear_source_backoff()
1072 # Cache the fresh handle so the next sync's _get_repo() reuses it
1073 # instead of reopening (or hitting a stale predecessor).
1074 GitPolicyFetcher.repos[str(self._repo_path)] = repo
1075 # A reclone just downloaded current remote state — record it so
1076 # _was_fetched_after() doesn't force a redundant fetch next cycle.
1077 GitPolicyFetcher.repos_last_fetched[self._source_id] = clone_started
1078 await self._notify_on_changes(repo)
1080 def _get_repo(self) -> Repository:
1081 path = str(self._repo_path)
1082 if path not in GitPolicyFetcher.repos:
1083 GitPolicyFetcher.repos[path] = Repository(path)
1084 return GitPolicyFetcher.repos[path]
1086 def _get_valid_repo(self) -> Optional[Repository]:
1087 try:
1088 repo = self._get_repo()
1089 RepoInterface.verify_found_repo_matches_remote(repo, self._source.url)
1090 # A clone can be discoverable yet unusable: refs and config
1091 # intact but the object store gutted (crash mid-gc, disk
1092 # corruption). A fetch then negotiates "up to date" against the
1093 # intact refs and downloads nothing, so without this check the
1094 # scope serves 500s forever with no self-heal. Validate that the
1095 # tracked branch's head object is actually readable FROM DISK:
1096 # the check must use a short-lived fresh handle, because the
1097 # cached warm handle keeps deleted pack files readable through
1098 # its open mmaps (unlink does not invalidate them) and would
1099 # report the object as present. Partial corruption deeper in
1100 # the tree is NOT caught here (that would need fsck-grade
1101 # checks).
1102 probe = Repository(str(self._repo_path))
1103 try:
1104 try:
1105 ref = probe.lookup_reference(
1106 f"refs/remotes/{self._remote}/{self._source.branch}"
1107 )
1108 except KeyError:
1109 # Branch not fetched yet — the fetch path handles that.
1110 return repo
1111 if probe.get(ref.target) is None: 1111 ↛ 1112line 1111 didn't jump to line 1112 because the condition on line 1111 was never true
1112 logger.warning(
1113 "Repo at {path} has refs but an unreadable object "
1114 "store (missing head object) — treating as invalid",
1115 path=self._repo_path,
1116 )
1117 return None
1118 return repo
1119 finally:
1120 probe.free()
1121 except pygit2.GitError:
1122 logger.warning("Invalid repo at: {path}", path=self._repo_path)
1123 return None
1125 async def _should_fetch(
1126 self,
1127 repo: Repository,
1128 hinted_hash: Optional[str] = None,
1129 force_fetch: bool = False,
1130 req_time: datetime.datetime = None,
1131 ) -> bool:
1132 if force_fetch:
1133 if req_time is not None and await self._was_fetched_after(req_time): 1133 ↛ 1134line 1133 didn't jump to line 1134 because the condition on line 1133 was never true
1134 logger.info(
1135 "Repo was fetched after refresh request, override force_fetch with False"
1136 )
1137 else:
1138 return True # must fetch
1140 if not RepoInterface.has_remote_branch(repo, self._source.branch, self._remote):
1141 logger.info(
1142 "Target branch was not found in local clone, re-fetching the remote"
1143 )
1144 return True # missing branch
1146 if hinted_hash is not None:
1147 try:
1148 _ = repo.revparse_single(hinted_hash)
1149 return False # hinted commit was found, no need to fetch
1150 except KeyError:
1151 logger.info(
1152 "Hinted commit hash was not found in local clone, re-fetching the remote"
1153 )
1154 return True # hinted commit was not found
1156 # by default, we try to avoid re-fetching the repo for performance
1157 return False
1159 @property
1160 def local_branch_name(self) -> str:
1161 # Use the scope id as local branch name, so different scopes could track the same remote branch separately
1162 branch_name_unescaped = f"scopes/{self._scope_id}"
1163 if reference_is_valid_name(branch_name_unescaped):
1164 return branch_name_unescaped
1166 # if scope id can't be used as a gitref (e.g invalid chars), use its hex representation
1167 return f"scopes/{self._scope_id.encode().hex()}"
1169 async def _notify_on_changes(self, repo: Repository):
1170 # Get the latest commit hash of the target branch
1171 new_revision = RepoInterface.get_commit_hash(
1172 repo, self._source.branch, self._remote
1173 )
1174 if new_revision is None:
1175 logger.error(f"Did not find target branch on remote: {self._source.branch}")
1176 return
1178 # Get the previous commit hash of the target branch
1179 local_branch = RepoInterface.get_local_branch(repo, self.local_branch_name)
1180 if local_branch is None:
1181 # First sync of a new branch (the first synced branch in this repo was set by the clone (see `checkout_branch`))
1182 old_revision = None
1183 local_branch = RepoInterface.create_local_branch_ref(
1184 repo, self.local_branch_name, self._remote, self._source.branch
1185 )
1186 else:
1187 old_revision = local_branch.target.hex
1189 await self.callbacks.on_update(old_revision, new_revision)
1191 # Bring forward local branch (a bit like "pull"), so we won't detect changes again
1192 local_branch.set_target(new_revision)
1194 def _get_current_branch_head(self) -> str:
1195 # Opened fresh per call instead of using the shared cached handle:
1196 # this runs on executor threads (run_sync(make_bundle) in the policy-
1197 # bundle route) and outside lock_source, where the cached handle can
1198 # be free()'d concurrently by a scope delete or invalid-repo recovery.
1199 # asyncio locks don't exclude executor threads — sharing the handle
1200 # here is a use-after-free. Same fresh-probe pattern as
1201 # _get_valid_repo's disk-truth check.
1202 repo = Repository(str(self._repo_path))
1203 try:
1204 # Resolve the branch ref inline rather than via
1205 # RepoInterface.get_commit_hash, which collapses BOTH failure modes
1206 # to None. On the serving path we must tell them apart:
1207 # * KeyError -> the branch ref genuinely does not exist: a
1208 # PERMANENT misconfiguration (wrong/deleted branch). Surface it
1209 # as BranchHeadNotFoundError -> the bundle route's 409
1210 # "not retryable".
1211 # * pygit2.GitError -> the ref is present but its object can't be
1212 # resolved right now (object store transiently gutted by a
1213 # concurrent re-clone/fetch). This is TRANSIENT and the sync
1214 # path is already self-healing it, so let it propagate: the
1215 # bundle route's own pygit2.GitError handler turns it into a
1216 # retryable 503, instead of telling the client "not retryable"
1217 # for a scope that will recover on its own.
1218 try:
1219 commit, _ = repo.resolve_refish(f"{self._remote}/{self._source.branch}")
1220 head_commit_hash = commit.hex
1221 except KeyError:
1222 # Split the KeyError by DISK STATE, not by the in-flight marker:
1223 # an empty remote-tracking namespace means the clone has not been
1224 # populated yet (transient), whereas siblings present but ours
1225 # missing means the configured branch is wrong (permanent).
1226 prefix = f"refs/remotes/{self._remote}/"
1227 if not any(ref.startswith(prefix) for ref in repo.listall_references()):
1228 raise CloneNotPopulatedError(
1229 f"No {prefix}* refs yet at {self._repo_path}"
1230 )
1231 head_commit_hash = None
1232 finally:
1233 free = getattr(repo, "free", None)
1234 if callable(free):
1235 free()
1236 if not head_commit_hash:
1237 logger.error("Could not find current branch head")
1238 raise BranchHeadNotFoundError("Could not find current branch head")
1239 return head_commit_hash
1241 @tracer.wrap("git_policy_fetcher.make_bundle")
1242 def make_bundle(self, base_hash: Optional[str] = None) -> PolicyBundle:
1243 repo = Repo(str(self._repo_path))
1244 bundle_maker = BundleMaker(
1245 repo,
1246 {Path(p) for p in self._source.directories},
1247 extensions=self._source.extensions,
1248 root_manifest_path=self._source.manifest,
1249 bundle_ignore=self._source.bundle_ignore,
1250 )
1251 current_head_commit = repo.commit(self._get_current_branch_head())
1253 if not base_hash:
1254 return bundle_maker.make_bundle(current_head_commit)
1255 else:
1256 try:
1257 base_commit = repo.commit(base_hash)
1258 return bundle_maker.make_diff_bundle(base_commit, current_head_commit)
1259 except ValueError:
1260 return bundle_maker.make_bundle(current_head_commit)
1262 @staticmethod
1263 def source_id(source: GitPolicyScopeSource) -> str:
1264 base = hashlib.sha256(source.url.encode("utf-8")).hexdigest()
1265 index = (
1266 hashlib.sha256(source.branch.encode("utf-8")).digest()[0]
1267 % opal_server_config.SCOPES_REPO_CLONES_SHARDS
1268 )
1269 return f"{base}-{index}"
1271 @staticmethod
1272 def base_dir(base_dir: Path) -> Path:
1273 return base_dir / "git_sources"
1275 @staticmethod
1276 def repo_clone_path(base_dir: Path, source: GitPolicyScopeSource) -> Path:
1277 return GitPolicyFetcher.base_dir(base_dir) / GitPolicyFetcher.source_id(source)
1279 @staticmethod
1280 def forget_repo(path: str) -> None:
1281 """Drop the cached repository for a clone path and release its handles.
1283 The cached ``pygit2.Repository`` keeps OS file descriptors and mmapped
1284 pack indexes open; without this, a deleted scope's repo pins memory and
1285 inodes for the lifetime of the process even after the clone is removed.
1286 ``Repository.free()`` is called only when available (the pinned pygit2
1287 always has it; the guard defends against test doubles and future API
1288 changes); otherwise the dropped reference is reclaimed by GC.
1289 """
1290 repo = GitPolicyFetcher.repos.pop(path, None)
1291 if repo is None:
1292 return
1293 free = getattr(repo, "free", None)
1294 if callable(free): 1294 ↛ exitline 1294 didn't return from function 'forget_repo' because the condition on line 1294 was always true
1295 try:
1296 free()
1297 except Exception as e:
1298 logger.warning(
1299 f"pygit2 Repository.free() failed for {path}: {e!r}; "
1300 "relying on GC to release the handles"
1301 )
1303 @staticmethod
1304 def reset_caches() -> None:
1305 """Free and drop every cached repo handle, lock, and timestamp.
1307 Called in the gunicorn master after preload and before fork so
1308 no fetcher state is inherited by workers. A forked worker that
1309 inherited a handle for a scope it never syncs (sync is leader-
1310 only) could never purge it — the fleet-wide purge broadcast only
1311 reaches workers whose broadcaster reader is running — so it
1312 would pin that handle for life. Workers re-open handles lazily
1313 from the on-disk clones (preserved). Inherited repo_locks are
1314 asyncio.Locks bound to the master's event loop and meaningless
1315 post-fork regardless.
1317 A source whose git op is still in flight (lingering past its
1318 timeout on a daemon thread) is skipped: its handle is only
1319 dropped from the cache, never free()'d, since the pool thread
1320 may still be reading from it — free()'ing it here would be a
1321 use-after-free. GC reclaims it once the blocking call actually
1322 returns.
1324 ``source_backoff`` is deliberately NOT cleared. It holds no handles,
1325 no fds and no loop-bound objects, so none of the reasons above apply —
1326 and the preload this runs after is exactly where a dead repo's clone
1327 failures are discovered. Letting the forked leader inherit them is
1328 the point when a periodic pass follows: otherwise the leader starts
1329 by re-hammering the same unreachable repos, which is the boot storm the
1330 backoff exists to collapse. (When no periodic pass follows,
1331 ScopesPolicyWatcherTask.start() drops the inherited entries itself, so
1332 the one boot sync still attempts every source once.)
1333 ``_reset_git_executor_after_fork`` leaves it alone for the same
1334 reason.
1335 """
1336 for path in list(GitPolicyFetcher.repos):
1337 source_id = os.path.basename(path.rstrip("/"))
1338 if git_op_in_flight(source_id):
1339 # Still in use on a pool thread — drop the reference, never free
1340 # (free()'ing a handle a daemon thread holds is a use-after-free).
1341 # GC reclaims it once the blocking call returns. Mirrors
1342 # purge_local_memory's guard.
1343 GitPolicyFetcher.repos.pop(path, None)
1344 continue
1345 GitPolicyFetcher.forget_repo(path)
1346 GitPolicyFetcher.repos.clear()
1347 GitPolicyFetcher.repos_last_fetched.clear()
1348 GitPolicyFetcher.repo_locks.clear()
1351class GitCallback(RemoteCallbacks):
1352 def __init__(self, source: GitPolicyScopeSource):
1353 super().__init__()
1354 self._source = source
1356 def credentials(self, url, username_from_url, allowed_types):
1357 if isinstance(self._source.auth, SSHAuthData):
1358 auth = cast(SSHAuthData, self._source.auth)
1360 ssh_key = dict(
1361 username=username_from_url,
1362 pubkey=auth.public_key or "",
1363 privkey=auth.private_key,
1364 passphrase="",
1365 )
1366 return KeypairFromMemory(**ssh_key)
1367 if isinstance(self._source.auth, GitHubTokenAuthData):
1368 auth = cast(GitHubTokenAuthData, self._source.auth)
1370 return UserPass(username="git", password=auth.token)
1372 return Username(username_from_url)