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    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
82impl EntityEvent {
83    /// Create a new entity event
84    #[must_use]
85    pub fn new(
86        entity_type: impl Into<String>,
87        entity_id: impl Into<String>,
88        operation: impl Into<String>,
89        data: serde_json::Value,
90    ) -> Self {
91        Self {
92            entity_type: entity_type.into(),
93            entity_id: entity_id.into(),
94            operation: operation.into(),
95            data,
96            old_data: None,
97            tenant_id: None,
98        }
99    }
100
101    /// Add old data for UPDATE operations
102    #[must_use]
103    pub fn with_old_data(mut self, old_data: serde_json::Value) -> Self {
104        self.old_data = Some(old_data);
105        self
106    }
107
108    /// Set tenant identifier for multi-tenant filtering.
109    #[must_use]
110    pub fn with_tenant_id(mut self, tenant_id: impl Into<String>) -> Self {
111        self.tenant_id = Some(tenant_id.into());
112        self
113    }
114}
115
116/// `EventBridge` that connects `ChangeLogListener` with `SubscriptionManager`
117pub struct EventBridge {
118    /// Subscription manager for broadcasting events
119    manager: Arc<SubscriptionManager>,
120
121    /// Receiver for entity events from `ChangeLogListener`
122    receiver: mpsc::Receiver<EntityEvent>,
123
124    /// Sender for entity events (used to send events to bridge)
125    sender: mpsc::Sender<EntityEvent>,
126}
127
128impl EventBridge {
129    /// Create a new `EventBridge`
130    #[must_use]
131    pub fn new(manager: Arc<SubscriptionManager>, config: EventBridgeConfig) -> Self {
132        let (sender, receiver) = mpsc::channel(config.channel_capacity);
133
134        Self {
135            manager,
136            receiver,
137            sender,
138        }
139    }
140
141    /// Get a sender for publishing entity events
142    #[must_use]
143    pub fn sender(&self) -> mpsc::Sender<EntityEvent> {
144        self.sender.clone()
145    }
146
147    /// Convert `EntityEvent` to `SubscriptionEvent`
148    pub fn convert_event(entity_event: EntityEvent) -> SubscriptionEvent {
149        // Convert operation string to SubscriptionOperation
150        let operation = match entity_event.operation.to_uppercase().as_str() {
151            "INSERT" => SubscriptionOperation::Create,
152            "UPDATE" => SubscriptionOperation::Update,
153            "DELETE" => SubscriptionOperation::Delete,
154            _ => {
155                // Default to Create for unknown operations
156                debug!("Unknown operation: {}, defaulting to Create", entity_event.operation);
157                SubscriptionOperation::Create
158            },
159        };
160
161        let mut event = SubscriptionEvent::new(
162            entity_event.entity_type,
163            entity_event.entity_id,
164            operation,
165            entity_event.data,
166        );
167
168        // Add old data if present
169        if let Some(old_data) = entity_event.old_data {
170            event = event.with_old_data(old_data);
171        }
172
173        // Propagate tenant_id for multi-tenant filtering
174        if let Some(tenant_id) = entity_event.tenant_id {
175            event = event.with_tenant_id(tenant_id);
176        }
177
178        event
179    }
180
181    /// Run the event bridge loop (spawned in background)
182    #[allow(clippy::cognitive_complexity)] // Reason: event loop with multi-source message routing and reconnection handling
183    pub async fn run(mut self) {
184        info!("EventBridge started");
185
186        while let Some(entity_event) = self.receiver.recv().await {
187            debug!("EventBridge received entity event: {}", entity_event.entity_type);
188
189            // Convert entity event to subscription event
190            let subscription_event = Self::convert_event(entity_event);
191
192            // Publish to subscription manager
193            let matched = self.manager.publish_event(subscription_event);
194
195            if matched > 0 {
196                debug!("EventBridge matched {} subscriptions", matched);
197            }
198        }
199
200        info!("EventBridge stopped");
201    }
202
203    /// Spawn `EventBridge` as a background task.
204    ///
205    /// Returns a `JoinHandle` that must not be silently dropped — callers
206    /// should either `.await` it for a clean shutdown or explicitly `.abort()`
207    /// it when the bridge is no longer needed.  Dropping the handle detaches
208    /// the task, making it impossible to observe panics or coordinate shutdown.
209    #[must_use = "dropping the JoinHandle detaches the task; store or abort it to control lifecycle"]
210    pub fn spawn(self) -> tokio::task::JoinHandle<()> {
211        tokio::spawn(self.run())
212    }
213
214    /// Get the sender for sending events to the bridge
215    #[must_use]
216    pub fn get_sender(&self) -> mpsc::Sender<EntityEvent> {
217        self.sender.clone()
218    }
219
220    /// Get the subscription manager (for testing)
221    #[must_use]
222    pub fn manager(&self) -> Arc<SubscriptionManager> {
223        Arc::clone(&self.manager)
224    }
225}