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