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

1import asyncio 

2import datetime 

3from pathlib import Path 

4from typing import Any 

5 

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 

22 

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 

33 

34 

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 

41 

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 ) 

53 

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

89 

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

92 

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

110 

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

118 

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 

137 

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

154 

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

171 

172 except asyncio.CancelledError: 

173 logger.info("Periodic sync cancelled") 

174 raise 

175 

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) 

200 

201 @staticmethod 

202 def preload_scopes(): 

203 """Clone all scopes repositories as part as server startup. 

204 

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

216 

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

223 

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 ) 

241 

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

247 

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

260 

261 logger.warning("Finished preloading repo clones for scopes.")