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

289 statements  

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

1import asyncio 

2import math 

3import os 

4import pathlib 

5from typing import List, Optional, cast 

6 

7import pygit2 

8from fastapi import ( 

9 APIRouter, 

10 Depends, 

11 Header, 

12 HTTPException, 

13 Path, 

14 Query, 

15 Request, 

16 Response, 

17 status, 

18) 

19from fastapi.responses import RedirectResponse 

20from fastapi_websocket_pubsub import PubSubEndpoint 

21from git import InvalidGitRepositoryError, NoSuchPathError 

22from opal_common.async_utils import run_sync 

23from opal_common.authentication.authz import ( 

24 require_peer_type, 

25 restrict_optional_topics_to_publish, 

26) 

27from opal_common.authentication.casting import cast_private_key 

28from opal_common.authentication.deps import JWTAuthenticator, get_token_from_header 

29from opal_common.authentication.types import EncryptionKeyFormat, JWTClaims 

30from opal_common.authentication.verifier import Unauthorized 

31from opal_common.logger import logger 

32from opal_common.monitoring import metrics 

33from opal_common.schemas.data import ( 

34 DataSourceConfig, 

35 DataUpdate, 

36 ServerDataSourceConfig, 

37) 

38from opal_common.schemas.policy import PolicyBundle, PolicyUpdateMessageNotification 

39from opal_common.schemas.policy_source import GitPolicyScopeSource, SSHAuthData 

40from opal_common.schemas.scopes import Scope 

41from opal_common.schemas.security import PeerType 

42from opal_common.topics.publisher import ( 

43 ScopedServerSideTopicPublisher, 

44 ServerSideTopicPublisher, 

45) 

46from opal_common.urls import set_url_query_param 

47from opal_server.config import opal_server_config 

48from opal_server.data.data_update_publisher import DataUpdatePublisher 

49from opal_server.git_fetcher import ( 

50 BranchHeadNotFoundError, 

51 CloneNotPopulatedError, 

52 GitPolicyFetcher, 

53) 

54from opal_server.scopes.purge import ScopePurgeCommand 

55from opal_server.scopes.scope_repository import ScopeNotFoundError, ScopeRepository 

56from opal_server.scopes.service import ScopesService 

57 

58# Retry-After hints, in seconds. Two constants rather than one escalating 

59# value: escalation would need per-client retry state on a stateless endpoint, 

60# and the expected wait genuinely differs between "a sync will re-create this 

61# on its next tick" and "a clone is running right now". 

62_RETRY_AFTER_CLONE_UNAVAILABLE = "5" 

63_RETRY_AFTER_CLONE_IN_PROGRESS = "30" 

64 

65# How often the clone wait re-checks the clone. A module constant rather than a 

66# second config key: the operator-visible quantity is the total hold 

67# (SCOPES_POLICY_CLONE_WAIT_SECONDS), while every SCOPES_* key is permanent 

68# public surface — config_docs_drift_test pins each one verbatim into the 

69# published reference. One poll is a cheap disk read (open the repo, list its 

70# refs), so a second between polls costs at most one such read per waiting 

71# request per second and still returns within a second of the clone landing. 

72# The published description says "once a second"; a test couples the two. 

73_CLONE_WAIT_POLL_SECONDS = 1.0 

74 

75# Ceiling on the configured hold. Above the load balancer's 60s idle timeout a 

76# hold stops being a hold and becomes a 504 — the exact failure the wait exists 

77# to prevent — so an over-large value is clamped rather than honoured. 

78_CLONE_WAIT_MAX_SECONDS = 55.0 

79 

80# Requests this process is currently holding in the wait. A plain int, no lock: 

81# it is read and written only from the event loop thread, between awaits, so 

82# the increment and the cap check cannot interleave with another request's. 

83_clone_wait_inflight = 0 

84 

85# The clamp warning latches: on a misconfigured fleet it would otherwise be one 

86# identical line per request. 

87_clone_wait_clamp_logged = False 

88 

89_CLONE_WAIT_METRIC = "opal_server.scopes.policy_clone_wait" 

90_CLONE_WAIT_INFLIGHT_METRIC = "opal_server.scopes.policy_clone_wait_inflight" 

91_CLONE_WAIT_SECONDS_METRIC = "opal_server.scopes.policy_clone_wait_seconds" 

92 

93 

94def _bounded_clone_wait() -> float: 

