Coverage for /usr/local/lib/python3.10/site-packages/opal_server-0.0.0-py3.10.egg/opal_server/scopes/task.py: 61%
86 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 datetime
3from pathlib import Path
4from typing import Any
6from fastapi_websocket_pubsub import Topic
7from opal_common.logger import logger
8from opal_common.monitoring import metrics
9from opal_server.config import opal_server_config
10from opal_server.git_fetcher import (
11 GitPolicyFetcher,
12 drain_git_ops,
13 git_busy_count,
14 shutdown_git_executor,
15)
16from opal_server.metrics_setup import configure_server_metrics
17from opal_server.policy.watcher.task import BasePolicyWatcherTask
18from opal_server.redis_utils import RedisDB
19from opal_server.scopes.purge import LeaderScopePurger
20from opal_server.scopes.scope_repository import ScopeRepository
21from opal_server.scopes.service import ScopesService
23# Upper bound on the shutdown drain of in-flight scope purges. The drain is
24# best-effort: a purge's rmtree runs on a worker thread and completes whether or
25# not we are still awaiting it. What is abandoned before it starts is NOT
26# recovered — no reconciliation sweep exists in this PR (split out, PER-15612),
27# and no later purge will name a deleted scope's source. Blocking shutdown
28# longer would be strictly worse —
29# stop() runs while the leadership lock is still held, so no other worker can
30# take over, and k8s's terminationGracePeriodSeconds (30s by default) would
31# SIGKILL us anyway.
32_PURGE_DRAIN_TIMEOUT = 5.0
35class ScopesPolicyWatcherTask(BasePolicyWatcherTask):
36 def __init__(self, *args, **kwargs):
37 super().__init__(*args, **kwargs)
38 # Set in start(); None means "nothing to unsubscribe" (stop() without a
39 # successful start, or a second stop()).
40 self._purger_sub_id = None
42 self._scopes = ScopeRepository(RedisDB(opal_server_config.REDIS_URL))
43 self._service = ScopesService(
44 base_dir=Path(opal_server_config.BASE_DIR),
45 scopes=self._scopes,
46 pubsub_endpoint=self._pubsub_endpoint,
47 )
48 self._purger = LeaderScopePurger(
49 base_dir=Path(opal_server_config.BASE_DIR),
50 scopes=self._scopes,
51 pubsub_endpoint=self._pubsub_endpoint,
52 )
54 async def start(self):
55 await super().start()
56 # Leader-only purge authorization: this task starts only on the leader, so
57 # registering here (not at worker boot) preserves the invariant that
58 # the leader is the only mutator on sync paths (see the note in
59 # scopes/purge.py — a delete also removes on the serving worker).
60 #
61 # Registered under a DEDICATED subscriber id rather than through
62 # PubSubEndpoint.subscribe, which files every server-side subscription
63 # under one shared endpoint id: EventNotifier.unsubscribe deletes that
64 # id's WHOLE callback list for the topic, so unsubscribing by topic in
65 # stop() would also drop the every-worker handle_purge_message
66 # subscription (registered once at boot in server.py, never re-added) —
67 # leaving this process deaf to fleet purges for the rest of its life.
68 # Our own id makes the leader subscription independently removable.
69 self._purger_sub_id = self._pubsub_endpoint.notifier.gen_subscriber_id()
70 await self._pubsub_endpoint.notifier.subscribe(
71 self._purger_sub_id,
72 [opal_server_config.SCOPES_PURGE_CHANNEL],
73 self._purger.handle,
74 )
75 # With POLICY_REFRESH_INTERVAL <= 0 this boot sync is the ONLY
76 # pass-originated sync this process ever runs, so a source that failed
77 # transiently during the pre-fork preload (whose entry survives
78 # reset_caches on purpose) must not be inherited — it would never be
79 # attempted again. Clearing the inherited entries (rather than not
80 # honouring the backoff) keeps the within-pass property: phase-2
81 # duplicates of a source that fails in THIS pass are still collapsed
82 # to one attempt.
83 # Clearing the whole dict is safe here: this runs once, right after
84 # this process won leadership, and a freshly forked worker cannot have
85 # recorded anything of its own before that (only leaders sync).
86 if opal_server_config.POLICY_REFRESH_INTERVAL <= 0: 86 ↛ 88line 86 didn't jump to line 88 because the condition on line 86 was always true
87 GitPolicyFetcher.source_backoff.clear()
88 self._tasks.append(asyncio.create_task(self._sync_all()))
90 if opal_server_config.POLICY_REFRESH_INTERVAL > 0: 90 ↛ 91line 90 didn't jump to line 91 because the condition on line 90 was never true
91 self._tasks.append(asyncio.create_task(self._periodic_polling()))
93 async def stop(self):
94 # stop() runs TWICE on the normal path — once from
95 # BasePolicyWatcherTask.__aexit__ (server.py's `async with self.watcher`)
96 # and again from stop_server_background_tasks — so every step here is
97 # idempotent: the subscriber id is consumed on first use, signal_stop and
98 # the drain are no-ops once done.
99 #
100 # Stop accepting new purge messages first (cheap, non-blocking).
101 if self._purger_sub_id is not None: 101 ↛ 109line 101 didn't jump to line 109 because the condition on line 101 was always true
102 sub_id, self._purger_sub_id = self._purger_sub_id, None
103 try:
104 await self._pubsub_endpoint.notifier.unsubscribe(
105 sub_id, [opal_server_config.SCOPES_PURGE_CHANNEL]
106 )
107 except Exception:
108 logger.exception("Failed to unsubscribe scope purge handler on stop")
109 self._purger.signal_stop()
111 # Cancel our tasks BEFORE draining the purges, not after: a queued purge
112 # first waits on GitPolicyFetcher.lock_source, which a sync task holds
113 # across a clone/fetch (up to SCOPES_GIT_FETCH_TIMEOUT, and unbounded
114 # when that is 0). Draining first would wait on a lock whose release
115 # requires the very cancellation the drain is blocking — an unrecoverable
116 # shutdown hang, taken while the leadership lock is still held.
117 result = await super().stop()
119 # Best-effort, bounded (see _PURGE_DRAIN_TIMEOUT).
120 #
121 # ONLY the purger is drained here. The DELETE floor's tasks live on the
122 # ScopesService that init_scope_router received (built in server.py), NOT
123 # on self._service — the watcher constructs its own, and only ever syncs
124 # with it. Draining self._service gathered an empty set. This is also the
125 # wrong place structurally: the watcher exists only on the leader, while
126 # a DELETE usually lands on a non-leader. That drain is per-worker, on
127 # OpalServer.stop_server_background_tasks.
128 try:
129 await asyncio.wait_for(self._purger.stop(), timeout=_PURGE_DRAIN_TIMEOUT)
130 except asyncio.TimeoutError:
131 logger.warning(
132 "Abandoned in-flight scope purges at shutdown after {timeout}s; "
133 "their clone dirs stay on disk until the source is purged again",
134 timeout=_PURGE_DRAIN_TIMEOUT,
135 )
136 return result
138 async def _sync_all(self, honor_backoff: bool = True):
139 # sync_scopes must be wrapped: this coroutine is launched fire-and-forget
140 # from start() (boot), so an unhandled raise here would die silently —
141 # the exception is never retrieved (stop() gathers with
142 # return_exceptions=True and discards it), not even asyncio's
143 # "never retrieved" warning until GC.
144 #
145 # honor_backoff defaults to True for the boot call in start(): this
146 # process may have been forked from a master whose preload already
147 # discovered which repos are unreachable, and re-attempting all of them
148 # at boot is the storm the backoff exists to prevent. trigger() passes
149 # False for the operator-driven refresh-all — see there.
150 try:
151 await self._service.sync_scopes(honor_backoff=honor_backoff)
152 except Exception:
153 logger.exception("Scope sync (sync_scopes) failed")
155 async def _periodic_polling(self):
156 try:
157 while True:
158 await asyncio.sleep(opal_server_config.POLICY_REFRESH_INTERVAL)
159 # Leader heartbeat. This loop runs only inside the leadership
160 # lock, so `sum by env` reaching 0 means no worker holds it
161 # anywhere and scope syncing has silently stopped — pods stay
162 # Ready and /healthcheck stays 200 throughout.
163 metrics.gauge("opal_server.scopes.leader", 1)
164 logger.info("Periodic sync")
165 try:
166 await self._service.sync_scopes(only_poll_updates=True)
167 except asyncio.CancelledError:
168 raise
169 except Exception as e:
170 logger.exception(f"Periodic sync (sync_scopes) failed")
172 except asyncio.CancelledError:
173 logger.info("Periodic sync cancelled")
174 raise
176 async def trigger(self, topic: Topic, data: Any):
177 if data is not None and isinstance(data, dict):
178 # Refresh single scope
179 try:
180 await self._service.sync_scope(
181 scope_id=data["scope_id"],
182 force_fetch=data.get("force_fetch", False),
183 hinted_hash=data.get("hinted_hash"),
184 req_time=datetime.datetime.now(),
185 )
186 except KeyError:
187 logger.warning(
188 "Got invalid keyword args for single scope refresh: %s", data
189 )
190 else:
191 # Refresh all scopes. This branch is reached only from something
192 # asking for a sync NOW — POST /scopes/refresh (the git-provider
193 # webhook publishes on this topic too, but only when
194 # POLICY_REPO_URL is set, which scopes mode does not use) — so it
195 # does not honour the per-source backoff: the
196 # sources an operator hits this endpoint for are precisely the ones
197 # they have just repaired, and answering them with a silent skip
198 # makes the endpoint useless in the only situation it is used.
199 await self._sync_all(honor_backoff=False)
201 @staticmethod
202 def preload_scopes():
203 """Clone all scopes repositories as part as server startup.
205 This speeds up the first sync of scopes after workers are
206 started.
207 """
208 if opal_server_config.SCOPES:
209 # This runs in the gunicorn MASTER (scripts/gunicorn_conf.py:when_ready)
210 # before any worker configured monitoring. Without this the git
211 # metrics emitted during preload (git_ops_in_flight, scopes.count,
212 # sources_in_backoff, git_op_failures...) go out un-namespaced and
213 # never reach the dashboards/monitors built on permit.opal.*.
214 configure_server_metrics()
215 logger.info("Preloading repo clones for scopes")
217 service = ScopesService(
218 base_dir=Path(opal_server_config.BASE_DIR),
219 scopes=ScopeRepository(RedisDB(opal_server_config.REDIS_URL)),
220 pubsub_endpoint=None,
221 )
222 asyncio.run(service.sync_scopes(notify_on_changes=False))
224 # Bounded window for a just-finished clone/fetch to clear its in-flight
225 # marker before teardown+fork. Ops still lingering (hung remote) are
226 # left running; reset_caches's guard then skips freeing their handles.
227 # A False return means the drain timed out with git ops STILL running:
228 # those threads persist in the master across the fork, so a forked
229 # worker can race them on the shared clone dir. Log it — this is the
230 # one condition that carries that risk, and it must not be silent.
231 drained = drain_git_ops(opal_server_config.SCOPES_GIT_PRELOAD_DRAIN_TIMEOUT)
232 if not drained:
233 logger.warning(
234 "Preload drain timed out ({timeout}s) with git ops still "
235 "in flight ({in_flight}); they persist in the master across "
236 "fork. Consider raising SCOPES_GIT_PRELOAD_DRAIN_TIMEOUT or "
237 "lowering SCOPES_GIT_FETCH_TIMEOUT.",
238 timeout=opal_server_config.SCOPES_GIT_PRELOAD_DRAIN_TIMEOUT,
239 in_flight=git_busy_count(),
240 )
242 # Clear git-op bookkeeping built during preload (in-flight markers
243 # and the loop-bound live-op semaphore) so the gunicorn master does
244 # not carry stale state into forked workers. Git ops run on per-op
245 # daemon threads; there is no shared pool to tear down.
246 shutdown_git_executor()
248 # Drop every cached repo handle/lock/timestamp built during preload
249 # so none of it is inherited by forked workers. Sync (the only path
250 # that populates these caches) is leader-only, so a non-leader worker
251 # that inherited a handle would only ever release it through the
252 # fleet-wide purge broadcast. Since 0.9.9-rc.3 every worker keeps a
253 # backbone reader when SCOPES is on and the broadcaster is the
254 # reconnecting one, so that broadcast does reach client-less workers;
255 # with BROADCAST_RECONNECT_ENABLED=false (legacy broadcaster) the reader
256 # runs only while a client is connected and an inherited handle could
257 # still be pinned for life — clearing here keeps both cases safe. The
258 # on-disk clones remain; workers re-open handles lazily.
259 GitPolicyFetcher.reset_caches()
261 logger.warning("Finished preloading repo clones for scopes.")