#![cfg(feature = "kafka")]
use std::time::Duration;
use rdkafka::ClientConfig;
use rdkafka::error::KafkaError;
use rdkafka::producer::{FutureProducer, FutureRecord, Producer};
use super::{CdcConfig, CdcExactlyOnceMode};
#[derive(Debug, Clone)]
pub struct KafkaTxConfig {
pub transactional_id: String,
pub transaction_timeout: Duration,
pub brokers: String,
}
impl KafkaTxConfig {
pub fn from_cdc_config(brokers: &str, cdc: &CdcConfig) -> Self {
Self {
transactional_id: cdc.transactional_id(),
transaction_timeout: Duration::from_secs(cdc.kafka_tx_timeout_secs.max(1)),
brokers: brokers.to_string(),
}
}
}
pub fn build_transactional_producer(cfg: &KafkaTxConfig) -> Result<FutureProducer, String> {
let producer: FutureProducer = ClientConfig::new()
.set("bootstrap.servers", &cfg.brokers)
.set("enable.idempotence", "true")
.set("acks", "all")
.set("transactional.id", &cfg.transactional_id)
.set("compression.type", "lz4")
.set("message.timeout.ms", "30000")
.set("retries", "10")
.set("retry.backoff.ms", "100")
.create()
.map_err(|e| format!("kafka transactional producer build failed: {e}"))?;
producer
.init_transactions(cfg.transaction_timeout)
.map_err(|e| {
format!(
"kafka init_transactions failed for transactional.id '{}': {e}",
cfg.transactional_id
)
})?;
Ok(producer)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum KafkaTxPublishOutcome {
Committed { partition: i32, offset: i64 },
Aborted { reason: String },
}
pub async fn run_in_transaction(
producer: &FutureProducer,
timeout: Duration,
topic: &str,
partition_key: &str,
payload: &str,
) -> Result<KafkaTxPublishOutcome, KafkaError> {
producer.begin_transaction()?;
let record = FutureRecord::to(topic).key(partition_key).payload(payload);
let (partition, offset) = match producer.send(record, timeout).await {
Ok(delivery) => delivery,
Err((err, _msg)) => {
let _ = producer.abort_transaction(timeout);
return Ok(KafkaTxPublishOutcome::Aborted {
reason: format!("send queue failed: {err}"),
});
}
};
match producer.commit_transaction(timeout) {
Ok(()) => Ok(KafkaTxPublishOutcome::Committed { partition, offset }),
Err(err) => {
let _ = producer.abort_transaction(timeout);
Ok(KafkaTxPublishOutcome::Aborted {
reason: format!("commit_transaction failed: {err}"),
})
}
}
}
pub fn needs_transactional_producer(cdc: &CdcConfig) -> bool {
matches!(
cdc.exactly_once_mode,
CdcExactlyOnceMode::KafkaTransactional
)
}
#[cfg(test)]
mod tests {
use super::*;
fn cfg_with_mode(mode: CdcExactlyOnceMode) -> CdcConfig {
CdcConfig {
exactly_once_mode: mode,
transactional_id_prefix: "udb-cdc".to_string(),
slot_name: "udb_outbox_slot".to_string(),
..Default::default()
}
}
#[test]
fn needs_transactional_only_for_kafka_transactional() {
assert!(!needs_transactional_producer(&cfg_with_mode(
CdcExactlyOnceMode::AtLeastOnce
)));
assert!(!needs_transactional_producer(&cfg_with_mode(
CdcExactlyOnceMode::StateMachine
)));
assert!(needs_transactional_producer(&cfg_with_mode(
CdcExactlyOnceMode::KafkaTransactional
)));
}
#[test]
fn transactional_id_is_stable_across_restarts() {
let cfg = cfg_with_mode(CdcExactlyOnceMode::KafkaTransactional);
let tx_cfg_1 = KafkaTxConfig::from_cdc_config("broker:9092", &cfg);
let tx_cfg_2 = KafkaTxConfig::from_cdc_config("broker:9092", &cfg);
assert_eq!(
tx_cfg_1.transactional_id, tx_cfg_2.transactional_id,
"transactional.id must be reproducible across constructions"
);
assert_eq!(tx_cfg_1.transactional_id, "udb-cdc-udb_outbox_slot");
}
#[test]
fn transaction_timeout_defaults_and_overrides() {
let cfg = cfg_with_mode(CdcExactlyOnceMode::KafkaTransactional);
let tx_cfg = KafkaTxConfig::from_cdc_config("broker:9092", &cfg);
assert_eq!(tx_cfg.transaction_timeout, Duration::from_secs(30));
}
#[test]
fn publish_outcome_variants_are_typed() {
let committed = KafkaTxPublishOutcome::Committed {
partition: 0,
offset: 42,
};
let aborted = KafkaTxPublishOutcome::Aborted {
reason: "broker fenced us".to_string(),
};
assert_ne!(committed, aborted);
match aborted {
KafkaTxPublishOutcome::Aborted { reason } => {
assert!(reason.contains("fenced"));
}
_ => panic!("variant mismatch"),
}
}
}