fraiseql_server/realtime/observer.rs
1//! `RealtimeBroadcastObserver` — bridge from mutation events to the realtime delivery pipeline.
2//!
3//! `RealtimeBroadcastObserver.on_mutation_complete()` is called on the mutation path and
4//! returns immediately (non-blocking), handing the event off to a bounded channel that
5//! the `EventDeliveryPipeline` consumes in a background task.
6//!
7//! If the bounded channel is full (delivery pipeline under backpressure), the event is
8//! dropped and a metric counter is incremented. This ensures mutation response latency
9//! is never affected by the realtime delivery path.
10
11use std::sync::atomic::{AtomicU64, Ordering};
12
13use tokio::sync::mpsc;
14
15use super::delivery::EntityEvent;
16
17/// Observer that forwards entity change events to the realtime delivery pipeline.
18///
19/// Create with [`RealtimeBroadcastObserver::new`], wire the returned receiver into an
20/// [`super::delivery::EventDeliveryPipeline`], then call
21/// [`on_mutation_complete`][RealtimeBroadcastObserver::on_mutation_complete] from the
22/// mutation path.
23pub struct RealtimeBroadcastObserver {
24 /// Sender half of the observer-to-pipeline channel.
25 event_tx: mpsc::Sender<EntityEvent>,
26 /// Count of events dropped due to backpressure (channel full).
27 ///
28 /// Maps to the `realtime_events_dropped_backpressure_total` metric.
29 events_dropped: AtomicU64,
30}
31
32impl RealtimeBroadcastObserver {
33 /// Create a new observer and its corresponding event receiver.
34 ///
35 /// The `capacity` controls how many events can be buffered before backpressure
36 /// causes events to be dropped. Pass the receiver to an
37 /// [`super::delivery::EventDeliveryPipeline`].
38 #[must_use]
39 pub fn new(capacity: usize) -> (Self, mpsc::Receiver<EntityEvent>) {
40 let (tx, rx) = mpsc::channel(capacity);
41 (
42 Self {
43 event_tx: tx,
44 events_dropped: AtomicU64::new(0),
45 },
46 rx,
47 )
48 }
49
50 /// Called when a mutation completes. Non-blocking.
51 ///
52 /// Tries to enqueue the event on the delivery-pipeline channel. If the channel
53 /// is full (pipeline under backpressure), the event is dropped and
54 /// `realtime_events_dropped_backpressure_total` is incremented. This keeps the
55 /// mutation response path free from realtime delivery latency.
56 pub fn on_mutation_complete(&self, event: EntityEvent) {
57 if self.event_tx.try_send(event).is_err() {
58 // Channel full — drop event and track for observability.
59 // Metric: realtime_events_dropped_backpressure_total
60 self.events_dropped.fetch_add(1, Ordering::Relaxed);
61 }
62 }
63
64 /// Total number of events dropped due to delivery-pipeline backpressure.
65 ///
66 /// Used by metrics exporters and health checks.
67 #[must_use]
68 pub fn events_dropped_total(&self) -> u64 {
69 self.events_dropped.load(Ordering::Relaxed)
70 }
71}