Coverage for /usr/local/lib/python3.10/site-packages/opal_server-0.0.0-py3.10.egg/opal_server/pubsub_resilience.py: 17%

279 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 11:54 +0000

1"""Resilience wrappers around the ``fastapi_websocket_pubsub`` / 

2``fastapi_websocket_rpc`` pub/sub layer used by the OPAL server. 

3 

4Two upstream issues let a *transient* broadcaster-backbone disconnect 

5(Postgres ``LISTEN/NOTIFY``, Redis, Kafka) escalate into a self-sustaining, 

6fleet-wide client connection-drop storm that only a worker restart clears: 

7 

81. ``EventBroadcaster``'s reader task runs to completion when the backbone 

9 connection drops and is never restarted while clients remain connected 

10 (``_subscription_task`` is reset only when the listener count reaches 0). 

11 Because OPAL runs with ``ignore_broadcaster_disconnected=False``, every 

12 client websocket waits on that shared reader task; once it is *done* every 

13 client is cancelled and dropped, indefinitely. 

142. ``ConnectionManager.disconnect`` is not idempotent; the RPC endpoint can 

15 call it twice for the same socket, raising 

16 ``ValueError('list.remove(x): x not in list')``. 

17 

18``ReconnectingBroadcaster`` keeps the reader task *pending* across transient 

19backbone losses by reconnecting with bounded exponential backoff, and adds two 

20consistency mechanisms so a backbone gap does not silently desync instances: 

21 

22* **(B) outbound replay buffer** — broadcasts that fail to reach the backbone 

23 while it is down are kept in a bounded FIFO and replayed once it reconnects, 

24 so peers that re-subscribe in time catch up without a refetch. This narrows 

25 the staleness window; it is *not* a delivery guarantee (the backbone keeps no 

26 replay of its own and a slow peer may not have re-subscribed at flush time). 

27* **(A) resync on recovery** — the *guarantee*. After **any** gap the broadcaster 

28 fires the registered ``on_reconnect`` callback. OPAL uses it to make this 

29 worker's own clients re-run their full (scope-aware) policy + data 

30 reconciliation. Every worker experiences the same gap, so every worker 

31 reconciles its own clients (a worker's clients may have missed *incoming* peer 

32 updates during the gap — independent of what this worker published), and the 

33 fleet converges to current truth. 

34 

35``SafeConnectionManager`` makes ``disconnect`` idempotent and can close a 

36worker's client connections (staggered) to drive the resync. All of this is a 

37stop-gap until the fixes land in the upstream libraries (see Phase 2 of the plan). 

38""" 

39import asyncio 

40import random 

41from collections import deque 

42from typing import Awaitable, Callable, Optional 

43 

44from fastapi import WebSocket 

45from fastapi_websocket_pubsub import EventBroadcaster, PubSubEndpoint 

46from fastapi_websocket_pubsub.event_broadcaster import BroadcastNotification 

47from fastapi_websocket_pubsub.event_notifier import Subscription 

48from fastapi_websocket_pubsub.util import pydantic_serialize 

49from fastapi_websocket_rpc.connection_manager import ConnectionManager 

50from opal_common.logger import logger 

51 

52ReconnectCallback = Callable[[], Awaitable[None]] 

53 

54 

55class SafeConnectionManager(ConnectionManager): 

56 """A ``ConnectionManager`` whose ``disconnect`` is idempotent, able to drop 

57 its connections on demand. 

58 

59 The upstream ``disconnect`` calls ``self.active_connections.remove(websocket)`` 

60 unconditionally, so a second disconnect for the same socket raises 

61 ``ValueError('list.remove(x): x not in list')`` — which escapes 

62 ``WebsocketRPCEndpoint.main_loop`` as an unretrieved task exception and, under a 

63 reconnect storm, is logged thousands of times. Guarding the removal turns a 

64 double disconnect into a no-op. ``close_all_staggered`` additionally lets the 

65 server intentionally recycle its client connections (e.g. to force a post-outage 

66 resync) without a thundering herd. 

67 """ 

68 

69 def disconnect(self, websocket: WebSocket): 

70 try: 

71 self.active_connections.remove(websocket) 

72 except ValueError: 

73 logger.debug("Ignoring duplicate websocket disconnect") 

74 

75 async def close_all_staggered( 

76 self, min_interval: float = 0.0, max_interval: float = 0.2 

77 ) -> int: 

