Skip to main content

crabka_replicator/
checkpoint_store.rs

1//! Position recovery: persist the worker's [`SourceOffset`] to a compacted
2//! internal topic on the TARGET cluster, keyed by flow name.
3//!
4//! On restart, [`InternalTopicCheckpointStore::load`] reads the last value for
5//! the flow's key out of the compacted topic, recovering the exact position the
6//! worker had reached before it stopped.
7
8use async_trait::async_trait;
9use bytes::Bytes;
10use crabka_connect::{CheckpointStore, ConnectError, SourceOffset};
11
12/// The internal compacted topic used to store replicator checkpoints.
13const STATE_TOPIC: &str = "crabka-replicator-offsets";
14
15/// A [`CheckpointStore`] backed by a compacted internal Kafka topic on the
16/// target cluster.
17///
18/// Each flow gets its own key (`flow_name`) within the shared compacted topic
19/// [`STATE_TOPIC`]. On restart, [`load`](Self::load) fetches the last value for
20/// that key, recovering the exact partition offsets the worker had reached.
21pub struct InternalTopicCheckpointStore {
22    producer: crabka_client_producer::Producer,
23    target_bootstrap: String,
24    topic: String,
25    key: String,
26    security: Option<crabka_client_core::security::ClientSecurity>,
27}
28
29impl InternalTopicCheckpointStore {
30    /// Ensure the compacted offset topic exists on the target cluster, build a
31    /// producer, and return a store keyed by `flow_name`.
32    ///
33    /// # Errors
34    ///
35    /// Returns [`ConnectError::Offset`] if the topic cannot be created or the
36    /// producer cannot connect.
37    pub async fn start(
38        target_bootstrap: &str,
39        flow_name: &str,
40        security: Option<crabka_client_core::security::ClientSecurity>,
41    ) -> Result<Self, ConnectError> {
42        crate::admin_util::ensure_compacted_topic(target_bootstrap, STATE_TOPIC, security.clone())
43            .await
44            .map_err(ConnectError::Offset)?;
45
46        let builder = crabka_client_producer::Producer::builder()
47            .bootstrap(target_bootstrap)
48            .enable_idempotence(false)
49            .acks(crabka_client_producer::Acks::All);
50
51        let producer = match security.clone() {
52            Some(s) => builder.security(s).build().await,
53            None => builder.build().await,
54        }
55        .map_err(|e| ConnectError::Offset(e.to_string()))?;
56
57        Ok(Self {
58            producer,
59            target_bootstrap: target_bootstrap.to_string(),
60            topic: STATE_TOPIC.into(),
61            key: flow_name.into(),
62            security,
63        })
64    }
65}
66
67#[async_trait]
68impl CheckpointStore for InternalTopicCheckpointStore {
69    async fn save(&self, offset: &SourceOffset) -> Result<(), ConnectError> {
70        let bytes = serde_json::to_vec(offset).map_err(|e| ConnectError::Offset(e.to_string()))?;
71
72        self.producer
73            .send(crabka_client_producer::ProducerRecord {
74                topic: self.topic.clone(),
75                partition: None,
76                key: Some(Bytes::copy_from_slice(self.key.as_bytes())),
77                value: Some(Bytes::from(bytes)),
78                headers: vec![],
79                timestamp_ms: None,
80            })
81            .await
82            .await
83            .map_err(|e| ConnectError::Offset(e.to_string()))?
84            .map_err(|e| ConnectError::Offset(e.to_string()))?;
85
86        self.producer
87            .flush()
88            .await
89            .map_err(|e| ConnectError::Offset(e.to_string()))?;
90
91        Ok(())
92    }
93
94    async fn load(&self) -> Result<Option<SourceOffset>, ConnectError> {
95        let latest = crate::admin_util::read_last_value_for_key(
96            &self.target_bootstrap,
97            &self.topic,
98            self.key.as_bytes(),
99            self.security.clone(),
100        )
101        .await
102        .map_err(ConnectError::Offset)?;
103
104        match latest {
105            Some(bytes) => Ok(Some(
106                serde_json::from_slice(&bytes).map_err(|e| ConnectError::Offset(e.to_string()))?,
107            )),
108            None => Ok(None),
109        }
110    }
111}
112
113#[cfg(test)]
114mod tests {
115    use std::collections::BTreeMap;
116
117    use assert2::assert;
118    use crabka_connect::{CheckpointStore, OffsetValue, SourceOffset};
119
120    use super::*;
121
122    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
123    async fn persists_and_reloads_position_from_target() {
124        let dir = tempfile::TempDir::new().unwrap();
125        let broker = crabka_broker::Broker::start(crabka_broker::BrokerConfig::for_tests(
126            dir.path().to_path_buf(),
127        ))
128        .await
129        .unwrap();
130        let target = broker.listen_addr().to_string();
131
132        let store = InternalTopicCheckpointStore::start(&target, "flow1", None)
133            .await
134            .unwrap();
135
136        let mut pos = BTreeMap::new();
137        pos.insert("orders-0".to_string(), OffsetValue::Long(42));
138        let off = SourceOffset::new(BTreeMap::new(), pos);
139
140        store.save(&off).await.unwrap();
141
142        let store2 = InternalTopicCheckpointStore::start(&target, "flow1", None)
143            .await
144            .unwrap();
145        let loaded = store2.load().await.unwrap().unwrap();
146        assert!(loaded == off);
147    }
148}