Skip to main content

reifydb_cdc/consume/
checkpoint.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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}