78 """Close every currently-tracked client websocket, spaced out with 

79 jitter. 

80 

81 Clients reconnect on their own and re-run their on-connect reconciliation, 

82 so this is how a worker forces its clients back to a consistent state after 

83 a broadcaster gap. Close code 1012 ("Service Restart") signals a reconnect. 

84 

85 Note: the broadcaster's reader task is tied to the per-client listening 

86 context, so the caller MUST pin a listening context around this call (see 

87 ``opal_server.pubsub``) — otherwise closing the last client cancels the 

88 reader we just worked to keep alive. 

89 """ 

90 connections = list(self.active_connections) 

91 if not connections: 

92 return 0 

93 logger.info( 

94 f"Resync: closing {len(connections)} client connection(s) to trigger " 

95 "client-side reconciliation" 

96 ) 

97 closed = 0 

98 for index, websocket in enumerate(connections): 

99 try: 

100 await websocket.close(code=1012) 

101 closed += 1 

102 except Exception as e: 

103 logger.warning(f"Error closing a client websocket during resync: {e!r}") 

104 # Jitter between closes (but not after the last one). 

105 if max_interval > 0 and index < len(connections) - 1: 

106 await asyncio.sleep(random.uniform(min_interval, max_interval)) 

107 return closed 

108 

109 

110class ReconnectingBroadcaster(EventBroadcaster): 

111 """An ``EventBroadcaster`` whose listener reconnects instead of dying, 

112 buffers failed outbound broadcasts, and fires a resync hook after a gap. 

113 

114 The base reader coroutine (``__read_notifications__``) returns when the backbone 

115 connection closes and is not restarted while clients stay connected, leaving 

116 ``get_reader_task()`` permanently *done* — which cancels every client websocket 

117 loop. This subclass wraps the connect/subscribe/read cycle in a reconnect loop 

118 with bounded exponential backoff, so the reader task stays *pending* across 

119 transient outages. The task only completes on clean shutdown (cancellation) or 

120 after ``reconnect_max_retries`` consecutive failures (give-up), in which case it 

121 fires the registered give-up hook so OPAL graceful-restarts the worker. 

122 """ 

123 

124 def __init__( 

125 self, 

126 *args, 

127 reconnect_max_retries: int = 0, 

128 reconnect_backoff_min: float = 0.5, 

129 reconnect_backoff_max: float = 30.0, 

130 replay_buffer_size: int = 10000, 

131 resync_settle_seconds: float = 2.0, 

132 **kwargs, 

133 ): 

134 super().__init__(*args, **kwargs) 

135 self._reconnect_max_retries = reconnect_max_retries 

136 self._reconnect_backoff_min = reconnect_backoff_min 

137 self._reconnect_backoff_max = reconnect_backoff_max 

138 self._replay_buffer_size = replay_buffer_size 

139 self._resync_settle_seconds = resync_settle_seconds 

140 # (B) bounded outbound replay buffer; deque(maxlen) drops the oldest on overflow. 

141 self._outbound_buffer: deque = deque( 

142 maxlen=replay_buffer_size if replay_buffer_size > 0 else None 

143 ) 

144 # Single lock serialises buffer mutation, the overflow flag, and the flush, 

145 # so a concurrent broadcast cannot race the flush's drain/flag-reset. 

146 self._buffer_lock = asyncio.Lock() 

147 # (A) single-flight resync hook. Gap detection itself is reader-task-local 

148 # (see __read_notifications__), not an instance flag. 

149 self._recovery_task: Optional[asyncio.Task] = None 

150 # Set when a gap arrives while a recovery is already in flight (single-flight), 

151 # so the in-flight recovery loops once more rather than dropping that gap. 

152 self._recovery_rerun_requested = False 

153 self._on_reconnect: Optional[ReconnectCallback] = None 

154 # Fired once if the reader gives up (exhausts reconnect retries) and returns, 

155 # so OPAL can graceful-restart the worker even with statistics disabled. 

156 self._on_give_up: Optional[ReconnectCallback] = None 

157 # Live backbone-subscription state; see is_in_backbone_gap(). True only while 

158 # actively subscribed — i.e. a publish right now would actually reach peer 

159 # workers. Deliberately a separate signal from is_reader_healthy(): that one 

160 # stays True across a transient reconnect (so the k8s probe does not flap the 

161 # pod), which is exactly the wrong signal for gating delivery. 

162 # Instance attrs (not task-local): FreezablePubSubEndpoint reads them from outside 

163 # the reader task. 

164 self._backbone_connected = False 

165 # Monotonic gap counter, bumped on every connected -> disconnected edge in the 

166 # reader; see backbone_gap_generation(). 

167 self._gap_generation = 0 

