use super::block::{HeaderBlock, PrimitiveBlock};
use crate::error::{new_blob_error, new_error, BlobError, ErrorKind, Result};
use bytes::Bytes;
use super::file_reader::FileReader;
use std::fs::File;
use std::io::{BufReader, Cursor, Read, Seek, SeekFrom};
use std::path::Path;
use std::sync::Arc;
pub(crate) use super::decompress::{
decompress_blob, decompress_blob_data_into, decompress_blob_raw,
decompress_wire_blob_into, pool_get_pub, pool_wrap, DecompressPool,
};
use super::decompress::{decompress_parsed_blob_into, pool_get};
pub(crate) use super::blob_wire::{
parse_blob_header_with_index, BlobData, BlobKind, WireBlob, WireBlobHeader,
};
pub use super::blob_wire::{MAX_BLOB_HEADER_SIZE, MAX_BLOB_MESSAGE_SIZE};
#[derive(Clone, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum BlobType<'a> {
OsmHeader,
OsmData,
Unknown(&'a str),
}
impl<'a> BlobType<'a> {
#[inline]
pub const fn as_str(&self) -> &'a str {
match self {
Self::OsmHeader => "OSMHeader",
Self::OsmData => "OSMData",
Self::Unknown(x) => x,
}
}
}
#[derive(Debug)]
#[non_exhaustive]
pub enum BlobDecode<'a> {
OsmHeader(Box<HeaderBlock>),
OsmData(PrimitiveBlock),
Unknown(&'a str),
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct ByteOffset(pub u64);
#[derive(Clone, Debug)]
pub struct Blob {
header: WireBlobHeader,
blob: WireBlob,
offset: Option<ByteOffset>,
}
impl Blob {
fn new(
header: WireBlobHeader,
blob: WireBlob,
offset: Option<ByteOffset>,
) -> Blob {
Blob {
header,
blob,
offset,
}
}
pub fn decode(&self) -> Result<BlobDecode<'_>> {
match self.get_type() {
BlobType::OsmHeader => {
let block = Box::new(self.to_headerblock()?);
Ok(BlobDecode::OsmHeader(block))
}
BlobType::OsmData => {
let block = self.to_primitiveblock()?;
Ok(BlobDecode::OsmData(block))
}
BlobType::Unknown(x) => Ok(BlobDecode::Unknown(x)),
}
}
#[inline]
pub fn get_type(&self) -> BlobType<'_> {
match &self.header.blob_type {
BlobKind::OsmHeader => BlobType::OsmHeader,
BlobKind::OsmData => BlobType::OsmData,
BlobKind::Unknown(s) => BlobType::Unknown(s),
}
}
#[inline]
pub fn offset(&self) -> Option<ByteOffset> {
self.offset
}
pub fn to_headerblock(&self) -> Result<HeaderBlock> {
decode_headerblock(&self.blob, None).map(HeaderBlock::new)
}
pub fn to_primitiveblock(&self) -> Result<PrimitiveBlock> {
decompress_blob(&self.blob, None).and_then(PrimitiveBlock::new)
}
#[hotpath::measure]
pub(crate) fn decompress_into(&self, buf: &mut Vec<u8>) -> Result<()> {
decompress_wire_blob_into(&self.blob, buf)
}
pub(crate) fn index(&self) -> Option<crate::blob_meta::BlobIndex> {
self.header
.indexdata
.as_ref()
.and_then(|d| crate::blob_meta::BlobIndex::deserialize(d))
}
pub(crate) fn compressed_data(&self) -> Option<(u8, &[u8])> {
match &self.blob.data {
Some(BlobData::Raw(b)) => Some((0, b)),
Some(BlobData::Zlib(b)) => Some((1, b)),
Some(BlobData::Zstd(b)) => Some((2, b)),
None => None,
}
}
pub(crate) fn tag_index(&self) -> Option<crate::blob_meta::TagIndex> {
self.header
.tagdata
.as_ref()
.and_then(|d| crate::blob_meta::TagIndex::deserialize(d))
}
pub(crate) fn to_primitiveblock_inline_with_scratch(
&self,
pool: &Arc<DecompressPool>,
st_scratch: &mut Vec<(u32, u32)>,
gr_scratch: &mut Vec<(u32, u32)>,
) -> Result<PrimitiveBlock> {
let mut buf = pool_get(Some(pool), self.blob.estimated_capacity());
decompress_parsed_blob_into(&self.blob, &mut buf)?;
PrimitiveBlock::from_vec_pooled_with_scratch(buf, pool, st_scratch, gr_scratch)
}
}
#[derive(Clone, Debug)]
pub struct BlobHeader {
header: WireBlobHeader,
}
impl BlobHeader {
fn new(header: WireBlobHeader) -> Self {
BlobHeader { header }
}
#[inline]
pub fn blob_type(&self) -> BlobType<'_> {
match &self.header.blob_type {
BlobKind::OsmHeader => BlobType::OsmHeader,
BlobKind::OsmData => BlobType::OsmData,
BlobKind::Unknown(s) => BlobType::Unknown(s),
}
}
#[inline]
pub fn get_blob_size(&self) -> i32 {
self.header.datasize
}
}
pub trait BlobReaderSource: Read + Seek {
fn skip_relative(&mut self, offset: i64) -> std::io::Result<()> {
self.seek(SeekFrom::Current(offset)).map(|_| ())
}
}
impl<R: Read + Seek> BlobReaderSource for BufReader<R> {
fn skip_relative(&mut self, offset: i64) -> std::io::Result<()> {
BufReader::seek_relative(self, offset)
}
}
impl BlobReaderSource for File {}
impl<T: AsRef<[u8]>> BlobReaderSource for Cursor<T> {}
#[derive(Clone, Debug)]
pub struct BlobReader<R: Read + Send> {
reader: R,
offset: Option<ByteOffset>,
last_blob_ok: bool,
header_buf: Vec<u8>,
parse_tagdata: bool,
parse_indexdata: bool,
#[cfg(target_os = "linux")]
evict_fd: Option<std::os::unix::io::RawFd>,
}
impl<R: Read + Send> BlobReader<R> {
pub fn new(reader: R) -> BlobReader<R> {
BlobReader {
reader,
offset: None,
last_blob_ok: true,
header_buf: Vec::new(),
parse_tagdata: false,
parse_indexdata: true,
#[cfg(target_os = "linux")]
evict_fd: None,
}
}
fn handle_error<T>(&mut self, error: crate::error::Error) -> Option<Result<T>> {
self.offset = None;
self.last_blob_ok = false;
Some(Err(error))
}
pub(crate) fn set_parse_tagdata(&mut self, enable: bool) {
self.parse_tagdata = enable;
}
pub(crate) fn set_parse_indexdata(&mut self, enable: bool) {
self.parse_indexdata = enable;
}
#[allow(clippy::cast_possible_truncation)]
fn read_blob_header(&mut self) -> Option<Result<WireBlobHeader>> {
let header_size: u64 = {
let mut buf = [0u8; 4];
match self.reader.read_exact(&mut buf[..1]) {
Ok(()) => {}
Err(e) if e.kind() == ::std::io::ErrorKind::UnexpectedEof => {
return None;
}
Err(e) => {
self.offset = None;
self.last_blob_ok = false;
return Some(Err(e.into()));
}
}
match self.reader.read_exact(&mut buf[1..]) {
Ok(()) => {
self.offset = self.offset.map(|x| ByteOffset(x.0 + 4));
u64::from(u32::from_be_bytes(buf))
}
Err(e) if e.kind() == ::std::io::ErrorKind::UnexpectedEof => {
return None;
}
Err(e) => {
self.offset = None;
self.last_blob_ok = false;
return Some(Err(e.into()));
}
}
};
if header_size >= MAX_BLOB_HEADER_SIZE {
self.last_blob_ok = false;
return Some(Err(new_blob_error(BlobError::HeaderTooBig {
size: header_size,
})));
}
let mut reader = self.reader.by_ref().take(header_size);
self.header_buf.clear();
self.header_buf.reserve(header_size as usize);
if let Err(e) = reader.read_to_end(&mut self.header_buf) {
return self.handle_error(e.into());
}
if self.header_buf.len() as u64 != header_size {
let header_start = self.offset.map_or(0, |x| x.0);
let got = self.header_buf.len() as u64;
let trunc_at = header_start + got;
return self.handle_error(new_error(ErrorKind::Io(
::std::io::Error::new(
::std::io::ErrorKind::UnexpectedEof,
format!(
"BlobHeader truncated at byte {trunc_at} (shape 3): \
declared {header_size} bytes from offset \
{header_start}, got {got}"
),
),
)));
}
let header = match WireBlobHeader::parse(&self.header_buf, self.parse_tagdata, self.parse_indexdata) {
Ok(header) => header,
Err(e) => {
return self.handle_error(e)
}
};
if header.datasize < 0 {
return self.handle_error(new_blob_error(BlobError::InvalidDataSize {
size: header.datasize,
}));
}
self.offset = self.offset.map(|x| ByteOffset(x.0 + header_size));
Some(Ok(header))
}
}
impl BlobReader<FileReader> {
pub fn from_path<P: AsRef<Path>>(path: P) -> Result<Self> {
let reader = FileReader::buffered(path.as_ref())?;
#[cfg(target_os = "linux")]
let evict_fd = Some({
use std::os::unix::io::AsRawFd;
match &reader {
FileReader::Buffered(r) => r.get_ref().as_raw_fd(),
#[cfg(feature = "linux-direct-io")]
FileReader::Direct(r) => r.raw_fd(),
}
});
Ok(BlobReader {
reader,
offset: Some(ByteOffset(0)),
last_blob_ok: true,
header_buf: Vec::new(),
parse_tagdata: false,
parse_indexdata: true,
#[cfg(target_os = "linux")]
evict_fd,
})
}
#[cfg(feature = "linux-direct-io")]
pub fn from_path_direct<P: AsRef<Path>>(path: P) -> Result<Self> {
let reader = FileReader::direct(path.as_ref())?;
Ok(BlobReader {
reader,
offset: Some(ByteOffset(0)),
last_blob_ok: true,
header_buf: Vec::new(),
parse_tagdata: false,
parse_indexdata: true,
evict_fd: None, })
}
pub fn open<P: AsRef<Path>>(path: P, direct: bool) -> Result<Self> {
let reader = FileReader::open(path.as_ref(), direct)?;
#[cfg(target_os = "linux")]
let evict_fd = if direct { None } else {
Some({
use std::os::unix::io::AsRawFd;
match &reader {
FileReader::Buffered(r) => r.get_ref().as_raw_fd(),
#[cfg(feature = "linux-direct-io")]
FileReader::Direct(r) => r.raw_fd(),
}
})
};
Ok(BlobReader {
reader,
offset: Some(ByteOffset(0)),
last_blob_ok: true,
header_buf: Vec::new(),
parse_tagdata: false,
parse_indexdata: true,
#[cfg(target_os = "linux")]
evict_fd,
})
}
}
impl<R: Read + Send> Iterator for BlobReader<R> {
type Item = Result<Blob>;
#[allow(clippy::cast_sign_loss, clippy::cast_possible_truncation)]
fn next(&mut self) -> Option<Self::Item> {
if !self.last_blob_ok {
return None;
}
let prev_offset = self.offset;
let header = match self.read_blob_header() {
Some(Ok(header)) => header,
Some(Err(err)) => return Some(Err(err)),
None => return None,
};
let mut reader = self.reader.by_ref().take(header.datasize as u64);
let mut blob_data = Vec::with_capacity(header.datasize as usize);
if let Err(e) = reader.read_to_end(&mut blob_data) {
return self.handle_error(e.into());
}
if blob_data.len() as u64 != header.datasize as u64 {
let payload_start = self.offset.map_or(0, |x| x.0);
let got = blob_data.len() as u64;
let trunc_at = payload_start + got;
return self.handle_error(new_error(ErrorKind::Io(
::std::io::Error::new(
::std::io::ErrorKind::UnexpectedEof,
format!(
"Blob payload truncated at byte {trunc_at} (shape 4): \
declared {} bytes from offset {payload_start}, got {got}",
header.datasize
),
),
)));
}
let blob_bytes = Bytes::from(blob_data);
let blob = match WireBlob::parse(&blob_bytes) {
Ok(blob) => blob,
Err(e) => return self.handle_error(e),
};
self.offset = self
.offset
.map(|x| ByteOffset(x.0 + header.datasize as u64));
#[cfg(target_os = "linux")]
if let Some(fd) = self.evict_fd
&& let Some(offset) = self.offset
{
unsafe {
libc::posix_fadvise(
fd,
0,
offset.0.try_into().unwrap_or(i64::MAX),
libc::POSIX_FADV_DONTNEED,
)
};
}
Some(Ok(Blob::new(header, blob, prev_offset)))
}
}
impl<R: BlobReaderSource + Send> BlobReader<R> {
pub fn new_seekable(mut reader: R) -> Result<BlobReader<R>> {
let pos = reader.stream_position()?;
Ok(BlobReader {
reader,
offset: Some(ByteOffset(pos)),
last_blob_ok: true,
header_buf: Vec::new(),
parse_tagdata: false,
parse_indexdata: true,
#[cfg(target_os = "linux")]
evict_fd: None,
})
}
pub fn blob_from_offset(&mut self, pos: ByteOffset) -> Result<Blob> {
self.seek(pos)?;
self.next().unwrap_or_else(|| {
Err(new_error(ErrorKind::Io(::std::io::Error::new(
::std::io::ErrorKind::UnexpectedEof,
"no blob at this stream position",
))))
})
}
pub fn seek(&mut self, pos: ByteOffset) -> Result<()> {
match self.reader.seek(SeekFrom::Start(pos.0)) {
Ok(offset) => {
self.offset = Some(ByteOffset(offset));
Ok(())
}
Err(e) => {
self.offset = None;
Err(e.into())
}
}
}
pub fn seek_raw(&mut self, pos: SeekFrom) -> Result<u64> {
match self.reader.seek(pos) {
Ok(offset) => {
self.offset = Some(ByteOffset(offset));
self.last_blob_ok = true;
Ok(offset)
}
Err(e) => {
self.offset = None;
Err(e.into())
}
}
}
fn skip_blob_body(&mut self, n: u64) -> Result<()> {
if n == 0 {
return Ok(());
}
#[allow(clippy::cast_possible_wrap)]
let signed = (n - 1) as i64;
if let Err(e) = self.reader.skip_relative(signed) {
self.offset = None;
return Err(e.into());
}
let mut sentinel = [0u8; 1];
match self.reader.read_exact(&mut sentinel) {
Ok(()) => {
self.offset = self.offset.map(|x| ByteOffset(x.0 + n));
Ok(())
}
Err(e) if e.kind() == ::std::io::ErrorKind::UnexpectedEof => {
let payload_start = self.offset.map_or(0, |x| x.0);
let trunc_at = payload_start + n - 1;
self.offset = None;
Err(new_error(ErrorKind::Io(::std::io::Error::new(
::std::io::ErrorKind::UnexpectedEof,
format!(
"Blob payload truncated at byte {trunc_at} (shape 4): \
declared {n} bytes from offset {payload_start}, \
file ended early"
),
))))
}
Err(e) => {
self.offset = None;
Err(e.into())
}
}
}
#[allow(clippy::cast_sign_loss)]
#[hotpath::measure]
pub fn next_header_skip_blob(&mut self) -> Option<Result<(BlobHeader, Option<ByteOffset>)>> {
if !self.last_blob_ok {
return None;
}
let prev_offset = self.offset;
let header = match self.read_blob_header() {
Some(Ok(header)) => header,
Some(Err(err)) => return Some(Err(err)),
None => return None,
};
#[allow(clippy::cast_sign_loss)]
if let Err(err) = self.skip_blob_body(header.datasize as u64) {
self.last_blob_ok = false;
return Some(Err(err));
}
Some(Ok((BlobHeader::new(header), prev_offset)))
}
}
impl BlobReader<BufReader<File>> {
pub fn seekable_from_path<P: AsRef<Path>>(path: P) -> Result<BlobReader<BufReader<File>>> {
let f = File::open(path.as_ref())?;
let buf_reader = BufReader::with_capacity(256 * 1024, f);
Self::new_seekable(buf_reader)
}
}
pub(crate) fn decode_blob_to_primitiveblock(blob_bytes: &[u8]) -> Result<crate::PrimitiveBlock> {
let blob = WireBlob::parse_slice(blob_bytes)?;
decompress_blob(&blob, None).and_then(crate::PrimitiveBlock::new)
}
pub(crate) fn parse_primitive_block_from_bytes_owned(raw: &Bytes) -> Result<crate::PrimitiveBlock> {
crate::PrimitiveBlock::new(raw.clone())
}
pub(crate) fn decode_blob_to_headerblock(blob_bytes: &[u8]) -> Result<crate::HeaderBlock> {
decode_blob_to_headerblock_from_bytes(&Bytes::copy_from_slice(blob_bytes))
}
pub(crate) fn decode_blob_to_headerblock_from_bytes(blob_bytes: &Bytes) -> Result<crate::HeaderBlock> {
let blob = WireBlob::parse(blob_bytes)?;
let raw = decompress_blob(&blob, None)?;
crate::HeaderBlock::parse_from_bytes(&raw)
}
pub(crate) fn decode_headerblock(
blob: &WireBlob,
pool: Option<&Arc<DecompressPool>>,
) -> Result<super::block::WireHeaderBlock> {
let raw = decompress_blob(blob, pool)?;
super::block::WireHeaderBlock::parse(&raw)
}
#[cfg(test)]
#[allow(clippy::unwrap_used)]
mod tests {
use super::*;
#[test]
fn test_get_type() {
let pairs: &[(BlobKind, BlobType<'_>)] = &[
(BlobKind::Unknown(String::new()), BlobType::Unknown("")),
(BlobKind::Unknown("abc".to_string()), BlobType::Unknown("abc")),
(BlobKind::OsmHeader, BlobType::OsmHeader),
(BlobKind::OsmData, BlobType::OsmData),
];
for (kind, expected_type) in pairs {
let ff_header = WireBlobHeader {
blob_type: kind.clone(),
datasize: 0,
indexdata: None,
tagdata: None,
};
let ff_blob = WireBlob { data: None, raw_size: None };
let blob = Blob::new(ff_header, ff_blob, None);
assert_eq!(blob.get_type(), *expected_type);
}
}
}