95 """The configured hold, validated and clamped. 0.0 means "do not wait". 

96 

97 The non-finite trio is what this exists for: `nan`, `inf` and `-inf` all 

98 parse cleanly, so a process configured with one of them starts normally 

99 and reaches here. `inf` would silently become the clamped maximum on every 

100 clone-in-progress request — a 55s hold nobody asked for — and `nan` makes 

101 every comparison against the deadline False, which the loop is written to 

102 survive but which is not a budget anyone meant to set. 

103 """ 

104 global _clone_wait_clamp_logged 

105 

106 try: 

107 wait = float(opal_server_config.SCOPES_POLICY_CLONE_WAIT_SECONDS) 

108 except (TypeError, ValueError): 

109 # Belt and braces, and NOT the load-bearing half: Confi parses the 

110 # environment once, when this module is imported, so 

111 # OPAL_SCOPES_POLICY_CLONE_WAIT_SECONDS=abc already fails the process 

112 # at startup and never reaches this line. What this covers is a value 

113 # assigned to the config object at runtime. 

114 return 0.0 

115 if not math.isfinite(wait) or wait <= 0: 

116 return 0.0 

117 if wait > _CLONE_WAIT_MAX_SECONDS: 

118 if not _clone_wait_clamp_logged: 

119 _clone_wait_clamp_logged = True 

120 logger.warning( 

121 "SCOPES_POLICY_CLONE_WAIT_SECONDS={configured}s exceeds the " 

122 "{ceiling}s ceiling and is clamped: a hold longer than the load " 

123 "balancer's idle timeout is served as a 504, not as a bundle", 

124 configured=wait, 

125 ceiling=_CLONE_WAIT_MAX_SECONDS, 

126 ) 

127 return _CLONE_WAIT_MAX_SECONDS 

128 return wait 

129 

130 

131def _publish_clone_wait_inflight() -> None: 

132 """Publish the held-request count for this process. 

133 

134 Tagged by pid only. A pod's workers each hold their own count, so an 

135 untagged series would be last-write-wins across them and a saturated 

136 worker would be invisible. Deliberately NOT tagged by scope_id or 

137 source_id: the cap is a per-process resource, and those tags are unbounded 

138 cardinality. 

139 """ 

140 metrics.gauge( 

141 _CLONE_WAIT_INFLIGHT_METRIC, 

142 _clone_wait_inflight, 

143 tags={"pid": str(os.getpid())}, 

144 ) 

145 

146 

147async def _make_bundle_waiting_for_clone( 

148 fetcher: GitPolicyFetcher, 

149 base_hash: Optional[str], 

150 scope_id: str, 

151 request: Optional[Request] = None, 

152) -> PolicyBundle: 

153 """Build the bundle, holding the request while the clone is populated. 

154 

155 The immediate 503 this replaces is honest and useless: opal-client 

156 ignores Retry-After, makes five attempts with random-exponential backoff 

157 capped at 10s, then stays quiet until the next pub/sub policy message or 

158 a reconnect. A clone that outlives those attempts leaves that PDP with no 

159 policy and nothing scheduled to fix it — the update-all published when the 

160 clone completes names only the scope that was syncing, so siblings sharing 

161 the clone are never woken. 

162 

163 Readiness comes from CloneNotPopulatedError, which is derived from disk, 

164 so this works on the workers that are not running the clone — which is all 

165 of them but one. 

166 

167 What is bounded is the WAIT plus at most one more bundle attempt. The 

168 attempt itself runs on the loop's shared default executor, so time spent 

169 queued behind other builds is outside the deadline; that queue is what 

170 SCOPES_POLICY_CLONE_WAIT_MAX_INFLIGHT bounds, by capping how many requests 

171 can be released into it at once. Excess requests are shed with the answer 

172 they would have got before the wait existed. 

173 

174 The hold takes no lock, touches no cache and occupies no thread between 

175 polls, so it is cancellation-safe: a client that hangs up mid-wait unwinds 

176 at the next await with nothing to undo. It is also abandoned as soon as the 

177 client is seen to have disconnected — nobody is waiting for that bundle, 

178 and the slot is worth more to a caller that is still listening. 

179 

180 Returns the bundle, or re-raises CloneNotPopulatedError once the budget is 

181 spent, so EACH caller's own handler shapes that answer (the primary path 

182 answers Retry-After 30, the default-scope path 5). Every OTHER exception 

183 propagates untouched from whichever attempt raised it: a clone can finish 

184 and still fail to build a bundle (an absent branch is a 409, a gutted 

185 object store a retryable 503), and the wait must not re-label those. 

186 

187 Exactly one `policy_clone_wait` count is emitted per request that reaches 

188 the wait, tagged with how it ended: served, timeout, shed, disconnected, 

189 cancelled or error. 

190 """ 