168 # Whether this broadcaster ever held a backbone subscription — distinguishes a real 

169 # GAP (had a session, lost it) from "never connected yet" (boot, or backbone down 

170 # from the start), where freezing would be wrong: no resync fires on a FIRST 

171 # connect, so anything frozen before it would be lost, not deferred. 

172 self._had_backbone_connection = False 

173 

174 def set_reconnect_callback(self, callback: Optional[ReconnectCallback]): 

175 """Register an ``async () -> None`` callback fired once after each gap 

176 (a reconnect that follows a previously-established connection).""" 

177 self._on_reconnect = callback 

178 

179 def set_give_up_callback(self, callback: Optional[ReconnectCallback]): 

180 """Register an ``async () -> None`` callback fired once if the reader 

181 gives up after exhausting ``reconnect_max_retries``. 

182 

183 Fires only when the reader task completes by *returning* (give- 

184 up), never on cancellation (clean shutdown), so OPAL can wire it 

185 to a graceful worker restart without normal shutdown re- 

186 triggering it. 

187 """ 

188 self._on_give_up = callback 

189 

190 async def start_reader_task(self): 

191 """Spawn the reconnecting reader task once. 

192 

193 Unlike the base implementation we do not connect the channel 

194 here — the reader loop owns (re)connection, so a backbone that 

195 is already down at startup is retried rather than raised to the 

196 first connecting client. 

197 """ 

198 if self._subscription_task is not None: 

199 logger.debug("No need for listen task, already started") 

200 return self._subscription_task 

201 logger.debug("Spawning reconnecting broadcast listen task") 

202 # Scope gap detection to THIS reader task: a stale True from a previous task 

203 # (reader cancelled when the last listener left, then restarted later) would make 

204 # is_in_backbone_gap() freeze publishes before the new task's FIRST subscribe — 

205 # but a first connect fires no resync, so those publishes would be lost, not 

206 # deferred. Same task-scoping rationale as ``had_prior_connection`` in 

207 # ``__read_notifications__``. 

208 self._had_backbone_connection = False 

209 self._subscription_task = asyncio.create_task(self.__read_notifications__()) 

210 return self._subscription_task 

211 

212 def is_reader_healthy(self) -> bool: 

213 """Report whether the reconnecting reader can still serve connected 

214 clients. 

215 

216 Used by the server ``/healthcheck`` so a k8s readiness/liveness probe can 

217 route away from (or restart) a worker whose reader is wedged while clients 

218 depend on it — defense in depth on top of the reconnect loop itself. 

219 

220 Health is judged against the listener count, not the backbone connection: 

221 

222 * No listeners (``_listen_count <= 0``): healthy. The reader is started lazily 

223 when the count goes 0->1 and reset to ``None`` when it returns to 0, so its 

224 absence here is expected idleness, not a fault — nothing depends on it. 

225 * Listeners present: healthy only while the reader task is a live, *pending* 

226 task. A pending task INCLUDES the case where it is mid-reconnect through a 

227 backbone outage: ``ReconnectingBroadcaster`` deliberately keeps the reader 

228 pending across a drop, so a transient reconnect must NOT read as unhealthy 

229 (otherwise the probe would flap the pod during every normal backbone blip). 

230 It reads unhealthy only in the two wedged states — the task is ``None`` 

231 (never started / leaked listen-count) or ``done`` (crashed, or gave up after 

232 ``reconnect_max_retries``) — which is exactly when clients are stuck. 

233 

234 This is a cheap, non-blocking attribute read (no await, no lock); any rare race 

235 against a reconnect is absorbed by the probe's ``failureThreshold``. 

236 

237 Returns: 

238 bool: True if idle or the reader is live and pending; False if listeners 

239 depend on a missing or completed reader task. 

240 """ 

241 if self._listen_count <= 0: 

242 return True 

243 return ( 

244 self._subscription_task is not None and not self._subscription_task.done() 

245 ) 

246 

247 def backbone_gap_generation(self) -> int: 

248 """Monotonic count of backbone gaps: bumped each time an established 

249 subscription drops (the connected -> disconnected edge in the reader). 

250 

251 Lets ``FreezablePubSubEndpoint`` tell two back-to-back gaps apart even when no 

252 publish is delivered between them (recovery itself never publishes — the resync 

253 closes client sockets and clients refetch), so each gap opens its own freeze 

254 episode instead of merging into the previous one. 

255 """ 

256 return self._gap_generation 

257 

258 def is_in_backbone_gap(self) -> bool: 

