use crate::BoxResult as Result;
pub(crate) type ScheduleEntry = (usize, u64, usize);
const WILLNEED_COALESCE_GAP_BYTES: u64 = 4 * 1024;
const WILLNEED_WINDOW_BYTES: u64 = 256 * 1024 * 1024;
fn coalesce_willneed_ranges(schedule: &[ScheduleEntry]) -> Vec<(u64, u64)> {
let mut ranges: Vec<(u64, u64)> = Vec::new();
for &(_, offset, size) in schedule {
let size = size as u64;
if size == 0 {
continue;
}
let Some(end) = offset.checked_add(size) else {
continue;
};
if let Some((previous_offset, previous_size)) = ranges.last_mut()
&& previous_offset
.checked_add(*previous_size)
.and_then(|previous_end| previous_end.checked_add(WILLNEED_COALESCE_GAP_BYTES))
.is_some_and(|end_with_gap| offset <= end_with_gap)
&& let Some(previous_end) = previous_offset.checked_add(*previous_size)
&& let Some(extension) = end.checked_sub(previous_end)
&& let Some(extended_size) = previous_size.checked_add(extension)
{
*previous_size = extended_size;
} else {
ranges.push((offset, size));
}
}
ranges
}
pub(crate) struct WillneedPrefetch {
ranges: Vec<(u64, u64)>,
next_range: usize,
next_offset: u64,
}
impl WillneedPrefetch {
pub(crate) fn new(schedule: &[ScheduleEntry]) -> Self {
Self {
ranges: coalesce_willneed_ranges(schedule),
next_range: 0,
next_offset: 0,
}
}
pub(crate) fn before_dispatch(
&mut self,
shared_file: &std::sync::Arc<std::fs::File>,
entry: ScheduleEntry,
) {
let Some(window_end) = entry
.1
.checked_add(entry.2 as u64)
.and_then(|end| end.checked_add(WILLNEED_WINDOW_BYTES))
else {
return;
};
while let Some(&(range_offset, range_size)) = self.ranges.get(self.next_range) {
let Some(range_end) = range_offset.checked_add(range_size) else {
self.next_range += 1;
self.next_offset = 0;
continue;
};
let offset = if self.next_offset == 0 {
range_offset
} else {
self.next_offset
};
if offset >= window_end {
break;
}
let end = range_end.min(window_end);
let Some(size) = end.checked_sub(offset) else {
self.next_range += 1;
self.next_offset = 0;
continue;
};
if end == range_end {
self.next_range += 1;
self.next_offset = 0;
} else {
self.next_offset = end;
}
#[cfg(target_os = "linux")]
{
use std::os::unix::io::AsRawFd as _;
let (Ok(offset), Ok(size)) = (offset.try_into(), size.try_into()) else {
continue;
};
unsafe {
libc::posix_fadvise(
shared_file.as_ref().as_raw_fd(),
offset,
size,
libc::POSIX_FADV_WILLNEED,
);
}
}
#[cfg(not(target_os = "linux"))]
let _ = shared_file;
}
}
}
fn resolve_thread_count(threads: Option<usize>) -> usize {
match threads {
Some(n) if n > 0 => n,
_ => std::thread::available_parallelism()
.map(|n| n.get().saturating_sub(2).max(1))
.unwrap_or(4),
}
}
#[cfg_attr(feature = "hotpath", hotpath::measure)]
pub(crate) fn build_classify_schedule(
input: &std::path::Path,
kind_filter: Option<crate::blob_meta::ElemKind>,
) -> Result<(Vec<ScheduleEntry>, std::sync::Arc<std::fs::File>)> {
crate::debug::emit_marker("SCHEDULE_SCANNER_OPEN_START");
let mut walker = crate::read::header_walker::HeaderWalker::open(input)?;
let _ = walker
.next_header()?
.ok_or_else(|| crate::error::new_error(crate::error::ErrorKind::MissingHeader))?;
crate::debug::emit_marker("SCHEDULE_SCANNER_OPEN_END");
crate::debug::emit_marker("SCHEDULE_SCAN_LOOP_START");
let file_size = walker.file_size();
let mut schedule: Vec<ScheduleEntry> = Vec::new();
let mut seq: usize = 0;
while let Some(meta) = walker.next_header()? {
if !matches!(meta.blob_type, crate::blob::BlobKind::OsmData) {
continue;
}
if let Some(filter_kind) = kind_filter
&& let Some(idx) = &meta.index
&& idx.kind != filter_kind
{
continue;
}
if meta.data_offset + meta.data_size as u64 > file_size {
return Err(format!(
"blob at offset {} claims data_size {} but file is only {} bytes",
meta.data_offset, meta.data_size, file_size,
)
.into());
}
schedule.push((seq, meta.data_offset, meta.data_size));
seq += 1;
}
crate::debug::emit_marker("SCHEDULE_SCAN_LOOP_END");
crate::debug::emit_marker("SCHEDULE_SCANNER_DROP_START");
let shared_file = std::sync::Arc::clone(walker.shared_file());
drop(walker);
crate::debug::emit_marker("SCHEDULE_SCANNER_DROP_END");
#[allow(clippy::cast_possible_wrap)]
crate::debug::emit_counter("schedule_blobs", schedule.len() as i64);
Ok((schedule, shared_file))
}
#[cfg_attr(feature = "hotpath", hotpath::measure)]
#[allow(clippy::type_complexity)]
pub(crate) fn build_classify_schedules_split(
input: &std::path::Path,
) -> Result<(
Vec<ScheduleEntry>,
Vec<ScheduleEntry>,
Vec<ScheduleEntry>,
std::sync::Arc<std::fs::File>,
)> {
crate::debug::emit_marker("SCHEDULE_SCANNER_OPEN_START");
let mut walker = crate::read::header_walker::HeaderWalker::open(input)?;
let _ = walker
.next_header()?
.ok_or_else(|| crate::error::new_error(crate::error::ErrorKind::MissingHeader))?;
crate::debug::emit_marker("SCHEDULE_SCANNER_OPEN_END");
crate::debug::emit_marker("SCHEDULE_SCAN_LOOP_START");
let file_size = walker.file_size();
let mut nodes: Vec<ScheduleEntry> = Vec::new();
let mut ways: Vec<ScheduleEntry> = Vec::new();
let mut rels: Vec<ScheduleEntry> = Vec::new();
while let Some(meta) = walker.next_header()? {
if !matches!(meta.blob_type, crate::blob::BlobKind::OsmData) {
continue;
}
if meta.data_offset + meta.data_size as u64 > file_size {
return Err(format!(
"blob at offset {} claims data_size {} but file is only {} bytes",
meta.data_offset, meta.data_size, file_size,
)
.into());
}
match meta.index.as_ref().map(|i| i.kind) {
Some(crate::blob_meta::ElemKind::Node) => {
nodes.push((nodes.len(), meta.data_offset, meta.data_size));
}
Some(crate::blob_meta::ElemKind::Way) => {
ways.push((ways.len(), meta.data_offset, meta.data_size));
}
Some(crate::blob_meta::ElemKind::Relation) => {
rels.push((rels.len(), meta.data_offset, meta.data_size));
}
None => {
nodes.push((nodes.len(), meta.data_offset, meta.data_size));
ways.push((ways.len(), meta.data_offset, meta.data_size));
rels.push((rels.len(), meta.data_offset, meta.data_size));
}
}
}
crate::debug::emit_marker("SCHEDULE_SCAN_LOOP_END");
crate::debug::emit_marker("SCHEDULE_SCANNER_DROP_START");
let shared_file = std::sync::Arc::clone(walker.shared_file());
drop(walker);
crate::debug::emit_marker("SCHEDULE_SCANNER_DROP_END");
#[allow(clippy::cast_possible_wrap)]
{
crate::debug::emit_counter("schedule_node_blobs", nodes.len() as i64);
crate::debug::emit_counter("schedule_way_blobs", ways.len() as i64);
crate::debug::emit_counter("schedule_relation_blobs", rels.len() as i64);
}
Ok((nodes, ways, rels, shared_file))
}
#[cfg_attr(feature = "hotpath", hotpath::measure)]
pub(crate) fn parallel_classify_phase<S: Send, R: Send>(
shared_file: &std::sync::Arc<std::fs::File>,
schedule: &[ScheduleEntry],
threads: Option<usize>,
worker_init: impl Fn() -> S + Send + Sync,
classify: impl Fn(&crate::PrimitiveBlock, &mut S) -> R + Send + Sync,
mut merge: impl FnMut(usize, R),
) -> Result<()> {
use std::os::unix::fs::FileExt as _;
let mut willneed = WillneedPrefetch::new(schedule);
if schedule.is_empty() {
return Ok(());
}
let decode_threads = resolve_thread_count(threads);
let (desc_tx, desc_rx) = std::sync::mpsc::sync_channel::<ScheduleEntry>(16);
let desc_rx = std::sync::Arc::new(std::sync::Mutex::new(desc_rx));
let (result_tx, result_rx) =
std::sync::mpsc::sync_channel::<(usize, crate::error::Result<R>)>(32);
std::thread::scope(|scope| -> Result<()> {
scope.spawn(move || {
for &item in schedule {
willneed.before_dispatch(shared_file, item);
if desc_tx.send(item).is_err() {
break;
}
}
});
for _ in 0..decode_threads {
let rx = std::sync::Arc::clone(&desc_rx);
let tx = result_tx.clone();
let file = std::sync::Arc::clone(shared_file);
let classify_ref = &classify;
let worker_init_ref = &worker_init;
scope.spawn(move || {
let mut read_buf: Vec<u8> = Vec::new();
let worker_pool = crate::blob::DecompressPool::new();
let mut st_scratch: Vec<(u32, u32)> = Vec::new();
let mut gr_scratch: Vec<(u32, u32)> = Vec::new();
let mut state = worker_init_ref();
loop {
let (s, data_offset, data_size) = {
let guard = rx.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
match guard.recv() {
Ok(d) => d,
Err(_) => break,
}
};
let r: crate::error::Result<R> = (|| {
read_buf.resize(data_size, 0);
file.read_exact_at(&mut read_buf, data_offset)
.map_err(|e| crate::error::new_error(crate::error::ErrorKind::Io(e)))?;
let mut buf = crate::blob::pool_get_pub(&worker_pool, data_size * 4);
crate::blob::decompress_blob_raw(&read_buf, &mut buf)?;
let block = crate::block::PrimitiveBlock::from_vec_pooled_with_scratch(
buf,
&worker_pool,
&mut st_scratch,
&mut gr_scratch,
)?;
Ok(classify_ref(&block, &mut state))
})();
if tx.send((s, r)).is_err() {
break;
}
}
});
}
drop(desc_rx);
drop(result_tx);
for (seq, result) in result_rx {
merge(seq, result?);
}
Ok(())
})?;
Ok(())
}
#[cfg_attr(feature = "hotpath", hotpath::measure)]
pub(crate) fn parallel_scan_blobs_raw<S: Send, R: Send>(
shared_file: &std::sync::Arc<std::fs::File>,
schedule: &[ScheduleEntry],
threads: Option<usize>,
worker_init: impl Fn() -> S + Send + Sync,
classify: impl Fn(&[u8], &mut S) -> crate::error::Result<R> + Send + Sync,
mut merge: impl FnMut(usize, R),
) -> Result<()> {
use std::os::unix::fs::FileExt as _;
let mut willneed = WillneedPrefetch::new(schedule);
if schedule.is_empty() {
return Ok(());
}
let decode_threads = resolve_thread_count(threads);
let (desc_tx, desc_rx) = std::sync::mpsc::sync_channel::<ScheduleEntry>(16);
let desc_rx = std::sync::Arc::new(std::sync::Mutex::new(desc_rx));
let (result_tx, result_rx) =
std::sync::mpsc::sync_channel::<(usize, crate::error::Result<R>)>(32);
std::thread::scope(|scope| -> Result<()> {
scope.spawn(move || {
for &item in schedule {
willneed.before_dispatch(shared_file, item);
if desc_tx.send(item).is_err() {
break;
}
}
});
for _ in 0..decode_threads {
let rx = std::sync::Arc::clone(&desc_rx);
let tx = result_tx.clone();
let file = std::sync::Arc::clone(shared_file);
let classify_ref = &classify;
let worker_init_ref = &worker_init;
scope.spawn(move || {
let mut read_buf: Vec<u8> = Vec::new();
let mut decompress_buf: Vec<u8> = Vec::new();
let mut state = worker_init_ref();
loop {
let (s, data_offset, data_size) = {
let guard = rx.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
match guard.recv() {
Ok(d) => d,
Err(_) => break,
}
};
let r: crate::error::Result<R> = (|| {
read_buf.resize(data_size, 0);
file.read_exact_at(&mut read_buf, data_offset)
.map_err(|e| crate::error::new_error(crate::error::ErrorKind::Io(e)))?;
crate::blob::decompress_blob_raw(&read_buf, &mut decompress_buf)?;
classify_ref(&decompress_buf, &mut state)
})();
if tx.send((s, r)).is_err() {
break;
}
}
});
}
drop(desc_rx);
drop(result_tx);
for (seq, result) in result_rx {
merge(seq, result?);
}
Ok(())
})?;
Ok(())
}
#[cfg_attr(feature = "hotpath", hotpath::measure)]
pub(crate) fn parallel_classify_accumulate<S: Send>(
shared_file: &std::sync::Arc<std::fs::File>,
schedule: &[ScheduleEntry],
threads: Option<usize>,
worker_init: impl Fn() -> S + Send + Sync,
classify: impl Fn(&crate::PrimitiveBlock, &mut S) + Send + Sync,
mut merge: impl FnMut(S),
) -> Result<()> {
use std::os::unix::fs::FileExt as _;
let mut willneed = WillneedPrefetch::new(schedule);
if schedule.is_empty() {
return Ok(());
}
let decode_threads = resolve_thread_count(threads);
let (desc_tx, desc_rx) = std::sync::mpsc::sync_channel::<ScheduleEntry>(16);
let desc_rx = std::sync::Arc::new(std::sync::Mutex::new(desc_rx));
let (result_tx, result_rx) =
std::sync::mpsc::sync_channel::<crate::error::Result<S>>(decode_threads);
std::thread::scope(|scope| -> Result<()> {
scope.spawn(move || {
for &item in schedule {
willneed.before_dispatch(shared_file, item);
if desc_tx.send(item).is_err() {
break;
}
}
});
for _ in 0..decode_threads {
let rx = std::sync::Arc::clone(&desc_rx);
let tx = result_tx.clone();
let file = std::sync::Arc::clone(shared_file);
let classify_ref = &classify;
let worker_init_ref = &worker_init;
scope.spawn(move || {
let mut read_buf: Vec<u8> = Vec::new();
let worker_pool = crate::blob::DecompressPool::new();
let mut st_scratch: Vec<(u32, u32)> = Vec::new();
let mut gr_scratch: Vec<(u32, u32)> = Vec::new();
let mut state = worker_init_ref();
let result: crate::error::Result<()> = (|| {
loop {
let (_s, data_offset, data_size) = {
let guard =
rx.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
match guard.recv() {
Ok(d) => d,
Err(_) => return Ok(()),
}
};
read_buf.resize(data_size, 0);
file.read_exact_at(&mut read_buf, data_offset)
.map_err(|e| crate::error::new_error(crate::error::ErrorKind::Io(e)))?;
let mut buf = crate::blob::pool_get_pub(&worker_pool, data_size * 4);
crate::blob::decompress_blob_raw(&read_buf, &mut buf)?;
let block = crate::block::PrimitiveBlock::from_vec_pooled_with_scratch(
buf,
&worker_pool,
&mut st_scratch,
&mut gr_scratch,
)?;
classify_ref(&block, &mut state);
}
})();
match result {
Ok(()) => {
tx.send(Ok(state)).ok();
}
Err(e) => {
tx.send(Err(e)).ok();
}
}
});
}
drop(desc_rx);
drop(result_tx);
for result in result_rx {
merge(result?);
}
Ok(())
})?;
Ok(())
}
#[cfg(test)]
#[allow(clippy::unwrap_used)]
mod tests {
use super::{ScheduleEntry, WillneedPrefetch, coalesce_willneed_ranges};
#[test]
fn willneed_prefetch_defaults_to_coalesced_schedule() {
let schedule: [ScheduleEntry; 2] = [(0, 100, 20), (1, 144, 5)];
let prefetch = WillneedPrefetch::new(&schedule);
assert_eq!(prefetch.ranges, vec![(100, 49)]);
assert_eq!(prefetch.next_range, 0);
assert_eq!(prefetch.next_offset, 0);
}
#[test]
fn willneed_ranges_coalesce_abutting_body_ranges() {
let schedule: [ScheduleEntry; 3] = [(0, 100, 20), (1, 120, 30), (2, 5_000, 5)];
assert_eq!(
coalesce_willneed_ranges(&schedule),
vec![(100, 50), (5_000, 5)]
);
}
#[test]
fn willneed_ranges_coalesce_realistic_blob_header_gaps() {
let schedule: [ScheduleEntry; 4] = [
(0, 100, 20),
(1, 144, 5),
(2, 4_245, 10),
(3, 8_352, 10),
];
assert_eq!(
coalesce_willneed_ranges(&schedule),
vec![(100, 4_155), (8_352, 10)]
);
}
}