Skip to main content

fraiseql_server/subscriptions/
event_bridge.rs

1//! `EventBridge` that connects `ChangeLogListener` with `SubscriptionManager`.
2//!
3//! The `EventBridge` is responsible for:
4//! 1. Spawning `ChangeLogListener` in background
5//! 2. Receiving `EntityEvent` via `mpsc::channel`
6//! 3. Converting `EntityEvent` to `SubscriptionEvent`
7//! 4. Publishing events to `SubscriptionManager`
8//!
9//! Architecture:
10//! ```text
11//! Database (tb_entity_change_log)
12//!     ↓
13//! ChangeLogListener (polls & converts)
14//!     ↓
15//! EventBridge (routes & converts)
16//!     ↓
17//! SubscriptionManager (broadcasts to subscribers)
18//!     ↓
19//! WebSocket Handler (delivers to clients)
20//! ```
21
22use std::sync::Arc;
23
24use fraiseql_core::runtime::subscription::{
25    ChangeSpineEnvelope, SubscriptionEvent, SubscriptionManager, SubscriptionOperation,
26};
27use tokio::sync::mpsc;
28use tracing::{debug, info};
29
30/// Configuration for the `EventBridge`
31#[derive(Debug, Clone, Copy)]
32pub struct EventBridgeConfig {
33    /// Channel capacity for event routing
34    pub channel_capacity: usize,
35}
36
37impl EventBridgeConfig {
38    /// Create config with defaults
39    #[must_use]
40    pub const fn new() -> Self {
41        Self {
42            channel_capacity: 100,
43        }
44    }
45
46    /// Set channel capacity
47    #[must_use]
48    pub const fn with_channel_capacity(mut self, capacity: usize) -> Self {
49        self.channel_capacity = capacity;
50        self
51    }
52}
53
54impl Default for EventBridgeConfig {
55    fn default() -> Self {
56        Self::new()
57    }
58}
59
60/// A simple event that `EventBridge` receives from `ChangeLogListener`
61#[derive(Debug, Clone)]
62pub struct EntityEvent {
63    /// Entity type (e.g., "Order", "User")
64    pub entity_type: String,
65
66    /// Entity ID (primary key)
67    pub entity_id: String,
68
69    /// Operation type ("INSERT", "UPDATE", "DELETE")
70    pub operation: String,
71
72    /// Entity data as JSON
73    pub data: serde_json::Value,
74
75    /// Optional old data (for UPDATE operations)
76    pub old_data: Option<serde_json::Value>,
77
78    /// Tenant identifier for multi-tenant filtering (`fk_customer_org`).
79    pub tenant_id: Option<String>,
80
81    /// Change-Spine envelope metadata for client delivery (#425). Propagated
82    /// through to the `SubscriptionEvent` and emitted in the `next` payload's
83    /// `extensions.changeSpine`; not used for filtering.
84    pub change_spine: Option<ChangeSpineEnvelope>,
85}
86
87impl EntityEvent {
88    /// Create a new entity event
89    #[must_use]
90    pub fn new(
91        entity_type: impl Into<String>,
92        entity_id: impl Into<String>,
93        operation: impl Into<String>,
94        data: serde_json::Value,
95    ) -> Self {
96        Self {
97            entity_type: entity_type.into(),
98            entity_id: entity_id.into(),
99            operation: operation.into(),
100            data,
101            old_data: None,
102            tenant_id: None,
103            change_spine: None,
104        }
105    }
106
107    /// Add old data for UPDATE operations
108    #[must_use]
109    pub fn with_old_data(mut self, old_data: serde_json::Value) -> Self {
110        self.old_data = Some(old_data);
111        self
112    }
113
114    /// Set tenant identifier for multi-tenant filtering.
115    #[must_use]
116    pub fn with_tenant_id(mut self, tenant_id: impl Into<String>) -> Self {
117        self.tenant_id = Some(tenant_id.into());
118        self
119    }
120
121    /// Attach the Change-Spine envelope for client delivery (#425).
122    #[must_use]
123    pub fn with_change_spine(mut self, envelope: ChangeSpineEnvelope) -> Self {
124        self.change_spine = Some(envelope);
125        self
126    }
127}
128
129/// `EventBridge` that connects `ChangeLogListener` with `SubscriptionManager`
130pub struct EventBridge {
131    /// Subscription manager for broadcasting events
132    manager: Arc<SubscriptionManager>,
133
134    /// Receiver for entity events from `ChangeLogListener`
135    receiver: mpsc::Receiver<EntityEvent>,
136
137    /// Sender for entity events (used to send events to bridge)
138    sender: mpsc::Sender<EntityEvent>,
139}
140
141impl EventBridge {
142    /// Create a new `EventBridge`
143    #[must_use]
144    pub fn new(manager: Arc<SubscriptionManager>, config: EventBridgeConfig) -> Self {
145        let (sender, receiver) = mpsc::channel(config.channel_capacity);
146
147        Self {
148            manager,
149            receiver,
150            sender,
151        }
152    }
153
154    /// Get a sender for publishing entity events
155    #[must_use]
156    pub fn sender(&self) -> mpsc::Sender<EntityEvent> {
157        self.sender.clone()
158    }
159
160    /// Convert `EntityEvent` to `SubscriptionEvent`
161    pub fn convert_event(entity_event: EntityEvent) -> SubscriptionEvent {
162        // Convert operation string to SubscriptionOperation
163        let operation = match entity_event.operation.to_uppercase().as_str() {
164            "INSERT" => SubscriptionOperation::Create,
165            "UPDATE" => SubscriptionOperation::Update,
166            "DELETE" => SubscriptionOperation::Delete,
167            _ => {
168                // Default to Create for unknown operations
169                debug!("Unknown operation: {}, defaulting to Create", entity_event.operation);
170                SubscriptionOperation::Create
171            },
172        };
173
174        let mut event = SubscriptionEvent::new(
175            entity_event.entity_type,
176            entity_event.entity_id,
177            operation,
178            entity_event.data,
179        );
180
181        // Add old data if present
182        if let Some(old_data) = entity_event.old_data {
183            event = event.with_old_data(old_data);
184        }
185
186        // Propagate tenant_id for multi-tenant filtering
187        if let Some(tenant_id) = entity_event.tenant_id {
188            event = event.with_tenant_id(tenant_id);
189        }
190
191        // Propagate the Change-Spine envelope for client delivery (#425)
192        if let Some(envelope) = entity_event.change_spine {
193            event = event.with_change_spine(envelope);
194        }
195
196        event
197    }
198
199    /// Run the event bridge loop (spawned in background)
200    #[allow(clippy::cognitive_complexity)] // Reason: event loop with multi-source message routing and reconnection handling
201    pub async fn run(mut self) {
202        info!("EventBridge started");
203
204        while let Some(entity_event) = self.receiver.recv().await {
205            debug!("EventBridge received entity event: {}", entity_event.entity_type);
206
207            // Convert entity event to subscription event
208            let subscription_event = Self::convert_event(entity_event);
209
210            // Publish to subscription manager
211            let matched = self.manager.publish_event(subscription_event);
212
213            if matched > 0 {
214                debug!("EventBridge matched {} subscriptions", matched);
215            }
216        }
217
218        info!("EventBridge stopped");
219    }
220
221    /// Spawn `EventBridge` as a background task.
222    ///
223    /// Returns a `JoinHandle` that must not be silently dropped — callers
224    /// should either `.await` it for a clean shutdown or explicitly `.abort()`
225    /// it when the bridge is no longer needed.  Dropping the handle detaches
226    /// the task, making it impossible to observe panics or coordinate shutdown.
227    #[must_use = "dropping the JoinHandle detaches the task; store or abort it to control lifecycle"]
228    pub fn spawn(self) -> tokio::task::JoinHandle<()> {
229        tokio::spawn(self.run())
230    }
231
232    /// Get the sender for sending events to the bridge
233    #[must_use]
234    pub fn get_sender(&self) -> mpsc::Sender<EntityEvent> {
235        self.sender.clone()
236    }
237
238    /// Get the subscription manager (for testing)
239    #[must_use]
240    pub fn manager(&self) -> Arc<SubscriptionManager> {
241        Arc::clone(&self.manager)
242    }
243}