259 """Whether the broadcaster is mid-GAP: it *had* a live backbone subscription, 

260 lost it, and the reader is still trying to get it back. 

261 

262 This — not mere "not connected" — is the publish-freeze condition used by 

263 ``FreezablePubSubEndpoint``, because only a real gap has the recovery path the 

264 freeze relies on (the ``on_reconnect`` resync fires exclusively for reconnects 

265 that follow an established session). The two excluded states must NOT freeze: 

266 

267 * **Reader not running** (no listeners yet / worker idle / last client left and 

268 the upstream cancelled the reader): the backbone may be perfectly healthy — 

269 freezing here would silently drop publishes fleet-wide with nothing to ever 

270 reconcile them. Delegating preserves the pre-freeze behavior (share-context 

271 broadcast + local delivery). 

272 * **Never connected in this reader's lifetime** (boot, or backbone already down 

273 at startup): a FIRST successful connect fires no gap recovery, so a publish 

274 frozen in this window would be lost, not deferred. Pre-freeze behavior 

275 (deliver locally, buffer outbound for replay) is strictly better here. 

276 """ 

277 return ( 

278 self._subscription_task is not None 

279 and not self._subscription_task.done() 

280 and self._had_backbone_connection 

281 and not self._backbone_connected 

282 ) 

283 

284 async def __broadcast_notifications__(self, subscription: Subscription, data): 

285 """Share a local notification with the backbone; buffer it on failure. 

286 

287 ``__broadcast_notifications__`` ends in a double underscore (not name-mangled), 

288 so this override is what ``_subscribe_to_all_topics`` dispatches to. The local 

289 delivery to this worker's own clients is independent of this publish (it runs 

290 as a sibling notifier callback), so a failure here means only "peers may miss 

291 this" — we keep it for replay rather than dropping it. 

292 """ 

293 try: 

294 await super().__broadcast_notifications__(subscription, data) 

295 except asyncio.CancelledError: 

296 raise 

297 except Exception as e: 

298 await self._buffer_outbound(subscription.topic, data, e) 

299 

300 async def _buffer_outbound(self, topic, data, error: Exception): 

301 async with self._buffer_lock: 

302 if self._replay_buffer_size <= 0: 

303 # buffering disabled — the gap resync is the only recovery path. 

304 logger.warning( 

305 f"Broadcast to backbone failed ({error!r}); replay buffer disabled" 

306 ) 

307 return 

308 # At capacity this append drops the oldest entry (bounded deque); the resync 

309 # on reconnect still reconciles clients, so the drop only widens the window. 

310 overflow = len(self._outbound_buffer) >= self._replay_buffer_size 

311 self._outbound_buffer.append((topic, data)) 

312 logger.warning( 

313 f"Broadcast to backbone failed ({error!r}); buffered for replay " 

314 f"({len(self._outbound_buffer)}/{self._replay_buffer_size}" 

315 f"{', OVERFLOW — oldest dropped' if overflow else ''})" 

316 ) 

317 

318 async def __read_notifications__(self): 

319 """Read incoming broadcasts, reconnecting on backbone disconnect. 

320 

321 ``had_prior_connection`` is deliberately a task-local, not an instance 

322 attribute: gap recovery must fire only for a reconnect *within this 

323 reader task's own loop* (a real backbone gap). When the last client 

324 disconnects, the upstream cancels this task and clears 

325 ``_subscription_task``; the next client starts a *fresh* reader task, 

326 which must start clean — an instance flag would carry a stale ``True`` 

327 into that fresh task and schedule a spurious full recovery (flush + 

328 client-recycling resync) on its very first connect, a churn loop that 

329 also hits stats-off deployments. 

330 """ 

331 attempt = 0 

332 had_prior_connection = False 

333 while True: 

334 try: 

335 channel = await self._ensure_connected() 

336 logger.info( 

337 f"Broadcaster listener connected to channel '{self._channel}'" 

338 ) 

339 async with channel.subscribe(channel=self._channel) as subscriber: 

340 # Subscribed: the backbone is reachable, so publishes will fan out to 

341 # peers again — reopen the publish gate (see FreezablePubSubEndpoint). 

342 self._backbone_connected = True 

343 self._had_backbone_connection = True 

344 # We are subscribed again; recover concurrently so we keep reading 

345 # (and can receive peers' replays) during the settle window. 

346 if had_prior_connection: 

347 self._schedule_gap_recovery() 

348 had_prior_connection = True 

349 # A connect that ends without sustaining (no read) is a flap, not a 

350 # healthy session, so only a sustained subscriber resets the attempt 

351 # counter — otherwise a connect-OK/instant-close loop would never 

