use crate::backends::mmap::MmapMosaicBackend;
use crate::backends::MosaicBackend;
use crate::creator::max_object_size;
use crate::ebml::{MosaicTag, DOCTYPE_READ_VERSION};
use crate::{CompressionMethod, IdxDescription, Position, Size};
use anyhow::{ensure, Result};
use ph::fmph::GOFunction;
use rayon::iter::{IntoParallelRefIterator, ParallelIterator};
use std::fmt;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use thiserror::Error;
use zstd::bulk::decompress;
#[derive(Error, Debug)]
pub enum MosaicReaderError {
#[error("Bad checksum for tag {tag}")]
BadChecksum { tag: MosaicTag },
#[error("Bad key size, expected {expected} bytes, got {actual} bytes")]
BadKeySize { expected: u64, actual: u64 },
#[error("{filename}: Error, {broken_tiles} corrupted tiles and {broken_indexes} corrupted indexes found")]
BrokenMosaic {
filename: String,
broken_tiles: i32,
broken_indexes: i32,
},
#[error("Element {tag} not found")]
ElementNotFound { tag: MosaicTag },
#[error("Index {idx_description} not found")]
IndexNotFound { idx_description: IdxDescription },
#[error("Missing checksum for element {tag}")]
MissingChecksum { tag: MosaicTag },
#[error("Unknown index description {description}")]
UnknownIdxDescription { description: String },
}
pub struct Cursor<B: MosaicBackend> {
pub backend: Arc<B>,
pub offset: Position,
mark: Option<Position>,
}
impl<B: MosaicBackend> Cursor<B> {
pub fn new(backend: Arc<B>) -> Self {
Cursor {
backend,
offset: Position(0),
mark: None,
}
}
pub fn new_at(backend: Arc<B>, offset: Position) -> Self {
Cursor {
backend,
offset,
mark: None,
}
}
pub fn next_tag_and_size(&mut self) -> Result<(MosaicTag, Size)> {
let (tag, size, new_offset) = self.backend.find_tag_and_size(self.offset)?;
self.offset = new_offset;
Ok((tag, size))
}
pub fn next_one_tag_and_size(&mut self) -> Result<(MosaicTag, Size)> {
let (element, size, new_offset) = self.backend.read_one_tag_and_size(self.offset)?;
self.offset = new_offset;
Ok((element, size))
}
pub fn next_binary(&mut self, size: Size) -> Result<&[u8]> {
let (binary, new_offset) = self.backend.read_binary(self.offset, size)?;
self.offset = new_offset;
Ok(binary)
}
pub fn next_u8(&mut self) -> Result<u8> {
let (byte, new_offset) = self.backend.read_u8(self.offset)?;
self.offset = new_offset;
Ok(byte)
}
pub fn next_uint(&mut self, size: Size) -> Result<u64> {
let (uint, new_offset) = self.backend.read_uint(self.offset, size)?;
self.offset = new_offset;
Ok(uint)
}
pub fn next_utf8(&mut self, size: Size) -> Result<String> {
let (string, new_offset) = self.backend.read_utf8(self.offset, size)?;
self.offset = new_offset;
Ok(string)
}
pub fn seek_to_element(&mut self, element_tag: MosaicTag) -> Result<Size> {
let (size, new_offset) = self.backend.find_element(self.offset, element_tag)?;
self.offset = new_offset;
Ok(size)
}
pub fn seek_to_element_pre(&mut self, element_tag: MosaicTag) -> Result<(Position, Size)> {
let (tag_offset, size, new_offset) =
self.backend.find_element_pre(self.offset, element_tag)?;
self.offset = new_offset;
Ok((tag_offset, size))
}
pub fn seek_to_and_check(&mut self, master_tag: MosaicTag) -> Result<Size> {
let (size, new_offset) = self.backend.find_and_check(self.offset, master_tag)?;
self.offset = new_offset;
Ok(size)
}
pub fn skip(&mut self, size: Size) -> Result<()> {
ensure!(self.offset + size <= self.backend.len());
self.offset += size;
Ok(())
}
pub fn mark_master(&mut self, size: Size) -> Result<()> {
self.mark = Some(self.offset + size);
Ok(())
}
pub fn skip_master(&mut self) -> Result<()> {
match self.mark {
Some(position) => {
self.mark = None;
self.offset = position;
Ok(())
}
None => panic!("skip_master called without mark_master"),
}
}
}
pub struct MosaicReader<B: MosaicBackend> {
pub path: PathBuf,
pub backend: Arc<B>, pub doctype: String,
pub doctype_version: u8,
pub doctype_read_version: u8,
pub objects_counter: u64,
pub objects_total_size: Size,
pub end_of_tiles_offset: Position,
pub comments: Vec<String>,
pub idx: Option<IdxDescription>,
mph: Option<Arc<GOFunction>>, idx_unrolled_entry_size: Size,
idx_unrolled_start: Position,
pub compression: CompressionMethod,
pub object_size_bound: usize,
}
impl<B: MosaicBackend> Clone for MosaicReader<B> {
fn clone(&self) -> Self {
let MosaicReader {
path,
backend,
doctype,
doctype_version,
doctype_read_version,
objects_counter,
objects_total_size,
end_of_tiles_offset,
comments,
idx,
mph,
idx_unrolled_entry_size,
idx_unrolled_start,
compression,
object_size_bound,
} = self;
MosaicReader {
path: path.clone(),
backend: backend.clone(),
doctype: doctype.clone(),
doctype_version: *doctype_version,
doctype_read_version: *doctype_read_version,
objects_counter: *objects_counter,
objects_total_size: *objects_total_size,
end_of_tiles_offset: *end_of_tiles_offset,
comments: comments.clone(),
idx: *idx,
mph: mph.clone(),
idx_unrolled_entry_size: *idx_unrolled_entry_size,
idx_unrolled_start: *idx_unrolled_start,
compression: *compression,
object_size_bound: *object_size_bound,
}
}
}
impl<B: MosaicBackend> fmt::Debug for MosaicReader<B> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> Result<(), fmt::Error> {
f.debug_struct("MosaicReader")
.field("path", &self.path)
.field("doctype_version", &self.doctype_version)
.field("objects_counter", &self.objects_counter)
.field("objects_total_size", &self.objects_total_size)
.finish_non_exhaustive()
}
}
impl<B: MosaicBackend> MosaicReader<B> {
pub fn with_backend(path: &Path, backend: Arc<B>) -> Result<Self> {
let mut cursor = Cursor::new(backend.clone());
let size = cursor.seek_to_element(MosaicTag::Ebml)?;
assert!(size >= 33);
cursor.seek_to_element(MosaicTag::EbmlVersion)?;
let ebml_version = cursor.next_u8()?;
cursor.seek_to_element(MosaicTag::EbmlReadVersion)?;
let ebml_read_version = cursor.next_u8()?;
assert!(
ebml_version >= ebml_read_version,
"EBML read version ({}) is bigger than EBML version ({})",
ebml_read_version,
ebml_version
);
cursor.seek_to_element(MosaicTag::EbmlMaxIdLength)?;
let _ebml_max_id_length = cursor.next_u8()?;
cursor.seek_to_element(MosaicTag::EbmlMaxSizeLength)?;
let _ebml_max_size_length = cursor.next_u8()?;
let size = cursor.seek_to_element(MosaicTag::DocType)?;
let doctype = cursor.next_utf8(size)?;
cursor.seek_to_element(MosaicTag::DocTypeVersion)?;
let doctype_version = cursor.next_u8()?;
cursor.seek_to_element(MosaicTag::DocTypeReadVersion)?;
let doctype_read_version = cursor.next_u8()?;
assert!(
DOCTYPE_READ_VERSION >= doctype_read_version,
"This version of swh-mosaic requires DocTypeReadVersion <= {},
this MOSAIC file is version {}. If possible, install a newer version of swh-mosaic",
DOCTYPE_READ_VERSION,
doctype_read_version
);
assert!(
doctype_version >= doctype_read_version,
"EBML read version ({}) is bigger than EBML version ({})",
doctype_version,
doctype_read_version
);
assert!(cursor.offset >= 38);
cursor.seek_to_element(MosaicTag::Mosaic)?;
cursor.seek_to_and_check(MosaicTag::ContainerMetaData)?;
let size = cursor.seek_to_element(MosaicTag::ObjectsCounter)?;
let objects_counter = cursor.next_uint(size)?;
let size = cursor.seek_to_element(MosaicTag::ObjectsTotalSize)?;
let objects_total_size = cursor.next_uint(size)?;
let size = cursor.seek_to_element(MosaicTag::EndOfTilesOffset)?;
let end_of_tiles_offset = cursor.next_uint(size)?;
let size = cursor.seek_to_element(MosaicTag::CompressionMethod)?;
let compression_name = cursor.next_utf8(size)?;
let compression = CompressionMethod::try_from(compression_name)?;
let mut comments = Vec::new();
let mut tag_and_size_or_err = cursor.seek_to_element(MosaicTag::Comment);
while tag_and_size_or_err.is_ok() {
let size = tag_and_size_or_err?;
let comment = cursor.next_utf8(size)?;
comments.push(comment);
tag_and_size_or_err = cursor.seek_to_element(MosaicTag::Comment)
}
let tile_size_or_err = cursor.seek_to_element(MosaicTag::Tile);
let object_size_bound = if let Ok(size) = tile_size_or_err {
max_object_size(size.0)?.try_into()?
} else {
0
};
Ok(MosaicReader {
path: path.to_path_buf(),
backend,
doctype,
doctype_version,
doctype_read_version,
objects_counter,
objects_total_size: objects_total_size.into(),
end_of_tiles_offset: (end_of_tiles_offset as usize).into(),
comments,
idx: None,
mph: None,
idx_unrolled_entry_size: 0.into(),
idx_unrolled_start: 0.into(),
compression,
object_size_bound,
})
}
pub fn load_index(&mut self, idx_description: IdxDescription) -> Result<()> {
let mut cursor = Cursor::new_at(self.backend.clone(), self.end_of_tiles_offset);
loop {
let idx_seek = cursor.seek_to_element(MosaicTag::Index);
if let Ok(idx_size) = idx_seek {
cursor.mark_master(idx_size)?;
} else {
return Err(MosaicReaderError::IndexNotFound { idx_description }.into());
}
let description_size = cursor.seek_to_element(MosaicTag::IdxDescription)?;
let description = cursor.next_utf8(description_size)?;
if description == idx_description.description() {
break;
} else {
cursor.skip_master()?;
}
}
self.idx = Some(idx_description);
let idx_unrolled_size = cursor.seek_to_element(MosaicTag::IdxUnrolled)?;
self.idx_unrolled_start = cursor.offset;
cursor.skip(idx_unrolled_size)?;
let idx_unrolled_entry_size_size =
cursor.seek_to_element(MosaicTag::IdxUnrolledEntrySize)?;
self.idx_unrolled_entry_size = cursor.next_uint(idx_unrolled_entry_size_size)?.into();
cursor.seek_to_and_check(MosaicTag::MapContainer)?;
let mph_size = cursor.seek_to_element(MosaicTag::Map)?;
let mut mph_vec = cursor.next_binary(mph_size)?;
self.mph = Some(Arc::new(GOFunction::read(&mut mph_vec)?));
Ok(())
}
pub fn list_indexes(&self) -> Result<Vec<(IdxDescription, Position)>> {
let mut indexes = Vec::new();
let mut cursor = Cursor::new_at(self.backend.clone(), self.end_of_tiles_offset);
while let Ok((index_offset, size)) = cursor.seek_to_element_pre(MosaicTag::Index) {
cursor.mark_master(size)?;
let description_size = cursor.seek_to_element(MosaicTag::IdxDescription)?;
let description = cursor.next_utf8(description_size)?;
let idx_description = IdxDescription::from_description(description.as_str());
match idx_description {
Ok(val) => indexes.push((val, index_offset)),
Err(err) => match err.downcast::<MosaicReaderError>()? {
MosaicReaderError::UnknownIdxDescription { description } => {
println!("Skipping unknown index description: {description}")
}
_ => panic!("Unexpected failure while trying to list indexes"),
},
}
cursor.skip_master()?;
}
Ok(indexes)
}
pub fn find_key(&self, key: &[u8]) -> Result<Option<(Position, Position)>> {
let mph_key = self
.mph
.as_ref()
.unwrap_or_else(|| panic!("No MPH loaded"))
.get(key);
let Some(mph_key) = mph_key else {
return Ok(None);
};
let offset = self.idx_unrolled_start + Size(mph_key * self.idx_unrolled_entry_size.0);
let (key_tag, key_size, key_offset) = self.backend.find_tag_and_size(offset)?;
if key_tag != MosaicTag::Key {
return Ok(None);
}
let (key_read, _) = self.backend.read_binary(key_offset, key_size)?;
if key_read != key {
return Ok(None);
}
Ok(Some((key_offset, offset)))
}
pub fn find_object_offset(
&self,
key_size: Size,
key_offset: Position,
) -> Result<Option<(Size, Position, Position)>> {
let (_, offset) = self.backend.read_binary(key_offset, key_size)?;
let (offset_tag, offset_size, offset) = self.backend.find_tag_and_size(offset)?;
if offset_tag != MosaicTag::Offset {
return Ok(None);
}
let offset = (self.backend.read_uint(offset, offset_size)?.0 as usize).into();
let (object_tag, object_size, object_offset) = self.backend.find_tag_and_size(offset)?;
if object_tag != MosaicTag::Object {
return Ok(None);
}
Ok(Some((object_size, object_offset, offset)))
}
pub fn lookup(&self, key: &[u8]) -> Result<Option<Vec<u8>>> {
let expected_key_len = match self.idx {
Some(val) => val.key_len(),
None => panic!("Trying to look up a key with no index loaded. Make sure to call MosaicReader::load_index() first."),
};
if u64::try_from(key.len())? != expected_key_len.0 {
return Err(MosaicReaderError::BadKeySize {
expected: expected_key_len.0,
actual: key.len().try_into()?,
}
.into());
}
let Some((key_offset, _)) = self.find_key(key)? else {
return Ok(None);
};
let key_size = Size(key.len().try_into()?);
let Some((object_size, offset, _)) = self.find_object_offset(key_size, key_offset)? else {
return Ok(None);
};
let object = self.read_object(offset, object_size)?;
Ok(Some(object))
}
pub fn get_batch(&self, keys: Vec<Vec<u8>>) -> Result<Vec<Option<Vec<u8>>>> {
keys.par_iter().map(|k| self.lookup(k)).collect()
}
pub fn find_object_tile(&self, offset: Position) -> Result<Position> {
let (_, mosaic_offset) = self.backend.find_element(0.into(), MosaicTag::Mosaic)?;
let (md_size, md_offset) = self
.backend
.find_element(mosaic_offset, MosaicTag::ContainerMetaData)?;
let mut tile_offset = md_offset + md_size;
let (mut next_tile_offset, _, _) = self
.backend
.find_element_pre(tile_offset, MosaicTag::Tile)?;
while tile_offset < offset && tile_offset != next_tile_offset {
tile_offset = next_tile_offset;
(next_tile_offset, _, _) = self
.backend
.find_element_pre(tile_offset, MosaicTag::Tile)?;
}
Ok(tile_offset)
}
pub fn read_object(&self, offset: Position, size: Size) -> Result<Vec<u8>> {
let raw = self.backend.read_binary(offset, size)?.0;
match self.compression {
CompressionMethod::Zstd => {
let decompressed = decompress(raw, self.object_size_bound)?;
Ok(decompressed)
}
CompressionMethod::None => Ok(raw.to_vec()),
_ => anyhow::bail!("Not implemented yet: decompression of {}", self.compression),
}
}
pub fn close(&self) -> Result<()> {
Ok(())
}
pub fn iter(&self) -> MosaicIterator<B> {
MosaicIterator::new(self.clone())
}
pub fn keys(&self) -> MosaicKeysIterator<B> {
MosaicKeysIterator::new(self.clone())
}
pub fn values(&self) -> MosaicValuesIterator<B> {
MosaicValuesIterator::new(self.clone())
}
}
pub struct MosaicIterator<B: MosaicBackend> {
reader: MosaicReader<B>,
offset: Position,
}
impl<B: MosaicBackend> MosaicIterator<B> {
pub fn new(reader: MosaicReader<B>) -> Self {
let offset = reader.idx_unrolled_start;
Self { reader, offset }
}
}
impl<B: MosaicBackend> Iterator for MosaicIterator<B> {
type Item = Result<(Vec<u8>, Vec<u8>)>;
fn next(&mut self) -> Option<Self::Item> {
while let Ok((tag, size, offset)) = self.reader.backend.find_tag_and_size(self.offset) {
if tag == MosaicTag::Key {
if let Ok((key, offset)) = self.reader.backend.read_binary(offset, size) {
if let Ok((tag, size, offset)) = self.reader.backend.find_tag_and_size(offset) {
if tag == MosaicTag::Offset {
if let Ok((obj_offset, offset)) =
self.reader.backend.read_uint(offset, size)
{
self.offset = offset;
if let Ok((object_tag, object_size, obj_offset)) = self
.reader
.backend
.find_tag_and_size((obj_offset as usize).into())
{
if object_tag == MosaicTag::Object {
let key = key.to_vec();
if let Ok(object) =
self.reader.read_object(obj_offset, object_size)
{
return Some(Ok((key, object)));
}
}
}
continue;
}
}
}
}
}
return None;
}
None
}
}
pub struct MosaicKeysIterator<B: MosaicBackend> {
cursor: Cursor<B>,
}
impl<B: MosaicBackend> MosaicKeysIterator<B> {
pub fn new(reader: MosaicReader<B>) -> Self {
let cursor = Cursor::new_at(reader.backend.clone(), reader.idx_unrolled_start);
Self { cursor }
}
}
impl<B: MosaicBackend> Iterator for MosaicKeysIterator<B> {
type Item = Result<Vec<u8>>;
fn next(&mut self) -> Option<Self::Item> {
let key_size = match self.cursor.next_tag_and_size() {
Ok((tag, size)) => {
if tag != MosaicTag::Key {
return None;
}
size
}
Err(_) => return None,
};
let key = match self.cursor.next_binary(key_size) {
Ok(key) => key.to_vec(),
Err(_) => return None,
};
let offset_size = match self.cursor.next_tag_and_size() {
Ok((tag, size)) => {
if tag != MosaicTag::Offset {
return None;
}
size
}
Err(_) => return None,
};
if let Err(e) = self.cursor.skip(offset_size) {
return Some(Err(e));
}
Some(Ok(key))
}
}
pub struct MosaicValuesIterator<B: MosaicBackend> {
reader: MosaicReader<B>,
cursor: Cursor<B>,
}
impl<B: MosaicBackend> MosaicValuesIterator<B> {
pub fn new(reader: MosaicReader<B>) -> Self {
let mut cursor = Cursor::new(reader.backend.clone());
if MosaicValuesIterator::prepare_cursor(&mut cursor).is_err() {
cursor.offset = reader.backend.len().into();
}
Self { reader, cursor }
}
fn prepare_cursor(cursor: &mut Cursor<B>) -> Result<()> {
if cursor.seek_to_element(MosaicTag::Mosaic).is_ok()
&& cursor.seek_to_element(MosaicTag::Tile).is_ok()
{
if let Ok((tag, size)) = cursor.next_tag_and_size() {
if tag == MosaicTag::Crc32 {
cursor.skip(size)?;
return Ok(());
}
}
}
anyhow::bail!("Failed to find a Tile");
}
}
impl<B: MosaicBackend> Iterator for MosaicValuesIterator<B> {
type Item = Result<Vec<u8>>;
fn next(&mut self) -> Option<Self::Item> {
let object_size = match self.cursor.next_tag_and_size() {
Ok((tag, size)) => match tag {
MosaicTag::Object => size,
MosaicTag::Tile => {
return self.next();
}
MosaicTag::Crc32 => {
if let Err(e) = self.cursor.skip(size) {
return Some(Err(e));
}
return self.next();
}
_ => return None,
},
Err(_) => return None,
};
let obj = self.reader.read_object(self.cursor.offset, object_size);
if let Ok(object) = obj {
if self.cursor.skip(object_size).is_ok() {
Some(Ok(object))
} else {
None
}
} else {
None
}
}
}
pub type MmapMosaicReader = MosaicReader<MmapMosaicBackend>;
impl MmapMosaicReader {
pub fn new(path: &Path, indexed_by: IdxDescription) -> Result<Self> {
let backend = MmapMosaicBackend::new(path)?;
let mut reader = MosaicReader::with_backend(path, Arc::new(backend))?;
reader.load_index(indexed_by)?;
Ok(reader)
}
}