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
« 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
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
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"
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
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
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
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
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"
94def _bounded_clone_wait() -> float:
95 """The configured hold, validated and clamped. 0.0 means "do not wait".
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
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
131def _publish_clone_wait_inflight() -> None:
132 """Publish the held-request count for this process.
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 )
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.
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.
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.
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.
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.
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.
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
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
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
214 loop = asyncio.get_running_loop()
215 started = loop.time()
216 deadline = started + wait
217 outcome = "error"
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))
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
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
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 )
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
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 )
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 )
338def init_scope_router(
339 scopes: ScopeRepository,
340 authenticator: JWTAuthenticator,
341 pubsub_endpoint: PubSubEndpoint,
342 scopes_service: ScopesService,
343):
344 router = APIRouter()
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
352 allowed_scopes = claims.get("allowed_scopes")
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)
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
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 )
397 verify_private_key_or_throw(scope_in)
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 )
426 force_fetch_str = " (force fetch)" if force_fetch else ""
427 logger.info(f"Sync scope: {scope_in.scope_id}{force_fetch_str}")
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 )
435 return Response(status_code=status.HTTP_201_CREATED)
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
449 return await scopes.all()
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
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 )
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
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
494 return Response(status_code=status.HTTP_204_NO_CONTENT)
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
513 try:
514 _ = await scopes.get(scope_id)
516 logger.info(f"Refresh scope: {scope_id}")
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
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 )
532 return Response(status_code=status.HTTP_200_OK)
534 except ScopeNotFoundError:
535 raise HTTPException(
536 status.HTTP_404_NOT_FOUND, detail=f"No such scope: {scope_id}"
537 )
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
548 # All server replicas (leaders) should sync all scopes.
549 await pubsub_endpoint.publish(opal_server_config.POLICY_REPO_WEBHOOK_TOPIC)
551 return Response(status_code=status.HTTP_200_OK)
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)
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 )
583 fetcher = GitPolicyFetcher(
584 pathlib.Path(opal_server_config.BASE_DIR),
585 scope.scope_id,
586 cast(GitPolicyScopeSource, scope.policy),
587 )
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 )
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 )
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 )
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
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))
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)
828 restrict_optional_topics_to_publish(authenticator, claims, update)
830 for entry in update.entries:
831 entry.topics = [f"data:{topic}" for topic in entry.topics]
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
840 return router