352 # increment ``attempt`` and ``reconnect_max_retries`` could never trip. 

353 sustained = False 

354 async for event in subscriber: 

355 if not sustained: 

356 sustained = True 

357 attempt = 0 

358 await self._handle_broadcast_event(event) 

359 if sustained: 

360 logger.warning( 

361 "Broadcast subscriber ended (backbone connection closed); " 

362 "reconnecting" 

363 ) 

364 else: 

365 attempt += 1 

366 logger.warning( 

367 "Broadcast subscriber ended immediately without sustaining " 

368 f"(attempt {attempt}); treating as a failed reconnect" 

369 ) 

370 if self._gave_up(attempt): 

371 await self._fire_give_up() 

372 return 

373 except asyncio.CancelledError: 

374 logger.info("Broadcaster listener cancelled; stopping") 

375 await self._cancel_pending_tasks() 

376 raise 

377 except Exception as e: 

378 attempt += 1 

379 logger.error(f"Broadcaster listener error (attempt {attempt}): {e!r}") 

380 if self._gave_up(attempt): 

381 await self._fire_give_up() 

382 return 

383 finally: 

384 # Any exit from the read cycle (backbone closed, error, or cancel) means we 

385 # are no longer subscribed — close the publish gate until we re-subscribe, so 

386 # a write during the gap is not applied on this worker alone. An established 

387 # subscription ending here is a NEW gap: bump the generation so freeze 

388 # episodes never span two gaps. 

389 if self._backbone_connected: 

390 self._gap_generation += 1 

391 self._backbone_connected = False 

392 await self._safe_disconnect_channel() 

393 await asyncio.sleep(self._backoff_seconds(attempt)) 

394 

395 def _gave_up(self, attempt: int) -> bool: 

396 """Return whether ``reconnect_max_retries`` is exhausted (0 = retry 

397 forever).""" 

398 if self._reconnect_max_retries and attempt >= self._reconnect_max_retries: 

399 logger.error( 

400 f"Broadcaster reconnect exhausted after {attempt} attempts; " 

401 "giving up so the worker can restart" 

402 ) 

403 return True 

404 return False 

405 

406 async def _cancel_pending_tasks(self): 

407 # Cancel AND join child notify/recovery tasks so they fully unwind (releasing any 

408 # pinned listening context) before the reader re-raises its own cancellation. 

409 # Exclude the current task defensively (the reader is not in _tasks, but be safe). 

410 current = asyncio.current_task() 

411 tasks = [task for task in self._tasks if task is not current] 

412 for task in tasks: 

413 task.cancel() 

414 if tasks: 

415 await asyncio.gather(*tasks, return_exceptions=True) 

416 

417 def _schedule_gap_recovery(self): 

418 # Single-flight: a flap during a recovery must not spawn a second concurrent 

419 # recovery (concurrent flushes corrupt the buffer; double resync re-storms). 

420 # Instead, request a rerun so a gap that lands while a recovery is in flight — 

421 # including during the late ``_fire_reconnect`` phase, after the flush — is still 

422 # flushed/resynced by one more loop of the in-flight recovery rather than dropped. 

423 if self._recovery_task is not None and not self._recovery_task.done(): 

424 logger.debug("Gap recovery already in progress; requesting a rerun") 

425 self._recovery_rerun_requested = True 

426 return 

427 self._recovery_rerun_requested = False 

428 self._recovery_task = asyncio.create_task(self._recover_after_gap()) 

429 self._tasks.add(self._recovery_task) 

430 self._recovery_task.add_done_callback(self._tasks.discard) 

431 

432 async def _recover_after_gap(self): 

433 """After reconnecting following a gap: let peers re-subscribe, replay the 

434 buffer (B), then fire the resync hook (A). 

435 

436 The whole recovery runs inside a pinned listening context so the reader task 

437 cannot be cancelled mid-recovery — neither by the resync hook closing every 

438 client nor by an unrelated drop to zero listeners during the settle window. 

439 

440 Loops once more whenever a gap arrives during an iteration (single-flight rerun): 

441 the rerun flag is cleared at the top of each iteration and re-checked after the 

442 full pin -> settle -> flush -> fire body, so a gap landing at any point during the 

443 body — clear happens first, check happens last — is captured and triggers exactly 

444 one more iteration (single-threaded asyncio, no lock needed). 

445 """ 

446 try: 

447 while True: 

448 self._recovery_rerun_requested = False 

449 async with self.get_listening_context(): 

450 if self._resync_settle_seconds > 0: 

451 await asyncio.sleep(self._resync_settle_seconds) 

