use std::sync::Arc;
use fraiseql_core::runtime::subscription::{
ChangeSpineEnvelope, SubscriptionEvent, SubscriptionManager, SubscriptionOperation,
};
use tokio::sync::{broadcast, mpsc};
use tracing::{debug, info};
#[derive(Debug, Clone, Copy)]
pub struct EventBridgeConfig {
pub channel_capacity: usize,
}
impl EventBridgeConfig {
#[must_use]
pub const fn new() -> Self {
Self {
channel_capacity: 100,
}
}
#[must_use]
pub const fn with_channel_capacity(mut self, capacity: usize) -> Self {
self.channel_capacity = capacity;
self
}
}
impl Default for EventBridgeConfig {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, Clone)]
pub struct EntityEvent {
pub entity_type: String,
pub entity_id: String,
pub operation: SubscriptionOperation,
pub data: serde_json::Value,
pub old_data: Option<serde_json::Value>,
pub tenant_id: Option<String>,
pub change_spine: Option<ChangeSpineEnvelope>,
}
impl EntityEvent {
#[must_use]
pub fn new(
entity_type: impl Into<String>,
entity_id: impl Into<String>,
operation: SubscriptionOperation,
data: serde_json::Value,
) -> Self {
Self {
entity_type: entity_type.into(),
entity_id: entity_id.into(),
operation,
data,
old_data: None,
tenant_id: None,
change_spine: None,
}
}
#[must_use]
pub fn with_old_data(mut self, old_data: serde_json::Value) -> Self {
self.old_data = Some(old_data);
self
}
#[must_use]
pub fn with_tenant_id(mut self, tenant_id: impl Into<String>) -> Self {
self.tenant_id = Some(tenant_id.into());
self
}
#[must_use]
pub fn with_change_spine(mut self, envelope: ChangeSpineEnvelope) -> Self {
self.change_spine = Some(envelope);
self
}
}
pub const DEFAULT_ENTITY_FANOUT_CAPACITY: usize = 256;
#[derive(Clone, Debug)]
pub struct EntityEventFanout {
sender: broadcast::Sender<EntityEvent>,
}
impl EntityEventFanout {
#[must_use]
pub fn new(capacity: usize) -> Self {
let (sender, _) = broadcast::channel(capacity);
Self { sender }
}
#[must_use]
pub fn subscribe(&self) -> broadcast::Receiver<EntityEvent> {
self.sender.subscribe()
}
#[must_use]
pub fn receiver_count(&self) -> usize {
self.sender.receiver_count()
}
#[must_use]
pub fn publish(&self, event: &EntityEvent) -> usize {
self.sender.send(event.clone()).unwrap_or(0)
}
}
impl Default for EntityEventFanout {
fn default() -> Self {
Self::new(DEFAULT_ENTITY_FANOUT_CAPACITY)
}
}
pub struct EventBridge {
manager: Arc<SubscriptionManager>,
receiver: mpsc::Receiver<EntityEvent>,
sender: mpsc::Sender<EntityEvent>,
entity_fanout: Option<EntityEventFanout>,
}
impl EventBridge {
#[must_use]
pub fn new(manager: Arc<SubscriptionManager>, config: EventBridgeConfig) -> Self {
let (sender, receiver) = mpsc::channel(config.channel_capacity);
Self {
manager,
receiver,
sender,
entity_fanout: None,
}
}
#[must_use]
pub fn with_entity_fanout(mut self, fanout: EntityEventFanout) -> Self {
self.entity_fanout = Some(fanout);
self
}
#[must_use]
pub fn sender(&self) -> mpsc::Sender<EntityEvent> {
self.sender.clone()
}
#[must_use]
pub fn convert_event(entity_event: EntityEvent) -> SubscriptionEvent {
let mut event = SubscriptionEvent::new(
entity_event.entity_type,
entity_event.entity_id,
entity_event.operation,
entity_event.data,
);
if let Some(old_data) = entity_event.old_data {
event = event.with_old_data(old_data);
}
if let Some(tenant_id) = entity_event.tenant_id {
event = event.with_tenant_id(tenant_id);
}
if let Some(envelope) = entity_event.change_spine {
event = event.with_change_spine(envelope);
}
event
}
#[allow(clippy::cognitive_complexity)] pub async fn run(self) {
let Self {
manager,
mut receiver,
sender,
entity_fanout,
} = self;
drop(sender);
info!("EventBridge started");
while let Some(entity_event) = receiver.recv().await {
debug!("EventBridge received entity event: {}", entity_event.entity_type);
if let Some(ref fanout) = entity_fanout {
let delivered = fanout.publish(&entity_event);
if delivered > 0 {
debug!(
entity_type = %entity_event.entity_type,
receivers = delivered,
"EventBridge fanned out entity event"
);
}
}
let subscription_event = Self::convert_event(entity_event);
let matched = manager.publish_event(subscription_event);
if matched > 0 {
debug!("EventBridge matched {} subscriptions", matched);
}
}
info!("EventBridge stopped");
}
#[must_use = "dropping the JoinHandle detaches the task; store or abort it to control lifecycle"]
pub fn spawn(self) -> tokio::task::JoinHandle<()> {
tokio::spawn(self.run())
}
#[must_use]
pub fn get_sender(&self) -> mpsc::Sender<EntityEvent> {
self.sender.clone()
}
#[must_use]
pub fn manager(&self) -> Arc<SubscriptionManager> {
Arc::clone(&self.manager)
}
}