191 global _clone_wait_inflight 

192 

193 try: 

194 return await run_sync(fetcher.make_bundle, base_hash) 

195 except CloneNotPopulatedError as first_exc: 

196 wait = _bounded_clone_wait() 

197 if wait <= 0: 

198 first_exc.waited_seconds = 0.0 

199 raise 

200 pending = first_exc 

201 

202 cap = opal_server_config.SCOPES_POLICY_CLONE_WAIT_MAX_INFLIGHT 

203 if 0 < cap <= _clone_wait_inflight: 

204 logger.info( 

205 "Scope {scope_id} clone wait is at its {cap}-request cap; " 

206 "answering 503 without waiting", 

207 scope_id=scope_id, 

208 cap=cap, 

209 ) 

210 metrics.increment(_CLONE_WAIT_METRIC, tags={"outcome": "shed"}) 

211 pending.waited_seconds = 0.0 

212 raise pending 

213 

214 loop = asyncio.get_running_loop() 

215 started = loop.time() 

216 deadline = started + wait 

217 outcome = "error" 

218 

219 _clone_wait_inflight += 1 

220 try: 

221 # Inside the try, not before it: this call ends in a metrics sink, and 

222 # a sink that raises between the increment and the try would leak the 

223 # slot for the life of the process — permanently lowering the cap. 

224 _publish_clone_wait_inflight() 

225 while True: 

226 remaining = deadline - loop.time() 

227 # `not (remaining > 0)` rather than `remaining <= 0`: NaN compares 

228 # False against BOTH, so the `<=` form does not break on a NaN 

229 # deadline — it polls forever, holding a capped slot for the life 

230 # of the process. _bounded_clone_wait already refuses a NaN budget, 

231 # so this is defence in depth: it makes the loop itself unable to 

232 # spin if that guard is ever weakened or bypassed. 

233 if not (remaining > 0): 

234 outcome = "timeout" 

235 break 

236 # Clamped to what is left of the budget, so the last poll cannot 

237 # overshoot the hold an operator configured. 

238 await asyncio.sleep(min(_CLONE_WAIT_POLL_SECONDS, remaining)) 

239 

240 if request is not None and await request.is_disconnected(): 

241 outcome = "disconnected" 

242 logger.info( 

243 "Scope {scope_id} clone wait abandoned after {waited:.1f}s: " 

244 "the client disconnected", 

245 scope_id=scope_id, 

246 waited=loop.time() - started, 

247 ) 

248 pending.waited_seconds = loop.time() - started 

249 pending.client_disconnected = True 

250 raise pending 

251 

252 try: 

253 bundle = await run_sync(fetcher.make_bundle, base_hash) 

254 except CloneNotPopulatedError as exc: 

255 pending = exc 

256 continue 

257 outcome = "served" 

258 logger.info( 

259 "Scope {scope_id} clone became available after {waited:.1f}s wait", 

260 scope_id=scope_id, 

261 waited=loop.time() - started, 

262 ) 

263 return bundle 

264 

265 # Carried on the exception so each caller can report the hold without 

266 # this function having to know how the 503 it falls through to is 

267 # shaped. 

268 pending.waited_seconds = loop.time() - started 

269 raise pending 

270 except asyncio.CancelledError: 

271 # A BaseException since 3.8, so no `except Exception` arm would see it. 

272 # Left uncounted, a fleet whose waits are all being torn down would 

273 # look exactly like one where nothing is waiting. 

274 outcome = "cancelled" 

275 logger.info( 

276 "Scope {scope_id} clone wait cancelled after {waited:.1f}s", 

277 scope_id=scope_id, 

278 waited=loop.time() - started, 

279 ) 

280 raise 

281 except Exception as exc: 

282 # Only the UNCLASSIFIED failure: `outcome` is already timeout or 

283 # disconnected when the exception being unwound is `pending`, and 

284 # those two are logged by whoever shapes the response. This arm is 

285 # what makes the `error` count readable — a bundle build that failed 

286 # for its own reasons AFTER the clone appeared. 

287 if outcome == "error": 

288 logger.info( 

289 "Scope {scope_id} held {waited:.1f}s waiting for its clone " 

290 "before failing: {exc!r}", 

291 scope_id=scope_id, 

292 waited=loop.time() - started, 

293 exc=exc, 

294 ) 

295 raise 

296 finally: 

297 _clone_wait_inflight -= 1 

298 _publish_clone_wait_inflight() 

299 metrics.increment(_CLONE_WAIT_METRIC, tags={"outcome": outcome}) 

300 if outcome in ("served", "timeout"): 

