fraiseql_server/subscriptions/
event_bridge.rs1use 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#[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 pub change_spine: Option<ChangeSpineEnvelope>,
85}
86
87impl EntityEvent {
88 #[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 #[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 #[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 #[must_use]
123 pub fn with_change_spine(mut self, envelope: ChangeSpineEnvelope) -> Self {
124 self.change_spine = Some(envelope);
125 self
126 }
127}
128
129pub struct EventBridge {
131 manager: Arc<SubscriptionManager>,
133
134 receiver: mpsc::Receiver<EntityEvent>,
136
137 sender: mpsc::Sender<EntityEvent>,
139}
140
141impl EventBridge {
142 #[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 #[must_use]
156 pub fn sender(&self) -> mpsc::Sender<EntityEvent> {
157 self.sender.clone()
158 }
159
160 pub fn convert_event(entity_event: EntityEvent) -> SubscriptionEvent {
162 let operation = match entity_event.operation.to_uppercase().as_str() {
164 "INSERT" => SubscriptionOperation::Create,
165 "UPDATE" => SubscriptionOperation::Update,
166 "DELETE" => SubscriptionOperation::Delete,
167 _ => {
168 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 if let Some(old_data) = entity_event.old_data {
183 event = event.with_old_data(old_data);
184 }
185
186 if let Some(tenant_id) = entity_event.tenant_id {
188 event = event.with_tenant_id(tenant_id);
189 }
190
191 if let Some(envelope) = entity_event.change_spine {
193 event = event.with_change_spine(envelope);
194 }
195
196 event
197 }
198
199 #[allow(clippy::cognitive_complexity)] 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 let subscription_event = Self::convert_event(entity_event);
209
210 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 #[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 #[must_use]
234 pub fn get_sender(&self) -> mpsc::Sender<EntityEvent> {
235 self.sender.clone()
236 }
237
238 #[must_use]
240 pub fn manager(&self) -> Arc<SubscriptionManager> {
241 Arc::clone(&self.manager)
242 }
243}