use std::{fs::remove_file, path::Path, sync::Arc};
use anyhow::Result;
use async_trait::async_trait;
use futures::lock::Mutex;
use r2d2::Pool;
use r2d2_sqlite::{SqliteConnectionManager, rusqlite::params};
use versatiles_core::{Blob, TileCoord, json::JsonObject};
use versatiles_derive::context;
use crate::{TileSource, TileSourceTraverseExt, TilesRuntime, TilesWriter, Traversal};
pub struct MBTilesWriter {
pool: Pool<SqliteConnectionManager>,
}
impl MBTilesWriter {
#[context("creating MBTilesWriter for '{}'", path.display())]
fn new(path: &Path) -> Result<Self> {
if path.exists() {
remove_file(path)?;
}
let manager = SqliteConnectionManager::file(path);
let pool = Pool::builder().max_size(10).build(manager)?;
pool.get()?.execute_batch(
"CREATE TABLE metadata (name TEXT, value TEXT, UNIQUE (name));
CREATE TABLE tiles (zoom_level INTEGER, tile_column INTEGER, tile_row INTEGER, tile_data BLOB, UNIQUE (zoom_level, tile_column, tile_row));
CREATE UNIQUE INDEX tile_index on tiles (zoom_level, tile_column, tile_row);",
)?;
Ok(MBTilesWriter { pool })
}
#[context("adding {} tiles to MBTiles database", tiles.len())]
fn add_tiles(&mut self, tiles: &Vec<(TileCoord, Blob)>) -> Result<()> {
let mut conn = self.pool.get()?;
let transaction = conn.transaction()?;
for (c, blob) in tiles {
let max_index = 2u32.pow(u32::from(c.level)) - 1;
transaction.execute(
"INSERT INTO tiles (zoom_level, tile_column, tile_row, tile_data) VALUES (?1, ?2, ?3, ?4)",
params![c.level, c.x, max_index - c.y, blob.as_slice()],
)?;
}
transaction.commit()?;
Ok(())
}
#[context("setting metadata key '{}' = '{}'", name, value)]
fn set_metadata(&self, name: &str, value: &str) -> Result<()> {
self.pool.get()?.execute(
"INSERT OR REPLACE INTO metadata (name, value) VALUES (?1, ?2)",
params![name, value],
)?;
Ok(())
}
}
#[async_trait]
impl TilesWriter for MBTilesWriter {
fn supports_data_writer() -> bool {
false
}
#[context("writing MBTiles to '{}'", path.display())]
async fn write_to_path(reader: &mut dyn TileSource, path: &Path, runtime: TilesRuntime) -> Result<()> {
let writer = MBTilesWriter::new(path)?;
let metadata = reader.metadata().clone();
let (format, compression) = super::required_encoding(*metadata.tile_format())?;
if metadata.tile_compression() != &compression {
log::info!(
"recompressing {} tiles from {} to {compression}, which is what MBTiles stores",
metadata.tile_format(),
metadata.tile_compression()
);
}
writer.set_metadata("format", format)?;
writer.set_metadata("type", "baselayer")?;
writer.set_metadata("version", "3.0")?;
let tilejson = reader.tilejson();
let pyramid = reader.tile_pyramid().await.ok();
let bounds = tilejson.bounds.or_else(|| pyramid.as_ref().and_then(|p| p.geo_bbox()));
if let Some(bbox) = bounds {
writer.set_metadata(
"bounds",
&format!("{},{},{},{}", bbox.x_min, bbox.y_min, bbox.x_max, bbox.y_max),
)?;
}
let center = tilejson
.center
.or_else(|| pyramid.as_ref().and_then(|p| p.geo_center()));
if let Some(center) = center {
writer.set_metadata("center", &format!("{},{},{}", center.0, center.1, center.2))?;
}
if let Some(zoom_min) = tilejson.zoom_min() {
writer.set_metadata("minzoom", &zoom_min.to_string())?;
}
if let Some(zoom_max) = tilejson.zoom_max() {
writer.set_metadata("maxzoom", &zoom_max.to_string())?;
}
if let Some(vector_layers) = tilejson.as_object().get("vector_layers") {
writer.set_metadata(
"json",
&JsonObject::from(vec![("vector_layers", vector_layers)]).stringify(),
)?;
}
for key in ["name", "author", "type", "description", "version", "license"] {
if let Some(value) = tilejson.str(key) {
writer.set_metadata(key, value)?;
}
}
let writer_mutex = Arc::new(Mutex::new(writer));
let tile_compression = compression;
reader
.traverse_all_tiles(
&Traversal::ANY,
|_bbox, stream| {
let writer_mutex = Arc::clone(&writer_mutex);
Box::pin(async move {
let mut writer = writer_mutex.lock().await;
stream
.map_parallel_try(move |_coord, tile| tile.into_blob(&tile_compression))
.unwrap_results()
.for_each_buffered_try(4096, |v| writer.add_tiles(&v))
.await?;
Ok(())
})
},
runtime.clone(),
)
.await?;
Ok(())
}
}
#[cfg(test)]
mod tests {
use assert_fs::NamedTempFile;
use versatiles_core::{TileCompression, TileFormat, TilePyramid};
use super::*;
use crate::{MBTilesReader, MockReader, MockWriter, TileSourceMetadata};
#[tokio::test]
async fn read_write() -> Result<()> {
let mut mock_reader = MockReader::new_mock(
TilePyramid::new_full_up_to(5),
TileSourceMetadata::new(TileFormat::MVT, TileCompression::Gzip, Traversal::ANY, None),
)?;
let filename = NamedTempFile::new("temp.mbtiles")?;
MBTilesWriter::write_to_path(&mut mock_reader, &filename, TilesRuntime::default()).await?;
let mut reader = MBTilesReader::open(&filename, TilesRuntime::default())?;
MockWriter::write(&mut reader).await?;
Ok(())
}
#[tokio::test]
async fn uncompressed_mvt_is_gzipped_rather_than_refused() -> Result<()> {
let file = NamedTempFile::new("uncompressed_mvt.mbtiles")?;
let mut reader = MockReader::new_mock(
TilePyramid::new_full_up_to(2),
TileSourceMetadata::new(TileFormat::MVT, TileCompression::Uncompressed, Traversal::ANY, None),
)?;
MBTilesWriter::write_to_path(&mut reader, file.path(), TilesRuntime::default()).await?;
let written = MBTilesReader::open(file.path(), TilesRuntime::default())?;
assert_eq!(written.metadata().tile_compression(), &TileCompression::Gzip);
assert_eq!(written.metadata().tile_format(), &TileFormat::MVT);
Ok(())
}
#[test]
fn an_unsupported_format_still_fails() {
let err = super::super::required_encoding(TileFormat::JSON).unwrap_err();
assert!(
format!("{err}").contains("supports jpg, png, webp and pbf"),
"unexpected: {err}"
);
}
}