301 metrics.gauge( 

302 _CLONE_WAIT_SECONDS_METRIC, 

303 loop.time() - started, 

304 tags={"outcome": outcome}, 

305 ) 

306 

307 

308def verify_private_key(private_key: str, key_format: EncryptionKeyFormat) -> bool: 

309 try: 

310 key = cast_private_key(private_key, key_format=key_format) 

311 return key is not None 

312 except Exception as e: 

313 return False 

314 

315 

316def verify_private_key_or_throw(scope_in: Scope): 

317 if isinstance(scope_in.policy.auth, SSHAuthData): 

318 auth = cast(SSHAuthData, scope_in.policy.auth) 

319 if not "\n" in auth.private_key: 

320 raise HTTPException( 

321 status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, 

322 detail={"error": "private key is expected to contain newlines!"}, 

323 ) 

324 

325 is_pem_key = verify_private_key( 

326 auth.private_key, key_format=EncryptionKeyFormat.pem 

327 ) 

328 is_ssh_key = verify_private_key( 

329 auth.private_key, key_format=EncryptionKeyFormat.ssh 

330 ) 

331 if not (is_pem_key or is_ssh_key): 331 ↛ exitline 331 didn't return from function 'verify_private_key_or_throw' because the condition on line 331 was always true

332 raise HTTPException( 

333 status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, 

334 detail={"error": "private key is invalid"}, 

335 ) 

336 

337 

338def init_scope_router( 

339 scopes: ScopeRepository, 

340 authenticator: JWTAuthenticator, 

341 pubsub_endpoint: PubSubEndpoint, 

342 scopes_service: ScopesService, 

343): 

344 router = APIRouter() 

345 

346 def _allowed_scoped_authenticator( 

347 claims: JWTClaims = Depends(authenticator), scope_id: str = Path(...) 

348 ): 

349 if not authenticator.enabled: 349 ↛ 350line 349 didn't jump to line 350 because the condition on line 349 was never true

350 return 

351 

352 allowed_scopes = claims.get("allowed_scopes") 

353 

354 if not allowed_scopes or scope_id not in allowed_scopes: 354 ↛ exitline 354 didn't return from function '_allowed_scoped_authenticator' because the condition on line 354 was always true

355 raise HTTPException(status.HTTP_403_FORBIDDEN) 

356 

357 @router.put("", status_code=status.HTTP_201_CREATED) 

358 async def put_scope( 

359 *, 

360 force_fetch: bool = Query( 

361 False, 

362 description="Whether the policy repo must be fetched from remote", 

363 ), 

364 scope_in: Scope, 

365 claims: JWTClaims = Depends(authenticator), 

366 ): 

367 try: 

368 require_peer_type(authenticator, claims, PeerType.datasource) 

369 except Unauthorized as ex: 

370 logger.error(f"Unauthorized to PUT scope: {repr(ex)}") 

371 raise 

372 

373 old_source_id = None 

374 old_clone_path = None 

375 try: 

376 old_scope = await scopes.get(scope_in.scope_id) 

377 if isinstance(old_scope.policy, GitPolicyScopeSource): 377 ↛ 397line 377 didn't jump to line 397 because the condition on line 377 was always true

378 old_source_id = GitPolicyFetcher.source_id(old_scope.policy) 

379 old_clone_path = str( 

380 GitPolicyFetcher.repo_clone_path( 

381 pathlib.Path(opal_server_config.BASE_DIR), 

382 old_scope.policy, 

383 ) 

384 ) 

385 except ScopeNotFoundError: 

386 pass # brand-new scope — nothing to repoint away from 

387 except Exception as e: 

388 # An unreadable old record must not block the overwrite that 

389 # fixes it. Its source_id is unknowable anyway, so no purge can be 

390 # published for it and nothing later will name it — the old clone 

391 # dir stays on disk until PER-15612's sweep lands. 

392 logger.warning( 

393 f"Could not read previous record for scope " 

394 f"{scope_in.scope_id}, skipping repoint purge: {e!r}" 

395 ) 

396 

397 verify_private_key_or_throw(scope_in) 

398 

399 new_source_id = ( 

400 GitPolicyFetcher.source_id(scope_in.policy) 

401 if isinstance(scope_in.policy, GitPolicyScopeSource) 

402 else None 

403 ) 

404 try: 

405 await scopes.put(scope_in) 

406 finally: 

407 # The repoint purge must stay reachable even when put() raises an 

408 # ambiguous outcome (committed server-side, error surfaced to the 

409 # client): a retry would see old_source_id == new_source_id 

410 # already (the store was updated) and never re-trigger the purge, 

