Skip to main content

icechunk_format/
transaction_log.rs

1//! Change records for commits, enabling conflict detection during rebase.
2
3use std::{
4    collections::{BTreeMap, BTreeSet},
5    iter,
6};
7
8use flatbuffers::{VerifierOptions, WIPOffset};
9use itertools::Either;
10
11use icechunk_types::Move;
12
13use crate::{
14    ChunkIndices, IcechunkResult, NodeId, Path, SnapshotId,
15    flatbuffers::generated::{self, MoveOperation, MoveOperationArgs},
16};
17use icechunk_types::ICResultExt as _;
18
19#[derive(Clone, Debug, PartialEq, Default)]
20pub struct TransactionLog {
21    buffer: Vec<u8>,
22}
23
24impl TransactionLog {
25    /// Low level method that creates a tx log from its parts
26    /// Intended to be used only by library creators
27    #[expect(clippy::too_many_arguments)]
28    pub fn new_from_parts(
29        id: &SnapshotId,
30        sorted_new_groups: impl ExactSizeIterator<Item = NodeId> + DoubleEndedIterator,
31        sorted_new_arrays: impl ExactSizeIterator<Item = NodeId> + DoubleEndedIterator,
32        sorted_deleted_groups: impl ExactSizeIterator<Item = NodeId> + DoubleEndedIterator,
33        sorted_deleted_arrays: impl ExactSizeIterator<Item = NodeId> + DoubleEndedIterator,
34        sorted_updated_groups: impl ExactSizeIterator<Item = NodeId> + DoubleEndedIterator,
35        sorted_updated_arrays: impl ExactSizeIterator<Item = NodeId> + DoubleEndedIterator,
36        sorted_updated_chunks: impl ExactSizeIterator<
37            Item = (NodeId, impl Iterator<Item = ChunkIndices>),
38        > + DoubleEndedIterator,
39        sorted_moves: impl Iterator<Item = Move>,
40    ) -> Self {
41        // TODO: what's a good capacity?
42        let mut builder = flatbuffers::FlatBufferBuilder::with_capacity(1_024 * 1_024);
43
44        //let new_groups = Some(builder.create_vector(sorted_new_groups.as_slice()));
45        let new_groups = Some(builder.create_vector_from_iter(
46            sorted_new_groups.map(|id| generated::ObjectId8::new(&id.0)),
47        ));
48        let new_arrays = Some(builder.create_vector_from_iter(
49            sorted_new_arrays.map(|id| generated::ObjectId8::new(&id.0)),
50        ));
51        let deleted_groups = Some(builder.create_vector_from_iter(
52            sorted_deleted_groups.map(|id| generated::ObjectId8::new(&id.0)),
53        ));
54        let deleted_arrays = Some(builder.create_vector_from_iter(
55            sorted_deleted_arrays.map(|id| generated::ObjectId8::new(&id.0)),
56        ));
57        let updated_groups = Some(builder.create_vector_from_iter(
58            sorted_updated_groups.map(|id| generated::ObjectId8::new(&id.0)),
59        ));
60        let updated_arrays = Some(builder.create_vector_from_iter(
61            sorted_updated_arrays.map(|id| generated::ObjectId8::new(&id.0)),
62        ));
63
64        let id = generated::ObjectId12::new(&id.0);
65        let id = Some(&id);
66        // TODO: very inefficient
67        let updated_chunks = sorted_updated_chunks
68            .map(|(node_id, chunks)| {
69                let node_id = generated::ObjectId8::new(&node_id.0);
70                let node_id = Some(&node_id);
71                let chunks = chunks
72                    .map(|indices| {
73                        let coords = Some(builder.create_vector(indices.0.as_slice()));
74                        generated::ChunkIndices::create(
75                            &mut builder,
76                            &generated::ChunkIndicesArgs { coords },
77                        )
78                    })
79                    .collect::<Vec<_>>();
80                let chunks = Some(builder.create_vector(chunks.as_slice()));
81                generated::ArrayUpdatedChunks::create(
82                    &mut builder,
83                    &generated::ArrayUpdatedChunksArgs { node_id, chunks },
84                )
85            })
86            .collect::<Vec<_>>();
87        let updated_chunks = builder.create_vector(updated_chunks.as_slice());
88        let updated_chunks = Some(updated_chunks);
89
90        let moved_nodes: Vec<_> = sorted_moves
91            .map(|Move { from, to }| {
92                let from = builder.create_string(from.to_string().as_str());
93                let to = builder.create_string(to.to_string().as_str());
94                let args = MoveOperationArgs { from: Some(from), to: Some(to) };
95                MoveOperation::create(&mut builder, &args)
96            })
97            .collect();
98
99        let moved_nodes = Some(builder.create_vector(moved_nodes.as_slice()));
100
101        let tx = generated::TransactionLog::create(
102            &mut builder,
103            &generated::TransactionLogArgs {
104                id,
105                new_groups,
106                new_arrays,
107                deleted_groups,
108                deleted_arrays,
109                updated_groups,
110                updated_arrays,
111                updated_chunks,
112                moved_nodes,
113                ..Default::default()
114            },
115        );
116
117        builder.finish(tx, Some("Ichk"));
118        let (mut buffer, offset) = builder.collapse();
119        buffer.drain(0..offset);
120        buffer.shrink_to_fit();
121        Self { buffer }
122    }
123
124    pub fn from_buffer(buffer: Vec<u8>) -> IcechunkResult<Self> {
125        let _ = flatbuffers::root_with_opts::<generated::TransactionLog<'_>>(
126            &ROOT_OPTIONS,
127            buffer.as_slice(),
128        )
129        .capture()?;
130        Ok(Self { buffer })
131    }
132
133    pub fn new_groups(
134        &self,
135    ) -> impl ExactSizeIterator<Item = NodeId> + DoubleEndedIterator + '_ {
136        self.root().new_groups().iter().map(From::from)
137    }
138
139    pub fn new_arrays(
140        &self,
141    ) -> impl ExactSizeIterator<Item = NodeId> + DoubleEndedIterator + '_ {
142        self.root().new_arrays().iter().map(From::from)
143    }
144
145    pub fn deleted_groups(
146        &self,
147    ) -> impl ExactSizeIterator<Item = NodeId> + DoubleEndedIterator + '_ {
148        self.root().deleted_groups().iter().map(From::from)
149    }
150
151    pub fn deleted_arrays(
152        &self,
153    ) -> impl ExactSizeIterator<Item = NodeId> + DoubleEndedIterator + '_ {
154        self.root().deleted_arrays().iter().map(From::from)
155    }
156
157    pub fn updated_groups(
158        &self,
159    ) -> impl ExactSizeIterator<Item = NodeId> + DoubleEndedIterator + '_ {
160        self.root().updated_groups().iter().map(From::from)
161    }
162
163    pub fn updated_arrays(
164        &self,
165    ) -> impl ExactSizeIterator<Item = NodeId> + DoubleEndedIterator + '_ {
166        self.root().updated_arrays().iter().map(From::from)
167    }
168
169    pub fn updated_chunks(
170        &self,
171    ) -> impl Iterator<Item = (NodeId, impl Iterator<Item = ChunkIndices> + '_)> + '_
172    {
173        self.root().updated_chunks().iter().map(|arr_chunks| {
174            let id: NodeId = arr_chunks.node_id().into();
175            let chunks = arr_chunks.chunks().iter().map(|idx| idx.into());
176            (id, chunks)
177        })
178    }
179
180    pub fn moves(&self) -> impl Iterator<Item = Move> + '_ {
181        let it = match self.root().moved_nodes() {
182            Some(it) => Either::Left(it.iter()),
183            None => Either::Right(iter::empty()),
184        };
185        it.filter_map(|m| {
186            let from = Path::new(m.from()).ok()?;
187            let to = Path::new(m.to()).ok()?;
188            Some(Move { from, to })
189        })
190    }
191
192    pub fn updated_chunks_for(
193        &self,
194        node: &NodeId,
195    ) -> impl Iterator<Item = ChunkIndices> + '_ + use<'_> {
196        let arr = self
197            .root()
198            .updated_chunks()
199            .lookup_by_key(node.0, |a, b| a.node_id().0.cmp(b));
200
201        match arr {
202            Some(arr) => Either::Left(arr.chunks().iter().map(From::from)),
203            None => Either::Right(iter::empty()),
204        }
205    }
206
207    pub fn updated_chunks_counts(
208        &self,
209    ) -> impl Iterator<Item = (NodeId, u64)> + '_ + use<'_> {
210        self.root().updated_chunks().iter().map(|arr_chunks| {
211            let id: NodeId = arr_chunks.node_id().into();
212            let n = arr_chunks.chunks().len();
213            (id, n as u64)
214        })
215    }
216
217    pub fn group_created(&self, id: &NodeId) -> bool {
218        self.root().new_groups().lookup_by_key(id.0, |a, b| a.0.cmp(b)).is_some()
219    }
220
221    pub fn array_created(&self, id: &NodeId) -> bool {
222        self.root().new_arrays().lookup_by_key(id.0, |a, b| a.0.cmp(b)).is_some()
223    }
224
225    pub fn group_deleted(&self, id: &NodeId) -> bool {
226        self.root().deleted_groups().lookup_by_key(id.0, |a, b| a.0.cmp(b)).is_some()
227    }
228
229    pub fn array_deleted(&self, id: &NodeId) -> bool {
230        self.root().deleted_arrays().lookup_by_key(id.0, |a, b| a.0.cmp(b)).is_some()
231    }
232
233    pub fn group_updated(&self, id: &NodeId) -> bool {
234        self.root().updated_groups().lookup_by_key(id.0, |a, b| a.0.cmp(b)).is_some()
235    }
236
237    pub fn array_updated(&self, id: &NodeId) -> bool {
238        self.root().updated_arrays().lookup_by_key(id.0, |a, b| a.0.cmp(b)).is_some()
239    }
240
241    pub fn chunks_updated(&self, id: &NodeId) -> bool {
242        self.root()
243            .updated_chunks()
244            .lookup_by_key(id.0, |a, b| a.node_id().0.cmp(b))
245            .is_some()
246    }
247
248    #[expect(unsafe_code)]
249    fn root(&self) -> generated::TransactionLog<'_> {
250        // SAFETY: self.buffer was serialized by our own flatbuffers serialization code.
251        // We skip validation for performance; a corrupt buffer here indicates
252        // file corruption or a bad Icechunk implementation, not a caller error.
253        unsafe {
254            flatbuffers::root_unchecked::<generated::TransactionLog<'_>>(&self.buffer)
255        }
256    }
257
258    pub fn bytes(&self) -> &[u8] {
259        self.buffer.as_slice()
260    }
261
262    pub fn len(&self) -> usize {
263        let root = self.root();
264        root.new_groups().len()
265            + root.new_arrays().len()
266            + root.deleted_groups().len()
267            + root.deleted_arrays().len()
268            + root.updated_groups().len()
269            + root.updated_arrays().len()
270            + root.updated_chunks().iter().map(|s| s.chunks().len()).sum::<usize>()
271            + root.moved_nodes().map(|v| v.len()).unwrap_or_default()
272    }
273
274    #[must_use]
275    pub fn is_empty(&self) -> bool {
276        self.len() == 0
277    }
278
279    pub fn has_moves(&self) -> bool {
280        self.root().moved_nodes().map(|v| !v.is_empty()).unwrap_or(false)
281    }
282
283    pub fn merge<'a, T: IntoIterator<Item = &'a TransactionLog>>(
284        id: &SnapshotId,
285        iter: T,
286    ) -> Self {
287        let txs = Vec::from_iter(iter);
288
289        // TODO: what's a good capacity?
290        let mut builder = flatbuffers::FlatBufferBuilder::with_capacity(1_024 * 1_024);
291
292        let new_groups = {
293            let new_groups =
294                BTreeSet::from_iter(txs.iter().flat_map(|tx| tx.new_groups()));
295            let new_groups = Vec::from_iter(
296                new_groups.into_iter().map(|id| generated::ObjectId8::new(&id.0)),
297            );
298            Some(builder.create_vector(new_groups.as_slice()))
299        };
300        let new_arrays = {
301            let new_arrays =
302                BTreeSet::from_iter(txs.iter().flat_map(|tx| tx.new_arrays()));
303            let new_arrays = Vec::from_iter(
304                new_arrays.into_iter().map(|id| generated::ObjectId8::new(&id.0)),
305            );
306            Some(builder.create_vector(new_arrays.as_slice()))
307        };
308        let deleted_groups = {
309            let deleted_groups =
310                BTreeSet::from_iter(txs.iter().flat_map(|tx| tx.deleted_groups()));
311            let deleted_groups = Vec::from_iter(
312                deleted_groups.into_iter().map(|id| generated::ObjectId8::new(&id.0)),
313            );
314            Some(builder.create_vector(deleted_groups.as_slice()))
315        };
316        let deleted_arrays = {
317            let deleted_arrays =
318                BTreeSet::from_iter(txs.iter().flat_map(|tx| tx.deleted_arrays()));
319            let deleted_arrays = Vec::from_iter(
320                deleted_arrays.into_iter().map(|id| generated::ObjectId8::new(&id.0)),
321            );
322            Some(builder.create_vector(deleted_arrays.as_slice()))
323        };
324        let updated_groups = {
325            let updated_groups =
326                BTreeSet::from_iter(txs.iter().flat_map(|tx| tx.updated_groups()));
327            let updated_groups = Vec::from_iter(
328                updated_groups.into_iter().map(|id| generated::ObjectId8::new(&id.0)),
329            );
330            Some(builder.create_vector(updated_groups.as_slice()))
331        };
332        let updated_arrays = {
333            let updated_arrays =
334                BTreeSet::from_iter(txs.iter().flat_map(|tx| tx.updated_arrays()));
335            let updated_arrays = Vec::from_iter(
336                updated_arrays.into_iter().map(|id| generated::ObjectId8::new(&id.0)),
337            );
338            Some(builder.create_vector(updated_arrays.as_slice()))
339        };
340        let updated_chunks = {
341            let updated_chunks = txs.iter().fold(BTreeMap::new(), |res, tx| {
342                tx.updated_chunks().fold(res, |mut res, (node_id, chunks_it)| {
343                    let set: &mut BTreeSet<_> = res.entry(node_id).or_default();
344                    set.extend(chunks_it);
345                    res
346                })
347            });
348
349            let updated_chunks = updated_chunks
350                .into_iter()
351                .map(|(node_id, chunks)| {
352                    let node_id = generated::ObjectId8::new(&node_id.0);
353                    let node_id = Some(&node_id);
354                    let chunks = chunks
355                        .into_iter()
356                        .map(|indices| {
357                            let coords =
358                                Some(builder.create_vector(indices.0.as_slice()));
359                            generated::ChunkIndices::create(
360                                &mut builder,
361                                &generated::ChunkIndicesArgs { coords },
362                            )
363                        })
364                        .collect::<Vec<_>>();
365                    let chunks = Some(builder.create_vector(chunks.as_slice()));
366                    generated::ArrayUpdatedChunks::create(
367                        &mut builder,
368                        &generated::ArrayUpdatedChunksArgs { node_id, chunks },
369                    )
370                })
371                .collect::<Vec<_>>();
372
373            let updated_chunks = builder.create_vector(updated_chunks.as_slice());
374            Some(updated_chunks)
375        };
376
377        let id = generated::ObjectId12::new(&id.0);
378        let id = Some(&id);
379        // FIXME: verify no moves in origins
380        let moved_nodes: &[WIPOffset<_>] = &[];
381        let moved_nodes = Some(builder.create_vector(moved_nodes)); // FIXME:
382        let tx = generated::TransactionLog::create(
383            &mut builder,
384            &generated::TransactionLogArgs {
385                id,
386                new_groups,
387                new_arrays,
388                deleted_groups,
389                deleted_arrays,
390                updated_groups,
391                updated_arrays,
392                updated_chunks,
393                moved_nodes,
394                ..Default::default()
395            },
396        );
397
398        builder.finish(tx, Some("Ichk"));
399        let (mut buffer, offset) = builder.collapse();
400        buffer.drain(0..offset);
401        buffer.shrink_to_fit();
402        Self { buffer }
403    }
404}
405
406static ROOT_OPTIONS: VerifierOptions = VerifierOptions {
407    max_depth: 64,
408    max_tables: 50_000_000,
409    max_apparent_size: 1 << 31, // taken from the default
410    ignore_missing_null_terminator: true,
411};
412
413// Tests for TransactionLog depend on ChangeSet which lives in the icechunk crate.
414// They are kept in icechunk's test suite instead.
415#[cfg(any())]
416mod tests {
417    use std::collections::HashSet;
418
419    use bytes::Bytes;
420    use itertools::Itertools as _;
421
422    use crate::{
423        change_set::{ArrayData, ChangeSet, transaction_log_from_change_set},
424        format::{
425            ChunkIndices, NodeId, SnapshotId, manifest::ChunkPayload,
426            snapshot::ArrayShape, transaction_log::TransactionLog,
427        },
428    };
429
430    #[icechunk_macros::test]
431    fn test_merge() -> Result<(), Box<dyn std::error::Error>> {
432        let mut cs1 = ChangeSet::for_edits();
433        let added_group = NodeId::random();
434        let added_array = NodeId::random();
435        let deleted_group = NodeId::random();
436        let deleted_array = NodeId::random();
437        let updated_group = NodeId::random();
438        let chunk_added = NodeId::random();
439        cs1.add_group("/g1".try_into().unwrap(), added_group.clone(), Bytes::new())?;
440        cs1.delete_group("/g2".try_into().unwrap(), &deleted_group)?;
441        cs1.add_array(
442            "/a1".try_into().unwrap(),
443            added_array.clone(),
444            ArrayData {
445                shape: ArrayShape::new([(0, 10)]).unwrap(),
446                dimension_names: None,
447                user_data: Bytes::new(),
448            },
449        )?;
450        cs1.delete_array("/a2".try_into().unwrap(), &deleted_array)?;
451        cs1.update_group(&updated_group, &"/g3".try_into().unwrap(), Bytes::new())?;
452        cs1.set_chunk_ref(
453            chunk_added.clone(),
454            ChunkIndices(vec![0]),
455            Some(ChunkPayload::Inline(Bytes::new())),
456        )?;
457
458        let t1 = transaction_log_from_change_set(&SnapshotId::random(), &cs1);
459        let t2 = transaction_log_from_change_set(&SnapshotId::random(), &cs1);
460
461        let tx = TransactionLog::merge(&SnapshotId::random(), [&t1, &t2]);
462        assert!(tx.new_groups().eq([added_group.clone()]));
463        assert!(tx.new_arrays().eq([added_array.clone()]));
464        assert!(tx.deleted_groups().eq([deleted_group.clone()]));
465        assert!(tx.deleted_arrays().eq([deleted_array.clone()]));
466        assert!(tx.updated_groups().eq([updated_group.clone()]));
467        let chunks =
468            Vec::from_iter(tx.updated_chunks().map(|(id, it)| (id, Vec::from_iter(it))));
469        assert_eq!(chunks, vec![(chunk_added.clone(), vec![ChunkIndices(vec![0])])]);
470        assert_eq!(
471            tx.updated_chunks_counts().collect::<Vec<_>>(),
472            vec![(chunk_added.clone(), 1)]
473        );
474
475        let added_group2 = NodeId::random();
476        let deleted_group2 = NodeId::random();
477        let deleted_array2 = NodeId::random();
478        let updated_group2 = NodeId::random();
479        let chunk_added2 = NodeId::random();
480        let mut cs2 = ChangeSet::for_edits();
481        cs2.add_group("/g1".try_into().unwrap(), added_group2.clone(), Bytes::new())?;
482        cs2.delete_group("/g2".try_into().unwrap(), &deleted_group2)?;
483        cs2.add_array(
484            "/a1".try_into().unwrap(),
485            added_array.clone(),
486            ArrayData {
487                shape: ArrayShape::new([(0, 10)]).unwrap(),
488                dimension_names: None,
489                user_data: Bytes::new(),
490            },
491        )?;
492        cs2.delete_array("/a2".try_into().unwrap(), &deleted_array2)?;
493        cs2.update_group(&updated_group2, &"/g3".try_into().unwrap(), Bytes::new())?;
494        cs2.set_chunk_ref(chunk_added.clone(), ChunkIndices(vec![0]), None)?;
495        cs2.set_chunk_ref(
496            chunk_added.clone(),
497            ChunkIndices(vec![1]),
498            Some(ChunkPayload::Inline(Bytes::new())),
499        )?;
500        cs2.set_chunk_ref(chunk_added.clone(), ChunkIndices(vec![42]), None)?;
501        cs2.set_chunk_ref(
502            chunk_added2.clone(),
503            ChunkIndices(vec![7]),
504            Some(ChunkPayload::Inline(Bytes::new())),
505        )?;
506
507        let t3 = transaction_log_from_change_set(&SnapshotId::random(), &cs2);
508        let tx_id = SnapshotId::random();
509        let tx = TransactionLog::merge(&tx_id, [&t1, &t2, &t3]);
510
511        assert!(
512            tx.new_groups()
513                .sorted()
514                .eq([added_group.clone(), added_group2.clone()].into_iter().sorted())
515        );
516        assert!(tx.new_arrays().eq([added_array.clone()]));
517        assert!(
518            tx.deleted_groups().sorted().eq([
519                deleted_group.clone(),
520                deleted_group2.clone()
521            ]
522            .into_iter()
523            .sorted())
524        );
525        assert!(
526            tx.deleted_arrays().sorted().eq([
527                deleted_array.clone(),
528                deleted_array2.clone()
529            ]
530            .into_iter()
531            .sorted())
532        );
533        assert!(
534            tx.updated_groups().sorted().eq([
535                updated_group.clone(),
536                updated_group2.clone()
537            ]
538            .into_iter()
539            .sorted())
540        );
541        let chunks = HashSet::from_iter(
542            tx.updated_chunks().map(|(id, it)| (id, Vec::from_iter(it))),
543        );
544        assert_eq!(
545            chunks,
546            HashSet::from([
547                (
548                    chunk_added.clone(),
549                    vec![
550                        ChunkIndices(vec![0]),
551                        ChunkIndices(vec![1]),
552                        ChunkIndices(vec![42])
553                    ]
554                ),
555                (chunk_added2.clone(), vec![ChunkIndices(vec![7]),])
556            ])
557        );
558
559        assert_eq!(
560            tx.updated_chunks_counts().collect::<HashSet<_>>(),
561            HashSet::from([(chunk_added.clone(), 3), (chunk_added2, 1)])
562        );
563        Ok(())
564    }
565}