swh-mosaic 0.3.1

MOdular Storage of Archived and Indexed Contents from Software Heritage
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
// Copyright (C) 2026  The Software Heritage developers
// See the AUTHORS file at the top-level directory of this distribution
// License: GNU General Public License version 3, or any later version
// See top-level LICENSE file for more information

//!
//! Tools to create MOSAIC files
//!

use crate::ebml::{
    crc32, read_vint, MosaicTag, ShortestBeBytes, VIntError, Vint, CRC32_SIZE, DOCTYPE,
    DOCTYPE_READ_VERSION, DOCTYPE_VERSION, EBML_MAX_ID_LENGTH, EBML_MAX_SIZE_LENGTH,
};
use crate::writer::MosaicWriter;
use crate::{CompressionMethod, IdxDescription, Position, Size};
use anyhow::{Error, Result};
use bytes::Bytes;
use crossbeam_channel;
use log::error;
use ph::fmph::keyset::CachedKeySet;
use ph::fmph::GOFunction;
use rayon::prelude::*;
use std::collections::BTreeMap;
use std::fs::File;
use std::io::{Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
use std::sync::OnceLock;
use std::thread;
use thiserror::Error;
use zstd::bulk::Compressor;

static OBJECTS_QUEUE_SIZE: OnceLock<usize> = OnceLock::new();
static NB_COMPRESSION_THREADS: OnceLock<usize> = OnceLock::new();

const DEFAULT_OBJECTS_QUEUE_SIZE: usize = 100000;

fn get_objects_queue_size() -> usize {
    *OBJECTS_QUEUE_SIZE.get_or_init(|| {
        std::env::var("MOSAIC_OBJECTS_QUEUE_SIZE")
            .ok()
            .and_then(|v| v.parse().ok())
            .unwrap_or(DEFAULT_OBJECTS_QUEUE_SIZE)
    })
}

fn get_nb_compression_threads() -> usize {
    *NB_COMPRESSION_THREADS.get_or_init(|| {
        std::env::var("RAYON_NUM_THREADS")
            .ok()
            .and_then(|v| v.parse().ok())
            .unwrap_or_else(|| {
                std::thread::available_parallelism()
                    .unwrap_or_else(|_| 1.try_into().unwrap())
                    .get()
            })
    })
}

/// Typed errors for MosaicCreator
#[derive(Error, Debug)]
pub enum MosaicCreatorError {
    #[error("Object bigger than a tile's maximum size ({0} bytes). Maybe increase tile_threshold to increase this maximum size.")]
    ObjectTooBig(u64),
}

/// This provides EBML writing primitives and allows to encapsulate master elements.
#[derive(Default)]
pub struct BufferedMasterElement {
    pub buffer: Vec<u8>,
}

impl BufferedMasterElement {
    pub fn new() -> Self {
        Self::default()
    }

    pub fn with_capacity(capacity: usize) -> Self {
        BufferedMasterElement {
            buffer: Vec::with_capacity(capacity),
        }
    }

    pub fn with_padding(padding: usize) -> Self {
        BufferedMasterElement {
            buffer: vec![0; padding],
        }
    }

    pub fn len(&self) -> usize {
        self.buffer.len()
    }

    pub fn is_empty(&self) -> bool {
        self.buffer.len() == 0
    }

    /// Returns the maximum number of bytes that [`Self::append_binary`] would insert
    pub fn binary_len(value: &[u8]) -> Result<usize, Error> {
        Ok(value.len()
            + usize::try_from(EBML_MAX_ID_LENGTH.0).unwrap()
            + usize::try_from(EBML_MAX_SIZE_LENGTH.0).unwrap())
    }
    pub fn append_binary(&mut self, tag: MosaicTag, value: &[u8]) -> Result<(), Error> {
        self.buffer.extend_from_slice(tag.to_be_bytes());
        let size: u64 = value.len().try_into()?;
        self.buffer.extend(&size.as_vint()?);
        self.buffer.extend(value);
        Ok(())
    }

    pub fn append_u8(&mut self, tag: MosaicTag, value: u8) -> Result<(), Error> {
        self.buffer.extend_from_slice(tag.to_be_bytes());
        self.buffer.extend([0x81]); // 1 as a VInt
        self.buffer.push(value);
        Ok(())
    }

    /// this shortcuts size appending because length will always fit in a 1-byte VInt
    pub fn append_uint(&mut self, tag: MosaicTag, value: u64) -> Result<(), Error> {
        self.buffer.extend_from_slice(tag.to_be_bytes());
        let sliced = value.shortest_be_bytes();
        let size = sliced.len() as u8 | 0x80;
        self.buffer.push(size);
        if !sliced.is_empty() {
            self.buffer.extend(sliced);
        }
        Ok(())
    }

    /// this forces appending an 8-bytes unsigned integer
    pub fn append_full_uint(&mut self, tag: MosaicTag, value: u64) -> Result<(), Error> {
        self.buffer.extend_from_slice(tag.to_be_bytes());
        self.buffer.push(0x88); // "8" as VInt
        self.buffer.extend(value.to_be_bytes());
        Ok(())
    }

    /// this forces appending an unsigned integer of length `size`
    pub fn append_sized_uint(
        &mut self,
        tag: MosaicTag,
        value: u64,
        size: usize,
    ) -> Result<(), Error> {
        self.buffer.extend_from_slice(tag.to_be_bytes());
        self.buffer
            .extend_from_slice(&u64::try_from(size)?.as_vint()?);
        let sized_value = &value.to_be_bytes()[8 - size..]; // last `size` bytes of value
        debug_assert_eq!(sized_value.len(), size);
        self.buffer.extend(sized_value);
        Ok(())
    }

    pub fn append_utf8(&mut self, tag: MosaicTag, value: String) -> Result<(), Error> {
        self.buffer.extend_from_slice(tag.to_be_bytes());
        let size: u64 = value.len().try_into()?;
        self.buffer.extend(&size.as_vint()?);
        self.buffer.extend(value.as_bytes());
        Ok(())
    }

    pub fn append_crc32(&mut self, content: &BufferedMasterElement) -> Result<(), Error> {
        let crc32 = crc32(&content.buffer);
        self.buffer
            .extend_from_slice(MosaicTag::Crc32.to_be_bytes());
        self.buffer.extend_from_slice(&4u8.as_vint()?);
        self.buffer.extend_from_slice(&crc32);
        Ok(())
    }

    pub fn append_master(
        &mut self,
        tag: MosaicTag,
        content: BufferedMasterElement,
        crc32: bool,
    ) -> Result<()> {
        self.buffer.extend_from_slice(tag.to_be_bytes());
        let mut size: Size = content.len().try_into()?;
        if crc32 {
            size += CRC32_SIZE;
            self.buffer.extend(&size.0.as_vint()?);
            self.append_crc32(&content)?;
        } else {
            self.buffer.extend_from_slice(&size.0.as_vint()?);
        }
        self.buffer.extend_from_slice(&content.buffer);
        Ok(())
    }
}

/// Convenience wrapper for passing objects around threads
struct ObjectInput {
    keys: Vec<Bytes>,
    object: Bytes,
}

/// MOSAIC file creator
/// -------------------
///
/// Indexes and some important header fields are only written on `close()`, so do not
/// forget to call it. Objects' keys must be computed by the caller.
///
/// Indexes' construction is multi-threaded, you can set the `RAYON_NUM_THREADS`
/// environment variable to restrict the number of threads used.
/// When compression is activated is is multi-threaded too, following the same setting.
///
/// The file is written by a separate thread, fed from a queue whose size can be
/// overridden by the `MOSAIC_OBJECTS_QUEUE_SIZE` environment variable.
pub struct MosaicCreator {
    pub path: PathBuf,

    /// current offset in the written file
    cursor: Position,

    /// offset of the ContainerMetaData. Useful to rewrite the CRC32 when closing a MOSAIC file.
    offset_container_metadata: Position,

    /// offset of the VInt measuring the total length of the MOSAIC root element. When starting a
    /// file we set it a zero and seek() back before closing to write the real value.
    offset_root_size: Position,

    /// offset of the ObjectsCounter element. When starting a file we set it a zero over 8 bytes,
    /// and seek() back before closing to write the real value.
    /// We assume the ObjectsTotalSize and EndOfTilesOffset elements follow, so we write the three
    /// at the same time.
    offset_object_counter: Position,

    pub max_object_size: Size,
    pub objects_counter: u64,
    pub objects_total_size: Size,
    pub comments: Vec<String>,

    /// list of indexes that will be written (order matters)
    idx_descriptions: Vec<IdxDescription>,

    /// Sender to the compression threads, or directly to the writer thread
    /// if compression is disabled
    compress_tx: Option<crossbeam_channel::Sender<ObjectInput>>,

    /// Compression worker threads. Empty when compression is disabled.
    compression_workers: Vec<thread::JoinHandle<()>>,

    /// objects ready to be written are handled in a separate thread.
    compressed_objects_thread: Option<std::thread::JoinHandle<TileCreator>>,
}

/// Computes the maximum object size for a given tile threshold.
///
/// Objects larger than this value cannot be inserted into a MOSAIC file.
pub fn max_object_size(tile_threshold: u64) -> Result<u64> {
    let max_object_size_vint = max_equivalent_vint(tile_threshold)?;
    let tile_size_length: u64 = max_object_size_vint.len().try_into()?;
    let mut tile_prefix_length: u64 = MosaicTag::Tile.to_be_bytes().len().try_into()?;
    tile_prefix_length += tile_size_length;
    tile_prefix_length += CRC32_SIZE.0;
    Ok(read_vint(&max_object_size_vint)?.0 - tile_prefix_length)
}

/// Helper: computes the greatest VInt having the same length as `number` as VInt
fn max_equivalent_vint(number: u64) -> Result<Vec<u8>, VIntError> {
    let mut max_equivalent_vint = number.as_vint()?;

    // in the first byte the first 1 is the length mark; we set to 1 all bits after
    let mut first = max_equivalent_vint[0];
    let mut mask = max_equivalent_vint[0] >> 1;
    while mask > 0 {
        first |= mask;
        mask >>= 1;
    }
    max_equivalent_vint[0] = first;

    for b in max_equivalent_vint.iter_mut().skip(1) {
        *b = 0xff;
    }

    Ok(max_equivalent_vint)
}

impl MosaicCreator {
    /// Creates a MosaicCreator that will write to `path` (and fail if it does already
    /// exists).
    ///
    /// `tile_threshold` is the minimum number of bytes before the creator will open a
    /// new `Tile` element. This approximates well future Tiles' size and allows to
    /// compute the `max_object_size` attribute. From the number of bytes required to
    /// represent `tile_threshold` as an EBML VInt, we derive `max_object_size` as the
    /// greatest value of a VInt of that size.
    /// For example, the recommended tile_threshold=32M requires a 4-bytes VInt, the
    /// greatest 4-bytes VInt is 0x1FFFFFFF, whose contained valued is
    /// 0xFFFFFFF = 268_435_455 so the computed maximum object size will be ~268MB.
    /// Use [`creator::max_object_size`][crate::creator::max_object_size] to get this
    /// value before construction.
    ///
    /// Compression is enabled if provided. Depending on your CPU budget we advise to
    /// use between 3 and 15. `max_object_size` is applied on the uncompressed object's
    /// size so behavior is more easily predictable. Compression is multi-threaded
    /// using all available CPUs, or set the `RAYON_NUM_THREADS` environment
    /// variable to choose the number of compression threads.
    pub fn new(
        path: &Path,
        tile_threshold: usize,
        comments: Vec<String>,
        idx_descriptions: Vec<IdxDescription>,
        compression_level: Option<u8>,
    ) -> Result<Self> {
        let max_object_size_vint = max_equivalent_vint(tile_threshold.try_into()?)?;
        let tile_size_length: Size = max_object_size_vint.len().try_into()?;
        let tile_prefix_length = Size(MosaicTag::Tile.to_be_bytes().len().try_into()?)
            + tile_size_length.0.into()
            + CRC32_SIZE;
        let max_object_size = Size(read_vint(&max_object_size_vint)?.0) - tile_prefix_length;

        let compression_method = if compression_level.is_none() {
            CompressionMethod::None
        } else {
            CompressionMethod::Zstd
        };

        let (write_tx, write_rx) =
            crossbeam_channel::bounded::<ObjectInput>(get_objects_queue_size());

        let (compress_tx, compression_workers) = match compression_level {
            Some(level) => {
                let (compress_tx, compress_rx) =
                    crossbeam_channel::bounded::<ObjectInput>(get_objects_queue_size());
                let num_threads = get_nb_compression_threads();
                let workers: Vec<_> = (0..num_threads)
                    .map(|_| {
                        let compress_rx = compress_rx.clone();
                        let write_tx = write_tx.clone();
                        thread::spawn(move || {
                            let mut compressor = Compressor::new(level.into())
                                .expect("failed to create zstd compressor");
                            while let Ok(ObjectInput { keys, object }) = compress_rx.recv() {
                                let compressed =
                                    compressor.compress(&object).expect("compression failed");
                                write_tx
                                    .send(ObjectInput {
                                        keys,
                                        object: Bytes::from(compressed),
                                    })
                                    .expect("writer thread terminated prematurely");
                            }
                        })
                    })
                    .collect();
                (compress_tx, workers)
            }
            None => (write_tx, Vec::new()), // send directly to writer thread, not compression thread
        };

        let mut writer = MosaicWriter::new(File::create_new(path)?)?;

        let mut creator = MosaicCreator {
            path: path.to_path_buf(),
            cursor: 0.into(),
            offset_container_metadata: 0.into(),
            offset_root_size: 0.into(),
            offset_object_counter: 0.into(),
            max_object_size,
            objects_counter: 0,
            objects_total_size: 0.into(),
            comments,
            idx_descriptions: idx_descriptions.clone(),
            compress_tx: Some(compress_tx),
            compression_workers,
            compressed_objects_thread: None,
        };

        creator.write_header(&mut writer)?;

        let mut mosaic_root = MosaicTag::Mosaic.to_be_bytes().to_vec();
        creator.offset_root_size = creator.cursor + mosaic_root.len().try_into()?;
        mosaic_root.extend(0u64.as_vint_sized(EBML_MAX_SIZE_LENGTH)?);
        writer.file.write_all(&mosaic_root)?;
        creator.cursor += mosaic_root.len().try_into()?;

        creator.write_container_metadata(compression_method, &mut writer)?;

        // launch the tiles writer thread
        let cursor_clone = creator.cursor;
        let writer_thread = std::thread::Builder::new()
            .name("mosaic-tile-creator".to_string())
            .spawn(move || {
                let mut tile_creator = TileCreator::new(
                    writer,
                    idx_descriptions,
                    tile_threshold,
                    tile_prefix_length,
                    Some(tile_size_length),
                    cursor_clone,
                );
                while let Ok(ObjectInput { keys, object }) = write_rx.recv() {
                    if let Err(e) = tile_creator.stack_object(keys, &object) {
                        error!("Critical error, writing stopped: {e}");
                        break;
                    }
                }
                tile_creator
            })?;

        creator.compressed_objects_thread = Some(writer_thread);
        Ok(creator)
    }

    /// Writes EBML header.
    fn write_header(&mut self, writer: &mut MosaicWriter) -> Result<(), Error> {
        let mut buffered_master = BufferedMasterElement::with_capacity(100);
        buffered_master.append_uint(MosaicTag::EbmlVersion, 1)?;
        buffered_master.append_uint(MosaicTag::EbmlReadVersion, 1)?;
        buffered_master.append_uint(MosaicTag::EbmlMaxIdLength, EBML_MAX_ID_LENGTH.0)?;
        buffered_master.append_uint(MosaicTag::EbmlMaxSizeLength, EBML_MAX_SIZE_LENGTH.0)?;
        buffered_master.append_utf8(MosaicTag::DocType, DOCTYPE.to_string())?;
        buffered_master.append_uint(MosaicTag::DocTypeVersion, DOCTYPE_VERSION.into())?;
        buffered_master.append_uint(MosaicTag::DocTypeReadVersion, DOCTYPE_READ_VERSION.into())?;
        self.cursor +=
            writer.write_master(MosaicTag::Ebml, &buffered_master.buffer, false, None)?;
        Ok(())
    }

    /// Writes Mosaic's ContainerMetaData element, sets offset_object_counter by the way.
    ///
    /// At that offset one will find the 3 elements that have to be updated when closing,
    /// so the order of elements should match the one in close().
    fn write_container_metadata(
        &mut self,
        compression: CompressionMethod,
        writer: &mut MosaicWriter,
    ) -> Result<(), Error> {
        let mut buffered_master = BufferedMasterElement::with_capacity(1000);
        buffered_master.append_full_uint(MosaicTag::ObjectsCounter, 0)?;
        buffered_master.append_full_uint(MosaicTag::ObjectsTotalSize, 0)?;
        buffered_master.append_full_uint(MosaicTag::EndOfTilesOffset, 0)?;
        buffered_master.append_utf8(MosaicTag::CompressionMethod, compression.to_string())?;
        for comment in self.comments.clone() {
            buffered_master.append_utf8(MosaicTag::Comment, comment.to_string())?;
        }
        self.offset_container_metadata = self.cursor;
        self.cursor += writer.write_master(
            MosaicTag::ContainerMetaData,
            &buffered_master.buffer,
            true,
            None,
        )?;
        self.offset_object_counter = self.cursor - Size(buffered_master.len().try_into()?);
        Ok(())
    }

    /// Inserts an object, with keys for configured indexes.
    ///
    /// `keys[i]` corresponds to the `i`‑th index given at construction.
    ///
    /// If that object would increase the buffer above `tile_threshold` (or if it was
    /// longer already), we write the Tile element and open a buffer for new one. The
    /// heuristic is that `tile_threshold` will be orders of magnitude greater than the
    /// median object, but some objects will be (much) greater than `tile_threshold` so
    /// this should approximate `tile_threshold` enough (we don't plan to exploit any
    /// kind of block alignment).
    pub fn add(&mut self, keys: Vec<impl Into<Bytes>>, object: impl Into<Bytes>) -> Result<()> {
        let object = object.into();
        let keys: Vec<_> = keys.into_iter().map(|key| key.into()).collect();
        if object.len() >= self.max_object_size.0 as usize {
            return Err(MosaicCreatorError::ObjectTooBig(self.max_object_size.0).into());
        }
        if keys.len() != self.idx_descriptions.len() {
            anyhow::bail!(
                "expected {} keys, got {}",
                self.idx_descriptions.len(),
                keys.len()
            );
        }
        for (i, key) in keys.iter().enumerate() {
            let key_len = self.idx_descriptions[i].key_len();
            if key.len() != key_len.0 as usize {
                anyhow::bail!(
                    "key #{} ({:?}) length {} does not match expected {}",
                    i,
                    key,
                    key.len(),
                    key_len
                );
            }
        }
        self.objects_total_size += object.len().try_into()?;
        self.objects_counter += 1;

        self.compress_tx
            .as_ref()
            .expect("add() called after compress_tx was closed")
            .send(ObjectInput { keys, object })
            .expect("compression worker channel disconnected");

        Ok(())
    }

    /// Finish a MOSAIC file; once this is done, any call to the MosaicCreator will have
    /// an undefined behavior.
    pub fn close(&mut self) -> Result<()> {
        drop(self.compress_tx.take());
        for worker in self.compression_workers.drain(..) {
            worker.join().expect("compression worker panicked");
        }
        let thread = self
            .compressed_objects_thread
            .take()
            .expect("Writer's JoinHandle is None: maybe this has been closed already?");
        let mut tile_creator = thread
            .join()
            .expect("A fatal error occurred in the tile writer thread");
        self.cursor = tile_creator.close(0)?;
        let mut writer = tile_creator.writer;

        let end_of_tiles_offset = self.cursor;
        let max_offset_len = Size(u64::try_from(self.cursor.0)?.as_vint()?.len().try_into()?);

        for (idx, key_type) in self.idx_descriptions.clone().iter().enumerate() {
            let mut buffered_master = BufferedMasterElement::with_capacity(1000);
            buffered_master.append_utf8(
                MosaicTag::IdxDescription,
                key_type.description().to_string(),
            )?;

            // 4 = key tag + key len VInt + offset tag + offset len VInt
            // that's OK as long we have keys or offsets smaller than 127 bytes
            let idx_unrolled_entry_size = max_offset_len + key_type.key_len() + Size(4);

            let get_keys = || tile_creator.indexes[idx].iter().map(|m| m.0);
            let par_get_keys = || tile_creator.indexes[idx].par_iter().map(|m| m.0);
            let keys = (get_keys, par_get_keys);
            let num_keys = tile_creator.indexes[idx].len();
            let clone_threshold = 1_000; // arbitrary
            let key_set = CachedKeySet::dynamic_with_len(keys, num_keys, clone_threshold);

            let mph = GOFunction::new(key_set);

            let idx_unrolled = MosaicCreator::unroll_index(
                &tile_creator.indexes[idx],
                max_offset_len,
                &mph,
                idx_unrolled_entry_size,
            )?;
            buffered_master.append_master(MosaicTag::IdxUnrolled, idx_unrolled, false)?;

            buffered_master
                .append_uint(MosaicTag::IdxUnrolledEntrySize, idx_unrolled_entry_size.0)?;

            let mut mph_vec = Vec::new();
            mph.write(&mut mph_vec)?;

            let mut map_container =
                BufferedMasterElement::with_capacity(BufferedMasterElement::binary_len(&mph_vec)?);
            map_container.append_binary(MosaicTag::Map, &mph_vec)?;
            buffered_master.append_master(MosaicTag::MapContainer, map_container, true)?;

            self.cursor +=
                writer.write_master(MosaicTag::Index, &buffered_master.buffer, true, None)?;
        }

        // get back to the beginning of the file to write total length and counters
        let root_length = self.cursor - self.offset_root_size - 8.into();
        let root_length_vint = root_length.0.as_vint_sized(EBML_MAX_SIZE_LENGTH)?;
        writer
            .file
            .seek(SeekFrom::Start(self.offset_root_size.0.try_into()?))?;
        writer.file.write_all(&root_length_vint)?;

        let mut buffered_master = BufferedMasterElement::with_capacity(100);
        buffered_master.append_full_uint(MosaicTag::ObjectsCounter, self.objects_counter)?;
        buffered_master.append_full_uint(MosaicTag::ObjectsTotalSize, self.objects_total_size.0)?;
        buffered_master.append_full_uint(
            MosaicTag::EndOfTilesOffset,
            end_of_tiles_offset.0.try_into()?,
        )?;
        writer
            .file
            .seek(SeekFrom::Start(self.offset_object_counter.0.try_into()?))?;
        writer.file.write_all(&buffered_master.buffer)?;

        writer.rewrite_crc32(self.offset_container_metadata)?;
        Ok(())
    }

    fn unroll_index(
        offsets: &BTreeMap<Bytes, Position>,
        vint_size: Size,
        mph: &GOFunction,
        idx_unrolled_entry_size: Size,
    ) -> Result<BufferedMasterElement> {
        let idx_unrolled_entry_size = idx_unrolled_entry_size.0.try_into()?;
        let mut idx_unrolled = BufferedMasterElement::new();
        idx_unrolled.buffer = vec![0; offsets.len() * idx_unrolled_entry_size];

        for (key, offset) in offsets {
            let start: usize = mph
                .get(key)
                .unwrap_or_else(|| panic!("Failed to read MPH for key {:?}", key))
                .try_into()?;
            let start = start * idx_unrolled_entry_size;

            let mut entry = BufferedMasterElement::with_capacity(idx_unrolled_entry_size);
            entry.append_binary(MosaicTag::Key, key)?;
            entry.append_sized_uint(
                MosaicTag::Offset,
                u64::try_from(offset.0)?,
                vint_size.0 as usize,
            )?;
            idx_unrolled.buffer[start..start + entry.len()].copy_from_slice(&entry.buffer);
        }

        Ok(idx_unrolled)
    }
}

/// This struct writes Tile elements in a separated thread, that returns the array
/// of indexes (for each index given to `new`, a map of each key to the object's
/// absolute offset in the file)
struct TileCreator {
    /// Handles target file
    writer: MosaicWriter,

    /// per‑index map a key to the object's absolute offset in the file
    indexes: Vec<BTreeMap<Bytes, Position>>,

    /// minimum number of bytes before the creator will open a new `Tile` element. This
    /// approximates well future Tiles' size and allows to compute the
    /// `max_object_size` attribute. Use [`creator::max_object_size`][crate::creator::max_object_size]
    /// to get this value before construction.
    pub tile_threshold: usize,

    /// This is deduced from `tile_threshold`, and allows the program to predict an
    /// object's offset as soon as it is appended in a tile.
    tile_prefix_length: Size,

    /// Number of bytes to use when writing all Tile elements' size VINT. This is
    /// deduced from `tile_threshold` too. Wrapped in an Option because it fits
    /// `write_master`.
    tile_size_length: Option<Size>,

    /// current offset in the written file
    cursor: Position,

    current_tile: BufferedMasterElement,
}

impl TileCreator {
    /// `receiver` receives `[keys], object` pairs, ready to be packed
    pub fn new(
        writer: MosaicWriter,
        idx_descriptions: Vec<IdxDescription>,
        tile_threshold: usize,
        tile_prefix_length: Size,
        tile_size_length: Option<Size>,
        cursor: Position,
    ) -> Self {
        let indexes = idx_descriptions.iter().map(|_| BTreeMap::new()).collect();

        TileCreator {
            writer,
            indexes,
            tile_threshold,
            tile_prefix_length,
            tile_size_length,
            current_tile: BufferedMasterElement::with_capacity(tile_threshold),
            cursor,
        }
    }

    /// The post-compression side of `MosaicWriter.add()`, where writing actually happens
    fn stack_object(&mut self, keys: Vec<Bytes>, compressed: &[u8]) -> Result<()> {
        // this test also ensures tile's size do not grow above max_object_size,
        // ensuring size fits in tile_prefix_length
        let binary_len = BufferedMasterElement::binary_len(compressed)?;
        if binary_len + self.current_tile.len() >= self.tile_threshold {
            // if case binary_len > self.tile_threshold, avoid a future resize by the way.
            self.close(binary_len.max(self.tile_threshold))?;
        }

        let offset = self.cursor + self.tile_prefix_length + self.current_tile.len().try_into()?;
        for (i, key) in keys.into_iter().enumerate() {
            self.indexes[i].insert(key, offset);
        }

        self.current_tile
            .append_binary(MosaicTag::Object, compressed)?;
        Ok(())
    }

    /// Closes the current tile unless it is empty and sets the current tile to be a new
    /// `BufferedMasterElement` of given capacity.
    fn close(&mut self, next_capacity: usize) -> Result<Position> {
        if !self.current_tile.is_empty() {
            self.cursor += self.writer.write_master(
                MosaicTag::Tile,
                &self.current_tile.buffer,
                true,
                self.tile_size_length,
            )?;
            self.current_tile = BufferedMasterElement::with_capacity(next_capacity);
        }
        Ok(self.cursor)
    }
}

#[cfg(test)]
mod tests {
    use tempfile::TempDir;

    use super::*;

    #[test]
    fn test_max_equivalent_vint() {
        // when converting those hex values, be careful that max_equivalent_vint
        // returns a VInt so its first bit at 1 is the length marker, set it to 0 to
        // convert to the actual integer.
        assert_eq!(max_equivalent_vint(100).unwrap(), vec![0xFF]);
        assert_eq!(max_equivalent_vint(1000).unwrap(), vec![0x7F, 0xFF]);
        assert_eq!(
            max_equivalent_vint(32000000).unwrap(),
            vec![0x1F, 0xFF, 0xFF, 0xFF]
        );
        assert_eq!(
            max_equivalent_vint(1000000000).unwrap(),
            vec![0x0F, 0xFF, 0xFF, 0xFF, 0xFF]
        );
    }

    #[test]
    fn test_max_object_size() -> Result<()> {
        // max_object_size should copy the computation occurring in MosaicCreator
        let temp_dir = TempDir::new()?;
        let mosaic_path = temp_dir.path().join("tmp.mosaic");
        let creator = MosaicCreator::new(&mosaic_path, 3200, vec![], vec![], None)?;

        assert_eq!(creator.max_object_size.0, max_object_size(3200)?);

        let mosaic_path = temp_dir.path().join("tmp2.mosaic");
        let creator = MosaicCreator::new(&mosaic_path, 32000000, vec![], vec![], None)?;

        assert_eq!(creator.max_object_size.0, max_object_size(32000000)?);
        Ok(())
    }
}