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

186 statements  

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

1import asyncio 

2import datetime 

3import shutil 

4from functools import partial 

5from pathlib import Path 

6from typing import List, Optional, Set, cast 

7 

8import git 

9from ddtrace import tracer 

10from fastapi_websocket_pubsub import PubSubEndpoint 

11from opal_common.async_utils import run_sync 

12from opal_common.git_utils.commit_viewer import VersionedFile 

13from opal_common.http_utils import redact_url 

14from opal_common.logger import logger 

15from opal_common.monitoring import metrics 

16from opal_common.schemas.policy import PolicyUpdateMessageNotification 

17from opal_common.schemas.policy_source import GitPolicyScopeSource 

18from opal_common.topics.publisher import ScopedServerSideTopicPublisher 

19from opal_server.config import opal_server_config 

20from opal_server.git_fetcher import ( 

21 GitConcurrencyLimitExceeded, 

22 GitPolicyFetcher, 

23 PolicyFetcherCallbacks, 

24 emit_sources_in_backoff, 

25 git_op_in_flight, 

26) 

27from opal_server.policy.watcher.callbacks import ( 

28 create_policy_update, 

29 create_update_all_directories_in_repo, 

30) 

31from opal_server.scopes.purge import ( 

32 ScopePurgeCommand, 

33 confined_clone_path, 

34 find_scope_sharing_source, 

35) 

36from opal_server.scopes.scope_repository import ( 

37 Scope, 

38 ScopeNotFoundError, 

39 ScopeRepository, 

40) 

41 

42 

43def is_rego_source_file( 

44 f: VersionedFile, extensions: Optional[List[str]] = None 

45) -> bool: 

46 """Filters only rego files or data.json files.""" 

47 REGO = ".rego" 

48 JSON = ".json" 

49 OPA_JSON = "data.json" 

50 

51 if extensions is None: 51 ↛ 52line 51 didn't jump to line 52 because the condition on line 51 was never true

52 extensions = [REGO, JSON] 

53 if JSON in extensions and f.path.suffix == JSON: 

54 return f.path.name == OPA_JSON 

55 return f.path.suffix in extensions 

56 

57 

58class NewCommitsCallbacks(PolicyFetcherCallbacks): 

59 def __init__( 

60 self, 

61 base_dir: Path, 

62 scope_id: str, 

63 source: GitPolicyScopeSource, 

64 pubsub_endpoint: PubSubEndpoint, 

65 ): 

66 self._scope_repo_dir = GitPolicyFetcher.repo_clone_path(base_dir, source) 

67 self._scope_id = scope_id 

68 self._source = source 

69 self._pubsub_endpoint = pubsub_endpoint 

70 

71 async def on_update(self, previous_head: str, head: str): 

72 if previous_head == head: 

73 logger.debug( 

74 f"scope '{self._scope_id}': No new commits, HEAD is at '{head}'" 

75 ) 

76 return 

77 

78 logger.info( 

79 f"scope '{self._scope_id}': Found new commits: old HEAD was '{previous_head}', new HEAD is '{head}'" 

80 ) 

81 if not self._scope_repo_dir.exists(): 81 ↛ 82line 81 didn't jump to line 82 because the condition on line 81 was never true

82 logger.error( 

83 f"on_update({self._scope_id}) was triggered, but repo path is not found: {self._scope_repo_dir}" 

84 ) 

85 return 

86 

87 try: 

88 repo = git.Repo(self._scope_repo_dir) 

89 except git.GitError as exc: 

90 logger.error( 

91 f"Got exception for repo in path: {self._scope_repo_dir}, scope_id: {self._scope_id}, error: {exc}" 

92 ) 

93 return 

94 

95 notification: Optional[PolicyUpdateMessageNotification] = None 

96 predicate = partial(is_rego_source_file, extensions=self._source.extensions) 

97 if previous_head is None: 97 ↛ 102line 97 didn't jump to line 102 because the condition on line 97 was always true

