use std::{
collections::HashMap,
hash::{DefaultHasher, Hash, Hasher},
path::{Path, PathBuf},
sync::Arc,
};
use anyhow::{Context, Result, ensure};
use async_trait::async_trait;
use futures::lock::Mutex;
use versatiles_core::{
compression::compress,
io::{DataReaderFile, DataReaderTrait, DataWriterFile, DataWriterTrait},
types::{Blob, ByteRange, TileCompression},
utils::HilbertIndex,
};
use versatiles_derive::context;
use super::types::{EntriesV3, EntryV3, HeaderV3, PMTilesCompression};
use crate::{
TileSource, TileSourceMetadata, TileSourceTraverseExt, TilesRuntime, TilesWriter,
traversal::{Traversal, TraversalOrder},
};
pub struct PMTilesWriter {}
const ALLOW_UNCLUSTERED: &str = "allow_unclustered";
const REORDER: &str = "reorder";
const TEMP_DIR: &str = "temp_dir";
#[async_trait]
impl TilesWriter for PMTilesWriter {
fn supported_options() -> &'static [&'static str] {
&[ALLOW_UNCLUSTERED, REORDER, TEMP_DIR]
}
async fn write_to_path(reader: &mut dyn TileSource, path: &Path, runtime: TilesRuntime) -> Result<()> {
let mut writer = DataWriterFile::from_path(path)?;
Self::write(reader, &mut writer, runtime, path.parent()).await
}
#[context("writing PMTiles to DataWriter")]
async fn write_to_writer(
reader: &mut dyn TileSource,
writer: &mut dyn DataWriterTrait,
runtime: TilesRuntime,
) -> Result<()> {
Self::write(reader, writer, runtime, None).await
}
}
impl PMTilesWriter {
#[context("writing PMTiles")]
async fn write(
reader: &mut dyn TileSource,
writer: &mut dyn DataWriterTrait,
runtime: TilesRuntime,
output_dir: Option<&Path>,
) -> Result<()> {
const INTERNAL_COMPRESSION: TileCompression = TileCompression::Gzip;
let parameters = reader.metadata().clone();
let plan = OrderPlan::negotiate(¶meters, &runtime)?;
let (traversal, reorder, clustered) = (plan.traversal, plan.reorder, plan.clustered);
let entries = EntriesV3::new();
writer.set_position(16384)?;
let tilejson = reader.tilejson();
let tile_pyramid = reader.tile_pyramid().await?;
let mut header = HeaderV3::from_parameters(¶meters, tile_pyramid.as_ref(), tilejson);
let mut metadata: Blob = tilejson.into();
metadata = compress(metadata, &INTERNAL_COMPRESSION)?;
header.metadata = writer.append(&metadata)?;
let tile_data_start = writer.position()?;
let temp = if reorder {
Some(TempTileData::create(&runtime, output_dir)?)
} else {
None
};
let mut temp_writer = match &temp {
Some(temp) => Some(DataWriterFile::from_path(temp.path())?),
None => None,
};
let tile_data_writer: &mut dyn DataWriterTrait = match temp_writer.as_mut() {
Some(w) => w,
None => &mut *writer,
};
let pass1_origin = if reorder { 0 } else { tile_data_start };
let writer_mutex = Arc::new(Mutex::new(tile_data_writer));
let entries_mutex = Arc::new(Mutex::new(entries));
let dedup_map: Arc<Mutex<HashMap<u64, ByteRange>>> = Arc::new(Mutex::new(HashMap::new()));
let tile_compression = *reader.metadata().tile_compression();
reader
.traverse_all_tiles(
&traversal,
|_bbox, stream| {
let writer_mutex = Arc::clone(&writer_mutex);
let entries_mutex = Arc::clone(&entries_mutex);
let dedup_map = Arc::clone(&dedup_map);
Box::pin(async move {
let stream = stream.map_parallel_try(move |_coord, mut tile| {
tile.as_blob(&tile_compression)?;
Ok(tile)
});
let mut tiles = Vec::new();
for (coord, result) in stream.to_vec().await {
tiles.push((coord, result?));
}
tiles
.sort_by_key(|(coord, _)| coord.get_hilbert_index().expect("valid tile coord has hilbert index"));
let mut writer = writer_mutex.lock().await;
let mut entries = entries_mutex.lock().await;
let mut dedup = dedup_map.lock().await;
for (coord, mut tile) in tiles {
let id = coord.get_hilbert_index()?;
let blob = tile.as_blob(&tile_compression)?;
let mut hasher = DefaultHasher::new();
blob.as_slice().hash(&mut hasher);
let hash = hasher.finish();
let range = if let Some(&existing) = dedup.get(&hash) {
existing
} else {
let range = writer.append(blob)?.shifted_backward(pass1_origin)?;
dedup.insert(hash, range);
range
};
entries.push(EntryV3::new(id, range, 1));
}
Ok(())
})
},
runtime.clone(),
)
.await?;
let mut entries = entries_mutex.lock().await;
let tile_contents_count = dedup_map.lock().await.len() as u64;
drop(writer_mutex.lock().await);
drop(temp_writer);
if let Some(temp) = &temp {
copy_tiles_in_id_order(&mut entries, temp, writer, tile_data_start, &runtime).await?;
}
let tile_data_end = writer.position()?;
header.tile_data = ByteRange::new(tile_data_start, tile_data_end - tile_data_start);
writer.set_position(HeaderV3::len())?;
let directory = entries.build_directory(16384 - HeaderV3::len(), INTERNAL_COMPRESSION)?;
header.root_dir = writer.append(&directory.root_bytes)?;
writer.set_position(tile_data_end)?;
header.leaf_dirs = writer.append(&directory.leaves_bytes)?;
header.clustered = clustered;
header.internal_compression = PMTilesCompression::from_value(INTERNAL_COMPRESSION)?;
header.addressed_tiles_count = entries.tile_count();
header.tile_entries_count = entries.len() as u64;
header.tile_contents_count = tile_contents_count;
writer.write_start(&header.serialize()?)?;
Ok(())
}
}
struct OrderPlan {
traversal: Traversal,
reorder: bool,
clustered: bool,
}
impl OrderPlan {
fn negotiate(parameters: &TileSourceMetadata, runtime: &TilesRuntime) -> Result<Self> {
let required = Traversal::new(TraversalOrder::PMTiles, 1, 64)?;
let available = parameters.traversal();
let source_is_ordered = available.can_translate_to(&required);
let allow_unclustered = runtime.writer_option_bool(ALLOW_UNCLUSTERED)?.unwrap_or(false);
let reorder_requested = runtime.writer_option_bool(REORDER)?.unwrap_or(false);
ensure!(
!(allow_unclustered && reorder_requested),
"`{ALLOW_UNCLUSTERED}` and `{REORDER}` ask for different things: one gives up clustering \
to write in a single pass, the other keeps it by writing the tile data twice. Pick one."
);
let reorder = reorder_requested && !source_is_ordered;
if reorder_requested && source_is_ordered {
log::info!(
"`{REORDER}` was requested but this source already produces Hilbert order; skipping the extra pass"
);
}
let clustered = source_is_ordered || reorder;
ensure!(
clustered || allow_unclustered,
"cannot write PMTiles: this source produces tiles in {} order, and PMTiles stores them \
in Hilbert order. Operations that build lower zoom levels from higher ones, such as \
`raster_overview`, force this order, and reordering while writing would mean holding \
the whole tileset in memory.\n\
\n\
Three ways forward:\n\
- write .versatiles or .mbtiles instead — they accept any order, and cost nothing\n\
- `--writer-option {REORDER}=true` — keeps the archive clustered, at the cost of \
writing the tile data twice and temporary disk the size of the output (see \
`{TEMP_DIR}`)\n\
- `--writer-option {ALLOW_UNCLUSTERED}=true` — a single pass, but the tile data is not \
physically clustered, so serving it needs more range requests\n\
\n\
The archive is valid and reads correctly whichever you pick. None is chosen for you \
because which cost is acceptable depends on the job.",
available.order()
);
Ok(Self {
traversal: if source_is_ordered { required } else { Traversal::ANY },
reorder,
clustered,
})
}
}
struct TempTileData {
file: tempfile::NamedTempFile,
}
impl TempTileData {
fn create(runtime: &TilesRuntime, output_dir: Option<&Path>) -> Result<Self> {
let dir: PathBuf = match runtime.writer_option(TEMP_DIR) {
Some(dir) => PathBuf::from(dir),
None => match output_dir {
Some(dir) if dir.as_os_str().is_empty() => PathBuf::from("."),
Some(dir) => dir.to_path_buf(),
None => std::env::temp_dir(),
},
};
let file = tempfile::Builder::new()
.prefix(".versatiles-reorder-")
.suffix(".tmp")
.tempfile_in(&dir)
.with_context(|| format!("could not create the reorder temporary file in {dir:?}"))?;
log::debug!("reorder: writing tile data to {:?} first", file.path());
Ok(Self { file })
}
fn path(&self) -> &Path {
self.file.path()
}
}
async fn copy_tiles_in_id_order(
entries: &mut EntriesV3,
temp: &TempTileData,
writer: &mut dyn DataWriterTrait,
tile_data_start: u64,
runtime: &TilesRuntime,
) -> Result<()> {
let reader = DataReaderFile::open(temp.path())?;
entries.sort_by_tile_id();
let progress = runtime.create_progress("reordering tiles", entries.len() as u64);
let mut moved: HashMap<ByteRange, ByteRange> = HashMap::new();
for entry in entries.iter_mut() {
let destination = if let Some(&range) = moved.get(&entry.range) {
range
} else {
let blob = reader.read_range(&entry.range).await?;
let range = writer.append(&blob)?.shifted_backward(tile_data_start)?;
moved.insert(entry.range, range);
range
};
entry.range = destination;
progress.inc(1);
}
progress.finish();
log::debug!(
"reorder: copied {} unique blob(s) for {} entries",
moved.len(),
entries.len()
);
Ok(())
}
#[cfg(test)]
mod tests {
use versatiles_core::{TileBBox, TileFormat, TilePyramid, io::*};
use versatiles_derive::context;
use super::*;
use crate::{
TileSourceMetadata,
container::{
mock::{MockReader, MockWriter},
pmtiles::PMTilesReader,
},
};
fn depth_first_source() -> Result<MockReader> {
MockReader::new_mock(
TilePyramid::new_full_up_to(3),
TileSourceMetadata::new(
TileFormat::PNG,
TileCompression::Uncompressed,
Traversal::new(TraversalOrder::DepthFirst, 1, 16)?,
None,
),
)
}
fn read_header(blob: &Blob) -> Result<HeaderV3> {
let len = usize::try_from(HeaderV3::len())?;
HeaderV3::deserialize(&Blob::from(&blob.as_slice()[..len]))
}
async fn write_to_blob(reader: &mut MockReader, runtime: TilesRuntime) -> Result<Blob> {
let mut writer = DataWriterBlob::new()?;
PMTilesWriter::write_to_writer(reader, &mut writer, runtime).await?;
Ok(writer.into_blob())
}
#[tokio::test]
async fn unclustered_source_is_refused_by_default() -> Result<()> {
let err = write_to_blob(&mut depth_first_source()?, TilesRuntime::new_silent())
.await
.unwrap_err();
let err = format!("{err:#}");
assert!(err.contains("DepthFirst"), "should name the source's order: {err}");
assert!(err.contains("raster_overview"), "should name the usual cause: {err}");
assert!(err.contains(".versatiles"), "should offer another container: {err}");
assert!(err.contains("reorder=true"), "should offer the reordering pass: {err}");
assert!(
err.contains("allow_unclustered=true"),
"should offer the unclustered opt-in: {err}"
);
assert!(err.contains("twice"), "should state reorder's cost: {err}");
assert!(
err.contains("range requests"),
"should state the unclustered cost: {err}"
);
Ok(())
}
#[tokio::test]
async fn allow_unclustered_accepts_the_source_and_says_so_in_the_header() -> Result<()> {
let runtime = TilesRuntime::builder()
.silent_progress(true)
.writer_option("allow_unclustered", "true")
.build();
let blob = write_to_blob(&mut depth_first_source()?, runtime).await?;
let header = read_header(&blob)?;
assert!(
!header.clustered,
"an archive written out of order must not claim to be clustered"
);
Ok(())
}
#[tokio::test]
async fn a_clustered_source_still_reports_clustered() -> Result<()> {
let mut reader = MockReader::new_mock(
TilePyramid::new_full_up_to(3),
TileSourceMetadata::new(TileFormat::PNG, TileCompression::Uncompressed, Traversal::ANY, None),
)?;
let runtime = TilesRuntime::builder()
.silent_progress(true)
.writer_option("allow_unclustered", "true")
.build();
let blob = write_to_blob(&mut reader, runtime).await?;
assert!(read_header(&blob)?.clustered);
Ok(())
}
#[tokio::test]
async fn allow_unclustered_false_still_refuses() -> Result<()> {
let runtime = TilesRuntime::builder()
.silent_progress(true)
.writer_option("allow_unclustered", "false")
.build();
assert!(write_to_blob(&mut depth_first_source()?, runtime).await.is_err());
Ok(())
}
#[tokio::test]
async fn a_non_boolean_option_value_is_an_error() -> Result<()> {
let runtime = TilesRuntime::builder()
.silent_progress(true)
.writer_option("allow_unclustered", "maybe")
.build();
let err = write_to_blob(&mut depth_first_source()?, runtime).await.unwrap_err();
assert!(format!("{err:#}").contains("expects true or false"));
Ok(())
}
async fn write_with(reader: &mut MockReader, options: &[(&str, &str)]) -> Result<Blob> {
let mut builder = TilesRuntime::builder().silent_progress(true);
for (key, value) in options {
builder = builder.writer_option(key, value);
}
write_to_blob(reader, builder.build()).await
}
#[tokio::test]
async fn reorder_accepts_an_unordered_source_and_keeps_it_clustered() -> Result<()> {
let blob = write_with(&mut depth_first_source()?, &[("reorder", "true")]).await?;
let header = read_header(&blob)?;
assert!(
header.clustered,
"a reordered archive is physically in tile-id order, so it may say so"
);
assert!(header.addressed_tiles_count > 0);
Ok(())
}
#[tokio::test]
async fn reorder_preserves_deduplication() -> Result<()> {
let unclustered = write_with(&mut depth_first_source()?, &[("allow_unclustered", "true")]).await?;
let reordered = write_with(&mut depth_first_source()?, &[("reorder", "true")]).await?;
let a = read_header(&unclustered)?;
let b = read_header(&reordered)?;
assert_eq!(a.addressed_tiles_count, b.addressed_tiles_count);
assert_eq!(
a.tile_contents_count, b.tile_contents_count,
"reordering must not change how many distinct blobs are stored"
);
assert_eq!(a.tile_entries_count, b.tile_entries_count);
assert_eq!(
unclustered.len(),
reordered.len(),
"expanding dedup would make the reordered archive much larger"
);
Ok(())
}
#[tokio::test]
async fn reorder_and_allow_unclustered_together_are_refused() -> Result<()> {
let err = write_with(
&mut depth_first_source()?,
&[("reorder", "true"), ("allow_unclustered", "true")],
)
.await
.unwrap_err();
let err = format!("{err:#}");
assert!(err.contains("ask for different things"), "unexpected: {err}");
assert!(err.contains("Pick one"), "unexpected: {err}");
Ok(())
}
#[tokio::test]
async fn reorder_is_skipped_when_the_source_is_already_ordered() -> Result<()> {
let mut reader = MockReader::new_mock(
TilePyramid::new_full_up_to(3),
TileSourceMetadata::new(TileFormat::PNG, TileCompression::Uncompressed, Traversal::ANY, None),
)?;
let with_reorder = write_with(&mut reader, &[("reorder", "true")]).await?;
let mut reader = MockReader::new_mock(
TilePyramid::new_full_up_to(3),
TileSourceMetadata::new(TileFormat::PNG, TileCompression::Uncompressed, Traversal::ANY, None),
)?;
let without = write_with(&mut reader, &[]).await?;
assert!(read_header(&with_reorder)?.clustered);
assert_eq!(with_reorder.len(), without.len(), "the extra pass should not have run");
Ok(())
}
#[tokio::test]
async fn reorder_false_still_refuses_an_unordered_source() -> Result<()> {
assert!(
write_with(&mut depth_first_source()?, &[("reorder", "false")])
.await
.is_err()
);
Ok(())
}
#[tokio::test]
async fn an_unusable_temp_dir_names_itself() -> Result<()> {
let err = write_with(
&mut depth_first_source()?,
&[("reorder", "true"), ("temp_dir", "/nonexistent-versatiles-test")],
)
.await
.unwrap_err();
let err = format!("{err:#}");
assert!(err.contains("/nonexistent-versatiles-test"), "unexpected: {err}");
Ok(())
}
#[tokio::test]
async fn reorder_leaves_no_temporary_behind() -> Result<()> {
let dir = assert_fs::TempDir::new()?;
let runtime = TilesRuntime::builder()
.silent_progress(true)
.writer_option("reorder", "true")
.writer_option("temp_dir", dir.path().to_str().unwrap())
.build();
write_to_blob(&mut depth_first_source()?, runtime).await?;
let leftovers: Vec<_> = std::fs::read_dir(dir.path())?.collect();
assert!(
leftovers.is_empty(),
"the temporary must be removed after a successful write"
);
Ok(())
}
#[test]
fn the_option_is_declared_so_the_registry_accepts_it() {
for option in ["allow_unclustered", "reorder", "temp_dir"] {
assert!(
PMTilesWriter::supported_options().contains(&option),
"{option} not declared"
);
}
}
#[context("test: PMTiles read↔write roundtrip")]
#[tokio::test]
async fn read_write() -> Result<()> {
let mut mock_reader = MockReader::new_mock(
TilePyramid::new_full_up_to(4),
TileSourceMetadata::new(TileFormat::MVT, TileCompression::Gzip, Traversal::ANY, None),
)?;
let runtime = TilesRuntime::default();
let mut data_writer = DataWriterBlob::new()?;
PMTilesWriter::write_to_writer(&mut mock_reader, &mut data_writer, runtime.clone()).await?;
let data_reader = DataReaderBlob::from(data_writer);
let mut reader = PMTilesReader::open_data(Box::new(data_reader), runtime).await?;
MockWriter::write(&mut reader).await?;
Ok(())
}
#[context("test: PMTiles tile ordering (Hilbert & offsets)")]
#[tokio::test]
async fn tiles_written_in_order() -> Result<()> {
let mut tile_pyramid = TilePyramid::new_empty();
tile_pyramid.insert_bbox(&TileBBox::from_min_and_max(15, 4090, 4090, 4139, 4139)?)?;
tile_pyramid.insert_bbox(&TileBBox::from_min_and_max(14, 250, 250, 260, 260)?)?;
let mut mock_reader = MockReader::new_mock(
tile_pyramid,
TileSourceMetadata::new(TileFormat::MVT, TileCompression::Uncompressed, Traversal::ANY, None),
)?;
let runtime = TilesRuntime::default();
let mut data_writer = DataWriterBlob::new()?;
PMTilesWriter::write_to_writer(&mut mock_reader, &mut data_writer, runtime.clone()).await?;
let data_reader = DataReaderBlob::from(data_writer);
let reader = PMTilesReader::open_data(Box::new(data_reader), runtime).await?;
let entries = reader.tile_entries()?;
let entries_vec = entries.iter().collect::<Vec<_>>();
let mut prev_tile_id = 0;
for (i, entry) in entries_vec.iter().enumerate() {
if i > 0 {
assert!(
entry.tile_id > prev_tile_id,
"Tile IDs are not in order: {} <= {}",
entry.tile_id,
prev_tile_id
);
}
prev_tile_id = entry.tile_id + u64::from(entry.run_length.max(1)) - 1;
}
assert_eq!(entries.tile_count(), 2621); Ok(())
}
}