452 await self._flush_outbound_buffer() 

453 await self._fire_reconnect() 

454 if not self._recovery_rerun_requested: 

455 break 

456 logger.info( 

457 "Gap arrived during recovery; rerunning flush + resync once" 

458 ) 

459 except asyncio.CancelledError: 

460 raise 

461 except Exception: 

462 logger.exception("Error during post-reconnect broadcast recovery") 

463 

464 async def _flush_outbound_buffer(self): 

465 """Replay buffered broadcasts to the backbone (best-effort). 

466 

467 The buffer is drained into a local snapshot under ``_buffer_lock`` and then 

468 published WITHOUT the lock, so a concurrent failed broadcast can still buffer 

469 (and a slow/hung publish can't wedge the buffer lock). Items that can no longer 

470 be serialized are dropped (so one poison payload can't wedge the buffer); a 

471 transport failure re-enqueues the unsent (older) tail ahead of any concurrent 

472 refill (newer) for the next recovery (see ``_requeue_unsent``). 

473 """ 

474 async with self._buffer_lock: 

475 if not self._outbound_buffer: 

476 return 

477 pending = list(self._outbound_buffer) 

478 self._outbound_buffer.clear() 

479 logger.info(f"Replaying {len(pending)} buffered broadcast(s) after recovery") 

480 unsent = list(pending) 

481 try: 

482 async with self._broadcast_type(self._broadcast_url) as channel: 

483 while unsent: 

484 topic, data = unsent[0] 

485 try: 

486 payload = pydantic_serialize( 

487 BroadcastNotification( 

488 notifier_id=self._id, topics=[topic], data=data 

489 ) 

490 ) 

491 except Exception as e: 

492 logger.error( 

493 f"Dropping un-serializable buffered broadcast on " 

494 f"topic '{topic}': {e!r}" 

495 ) 

496 unsent.pop(0) 

497 continue 

498 await channel.publish(self._channel, payload) 

499 unsent.pop(0) 

500 except asyncio.CancelledError: 

501 raise 

502 except Exception as e: 

503 logger.error( 

504 f"Failed to replay buffered broadcasts ({len(unsent)} left, will retry " 

505 f"on next recovery): {e!r}" 

506 ) 

507 if unsent: 

508 await self._requeue_unsent(unsent) 

509 

510 async def _requeue_unsent(self, unsent: list): 

511 """Re-enqueue the unsent (older) tail ahead of any concurrent refill 

512 (newer). 

513 

514 The publish above ran WITHOUT the lock, so concurrent failed broadcasts may have 

515 refilled the buffer meanwhile. Rebuild under the lock as ``unsent + refill`` and 

516 rely on the bounded deque dropping from the FRONT (oldest) on overflow — a plain 

517 ``extendleft`` would instead evict the newest refill, inverting the drop-oldest 

518 policy on a deque that is already at ``maxlen``. 

519 """ 

520 async with self._buffer_lock: 

521 refill = list(self._outbound_buffer) 

522 self._outbound_buffer.clear() 

523 self._outbound_buffer.extend(unsent) 

524 self._outbound_buffer.extend(refill) 

525 

526 async def _fire_reconnect(self): 

527 if self._on_reconnect is None: 

528 return 

529 try: 

530 await self._on_reconnect() 

531 except asyncio.CancelledError: 

532 raise 

533 except Exception: 

534 logger.exception("Broadcaster on_reconnect callback failed") 

535 

536 async def _fire_give_up(self): 

537 """Fire the give-up hook so OPAL can graceful-restart the worker. 

538 

539 Called only from the give-up (returning) path of the reader 

540 loop, never on cancellation, so a clean shutdown does not re- 

541 trigger the restart. If no hook is wired, log loudly that this 

542 worker now depends on the liveness probe. 

543 """ 

544 if self._on_give_up is None: 

545 logger.error( 

546 "Broadcaster gave up reconnecting and no give-up hook is wired; this " 

547 "worker now requires the liveness probe (/healthcheck) to be restarted" 

548 ) 

549 return 

550 try: 

551 await self._on_give_up() 

552 except asyncio.CancelledError: 

553 raise 

554 except Exception: 

555 logger.exception("Broadcaster on_give_up callback failed") 

556 

557 async def _ensure_connected(self): 

558 if self.listening_broadcast_channel is None: 

559 self.listening_broadcast_channel = self._broadcast_type(self._broadcast_url) 

560 await self.listening_broadcast_channel.connect() 

561 return self.listening_broadcast_channel 

562 

563 async def _handle_broadcast_event(self, event): 

