#[cfg(test)]
mod tests;
use std::time::Duration;
use fraiseql_cdc_sinks::{
CdcSink, CdcSinkConfig, ChangeEvent, DrainWorker, NatsJetStreamSink, PublishOutcome, SinkKind,
};
use sqlx::PgPool;
use tracing::{error, info, warn};
use crate::server_config::cdc_outbound::{CdcOutboundConfig, CdcSinkSectionConfig};
#[allow(clippy::large_enum_variant)]
pub enum ConfiguredSink {
NatsJetStream(NatsJetStreamSink),
#[cfg(feature = "cdc-kafka")]
Kafka(fraiseql_cdc_sinks::KafkaSink),
#[cfg(feature = "cdc-kinesis")]
Kinesis(fraiseql_cdc_sinks::KinesisSink),
}
impl CdcSink for ConfiguredSink {
fn name(&self) -> &str {
match self {
Self::NatsJetStream(sink) => sink.name(),
#[cfg(feature = "cdc-kafka")]
Self::Kafka(sink) => sink.name(),
#[cfg(feature = "cdc-kinesis")]
Self::Kinesis(sink) => sink.name(),
}
}
fn kind(&self) -> SinkKind {
match self {
Self::NatsJetStream(sink) => sink.kind(),
#[cfg(feature = "cdc-kafka")]
Self::Kafka(sink) => sink.kind(),
#[cfg(feature = "cdc-kinesis")]
Self::Kinesis(sink) => sink.kind(),
}
}
fn matches(&self, ev: &ChangeEvent) -> bool {
match self {
Self::NatsJetStream(sink) => sink.matches(ev),
#[cfg(feature = "cdc-kafka")]
Self::Kafka(sink) => sink.matches(ev),
#[cfg(feature = "cdc-kinesis")]
Self::Kinesis(sink) => sink.matches(ev),
}
}
async fn publish(&self, ev: &ChangeEvent) -> PublishOutcome {
match self {
Self::NatsJetStream(sink) => sink.publish(ev).await,
#[cfg(feature = "cdc-kafka")]
Self::Kafka(sink) => sink.publish(ev).await,
#[cfg(feature = "cdc-kinesis")]
Self::Kinesis(sink) => sink.publish(ev).await,
}
}
}
pub struct SinkDrain {
name: String,
worker: DrainWorker<ConfiguredSink>,
tick: Duration,
}
impl SinkDrain {
pub async fn run_forever(self) {
let mut ticker = tokio::time::interval(self.tick);
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
info!(sink = %self.name, tick_secs = self.tick.as_secs(), "cdc outbound drain started");
loop {
ticker.tick().await;
match self.worker.tick().await {
Ok(stats) => {
if stats.dead > 0 {
warn!(
sink = %self.name,
dead = stats.dead,
published = stats.published,
"cdc drain dead-lettered rows"
);
}
if stats.late_recovered > 0 {
warn!(
sink = %self.name,
late_recovered = stats.late_recovered,
"cdc drain recovered rows that committed after the lag window"
);
}
},
Err(error) => {
error!(
sink = %self.name,
%error,
"cdc drain tick failed — retrying on the next tick"
);
},
}
}
}
}
pub async fn build_drains(
config: Option<&CdcOutboundConfig>,
db_pool: Option<&PgPool>,
) -> Result<Option<Vec<SinkDrain>>, String> {
let Some(cfg) = config else {
return Ok(None);
};
cfg.validate()?;
let pool = db_pool.cloned().ok_or_else(|| {
"[cdc_outbound] requires a database pool — the change-log outbox and the per-sink \
delivery state are database-resident. The binary provides one when database_url is \
set; library embedders must pass a PgPool to Server::new."
.to_string()
})?;
crate::migration_lock::run_migration(
&pool,
fraiseql_cdc_sinks::outbox_sink_state_migration_sql(),
)
.await
.map_err(|e| {
format!(
"[cdc_outbound] could not create the delivery-state table \
(core.tb_cdc_sink_state): {e}. Refusing to boot rather than draining without \
durable delivery state, which would re-publish every event on restart."
)
})?;
let tick = Duration::from_secs(cfg.tick_interval_secs);
let mut drains = Vec::with_capacity(cfg.sinks.len());
for section in &cfg.sinks {
drains.push(build_one(section, &pool, tick, cfg.batch_size).await?);
}
info!(sinks = drains.len(), "cdc outbound configured");
Ok(Some(drains))
}
async fn build_one(
section: &CdcSinkSectionConfig,
pool: &PgPool,
tick: Duration,
batch_size: i64,
) -> Result<SinkDrain, String> {
let mut sink_config = CdcSinkConfig::new(§ion.name, §ion.subject_template);
sink_config.tables.clone_from(§ion.tables);
sink_config.tenants.clone_from(§ion.tenants);
if let Some(max_attempts) = section.max_attempts {
sink_config.max_attempts = max_attempts;
}
let connect_failed = |e: fraiseql_cdc_sinks::CdcError| {
format!(
"[cdc_outbound] sink {:?} could not connect to {}: {e}. Refusing to boot \
rather than starting a drain that publishes nowhere.",
section.name, section.endpoint
)
};
let sink = match section.kind.to_ascii_lowercase().as_str() {
#[cfg(feature = "cdc-kafka")]
"kafka" => ConfiguredSink::Kafka(
fraiseql_cdc_sinks::KafkaSink::connect(§ion.endpoint, sink_config.clone())
.map_err(connect_failed)?,
),
#[cfg(feature = "cdc-kinesis")]
"kinesis" => ConfiguredSink::Kinesis(
fraiseql_cdc_sinks::KinesisSink::connect(§ion.endpoint, sink_config.clone())
.await
.map_err(connect_failed)?,
),
_ => {
let sink = NatsJetStreamSink::connect(§ion.endpoint, sink_config.clone())
.await
.map_err(connect_failed)?;
if let Some(ref stream) = section.ensure_stream {
sink.ensure_stream(stream, vec![subject_wildcard(§ion.subject_template)])
.await
.map_err(|e| {
format!(
"[cdc_outbound] sink {:?}: could not ensure JetStream stream \
{stream:?}: {e}",
section.name
)
})?;
}
ConfiguredSink::NatsJetStream(sink)
},
};
Ok(SinkDrain {
name: section.name.clone(),
worker: DrainWorker::new(pool.clone(), sink, sink_config).with_batch_size(batch_size),
tick,
})
}
fn subject_wildcard(template: &str) -> String {
let literal_prefix = template.split('{').next().unwrap_or("");
let trimmed = literal_prefix.trim_end_matches('.');
if trimmed.is_empty() {
">".to_string()
} else {
format!("{trimmed}.>")
}
}
pub fn spawn_all(drains: Vec<SinkDrain>, tasks: &mut tokio::task::JoinSet<()>) {
for drain in drains {
tasks.spawn(async move { drain.run_forever().await });
}
}