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

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 

18 

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) 

46 

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

55 

56 

57class GitConcurrencyLimitExceeded(RuntimeError): 

58 """Raised when in-flight (live + zombie) git ops reach 

59 SCOPES_GIT_MAX_ZOMBIES.""" 

60 

61 

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. 

65 

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. 

71 

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. 

78 

79 Subclasses ValueError so broad handlers still catch it. 

80 

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

89 

90 waited_seconds: float = 0.0 

91 client_disconnected: bool = False 

92 

93 

94class BranchHeadNotFoundError(ValueError): 

95 """Configured branch has no resolvable HEAD (permanent misconfig), NOT a 

96 transient clone gap. 

97 

98 Subclasses ValueError so broad handlers still catch it. 

99 """ 

100 

101 

102_zombie_cap_logged = False 

103 

104 

105class _DaemonThreadPoolExecutor(ThreadPoolExecutor): 

106 """A ``ThreadPoolExecutor`` whose worker threads are daemon threads. 

107 

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

116 

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

121 

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 

166 

167 def weakref_cb(_, q=self._work_queue): 

168 q.put(None) 

169 

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. 

193 

194 

195def shutdown_git_executor() -> None: 

196 """Drop loop-bound live-op accounting before fork. 

197 

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

204 

205 

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

219 

220 

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 ) 

227 

228 

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 ) 

243 

244 

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) 

250 

251 

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) 

263 

264 

265def git_op_in_flight(key: str) -> bool: 

266 """True while a git op for ``key`` is still running on a pool thread. 

267 

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 

273 

274 

275def drain_git_ops(timeout: float) -> bool: 

276 """Block up to `timeout`s for all in-flight git ops to finish. 

277 

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) 

293 

294 

295def git_busy_count() -> int: 

296 """Number of scope git ops holding a pool thread (incl. 

297 

298 timed-out zombies). 

299 """ 

300 with _git_busy_lock: 

301 return len(_git_busy) 

302 

303 

304@dataclass 

305class SourceBackoff: 

306 """How long a repeatedly-failing source is skipped by the periodic pass. 

307 

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

312 

313 consecutive_failures: int 

314 next_attempt_at: float 

315 last_error: str 

316 

317 

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 

324 

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 

328 

329 

330def _finite_positive_or_zero(value) -> float: 

331 """Read a config number as a positive finite float, else 0.0. 

332 

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 

345 

346 

347def _backoff_base_seconds() -> float: 

348 """SCOPES_GIT_BACKOFF_BASE_SECONDS, validated. 0.0 means "disabled". 

349 

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) 

360 

361 

362def _backoff_max_seconds() -> float: 

363 """SCOPES_GIT_BACKOFF_MAX_SECONDS, validated. 0.0 means "no cap". 

364 

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) 

372 

373 

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 

385 

386 

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 ) 

417 

418 

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 

422 

423 

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 

432 

433 

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

440 

441 

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 

449 

450 

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. 

453 

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. 

459 

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. 

468 

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 ) 

500 

501 loop = asyncio.get_running_loop() 

502 

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) 

509 

510 sem = _get_live_ops_semaphore() 

511 await sem.acquire() 

512 released = False 

513 

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

519 

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

538 

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 

541 

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

554 

555 

556class PolicyFetcherCallbacks: 

557 async def on_update(self, old_head: Optional[str], head: str): 

558 pass 

559 

560 

561class PolicyFetcher: 

562 def __init__(self, callbacks): 

563 self.callbacks = callbacks 

564 

565 def fetch(self, hinted_hash: Optional[str] = None): 

566 raise NotImplementedError() 

567 

568 

569class RepoInterface: 

570 """Manages a git repo with pygit2.""" 

571 

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

594 

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 

602 

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 

609 

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 

617 

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) 

634 

635 

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] = {} 

650 

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 ) 

672 

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. 

677 

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 

691 

692 def _backoff_entry(self) -> Optional[SourceBackoff]: 

693 """This source's live backoff entry, or None if it may be attempted. 

694 

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 

706 

707 def _record_source_failure(self, err: BaseException) -> None: 

708 """Count one failed clone/fetch against this source and re-arm the 

709 delay. 

710 

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

768 

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

781 

782 @staticmethod 

783 def forget_source_backoff(source_id: str) -> None: 

784 """Drop a source's backoff entry when the source itself goes away. 

785 

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

793 

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 

799 

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. 

809 

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. 

815 

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 ) 

959 

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) 

986 

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

1012 

1013 def _discover_repository(self, path: Path) -> bool: 

1014 git_path: Path = path / ".git" 

1015 return discover_repository(str(path)) and git_path.exists() 

1016 

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) 

1079 

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] 

1085 

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 

1124 

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 

1139 

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 

1145 

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 

1155 

1156 # by default, we try to avoid re-fetching the repo for performance 

1157 return False 

1158 

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 

1165 

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

1168 

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 

1177 

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 

1188 

1189 await self.callbacks.on_update(old_revision, new_revision) 

1190 

1191 # Bring forward local branch (a bit like "pull"), so we won't detect changes again 

1192 local_branch.set_target(new_revision) 

1193 

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 

1240 

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

1252 

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) 

1261 

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

1270 

1271 @staticmethod 

1272 def base_dir(base_dir: Path) -> Path: 

1273 return base_dir / "git_sources" 

1274 

1275 @staticmethod 

1276 def repo_clone_path(base_dir: Path, source: GitPolicyScopeSource) -> Path: 

1277 return GitPolicyFetcher.base_dir(base_dir) / GitPolicyFetcher.source_id(source) 

1278 

1279 @staticmethod 

1280 def forget_repo(path: str) -> None: 

1281 """Drop the cached repository for a clone path and release its handles. 

1282 

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 ) 

1302 

1303 @staticmethod 

1304 def reset_caches() -> None: 

1305 """Free and drop every cached repo handle, lock, and timestamp. 

1306 

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. 

1316 

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. 

1323 

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

1349 

1350 

1351class GitCallback(RemoteCallbacks): 

1352 def __init__(self, source: GitPolicyScopeSource): 

1353 super().__init__() 

1354 self._source = source 

1355 

1356 def credentials(self, url, username_from_url, allowed_types): 

1357 if isinstance(self._source.auth, SSHAuthData): 

1358 auth = cast(SSHAuthData, self._source.auth) 

1359 

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) 

1369 

1370 return UserPass(username="git", password=auth.token) 

1371 

1372 return Username(username_from_url)