mod common;
use std::collections::BTreeMap;
use std::time::Duration;
use crabka_replicator::config::{
ClusterConfig, Delivery, FlowConfig, NamingPolicy, ReplicatorConfig, Selectors,
};
use crabka_replicator::mm2::{Checkpoint, OffsetSync};
use crabka_replicator::offset_sync_store::OffsetSyncStore;
use crabka_replicator::selector::Selector;
use crabka_replicator::supervisor::FlowSupervisor;
use crabka_replicator::tasks::checkpoint::{CheckpointParams, run_once};
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[allow(clippy::too_many_lines)]
async fn offset_translation_never_skips_unreplicated_data() {
let source = common::start_broker().await;
let target = common::start_broker().await;
common::create_topic(&source.bootstrap, "orders", 1).await;
for i in 0..50_u32 {
let key = format!("k{i}");
let val = format!("v{i}");
common::produce(&source.bootstrap, "orders", key.as_bytes(), val.as_bytes()).await;
}
let mut clusters = BTreeMap::new();
clusters.insert(
"us-east".to_string(),
ClusterConfig {
bootstrap: source.bootstrap.clone(),
region: "us".into(),
zones: vec!["us".into()],
},
);
clusters.insert(
"eu-west".to_string(),
ClusterConfig {
bootstrap: target.bootstrap.clone(),
region: "eu".into(),
zones: vec!["eu".into()],
},
);
let config = ReplicatorConfig {
clusters,
flows: vec![FlowConfig {
from: "us-east".into(),
to: "eu-west".into(),
topics: Selectors {
include: vec!["orders".into()],
exclude: vec![],
},
groups: Selectors::default(),
naming: NamingPolicy::Default,
delivery: Delivery::AtLeastOnce,
}],
policies: vec![],
};
let sup = FlowSupervisor::run(config).await.expect("supervisor start");
common::await_count(
&target.bootstrap,
"us-east.orders",
50,
Duration::from_secs(30),
)
.await;
common::consume_and_commit(&source.bootstrap, "analytics", "orders").await;
let sync_records = crabka_replicator::admin_util::read_all(
&target.bootstrap,
&OffsetSync::topic_name("us-east"),
None,
)
.await
.expect("read offset-syncs topic");
println!(
"[offset_translation] offset-sync records on target: {}",
sync_records.len()
);
let mut store = OffsetSyncStore::default();
for (k, v) in &sync_records {
if let (Some(k), Some(v)) = (k.as_deref(), v.as_deref()) {
match OffsetSync::from_bytes(k, v) {
Ok(os) => {
println!(
"[offset_translation] ingest OffsetSync: topic={} part={} up={} down={}",
os.topic, os.partition, os.upstream, os.downstream
);
store.ingest(os);
}
Err(e) => eprintln!("[offset_translation] skip malformed offset-sync: {e}"),
}
}
}
run_once(
&CheckpointParams {
source_bootstrap: source.bootstrap.clone(),
target_bootstrap: target.bootstrap.clone(),
source_alias: "us-east".into(),
naming: crabka_replicator::config::NamingPolicy::Default,
group_selector: Selector::compile(&["analytics".into()], &[]).unwrap(),
security: None,
},
&store,
)
.await
.expect("run_once checkpoint");
let checkpoint_records = crabka_replicator::admin_util::read_all(
&target.bootstrap,
&Checkpoint::topic_name("us-east"),
None,
)
.await
.expect("read checkpoints topic");
println!(
"[offset_translation] checkpoint records on target: {}",
checkpoint_records.len()
);
let checkpoint = checkpoint_records
.into_iter()
.filter_map(|(k, v)| {
if let (Some(k), Some(v)) = (k, v) {
Checkpoint::from_bytes(&k, &v).ok()
} else {
None
}
})
.filter(|c| c.group == "analytics" && c.topic == "us-east.orders" && c.partition == 0)
.last()
.expect("no checkpoint found for (analytics, us-east.orders, 0)");
#[allow(clippy::cast_possible_wrap)]
let replicated_count = common::count(&target.bootstrap, "us-east.orders").await as i64;
println!(
"[offset_translation] upstream={} downstream={} replicated={}",
checkpoint.upstream, checkpoint.downstream, replicated_count
);
assert_eq!(
checkpoint.upstream, 50,
"upstream committed offset must be 50 (all records consumed); got {}",
checkpoint.upstream
);
assert!(
checkpoint.downstream >= 0,
"downstream offset must be >= 0; got {}",
checkpoint.downstream
);
assert!(
checkpoint.downstream <= replicated_count,
"NEVER-SKIP VIOLATED: downstream={} exceeds replicated_count={}",
checkpoint.downstream,
replicated_count
);
let translated = store.translate("orders", 0, 50);
assert_eq!(
translated,
Some(checkpoint.downstream),
"store.translate(\"orders\", 0, 50) = {:?} but checkpoint.downstream = {}",
translated,
checkpoint.downstream
);
println!(
"[offset_translation] PASS — upstream={} downstream={} replicated={} (never-skip held)",
checkpoint.upstream, checkpoint.downstream, replicated_count
);
sup.shutdown().await;
}