564 """Forward one incoming broadcast to the internal notifier. 

565 

566 Mirrors the base class' per-event handling; kept here so the 

567 reconnect loop above stays readable. 

568 """ 

569 try: 

570 notification = BroadcastNotification.parse_raw(event.message) 

571 # Avoid re-publishing our own broadcasts 

572 if notification.notifier_id != self._id: 

573 logger.debug( 

574 "Handling incoming broadcast event: {}".format( 

575 {"topics": notification.topics, "src": notification.notifier_id} 

576 ) 

577 ) 

578 task = asyncio.create_task( 

579 self._notifier.notify( 

580 notification.topics, 

581 notification.data, 

582 notifier_id=self._id, 

583 ) 

584 ) 

585 self._tasks.add(task) 

586 task.add_done_callback(self._tasks.discard) 

587 except Exception: 

588 logger.exception("Failed handling incoming broadcast") 

589 

590 async def _safe_disconnect_channel(self): 

591 channel = self.listening_broadcast_channel 

592 self.listening_broadcast_channel = None 

593 if channel is not None: 

594 try: 

595 await channel.disconnect() 

596 except Exception: 

597 logger.debug("Error while disconnecting broadcast channel; ignoring") 

598 

599 def _backoff_seconds(self, attempt: int) -> float: 

600 if attempt <= 0: 

601 base = self._reconnect_backoff_min 

602 else: 

603 base = self._reconnect_backoff_min * (2 ** (attempt - 1)) 

604 base = min(base, self._reconnect_backoff_max) 

605 # Equal jitter, so a fleet of pods does not reconnect to the backbone in lockstep. 

606 return base / 2 + random.uniform(0, base / 2) 

607 

608 

609class FreezablePubSubEndpoint(PubSubEndpoint): 

610 """A ``PubSubEndpoint`` that *freezes* client-facing publishes during a 

611 broadcaster backbone GAP, to keep a multi-worker fleet consistent through 

612 an outage. 

613 

614 The problem: a server-side ``publish`` fans out two independent ways — local in-process 

615 delivery to *this* worker's own clients, and (via the broadcaster) to peer workers. Only 

616 the outbound path is buffered when the backbone is down; local delivery still fires. So a 

617 data/policy update that reaches one worker during a backbone gap is applied to that 

618 worker's clients but not the fleet — a transient split (some PDPs new, others old) that 

619 lasts the whole outage. 

620 

621 With freeze enabled, while the ``ReconnectingBroadcaster`` reports a real gap 

622 (``is_in_backbone_gap()`` — an established backbone session was lost and is being 

623 re-acquired; see its docstring for why "never connected" and "reader not running" must 

624 NOT freeze), ``publish`` is skipped entirely: neither local clients nor the outbound 

625 buffer see it. The write still lands in the source of truth, and the reconnect *resync* 

626 makes every worker's clients refetch on recovery — the whole fleet moves together. 

627 Skipping the whole publish also means nothing is buffered for replay during the freeze, 

628 so recovery converges purely via the resync refetch. 

629 

630 **Exempt topics** keep the pre-freeze behavior (local delivery + outbound replay buffer) 

631 even mid-gap: topics prefixed ``__`` (the statistics protocol and the broadcaster 

632 keepalive under their default names — dropping those corrupts server-to-server state 

633 that no resync rebuilds: ghost clients, workers that never stat-sync) and any topic in 

634 ``freeze_exempt_topics``. OPAL passes the git-webhook trigger topic (it targets the 

635 server-side policy watcher, not clients, and a dropped trigger means the repo pull it 

636 requests simply never happens — the resync would then refetch from a clone that was 

637 never advanced) plus the *configured* statistics/keepalive channel names, since those 

638 are operator-overridable and the ``__`` prefix rule only covers the defaults. 

639 

640 Delegates straight to the base when: freeze is disabled; there is no broadcaster 

641 (single worker); or the broadcaster is the stock ``EventBroadcaster``. 

642 

643 **Recovery scope** (what "reconciled by the resync" actually covers): data the clients 

644 re-fetch on reconnect, i.e. their configured data sources (``OPAL_DATA_CONFIG_SOURCES`` 

645 or scope config) and the policy bundle. One-off updates outside that set — an inline 

646 ``data`` payload, or a fetch URL that is not part of the configured sources — are 

647 dropped by a freeze, not deferred. Accepted trade: consistency over freshness, and such 

648 updates are the legacy path. 

649 

650 **Known limitations** (all degrade to the PRE-freeze behavior, never worse): 

651 * If the reader's subscription is alive but an individual outbound broadcast fails 

652 (separate per-publish channel), the gate does not engage — that publish is delivered 

653 locally and buffered for replay, the pre-existing split-until-replay behavior. 

654 * The gate reopens when the subscription is re-established, before the session proves 

655 "sustained" — during a rare connect-then-instant-close flap a publish can slip 

656 through (deliver locally + buffer). Gating on sustained instead would wrongly freeze 

657 quiet channels forever (a session only proves sustained on its first inbound event). 

658 * Client-originated RPC publishes (``RpcEventServerMethods.publish``) notify the local 

659 subscribers directly, bypassing this override. 

660 """ 

