1use 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 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 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 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 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 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}