Skip to main content

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}