use bytes::Bytes;
use crabka_client_admin::AdminClient;
use crabka_client_producer::{Acks, Producer, ProducerRecord};
use tracing::warn;
use crate::config::NamingPolicy;
use crate::error::ReplicatorError;
use crate::mm2::{Checkpoint, OffsetSync};
use crate::naming::Renamer;
use crate::offset_sync_store::OffsetSyncStore;
use crate::selector::Selector;
pub struct CheckpointParams {
pub source_bootstrap: String,
pub target_bootstrap: String,
pub source_alias: String,
pub naming: NamingPolicy,
pub group_selector: Selector,
pub security: Option<crabka_client_core::security::ClientSecurity>,
}
pub async fn run_once(
params: &CheckpointParams,
store: &OffsetSyncStore,
) -> Result<(), ReplicatorError> {
let checkpoint_topic = Checkpoint::topic_name(¶ms.source_alias);
crate::admin_util::ensure_topic(
¶ms.target_bootstrap,
&checkpoint_topic,
1,
params.security.clone(),
)
.await
.map_err(|e| ReplicatorError::Client(format!("ensure checkpoints topic: {e}")))?;
let mut admin = AdminClient::connect_secured(
std::slice::from_ref(¶ms.source_bootstrap),
None, )
.await
.map_err(|e| ReplicatorError::Client(format!("source admin connect: {e}")))?;
let all_groups = admin
.list_groups()
.await
.map_err(|e| ReplicatorError::Client(format!("list_groups: {e}")))?;
let groups: Vec<String> = all_groups
.into_iter()
.filter(|g| !g.starts_with("crabka-replicator-"))
.filter(|g| params.group_selector.matches(g))
.collect();
let producer = build_producer(¶ms.target_bootstrap, params.security.clone())
.await
.map_err(|e| ReplicatorError::Client(format!("build target producer: {e}")))?;
let renamer = Renamer::new(params.naming, ¶ms.source_alias);
for group in &groups {
let offsets = match admin.list_consumer_group_offsets(group).await {
Ok(o) => o,
Err(e) => {
warn!(group = %group, error = %e, "list_consumer_group_offsets failed; skipping group");
continue;
}
};
for ((topic, partition), committed) in offsets {
let Some(downstream) = store.translate(&topic, partition, committed) else {
continue;
};
let checkpoint = Checkpoint {
group: group.clone(),
topic: renamer.target_name(&topic),
partition,
upstream: committed,
downstream,
metadata: String::new(),
};
let _rx = producer
.send(ProducerRecord {
topic: checkpoint_topic.clone(),
partition: None,
key: Some(Bytes::from(checkpoint.key_bytes())),
value: Some(Bytes::from(checkpoint.value_bytes())),
headers: Vec::new(),
timestamp_ms: None,
})
.await;
}
}
producer
.flush()
.await
.map_err(|e| ReplicatorError::Client(format!("producer flush: {e}")))?;
Ok(())
}
pub struct CheckpointTask {
handle: tokio::task::JoinHandle<()>,
shutdown: tokio::sync::watch::Sender<bool>,
}
impl CheckpointTask {
#[allow(clippy::unused_async)]
pub async fn start(
params: CheckpointParams,
interval: std::time::Duration,
) -> Result<Self, ReplicatorError> {
let (shutdown_tx, mut shutdown_rx) = tokio::sync::watch::channel(false);
let handle = tokio::spawn(async move {
loop {
let mut store = OffsetSyncStore::default();
let offset_syncs_topic = OffsetSync::topic_name(¶ms.source_alias);
match crate::admin_util::read_all(
¶ms.target_bootstrap,
&offset_syncs_topic,
params.security.clone(),
)
.await
{
Ok(records) => {
for (k, v) in records {
if let (Some(k), Some(v)) = (k, v) {
match OffsetSync::from_bytes(&k, &v) {
Ok(os) => store.ingest(os),
Err(e) => {
warn!(error = %e, "failed to decode offset-sync record; skipping");
}
}
}
}
}
Err(e) => {
warn!(error = %e, "failed to read offset-syncs topic; skipping cycle");
}
}
if let Err(e) = run_once(¶ms, &store).await {
warn!(error = %e, "checkpoint run_once failed");
}
tokio::select! {
() = tokio::time::sleep(interval) => {}
_ = shutdown_rx.changed() => {
if *shutdown_rx.borrow() {
break;
}
}
}
}
});
Ok(Self {
handle,
shutdown: shutdown_tx,
})
}
pub async fn shutdown(self) {
let _ = self.shutdown.send(true);
let _ = self.handle.await;
}
}
async fn build_producer(
bootstrap: &str,
security: Option<crabka_client_core::security::ClientSecurity>,
) -> Result<Producer, crabka_client_producer::ProducerError> {
let builder = Producer::builder()
.bootstrap(bootstrap)
.enable_idempotence(false)
.acks(Acks::All);
match security {
Some(s) => builder.security(s).build().await,
None => builder.build().await,
}
}
#[cfg(test)]
mod tests {
use assert2::assert;
use super::*;
use crate::mm2::{Checkpoint, OffsetSync};
use crate::offset_sync_store::OffsetSyncStore;
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn writes_translated_checkpoints() {
let s_dir = tempfile::TempDir::new().unwrap();
let t_dir = tempfile::TempDir::new().unwrap();
let source = crabka_broker::Broker::start(crabka_broker::BrokerConfig::for_tests(
s_dir.path().to_path_buf(),
))
.await
.unwrap();
let target = crabka_broker::Broker::start(crabka_broker::BrokerConfig::for_tests(
t_dir.path().to_path_buf(),
))
.await
.unwrap();
let sb = source.listen_addr().to_string();
let tb = target.listen_addr().to_string();
crate::test_util::create_topic(&sb, "orders", 1).await;
for _ in 0..200 {
crate::test_util::produce(&sb, "orders", b"k", b"v").await;
}
crate::test_util::commit_group(&sb, "g", "orders").await;
let mut syncs = OffsetSyncStore::default();
syncs.ingest(OffsetSync {
topic: "orders".into(),
partition: 0,
upstream: 0,
downstream: 0,
});
syncs.ingest(OffsetSync {
topic: "orders".into(),
partition: 0,
upstream: 200,
downstream: 165,
});
run_once(
&CheckpointParams {
source_bootstrap: sb,
target_bootstrap: tb.clone(),
source_alias: "us-east".into(),
naming: NamingPolicy::Default,
group_selector: Selector::compile(&["g".into()], &[]).unwrap(),
security: None,
},
&syncs,
)
.await
.unwrap();
let raw = crate::admin_util::read_last_value_for_key(
&tb,
&Checkpoint::topic_name("us-east"),
b"",
None,
)
.await
.unwrap()
.unwrap();
assert!(!raw.is_empty());
}
}