1use 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 #[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 let mut builder = flatbuffers::FlatBufferBuilder::with_capacity(1_024 * 1_024);
42
43 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 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 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 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 let moved_nodes: &[WIPOffset<_>] = &[];
412 let moved_nodes = Some(builder.create_vector(moved_nodes)); 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, ignore_missing_null_terminator: true,
442};
443
444#[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}