Skip to main content

reifydb_cdc/
rebuild.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::collections::{BTreeMap, BTreeSet};
5
6use reifydb_catalog::catalog::Catalog;
7use reifydb_codec::{
8	key::encoded::EncodedKey,
9	row::{
10		bytes::{EncodedBytes, read_fingerprint},
11		shape::{RowShape, fingerprint::RowShapeFingerprint},
12	},
13};
14use reifydb_core::{
15	error::diagnostic::internal::internal,
16	interface::{
17		catalog::object::ObjectId,
18		cdc::{Cdc, CdcChange},
19		change::{Change, ChangeOrigin, Diff, Diffs},
20	},
21	key::{
22		row::{PartitionedRowKey, PartitionedSortedViewRowKey, RowKey, SortedViewRowKey},
23		series::{PartitionedSeriesRowKey, SeriesRowKey},
24		tag::KeyTag,
25	},
26	value::column::columns::Columns,
27};
28use reifydb_transaction::transaction::Transaction;
29use reifydb_value::{Result, error::Error, value::row_number::RowNumber};
30
31pub struct RowTarget {
32	pub object: ObjectId,
33	pub row: RowNumber,
34}
35
36#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
37enum RebuiltKind {
38	Insert,
39	Update,
40	Remove,
41}
42
43#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
44struct BucketKey {
45	kind: RebuiltKind,
46	post_shape: RowShapeFingerprint,
47	pre_shape: RowShapeFingerprint,
48}
49
50#[derive(Default)]
51struct Bucket {
52	ids: Vec<RowNumber>,
53	pre: Vec<EncodedBytes>,
54	post: Vec<EncodedBytes>,
55}
56
57pub fn row_target(key: &EncodedKey) -> Option<RowTarget> {
58	match KeyTag::of(key)? {
59		KeyTag::Row => {
60			let row_key = RowKey::decode(key)?;
61			Some(RowTarget {
62				object: ObjectId::from(row_key.storage),
63				row: row_key.row,
64			})
65		}
66		KeyTag::SeriesRow => {
67			let series_key = SeriesRowKey::decode(key)?;
68			Some(RowTarget {
69				object: ObjectId::from(series_key.storage),
70				row: RowNumber(series_key.sequence),
71			})
72		}
73		KeyTag::PartitionedRow => {
74			let partitioned = PartitionedRowKey::decode(key)?;
75			Some(RowTarget {
76				object: ObjectId::from(partitioned.storage),
77				row: partitioned.row,
78			})
79		}
80		KeyTag::PartitionedSeriesRow => {
81			let partitioned = PartitionedSeriesRowKey::decode(key)?;
82			Some(RowTarget {
83				object: ObjectId::from(partitioned.storage),
84				row: RowNumber(partitioned.sequence),
85			})
86		}
87		KeyTag::SortedViewRow => Some(RowTarget {
88			object: ObjectId::from(SortedViewRowKey::storage_of(key)?),
89			row: SortedViewRowKey::row_of(key)?,
90		}),
91		KeyTag::PartitionedSortedViewRow => Some(RowTarget {
92			object: ObjectId::from(PartitionedSortedViewRowKey::storage_of(key)?),
93			row: PartitionedSortedViewRowKey::row_of(key)?,
94		}),
95		_ => None,
96	}
97}
98
99fn tracked_target(key: &EncodedKey) -> Option<RowTarget> {
100	let target = row_target(key)?;
101	if matches!(target.object, ObjectId::Queue(_)) {
102		return None;
103	}
104	Some(target)
105}
106
107pub fn changed_objects(cdc: &Cdc) -> BTreeSet<ObjectId> {
108	cdc.changes.iter().filter_map(|change| tracked_target(change.key())).map(|t| t.object).collect()
109}
110
111pub fn rebuild_changes(cdc: &Cdc, catalog: &Catalog, txn: &mut Transaction<'_>) -> Result<Vec<Change>> {
112	rebuild_selected_changes(cdc, catalog, txn, |_| true)
113}
114
115pub fn rebuild_selected_changes(
116	cdc: &Cdc,
117	catalog: &Catalog,
118	txn: &mut Transaction<'_>,
119	accept: impl Fn(ObjectId) -> bool,
120) -> Result<Vec<Change>> {
121	let mut grouped: BTreeMap<ObjectId, BTreeMap<BucketKey, Bucket>> = BTreeMap::new();
122
123	for cdc_change in &cdc.changes {
124		let Some(target) = tracked_target(cdc_change.key()) else {
125			continue;
126		};
127		if !accept(target.object) {
128			continue;
129		}
130		let (key, pre, post) = match cdc_change {
131			CdcChange::Insert {
132				post,
133				..
134			} => {
135				let fingerprint = read_fingerprint(post);
136				(
137					BucketKey {
138						kind: RebuiltKind::Insert,
139						post_shape: fingerprint,
140						pre_shape: fingerprint,
141					},
142					None,
143					Some(post.clone()),
144				)
145			}
146			CdcChange::Update {
147				pre,
148				post,
149				..
150			} => (
151				BucketKey {
152					kind: RebuiltKind::Update,
153					post_shape: read_fingerprint(post),
154					pre_shape: read_fingerprint(pre),
155				},
156				Some(pre.clone()),
157				Some(post.clone()),
158			),
159			CdcChange::Delete {
160				visible: false,
161				..
162			} => continue,
163			CdcChange::Delete {
164				key,
165				pre,
166				visible: true,
167			} => {
168				let pre = pre.as_ref().ok_or_else(|| {
169					Error(Box::new(internal(format!(
170						"CDC delete for key {:?} at version {} carries no pre-image, so its \
171						 change cannot be rebuilt",
172						key.as_slice(),
173						cdc.version.0
174					))))
175				})?;
176				let fingerprint = read_fingerprint(pre);
177				(
178					BucketKey {
179						kind: RebuiltKind::Remove,
180						post_shape: fingerprint,
181						pre_shape: fingerprint,
182					},
183					Some(pre.clone()),
184					None,
185				)
186			}
187		};
188
189		let bucket = grouped.entry(target.object).or_default().entry(key).or_default();
190		bucket.ids.push(target.row);
191		if let Some(pre) = pre {
192			bucket.pre.push(pre);
193		}
194		if let Some(post) = post {
195			bucket.post.push(post);
196		}
197	}
198
199	let mut shapes: BTreeMap<RowShapeFingerprint, RowShape> = BTreeMap::new();
200	let mut changes: Vec<Change> = Vec::with_capacity(grouped.len());
201
202	for (object, buckets) in grouped {
203		let mut diffs: Diffs = Diffs::new();
204		for (key, bucket) in buckets {
205			let diff = match key.kind {
206				RebuiltKind::Insert => {
207					let shape = load_shape(catalog, txn, &mut shapes, key.post_shape)?;
208					Diff::insert(Columns::from_encoded_bytes(&shape, &bucket.ids, &bucket.post))
209				}
210				RebuiltKind::Update => {
211					let pre_shape = load_shape(catalog, txn, &mut shapes, key.pre_shape)?;
212					let post_shape = load_shape(catalog, txn, &mut shapes, key.post_shape)?;
213					Diff::update(
214						Columns::from_encoded_bytes(&pre_shape, &bucket.ids, &bucket.pre),
215						Columns::from_encoded_bytes(&post_shape, &bucket.ids, &bucket.post),
216					)
217				}
218				RebuiltKind::Remove => {
219					let shape = load_shape(catalog, txn, &mut shapes, key.pre_shape)?;
220					Diff::remove(Columns::from_encoded_bytes(&shape, &bucket.ids, &bucket.pre))
221				}
222			};
223			diffs.push(diff);
224		}
225		changes.push(Change {
226			origin: ChangeOrigin::Object(object),
227			diffs,
228			version: cdc.version,
229			changed_at: cdc.timestamp,
230		});
231	}
232
233	Ok(changes)
234}
235
236fn load_shape(
237	catalog: &Catalog,
238	txn: &mut Transaction<'_>,
239	cache: &mut BTreeMap<RowShapeFingerprint, RowShape>,
240	fingerprint: RowShapeFingerprint,
241) -> Result<RowShape> {
242	if let Some(shape) = cache.get(&fingerprint) {
243		return Ok(shape.clone());
244	}
245	let shape = catalog.get_or_load_row_shape(fingerprint, txn)?.ok_or_else(|| {
246		Error(Box::new(internal(format!(
247			"RowShape with fingerprint {:?} not found while rebuilding CDC changes",
248			fingerprint
249		))))
250	})?;
251	cache.insert(fingerprint, shape.clone());
252	Ok(shape)
253}
254
255#[cfg(test)]
256mod tests {
257	use reifydb_core::{
258		interface::catalog::{
259			id::{SeriesId, TableId, ViewId},
260			storage::StorageId,
261		},
262		key::{
263			row::{PartitionedRowKey, RowKey},
264			series::{PartitionedSeriesRowKey, SeriesRowKey},
265		},
266	};
267	use reifydb_value::value::partition::Partition;
268
269	use super::*;
270
271	#[test]
272	fn test_row_key_maps_to_its_storage_object() {
273		let target = row_target(&RowKey::encoded(StorageId::table(3), RowNumber(9))).expect("row target");
274		assert_eq!(target.object, ObjectId::Table(TableId(3)));
275		assert_eq!(target.row, RowNumber(9));
276	}
277
278	#[test]
279	fn test_view_row_key_maps_to_the_view_and_never_to_its_former_backing_table() {
280		let target = row_target(&RowKey::encoded(StorageId::view(42), RowNumber(1))).expect("row target");
281		assert_eq!(target.object, ObjectId::View(ViewId(42)));
282	}
283
284	#[test]
285	fn test_partitioned_row_key_maps_to_its_storage_object() {
286		let key = PartitionedRowKey::encoded(StorageId::view(8), Partition(5), RowNumber(2));
287		let target = row_target(&key).expect("row target");
288		assert_eq!(target.object, ObjectId::View(ViewId(8)));
289		assert_eq!(target.row, RowNumber(2));
290	}
291
292	#[test]
293	fn test_series_row_key_uses_the_sequence_as_row_number_never_the_series_key() {
294		// A RowKey decode of the longer series suffix would invent RowNumber(1_000) out of the key bytes.
295		let key = SeriesRowKey {
296			storage: StorageId::series(4),
297			variant_tag: None,
298			key: 1_000,
299			sequence: 7,
300		}
301		.encode();
302		assert!(RowKey::decode(&key).is_none(), "a series row key must no longer decode as a plain row key");
303		let target = row_target(&key).expect("row target");
304		assert_eq!(target.object, ObjectId::Series(SeriesId(4)));
305		assert_eq!(target.row, RowNumber(7));
306	}
307
308	#[test]
309	fn test_tagged_series_row_key_maps_to_its_series() {
310		// The variant tag shifts the key and sequence by one byte, so the tagged layout needs its own cover.
311		let key = SeriesRowKey {
312			storage: StorageId::series(9),
313			variant_tag: Some(3),
314			key: 1_000,
315			sequence: 11,
316		}
317		.encode();
318		let target = row_target(&key).expect("row target");
319		assert_eq!(target.object, ObjectId::Series(SeriesId(9)));
320		assert_eq!(target.row, RowNumber(11));
321	}
322
323	#[test]
324	fn test_series_row_key_on_a_view_storage_maps_to_the_view_never_to_a_series() {
325		// A series-backed view writes series keys under a View storage id; reading the object id back as a
326		// series would attribute every rebuilt change of that view to a series that does not exist.
327		let key = SeriesRowKey {
328			storage: StorageId::view(8),
329			variant_tag: None,
330			key: 1_000,
331			sequence: 7,
332		}
333		.encode();
334		let target = row_target(&key).expect("row target");
335		assert_eq!(target.object, ObjectId::View(ViewId(8)));
336		assert_eq!(target.row, RowNumber(7));
337	}
338
339	#[test]
340	fn test_partitioned_series_row_key_uses_the_sequence_as_row_number() {
341		// The partitioned series kind carries its row identity in the sequence, exactly like the unpartitioned
342		// one; taking the series key instead would invent RowNumber(1_000) and misjoin every pre-image.
343		let key = PartitionedSeriesRowKey::encoded(StorageId::series(4), Partition(1), None, 1_000, 7);
344		let target = row_target(&key).expect("row target");
345		assert_eq!(target.object, ObjectId::Series(SeriesId(4)));
346		assert_eq!(target.row, RowNumber(7));
347	}
348
349	#[test]
350	fn test_partitioned_series_row_key_on_a_view_storage_maps_to_the_view() {
351		// Same widening as the unpartitioned case: the storage id decides the object, not the key kind.
352		let key = PartitionedSeriesRowKey::encoded(StorageId::view(8), Partition(1), Some(2), 5, 3);
353		let target = row_target(&key).expect("row target");
354		assert_eq!(target.object, ObjectId::View(ViewId(8)));
355		assert_eq!(target.row, RowNumber(3));
356	}
357
358	#[test]
359	fn test_catalog_keys_are_skipped() {
360		assert!(row_target(&EncodedKey::new(b"not a row key".to_vec())).is_none());
361	}
362}