98 notification = await create_update_all_directories_in_repo( 

99 repo.commit(head), repo.commit(head), predicate=predicate 

100 ) 

101 else: 

102 notification = await create_policy_update( 

103 repo.commit(previous_head), 

104 repo.commit(head), 

105 self._source.extensions, 

106 predicate=predicate, 

107 ) 

108 

109 if notification is not None: 109 ↛ exitline 109 didn't return from function 'on_update' because the condition on line 109 was always true

110 await self.trigger_notification(notification) 

111 

112 async def trigger_notification(self, notification: PolicyUpdateMessageNotification): 

113 logger.info( 

114 f"Triggering policy update for scope {self._scope_id}: {notification.dict()}" 

115 ) 

116 async with ScopedServerSideTopicPublisher( 

117 self._pubsub_endpoint, self._scope_id 

118 ) as publisher: 

119 await publisher.publish(notification.topics, notification.update) 

120 

121 

122class ScopesService: 

123 def __init__( 

124 self, 

125 base_dir: Path, 

126 scopes: ScopeRepository, 

127 pubsub_endpoint: PubSubEndpoint, 

128 ): 

129 self._base_dir = base_dir 

130 self._scopes = scopes 

131 self._pubsub_endpoint = pubsub_endpoint 

132 # Strong refs to the best-effort local clone purges delete_scope spawns 

133 # (create_task results are otherwise GC-able); discarded on completion. 

134 self._local_purges: Set[asyncio.Task] = set() 

135 

136 async def sync_scope( 

137 self, 

138 scope_id: str = None, 

139 scope: Scope = None, 

140 hinted_hash: Optional[str] = None, 

141 force_fetch: bool = False, 

142 notify_on_changes: bool = True, 

143 req_time: datetime.datetime = None, 

144 honor_backoff: bool = False, 

145 ): 

146 """Sync one scope's policy source. 

147 

148 ``honor_backoff`` marks this call as pass-originated, letting the 

149 fetcher skip a source that keeps failing (see 

150 SCOPES_GIT_BACKOFF_BASE_SECONDS). It defaults to False so the explicit 

151 callers — POST /scopes/{id}/refresh and PUT /scopes, both of which 

152 arrive here through the watcher's ``trigger`` — always attempt the 

153 source: a customer who has just repaired credentials must not be told 

154 200 OK and then wait out hours or days of backoff. 

155 """ 

156 if scope is None: 156 ↛ 160line 156 didn't jump to line 160 because the condition on line 156 was always true

157 assert scope_id, ValueError("scope_id not set for sync_scope") 

158 scope = await self._scopes.get(scope_id) 

159 

160 with tracer.trace("scopes_service.sync_scope", resource=scope.scope_id): 

161 if not isinstance(scope.policy, GitPolicyScopeSource): 161 ↛ 162line 161 didn't jump to line 162 because the condition on line 161 was never true

162 logger.warning("Non-git scopes are currently not supported!") 

163 return 

164 source = cast(GitPolicyScopeSource, scope.policy) 

165 

166 logger.debug( 

167 f"Sync scope: {scope.scope_id} (remote: {redact_url(source.url)}, branch: {source.branch}, req_time: {req_time})" 

168 ) 

169 

170 callbacks = PolicyFetcherCallbacks() 

171 if notify_on_changes: 171 ↛ 179line 171 didn't jump to line 179 because the condition on line 171 was always true

172 callbacks = NewCommitsCallbacks( 

173 base_dir=self._base_dir, 

174 scope_id=scope.scope_id, 

175 source=source, 

176 pubsub_endpoint=self._pubsub_endpoint, 

177 ) 

178 

179 source_id = GitPolicyFetcher.source_id(source) 

180 

181 async def _scope_still_exists() -> bool: 

182 # Also confirms the scope still points at the source this 

183 # fetcher is syncing: after a repoint the old source was 

184 # already purged, so cloning it again would strand a dir that 