411 # orphaning the old source permanently. Same channel/handlers as 

412 # delete — over-publishing self-heals, since the leader 

413 # sibling-checks and a source still shared by another scope 

414 # survives. 

415 if old_source_id is not None and old_source_id != new_source_id: 

416 await pubsub_endpoint.publish( 

417 [opal_server_config.SCOPES_PURGE_CHANNEL], 

418 ScopePurgeCommand( 

419 source_id=old_source_id, 

420 clone_path=old_clone_path, 

421 scope_id=scope_in.scope_id, 

422 reason="repoint", 

423 ).dict(), 

424 ) 

425 

426 force_fetch_str = " (force fetch)" if force_fetch else "" 

427 logger.info(f"Sync scope: {scope_in.scope_id}{force_fetch_str}") 

428 

429 # All server replicas (leaders) should sync the scope. 

430 await pubsub_endpoint.publish( 

431 opal_server_config.POLICY_REPO_WEBHOOK_TOPIC, 

432 {"scope_id": scope_in.scope_id, "force_fetch": force_fetch}, 

433 ) 

434 

435 return Response(status_code=status.HTTP_201_CREATED) 

436 

437 @router.get( 

438 "", 

439 response_model=List[Scope], 

440 response_model_exclude={"policy": {"auth"}}, 

441 ) 

442 async def get_all_scopes(*, claims: JWTClaims = Depends(authenticator)): 

443 try: 

444 require_peer_type(authenticator, claims, PeerType.datasource) 

445 except Unauthorized as ex: 

446 logger.error(f"Unauthorized to get scopes: {repr(ex)}") 

447 raise 

448 

449 return await scopes.all() 

450 

451 @router.get( 

452 "/{scope_id}", 

453 response_model=Scope, 

454 response_model_exclude={"policy": {"auth"}}, 

455 ) 

456 async def get_scope(*, scope_id: str, claims: JWTClaims = Depends(authenticator)): 

457 try: 

458 require_peer_type(authenticator, claims, PeerType.datasource) 

459 except Unauthorized as ex: 

460 logger.error(f"Unauthorized to get scope: {repr(ex)}") 

461 raise 

462 

463 try: 

464 scope = await scopes.get(scope_id) 

465 return scope 

466 except ScopeNotFoundError: 

467 raise HTTPException( 

468 status.HTTP_404_NOT_FOUND, detail=f"No such scope: {scope_id}" 

469 ) 

470 

471 @router.delete( 

472 "/{scope_id}", 

473 status_code=status.HTTP_204_NO_CONTENT, 

474 ) 

475 async def delete_scope( 

476 *, scope_id: str, claims: JWTClaims = Depends(authenticator) 

477 ): 

478 try: 

479 require_peer_type(authenticator, claims, PeerType.datasource) 

480 except Unauthorized as ex: 

481 logger.error(f"Unauthorized to delete scope: {repr(ex)}") 

482 raise 

483 

484 try: 

485 # Deletes the record and broadcasts a ScopePurgeCommand; every worker 

486 # drops its in-memory caches when the leader's confirmation broadcast 

487 # arrives. The clone dir is removed by THIS worker's best-effort 

488 # floor, not by the leader — the leader does no disk work (PER-15612). 

489 await scopes_service.delete_scope(scope_id) 

490 except ScopeNotFoundError: 

491 # Deleting a missing scope was always a silent no-op (204); keep it. 

492 pass 

493 

494 return Response(status_code=status.HTTP_204_NO_CONTENT) 

495 

496 @router.post("/{scope_id}/refresh", status_code=status.HTTP_200_OK) 

497 async def refresh_scope( 

498 scope_id: str, 

499 hinted_hash: Optional[str] = Query( 

500 None, 

501 description="Commit hash that should exist in the repo. " 

502 + "If the commit is missing from the local clone, OPAL " 

503 + "understands it as a hint that the repo should be fetched from remote.", 

504 ), 

505 claims: JWTClaims = Depends(authenticator), 

506 ): 

507 try: 

508 require_peer_type(authenticator, claims, PeerType.datasource) 

509 except Unauthorized as ex: 

510 logger.error(f"Unauthorized to delete scope: {repr(ex)}") 

511 raise 

512 

513 try: 

514 _ = await scopes.get(scope_id) 

515 

516 logger.info(f"Refresh scope: {scope_id}") 

517 

518 # If the hinted hash is None, we have no way to know whether we should 

519 # re-fetch the remote, so we force fetch, just in case. 

520 force_fetch = hinted_hash is None 

521 

522 # All server replicas (leaders) should sync the scope. 

