use std::sync::Arc;
use fraiseql_core::runtime::{
SubscriptionManager, TransportAdapter,
subscription::{KafkaAdapter, KafkaConfig},
};
use tokio::sync::broadcast::error::RecvError;
use tracing::{error, info, warn};
use crate::server_config::SubscriptionKafkaConfig;
#[derive(Debug)]
pub struct SubscriptionKafkaMirror {
adapter: Arc<KafkaAdapter>,
topic: String,
}
pub fn build_mirror(
config: Option<&SubscriptionKafkaConfig>,
) -> Result<Option<SubscriptionKafkaMirror>, String> {
let Some(config) = config else {
return Ok(None);
};
config.validate()?;
let mut adapter_config = KafkaConfig::new(&config.endpoint, &config.default_topic)
.with_client_id(&config.client_id)
.with_timeout(config.timeout_ms);
if let Some(ref compression) = config.compression {
adapter_config = adapter_config.with_compression(compression);
}
let adapter =
KafkaAdapter::new(adapter_config).map_err(|e| format!("[subscription_kafka] {e}"))?;
info!(
topic = %config.default_topic,
"subscription Kafka mirror configured"
);
Ok(Some(SubscriptionKafkaMirror {
adapter: Arc::new(adapter),
topic: config.default_topic.clone(),
}))
}
pub fn spawn(
mirror: SubscriptionKafkaMirror,
manager: &Arc<SubscriptionManager>,
tasks: &mut tokio::task::JoinSet<()>,
) {
let mut events = manager.receiver();
let SubscriptionKafkaMirror { adapter, topic } = mirror;
tasks.spawn(async move {
info!(topic = %topic, "subscription Kafka mirror started");
loop {
match events.recv().await {
Ok(payload) => {
let _ = adapter.deliver(&payload.event, &payload.subscription_name).await;
},
Err(RecvError::Lagged(missed)) => {
error!(
missed,
topic = %topic,
"subscription Kafka mirror fell behind; those deliveries were \
dropped by the broadcast channel and are not recoverable"
);
},
Err(RecvError::Closed) => {
warn!(topic = %topic, "subscription channel closed; Kafka mirror stopping");
break;
},
}
}
});
}
#[cfg(test)]
mod tests;