185 # nothing reclaims (no later purge names that source). 

186 try: 

187 fresh = await self._scopes.get(scope.scope_id) 

188 except ScopeNotFoundError: 

189 return False 

190 if not isinstance(fresh.policy, GitPolicyScopeSource): 190 ↛ 191line 190 didn't jump to line 191 because the condition on line 190 was never true

191 return False 

192 return GitPolicyFetcher.source_id(fresh.policy) == source_id 

193 

194 fetcher = GitPolicyFetcher( 

195 self._base_dir, 

196 scope.scope_id, 

197 source, 

198 callbacks=callbacks, 

199 liveness_probe=_scope_still_exists, 

200 ) 

201 

202 try: 

203 await fetcher.fetch_and_notify_on_changes( 

204 hinted_hash=hinted_hash, 

205 force_fetch=force_fetch, 

206 req_time=req_time, 

207 honor_backoff=honor_backoff, 

208 ) 

209 except GitConcurrencyLimitExceeded as e: 

210 # Expected backpressure, not a fault: the zombie cap 

211 # (SCOPES_GIT_MAX_ZOMBIES) is refusing new git ops because too 

212 # many are stuck on unreachable remotes. Log it cleanly at 

213 # warning — a full stack trace per refused scope per pass would 

214 # bury the one cap-reached signal under noise during an outage. 

215 logger.warning( 

216 "Skipping scope {scope_id} this pass: {err}", 

217 scope_id=scope.scope_id, 

218 err=e, 

219 ) 

220 except Exception as e: 

221 # ERROR without the traceback: with a broken tail of dozens of 

222 # sources this line is emitted per source per pass, and the 

223 # ~40-line traceback that used to ride along rotated the 

224 # container log (10 MiB) within minutes of a boot — the boot 

225 # markers were gone before anyone could read them, and in prod 

226 # it is the mechanism behind opal-server being 56% of the org's 

227 # log bytes. The traceback is still there at DEBUG for anyone 

228 # chasing a specific source; the source, scope and reason are 

229 # in the ERROR line. 

230 logger.error( 

231 "Could not fetch policy for scope {scope_id} " 

232 "(remote: {url}): {etype}: {err}", 

233 scope_id=scope.scope_id, 

234 url=redact_url(scope.policy.url), 

235 etype=type(e).__name__, 

236 err=e, 

237 ) 

238 logger.opt(exception=True).debug( 

239 "traceback for the failed sync of scope {scope_id}", 

240 scope_id=scope.scope_id, 

241 ) 

242 

243 async def delete_scope(self, scope_id: str): 

244 with tracer.trace("scopes_service.delete_scope", resource=scope_id): 

245 logger.info(f"Delete scope: {scope_id}") 

246 scope = await self._scopes.get(scope_id) 

247 

248 if not isinstance(scope.policy, GitPolicyScopeSource): 248 ↛ 251line 248 didn't jump to line 251 because the condition on line 248 was never true

249 # Mirrors sync_scope: only git sources have a clone dir and 

250 # fetcher caches to clean up. 

251 logger.warning( 

252 f"Scope {scope_id} has a non-git policy source, " 

253 "deleting the scope record only" 

254 ) 

255 await self._scopes.delete(scope_id) 

256 return 

257 

258 deleted_source_id = GitPolicyFetcher.source_id(scope.policy) 

259 scope_dir = GitPolicyFetcher.repo_clone_path(self._base_dir, scope.policy) 

260 

261 try: 

262 await self._scopes.delete(scope_id) 

263 finally: 

264 # The publish must stay reachable even when the record delete 

265 # raises an ambiguous outcome (committed server-side, error 

266 # surfaced to the client): the retry is a 204 no-op 

267 # (ScopeNotFoundError), so a publish gated on a clean delete 

268 # would orphan the purge permanently. Over-publishing 

269 # self-heals: the leader's sibling-check sees a still-live 

