Coverage for /usr/local/lib/python3.10/site-packages/opal_server-0.0.0-py3.10.egg/opal_server/server.py: 67%
241 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 os
3import signal
4import sys
5import traceback
6from functools import partial
7from pathlib import Path
8from typing import List, Optional
10from fastapi import Depends, FastAPI, Request
11from fastapi.openapi.docs import get_redoc_html
12from fastapi.responses import HTMLResponse, JSONResponse
13from fastapi_websocket_pubsub.event_broadcaster import EventBroadcasterContextManager
14from opal_common.authentication.deps import JWTAuthenticator, StaticBearerAuthenticator
15from opal_common.authentication.signer import JWTSigner
16from opal_common.confi.confi import load_conf_if_none
17from opal_common.config import opal_common_config
18from opal_common.logger import configure_logs, logger
19from opal_common.middleware import configure_middleware
20from opal_common.monitoring import apm, metrics
21from opal_common.schemas.data import ServerDataSourceConfig
22from opal_common.synchronization.named_lock import NamedLock
23from opal_common.topics.publisher import (
24 PeriodicPublisher,
25 ServerSideTopicPublisher,
26 TopicPublisher,
27)
28from opal_server.config import opal_server_config
29from opal_server.data.api import init_data_updates_router
30from opal_server.data.data_update_publisher import DataUpdatePublisher
31from opal_server.debug_stats import register_internal_stats_route
32from opal_server.loadlimiting import init_loadlimit_router
33from opal_server.metrics_setup import configure_server_metrics
34from opal_server.policy.bundles.api import router as bundles_router
35from opal_server.policy.watcher.factory import setup_watcher_task
36from opal_server.policy.watcher.task import PolicyWatcherTask
37from opal_server.policy.webhook.api import init_git_webhook_router
38from opal_server.publisher import setup_broadcaster_keepalive_task
39from opal_server.pubsub import PubSub
40from opal_server.pubsub_resilience import ReconnectingBroadcaster
41from opal_server.redis_utils import RedisDB
42from opal_server.scopes.api import init_scope_router
43from opal_server.scopes.loader import load_scopes
44from opal_server.scopes.purge import subscribe_worker_purge_handler
45from opal_server.scopes.scope_repository import ScopeRepository
46from opal_server.scopes.service import ScopesService
47from opal_server.security.api import init_security_router
48from opal_server.security.jwks import JwksStaticEndpoint
49from opal_server.statistics import OpalStatistics, init_statistics_router
51# Upper bound on the shutdown drain of the DELETE floor's clone purges. Same
52# rationale and same order of magnitude as the watcher's _PURGE_DRAIN_TIMEOUT:
53# stop() runs while k8s's terminationGracePeriodSeconds (30s default) is
54# counting down, so blocking longer just converts a clean exit into a SIGKILL.
55_SCOPES_DRAIN_TIMEOUT = 5.0
58class OpalServer:
59 def __init__(
60 self,
61 init_policy_watcher: bool = None,
62 policy_remote_url: str = None,
63 init_publisher: bool = None,
64 data_sources_config: Optional[ServerDataSourceConfig] = None,
65 broadcaster_uri: str = None,
66 signer: Optional[JWTSigner] = None,
67 enable_jwks_endpoint=True,
68 jwks_url: str = None,
69 jwks_static_dir: str = None,
70 master_token: str = None,
71 loadlimit_notation: str = None,
72 ) -> None:
73 """
74 Args:
75 policy_remote_url (str, optional): the url of the repo watched by policy watcher.
76 init_publisher (bool, optional): whether or not to launch a publisher pub/sub client.
77 this publisher is used by the server processes to publish data to the client.
78 data_sources_config (ServerDataSourceConfig, optional): base data configuration, that opal
79 clients should get the data from.
80 broadcaster_uri (str, optional): Which server/medium should the PubSub use for broadcasting.
81 Defaults to BROADCAST_URI.
82 loadlimit_notation (str, optional): Rate limit configuration for opal client connections.
83 Defaults to None, in that case no rate limit is enforced
85 The server can run in multiple workers (by gunicorn or uvicorn).
87 Every worker of the server launches the following internal components:
88 publisher (PubSubClient): a client that is used to publish updates to the client.
89 data_update_publisher (DataUpdatePublisher): a specialized publisher for data updates.
91 Besides the components above, the works are also deciding among themselves
92 on a *leader* worker (the first worker to obtain a file-lock) that also
93 launches the following internal components:
95 watcher (PolicyWatcherTask): run by the leader, monitors the policy git repository
96 by polling on it or by being triggered by the callback subscribed on the "webhook"
97 topic. upon being triggered, will detect updates to the policy (new commits) and
98 will update the opal client via pubsub.
99 """
100 # load defaults
101 init_publisher: bool = load_conf_if_none(
102 init_publisher, opal_server_config.PUBLISHER_ENABLED
103 )
104 broadcaster_uri: str = load_conf_if_none(
105 broadcaster_uri, opal_server_config.BROADCAST_URI
106 )
107 jwks_url: str = load_conf_if_none(jwks_url, opal_server_config.AUTH_JWKS_URL)
108 jwks_static_dir: str = load_conf_if_none(
109 jwks_static_dir, opal_server_config.AUTH_JWKS_STATIC_DIR
110 )
111 master_token: str = load_conf_if_none(
112 master_token, opal_server_config.AUTH_MASTER_TOKEN
113 )
114 self._init_policy_watcher: bool = load_conf_if_none(
115 init_policy_watcher, opal_server_config.REPO_WATCHER_ENABLED
116 )
117 self.loadlimit_notation: str = load_conf_if_none(
118 loadlimit_notation, opal_server_config.CLIENT_LOAD_LIMIT_NOTATION
119 )
120 self._policy_remote_url = policy_remote_url
122 self._configure_monitoring()
123 metrics.increment("startup")
125 self.data_sources_config: ServerDataSourceConfig = (
126 data_sources_config
127 if data_sources_config is not None
128 else opal_server_config.DATA_CONFIG_SOURCES
129 )
131 self.broadcaster_uri = broadcaster_uri
132 self.master_token = master_token
134 if signer is not None: 134 ↛ 135line 134 didn't jump to line 135 because the condition on line 134 was never true
135 self.signer = signer
136 else:
137 self.signer = JWTSigner(
138 private_key=opal_server_config.AUTH_PRIVATE_KEY,
139 public_key=opal_common_config.AUTH_PUBLIC_KEY,
140 algorithm=opal_common_config.AUTH_JWT_ALGORITHM,
141 audience=opal_common_config.AUTH_JWT_AUDIENCE,
142 issuer=opal_common_config.AUTH_JWT_ISSUER,
143 )
144 if self.signer.enabled: 144 ↛ 149line 144 didn't jump to line 149 because the condition on line 144 was always true
145 logger.info(
146 "OPAL is running in secure mode - will verify API requests with JWT tokens."
147 )
148 else:
149 logger.info(
150 "OPAL was not provided with JWT encryption keys, cannot verify api requests!"
151 )
153 if enable_jwks_endpoint: 153 ↛ 158line 153 didn't jump to line 158 because the condition on line 153 was always true
154 self.jwks_endpoint = JwksStaticEndpoint(
155 signer=self.signer, jwks_url=jwks_url, jwks_static_dir=jwks_static_dir
156 )
157 else:
158 self.jwks_endpoint = None
160 self.pubsub = PubSub(signer=self.signer, broadcaster_uri=broadcaster_uri)
161 self._wire_broadcaster_give_up()
163 self.publisher: Optional[TopicPublisher] = None
164 self.broadcast_keepalive: Optional[PeriodicPublisher] = None
165 if init_publisher: 165 ↛ 178line 165 didn't jump to line 178 because the condition on line 165 was always true
166 self.publisher = ServerSideTopicPublisher(self.pubsub.endpoint)
168 if ( 168 ↛ 172line 168 didn't jump to line 172 because the condition on line 168 was never true
169 opal_server_config.BROADCAST_KEEPALIVE_INTERVAL > 0
170 and self.broadcaster_uri is not None
171 ):
172 self.broadcast_keepalive = setup_broadcaster_keepalive_task(
173 self.publisher,
174 time_interval=opal_server_config.BROADCAST_KEEPALIVE_INTERVAL,
175 topic=opal_server_config.BROADCAST_KEEPALIVE_TOPIC,
176 )
178 if opal_common_config.STATISTICS_ENABLED: 178 ↛ 181line 178 didn't jump to line 181 because the condition on line 178 was always true
179 self.opal_statistics = OpalStatistics(self.pubsub.endpoint)
180 else:
181 self.opal_statistics = None
183 # A worker's backbone READER (EventBroadcaster listening context) is
184 # otherwise only entered while a WebSocket client is connected to that
185 # worker: server-side subscriptions on a worker with zero clients hear
186 # nothing from the fleet. Two features need every worker listening for
187 # its own sake, not its clients':
188 # - statistics: workers synchronise their own view over the backbone
189 # - scopes: the fleet purge broadcast (scopes/purge.py) must reach
190 # EVERY worker so it drops its GitPolicyFetcher caches for a
191 # deleted/repointed source. Measured on staging (~23 WS conns over
192 # 16 workers): 5 of 8 workers per pod never received a single purge
193 # and kept stale repo_locks entries for the life of the process.
194 # So the "global" listening context is held whenever either is on and a
195 # broadcaster exists (single-process deployments have nothing to read).
196 self.broadcast_listening_context: Optional[
197 EventBroadcasterContextManager
198 ] = None
199 if self.broadcaster_uri is not None: 199 ↛ 212line 199 didn't jump to line 212 because the condition on line 199 was never true
200 # Legacy EventBroadcaster (BROADCAST_RECONNECT_ENABLED=false): its
201 # reader connects EAGERLY in __aenter__ and re-raises, so a backbone
202 # that is down at boot would abort the background task before the
203 # purge subscription and the leadership lock. The reconnecting
204 # broadcaster connects lazily and never raises here, so the scopes
205 # reason is honoured only for it.
206 # Rollout cost (Postgres backbone): entering the context makes THIS
207 # worker open the broadcaster pool - asyncpg's default min_size is
208 # 10 eager connections at boot, decaying to ~1 after ~300 s
209 # (BROADCASTER_PG_MAX_POOL_SIZE can only raise the ceiling, not
210 # lower the floor). Budget the broadcast DB's max_connections for
211 # 10 x workers x pods at a rolling restart, not the steady state.
212 wants_reader_for_scopes = bool(opal_server_config.SCOPES) and isinstance(
213 self.pubsub.broadcaster, ReconnectingBroadcaster
214 )
215 if opal_common_config.STATISTICS_ENABLED or wants_reader_for_scopes:
216 self.broadcast_listening_context = (
217 self.pubsub.endpoint.broadcaster.get_listening_context()
218 )
219 if opal_server_config.SCOPES and self.broadcast_listening_context is None:
220 # Accurate in all four combinations: statistics on the legacy
221 # broadcaster DOES arm the context (for its own reason), and
222 # then purges are delivered too.
223 logger.info(
224 "OPAL_SCOPES is on but BROADCAST_RECONNECT_ENABLED is off: "
225 "workers without a WebSocket client keep no backbone reader, "
226 "so fleet purge delivery to them is not guaranteed"
227 )
229 self.watcher: PolicyWatcherTask = None
230 self.leadership_lock: Optional[NamedLock] = None
232 if opal_server_config.SCOPES: 232 ↛ 263line 232 didn't jump to line 263 because the condition on line 232 was always true
233 self._redis_db = RedisDB(opal_server_config.REDIS_URL)
234 self._scopes = ScopeRepository(self._redis_db)
235 logger.info("OPAL Scopes: server is connected to scopes repository")
236 if not self._init_policy_watcher: 236 ↛ 250line 236 didn't jump to line 250 because the condition on line 236 was never true
237 # The flag's name reads as "single-repo watcher", but in scopes
238 # mode it gates setup_watcher_task() -> ScopesPolicyWatcherTask:
239 # the periodic sync_scopes pass, the leader's post-leadership
240 # sync-all and the fleet purger. It does NOT gate the gunicorn
241 # master's pre-fork preload (gunicorn_conf.py when_ready ->
242 # preload_scopes), which clones/fetches every registered scope
243 # on every boot regardless - the flag cannot mean "no git
244 # activity", only "no ongoing sync". With it off the leader
245 # parks on the keepalive forever and the pods still report
246 # Ready. A WARNING rather than a hard refusal: an existing
247 # deployment must not stop booting on upgrade over a flag it may
248 # have set deliberately (e.g. a read-only replica behind another
249 # OPAL that does the syncing).
250 logger.warning(
251 "OPAL_REPO_WATCHER_ENABLED is off while OPAL_SCOPES is on: the "
252 "leader worker will never SYNC or PURGE scopes (no periodic "
253 "sync_scopes pass, no post-leadership sync-all, no fleet "
254 "purge). NOTE the gunicorn master's PRE-FORK PRELOAD still "
255 "clones/fetches every registered scope on every boot - this "
256 "flag does not gate it, so git traffic and clone-tree growth "
257 "at boot are expected regardless. If ongoing sync is "
258 "intended, set OPAL_REPO_WATCHER_ENABLED=true."
259 )
261 # Set BEFORE _init_fast_api_app(): _configure_api_routes assigns it,
262 # and it must exist even when SCOPES is off (shutdown reads it).
263 self._scopes_service: Optional[ScopesService] = None
265 # init fastapi app
266 self.app: FastAPI = self._init_fast_api_app()
268 def _init_fast_api_app(self):
269 """Inits the fastapi app object."""
270 app = FastAPI(
271 title="Opal Server",
272 description="OPAL is an administration layer for Open Policy Agent (OPA), detecting changes"
273 + " to both policy and data and pushing live updates to your agents. The opal server creates"
274 + " a pub/sub channel clients can subscribe to (i.e: acts as coordinator). The server also"
275 + " tracks a git repository (via webhook) for updates to policy (or static data) and accepts"
276 + " continuous data update notifications via REST api, which are then pushed to clients.",
277 version="0.1.0",
278 # disabled because of an issue with redoc_js_url in fastapi versions prior to 0.115.13
279 redoc_url=None,
280 )
282 configure_middleware(app)
283 self._configure_api_routes(app)
284 self._configure_lifecycle_callbacks(app)
286 return app
288 def _configure_monitoring(self):
289 configure_logs()
291 apm.configure_apm(opal_server_config.ENABLE_DATADOG_APM, "opal-server")
293 configure_server_metrics()
295 def _configure_api_routes(self, app: FastAPI):
296 """Mounts the api routes on the app object."""
297 authenticator = JWTAuthenticator(self.signer)
299 data_update_publisher: Optional[DataUpdatePublisher] = None
300 if self.publisher is not None: 300 ↛ 304line 300 didn't jump to line 304 because the condition on line 300 was always true
301 data_update_publisher = DataUpdatePublisher(self.publisher)
303 # Init api routers with required dependencies
304 data_updates_router = init_data_updates_router(
305 data_update_publisher, self.data_sources_config, authenticator
306 )
307 webhook_router = init_git_webhook_router(self.pubsub.endpoint, authenticator)
308 security_router = init_security_router(
309 self.signer, StaticBearerAuthenticator(self.master_token)
310 )
311 statistics_router = init_statistics_router(self.opal_statistics)
312 loadlimit_router = init_loadlimit_router(self.loadlimit_notation)
314 # mount the api routes on the app object
315 app.include_router(
316 bundles_router,
317 tags=["Bundle Server"],
318 dependencies=[Depends(authenticator)],
319 )
320 app.include_router(data_updates_router, tags=["Data Updates"])
321 app.include_router(webhook_router, tags=["Github Webhook"])
322 app.include_router(security_router, tags=["Security"])
323 app.include_router(self.pubsub.pubsub_router, tags=["Pub/Sub"])
324 app.include_router(
325 self.pubsub.api_router,
326 tags=["Pub/Sub"],
327 dependencies=[Depends(authenticator)],
328 )
329 app.include_router(
330 statistics_router,
331 tags=["Server Statistics"],
332 dependencies=[Depends(authenticator)],
333 )
334 app.include_router(
335 loadlimit_router,
336 tags=["Client Load Limiting"],
337 dependencies=[Depends(authenticator)],
338 )
340 if opal_server_config.SCOPES: 340 ↛ 362line 340 didn't jump to line 362 because the condition on line 340 was always true
341 scopes_service = ScopesService(
342 base_dir=Path(opal_server_config.BASE_DIR),
343 scopes=self._scopes,
344 pubsub_endpoint=self.pubsub.endpoint,
345 )
346 # Held so shutdown can drain it. THIS is the instance DELETE runs on
347 # (api.py -> scopes_service.delete_scope), so this is the one whose
348 # _local_purges accumulates the backgrounded clone-dir purges. The
349 # watcher builds its own separate ScopesService and only ever syncs
350 # with it — draining that one drains an empty set. The watcher is
351 # also leader-only, and a DELETE usually lands on a non-leader, so
352 # it could not be the drain point even if it shared the object.
353 self._scopes_service = scopes_service
354 app.include_router(
355 init_scope_router(
356 self._scopes, authenticator, self.pubsub.endpoint, scopes_service
357 ),
358 tags=["Scopes"],
359 prefix="/scopes",
360 )
362 if self.jwks_endpoint is not None: 362 ↛ 366line 362 didn't jump to line 366 because the condition on line 362 was always true
363 # mount jwts (static) route
364 self.jwks_endpoint.configure_app(app)
366 @app.get("/redoc", include_in_schema=False)
367 async def redoc_html(req: Request) -> HTMLResponse:
368 root_path = req.scope.get("root_path", "").rstrip("/")
369 openapi_url = root_path + app.openapi_url
370 return get_redoc_html(
371 openapi_url=openapi_url,
372 title=app.title + " - ReDoc",
373 redoc_js_url="https://cdn.jsdelivr.net/npm/redoc@2/bundles/redoc.standalone.js",
374 )
376 # top level routes (i.e: healthchecks)
377 @app.get("/", include_in_schema=False)
378 def root():
379 # Liveness / process-up: trivially ok as long as the process serves.
380 return {"status": "ok"}
382 @app.get("/healthcheck", include_in_schema=False)
383 def healthcheck():
384 # Readiness: also reflect broadcaster-reader health so a wedged reader
385 # (reader task absent/dead while clients depend on it) fails the probe
386 # and k8s can route away from / restart this worker. Stays ok through a
387 # normal transient reconnect (see ReconnectingBroadcaster.is_reader_healthy).
388 broadcaster = self.pubsub.broadcaster
389 if opal_server_config.BROADCAST_HEALTHCHECK_ENABLED and isinstance(
390 broadcaster, ReconnectingBroadcaster
391 ):
392 healthy = broadcaster.is_reader_healthy()
393 # Publish what the probe already decided. A wedged reader is
394 # otherwise invisible outside this handler: staging runs no
395 # liveness probe, so nothing acts on the 503 below.
396 metrics.gauge(
397 "opal_server.broadcaster_reader_healthy",
398 1 if healthy else 0,
399 tags={"pid": str(os.getpid())},
400 )
401 if not healthy:
402 return JSONResponse(
403 status_code=503,
404 content={"status": "error", "broadcaster": "unhealthy"},
405 )
406 return {"status": "ok"}
408 register_internal_stats_route(
409 app,
410 enabled=opal_server_config.DEBUG_INTERNAL_STATS,
411 dependencies=[Depends(authenticator)],
412 )
414 return app
416 def _configure_lifecycle_callbacks(self, app: FastAPI):
417 """Registers callbacks on app startup and shutdown.
419 on app startup we launch our long running processes (async
420 tasks) on the event loop. on app shutdown we stop these long
421 running tasks.
422 """
424 @app.on_event("startup")
425 async def startup_event():
426 logger.info("*** OPAL Server Startup ***")
428 try:
429 self._task = asyncio.create_task(self.start_server_background_tasks())
431 except Exception:
432 logger.critical("Exception while starting OPAL")
433 traceback.print_exc()
435 sys.exit(1)
437 @app.on_event("shutdown")
438 async def shutdown_event():
439 logger.info("triggered shutdown event")
440 await self.stop_server_background_tasks()
442 return app
444 async def start_server_background_tasks(self):
445 """Starts the background processes (as asyncio tasks) if such are
446 configured.
448 all workers will start these tasks:
449 - publisher: a client that is used to publish updates to the client.
451 only the leader worker (first to obtain leadership lock) will start these tasks:
452 - (repo) watcher: monitors the policy git repository for changes.
453 """
454 if self.publisher is not None: 454 ↛ exitline 454 didn't return from function 'start_server_background_tasks' because the condition on line 454 was always true
455 async with self.publisher:
456 if self.broadcast_listening_context is not None: 456 ↛ 460line 456 didn't jump to line 460 because the condition on line 456 was never true
457 # Entered ONCE per worker, whatever the reason(s) it exists for
458 # (statistics and/or scopes — see __init__). Entering it twice
459 # would double the listener count and leak a reader on exit.
460 logger.info(
461 "listening on the broadcast channel on this worker "
462 "(statistics={stats}, scopes={scopes})",
463 stats=self.opal_statistics is not None,
464 scopes=bool(opal_server_config.SCOPES),
465 )
466 try:
467 await self.broadcast_listening_context.__aenter__()
468 except (
469 Exception
470 ) as exc: # noqa: BLE001 — everything below must still run
471 # Belt and braces for the eager-connect case the isinstance
472 # gate above should already exclude: an unreadable backbone
473 # must not cost this worker its purge subscription, the
474 # leadership lock and the watcher. Leave the context None
475 # so the exit path does not touch a half-entered manager.
476 logger.warning(
477 "Could not start listening on the broadcast channel on "
478 "this worker ({etype}: {err}); continuing without a "
479 "global reader — fleet purges reach this worker only "
480 "while it has a WebSocket client",
481 etype=type(exc).__name__,
482 err=exc,
483 )
484 # UNWIND before dropping the reference: the library's
485 # __aenter__ increments _listen_count BEFORE it starts the
486 # reader, so a raise leaves the count at 1. Left there,
487 # every later client context takes it to 2, 3, ... and
488 # start_reader_task() never fires again — that worker
489 # would serve clients that never receive a broadcast,
490 # forever, with /healthcheck 200. __aexit__ decrements to
491 # 0, skips the cancel (no task) and swallows its own errors.
492 try:
493 await self.broadcast_listening_context.__aexit__(
494 None, None, None
495 )
496 except Exception: # noqa: BLE001 — best effort, already failing
497 logger.exception(
498 "Could not unwind the broadcast listening context "
499 "after a failed enter"
500 )
501 self.broadcast_listening_context = None
502 if (
503 self.broadcast_listening_context is not None
504 and self.opal_statistics is not None
505 ):
506 # If the broadcast channel DIES, statistics on this worker
507 # can't be trusted anymore - restart. Guarded against the
508 # exit path's own clean cancellation (__aexit__ cancels the
509 # reader): on master the zero-arg __aexit__ raised TypeError
510 # before the reader was ever cancelled, so this callback
511 # never saw a cancelled task - firing on one now would be a
512 # NEW self-SIGTERM at the end of every fall-through boot
513 # (statistics on + keepalive off + watcher off = a restart
514 # loop of the leader slot). Give-up (the task RETURNING) is
515 # also covered fleet-wide by _wire_broadcaster_give_up; the
516 # duplicate SIGTERM in that overlap is a harmless no-op.
517 def _restart_if_reader_died(task):
518 if not task.cancelled():
519 self._graceful_shutdown()
521 self.broadcast_listening_context._event_broadcaster.get_reader_task().add_done_callback(
522 _restart_if_reader_died
523 )
524 # No done-callback for the scopes reason: the context exists
525 # for scopes only on a ReconnectingBroadcaster (see __init__),
526 # whose reader never completes on its own except by GIVING UP
527 # after BROADCAST_RECONNECT_MAX_RETRIES — and that case already
528 # restarts every worker through _wire_broadcaster_give_up
529 # (fires on give-up, never on clean cancellation). The legacy
530 # broadcaster is skipped for scopes altogether.
531 if self.opal_statistics is not None: 531 ↛ 537line 531 didn't jump to line 537 because the condition on line 531 was always true
532 asyncio.create_task(self.opal_statistics.run())
533 self.pubsub.endpoint.notifier.register_unsubscribe_event(
534 self.opal_statistics.remove_client
535 )
537 if opal_server_config.SCOPES: 537 ↛ 547line 537 didn't jump to line 547 because the condition on line 537 was always true
538 # Every worker (leader or not) must drop its in-memory
539 # GitPolicyFetcher caches when a scope is deleted/repointed
540 # anywhere in the fleet. Subscribed before the leadership
541 # lock on purpose: non-leaders block on that lock forever.
542 await subscribe_worker_purge_handler(self.pubsub.endpoint)
544 # We want only one worker to run repo watchers
545 # (otherwise for each new commit, we will publish multiple updates via pub/sub).
546 # leadership is determined by the first worker to obtain a lock
547 self.leadership_lock = NamedLock(
548 opal_server_config.LEADER_LOCK_FILE_PATH
549 )
550 async with self.leadership_lock:
551 # only one worker gets here, the others block. in case the leader worker
552 # is terminated, another one will obtain the lock and become leader.
553 logger.info(
554 "leadership lock acquired, leader pid: {pid}",
555 pid=os.getpid(),
556 )
558 if opal_server_config.SCOPES: 558 ↛ 561line 558 didn't jump to line 561 because the condition on line 558 was always true
559 await load_scopes(self._scopes)
561 if self.broadcast_keepalive is not None: 561 ↛ 562line 561 didn't jump to line 562 because the condition on line 561 was never true
562 self.broadcast_keepalive.start()
563 if not self._init_policy_watcher:
564 # Wait on keepalive instead to keep leadership lock acquired
565 await self.broadcast_keepalive.wait_until_done()
567 if self._init_policy_watcher: 567 ↛ 578line 567 didn't jump to line 578
568 self.watcher = setup_watcher_task(
569 self.publisher, self.pubsub.endpoint
570 )
571 # running the watcher, and waiting until it stops (until self.watcher.signal_stop() is called)
572 async with self.watcher:
573 await self.watcher.wait_until_should_stop()
575 # Worker should restart when watcher stops
576 self._graceful_shutdown()
578 if self.broadcast_listening_context is not None:
579 # __aexit__(exc_type, exc, tb): the real
580 # EventBroadcasterContextManager has no defaults, and a
581 # TypeError here (in an un-awaited background task) leaves
582 # _listen_count at 1 and the reader never cancelled.
583 await self.broadcast_listening_context.__aexit__(None, None, None)
584 logger.info(
585 "stopped listening on the broadcast channel on this worker"
586 )
588 async def stop_server_background_tasks(self):
589 logger.info("stopping background tasks...")
591 tasks: List[asyncio.Task] = []
593 if self._scopes_service is not None: 593 ↛ 602line 593 didn't jump to line 602 because the condition on line 593 was always true
594 # Bounded: a floor task's first act is to take lock_source, which a
595 # sync can hold across a whole clone/fetch. Unbounded here would
596 # hang shutdown; abandoning is the same outcome as not draining at
597 # all, so the timeout is the safe direction.
598 tasks.append(
599 asyncio.create_task(self._drain_scopes_service(_SCOPES_DRAIN_TIMEOUT))
600 )
602 if self.watcher is not None: 602 ↛ 604line 602 didn't jump to line 604 because the condition on line 602 was always true
603 tasks.append(asyncio.create_task(self.watcher.stop()))
604 if self.publisher is not None: 604 ↛ 606line 604 didn't jump to line 606 because the condition on line 604 was always true
605 tasks.append(asyncio.create_task(self.publisher.stop()))
606 if self.broadcast_keepalive is not None: 606 ↛ 607line 606 didn't jump to line 607 because the condition on line 606 was never true
607 tasks.append(asyncio.create_task(self.broadcast_keepalive.stop()))
608 if self.opal_statistics is not None: 608 ↛ 611line 608 didn't jump to line 611 because the condition on line 608 was always true
609 tasks.append(asyncio.create_task(self.opal_statistics.stop()))
611 try:
612 await asyncio.gather(*tasks)
613 except Exception:
614 logger.exception("exception while shutting down background tasks")
616 async def _drain_scopes_service(self, timeout: float) -> None:
617 """Await the DELETE floor's backgrounded clone-dir purges, bounded.
619 Without this, a DELETE that returns 204 followed by SIGTERM loses its
620 floor: the task is detached, nothing else references it, and the clone
621 dir it was about to remove survives with nothing left to reclaim it (no
622 reconciliation in this PR — PER-15612). Master removed the dir inline,
623 before returning 204, so an undrained floor is a regression against the
624 merge base rather than merely a missing improvement.
625 """
626 try:
627 await asyncio.wait_for(self._scopes_service.stop(), timeout=timeout)
628 except asyncio.TimeoutError:
629 logger.warning(
630 "Abandoned in-flight scope clone purges at shutdown after "
631 "{timeout}s; their clone dirs stay on disk (PER-15612)",
632 timeout=timeout,
633 )
634 except Exception:
635 logger.exception("Failed to drain scope clone purges at shutdown")
637 def _wire_broadcaster_give_up(self):
638 """Graceful-restart the worker if the reconnecting broadcaster gives
639 up.
641 When ``BROADCAST_RECONNECT_MAX_RETRIES`` is exhausted the reader task
642 completes by *returning*. With ``ignore_broadcaster_disconnected=False`` a
643 done reader cancels every client websocket, re-creating the drop storm. The
644 statistics path already restarts the worker on reader completion, but only
645 when ``STATISTICS_ENABLED`` — so with statistics off nothing restarts the
646 worker until the liveness probe acts. This broadcaster-level give-up hook
647 triggers the same graceful shutdown regardless of the statistics flag. It
648 fires only on give-up (return), never on cancellation (clean shutdown), so
649 normal shutdown does not re-trigger it. (When statistics are enabled both
650 this hook and the stats done-callback may fire; a second SIGTERM to an
651 already-terminating worker is a harmless no-op.)
652 """
653 broadcaster = self.pubsub.broadcaster
654 if not isinstance(broadcaster, ReconnectingBroadcaster): 654 ↛ 657line 654 didn't jump to line 657 because the condition on line 654 was always true
655 return
657 async def _on_broadcaster_give_up():
658 logger.error(
659 "Broadcaster gave up reconnecting to the backbone; "
660 "restarting this worker"
661 )
662 self._graceful_shutdown()
664 broadcaster.set_give_up_callback(_on_broadcaster_give_up)
666 def _graceful_shutdown(self):
667 logger.info("Trigger worker graceful shutdown")
668 os.kill(os.getpid(), signal.SIGTERM)