use crate::{SharedTileSource, SourceType, Tile, TileSource, TileSourceMetadata, TilesRuntime};
use anyhow::{Result, bail};
use async_trait::async_trait;
use std::{path::Path, sync::Arc};
use versatiles_core::{GeoBBox, TileBBox, TileCompression, TileCoord, TileFormat, TileJSON, TilePyramid, TileStream};
use versatiles_derive::context;
#[derive(Debug)]
pub struct TilesConverterParameters {
pub tile_pyramid: Option<TilePyramid>,
pub geo_bbox: Option<GeoBBox>,
pub tile_compression: Option<TileCompression>,
pub tile_format: Option<TileFormat>,
pub format_quality: Option<u8>,
pub format_effort: Option<u8>,
pub flip_y: bool,
pub swap_xy: bool,
}
impl Default for TilesConverterParameters {
fn default() -> Self {
TilesConverterParameters {
tile_pyramid: None,
geo_bbox: None,
tile_compression: None,
tile_format: None,
format_quality: None,
format_effort: None,
flip_y: false,
swap_xy: false,
}
}
}
#[context("Converting tiles from reader to file")]
pub async fn convert_tiles_container(
reader: SharedTileSource,
cp: TilesConverterParameters,
path: &Path,
runtime: TilesRuntime,
) -> Result<()> {
runtime.events().step("Starting conversion".to_string());
let converter = TilesConvertReader::new_from_reader(reader, cp).await?;
runtime.write_to_path(converter.into_shared(), path).await?;
if runtime.had_errors() {
bail!(
"conversion completed with {} read error(s) — output may be incomplete",
runtime.error_count()
);
}
runtime.events().step("Conversion complete".to_string());
Ok(())
}
#[context("Converting tiles from reader to destination")]
pub async fn convert_tiles_container_to_str(
reader: SharedTileSource,
cp: TilesConverterParameters,
destination: &str,
runtime: TilesRuntime,
) -> Result<()> {
runtime.events().step("Starting conversion".to_string());
let converter = TilesConvertReader::new_from_reader(reader, cp).await?;
runtime.write_to_str(converter.into_shared(), destination).await?;
if runtime.had_errors() {
bail!(
"conversion completed with {} read error(s) — output may be incomplete",
runtime.error_count()
);
}
runtime.events().step("Conversion complete".to_string());
Ok(())
}
#[derive(Debug)]
pub struct TilesConvertReader {
reader: SharedTileSource,
converter_parameters: TilesConverterParameters,
reader_metadata: TileSourceMetadata,
tilejson: TileJSON,
tile_pyramid: Arc<TilePyramid>,
}
impl TilesConvertReader {
#[context("Creating converter reader from existing reader")]
pub async fn new_from_reader(reader: SharedTileSource, cp: TilesConverterParameters) -> Result<TilesConvertReader> {
let mut tile_pyramid = reader.tile_pyramid().await?.as_ref().clone();
if cp.flip_y {
tile_pyramid.flip_y();
}
if cp.swap_xy {
tile_pyramid.swap_xy();
}
if let Some(filter_pyramid) = &cp.tile_pyramid {
tile_pyramid.intersect_pyramid(filter_pyramid);
}
let mut new_rp: TileSourceMetadata = reader.metadata().clone();
if let Some(tile_format) = cp.tile_format {
new_rp.set_tile_format(tile_format);
}
if let Some(tile_compression) = cp.tile_compression {
new_rp.set_tile_compression(tile_compression);
}
new_rp.set_tile_pyramid(tile_pyramid.clone());
let mut tilejson = reader.tilejson().clone();
if let Some(ref geo_bbox) = cp.geo_bbox {
if let Some(ref mut bounds) = tilejson.bounds {
bounds.intersect(geo_bbox);
} else {
tilejson.bounds = Some(*geo_bbox);
}
tilejson.center = None;
}
new_rp.update_tilejson(&mut tilejson);
Ok(TilesConvertReader {
reader,
converter_parameters: cp,
reader_metadata: new_rp,
tilejson,
tile_pyramid: Arc::new(tile_pyramid),
})
}
}
#[async_trait]
impl TileSource for TilesConvertReader {
fn source_type(&self) -> Arc<SourceType> {
SourceType::new_processor("TilesConvertReader", self.reader.source_type())
}
fn metadata(&self) -> &TileSourceMetadata {
&self.reader_metadata
}
fn tilejson(&self) -> &TileJSON {
&self.tilejson
}
async fn tile(&self, coord: &TileCoord) -> Result<Option<Tile>> {
let mut coord = *coord;
if self.converter_parameters.flip_y {
coord.flip_y();
}
if self.converter_parameters.swap_xy {
coord.swap_xy();
}
let tile = self.reader.tile(&coord).await?;
let Some(mut tile) = tile else { return Ok(None) };
if let Some(tile_format) = self.converter_parameters.tile_format {
tile.change_format(
tile_format,
self.converter_parameters.format_quality,
self.converter_parameters.format_effort,
)?;
}
if let Some(compression) = self.converter_parameters.tile_compression {
tile.change_compression(&compression)?;
}
Ok(Some(tile))
}
async fn tile_coord_stream(&self, mut bbox: TileBBox) -> Result<TileStream<'static, ()>> {
if self.converter_parameters.swap_xy {
bbox.swap_xy();
}
if self.converter_parameters.flip_y {
bbox.flip_y();
}
let mut stream = self.reader.tile_coord_stream(bbox).await?;
let flip_y = self.converter_parameters.flip_y;
let swap_xy = self.converter_parameters.swap_xy;
if flip_y || swap_xy {
stream = stream.map_coord(move |mut coord| {
if flip_y {
coord.flip_y();
}
if swap_xy {
coord.swap_xy();
}
coord
});
}
Ok(stream)
}
async fn tile_stream(&self, mut bbox: TileBBox) -> Result<TileStream<'static, Tile>> {
log::trace!("converter::tile_stream {bbox:?}");
if self.converter_parameters.swap_xy {
bbox.swap_xy();
}
if self.converter_parameters.flip_y {
bbox.flip_y();
}
let mut stream = self.reader.tile_stream(bbox).await?;
let flip_y = self.converter_parameters.flip_y;
let swap_xy = self.converter_parameters.swap_xy;
if flip_y || swap_xy {
stream = stream.map_coord(move |mut coord| {
if flip_y {
coord.flip_y();
}
if swap_xy {
coord.swap_xy();
}
coord
});
}
if let Some(tile_format) = self.converter_parameters.tile_format {
let quality = self.converter_parameters.format_quality;
let effort = self.converter_parameters.format_effort;
stream = stream
.map_parallel_try(move |_coord, mut tile| {
tile.change_format(tile_format, quality, effort)?;
Ok(tile)
})
.unwrap_results();
}
if let Some(tile_compression) = self.converter_parameters.tile_compression {
stream = stream
.map_parallel_try(move |_coord, mut tile| {
tile.change_compression(&tile_compression)?;
Ok(tile)
})
.unwrap_results();
}
Ok(stream)
}
async fn tile_pyramid(&self) -> Result<Arc<TilePyramid>> {
Ok(self.tile_pyramid.clone())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{MockReader, Traversal, VersaTilesReader};
use assert_fs::NamedTempFile;
use rstest::rstest;
use versatiles_core::{
GeoBBox, TileBBox,
TileCompression::*,
TileFormat::{self, *},
TilePyramid,
};
fn new_pyramid(b: [u32; 4]) -> TilePyramid {
TilePyramid::from([TileBBox::from_min_and_max(3, b[0], b[1], b[2], b[3]).unwrap()].as_slice())
}
fn get_mock_reader(tf: TileFormat, tc: TileCompression) -> SharedTileSource {
let tile_pyramid = TilePyramid::new_full_up_to(4);
let reader_metadata = TileSourceMetadata::new(tf, tc, Traversal::ANY, None);
MockReader::new_mock(tile_pyramid, reader_metadata)
.unwrap()
.into_shared()
}
#[rstest]
#[case(false, false, [2, 3, 4, 5], "23 33 43 24 34 25 35 44 45")]
#[case(false, true, [2, 3, 5, 4], "32 33 34 35 42 43 44 45")]
#[case(true, false, [2, 3, 4, 6], "24 34 44 23 33 22 32 21 31 43 42 41")]
#[case(true, true, [2, 3, 6, 4], "35 34 33 32 31 45 44 43 42 41")]
#[tokio::test]
async fn bbox_and_tile_order(
#[case] flip_y: bool,
#[case] swap_xy: bool,
#[case] bbox_out: [u32; 4],
#[case] tile_list: &str,
) -> Result<()> {
let pyramid_in = new_pyramid([0, 1, 4, 5]);
let pyramid_convert = new_pyramid([2, 3, 7, 7]);
let pyramid_out = new_pyramid(bbox_out);
let reader_metadata = TileSourceMetadata::new(JSON, Uncompressed, Traversal::ANY, None);
let reader = MockReader::new_mock(pyramid_in, reader_metadata)?.into_shared();
let temp_file = NamedTempFile::new("test.versatiles")?;
let runtime = TilesRuntime::default();
let cp = TilesConverterParameters {
tile_pyramid: Some(pyramid_convert),
geo_bbox: None,
flip_y,
swap_xy,
tile_compression: None,
..Default::default()
};
convert_tiles_container(reader, cp, &temp_file, runtime.clone()).await?;
let reader_out = VersaTilesReader::open(&temp_file, runtime).await?;
let parameters_out = reader_out.metadata();
let tile_compression_out = parameters_out.tile_compression();
assert_eq!(reader_out.tile_pyramid().await?.as_ref(), &pyramid_out);
let bbox = pyramid_out.level_ref(3).to_bbox();
let mut tiles: Vec<String> = Vec::new();
for coord in bbox.iter_coords_zorder() {
let mut text = reader_out
.tile(&coord)
.await?
.unwrap()
.into_blob(tile_compression_out)?
.to_string();
text = text
.replace("{\"z\":3,\"x\":", "")
.replace(",\"y\":", "")
.replace('}', "");
tiles.push(text);
}
let tiles = tiles.join(" ");
assert_eq!(tiles, tile_list);
Ok(())
}
#[test]
fn test_tiles_converter_parameters_new() {
let cp = TilesConverterParameters {
tile_pyramid: Some(TilePyramid::new_full_up_to(1)),
geo_bbox: None,
flip_y: true,
swap_xy: true,
tile_compression: None,
..Default::default()
};
assert!(cp.tile_pyramid.is_some());
assert!(cp.flip_y);
assert!(cp.swap_xy);
}
#[test]
fn test_tiles_converter_parameters_default() {
let cp = TilesConverterParameters::default();
assert_eq!(cp.tile_pyramid, None);
assert!(!cp.flip_y);
assert!(!cp.swap_xy);
}
#[tokio::test]
async fn test_tiles_convert_reader_new_from_reader() -> Result<()> {
let reader = get_mock_reader(MVT, Uncompressed);
let cp = TilesConverterParameters::default();
let tcr = TilesConvertReader::new_from_reader(reader, cp).await?;
assert_eq!(tcr.reader.source_type().to_string(), "container 'dummy' ('dummy')");
assert_eq!(tcr.source_type().to_string(), "processor 'TilesConvertReader'");
Ok(())
}
#[tokio::test]
async fn test_tile() -> Result<()> {
let reader = get_mock_reader(MVT, Uncompressed);
let cp = TilesConverterParameters::default();
let tcr = TilesConvertReader::new_from_reader(reader, cp).await?;
let coord = TileCoord::new(0, 0, 0)?;
let data = tcr.tile(&coord).await?;
assert!(data.is_some());
Ok(())
}
#[tokio::test]
async fn test_flip_y_and_swap_xy() -> Result<()> {
let reader = get_mock_reader(MVT, Uncompressed);
let cp = TilesConverterParameters {
flip_y: true,
swap_xy: true,
..Default::default()
};
let tcr = TilesConvertReader::new_from_reader(reader, cp).await?;
let mut coord = TileCoord::new(4, 5, 6)?;
let data = tcr.tile(&coord).await?;
assert!(data.is_some());
coord.flip_y();
coord.swap_xy();
let data_flipped = tcr.tile(&coord).await?;
assert_eq!(data, data_flipped);
Ok(())
}
#[tokio::test]
async fn test_geo_bbox_intersects_existing_tilejson_bounds() -> Result<()> {
let source_bounds = GeoBBox::new(10.0, 50.0, 15.0, 55.0)?;
let tile_pyramid = TilePyramid::new_full_up_to(4);
let reader_metadata = TileSourceMetadata::new(MVT, Uncompressed, Traversal::ANY, None);
let mut reader = MockReader::new_mock(tile_pyramid, reader_metadata)?;
reader.tilejson_mut().bounds = Some(source_bounds);
let reader = reader.into_shared();
let filter_bbox = GeoBBox::new(12.0, 52.0, 20.0, 60.0)?;
let mut filter_pyramid = TilePyramid::new_full();
filter_pyramid.intersect_geo_bbox(&filter_bbox)?;
let cp = TilesConverterParameters {
tile_pyramid: Some(filter_pyramid),
geo_bbox: Some(filter_bbox),
..Default::default()
};
let tcr = TilesConvertReader::new_from_reader(reader, cp).await?;
let bounds = tcr.tilejson().bounds.ok_or_else(|| anyhow::anyhow!("missing bounds"))?;
assert_eq!(bounds.as_tuple(), (12.0, 52.0, 15.0, 55.0));
Ok(())
}
#[tokio::test]
async fn test_geo_bbox_without_existing_bounds_uses_pyramid() -> Result<()> {
let tile_pyramid = TilePyramid::new_full_up_to(4);
let reader_metadata = TileSourceMetadata::new(MVT, Uncompressed, Traversal::ANY, None);
let reader = MockReader::new_mock(tile_pyramid, reader_metadata)?.into_shared();
let filter_bbox = GeoBBox::new(12.0, 52.0, 14.0, 54.0)?;
let mut filter_pyramid = TilePyramid::new_full();
filter_pyramid.intersect_geo_bbox(&filter_bbox)?;
let cp = TilesConverterParameters {
tile_pyramid: Some(filter_pyramid),
geo_bbox: Some(filter_bbox),
..Default::default()
};
let tcr = TilesConvertReader::new_from_reader(reader, cp).await?;
assert!(tcr.tilejson().bounds.is_some());
Ok(())
}
#[tokio::test]
async fn test_format_conversion() -> Result<()> {
let reader = get_mock_reader(PNG, Uncompressed);
let cp = TilesConverterParameters {
tile_format: Some(WEBP),
format_quality: Some(80),
..Default::default()
};
let tcr = TilesConvertReader::new_from_reader(reader, cp).await?;
assert_eq!(tcr.metadata().tile_format(), &WEBP);
assert_eq!(tcr.metadata().tile_compression(), &Uncompressed);
let coord = TileCoord::new(0, 0, 0)?;
let tile = tcr.tile(&coord).await?.unwrap();
assert_eq!(tile.format(), WEBP);
Ok(())
}
#[tokio::test]
async fn test_format_conversion_with_compression() -> Result<()> {
let reader = get_mock_reader(PNG, Uncompressed);
let cp = TilesConverterParameters {
tile_format: Some(WEBP),
format_quality: Some(80),
tile_compression: Some(Gzip),
..Default::default()
};
let tcr = TilesConvertReader::new_from_reader(reader, cp).await?;
assert_eq!(tcr.metadata().tile_format(), &WEBP);
assert_eq!(tcr.metadata().tile_compression(), &Gzip);
let coord = TileCoord::new(0, 0, 0)?;
let tile = tcr.tile(&coord).await?.unwrap();
assert_eq!(tile.format(), WEBP);
Ok(())
}
#[tokio::test]
async fn swap_xy_does_not_corrupt_source_metadata() -> Result<()> {
use crate::MBTilesReader;
let path = std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.parent()
.unwrap()
.join("testdata/berlin.mbtiles");
let runtime = TilesRuntime::default();
let reader = MBTilesReader::open(&path, runtime.clone())?;
let shared = reader.into_shared();
let before = shared.tile_pyramid().await?.as_ref().clone();
let conv = TilesConvertReader::new_from_reader(
shared.clone(),
TilesConverterParameters {
swap_xy: true,
..Default::default()
},
)
.await?;
let after = shared.tile_pyramid().await?.as_ref().clone();
assert_eq!(before, after, "source pyramid was mutated by the converter");
let pyr = conv.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(&conv, bbox).await?;
}
Ok(())
}
#[tokio::test]
async fn converter_satisfies_stream_count_invariant() -> Result<()> {
use crate::MBTilesReader;
let path = std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.parent()
.unwrap()
.join("testdata/berlin.mbtiles");
let runtime = TilesRuntime::default();
let shared = MBTilesReader::open(&path, runtime.clone())?.into_shared();
let conv = TilesConvertReader::new_from_reader(shared, TilesConverterParameters::default()).await?;
let pyr = conv.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(&conv, bbox).await?;
}
Ok(())
}
}