523 await pubsub_endpoint.publish( 

524 opal_server_config.POLICY_REPO_WEBHOOK_TOPIC, 

525 { 

526 "scope_id": scope_id, 

527 "force_fetch": force_fetch, 

528 "hinted_hash": hinted_hash, 

529 }, 

530 ) 

531 

532 return Response(status_code=status.HTTP_200_OK) 

533 

534 except ScopeNotFoundError: 

535 raise HTTPException( 

536 status.HTTP_404_NOT_FOUND, detail=f"No such scope: {scope_id}" 

537 ) 

538 

539 @router.post("/refresh", status_code=status.HTTP_200_OK) 

540 async def sync_all_scopes(claims: JWTClaims = Depends(authenticator)): 

541 """Sync all scopes.""" 

542 try: 

543 require_peer_type(authenticator, claims, PeerType.datasource) 

544 except Unauthorized as ex: 

545 logger.error(f"Unauthorized to refresh all scopes: {repr(ex)}") 

546 raise 

547 

548 # All server replicas (leaders) should sync all scopes. 

549 await pubsub_endpoint.publish(opal_server_config.POLICY_REPO_WEBHOOK_TOPIC) 

550 

551 return Response(status_code=status.HTTP_200_OK) 

552 

553 @router.get( 

554 "/{scope_id}/policy", 

555 response_model=PolicyBundle, 

556 status_code=status.HTTP_200_OK, 

557 dependencies=[Depends(_allowed_scoped_authenticator)], 

558 ) 

559 async def get_scope_policy( 

560 *, 

561 request: Request, 

562 scope_id: str = Path(..., title="Scope ID"), 

563 base_hash: Optional[str] = Query( 

564 None, 

565 description="hash of previous bundle already downloaded, server will return a diff bundle.", 

566 ), 

567 ): 

568 try: 

569 scope = await scopes.get(scope_id) 

570 except ScopeNotFoundError: 

571 logger.warning( 

572 "Requested scope {scope_id} not found, returning default scope", 

573 scope_id=scope_id, 

574 ) 

575 return await _generate_default_scope_bundle(scope_id, request) 

576 

577 if not isinstance(scope.policy, GitPolicyScopeSource): 

578 raise HTTPException( 

579 status.HTTP_501_NOT_IMPLEMENTED, 

580 detail=f"policy source is not yet implemented: {scope_id}", 

581 ) 

582 

583 fetcher = GitPolicyFetcher( 

584 pathlib.Path(opal_server_config.BASE_DIR), 

585 scope.scope_id, 

586 cast(GitPolicyScopeSource, scope.policy), 

587 ) 

588 

589 try: 

590 return await _make_bundle_waiting_for_clone( 

591 fetcher, base_hash, scope_id, request 

592 ) 

593 except CloneNotPopulatedError as exc: 

594 # The clone has no refs/remotes/<remote>/* at all, so it is being 

595 # populated right now — _clone() rmtree's the destination and clones 

596 # INTO the final path, so this window is the whole clone. Telling a 

597 # client its configuration is permanently wrong during the recovery 

598 # that fixes it is the opposite of the truth. 

599 # 

600 # Derived from disk, NOT from the in-flight marker: that marker is a 

601 # per-process global written only by the leader's sync, while this 

602 # route is served by any worker, so keying on it answered 409 on 

603 # every non-leader — N-1 of N workers. 

604 # 

605 # Reached once the wait budget, if any, is spent — or straight 

606 # away when the wait is disabled, shed at the in-flight cap, or the 

607 # client has already hung up. When a budget was spent, this clone is 

608 # slower than a client's whole retry budget, so the 503 is a report 

609 # that waiting did not help rather than a first reflex. 

610 # 

611 # Both the line and the event are skipped when the caller has 

612 # already hung up: this 503 is shaped for a socket nobody is 

613 # reading, so counting it would inflate the very rate an operator 

614 # watches to decide whether clients are being served — and the 

615 # wait has already logged the abandonment once, with the hold. 

616 if not exc.client_disconnected: 

617 logger.info( 

618 "Scope {scope_id} clone is not populated yet ({exc!r}), " 

619 "returning 503 after waiting {waited:.1f}s", 

620 scope_id=scope_id, 

621 exc=exc, 

622 waited=exc.waited_seconds, 

623 ) 

624 metrics.event( 

625 "ScopePolicyUnavailable", 

626 message=f"Scope {scope_id} policy 503 (clone in progress)", 

627 tags={"scope_id": scope_id, "status": "503", "retryable": "true"}, 

628 ) 

