use crate::Tile;
use anyhow::Result;
use futures::stream::StreamExt;
use std::sync::Arc;
use versatiles_core::{Blob, ByteRange, TileCompression, TileCoord, TileFormat, TileStream, io::DataReader};
const MAX_CHUNK_SIZE: u64 = 256 * 1024 * 1024;
const MAX_CHUNK_GAP: u64 = 256 * 1024;
#[derive(Debug)]
pub struct Chunk {
tiles: Vec<(TileCoord, ByteRange)>,
range: ByteRange,
}
impl Chunk {
fn new(start: u64) -> Self {
Self {
tiles: Vec::new(),
range: ByteRange::new(start, 0),
}
}
fn push(&mut self, entry: (TileCoord, ByteRange)) {
assert!(
entry.1.offset >= self.range.offset,
"entry offset must be >= range offset"
);
self.range.length = self
.range
.length
.max(entry.1.offset + entry.1.length - self.range.offset);
self.tiles.push(entry);
}
async fn read(
&self,
reader: &DataReader,
tile_compression: TileCompression,
tile_format: TileFormat,
) -> Result<Vec<(TileCoord, Tile)>> {
let big_blob = reader.read_range(&self.range).await?;
Ok(self.slice_tiles(&big_blob, tile_compression, tile_format))
}
fn slice_tiles(
&self,
big_blob: &Blob,
tile_compression: TileCompression,
tile_format: TileFormat,
) -> Vec<(TileCoord, Tile)> {
let chunk_start = self.range.offset;
self
.tiles
.iter()
.map(|(coord, range)| {
let start =
usize::try_from(range.offset - chunk_start).expect("range offset difference should fit in usize");
let end = start + usize::try_from(range.length).expect("range length should fit in usize");
let blob = Blob::from(big_blob.range(start..end));
let tile = Tile::from_blob(blob, tile_compression, tile_format);
(*coord, tile)
})
.collect()
}
}
pub struct Chunks {
chunks: Vec<Chunk>,
}
impl Chunks {
fn new(chunks: Vec<Chunk>) -> Self {
Self { chunks }
}
pub fn new_empty() -> Self {
Self { chunks: Vec::new() }
}
fn coalesce(tile_ranges: &mut Vec<(TileCoord, ByteRange)>, max_size: u64, max_gap: u64) -> Chunks {
if tile_ranges.is_empty() {
return Chunks::new(Vec::new());
}
tile_ranges.sort_by_key(|e| e.1.offset);
let mut chunks: Vec<Chunk> = Vec::new();
let mut chunk = Chunk::new(tile_ranges[0].1.offset);
for entry in tile_ranges.drain(..) {
let chunk_start = chunk.range.offset;
let chunk_end = chunk.range.offset + chunk.range.length;
let tile_start = entry.1.offset;
let tile_end = entry.1.offset + entry.1.length;
if (chunk_start + max_size > tile_end) && (chunk_end + max_gap > tile_start) {
chunk.push(entry);
} else {
chunks.push(chunk);
chunk = Chunk::new(entry.1.offset);
chunk.push(entry);
}
}
if !chunk.tiles.is_empty() {
chunks.push(chunk);
}
Chunks::new(chunks)
}
pub fn from_tile_ranges(mut tile_ranges: Vec<(TileCoord, ByteRange)>) -> Chunks {
Chunks::coalesce(&mut tile_ranges, MAX_CHUNK_SIZE, MAX_CHUNK_GAP)
}
pub fn stream(
self,
reader: Arc<DataReader>,
tile_compression: TileCompression,
tile_format: TileFormat,
) -> TileStream<'static, Tile> {
TileStream::from_stream(
futures::stream::iter(self.chunks)
.then(move |chunk| {
let reader = Arc::clone(&reader);
async move {
let entries = chunk
.read(&reader, tile_compression, tile_format)
.await
.unwrap_or_else(|e| panic!("aborting to prevent corrupt output — {e:#}"));
futures::stream::iter(entries)
}
})
.flatten()
.boxed(),
)
}
}
impl IntoIterator for Chunks {
type Item = Chunk;
type IntoIter = std::vec::IntoIter<Self::Item>;
fn into_iter(self) -> Self::IntoIter {
self.chunks.into_iter()
}
}
impl FromIterator<Chunk> for Chunks {
fn from_iter<T: IntoIterator<Item = Chunk>>(iter: T) -> Self {
Chunks::new(iter.into_iter().collect())
}
}