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()
})
})
}
#[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),
}
#[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
}
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]); self.buffer.push(value);
Ok(())
}
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(())
}
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); self.buffer.extend(value.to_be_bytes());
Ok(())
}
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..]; 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(())
}
}
struct ObjectInput {
keys: Vec<Bytes>,
object: Bytes,
}
pub struct MosaicCreator {
pub path: PathBuf,
cursor: Position,
offset_container_metadata: Position,
offset_root_size: Position,
offset_object_counter: Position,
pub max_object_size: Size,
pub objects_counter: u64,
pub objects_total_size: Size,
pub comments: Vec<String>,
idx_descriptions: Vec<IdxDescription>,
compress_tx: Option<crossbeam_channel::Sender<ObjectInput>>,
compression_workers: Vec<thread::JoinHandle<()>>,
compressed_objects_thread: Option<std::thread::JoinHandle<TileCreator>>,
}
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)
}
fn max_equivalent_vint(number: u64) -> Result<Vec<u8>, VIntError> {
let mut max_equivalent_vint = number.as_vint()?;
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 {
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()), };
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)?;
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)
}
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(())
}
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(())
}
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(())
}
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(),
)?;
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; 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)?;
}
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)
}
}
struct TileCreator {
writer: MosaicWriter,
indexes: Vec<BTreeMap<Bytes, Position>>,
pub tile_threshold: usize,
tile_prefix_length: Size,
tile_size_length: Option<Size>,
cursor: Position,
current_tile: BufferedMasterElement,
}
impl TileCreator {
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,
}
}
fn stack_object(&mut self, keys: Vec<Bytes>, compressed: &[u8]) -> Result<()> {
let binary_len = BufferedMasterElement::binary_len(compressed)?;
if binary_len + self.current_tile.len() >= self.tile_threshold {
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(())
}
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() {
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<()> {
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(())
}
}