use std::os::unix::fs::FileExt as _;
use std::path::Path;
use std::sync::Mutex;
use std::sync::atomic::{AtomicUsize, Ordering};
use super::Result;
#[hotpath::measure]
pub(super) fn run_pass1_5(
way_schedule: &[(usize, u64, usize)],
max_node_id: i64,
shared_file: &std::sync::Arc<std::fs::File>,
needed_admin_ways: &crate::idset::IdSet,
) -> Result<crate::idset::IdSet> {
let mut referenced_nodes = crate::idset::IdSet::new();
referenced_nodes.pre_allocate(max_node_id);
crate::debug::emit_marker("GEOCODE_PASS1_5_SCAN_START");
{
let referenced_ref = &referenced_nodes;
let literals = crate::scan::way::GeocodeTagLiterals::standard();
let literals_ref = &literals;
let needed_admin_ways_ref = needed_admin_ways;
let next_idx = AtomicUsize::new(0);
let next_ref = &next_idx;
let first_err: Mutex<Option<String>> = Mutex::new(None);
let first_err_ref = &first_err;
let decode_threads = std::thread::available_parallelism()
.map(|n| n.get().saturating_sub(2).max(1))
.unwrap_or(4);
std::thread::scope(|scope| {
for _ in 0..decode_threads {
let file = std::sync::Arc::clone(shared_file);
scope.spawn(move || {
let mut read_buf: Vec<u8> = Vec::new();
let mut decompress_buf: Vec<u8> = Vec::new();
let mut refs_buf: Vec<i64> = Vec::new();
let mut group_starts: Vec<(usize, usize)> = Vec::new();
loop {
if first_err_ref
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.is_some()
{
return;
}
let idx = next_ref.fetch_add(1, Ordering::Relaxed);
if idx >= way_schedule.len() {
break;
}
let (_seq, offset, size) = way_schedule[idx];
let result = (|| -> std::result::Result<(), String> {
read_buf.resize(size, 0);
file.read_exact_at(&mut read_buf, offset)
.map_err(|e| format!("pread at {offset}: {e}"))?;
decompress_buf.clear();
crate::blob::decompress_blob_raw(&read_buf, &mut decompress_buf)
.map_err(|e| e.to_string())?;
crate::scan::way::scan_way_geocode_tagged_refs(
&decompress_buf,
literals_ref,
&mut refs_buf,
&mut group_starts,
|way_id, flags, refs| {
let is_admin = needed_admin_ways_ref.get(way_id);
if flags.is_street
|| flags.is_building_addr
|| flags.is_interp
|| is_admin
{
for &r in refs {
if r < 0 {
continue;
}
referenced_ref.set_atomic(r);
}
}
},
)
.map_err(|e| e.to_string())?;
Ok(())
})();
if let Err(e) = result {
let mut slot = first_err_ref
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if slot.is_none() {
*slot = Some(e);
}
return;
}
}
});
}
});
if let Some(e) = first_err
.into_inner()
.unwrap_or_else(std::sync::PoisonError::into_inner)
{
return Err(e.into());
}
}
crate::debug::emit_marker("GEOCODE_PASS1_5_SCAN_END");
Ok(referenced_nodes)
}
#[allow(clippy::type_complexity)]
pub(super) fn build_pass2_schedules(
input_path: &Path,
) -> Result<(
Vec<(usize, u64, usize)>, // node_schedule
Vec<(usize, u64, usize)>, // way_schedule
i64, // max_node_id
std::sync::Arc<std::fs::File>,
)> {
let mut walker = crate::read::header_walker::HeaderWalker::open(input_path)?;
let _ = walker
.next_header()?
.ok_or_else(|| crate::error::new_error(crate::error::ErrorKind::MissingHeader))?;
let mut node_schedule: Vec<(usize, u64, usize)> = Vec::new();
let mut way_schedule: Vec<(usize, u64, usize)> = Vec::new();
let mut node_seq: usize = 0;
let mut way_seq: usize = 0;
let mut max_node_id: i64 = 0;
while let Some(meta) = walker.next_header()? {
if !matches!(meta.blob_type, crate::blob::BlobKind::OsmData) {
continue;
}
let Some(idx) = meta.index.as_ref() else {
node_schedule.push((node_seq, meta.data_offset, meta.data_size));
node_seq += 1;
way_schedule.push((way_seq, meta.data_offset, meta.data_size));
way_seq += 1;
continue;
};
match idx.kind {
crate::blob_meta::ElemKind::Node => {
if idx.max_id > max_node_id {
max_node_id = idx.max_id;
}
node_schedule.push((node_seq, meta.data_offset, meta.data_size));
node_seq += 1;
}
crate::blob_meta::ElemKind::Way => {
way_schedule.push((way_seq, meta.data_offset, meta.data_size));
way_seq += 1;
}
crate::blob_meta::ElemKind::Relation => {}
}
}
let shared_file = std::sync::Arc::clone(walker.shared_file());
drop(walker);
#[allow(clippy::cast_possible_wrap)]
{
crate::debug::emit_counter("pass2_node_blobs", node_schedule.len() as i64);
crate::debug::emit_counter("pass2_way_blobs", way_schedule.len() as i64);
crate::debug::emit_counter("pass2_max_node_id", max_node_id);
}
Ok((node_schedule, way_schedule, max_node_id, shared_file))
}