use std::io::Read;
use std::path::Path;
use std::sync::Arc;
use crate::blob::{BlobKind, MAX_BLOB_HEADER_SIZE, parse_blob_header_with_index};
use crate::blob_meta::BlobIndex;
use crate::error::Result;
const HEADER_PROBE_SIZE: usize = 4096;
const SAMPLE_CAP: usize = 1_000;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) struct BlobCountEstimate {
pub(crate) osmdata_blobs: u64,
pub(crate) exact: bool,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum ScanArm {
Walker,
FullScan,
}
pub(crate) const FULL_SCAN_ARM_MIN_BLOBS: u64 = 150_000;
pub(crate) fn estimate_blob_count(path: &Path) -> Result<BlobCountEstimate> {
let mut walker = HeaderWalker::open(path)?;
let file_size = walker.file_size();
if file_size == 0 {
return Err(crate::error::new_error(
crate::error::ErrorKind::MissingHeader,
));
}
let mut frames = 0_usize;
let mut osmdata_blobs = 0_u64;
let mut sampled_osmdata_bytes = 0_u64;
let mut sampled_end = 0_u64;
while frames < SAMPLE_CAP {
let Some(meta) = walker.next_header()? else {
return Ok(BlobCountEstimate {
osmdata_blobs,
exact: true,
});
};
frames += 1;
if frames == 1 && meta.blob_type != BlobKind::OsmHeader {
return Err(crate::error::new_error(
crate::error::ErrorKind::MissingHeader,
));
}
if meta.blob_type == BlobKind::OsmData {
osmdata_blobs += 1;
sampled_osmdata_bytes += meta.frame_size as u64;
sampled_end = meta.frame_start + meta.frame_size as u64;
}
}
if osmdata_blobs == 0 {
return Ok(BlobCountEstimate {
osmdata_blobs: 0,
exact: false,
});
}
let mean_frame_bytes = sampled_osmdata_bytes / osmdata_blobs;
if mean_frame_bytes == 0 {
return Ok(BlobCountEstimate {
osmdata_blobs,
exact: false,
});
}
let remaining_bytes = file_size.saturating_sub(sampled_end);
Ok(BlobCountEstimate {
osmdata_blobs: osmdata_blobs + remaining_bytes / mean_frame_bytes,
exact: false,
})
}
pub(crate) fn choose_scan_arm_at(estimate: &BlobCountEstimate, min_blobs: u64) -> ScanArm {
if estimate.osmdata_blobs >= min_blobs {
ScanArm::FullScan
} else {
ScanArm::Walker
}
}
pub(crate) struct BlobHeaderMeta {
pub blob_type: BlobKind,
pub frame_start: u64,
pub data_offset: u64,
pub data_size: usize,
pub index: Option<BlobIndex>,
pub tagdata: Option<Box<[u8]>>,
pub frame_size: usize,
}
pub(crate) struct HeaderWalker {
file: Arc<std::fs::File>,
offset: u64,
file_size: u64,
header_buf: Vec<u8>,
}
impl HeaderWalker {
pub(crate) fn open(path: &Path) -> Result<Self> {
let file = std::fs::File::open(path).map_err(|e| {
crate::error::new_error(crate::error::ErrorKind::Io(std::io::Error::other(format!(
"failed to open {}: {e}",
path.display()
))))
})?;
let file_size = file
.metadata()
.map_err(|e| crate::error::new_error(crate::error::ErrorKind::Io(e)))?
.len();
#[cfg(target_os = "linux")]
{
use std::os::unix::io::AsRawFd;
unsafe {
libc::posix_fadvise(file.as_raw_fd(), 0, 0, libc::POSIX_FADV_RANDOM);
}
}
Ok(Self {
file: Arc::new(file),
offset: 0,
file_size,
header_buf: Vec::new(),
})
}
pub(crate) fn shared_file(&self) -> &Arc<std::fs::File> {
&self.file
}
pub(crate) fn file_size(&self) -> u64 {
self.file_size
}
pub(crate) fn next_header(&mut self) -> Result<Option<BlobHeaderMeta>> {
use std::os::unix::fs::FileExt as _;
if self.offset >= self.file_size {
return Ok(None);
}
let frame_start = self.offset;
let remaining = self.file_size - self.offset;
let probe_len = usize::try_from(remaining)
.unwrap_or(usize::MAX)
.min(HEADER_PROBE_SIZE);
if probe_len < 4 {
return Ok(None);
}
self.header_buf.resize(probe_len, 0);
self.file
.read_exact_at(&mut self.header_buf, self.offset)
.map_err(|e| crate::error::new_error(crate::error::ErrorKind::Io(e)))?;
let header_len = u32::from_be_bytes([
self.header_buf[0],
self.header_buf[1],
self.header_buf[2],
self.header_buf[3],
]) as usize;
if header_len as u64 >= MAX_BLOB_HEADER_SIZE {
return Err(crate::error::new_blob_error(
crate::error::BlobError::HeaderTooBig {
size: header_len as u64,
},
));
}
let header_end = 4 + header_len;
if header_end > probe_len {
self.header_buf.resize(header_end, 0);
let tail_offset = self.offset + probe_len as u64;
self.file
.read_exact_at(&mut self.header_buf[probe_len..header_end], tail_offset)
.map_err(|e| crate::error::new_error(crate::error::ErrorKind::Io(e)))?;
}
let (blob_type, data_size, raw_index, tagdata) =
parse_blob_header_with_index(&self.header_buf[4..header_end])?;
let index = raw_index.as_ref().and_then(|b| BlobIndex::deserialize(b));
let data_offset = self.offset + header_end as u64;
let payload_end = data_offset.checked_add(data_size as u64).ok_or_else(|| {
crate::error::new_error(crate::error::ErrorKind::Io(::std::io::Error::new(
::std::io::ErrorKind::InvalidData,
format!(
"blob at offset {} declares overflowing payload size {data_size}",
self.offset
),
)))
})?;
if payload_end > self.file_size {
return Err(crate::error::new_error(crate::error::ErrorKind::Io(
::std::io::Error::new(
::std::io::ErrorKind::UnexpectedEof,
format!(
"blob payload truncated: declared {data_size} bytes \
from offset {data_offset}, file_size {}",
self.file_size
),
),
)));
}
self.offset = payload_end;
let frame_size = 4 + header_len + data_size;
Ok(Some(BlobHeaderMeta {
blob_type,
frame_start,
data_offset,
data_size,
index,
tagdata,
frame_size,
}))
}
pub(crate) fn pread_data(&self, offset: u64, size: usize, buf: &mut Vec<u8>) -> Result<()> {
use std::os::unix::fs::FileExt as _;
buf.resize(size, 0);
self.file
.read_exact_at(buf, offset)
.map_err(|e| crate::error::new_error(crate::error::ErrorKind::Io(e)))
}
}
#[allow(dead_code)]
pub(crate) fn pread_exact(
file: &std::fs::File,
offset: u64,
size: usize,
buf: &mut Vec<u8>,
) -> Result<()> {
use std::os::unix::fs::FileExt as _;
buf.resize(size, 0);
file.read_exact_at(buf, offset)
.map_err(|e| crate::error::new_error(crate::error::ErrorKind::Io(e)))
}
#[allow(dead_code)]
pub(crate) fn read_blob_data<R: Read>(reader: &mut R, size: usize) -> Result<Vec<u8>> {
let mut buf = vec![0u8; size];
reader
.read_exact(&mut buf)
.map_err(|e| crate::error::new_error(crate::error::ErrorKind::Io(e)))?;
Ok(buf)
}
#[cfg(test)]
mod tests {
use super::{
BlobCountEstimate, FULL_SCAN_ARM_MIN_BLOBS, SAMPLE_CAP, ScanArm, choose_scan_arm_at,
estimate_blob_count,
};
use crate::block_builder::{BlockBuilder, HeaderBuilder};
use crate::writer::{Compression, PbfWriter};
fn write_fixture(path: &std::path::Path, with_data: bool) {
let file = std::fs::File::create(path).expect("create fixture");
let mut writer = PbfWriter::new(std::io::BufWriter::new(file), Compression::default());
writer
.write_header(&HeaderBuilder::new().build().expect("header"))
.expect("write header");
if with_data {
let mut block = BlockBuilder::new();
block.add_node(1, 0, 0, std::iter::empty::<(&str, &str)>(), None);
writer
.write_primitive_block(block.take().expect("take").expect("block"))
.expect("write block");
}
writer.flush().expect("flush fixture");
}
#[test]
fn chooser_switches_at_the_policy_boundary() {
assert_eq!(
choose_scan_arm_at(
&BlobCountEstimate {
osmdata_blobs: FULL_SCAN_ARM_MIN_BLOBS - 1,
exact: true,
},
FULL_SCAN_ARM_MIN_BLOBS
),
ScanArm::Walker
);
assert_eq!(
choose_scan_arm_at(
&BlobCountEstimate {
osmdata_blobs: FULL_SCAN_ARM_MIN_BLOBS,
exact: false,
},
FULL_SCAN_ARM_MIN_BLOBS
),
ScanArm::FullScan
);
}
#[test]
fn chooser_is_independent_of_estimate_exactness() {
let estimate = BlobCountEstimate {
osmdata_blobs: 1,
exact: false,
};
assert_eq!(choose_scan_arm_at(&estimate, 1), ScanArm::FullScan);
}
#[test]
fn estimator_is_exact_for_small_and_header_only_files() {
let dir = tempfile::tempdir().expect("tempdir");
let header_only = dir.path().join("header-only.pbf");
write_fixture(&header_only, false);
assert_eq!(
estimate_blob_count(&header_only).expect("estimate"),
BlobCountEstimate {
osmdata_blobs: 0,
exact: true,
}
);
let one_block = dir.path().join("one-block.pbf");
write_fixture(&one_block, true);
assert_eq!(
estimate_blob_count(&one_block).expect("estimate"),
BlobCountEstimate {
osmdata_blobs: 1,
exact: true,
}
);
}
#[test]
fn estimator_rejects_empty_input() {
let dir = tempfile::tempdir().expect("tempdir");
let empty = dir.path().join("empty.pbf");
std::fs::File::create(&empty).expect("create empty file");
assert!(estimate_blob_count(&empty).is_err());
}
#[test]
fn estimator_projects_within_tolerance_on_mixed_frame_sizes() {
let dir = tempfile::tempdir().expect("tempdir");
let mixed = dir.path().join("mixed.pbf");
let file = std::fs::File::create(&mixed).expect("create fixture");
let mut writer = PbfWriter::new(std::io::BufWriter::new(file), Compression::None);
writer
.write_header(&HeaderBuilder::new().build().expect("header"))
.expect("write header");
let actual_blobs = (SAMPLE_CAP + SAMPLE_CAP / 2) as u64;
let mut block = BlockBuilder::new();
for blob in 0..actual_blobs {
let nodes_in_blob = if blob % 2 == 0 { 1 } else { 40 };
for node in 0..nodes_in_blob {
#[allow(clippy::cast_possible_wrap)]
let id = (blob * 64 + node) as i64;
block.add_node(
id,
0,
0,
[("highway", "primary_link")].iter().copied(),
None,
);
}
writer
.write_primitive_block(block.take().expect("take").expect("block"))
.expect("write block");
}
writer.flush().expect("flush fixture");
let estimate = estimate_blob_count(&mixed).expect("estimate");
assert!(!estimate.exact);
let error = estimate.osmdata_blobs.abs_diff(actual_blobs);
assert!(
error * 100 < actual_blobs * 30,
"estimated {} of {actual_blobs} actual blobs; relative error over 30 %",
estimate.osmdata_blobs,
);
}
}