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
« 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.
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:
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')``.
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:
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.
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
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
52ReconnectCallback = Callable[[], Awaitable[None]]
55class SafeConnectionManager(ConnectionManager):
56 """A ``ConnectionManager`` whose ``disconnect`` is idempotent, able to drop
57 its connections on demand.
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 """
69 def disconnect(self, websocket: WebSocket):
70 try:
71 self.active_connections.remove(websocket)
72 except ValueError:
73 logger.debug("Ignoring duplicate websocket disconnect")
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.
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.
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
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.
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 """
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
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
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``.
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
190 async def start_reader_task(self):
191 """Spawn the reconnecting reader task once.
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
212 def is_reader_healthy(self) -> bool:
213 """Report whether the reconnecting reader can still serve connected
214 clients.
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.
220 Health is judged against the listener count, not the backbone connection:
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.
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``.
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 )
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).
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
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.
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:
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 )
284 async def __broadcast_notifications__(self, subscription: Subscription, data):
285 """Share a local notification with the backbone; buffer it on failure.
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)
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 )
318 async def __read_notifications__(self):
319 """Read incoming broadcasts, reconnecting on backbone disconnect.
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))
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
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)
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)
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).
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.
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")
464 async def _flush_outbound_buffer(self):
465 """Replay buffered broadcasts to the backbone (best-effort).
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)
510 async def _requeue_unsent(self, unsent: list):
511 """Re-enqueue the unsent (older) tail ahead of any concurrent refill
512 (newer).
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)
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")
536 async def _fire_give_up(self):
537 """Fire the give-up hook so OPAL can graceful-restart the worker.
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")
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
563 async def _handle_broadcast_event(self, event):
564 """Forward one incoming broadcast to the internal notifier.
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")
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")
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)
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.
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.
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.
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.
640 Delegates straight to the base when: freeze is disabled; there is no broadcaster
641 (single worker); or the broadcaster is the stock ``EventBroadcaster``.
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.
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 """
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
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 )
695 def _should_freeze(self, topics) -> bool:
696 return self._in_frozen_gap() and not self._is_exempt(topics)
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 )
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)
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 )
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