270 # record and keeps everything. Memory entries (all workers, 

271 # this one included) drop when the leader's confirmation 

272 # broadcast arrives. 

273 # FLOOR, not the primary path. The publish above is the fleet- 

274 # wide purge, but it is droppable at shipped defaults: a DELETE 

275 # usually lands on a non-leader worker (SERVER_WORKER_COUNT 

276 # defaults to the core count) and must traverse the broadcaster, 

277 # while a NON-LEADER worker only has a broadcaster reader if it 

278 # has a connected client or STATISTICS_ENABLED (default False). 

279 # (The LEADER always has one: its watcher enters a listening 

280 # context unconditionally — policy/watcher/task.py — so the 

281 # earlier claim that the leader could be deaf was backwards.) 

282 # A backbone outage still loses the message for everyone. If it 

283 # never arrives, nothing removes the dir — and on master 

284 # delete_scope removed it INLINE here, with no broadcast 

285 # involved, so without this the lost-broadcast case is a 

286 # regression against the merge base rather than parity with it. 

287 # 

288 # Backgrounded because DELETE's latency is bounded by contract 

289 # and this takes lock_source, which a sync holds across a whole 

290 # clone/fetch (unbounded when SCOPES_GIT_FETCH_TIMEOUT is 0). 

291 # 

292 # LOAD-BEARING ORDER: scheduled BEFORE the publish below, never 

293 # after it. publish() can raise a broadcaster error (see 

294 # LeaderScopePurger._purge_and_log), and SCOPES_PURGE_CHANNEL is 

295 # freeze-exempt, so during a backbone gap it is attempted and 

296 # fails rather than being deferred. Scheduling after it would 

297 # therefore skip the floor in precisely the degraded case the 

298 # floor exists to cover. create_task only schedules — the floor 

299 # cannot delay the publish or the response. 

300 # 

301 # It is a floor, not a guarantee: master serialized the record 

302 # delete and the sibling check under one lock, so the LAST of 

303 # two concurrent sibling deleters always saw no sharer. Here the 

304 # record delete is outside the lock, so an unlucky interleaving 

305 # can have both deleters see the other as still-live and both 

306 # skip. The leader's sibling-checked purge is the authoritative 

307 # path; this only guarantees that a delete reclaims at least the 

308 # serving pod's copy regardless of broadcaster state. 

309 task = asyncio.create_task( 

310 self._purge_local_clone_best_effort( 

311 deleted_source_id, scope_dir, scope_id 

312 ) 

313 ) 

314 self._local_purges.add(task) 

315 task.add_done_callback(self._local_purges.discard) 

316 if self._pubsub_endpoint is not None: 316 ↛ exitline 316 didn't jump to the function exit

317 await self._pubsub_endpoint.publish( 

318 [opal_server_config.SCOPES_PURGE_CHANNEL], 

319 ScopePurgeCommand( 

320 source_id=deleted_source_id, 

321 clone_path=str(scope_dir), 

322 scope_id=scope_id, 

323 reason="delete", 

324 ).dict(), 

325 ) 

326 

327 async def _purge_local_clone_best_effort( 

328 self, deleted_source_id: str, scope_dir: Path, scope_id: str 

329 ): 

330 """Remove THIS process's clone dir + fetcher cache entries for a 

331 deleted scope's source, unless a surviving scope still shares them. 

332 

333 Restores the floor master had (``_purge_source_cache_if_unshared``), 

334 with two changes master did not have: 

335 

336 - the store read is bounded by SCOPES_STORE_READ_TIMEOUT, because it is 

337 taken under ``lock_source`` and the Redis client has no socket timeout; 

338 - the removal is skipped while a git op is in flight for the source. 

339 Master freed the handle unconditionally; freeing one a lingering 

340 timed-out pygit2 call still holds on a pool thread is the 

341 use-after-free class 89e090be fixed. NOTHING else owns that case: 

342 the leader does no disk work and the deferred retry was cut, so a 

343 delete whose remote is hung leaves the dir until PER-15612. 

344 """ 