629 raise HTTPException( 

630 status.HTTP_503_SERVICE_UNAVAILABLE, 

631 detail=( 

632 f"Policy clone for scope {scope_id} is being created, " 

633 "retry shortly" 

634 ), 

635 headers={"Retry-After": _RETRY_AFTER_CLONE_IN_PROGRESS}, 

636 ) 

637 except BranchHeadNotFoundError as exc: 

638 logger.error( 

639 "Scope {scope_id} bundle unavailable: {exc!r} (non-retryable)", 

640 scope_id=scope_id, 

641 exc=exc, 

642 ) 

643 metrics.event( 

644 "ScopePolicyUnavailable", 

645 message=f"Scope {scope_id} policy 409 (branch unresolved)", 

646 tags={"scope_id": scope_id, "status": "409", "retryable": "false"}, 

647 ) 

648 raise HTTPException( 

649 status.HTTP_409_CONFLICT, 

650 detail=( 

651 f"Policy branch for scope {scope_id} could not be resolved " 

652 "(check the configured branch); not retryable" 

653 ), 

654 ) 

655 except ( 

656 InvalidGitRepositoryError, 

657 # A concurrent delete/recovery can rmtree the clone dir before 

658 # Repo() opens it (NoSuchPathError), or mid-tree-walk (raw 

659 # OSError). The record exists, so this is transient: recovery or 

660 # the next sync re-creates the clone. Serving the default 

661 # scope's bundle here would hand a live tenant another tenant's 

662 # policy — tell the client to retry instead. 

663 NoSuchPathError, 

664 pygit2.GitError, 

665 ValueError, 

666 OSError, 

667 ) as exc: 

668 logger.warning( 

669 "Scope {scope_id} is live but its clone is unavailable ({exc!r}), " 

670 "returning 503", 

671 scope_id=scope_id, 

672 exc=exc, 

673 ) 

674 metrics.event( 

675 "ScopePolicyUnavailable", 

676 message=f"Scope {scope_id} policy 503 (clone unavailable)", 

677 tags={"scope_id": scope_id, "status": "503", "retryable": "true"}, 

678 ) 

679 raise HTTPException( 

680 status.HTTP_503_SERVICE_UNAVAILABLE, 

681 detail=( 

682 f"Policy clone for scope {scope_id} is temporarily " 

683 "unavailable, retry shortly" 

684 ), 

685 headers={"Retry-After": _RETRY_AFTER_CLONE_UNAVAILABLE}, 

686 ) 

687 

688 async def _generate_default_scope_bundle( 

689 scope_id: str, request: Optional[Request] = None 

690 ) -> PolicyBundle: 

691 metrics.event( 

692 "ScopeNotFound", 

693 message=f"Scope {scope_id} not found. Serving default scope instead", 

694 tags={"scope_id": scope_id}, 

695 ) 

696 

697 try: 

698 scope = await scopes.get("default") 

699 fetcher = GitPolicyFetcher( 

700 pathlib.Path(opal_server_config.BASE_DIR), 

701 scope.scope_id, 

702 cast(GitPolicyScopeSource, scope.policy), 

703 ) 

704 # run_sync, like the primary path at the top of this route. Without 

705 # it a full bundle build — open the repo, walk the commit tree, read 

706 # and encode every matching file — runs ON THE EVENT LOOP, stalling 

707 # every other request this worker is serving, including other 

708 # tenants' bundles and the pub/sub websocket traffic. Reached by any 

709 # GET for an unknown scope, which a PDP with a stale id re-hits on 

710 # its normal poll cadence. 

711 # 

712 # Waits for the default clone on the same terms as the primary 

713 # path: this is the branch every PDP holding a stale scope id 

714 # takes, and the default scope's clone is populated by the same 

715 # recovery as any other. On expiry the re-raised 

716 # CloneNotPopulatedError lands in the broad tuple below, so this 

717 # path keeps ITS contract (Retry-After 5), not the primary path's. 

718 return await _make_bundle_waiting_for_clone( 

719 fetcher, None, scope.scope_id, request 

720 ) 

721 except ScopeNotFoundError: 

722 # 404, not a bare ScopeNotFoundError. Nothing registers an exception 

723 # handler for that, so it escaped the route as an unhandled 500 — 

724 # for the ordinary case of an unknown scope on a deployment that has 

725 # no "default" scope at all. get_scope and refresh_scope already 

726 # answer 404 here; this now matches them. 

727 raise HTTPException( 

728 status.HTTP_404_NOT_FOUND, detail=f"No such scope: {scope_id}" 

729 ) 

730 except CloneNotPopulatedError as exc: 

731 # Its own arm, ahead of the broad tuple that would otherwise catch 

