reifydb_cdc/consume/
checkpoint.rs1use reifydb_codec::row::bytes::EncodedBytes;
5use reifydb_core::{
6 common::CommitVersion,
7 interface::cdc::{CheckpointState, ConsumerClass},
8 key::cdc::ToConsumerKey,
9};
10use reifydb_transaction::transaction::{Transaction, command::CommandTransaction};
11use reifydb_value::{Result, util::cowvec::CowVec};
12
13#[derive(Debug, Clone, Copy, PartialEq, Eq)]
14pub struct CheckpointRow {
15 pub version: CommitVersion,
16 pub class: ConsumerClass,
17 pub state: CheckpointState,
18}
19
20impl CheckpointRow {
21 pub fn decode(row: &[u8]) -> Option<Self> {
22 if row.len() < 10 {
23 return None;
24 }
25 let mut buffer = [0u8; 8];
26 buffer.copy_from_slice(&row[0..8]);
27 Some(Self {
28 version: CommitVersion(u64::from_be_bytes(buffer)),
29 class: ConsumerClass::decode(row[8])?,
30 state: CheckpointState::decode(row[9])?,
31 })
32 }
33
34 fn encode(&self) -> EncodedBytes {
35 let mut bytes = Vec::with_capacity(10);
36 bytes.extend_from_slice(&self.version.0.to_be_bytes());
37 bytes.push(self.class.encode());
38 bytes.push(self.state.encode());
39 EncodedBytes(CowVec::new(bytes))
40 }
41}
42
43pub struct CdcCheckpoint {}
44
45impl CdcCheckpoint {
46 pub fn fetch<K: ToConsumerKey>(txn: &mut Transaction<'_>, consumer: &K) -> Result<CommitVersion> {
47 Ok(Self::fetch_opt(txn, consumer)?.unwrap_or(CommitVersion(1)))
48 }
49
50 pub fn fetch_opt<K: ToConsumerKey>(txn: &mut Transaction<'_>, consumer: &K) -> Result<Option<CommitVersion>> {
51 Ok(Self::fetch_row(txn, consumer)?.map(|row| row.version))
52 }
53
54 pub fn fetch_row<K: ToConsumerKey>(txn: &mut Transaction<'_>, consumer: &K) -> Result<Option<CheckpointRow>> {
55 let key = consumer.to_consumer_key();
56 Ok(txn.get(&key)?.and_then(|multi| CheckpointRow::decode(&multi.bytes)))
57 }
58
59 pub fn persist<K: ToConsumerKey>(
60 txn: &mut CommandTransaction,
61 consumer: &K,
62 version: CommitVersion,
63 class: ConsumerClass,
64 ) -> Result<()> {
65 let key = consumer.to_consumer_key();
66 let row = CheckpointRow {
67 version,
68 class,
69 state: CheckpointState::Valid,
70 };
71 txn.set(&key, row.encode())
72 }
73
74 pub fn invalidate<K: ToConsumerKey>(txn: &mut CommandTransaction, consumer: &K) -> Result<()> {
75 let key = consumer.to_consumer_key();
76 let Some(multi) = txn.get(&key)? else {
77 return Ok(());
78 };
79 let Some(mut bytes) = CheckpointRow::decode(&multi.bytes) else {
80 return Ok(());
81 };
82 assert_ne!(
83 bytes.class,
84 ConsumerClass::Pinning,
85 "a Pinning consumer checkpoint can never be invalidated: retention must never overtake it"
86 );
87 bytes.state = CheckpointState::Invalidated;
88 txn.set(&key, bytes.encode())
89 }
90
91 pub fn delete<K: ToConsumerKey>(txn: &mut CommandTransaction, consumer: &K) -> Result<()> {
92 let key = consumer.to_consumer_key();
93 txn.remove(&key)
94 }
95}