345 # Every other destructive path in this series derives its target from 

346 # source_id via confined_clone_path and refuses a malformed id; this one 

347 # took the caller's Path. It is not wire-controlled (it comes from the 

348 # stored record, not a pub/sub message), so this is consistency rather 

349 # than a live hole — but "the one rmtree that skips the check" is not a 

350 # sentence worth leaving in a series about unsafe deletes. 

351 safe_path = confined_clone_path(self._base_dir, deleted_source_id) 

352 if safe_path is None or safe_path != str(scope_dir): 352 ↛ 353line 352 didn't jump to line 353 because the condition on line 352 was never true

353 logger.warning( 

354 f"Skipping the local clone purge for scope {scope_id}: derived " 

355 f"path {safe_path!r} does not match {str(scope_dir)!r}" 

356 ) 

357 return 

358 try: 

359 async with GitPolicyFetcher.lock_source(deleted_source_id): 

360 # I4: drain the entry lock_source minted on every exit that 

361 # abandons this source — the early returns below included, which 

362 # previously leaked one each. NOT on the live-sibling path: that 

363 # source is still in use, and the bed asserts the entry survives 

364 # a sibling delete. Same rule and same lock-identity guard as 

365 # LeaderScopePurger.purge_source_if_unshared. 

366 minted = GitPolicyFetcher.repo_locks.get(deleted_source_id) 

367 try: 

368 try: 

369 timeout = opal_server_config.SCOPES_STORE_READ_TIMEOUT 

370 check = find_scope_sharing_source( 

371 self._scopes, deleted_source_id 

372 ) 

373 sharer = await ( 

374 asyncio.wait_for(check, timeout=timeout) 

375 if timeout > 0 

376 else check 

377 ) 

378 except Exception as e: 

379 # KEEP the clone, unlike master. Master purged defensively 

380 # on any scan failure, reasoning that over-purging 

381 # self-heals. It does not self-heal cheaply here: if a 

382 # sibling scope does share this source, deleting its clone 

383 # takes a LIVE tenant's policy offline until the re-clone 

384 # completes — and the trigger is a transient store blip. 

385 # The cost of keeping is an orphan dir (PER-15612): disk 

386 # against availability, and this is a best-effort floor, 

387 # so it takes the conservative branch when it cannot tell. 

388 logger.warning( 

389 f"Local sibling check for {deleted_source_id} failed " 

390 f"after deleting scope {scope_id}; keeping this " 

391 f"worker's clone (it stays until PER-15612's sweep " 

392 f"lands): {e!r}" 

393 ) 

394 return 

395 if sharer is not None: 

396 logger.info( 

397 f"Scope {sharer} still shares source " 

398 f"{deleted_source_id}, keeping this worker's clone" 

399 ) 

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

401 return 

402 if git_op_in_flight(deleted_source_id): 402 ↛ 403line 402 didn't jump to line 403 because the condition on line 402 was never true

403 logger.info( 

404 f"Skipping the local clone purge for " 

405 f"{deleted_source_id}: a git operation is still in " 

406 "flight; the dir stays until PER-15612's sweep lands" 

407 ) 

408 return 

409 GitPolicyFetcher.forget_repo(safe_path) 

410 GitPolicyFetcher.repos_last_fetched.pop(deleted_source_id, None) 

411 # Same reason as repos_last_fetched: no live scope on this 

412 # worker points at the source any more, so an entry kept 

413 # here is counted in the sources_in_backoff gauge for the 

414 # life of the process — and would suppress the first sync 

415 # of a scope later re-created against the same URL, on the 

416 # strength of a deleted scope's failure history. 

417 GitPolicyFetcher.forget_source_backoff(deleted_source_id) 

418 try: 

419 await run_sync(shutil.rmtree, safe_path) 

420 except FileNotFoundError: 

421 pass # never cloned (or already gone) — nothing to clean 

