crabka_replicator/
checkpoint_store.rs1use async_trait::async_trait;
9use bytes::Bytes;
10use crabka_connect::{CheckpointStore, ConnectError, SourceOffset};
11
12const STATE_TOPIC: &str = "crabka-replicator-offsets";
14
15pub 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 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}