661 

662 def __init__( 

663 self, 

664 *args, 

665 freeze_on_disconnect: bool = True, 

666 freeze_exempt_topics=(), 

667 **kwargs, 

668 ): 

669 super().__init__(*args, **kwargs) 

670 self._freeze_on_disconnect = freeze_on_disconnect 

671 self._freeze_exempt_topics = frozenset(freeze_exempt_topics) 

672 # Publishes suppressed in the current freeze episode — first one logs at WARNING, 

673 # the rest at DEBUG (a long outage would otherwise emit an unbounded WARNING per 

674 # frozen stats keepalive), and the first delivered publish afterwards logs a 

675 # summary count. 

676 self._frozen_in_episode = 0 

677 # Which backbone gap (backbone_gap_generation()) the open episode belongs to: 

678 # a gap can end and a NEW one open before any out-of-gap publish is delivered 

679 # (recovery itself never publishes), so publish() flushes the previous gap's 

680 # pending summary when it sees a frozen publish from a different generation. 

681 self._frozen_gap_generation: Optional[int] = None 

682 

683 def _is_exempt(self, topics) -> bool: 

684 if isinstance(topics, str): 

685 topics = [topics] 

686 if not topics: 

687 # all([]) is True — an empty (or None) topic list must not slip past the 

688 # gate as "exempt". 

689 return False 

690 return all( 

691 topic.startswith("__") or topic in self._freeze_exempt_topics 

692 for topic in topics 

693 ) 

694 

695 def _should_freeze(self, topics) -> bool: 

696 return self._in_frozen_gap() and not self._is_exempt(topics) 

697 

698 def _in_frozen_gap(self) -> bool: 

699 broadcaster = self.broadcaster 

700 return ( 

701 self._freeze_on_disconnect 

702 and isinstance(broadcaster, ReconnectingBroadcaster) 

703 and broadcaster.is_in_backbone_gap() 

704 ) 

705 

706 async def publish(self, topics, data=None): 

707 in_gap = self._in_frozen_gap() 

708 if self._should_freeze(topics): 708 ↛ 712line 708 didn't jump to line 712 because the condition on line 708 was never true

709 # A new gap can open before the previous episode's summary got out 

710 # (recovery itself never publishes) — flush it first, so each gap gets 

711 # its own leading WARNING and its own summary count. 

712 generation = self.broadcaster.backbone_gap_generation() 

713 if self._frozen_in_episode and generation != self._frozen_gap_generation: 

714 self._log_episode_summary() 

715 self._frozen_gap_generation = generation 

716 self._frozen_in_episode += 1 

717 log = logger.warning if self._frozen_in_episode == 1 else logger.debug 

718 log( 

719 "Broadcaster backbone gap; freezing publish to preserve fleet consistency " 

720 "(not delivered to clients; reconciled via resync on reconnect). " 

721 "topics={topics} (suppressed {count} so far this gap)", 

722 topics=topics, 

723 count=self._frozen_in_episode, 

724 ) 

725 return 

726 # Emit the episode summary only once the gap is actually over — an EXEMPT publish 

727 # mid-gap also reaches this point and must not reset the counter or claim recovery. 

728 if self._frozen_in_episode and not in_gap: 728 ↛ 729line 728 didn't jump to line 729 because the condition on line 728 was never true

729 self._log_episode_summary() 

730 return await super().publish(topics, data) 

731 

732 def _log_episode_summary(self): 

733 count, self._frozen_in_episode = self._frozen_in_episode, 0 

734 logger.warning( 

735 "Backbone recovered; froze {count} publish(es) during the gap — clients " 

736 "reconcile via the reconnect resync", 

737 count=count, 

738 ) 

739 

740 # The library aliases ``notify = publish`` at class level (backward-compat canonical 

741 # name), which binds the BASE publish — re-bind it here or ``endpoint.notify(...)`` 

742 # would silently bypass the freeze gate. 

743 notify = publish