1use bytes::Bytes;
8use crabka_client_admin::AdminClient;
9use crabka_client_producer::{Acks, Producer, ProducerRecord};
10use tracing::warn;
11
12use crate::config::NamingPolicy;
13use crate::error::ReplicatorError;
14use crate::mm2::{Checkpoint, OffsetSync};
15use crate::naming::Renamer;
16use crate::offset_sync_store::OffsetSyncStore;
17use crate::selector::Selector;
18
19pub struct CheckpointParams {
21 pub source_bootstrap: String,
23 pub target_bootstrap: String,
25 pub source_alias: String,
27 pub naming: NamingPolicy,
30 pub group_selector: Selector,
32 pub security: Option<crabka_client_core::security::ClientSecurity>,
34}
35
36pub async fn run_once(
42 params: &CheckpointParams,
43 store: &OffsetSyncStore,
44) -> Result<(), ReplicatorError> {
45 let checkpoint_topic = Checkpoint::topic_name(¶ms.source_alias);
46
47 crate::admin_util::ensure_topic(
49 ¶ms.target_bootstrap,
50 &checkpoint_topic,
51 1,
52 params.security.clone(),
53 )
54 .await
55 .map_err(|e| ReplicatorError::Client(format!("ensure checkpoints topic: {e}")))?;
56
57 let mut admin = AdminClient::connect_secured(
59 std::slice::from_ref(¶ms.source_bootstrap),
60 None, )
62 .await
63 .map_err(|e| ReplicatorError::Client(format!("source admin connect: {e}")))?;
64
65 let all_groups = admin
67 .list_groups()
68 .await
69 .map_err(|e| ReplicatorError::Client(format!("list_groups: {e}")))?;
70
71 let groups: Vec<String> = all_groups
72 .into_iter()
73 .filter(|g| !g.starts_with("crabka-replicator-"))
74 .filter(|g| params.group_selector.matches(g))
75 .collect();
76
77 let producer = build_producer(¶ms.target_bootstrap, params.security.clone())
79 .await
80 .map_err(|e| ReplicatorError::Client(format!("build target producer: {e}")))?;
81
82 let renamer = Renamer::new(params.naming, ¶ms.source_alias);
87 for group in &groups {
88 let offsets = match admin.list_consumer_group_offsets(group).await {
89 Ok(o) => o,
90 Err(e) => {
91 warn!(group = %group, error = %e, "list_consumer_group_offsets failed; skipping group");
92 continue;
93 }
94 };
95
96 for ((topic, partition), committed) in offsets {
97 let Some(downstream) = store.translate(&topic, partition, committed) else {
98 continue;
99 };
100
101 let checkpoint = Checkpoint {
102 group: group.clone(),
103 topic: renamer.target_name(&topic),
104 partition,
105 upstream: committed,
106 downstream,
107 metadata: String::new(),
108 };
109
110 let _rx = producer
111 .send(ProducerRecord {
112 topic: checkpoint_topic.clone(),
113 partition: None,
114 key: Some(Bytes::from(checkpoint.key_bytes())),
115 value: Some(Bytes::from(checkpoint.value_bytes())),
116 headers: Vec::new(),
117 timestamp_ms: None,
118 })
119 .await;
120 }
121 }
122
123 producer
125 .flush()
126 .await
127 .map_err(|e| ReplicatorError::Client(format!("producer flush: {e}")))?;
128
129 Ok(())
130}
131
132pub struct CheckpointTask {
135 handle: tokio::task::JoinHandle<()>,
136 shutdown: tokio::sync::watch::Sender<bool>,
137}
138
139impl CheckpointTask {
140 #[allow(clippy::unused_async)]
153 pub async fn start(
154 params: CheckpointParams,
155 interval: std::time::Duration,
156 ) -> Result<Self, ReplicatorError> {
157 let (shutdown_tx, mut shutdown_rx) = tokio::sync::watch::channel(false);
158
159 let handle = tokio::spawn(async move {
160 loop {
161 let mut store = OffsetSyncStore::default();
163 let offset_syncs_topic = OffsetSync::topic_name(¶ms.source_alias);
164
165 match crate::admin_util::read_all(
166 ¶ms.target_bootstrap,
167 &offset_syncs_topic,
168 params.security.clone(),
169 )
170 .await
171 {
172 Ok(records) => {
173 for (k, v) in records {
174 if let (Some(k), Some(v)) = (k, v) {
175 match OffsetSync::from_bytes(&k, &v) {
176 Ok(os) => store.ingest(os),
177 Err(e) => {
178 warn!(error = %e, "failed to decode offset-sync record; skipping");
179 }
180 }
181 }
182 }
183 }
184 Err(e) => {
185 warn!(error = %e, "failed to read offset-syncs topic; skipping cycle");
186 }
187 }
188
189 if let Err(e) = run_once(¶ms, &store).await {
190 warn!(error = %e, "checkpoint run_once failed");
191 }
192
193 tokio::select! {
195 () = tokio::time::sleep(interval) => {}
196 _ = shutdown_rx.changed() => {
197 if *shutdown_rx.borrow() {
198 break;
199 }
200 }
201 }
202 }
203 });
204
205 Ok(Self {
206 handle,
207 shutdown: shutdown_tx,
208 })
209 }
210
211 pub async fn shutdown(self) {
213 let _ = self.shutdown.send(true);
214 let _ = self.handle.await;
215 }
216}
217
218async fn build_producer(
220 bootstrap: &str,
221 security: Option<crabka_client_core::security::ClientSecurity>,
222) -> Result<Producer, crabka_client_producer::ProducerError> {
223 let builder = Producer::builder()
224 .bootstrap(bootstrap)
225 .enable_idempotence(false)
226 .acks(Acks::All);
227 match security {
228 Some(s) => builder.security(s).build().await,
229 None => builder.build().await,
230 }
231}
232
233#[cfg(test)]
234mod tests {
235 use assert2::assert;
236
237 use super::*;
238 use crate::mm2::{Checkpoint, OffsetSync};
239 use crate::offset_sync_store::OffsetSyncStore;
240
241 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
242 async fn writes_translated_checkpoints() {
243 let s_dir = tempfile::TempDir::new().unwrap();
244 let t_dir = tempfile::TempDir::new().unwrap();
245 let source = crabka_broker::Broker::start(crabka_broker::BrokerConfig::for_tests(
246 s_dir.path().to_path_buf(),
247 ))
248 .await
249 .unwrap();
250 let target = crabka_broker::Broker::start(crabka_broker::BrokerConfig::for_tests(
251 t_dir.path().to_path_buf(),
252 ))
253 .await
254 .unwrap();
255 let sb = source.listen_addr().to_string();
256 let tb = target.listen_addr().to_string();
257
258 crate::test_util::create_topic(&sb, "orders", 1).await;
259 for _ in 0..200 {
260 crate::test_util::produce(&sb, "orders", b"k", b"v").await;
261 }
262 crate::test_util::commit_group(&sb, "g", "orders").await;
263
264 let mut syncs = OffsetSyncStore::default();
265 syncs.ingest(OffsetSync {
266 topic: "orders".into(),
267 partition: 0,
268 upstream: 0,
269 downstream: 0,
270 });
271 syncs.ingest(OffsetSync {
272 topic: "orders".into(),
273 partition: 0,
274 upstream: 200,
275 downstream: 165,
276 });
277
278 run_once(
279 &CheckpointParams {
280 source_bootstrap: sb,
281 target_bootstrap: tb.clone(),
282 source_alias: "us-east".into(),
283 naming: NamingPolicy::Default,
284 group_selector: Selector::compile(&["g".into()], &[]).unwrap(),
285 security: None,
286 },
287 &syncs,
288 )
289 .await
290 .unwrap();
291
292 let raw = crate::admin_util::read_last_value_for_key(
293 &tb,
294 &Checkpoint::topic_name("us-east"),
295 b"",
296 None,
297 )
298 .await
299 .unwrap()
300 .unwrap();
301 assert!(!raw.is_empty());
302 }
303}