use super::types::{BlockBuilder, BlockIndex, FileHeader};
use crate::{TileSource, TileSourceTraverseExt, TilesRuntime, TilesWriter, Traversal};
use anyhow::{Result, anyhow};
use async_trait::async_trait;
use futures::lock::Mutex;
use std::sync::Arc;
use versatiles_core::{
compression::compress,
io::DataWriterTrait,
types::{Blob, ByteRange, TileCompression},
};
use versatiles_derive::context;
pub struct VersaTilesWriter {}
#[async_trait]
impl TilesWriter for VersaTilesWriter {
#[context("writing VersaTiles to DataWriter")]
async fn write_to_writer(
reader: &mut dyn TileSource,
writer: &mut dyn DataWriterTrait,
runtime: TilesRuntime,
) -> Result<()> {
let parameters = reader.metadata();
log::trace!("convert_from - reader.parameters: {parameters:?}");
let tile_compression = *parameters.tile_compression();
let tile_pyramid = reader.tile_pyramid().await?;
log::trace!("convert_from - tile_pyramid: {tile_pyramid:#}");
let tilejson = reader.tilejson();
let zoom_min = tilejson
.zoom_min()
.or(tile_pyramid.level_min())
.ok_or(anyhow!("invalid minzoom"))?;
let zoom_max = tilejson
.zoom_max()
.or(tile_pyramid.level_max())
.ok_or(anyhow!("invalid maxzoom"))?;
let bbox = tilejson
.bounds
.or(tile_pyramid.geo_bbox())
.ok_or(anyhow!("invalid geo bounding box"))?;
let mut header = FileHeader::new(*parameters.tile_format(), tile_compression, [zoom_min, zoom_max], &bbox)?;
let blob: Blob = header.to_blob()?;
log::trace!("write header");
writer.append(&blob)?;
log::trace!("write meta");
header.meta_range = Self::write_meta(reader, writer, tile_compression).await?;
log::trace!("write blocks");
header.blocks_range = Self::write_blocks(reader, writer, tile_compression, runtime).await?;
log::trace!("update header");
let blob: Blob = header.to_blob()?;
writer.write_start(&blob)?;
Ok(())
}
}
impl VersaTilesWriter {
#[context("Failed to write metadata")]
async fn write_meta(
reader: &dyn TileSource,
writer: &mut dyn DataWriterTrait,
compression: TileCompression,
) -> Result<ByteRange> {
let meta: Blob = reader.tilejson().into();
let compressed = compress(meta, &compression)?;
writer.append(&compressed)
}
#[context("Failed to write blocks")]
async fn write_blocks(
reader: &mut dyn TileSource,
writer: &mut dyn DataWriterTrait,
tile_compression: TileCompression,
runtime: TilesRuntime,
) -> Result<ByteRange> {
if reader.tile_pyramid().await?.is_empty() {
return Ok(ByteRange::empty());
}
let block_index_mutex = Arc::new(Mutex::new(BlockIndex::new_empty()));
let writer_mutex = Arc::new(Mutex::new(writer));
reader
.traverse_all_tiles(
&Traversal::new_any_size(256, 256)?,
|bbox, stream| {
let writer_mutex = Arc::clone(&writer_mutex);
let block_index_mutex = Arc::clone(&block_index_mutex);
Box::pin(async move {
log::trace!("start processing block at {bbox:?}");
let compressed_stream = stream
.map_parallel_try(move |_coord, tile| tile.into_blob(&tile_compression))
.unwrap_results();
let mut writer = writer_mutex.lock().await;
let mut block_builder = BlockBuilder::new(bbox.level(), &mut **writer)?;
compressed_stream
.for_each_try(|coord, blob| block_builder.write_tile(coord, blob))
.await?;
if let Some(block) = block_builder.finalize()? {
log::trace!("finish block {block:?}");
block_index_mutex.lock().await.insert_block(block);
} else {
log::trace!("skipping empty block at {bbox:?}");
}
Ok(())
})
},
runtime.clone(),
)
.await?;
let range = writer_mutex
.lock()
.await
.append(&block_index_mutex.lock().await.to_brotli_blob()?)?;
Ok(range)
}
}