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