use rayon::prelude::*;
use std::fs::File;
use std::io::{BufReader, Read};
use std::path::Path;
use std::rc::Rc;
use super::traits::{BlobData, PbfRandomRead};
use crate::codecs::blob::{BlobReader, DecodedBlob};
use crate::codecs::block_decorators::{HeaderReader, PrimitiveReader};
use crate::models::{Element, ElementType};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ReaderProgress {
pub bytes_read: u64,
pub total_bytes: Option<u64>,
}
impl ReaderProgress {
pub fn fraction(&self) -> Option<f64> {
self.total_bytes
.map(|total| self.bytes_read as f64 / total as f64)
}
pub fn percent(&self) -> Option<f64> {
self.fraction().map(|f| f * 100.0)
}
}
pub struct PbfReader<R: Read + Send> {
blob_reader: BlobReader<R>,
total_bytes: Option<u64>,
}
impl<R: Read + Send> PbfReader<R> {
pub fn new(reader: R) -> PbfReader<R> {
Self {
blob_reader: BlobReader::new(reader),
total_bytes: None,
}
}
pub fn progress(&self) -> ReaderProgress {
ReaderProgress {
bytes_read: self.blob_reader.offset,
total_bytes: self.total_bytes,
}
}
pub fn read_next_blob(&mut self) -> anyhow::Result<Option<BlobData>> {
if self.blob_reader.eof {
return Ok(None);
}
let offset = self.blob_reader.offset;
match self.blob_reader.next_blob()? {
Some(blob) => {
let data = match blob.decode()? {
Some(DecodedBlob::OsmHeader(header)) => {
HeaderReader::new(header).validate_features()?;
BlobData {
nodes: Vec::with_capacity(0),
ways: Vec::with_capacity(0),
relations: Vec::with_capacity(0),
offset,
}
}
Some(DecodedBlob::OsmData(data)) => {
let decorator = PrimitiveReader::new(data)?;
let (nodes, ways, relations) = decorator.get_all_elements()?;
BlobData {
nodes,
ways,
relations,
offset,
}
}
None => BlobData {
nodes: Vec::with_capacity(0),
ways: Vec::with_capacity(0),
relations: Vec::with_capacity(0),
offset,
},
};
Ok(Some(data))
}
None => Ok(None),
}
}
pub fn read<F>(&mut self, mut callback: F) -> anyhow::Result<()>
where
F: FnMut(Option<HeaderReader>, Option<Element>),
{
while let Some(blob) = self.blob_reader.next_blob()? {
match blob.decode()? {
Some(DecodedBlob::OsmHeader(b)) => {
let header_reader = HeaderReader::new(b);
header_reader.validate_features()?;
callback(Some(header_reader), None);
}
Some(DecodedBlob::OsmData(data)) => {
let decorator = PrimitiveReader::new(data)?;
decorator.for_each_element(|el| callback(None, Some(el)))?;
}
None => {}
}
}
Ok(())
}
pub fn par_find<F>(
self,
inclination: Option<&ElementType>,
callback: F,
) -> anyhow::Result<Vec<Element>>
where
F: Fn(&Element) -> bool + Send + Sync,
{
let result = self
.blob_reader
.par_bridge()
.map(|blob| -> anyhow::Result<Vec<Element>> {
let decoded = match blob?.decode()? {
Some(DecodedBlob::OsmData(b)) => Some(PrimitiveReader::new(b)?),
_ => None,
};
let Some(p) = decoded else {
return Ok(Vec::new());
};
if let Some(element_type) = inclination {
let result = match element_type {
ElementType::Node => p
.get_nodes()?
.into_iter()
.map(Element::Node)
.filter(&callback)
.collect::<Vec<Element>>(),
ElementType::Way => p
.get_ways()?
.into_iter()
.map(Element::Way)
.filter(&callback)
.collect::<Vec<Element>>(),
ElementType::Relation => p
.get_relations()?
.into_iter()
.map(Element::Relation)
.filter(&callback)
.collect::<Vec<Element>>(),
};
Ok(result)
} else {
let (nodes, ways, relations) = p.get_all_elements()?;
let mut result: Vec<Element> = nodes
.into_iter()
.map(Element::Node)
.filter(&callback)
.collect();
result.extend(ways.into_iter().map(Element::Way).filter(&callback));
result.extend(
relations
.into_iter()
.map(Element::Relation)
.filter(&callback),
);
Ok(result)
}
})
.reduce(
|| Ok(Vec::new()),
|acc: anyhow::Result<Vec<Element>>, item: anyhow::Result<Vec<Element>>| match (
acc, item,
) {
(Ok(mut a), Ok(b)) => {
a.extend(b);
Ok(a)
}
(Err(e), _) | (_, Err(e)) => Err(e),
},
)?;
Ok(result)
}
}
impl PbfReader<BufReader<File>> {
pub fn from_path<P: AsRef<Path>>(path: P) -> anyhow::Result<Self> {
let total_bytes = std::fs::metadata(path.as_ref())?.len();
let f = File::open(path)?;
let reader = BufReader::new(f);
Ok(Self {
blob_reader: BlobReader::new(reader),
total_bytes: Some(total_bytes),
})
}
pub fn rewind(&mut self) -> anyhow::Result<()> {
self.blob_reader.rewind()
}
}
impl PbfRandomRead for PbfReader<BufReader<File>> {
fn read_blob_by_offset(&mut self, offset: u64) -> anyhow::Result<Rc<BlobData>> {
self.blob_reader.seek(offset)?;
let data = self
.read_next_blob()?
.ok_or(anyhow!("no blob data found."))?;
Ok(Rc::new(data))
}
}