use std::path::Path;
use std::sync::Arc;
use super::getid::ElementIds;
use super::{
HeaderOverrides, ensure_node_capacity_local, ensure_relation_capacity_local,
ensure_way_capacity_local, flush_local, writer_from_header_bytes,
};
use crate::blob::{BlobKind, decode_blob_to_headerblock};
use crate::blob_meta::ElemKind;
use crate::block_builder::{BlockBuilder, MemberData, OwnedBlock};
use crate::owned::{dense_node_metadata, element_metadata};
use crate::read::header_walker::{FULL_SCAN_ARM_MIN_BLOBS, HeaderWalker, ScanArm};
use crate::reorder_buffer::ReorderBuffer;
use crate::writer::Compression;
use crate::{BlobFilter, Element, ElementReader, MemberId, PrimitiveBlock};
use super::Result;
pub struct GetparentsOptions {
pub add_self: bool,
}
pub struct GetparentsStats {
pub nodes_written: u64,
pub ways_written: u64,
pub relations_written: u64,
}
impl GetparentsStats {
pub fn print_summary(&self) {
let total = self.nodes_written + self.ways_written + self.relations_written;
eprintln!(
"Wrote {total} elements: {} nodes, {} ways, {} relations",
self.nodes_written, self.ways_written, self.relations_written,
);
}
}
#[hotpath::measure]
pub fn getparents(
input: &Path,
output: &Path,
ids: &ElementIds,
opts: &GetparentsOptions,
compression: Compression,
direct_io: bool,
overrides: &HeaderOverrides,
) -> Result<GetparentsStats> {
let (_, stats) = getparents_dispatched(
input,
output,
ids,
opts,
compression,
direct_io,
overrides,
FULL_SCAN_ARM_MIN_BLOBS,
)?;
Ok(stats)
}
#[doc(hidden)]
#[allow(clippy::too_many_arguments)]
pub fn getparents_with_min_blobs(
input: &Path,
output: &Path,
ids: &ElementIds,
opts: &GetparentsOptions,
compression: Compression,
direct_io: bool,
overrides: &HeaderOverrides,
min_blobs: u64,
) -> Result<GetparentsStats> {
getparents_dispatched(
input,
output,
ids,
opts,
compression,
direct_io,
overrides,
min_blobs,
)
.map(|(_, stats)| stats)
}
#[allow(clippy::too_many_arguments)]
fn getparents_dispatched(
input: &Path,
output: &Path,
ids: &ElementIds,
opts: &GetparentsOptions,
compression: Compression,
direct_io: bool,
overrides: &HeaderOverrides,
min_blobs: u64,
) -> Result<(ScanArm, GetparentsStats)> {
let arm = super::dispatch_scan_arm(input, super::has_indexdata(input, direct_io)?, min_blobs)?;
let stats = getparents_with_arm(
input,
output,
ids,
opts,
compression,
direct_io,
overrides,
arm,
)?;
Ok((arm, stats))
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn getparents_with_arm(
input: &Path,
output: &Path,
ids: &ElementIds,
opts: &GetparentsOptions,
compression: Compression,
direct_io: bool,
overrides: &HeaderOverrides,
arm: ScanArm,
) -> Result<GetparentsStats> {
match arm {
ScanArm::Walker => {
getparents_walker(input, output, ids, opts, compression, direct_io, overrides)
}
ScanArm::FullScan => {
getparents_pipelined(input, output, ids, opts, compression, direct_io, overrides)
}
}
}
fn needed_blob_kinds(ids: &ElementIds, opts: &GetparentsOptions) -> (bool, bool, bool) {
(
opts.add_self && ids.node_ids.has_any(),
ids.node_ids.has_any() || (opts.add_self && ids.way_ids.has_any()),
ids.node_ids.has_any() || ids.way_ids.has_any() || ids.relation_ids.has_any(),
)
}
#[allow(clippy::too_many_arguments)]
#[allow(clippy::too_many_lines)]
fn getparents_walker(
input: &Path,
output: &Path,
ids: &ElementIds,
opts: &GetparentsOptions,
compression: Compression,
direct_io: bool,
overrides: &HeaderOverrides,
) -> Result<GetparentsStats> {
let (need_node_blobs, need_way_blobs, need_relation_blobs) = needed_blob_kinds(ids, opts);
crate::debug::emit_marker("GETPARENTS_SCHEDULE_START");
let mut walker = HeaderWalker::open(input)?;
let file_size = walker.file_size();
let mut header_buf: Vec<u8> = Vec::new();
let mut header_block: Option<crate::HeaderBlock> = None;
let mut schedule: Vec<(usize, u64, usize)> = Vec::new();
let mut blobs_skipped: u64 = 0;
while let Some(meta) = walker.next_header()? {
match meta.blob_type {
BlobKind::OsmHeader if header_block.is_none() => {
walker.pread_data(meta.data_offset, meta.data_size, &mut header_buf)?;
header_block = Some(decode_blob_to_headerblock(&header_buf)?);
}
BlobKind::OsmData => {
let keep = match meta.index.as_ref().map(|i| i.kind) {
Some(ElemKind::Node) => need_node_blobs,
Some(ElemKind::Way) => need_way_blobs,
Some(ElemKind::Relation) => need_relation_blobs,
None => true,
};
if keep {
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((schedule.len(), meta.data_offset, meta.data_size));
} else {
blobs_skipped += 1;
}
}
_ => {}
}
}
crate::debug::emit_counter(
"walk_actual_osmdata_blobs",
i64::try_from(schedule.len())
.unwrap_or(i64::MAX)
.saturating_add(i64::try_from(blobs_skipped).unwrap_or(i64::MAX)),
);
let shared_file = Arc::clone(walker.shared_file());
drop(walker);
let header = header_block
.ok_or_else(|| crate::error::new_error(crate::error::ErrorKind::MissingHeader))?;
#[allow(clippy::cast_possible_wrap)]
{
crate::debug::emit_counter("getparents_schedule_blobs", schedule.len() as i64);
crate::debug::emit_counter("getparents_blobs_skipped", blobs_skipped as i64);
}
crate::debug::emit_marker("GETPARENTS_SCHEDULE_END");
super::warn_locations_on_ways_loss(&header);
let header_bytes = super::build_output_header(&header, true, overrides, |hb| hb)?;
let mut writer =
writer_from_header_bytes(output, compression, &header_bytes, direct_io, false)?;
let mut stats = GetparentsStats {
nodes_written: 0,
ways_written: 0,
relations_written: 0,
};
crate::debug::emit_marker("GETPARENTS_DECODE_START");
type ClassifyResult = std::result::Result<(Vec<OwnedBlock>, (u64, u64, u64)), String>;
let mut reorder: ReorderBuffer<ClassifyResult> = ReorderBuffer::with_capacity(32);
let mut write_err: Option<Box<dyn std::error::Error + Send + Sync>> = None;
crate::scan::classify::parallel_classify_phase(
&shared_file,
&schedule,
None,
|| (),
|block, _state| -> ClassifyResult {
let mut bb = BlockBuilder::new();
let mut output: Vec<OwnedBlock> = Vec::new();
let counts = process_block(block, &mut bb, &mut output, ids, opts.add_self)?;
flush_local(&mut bb, &mut output)?;
Ok((output, counts))
},
|seq, result| {
if write_err.is_some() {
return;
}
reorder.push(seq, result);
while let Some(item) = reorder.pop_ready() {
match item {
Ok((blocks, (n, w, r))) => {
for OwnedBlock {
bytes: block_bytes,
index,
tagdata,
way_members,
} in blocks
{
if let Err(e) = writer.write_primitive_block_owned(
block_bytes,
index,
tagdata.as_deref(),
way_members.as_deref(),
) {
write_err = Some(Box::new(e));
return;
}
}
stats.nodes_written += n;
stats.ways_written += w;
stats.relations_written += r;
}
Err(e) => {
write_err = Some(e.into());
return;
}
}
}
},
)?;
if let Some(e) = write_err {
return Err(e);
}
writer.flush()?;
crate::debug::emit_marker("GETPARENTS_DECODE_END");
Ok(stats)
}
#[allow(clippy::too_many_arguments)]
fn getparents_pipelined(
input: &Path,
output: &Path,
ids: &ElementIds,
opts: &GetparentsOptions,
compression: Compression,
direct_io: bool,
overrides: &HeaderOverrides,
) -> Result<GetparentsStats> {
let (need_node_blobs, need_way_blobs, need_relation_blobs) = needed_blob_kinds(ids, opts);
let reader = ElementReader::open(input, direct_io)?.with_blob_filter(BlobFilter::new(
need_node_blobs,
need_way_blobs,
need_relation_blobs,
));
super::warn_locations_on_ways_loss(reader.header());
let header_bytes = super::build_output_header(reader.header(), true, overrides, |hb| hb)?;
let mut writer =
writer_from_header_bytes(output, compression, &header_bytes, direct_io, false)?;
let mut stats = GetparentsStats {
nodes_written: 0,
ways_written: 0,
relations_written: 0,
};
crate::debug::emit_marker("GETPARENTS_DECODE_START");
reader.for_each_fused_block(
|block| {
let mut bb = BlockBuilder::new();
let mut output = Vec::new();
let counts = process_block(&block, &mut bb, &mut output, ids, opts.add_self)?;
flush_local(&mut bb, &mut output)?;
Ok((output, counts))
},
|(blocks, (nodes, ways, relations))| {
for OwnedBlock {
bytes,
index,
tagdata,
way_members,
} in blocks
{
writer.write_primitive_block_owned(
bytes,
index,
tagdata.as_deref(),
way_members.as_deref(),
)?;
}
stats.nodes_written += nodes;
stats.ways_written += ways;
stats.relations_written += relations;
Ok(())
},
)?;
writer.flush()?;
crate::debug::emit_marker("GETPARENTS_DECODE_END");
Ok(stats)
}
fn process_block(
block: &PrimitiveBlock,
bb: &mut BlockBuilder,
output: &mut Vec<OwnedBlock>,
ids: &ElementIds,
add_self: bool,
) -> std::result::Result<(u64, u64, u64), String> {
let mut nodes: u64 = 0;
let mut ways: u64 = 0;
let mut relations: u64 = 0;
let mut refs_buf: Vec<i64> = Vec::new();
let mut members_buf: Vec<MemberData<'_>> = Vec::new();
for element in block.elements() {
match &element {
Element::DenseNode(dn) => {
if add_self && ids.node_ids.get(dn.id()) {
ensure_node_capacity_local(bb, output)?;
let meta = dense_node_metadata(dn);
bb.add_node(
dn.id(),
dn.decimicro_lat(),
dn.decimicro_lon(),
dn.tags(),
meta.as_ref(),
);
nodes += 1;
}
}
Element::Node(n) => {
if add_self && ids.node_ids.get(n.id()) {
ensure_node_capacity_local(bb, output)?;
let meta = element_metadata(&n.info());
bb.add_node(
n.id(),
n.decimicro_lat(),
n.decimicro_lon(),
n.tags(),
meta.as_ref(),
);
nodes += 1;
}
}
Element::Way(w) => {
let is_parent = w.refs().any(|r| ids.node_ids.get(r));
let is_self = add_self && ids.way_ids.get(w.id());
if is_parent || is_self {
ensure_way_capacity_local(bb, output)?;
refs_buf.clear();
refs_buf.extend(w.refs());
let meta = element_metadata(&w.info());
bb.add_way(w.id(), w.tags(), &refs_buf, meta.as_ref());
ways += 1;
}
}
Element::Relation(r) => {
let is_parent = r.members().any(|m| match m.id {
MemberId::Node(id) => ids.node_ids.get(id),
MemberId::Way(id) => ids.way_ids.get(id),
MemberId::Relation(id) => ids.relation_ids.get(id),
MemberId::Unknown(..) => false,
});
let is_self = add_self && ids.relation_ids.get(r.id());
if is_parent || is_self {
ensure_relation_capacity_local(bb, output)?;
members_buf.clear();
members_buf.extend(r.members().map(|m| MemberData {
id: m.id,
role: m.role().unwrap_or(""),
}));
let meta = element_metadata(&r.info());
bb.add_relation(r.id(), r.tags(), &members_buf, meta.as_ref());
relations += 1;
}
}
}
}
Ok((nodes, ways, relations))
}
#[cfg(test)]
mod tests {
use super::{
FULL_SCAN_ARM_MIN_BLOBS, GetparentsOptions, ScanArm, getparents_dispatched,
getparents_with_arm,
};
use crate::block_builder::{BlockBuilder, HeaderBuilder, MemberData};
use crate::commands::getid::parse_ids;
use crate::writer::{Compression, PbfWriter};
use crate::{Element, ElementReader, HeaderOverrides, MemberId};
fn write_blocks(
path: &std::path::Path,
node_ids: &[i64],
way: Option<(i64, &[i64])>,
index_way_blob: 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().sorted().build().expect("header"))
.expect("write header");
let mut block = BlockBuilder::new();
for &id in node_ids {
block.add_node(id, 0, 0, std::iter::empty::<(&str, &str)>(), None);
}
writer
.write_primitive_block(block.take().expect("node block").expect("nodes"))
.expect("write nodes");
if let Some((way_id, refs)) = way {
block.add_way(way_id, std::iter::empty::<(&str, &str)>(), refs, None);
let bytes = block.take().expect("way block").expect("ways");
if index_way_blob {
writer.write_primitive_block(bytes).expect("write ways");
} else {
writer
.write_primitive_block_no_indexdata(bytes)
.expect("write ways");
}
}
writer.flush().expect("flush fixture");
}
fn fixture(path: &std::path::Path) {
write_blocks(path, &[1, 2], Some((10, &[1, 2])), true);
}
fn ids(path: &std::path::Path) -> Vec<i64> {
let mut ids = Vec::new();
ElementReader::from_path(path)
.expect("read output")
.for_each(|element| match element {
Element::DenseNode(node) => ids.push(node.id()),
Element::Node(node) => ids.push(node.id()),
Element::Way(way) => ids.push(way.id()),
Element::Relation(relation) => ids.push(relation.id()),
})
.expect("iterate output");
ids.sort_unstable();
ids
}
#[test]
fn full_scan_writes_parents() {
let dir = tempfile::tempdir().expect("tempdir");
let input = dir.path().join("input.pbf");
fixture(&input);
let query = parse_ids(&["n1".to_owned()]).expect("ids");
let opts = GetparentsOptions { add_self: true };
let output = dir.path().join("output.pbf");
getparents_with_arm(
&input,
&output,
&query,
&opts,
Compression::default(),
false,
&HeaderOverrides::default(),
ScanArm::FullScan,
)
.expect("getparents");
assert_eq!(ids(&output), vec![1, 10]);
}
fn ids_from_both_arms(input: &std::path::Path, query: &[&str], add_self: bool) -> Vec<i64> {
let dir = tempfile::tempdir().expect("tempdir");
let query: Vec<String> = query.iter().map(|s| (*s).to_owned()).collect();
let query = parse_ids(&query).expect("ids");
let opts = GetparentsOptions { add_self };
let walker = dir.path().join("walker.pbf");
let pipelined = dir.path().join("pipelined.pbf");
for (output, arm) in [(&walker, ScanArm::Walker), (&pipelined, ScanArm::FullScan)] {
getparents_with_arm(
input,
output,
&query,
&opts,
Compression::default(),
false,
&HeaderOverrides::default(),
arm,
)
.expect("getparents");
}
let walker_ids = ids(&walker);
assert_eq!(walker_ids, ids(&pipelined));
walker_ids
}
#[test]
fn walker_and_pipelined_arms_emit_the_same_elements() {
let dir = tempfile::tempdir().expect("tempdir");
let input = dir.path().join("input.pbf");
fixture(&input);
assert_eq!(ids_from_both_arms(&input, &["n1"], false), vec![10]);
assert_eq!(ids_from_both_arms(&input, &["n1"], true), vec![1, 10]);
}
#[test]
fn arms_agree_on_empty_query_result() {
let dir = tempfile::tempdir().expect("tempdir");
let input = dir.path().join("input.pbf");
fixture(&input);
assert!(ids_from_both_arms(&input, &["n999"], true).is_empty());
}
#[test]
fn arms_agree_on_single_data_blob_file() {
let dir = tempfile::tempdir().expect("tempdir");
let input = dir.path().join("input.pbf");
write_blocks(&input, &[1, 2], None, true);
assert_eq!(ids_from_both_arms(&input, &["n1"], true), vec![1]);
}
#[test]
fn arms_agree_when_a_blob_lacks_indexdata() {
let dir = tempfile::tempdir().expect("tempdir");
let input = dir.path().join("input.pbf");
write_blocks(&input, &[1, 2], Some((10, &[1, 2])), false);
assert_eq!(ids_from_both_arms(&input, &["n1"], false), vec![10]);
}
fn write_rich_fixture(path: &std::path::Path) {
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().sorted().build().expect("header"))
.expect("write header");
let mut block = BlockBuilder::new();
for id in [1i64, 2] {
block.add_node(id, 0, 0, std::iter::empty::<(&str, &str)>(), None);
}
writer
.write_primitive_block(block.take().expect("node block").expect("nodes"))
.expect("write nodes");
block.add_way(10, std::iter::empty::<(&str, &str)>(), &[1, 2], None);
writer
.write_primitive_block(block.take().expect("way block").expect("way"))
.expect("write way");
let rel100 = [
MemberData {
id: MemberId::Node(1),
role: "",
},
MemberData {
id: MemberId::Way(10),
role: "",
},
];
block.add_relation(100, std::iter::empty::<(&str, &str)>(), &rel100, None);
let rel101 = [MemberData {
id: MemberId::Relation(100),
role: "",
}];
block.add_relation(101, std::iter::empty::<(&str, &str)>(), &rel101, None);
writer
.write_primitive_block(block.take().expect("relation block").expect("relations"))
.expect("write relations");
writer.flush().expect("flush fixture");
}
#[test]
fn node_query_scans_relation_blobs_for_relation_parents() {
let dir = tempfile::tempdir().expect("tempdir");
let input = dir.path().join("input.pbf");
write_rich_fixture(&input);
assert_eq!(ids_from_both_arms(&input, &["n1"], false), vec![10, 100]);
}
#[test]
fn way_query_skips_way_blobs_but_finds_relation_parents() {
let dir = tempfile::tempdir().expect("tempdir");
let input = dir.path().join("input.pbf");
write_rich_fixture(&input);
assert_eq!(ids_from_both_arms(&input, &["w10"], false), vec![100]);
}
#[test]
fn relation_query_scans_relation_blobs_only() {
let dir = tempfile::tempdir().expect("tempdir");
let input = dir.path().join("input.pbf");
write_rich_fixture(&input);
assert_eq!(ids_from_both_arms(&input, &["r100"], false), vec![101]);
}
#[test]
fn add_self_node_query_emits_node_and_all_parent_kinds() {
let dir = tempfile::tempdir().expect("tempdir");
let input = dir.path().join("input.pbf");
write_rich_fixture(&input);
assert_eq!(ids_from_both_arms(&input, &["n1"], true), vec![1, 10, 100]);
}
#[test]
fn auto_dispatch_crosses_arms_at_the_injected_threshold() {
let dir = tempfile::tempdir().expect("tempdir");
let input = dir.path().join("input.pbf");
fixture(&input);
let query = parse_ids(&["n1".to_owned()]).expect("ids");
let opts = GetparentsOptions { add_self: true };
for (min_blobs, expected_arm) in [
(1, ScanArm::FullScan),
(FULL_SCAN_ARM_MIN_BLOBS, ScanArm::Walker),
] {
let output = dir.path().join(format!("out-{min_blobs}.pbf"));
let (arm, _) = getparents_dispatched(
&input,
&output,
&query,
&opts,
Compression::default(),
false,
&HeaderOverrides::default(),
min_blobs,
)
.expect("getparents");
assert_eq!(arm, expected_arm);
assert_eq!(ids(&output), vec![1, 10]);
}
}
}