422 except OSError as e: 

423 logger.warning( 

424 f"Failed to remove clone dir {safe_path} of deleted " 

425 f"scope {scope_id}: {e!r}" 

426 ) 

427 finally: 

428 if ( 

429 minted is not None 

430 and GitPolicyFetcher.repo_locks.get(deleted_source_id) is minted 

431 ): 

432 # Popped while the lock is held: lock_source waiters 

433 # re-check the dict entry after acquiring and retry on the 

434 # freshly-minted lock. 

435 GitPolicyFetcher.repo_locks.pop(deleted_source_id, None) 

436 except Exception: 

437 # Detached background task: without this an unexpected failure 

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

439 logger.exception( 

440 f"Best-effort local clone purge for source {deleted_source_id} " 

441 f"(scope {scope_id}) failed" 

442 ) 

443 

444 async def stop(self) -> None: 

445 """Await the best-effort local clone purges delete_scope spawned. 

446 

447 Without this a DELETE that returns 204 and is followed by SIGTERM loses 

448 its floor: the task is detached, nothing else references it, and the 

449 clone dir it was about to remove survives with nothing left to reclaim 

450 it (no reconciliation in this PR — PER-15612). 

451 

452 Best-effort and expected to be bounded by the caller: each task takes 

453 lock_source, which a sync can hold across a whole clone/fetch. The 

454 watcher calls this after cancelling its tasks and under a wait_for, the 

455 same discipline LeaderScopePurger.stop() already gets. 

456 """ 

457 if self._local_purges: 457 ↛ 458line 457 didn't jump to line 458 because the condition on line 457 was never true

458 await asyncio.gather(*list(self._local_purges), return_exceptions=True) 

459 

460 async def sync_scopes( 

461 self, only_poll_updates=False, notify_on_changes=True, honor_backoff=True 

462 ): 

463 """Sync every scope, in two phases. 

464 

465 ``honor_backoff`` defaults to True because a whole-fleet pass is what 

466 the per-source backoff exists for: the periodic poll and the pre-fork 

467 boot preload both land here, and both re-attempt every source they 

468 know about. Honouring it in BOTH phases is what collapses the 

469 duplicate storm — phase 2 visits every scope that merely reuses a 

470 source, and for a source with no local clone each of those goes 

471 straight to a clone of its own, so one dead repo shared by N scopes 

472 costs N attempts per pass without it. 

473 

474 The refresh-all endpoint passes False: see ScopesPolicyWatcherTask. 

475 """ 

476 with tracer.trace("scopes_service.sync_scopes"): 

477 scopes = await self._scopes.all() 

478 # Emitted before the poll-updates filter below, so this is always the 

479 # true total rather than flapping with the caller's filter. 

480 metrics.gauge("opal_server.scopes.count", len(scopes)) 

481 # Once per pass as well as on transitions: a DogStatsD gauge is 

482 # NO DATA between sends, and the steady state the backoff creates 

483 # (dead sources parked for the cap) has almost no transitions. 

484 emit_sources_in_backoff() 

485 if only_poll_updates: 485 ↛ 487line 485 didn't jump to line 487 because the condition on line 485 was never true

486 # Only sync scopes that have polling enabled (in a periodic check) 

487 scopes = [scope for scope in scopes if scope.policy.poll_updates] 

488 

489 logger.info( 

490 f"OPAL Scopes: syncing {len(scopes)} scopes in the background (polling updates: {only_poll_updates})" 

491 ) 

492 

493 # Partition into distinct repos (cloned/fetched once, with priority 

494 # so every repo is pulled asap) and the scopes that merely reuse an 

495 # already-handled repo (checked for changes only). 

496 unique_scopes = [] 

497 duplicate_scopes = [] 

498 seen_source_ids = set() 

499 for scope in scopes: 

500 src_id = GitPolicyFetcher.source_id(scope.policy) 

501 if src_id in seen_source_ids: 

