use std::time::Instant;
use anyhow::{Result, anyhow};
use async_trait::async_trait;
use futures::{SinkExt, StreamExt, channel::mpsc, try_join};
use versatiles_core::{
compression::compress,
io::DataWriterTrait,
types::{Blob, ByteRange, TileCompression, TileCoord},
};
use super::types::{BlockBuilder, BlockIndex, FileHeader};
use crate::{TileSource, TileSourceTraverseExt, TilesRuntime, TilesWriter, Traversal};
enum BlockMessage {
Start(u8),
Tile(TileCoord, Blob),
End,
}
use versatiles_derive::context;
fn tile_buffer_size() -> usize {
std::env::var("VERSATILES_WRITE_TILE_BUFFER")
.ok()
.and_then(|s| s.trim().parse::<usize>().ok())
.filter(|&n| n > 0)
.unwrap_or(64)
}
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 (tx, mut rx) = mpsc::channel::<BlockMessage>(tile_buffer_size());
let produce = async move {
reader
.traverse_all_tiles(
&Traversal::new_any_size(256, 256)?,
move |bbox, stream| {
let mut tx = tx.clone();
Box::pin(async move {
log::trace!("start processing block at {bbox:?}");
tx.send(BlockMessage::Start(bbox.level()))
.await
.map_err(|_| anyhow!("writer stopped accepting blocks"))?;
stream
.map_parallel_try(move |coord, tile| {
let started = Instant::now();
let blob = tile.into_blob(&tile_compression)?;
let elapsed = started.elapsed();
if elapsed.as_secs_f64() > 1.0 {
log::trace!(
"writer: slow tile encode {coord:?}: {elapsed:?} -> {} bytes",
blob.len()
);
}
Ok(blob)
})
.unwrap_results()
.for_each_async_try(|coord, blob| {
let mut tx = tx.clone();
async move {
tx.send(BlockMessage::Tile(coord, blob))
.await
.map_err(|_| anyhow!("writer stopped accepting tiles"))
}
})
.await?;
tx.send(BlockMessage::End)
.await
.map_err(|_| anyhow!("writer stopped accepting blocks"))?;
Ok(())
})
},
runtime.clone(),
)
.await?;
Ok::<(), anyhow::Error>(())
};
let consume = async {
let mut block_index = BlockIndex::new_empty();
while let Some(message) = rx.next().await {
let BlockMessage::Start(level) = message else {
return Err(anyhow!("writer protocol error: expected block start"));
};
let mut block_builder = BlockBuilder::new(level, writer)?;
let mut tile_count: u64 = 0;
loop {
match rx.next().await {
Some(BlockMessage::Tile(coord, blob)) => {
block_builder.write_tile(coord, blob)?;
tile_count += 1;
}
Some(BlockMessage::End) => break,
Some(BlockMessage::Start(_)) => return Err(anyhow!("writer protocol error: nested block start")),
None => return Err(anyhow!("producer dropped mid-block")),
}
}
if let Some(block) = block_builder.finalize()? {
log::trace!("writer: wrote block ({tile_count} tiles) to output");
block_index.insert_block(block);
}
}
log::trace!("writer: input drained, finalizing block index");
Ok::<BlockIndex, anyhow::Error>(block_index)
};
let ((), block_index) = try_join!(produce, consume)?;
let range = writer.append(&block_index.to_brotli_blob()?)?;
Ok(range)
}
}