use std::{fmt::Debug, mem::size_of, path::Path, sync::Arc};
use anyhow::{Result, bail};
use async_trait::async_trait;
use moka::future::Cache;
#[cfg(feature = "cli")]
use versatiles_core::utils::PrettyPrint;
use versatiles_core::{
Blob, ByteRange, GeoBBox, TileBBox, TileCompression, TileCoord, TileFormat, TileJSON, TilePyramid, TileStream,
compression::decompress,
io::{DataReader, DataReaderFile},
utils::HilbertIndex,
};
use versatiles_derive::context;
use super::types::{EntriesV3, HeaderV3};
use crate::{
SharedTileSource, SourceType, Tile, TileSource, TileSourceMetadata, TilesReader, TilesRuntime, Traversal,
TraversalOrder, TraversalSize, container::tile_chunking::Chunks,
};
#[derive(Debug)]
pub struct PMTilesReader {
pub data_reader: Arc<DataReader>,
pub header: HeaderV3,
pub internal_compression: TileCompression,
pub leaves_bytes: Arc<Blob>,
pub leaves_cache: Cache<ByteRange, Arc<EntriesV3>>,
pub tilejson: TileJSON,
pub metadata: TileSourceMetadata,
pub root_bytes_uncompressed: Blob,
pub root_entries: Arc<EntriesV3>,
}
impl PMTilesReader {
#[context("opening PMTiles at '{}'", path.display())]
pub async fn open(path: &Path, runtime: TilesRuntime) -> Result<PMTilesReader> {
PMTilesReader::open_data(DataReaderFile::open(path)?, runtime).await
}
#[context("opening PMTiles from reader")]
pub async fn open_data(data_reader: DataReader, _runtime: TilesRuntime) -> Result<PMTilesReader>
where
Self: Sized,
{
log::debug!("Opening PMTilesReader for {}", data_reader.name());
let header = HeaderV3::deserialize(&data_reader.read_range(&ByteRange::new(0, HeaderV3::len())).await?)?;
log::trace!("Header: {header:?}");
let internal_compression = header.internal_compression.as_value()?;
log::trace!("Internal compression: {internal_compression:?}");
let metadata_fut = data_reader.read_range(&header.metadata);
let root_dir_fut = data_reader.read_range(&header.root_dir);
let leaf_dirs_fut = async {
if header.leaf_dirs.length == 0 {
Ok::<Blob, anyhow::Error>(Blob::default())
} else {
data_reader.read_range(&header.leaf_dirs).await
}
};
let (meta, root_bytes, leaves_bytes) =
futures::future::try_join3(metadata_fut, root_dir_fut, leaf_dirs_fut).await?;
let meta = decompress(meta, &internal_compression)?;
let mut tilejson = TileJSON::try_from_blob_or_default(&meta);
log::trace!("TileJSON: {tilejson:?}");
log::trace!("Root directory bytes length: {}", root_bytes.len());
let root_bytes_uncompressed = decompress(root_bytes, &internal_compression)?;
log::trace!(
"Root directory bytes uncompressed length: {}",
root_bytes_uncompressed.len()
);
log::trace!("Leaf directories bytes length: {}", leaves_bytes.len());
if tilejson.bounds.is_none() {
tilejson.bounds = GeoBBox::new(
f64::from(header.min_lon_e7) / 1e7,
f64::from(header.min_lat_e7) / 1e7,
f64::from(header.max_lon_e7) / 1e7,
f64::from(header.max_lat_e7) / 1e7,
)
.ok();
}
if tilejson.zoom_min().is_none() {
tilejson.set_zoom_min(header.min_zoom);
}
if tilejson.zoom_max().is_none() {
tilejson.set_zoom_max(header.max_zoom);
}
let metadata = TileSourceMetadata::new(
header.tile_type.as_value()?,
header.tile_compression.as_value()?,
Traversal {
order: TraversalOrder::PMTiles,
size: TraversalSize::new_default(),
},
None,
);
log::trace!("Reader parameters: {metadata:?}");
let root_entries = Arc::new(EntriesV3::from_blob(&root_bytes_uncompressed)?);
Ok(PMTilesReader {
data_reader: Arc::new(data_reader),
header,
internal_compression,
leaves_bytes: Arc::new(leaves_bytes),
leaves_cache: Cache::builder()
.max_capacity(100_000_000)
.weigher(|_k, v: &Arc<EntriesV3>| {
let bytes = size_of::<ByteRange>() + size_of::<Arc<EntriesV3>>() + v.len() * 32;
u32::try_from(bytes).unwrap_or(u32::MAX)
})
.build(),
tilejson,
metadata,
root_bytes_uncompressed,
root_entries,
})
}
#[context("reading PMTiles root entries")]
pub fn tile_entries(&self) -> Result<EntriesV3> {
EntriesV3::from_blob(&self.root_bytes_uncompressed)
}
#[allow(clippy::too_many_arguments)]
async fn lookup_tile_by_id(
tile_id: u64,
data_reader: &DataReader,
root_entries: Arc<EntriesV3>,
leaves_cache: &Cache<ByteRange, Arc<EntriesV3>>,
leaves_bytes: &Arc<Blob>,
tile_data_offset: u64,
tile_compression: &TileCompression,
tile_format: &TileFormat,
internal_compression: &TileCompression,
) -> Result<Option<Tile>> {
let mut entries = root_entries;
for _depth in 0..3 {
let Some(entry) = entries.find_tile(tile_id) else {
return Ok(None);
};
if entry.range.length > 0 {
if entry.run_length > 0 {
return Ok(Some(Tile::from_blob(
data_reader
.read_range(&entry.range.shifted_forward(tile_data_offset))
.await?,
*tile_compression,
*tile_format,
)));
}
let range = entry.range;
let leaves = Arc::clone(leaves_bytes);
let compression = *internal_compression;
entries = leaves_cache
.try_get_with(range, async move {
let mut blob = leaves.read_range(&range)?;
blob = decompress(blob, &compression)?;
let parsed = EntriesV3::from_blob(&blob)?;
anyhow::Ok(Arc::new(parsed))
})
.await
.map_err(|e: Arc<anyhow::Error>| anyhow::anyhow!("{e:#}"))?;
} else {
return Ok(None);
}
}
bail!("not found")
}
async fn resolve_tile_range(
tile_id: u64,
root_entries: Arc<EntriesV3>,
leaves_cache: &Cache<ByteRange, Arc<EntriesV3>>,
leaves_bytes: &Arc<Blob>,
tile_data_offset: u64,
internal_compression: TileCompression,
) -> Result<Option<ByteRange>> {
let mut entries = root_entries;
for _depth in 0..3 {
let Some(entry) = entries.find_tile(tile_id) else {
return Ok(None);
};
if entry.range.length == 0 {
return Ok(None);
}
if entry.run_length > 0 {
return Ok(Some(entry.range.shifted_forward(tile_data_offset)));
}
let range = entry.range;
let leaves = Arc::clone(leaves_bytes);
entries = leaves_cache
.try_get_with(range, async move {
let mut blob = leaves.read_range(&range)?;
blob = decompress(blob, &internal_compression)?;
anyhow::Ok(Arc::new(EntriesV3::from_blob(&blob)?))
})
.await
.map_err(|e: Arc<anyhow::Error>| anyhow::anyhow!("{e:#}"))?;
}
Ok(None)
}
async fn get_chunks(&self, bbox: TileBBox) -> Result<Chunks> {
let mut tile_ranges: Vec<(TileCoord, ByteRange)> = Vec::new();
let coords: Vec<TileCoord> = bbox.iter_coords().collect();
for coord in coords {
let Ok(tile_id) = coord.get_hilbert_index() else {
continue;
};
if let Some(range) = Self::resolve_tile_range(
tile_id,
Arc::clone(&self.root_entries),
&self.leaves_cache,
&self.leaves_bytes,
self.header.tile_data.offset,
self.internal_compression,
)
.await?
{
tile_ranges.push((coord, range));
}
}
Ok(Chunks::from_tile_ranges(tile_ranges))
}
}
#[context("building tile pyramid from PMTiles directories")]
fn calc_tile_pyramid(
root_bytes_uncompressed: &Blob,
leaves_bytes: &Blob,
compression: TileCompression,
) -> Result<TilePyramid> {
let mut coords: Vec<TileCoord> = Vec::new();
parse_directories(&mut coords, root_bytes_uncompressed, leaves_bytes, compression)?;
fn parse_directories(
coords: &mut Vec<TileCoord>,
dir: &Blob,
leaves_bytes: &Blob,
compression: TileCompression,
) -> Result<u64> {
log::trace!("parse_directories");
let entries = EntriesV3::from_blob(dir)?;
let entries = entries.iter().collect::<Vec<_>>();
let mut total_entries = 0;
for entry in &entries {
if entry.range.length > 0 {
if entry.run_length > 0 {
for i in 0..u64::from(entry.run_length) {
coords.push(TileCoord::from_hilbert_index(i + entry.tile_id)?);
}
total_entries += u64::from(entry.run_length);
} else {
let range = entry.range;
let mut blob = leaves_bytes.read_range(&range)?;
blob = decompress(blob, &compression)?;
total_entries += parse_directories(coords, &blob, leaves_bytes, compression)?;
}
}
}
Ok(total_entries)
}
Ok(TilePyramid::from_tile_coords(coords.into_iter()))
}
#[async_trait]
impl TilesReader for PMTilesReader {
async fn open_reader(reader: DataReader, runtime: TilesRuntime) -> Result<SharedTileSource> {
Ok(Self::open_data(reader, runtime).await?.into_shared())
}
}
#[async_trait]
impl TileSource for PMTilesReader {
fn source_type(&self) -> Arc<SourceType> {
SourceType::new_container("pmtiles", self.data_reader.name())
}
fn metadata(&self) -> &TileSourceMetadata {
&self.metadata
}
fn tilejson(&self) -> &TileJSON {
&self.tilejson
}
async fn tile_pyramid(&self) -> Result<Arc<TilePyramid>> {
self.metadata.get_or_compute_tile_pyramid(|| {
calc_tile_pyramid(
&self.root_bytes_uncompressed,
&self.leaves_bytes,
self.internal_compression,
)
})
}
#[context("fetching tile {:?} from PMTiles", coord)]
async fn tile(&self, coord: &TileCoord) -> Result<Option<Tile>> {
log::trace!("tile {coord:?}");
let tile_id = coord.get_hilbert_index()?;
Self::lookup_tile_by_id(
tile_id,
&self.data_reader,
Arc::clone(&self.root_entries),
&self.leaves_cache,
&self.leaves_bytes,
self.header.tile_data.offset,
self.metadata.tile_compression(),
self.metadata.tile_format(),
&self.internal_compression,
)
.await
}
async fn tile_stream(&self, bbox: TileBBox) -> Result<TileStream<'static, Tile>> {
log::trace!("pmtiles::tile_stream {bbox:?}");
let bbox = self.metadata.intersection_bbox(&bbox);
let chunks = self.get_chunks(bbox).await?;
Ok(chunks.stream(
Arc::clone(&self.data_reader),
*self.metadata.tile_compression(),
*self.metadata.tile_format(),
))
}
#[context("streaming tile sizes for bbox {:?}", bbox)]
async fn tile_size_stream(&self, bbox: TileBBox) -> Result<TileStream<'static, u32>> {
let bbox = self.metadata.intersection_bbox(&bbox);
let mut tile_sizes: Vec<(TileCoord, u32)> = Vec::new();
let coords: Vec<TileCoord> = bbox.iter_coords().collect();
for coord in coords {
let Ok(tile_id) = coord.get_hilbert_index() else {
continue;
};
if let Some(range) = Self::resolve_tile_range(
tile_id,
Arc::clone(&self.root_entries),
&self.leaves_cache,
&self.leaves_bytes,
self.header.tile_data.offset,
self.internal_compression,
)
.await? && let Ok(size) = u32::try_from(range.length)
{
tile_sizes.push((coord, size));
}
}
Ok(TileStream::from_vec(tile_sizes))
}
async fn tile_coord_stream(&self, bbox: TileBBox) -> Result<TileStream<'static, ()>> {
Ok(self.tile_size_stream(bbox).await?.filter_map(|_, _| Some(())))
}
#[cfg(feature = "cli")]
#[context("probing PMTiles container metadata")]
async fn probe_container(&self, print: &mut PrettyPrint, _runtime: &TilesRuntime) -> Result<()> {
let h = &self.header;
print
.add_key_value("addressed tiles count", &h.addressed_tiles_count)
.await;
print.add_key_value("tile entries count", &h.tile_entries_count).await;
print.add_key_value("tile contents count", &h.tile_contents_count).await;
print.add_key_value("clustered", &h.clustered).await;
print
.add_key_value("internal compression", &h.internal_compression)
.await;
print.add_key_value("tile type", &h.tile_type).await;
print
.add_key_value("zoom range", &format!("{}..{}", h.min_zoom, h.max_zoom))
.await;
print.add_key_value("root dir size", &h.root_dir.length).await;
print.add_key_value("metadata size", &h.metadata.length).await;
print.add_key_value("leaf dirs size", &h.leaf_dirs.length).await;
print.add_key_value("tile data size", &h.tile_data.length).await;
Ok(())
}
}
#[cfg(test)]
mod tests {
use std::{env::current_dir, path::PathBuf, sync::LazyLock};
use versatiles_core::assert_wildcard;
use super::*;
static PATH: LazyLock<PathBuf> = LazyLock::new(|| current_dir().unwrap().join("../testdata/berlin.pmtiles"));
#[tokio::test]
async fn reader() -> Result<()> {
let reader = PMTilesReader::open(&PATH, TilesRuntime::default()).await?;
assert_wildcard!(
reader.source_type().to_string(),
"container 'pmtiles' ('*testdata?berlin.pmtiles')"
);
assert_eq!(
format!("{:?}", reader.header),
"HeaderV3 { root_dir: ByteRange[127,504], metadata: ByteRange[631,888], leaf_dirs: ByteRange[1519,0], tile_data: ByteRange[1519,11339914], addressed_tiles_count: 130, tile_entries_count: 130, tile_contents_count: 130, clustered: true, internal_compression: Gzip, tile_compression: Gzip, tile_type: MVT, min_zoom: 0, max_zoom: 14, min_lon_e7: 133000000, min_lat_e7: 524500000, max_lon_e7: 134600000, max_lat_e7: 525500000, center_zoom: 2, center_lon_e7: 133813476, center_lat_e7: 525028064 }"
);
assert_wildcard!(
reader.tilejson().stringify(),
"{\"author\":\"OpenStreetMap contributors\",*,\"version\":\"2.0\"}"
);
assert_eq!(
format!("{:?}", reader.metadata()),
"TileSourceMetadata { tile_compression: Gzip, tile_format: MVT, traversal: Traversal(PMTiles,full), tile_pyramid: RwLock { data: None, poisoned: false, .. } }"
);
assert_eq!(
reader
.tile(&TileCoord::new(0, 0, 0)?)
.await?
.unwrap()
.as_blob(reader.metadata.tile_compression())?
.len(),
90891
);
assert_eq!(
reader
.tile(&TileCoord::new(14, 8800, 5370)?)
.await?
.unwrap()
.as_blob(reader.metadata.tile_compression())?
.len(),
130653
);
assert!(reader.tile(&TileCoord::new(16, 0, 0)?).await?.is_none());
Ok(())
}
#[tokio::test]
async fn satisfies_stream_count_invariant() -> Result<()> {
let reader = PMTilesReader::open(&PATH, TilesRuntime::default()).await?;
let pyr = reader.tile_pyramid().await?;
for z in 0..=pyr.level_max().unwrap_or(0) {
let bbox = pyr.level_ref(z).to_bbox();
if bbox.is_empty() {
continue;
}
crate::testing::assert_stream_counts_agree(&reader, bbox).await?;
}
Ok(())
}
#[tokio::test]
async fn tile_size_stream_matches_tile_reads() -> Result<()> {
let reader = PMTilesReader::open(&PATH, TilesRuntime::default()).await?;
let bbox = TileBBox::from_min_and_max(9, 274, 167, 275, 168)?;
let compression = reader.metadata().tile_compression();
let mut sizes: Vec<(TileCoord, u32)> = reader.tile_size_stream(bbox).await?.to_vec().await;
sizes.sort_by_key(|(c, _)| (c.y, c.x));
assert_eq!(sizes.len(), 4);
for (coord, size) in &sizes {
let blob = reader
.tile(coord)
.await?
.expect("tile should exist")
.into_blob(compression)?;
assert_eq!(u64::from(*size), blob.len(), "size mismatch at {coord:?}");
}
Ok(())
}
#[tokio::test]
async fn tile_size_stream_single_tile() -> Result<()> {
let reader = PMTilesReader::open(&PATH, TilesRuntime::default()).await?;
let bbox = TileBBox::from_min_and_max(0, 0, 0, 0, 0)?;
let sizes: Vec<(TileCoord, u32)> = reader.tile_size_stream(bbox).await?.to_vec().await;
assert_eq!(sizes.len(), 1);
assert_eq!(sizes[0].0, TileCoord::new(0, 0, 0)?);
assert_eq!(sizes[0].1, 90891);
Ok(())
}
#[tokio::test]
async fn tile_size_stream_empty_for_missing_zoom() -> Result<()> {
let reader = PMTilesReader::open(&PATH, TilesRuntime::default()).await?;
let bbox = TileBBox::from_min_and_max(20, 0, 0, 3, 3)?;
let sizes: Vec<(TileCoord, u32)> = reader.tile_size_stream(bbox).await?.to_vec().await;
assert!(sizes.is_empty());
Ok(())
}
#[cfg(feature = "cli")]
#[tokio::test]
async fn probe() -> Result<()> {
use versatiles_core::utils::PrettyPrint;
let reader = PMTilesReader::open(&PATH, TilesRuntime::default()).await?;
let runtime = TilesRuntime::default();
let mut printer = PrettyPrint::new();
reader
.probe_container(&mut printer.category("container").await, &runtime)
.await?;
let output = printer.stringify().await;
assert!(
output.contains("addressed tiles count: 130"),
"unexpected output: {output}"
);
assert!(output.contains("clustered: true"), "unexpected output: {output}");
assert!(output.contains("tile data size:"), "unexpected output: {output}");
Ok(())
}
#[tokio::test]
async fn concurrent_tile_lookups_share_leaves_cache() -> Result<()> {
let reader = Arc::new(PMTilesReader::open(&PATH, TilesRuntime::default()).await?);
let coord = TileCoord::new(14, 8800, 5370)?;
let mut handles = Vec::new();
for _ in 0..16 {
let r = Arc::clone(&reader);
handles.push(tokio::spawn(async move { r.tile(&coord).await.unwrap() }));
}
for h in handles {
let tile = h.await.unwrap().expect("tile should exist");
assert!(!tile.into_blob(&TileCompression::Gzip)?.is_empty());
}
Ok(())
}
#[tokio::test]
async fn tile_stream_matches_individual_reads() -> Result<()> {
let reader = PMTilesReader::open(&PATH, TilesRuntime::default()).await?;
let bbox = TileBBox::from_min_and_max(9, 274, 167, 275, 168)?;
let stream = reader.tile_stream(bbox).await?;
let stream_tiles: Vec<_> = stream.to_vec().await;
assert_eq!(stream_tiles.len(), 4);
for (coord, mut tile) in stream_tiles {
let stream_blob = tile.as_blob(reader.metadata().tile_compression())?;
let single_blob = reader
.tile(&coord)
.await?
.expect("tile should exist")
.into_blob(reader.metadata().tile_compression())?;
assert_eq!(
stream_blob.as_slice(),
single_blob.as_slice(),
"blob mismatch at {coord:?}"
);
}
Ok(())
}
}