502 duplicate_scopes.append(scope) 

503 else: 

504 seen_source_ids.add(src_id) 

505 unique_scopes.append(scope) 

506 

507 # Phase 1 clones/fetches every distinct repo; phase 2 then checks the 

508 # duplicates against those now-present repos. 

509 # 

510 # The two phases have different cost profiles, so they get separate 

511 # bounds. Phase 1 does the network clone/fetch, so it is capped at 

512 # SCOPES_GIT_MAX_WORKERS: one unreachable repo then only stalls its 

513 # own slot (for the fetch timeout), not the whole pass. 

514 git_semaphore = asyncio.Semaphore( 

515 max(1, opal_server_config.SCOPES_GIT_MAX_WORKERS) 

516 ) 

517 await self._sync_scopes_concurrently( 

518 unique_scopes, 

519 git_semaphore, 

520 force_fetch=True, 

521 notify_on_changes=notify_on_changes, 

522 honor_backoff=honor_backoff, 

523 ) 

524 

525 # Phase 2 is local-only in the common case: the repos were just 

526 # handled in phase 1, so _should_fetch returns False and no network 

527 # fetch happens (only a disk open + change-check + notify). 

528 # It shares the loop's default executor (the same pool that serves 

529 # policy bundles), so bound it by the SAME SCOPES_GIT_MAX_WORKERS 

530 # knob as phase 1 rather than a hard-coded floor: an operator who 

531 # lowers the knob to protect a small pod must be able to lower phase 

532 # 2 too, and over-subscribing that ~min(32, cpu+4)-thread pool only 

533 # queues work and contends with bundle serving. (The earlier 

534 # max(..., 32) made 32 a floor the knob could never reduce below.) 

535 local_concurrency = max(1, opal_server_config.SCOPES_GIT_MAX_WORKERS) 

536 local_semaphore = asyncio.Semaphore(local_concurrency) 

537 await self._sync_scopes_concurrently( 

538 duplicate_scopes, 

539 local_semaphore, 

540 force_fetch=False, 

541 notify_on_changes=notify_on_changes, 

542 honor_backoff=honor_backoff, 

543 ) 

544 

545 async def _sync_scopes_concurrently( 

546 self, scopes, semaphore, *, force_fetch, notify_on_changes, honor_backoff=False 

547 ): 

548 """Sync ``scopes`` concurrently, bounded by ``semaphore``. 

549 

550 Each scope's failure is logged and isolated so one bad repo 

551 never fails the whole pass. 

552 

553 Passes scope_id (not the snapshot object) so sync_scope re-gets 

554 fresh state right before use — a delete that landed after the 

555 snapshot surfaces as ScopeNotFoundError instead of re-cloning a 

556 dead scope's repo and re-populating the fetcher caches for it. 

557 """ 

558 

559 async def _sync_one(scope): 

560 async with semaphore: 

561 try: 

562 await self.sync_scope( 

563 scope_id=scope.scope_id, 

564 force_fetch=force_fetch, 

565 notify_on_changes=notify_on_changes, 

566 honor_backoff=honor_backoff, 

567 ) 

568 except ScopeNotFoundError: 

569 logger.info( 

570 f"scope {scope.scope_id} was deleted while sync was queued, skipping" 

571 ) 

572 except Exception as e: 

573 # sync_scope already logs its own git failures without a 

574 # traceback (see there); this catches anything that escaped 

575 # it. Same rule: one ERROR line, traceback at DEBUG. 

576 logger.error( 

577 "sync_scope failed for {scope_id}: {etype}: {err}", 

578 scope_id=scope.scope_id, 

579 etype=type(e).__name__, 

580 err=e, 

581 ) 

582 logger.opt(exception=True).debug( 

583 "traceback for the failed sync of scope {scope_id}", 

584 scope_id=scope.scope_id, 

585 ) 

586 

587 await asyncio.gather(*(_sync_one(scope) for scope in scopes))