fraiseql_server/subscriptions/
event_bridge.rs1use std::sync::Arc;
23
24use fraiseql_core::runtime::subscription::{
25 SubscriptionEvent, SubscriptionManager, SubscriptionOperation,
26};
27use tokio::sync::mpsc;
28use tracing::{debug, info};
29
30#[derive(Debug, Clone, Copy)]
32pub struct EventBridgeConfig {
33 pub channel_capacity: usize,
35}
36
37impl EventBridgeConfig {
38 #[must_use]
40 pub const fn new() -> Self {
41 Self {
42 channel_capacity: 100,
43 }
44 }
45
46 #[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#[derive(Debug, Clone)]
62pub struct EntityEvent {
63 pub entity_type: String,
65
66 pub entity_id: String,
68
69 pub operation: String,
71
72 pub data: serde_json::Value,
74
75 pub old_data: Option<serde_json::Value>,
77
78 pub tenant_id: Option<String>,
80}
81
82impl EntityEvent {
83 #[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 #[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 #[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
116pub struct EventBridge {
118 manager: Arc<SubscriptionManager>,
120
121 receiver: mpsc::Receiver<EntityEvent>,
123
124 sender: mpsc::Sender<EntityEvent>,
126}
127
128impl EventBridge {
129 #[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 #[must_use]
143 pub fn sender(&self) -> mpsc::Sender<EntityEvent> {
144 self.sender.clone()
145 }
146
147 pub fn convert_event(entity_event: EntityEvent) -> SubscriptionEvent {
149 let operation = match entity_event.operation.to_uppercase().as_str() {
151 "INSERT" => SubscriptionOperation::Create,
152 "UPDATE" => SubscriptionOperation::Update,
153 "DELETE" => SubscriptionOperation::Delete,
154 _ => {
155 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 if let Some(old_data) = entity_event.old_data {
170 event = event.with_old_data(old_data);
171 }
172
173 if let Some(tenant_id) = entity_event.tenant_id {
175 event = event.with_tenant_id(tenant_id);
176 }
177
178 event
179 }
180
181 #[allow(clippy::cognitive_complexity)] 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 let subscription_event = Self::convert_event(entity_event);
191
192 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 #[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 #[must_use]
216 pub fn get_sender(&self) -> mpsc::Sender<EntityEvent> {
217 self.sender.clone()
218 }
219
220 #[must_use]
222 pub fn manager(&self) -> Arc<SubscriptionManager> {
223 Arc::clone(&self.manager)
224 }
225}