732 # it (CloneNotPopulatedError subclasses ValueError), for one 

733 # reason: only this exception carries the hold, and a 503 that does 

734 # not say how long the server waited cannot be told apart from one 

735 # that never waited. Same 503 + Retry-After 5 as the tuple below — 

736 # this path's contract is unchanged. Silent when the caller has 

737 # already hung up, like the primary path. 

738 if not exc.client_disconnected: 

739 logger.warning( 

740 "Default-scope bundle for {scope_id} is temporarily " 

741 "unavailable after waiting {waited:.1f}s ({exc!r}), " 

742 "returning 503", 

743 scope_id=scope_id, 

744 waited=exc.waited_seconds, 

745 exc=exc, 

746 ) 

747 raise HTTPException( 

748 status.HTTP_503_SERVICE_UNAVAILABLE, 

749 detail=( 

750 f"Policy clone for scope {scope_id} is temporarily " 

751 "unavailable, retry shortly" 

752 ), 

753 headers={"Retry-After": _RETRY_AFTER_CLONE_UNAVAILABLE}, 

754 ) 

755 except ( 

756 InvalidGitRepositoryError, 

757 NoSuchPathError, 

758 pygit2.GitError, 

759 OSError, 

760 ValueError, 

761 ) as exc: 

762 # A TRANSIENT fault building the default scope's bundle is not 

763 # "no such scope". These are the same exceptions the primary path 

764 # answers with 503 forty lines up, on the same reasoning: the clone 

765 # is being recovered and will be back. Folding them into the 404 

766 # told a client to stop asking about a condition that self-heals in 

767 # seconds — and §6 explicitly tells third-party consumers to act on 

768 # these codes, so it was wrong in the unsafe direction. 

769 logger.warning( 

770 "Default-scope bundle for {scope_id} is temporarily unavailable " 

771 "({exc!r}), returning 503", 

772 scope_id=scope_id, 

773 exc=exc, 

774 ) 

775 raise HTTPException( 

776 status.HTTP_503_SERVICE_UNAVAILABLE, 

777 detail=( 

778 f"Policy clone for scope {scope_id} is temporarily " 

779 "unavailable, retry shortly" 

780 ), 

781 headers={"Retry-After": _RETRY_AFTER_CLONE_UNAVAILABLE}, 

782 ) 

783 

784 @router.get( 

785 "/{scope_id}/data", 

786 response_model=DataSourceConfig, 

787 status_code=status.HTTP_200_OK, 

788 dependencies=[Depends(_allowed_scoped_authenticator)], 

789 ) 

790 async def get_scope_data_config( 

791 *, 

792 scope_id: str = Path(..., title="Scope ID"), 

793 authorization: Optional[str] = Header(None), 

794 ): 

795 logger.info( 

796 "Serving source configuration for scope {scope_id}", scope_id=scope_id 

797 ) 

798 try: 

799 scope = await scopes.get(scope_id) 

800 return scope.data 

801 except ScopeNotFoundError as ex: 

802 logger.warning( 

803 "Requested scope {scope_id} not found, returning OPAL_DATA_CONFIG_SOURCES", 

804 scope_id=scope_id, 

805 ) 

806 try: 

807 config: ServerDataSourceConfig = opal_server_config.DATA_CONFIG_SOURCES 

808 

809 if config.external_source_url: 

810 url = str(config.external_source_url) 

811 token = get_token_from_header(authorization) 

812 redirect_url = set_url_query_param(url, "token", token) 

813 return RedirectResponse(url=redirect_url) 

814 else: 

815 return config.config 

816 except ScopeNotFoundError: 

817 raise HTTPException(status.HTTP_404_NOT_FOUND, detail=str(ex)) 

818 

819 @router.post("/{scope_id}/data/update") 

820 async def publish_data_update_event( 

821 update: DataUpdate, 

822 claims: JWTClaims = Depends(authenticator), 

823 scope_id: str = Path(..., description="Scope ID"), 

824 ): 

825 try: 

826 require_peer_type(authenticator, claims, PeerType.datasource) 

827 

828 restrict_optional_topics_to_publish(authenticator, claims, update) 

829 

830 for entry in update.entries: 

831 entry.topics = [f"data:{topic}" for topic in entry.topics] 

832 

833 await DataUpdatePublisher( 

834 ScopedServerSideTopicPublisher(pubsub_endpoint, scope_id) 

835 ).publish_data_updates(update) 

836 except Unauthorized as ex: 

837 logger.error(f"Unauthorized to publish update: {repr(ex)}") 

838 raise 

839 

840 return router