use std::{collections::HashMap, fmt::Debug, io::Read, path::Path, sync::Arc};
use anyhow::{Result, anyhow, ensure};
use async_trait::async_trait;
use tar::{Archive, EntryType};
#[cfg(feature = "cli")]
use versatiles_core::utils::PrettyPrint;
use versatiles_core::{
Blob, ByteRange, TileBBox, TileCompression, TileCoord, TileFormat, TileJSON, TilePyramid, TileStream,
compression::decompress,
io::{DataReaderFile, DataReaderTrait},
};
use versatiles_derive::context;
use crate::{
ProgressHandle, SharedTileSource, SourceType, Tile, TileSource, TileSourceMetadata, TilesReader, TilesRuntime,
Traversal,
};
pub struct TarTilesReader {
tilejson: TileJSON,
name: String,
reader: Arc<DataReaderFile>,
tile_map: Arc<HashMap<TileCoord, ByteRange>>,
metadata: TileSourceMetadata,
}
struct ParsedTile {
coord: TileCoord,
format: TileFormat,
compression: TileCompression,
offset: u64,
length: u64,
}
fn try_parse_tile_path<R: std::io::Read>(entry: &tar::Entry<R>, path_vec: &[&str]) -> Result<Option<ParsedTile>> {
if path_vec.len() != 3 {
return Ok(None);
}
let level = path_vec[0].parse::<u8>()?;
let x = path_vec[1].parse::<u32>()?;
let mut filename = String::from(path_vec[2]);
let compression = TileCompression::from_filename(&mut filename);
let Some(format) = TileFormat::from_filename(&mut filename) else {
return Ok(None);
};
let y = filename.parse::<u32>()?;
let coord = TileCoord::new(level, x, y)?;
Ok(Some(ParsedTile {
coord,
format,
compression,
offset: entry.raw_file_position(),
length: entry.size(),
}))
}
impl TarTilesReader {
#[context("opening tar from path '{}'", path.display())]
pub fn open(path: &Path) -> Result<TarTilesReader> {
Self::open_with_progress(path, None)
}
#[context("opening tar from path '{}'", path.display())]
#[allow(clippy::too_many_lines)]
fn open_with_progress(path: &Path, progress: Option<&ProgressHandle>) -> Result<TarTilesReader> {
let mut reader = DataReaderFile::open(path)?;
if let Some(progress) = &progress
&& let Ok(file_meta) = std::fs::metadata(path)
{
progress.set_max_value(file_meta.len());
}
let mut archive = Archive::new(&mut reader);
let mut tilejson = TileJSON::default();
let mut tile_map = HashMap::new();
let mut tile_format: Option<TileFormat> = None;
let mut tile_compression: Option<TileCompression> = None;
for entry in archive.entries()? {
let mut entry = entry?;
let header = entry.header();
if header.entry_type() != EntryType::Regular {
continue;
}
if let Some(progress) = &progress {
progress.set_position(entry.raw_file_position());
}
let path = entry.path()?.clone();
let Some(mut path_tmp) = path.iter().map(std::ffi::OsStr::to_str).collect::<Option<Vec<&str>>>() else {
log::debug!("skipping tar entry with a non-UTF-8 name: {:?}", path.as_os_str());
continue;
};
if path_tmp.is_empty() {
log::debug!("skipping tar entry with an empty name");
continue;
}
if path_tmp[0] == "." {
path_tmp.remove(0);
}
let path_tmp_string = path_tmp.join("/");
drop(path);
let path_vec: Vec<&str> = path_tmp_string.split('/').collect();
if let Some(tile) = try_parse_tile_path(&entry, &path_vec)? {
if let Some(f) = &tile_format {
ensure!(
f == &tile.format,
"mixed tile formats in tar, found both {f:?} and {:?}",
tile.format
);
} else {
tile_format = Some(tile.format);
}
if let Some(c) = &tile_compression {
ensure!(
c == &tile.compression,
"mixed tile compressions in tar, found both {c:?} and {:?}",
tile.compression
);
} else {
tile_compression = Some(tile.compression);
}
tile_map.insert(
tile.coord,
ByteRange {
offset: tile.offset,
length: tile.length,
},
);
continue;
}
let mut read_to_end = || -> Result<Blob> {
let mut blob: Vec<u8> = Vec::new();
entry
.read_to_end(&mut blob)
.with_context(|| format!("reading tar entry {path_tmp_string:?}"))?;
Ok(Blob::from(blob))
};
if path_vec.len() == 1 {
match path_vec[0] {
"meta.json" | "tiles.json" | "metadata.json" => {
tilejson.merge(&TileJSON::try_from_blob_or_default(&read_to_end()?))?;
continue;
}
"meta.json.gz" | "tiles.json.gz" | "metadata.json.gz" => {
tilejson.merge(&TileJSON::try_from_blob_or_default(&decompress(
read_to_end()?,
&TileCompression::Gzip,
)?))?;
continue;
}
"meta.json.br" | "tiles.json.br" | "metadata.json.br" => {
tilejson.merge(&TileJSON::try_from_blob_or_default(&decompress(
read_to_end()?,
&TileCompression::Brotli,
)?))?;
continue;
}
&_ => {}
}
}
log::warn!("unknown file in tar: {path_tmp_string:?}");
}
if let Some(progress) = &progress {
progress.finish();
}
if tile_map.is_empty() {
return Err(anyhow!("no tiles found in tar"));
}
let metadata = TileSourceMetadata::new(
tile_format.ok_or(anyhow!("unknown tile format, can't detect format"))?,
tile_compression.ok_or(anyhow!("unknown tile compression, can't detect compression"))?,
Traversal::ANY,
None,
);
Ok(TarTilesReader {
tilejson,
name: path.to_string_lossy().into_owned(),
metadata,
reader: Arc::new(*reader),
tile_map: Arc::new(tile_map),
})
}
async fn lookup_tile(
coord: &TileCoord,
tile_map: &HashMap<TileCoord, ByteRange>,
reader: &DataReaderFile,
tile_compression: TileCompression,
tile_format: TileFormat,
) -> Result<Option<Tile>> {
if let Some(range) = tile_map.get(coord) {
let blob = reader.read_range(range).await?;
Ok(Some(Tile::from_blob(blob, tile_compression, tile_format)))
} else {
Ok(None)
}
}
}
#[async_trait]
impl TilesReader for TarTilesReader {
fn supports_data_reader() -> bool {
false
}
async fn open_path(path: &Path, runtime: TilesRuntime) -> Result<SharedTileSource> {
let progress = runtime.create_progress("scanning tar", 0);
Ok(Self::open_with_progress(path, Some(&progress))?.into_shared())
}
}
#[async_trait]
impl TileSource for TarTilesReader {
fn source_type(&self) -> Arc<SourceType> {
SourceType::new_container("tar", &self.name)
}
fn metadata(&self) -> &TileSourceMetadata {
&self.metadata
}
fn tilejson(&self) -> &TileJSON {
&self.tilejson
}
async fn tile_pyramid(&self) -> Result<Arc<TilePyramid>> {
self
.metadata
.get_or_compute_tile_pyramid(|| Ok(TilePyramid::from_tile_coords(self.tile_map.keys().copied())))
}
#[context("getting tile {:?}", coord)]
async fn tile(&self, coord: &TileCoord) -> Result<Option<Tile>> {
log::trace!("tile {coord:?}");
Self::lookup_tile(
coord,
&self.tile_map,
&self.reader,
*self.metadata.tile_compression(),
*self.metadata.tile_format(),
)
.await
}
async fn tile_stream(&self, bbox: TileBBox) -> Result<TileStream<'static, Tile>> {
log::trace!("tar::tile_stream {bbox:?}");
let bbox = self.metadata.intersection_bbox(&bbox);
let reader = Arc::clone(&self.reader);
let tile_map = Arc::clone(&self.tile_map);
let tile_compression = *self.metadata.tile_compression();
let tile_format = *self.metadata.tile_format();
Ok(TileStream::from_bbox_async_parallel(bbox, move |coord| {
let reader = Arc::clone(&reader);
let tile_map = Arc::clone(&tile_map);
async move {
let tile = TarTilesReader::lookup_tile(&coord, &tile_map, &reader, tile_compression, tile_format)
.await
.ok()??;
Some((coord, tile))
}
}))
}
async fn tile_coord_stream(&self, bbox: TileBBox) -> Result<TileStream<'static, ()>> {
let bbox = self.metadata.intersection_bbox(&bbox);
let tile_map = Arc::clone(&self.tile_map);
Ok(TileStream::from_bbox_parallel(bbox, move |coord| {
tile_map.get(&coord).map(|_| ())
}))
}
async fn tile_size_stream(&self, bbox: TileBBox) -> Result<TileStream<'static, u32>> {
let bbox = self.metadata.intersection_bbox(&bbox);
let tile_map = Arc::clone(&self.tile_map);
Ok(TileStream::from_bbox_parallel(bbox, move |coord| {
let range = tile_map.get(&coord)?;
u32::try_from(range.length).ok()
}))
}
#[cfg(feature = "cli")]
async fn probe_container(&self, print: &mut PrettyPrint, _runtime: &TilesRuntime) -> Result<()> {
print.add_key_value("tile count", &self.tile_map.len()).await;
let total_size: u64 = self.tile_map.values().map(|r| r.length).sum();
print.add_key_value("total stored size", &total_size).await;
Ok(())
}
}
impl Debug for TarTilesReader {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("TarTilesReader")
.field("parameters", &self.metadata())
.finish()
}
}
#[cfg(test)]
pub mod tests {
use versatiles_core::assert_wildcard;
#[cfg(feature = "cli")]
use versatiles_core::utils::PrettyPrint;
use super::*;
use crate::{MOCK_BYTES_PBF, MockWriter, make_test_file};
#[cfg(unix)]
#[tokio::test]
async fn hostile_entry_names_are_skipped_not_fatal() -> Result<()> {
use std::{ffi::OsStr, os::unix::ffi::OsStrExt};
use versatiles_core::compression::compress_gzip;
let dir = assert_fs::TempDir::new()?;
let path = dir.path().join("hostile.tar");
let tile = compress_gzip(&Blob::from(MOCK_BYTES_PBF.to_vec()))?;
let mut builder = tar::Builder::new(std::fs::File::create(&path)?);
let mut append = |name: &OsStr, data: &[u8]| -> Result<()> {
let mut header = tar::Header::new_gnu();
header.set_size(data.len() as u64);
header.set_mode(0o644);
header.set_cksum();
builder.append_data(&mut header, Path::new(name), data)?;
Ok(())
};
append(OsStr::new("5/3/4.pbf.gz"), tile.as_slice())?;
append(OsStr::from_bytes(b"5/3/\xff\xfe.pbf.gz"), tile.as_slice())?;
builder.finish()?;
drop(builder);
let reader = TarTilesReader::open(&path)?;
let coord = TileCoord::new(5, 3, 4)?;
assert!(reader.tile_map.contains_key(&coord), "the valid tile was lost");
assert_eq!(reader.tile_map.len(), 1, "a skipped entry was counted");
Ok(())
}
#[tokio::test]
async fn reader() -> Result<()> {
let temp_file = make_test_file(TileFormat::MVT, TileCompression::Gzip, 3, "tar").await?;
let reader = TarTilesReader::open(&temp_file)?;
assert_eq!(
format!("{reader:?}"),
"TarTilesReader { parameters: TileSourceMetadata { tile_compression: Gzip, tile_format: MVT, traversal: Traversal(AnyOrder,full), tile_pyramid: RwLock { data: None, poisoned: false, .. } } }"
);
assert_wildcard!(reader.source_type().to_string(), "container 'tar' ('*.tar')");
assert_eq!(
reader.tilejson().stringify(),
"{\"tilejson\":\"3.0.0\",\"type\":\"dummy\"}"
);
assert_eq!(
format!("{:?}", reader.metadata()),
"TileSourceMetadata { tile_compression: Gzip, tile_format: MVT, traversal: Traversal(AnyOrder,full), tile_pyramid: RwLock { data: None, poisoned: false, .. } }"
);
assert_eq!(reader.metadata().tile_compression(), &TileCompression::Gzip);
assert_eq!(reader.metadata().tile_format(), &TileFormat::MVT);
let blob = reader
.tile(&TileCoord::new(3, 6, 2)?)
.await?
.unwrap()
.into_blob(&TileCompression::Uncompressed)?;
assert_eq!(blob.as_slice(), MOCK_BYTES_PBF);
Ok(())
}
#[tokio::test]
async fn all_compressions() -> Result<()> {
async fn test_compression(compression: TileCompression) -> Result<()> {
let temp_file = make_test_file(TileFormat::MVT, compression, 2, "tar").await?;
let mut reader = TarTilesReader::open(&temp_file)?;
MockWriter::write(&mut reader).await?;
Ok(())
}
test_compression(TileCompression::Uncompressed).await?;
test_compression(TileCompression::Gzip).await?;
test_compression(TileCompression::Brotli).await?;
Ok(())
}
#[cfg(feature = "cli")]
#[tokio::test]
async fn probe() -> Result<()> {
let temp_file = make_test_file(TileFormat::MVT, TileCompression::Gzip, 4, "tar").await?;
let reader = TarTilesReader::open(&temp_file)?;
let runtime = TilesRuntime::default();
let mut printer = PrettyPrint::new();
reader
.probe_container(&mut printer.category("container").await, &runtime)
.await?;
assert_eq!(
printer.stringify().await.split('\n').collect::<Vec<_>>(),
["container:", " tile count: 341", " total stored size: 26_257", ""]
);
Ok(())
}
#[tokio::test]
async fn empty_tar_file() -> Result<()> {
let filename = assert_fs::NamedTempFile::new("empty_tar_file.tar")?;
let file = std::fs::File::create(&filename)?;
let mut a = tar::Builder::new(file);
a.finish()?;
assert_eq!(
TarTilesReader::open(&filename)
.unwrap_err()
.chain()
.last()
.unwrap()
.to_string(),
"no tiles found in tar"
);
Ok(())
}
#[tokio::test]
async fn correct_zxy_scheme() -> Result<()> {
let filename = assert_fs::NamedTempFile::new("correct_zxy_scheme.tar")?;
let file = std::fs::File::create(&filename)?;
let mut a = tar::Builder::new(file);
let mut header = tar::Header::new_gnu();
header.set_size(6);
header.set_cksum();
a.append_data(&mut header, "3/1/2.bin", [3, 1, 4, 1, 5, 9].as_ref())?;
a.finish()?;
let reader = TarTilesReader::open(&filename)?;
assert_eq!(reader.metadata().tile_format(), &TileFormat::BIN);
assert_eq!(reader.metadata().tile_compression(), &TileCompression::Uncompressed);
assert_eq!(reader.tile_pyramid().await?.count_tiles(), 1);
assert_eq!(
reader
.tile(&TileCoord::new(3, 1, 2)?)
.await?
.unwrap()
.as_blob(&TileCompression::Uncompressed)?
.as_slice(),
[3, 1, 4, 1, 5, 9].as_ref()
);
Ok(())
}
#[tokio::test]
async fn tile_stream_matches_individual_reads() -> Result<()> {
let temp_file = make_test_file(TileFormat::MVT, TileCompression::Gzip, 2, "tar").await?;
let reader = TarTilesReader::open(&temp_file)?;
let bbox = TileBBox::new_full(2)?;
let stream = reader.tile_stream(bbox).await?;
let stream_tiles: Vec<_> = stream.to_vec().await;
assert_eq!(stream_tiles.len(), 16);
for (coord, mut tile) in stream_tiles {
let stream_blob = tile.as_blob(reader.metadata().tile_compression())?;
let single_blob = reader
.tile(&coord)
.await?
.expect("tile should exist")
.into_blob(reader.metadata().tile_compression())?;
assert_eq!(
stream_blob.as_slice(),
single_blob.as_slice(),
"blob mismatch at {coord:?}"
);
}
Ok(())
}
}