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

1import asyncio 

2import os 

3import signal 

4import sys 

5import traceback 

6from functools import partial 

7from pathlib import Path 

8from typing import List, Optional 

9 

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 

50 

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 

56 

57 

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 

84 

85 The server can run in multiple workers (by gunicorn or uvicorn). 

86 

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. 

90 

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: 

94 

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 

121 

122 self._configure_monitoring() 

123 metrics.increment("startup") 

124 

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 ) 

130 

131 self.broadcaster_uri = broadcaster_uri 

132 self.master_token = master_token 

133 

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 ) 

152 

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 

159 

160 self.pubsub = PubSub(signer=self.signer, broadcaster_uri=broadcaster_uri) 

161 self._wire_broadcaster_give_up() 

162 

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) 

167 

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 ) 

177 

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 

182 

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 ) 

228 

229 self.watcher: PolicyWatcherTask = None 

230 self.leadership_lock: Optional[NamedLock] = None 

231 

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 ) 

260 

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 

264 

265 # init fastapi app 

266 self.app: FastAPI = self._init_fast_api_app() 

267 

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 ) 

281 

282 configure_middleware(app) 

283 self._configure_api_routes(app) 

284 self._configure_lifecycle_callbacks(app) 

285 

286 return app 

287 

288 def _configure_monitoring(self): 

289 configure_logs() 

290 

291 apm.configure_apm(opal_server_config.ENABLE_DATADOG_APM, "opal-server") 

292 

293 configure_server_metrics() 

294 

295 def _configure_api_routes(self, app: FastAPI): 

296 """Mounts the api routes on the app object.""" 

297 authenticator = JWTAuthenticator(self.signer) 

298 

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) 

302 

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) 

313 

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 ) 

339 

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 ) 

361 

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) 

365 

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 ) 

375 

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

381 

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

407 

408 register_internal_stats_route( 

409 app, 

410 enabled=opal_server_config.DEBUG_INTERNAL_STATS, 

411 dependencies=[Depends(authenticator)], 

412 ) 

413 

414 return app 

415 

416 def _configure_lifecycle_callbacks(self, app: FastAPI): 

417 """Registers callbacks on app startup and shutdown. 

418 

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

423 

424 @app.on_event("startup") 

425 async def startup_event(): 

426 logger.info("*** OPAL Server Startup ***") 

427 

428 try: 

429 self._task = asyncio.create_task(self.start_server_background_tasks()) 

430 

431 except Exception: 

432 logger.critical("Exception while starting OPAL") 

433 traceback.print_exc() 

434 

435 sys.exit(1) 

436 

437 @app.on_event("shutdown") 

438 async def shutdown_event(): 

439 logger.info("triggered shutdown event") 

440 await self.stop_server_background_tasks() 

441 

442 return app 

443 

444 async def start_server_background_tasks(self): 

445 """Starts the background processes (as asyncio tasks) if such are 

446 configured. 

447 

448 all workers will start these tasks: 

449 - publisher: a client that is used to publish updates to the client. 

450 

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

520 

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 ) 

536 

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) 

543 

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 ) 

557 

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) 

560 

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

566 

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

574 

575 # Worker should restart when watcher stops 

576 self._graceful_shutdown() 

577 

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 ) 

587 

588 async def stop_server_background_tasks(self): 

589 logger.info("stopping background tasks...") 

590 

591 tasks: List[asyncio.Task] = [] 

592 

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 ) 

601 

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

610 

611 try: 

612 await asyncio.gather(*tasks) 

613 except Exception: 

614 logger.exception("exception while shutting down background tasks") 

615 

616 async def _drain_scopes_service(self, timeout: float) -> None: 

617 """Await the DELETE floor's backgrounded clone-dir purges, bounded. 

618 

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

636 

637 def _wire_broadcaster_give_up(self): 

638 """Graceful-restart the worker if the reconnecting broadcaster gives 

639 up. 

640 

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 

656 

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

663 

664 broadcaster.set_give_up_callback(_on_broadcaster_give_up) 

665 

666 def _graceful_shutdown(self): 

667 logger.info("Trigger worker graceful shutdown") 

668 os.kill(os.getpid(), signal.SIGTERM)