use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct CdcOutboundConfig {
pub sinks: Vec<CdcSinkSectionConfig>,
pub tick_interval_secs: u64,
pub batch_size: i64,
}
impl Default for CdcOutboundConfig {
fn default() -> Self {
Self {
sinks: Vec::new(),
tick_interval_secs: 5,
batch_size: 256,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct CdcSinkSectionConfig {
pub name: String,
pub kind: String,
pub endpoint: String,
pub subject_template: String,
#[serde(default)]
pub tables: Option<Vec<String>>,
#[serde(default)]
pub tenants: Option<Vec<uuid::Uuid>>,
#[serde(default)]
pub max_attempts: Option<i32>,
#[serde(default)]
pub ensure_stream: Option<String>,
}
impl CdcOutboundConfig {
pub fn validate(&self) -> Result<(), String> {
if self.sinks.is_empty() {
return Err("[cdc_outbound] declares no sinks; remove the section to disable \
outbound CDC rather than configuring a drain with nowhere to go"
.to_string());
}
if self.tick_interval_secs == 0 {
return Err("[cdc_outbound] tick_interval_secs must be at least 1".to_string());
}
if self.batch_size < 1 {
return Err("[cdc_outbound] batch_size must be at least 1".to_string());
}
let mut seen = std::collections::HashSet::new();
for sink in &self.sinks {
if sink.name.trim().is_empty() {
return Err("[cdc_outbound] a sink has an empty name; the name is the \
delivery-state partition key"
.to_string());
}
if !seen.insert(sink.name.as_str()) {
return Err(format!(
"[cdc_outbound] duplicate sink name {:?}; two sinks sharing a name would \
share one delivery-state partition and each mark the other's rows published",
sink.name
));
}
if sink.endpoint.trim().is_empty() {
return Err(format!("[cdc_outbound] sink {:?} has an empty endpoint", sink.name));
}
if sink.subject_template.trim().is_empty() {
return Err(format!(
"[cdc_outbound] sink {:?} has an empty subject_template",
sink.name
));
}
if sink.max_attempts.is_some_and(|attempts| attempts < 1) {
return Err(format!(
"[cdc_outbound] sink {:?}: max_attempts must be at least 1",
sink.name
));
}
validate_kind(&sink.name, &sink.kind)?;
}
Ok(())
}
}
fn validate_kind(sink: &str, kind: &str) -> Result<(), String> {
match kind.to_ascii_lowercase().as_str() {
"nats-jetstream" => Ok(()),
#[cfg(feature = "cdc-kafka")]
"kafka" => Ok(()),
#[cfg(not(feature = "cdc-kafka"))]
"kafka" => Err(format!(
"[cdc_outbound] sink {sink:?}: kind \"kafka\" is implemented but not compiled into \
this binary. Rebuild with the `cdc-kafka` feature. Refusing to boot rather than \
silently draining nothing."
)),
#[cfg(feature = "cdc-kinesis")]
"kinesis" => Ok(()),
#[cfg(not(feature = "cdc-kinesis"))]
"kinesis" => Err(format!(
"[cdc_outbound] sink {sink:?}: kind \"kinesis\" is implemented but not compiled \
into this binary. Rebuild with the `cdc-kinesis` feature. Refusing to boot rather \
than silently draining nothing."
)),
"pulsar" => Err(format!(
"[cdc_outbound] sink {sink:?}: kind \"pulsar\" is not implemented yet (#382 tracks \
it). Refusing to boot rather than silently draining nothing."
)),
other => Err(format!(
"[cdc_outbound] sink {sink:?}: unknown kind {other:?}; expected \"nats-jetstream\"{}",
match (cfg!(feature = "cdc-kafka"), cfg!(feature = "cdc-kinesis")) {
(true, true) => ", \"kafka\" or \"kinesis\"",
(true, false) => " or \"kafka\"",
(false, true) => " or \"kinesis\"",
(false, false) => "",
}
)),
}
}