use crate::compression::{self, Compression};
use crate::dedup::{DeduplicationCache, DeduplicationStats, TileHasher};
use crate::tile::TileBounds;
use crate::{Error, Result};
use std::collections::{BTreeMap, HashMap};
use std::fs::File;
use std::io::{BufWriter, Write};
use std::path::Path;
const PMTILES_MAGIC: &[u8; 7] = b"PMTiles";
const PMTILES_VERSION: u8 = 3;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(u8)]
pub enum TileType {
Unknown = 0,
Mvt = 1,
Png = 2,
Jpeg = 3,
Webp = 4,
Avif = 5,
}
#[derive(Debug, Clone)]
pub struct Header {
pub root_dir_offset: u64,
pub root_dir_length: u64,
pub json_metadata_offset: u64,
pub json_metadata_length: u64,
pub leaf_dirs_offset: u64,
pub leaf_dirs_length: u64,
pub tile_data_offset: u64,
pub tile_data_length: u64,
pub addressed_tiles_count: u64,
pub tile_entries_count: u64,
pub tile_contents_count: u64,
pub clustered: bool,
pub internal_compression: Compression,
pub tile_compression: Compression,
pub tile_type: TileType,
pub min_zoom: u8,
pub max_zoom: u8,
pub min_lon: f64,
pub min_lat: f64,
pub max_lon: f64,
pub max_lat: f64,
pub center_zoom: u8,
pub center_lon: f64,
pub center_lat: f64,
}
impl Default for Header {
fn default() -> Self {
Self {
root_dir_offset: 127, root_dir_length: 0,
json_metadata_offset: 0,
json_metadata_length: 0,
leaf_dirs_offset: 0,
leaf_dirs_length: 0,
tile_data_offset: 0,
tile_data_length: 0,
addressed_tiles_count: 0,
tile_entries_count: 0,
tile_contents_count: 0,
clustered: true,
internal_compression: Compression::Gzip,
tile_compression: Compression::Gzip,
tile_type: TileType::Mvt,
min_zoom: 0,
max_zoom: 14,
min_lon: -180.0,
min_lat: -85.0,
max_lon: 180.0,
max_lat: 85.0,
center_zoom: 0,
center_lon: 0.0,
center_lat: 0.0,
}
}
}
impl Header {
pub fn to_bytes(&self) -> [u8; 127] {
let mut buf = [0u8; 127];
buf[0..7].copy_from_slice(PMTILES_MAGIC);
buf[7] = PMTILES_VERSION;
buf[8..16].copy_from_slice(&self.root_dir_offset.to_le_bytes());
buf[16..24].copy_from_slice(&self.root_dir_length.to_le_bytes());
buf[24..32].copy_from_slice(&self.json_metadata_offset.to_le_bytes());
buf[32..40].copy_from_slice(&self.json_metadata_length.to_le_bytes());
buf[40..48].copy_from_slice(&self.leaf_dirs_offset.to_le_bytes());
buf[48..56].copy_from_slice(&self.leaf_dirs_length.to_le_bytes());
buf[56..64].copy_from_slice(&self.tile_data_offset.to_le_bytes());
buf[64..72].copy_from_slice(&self.tile_data_length.to_le_bytes());
buf[72..80].copy_from_slice(&self.addressed_tiles_count.to_le_bytes());
buf[80..88].copy_from_slice(&self.tile_entries_count.to_le_bytes());
buf[88..96].copy_from_slice(&self.tile_contents_count.to_le_bytes());
buf[96] = if self.clustered { 1 } else { 0 };
buf[97] = self.internal_compression as u8;
buf[98] = self.tile_compression as u8;
buf[99] = self.tile_type as u8;
buf[100] = self.min_zoom;
buf[101] = self.max_zoom;
let encode_coord = |v: f64| -> [u8; 4] { ((v * 10_000_000.0) as i32).to_le_bytes() };
buf[102..106].copy_from_slice(&encode_coord(self.min_lon));
buf[106..110].copy_from_slice(&encode_coord(self.min_lat));
buf[110..114].copy_from_slice(&encode_coord(self.max_lon));
buf[114..118].copy_from_slice(&encode_coord(self.max_lat));
buf[118] = self.center_zoom;
buf[119..123].copy_from_slice(&encode_coord(self.center_lon));
buf[123..127].copy_from_slice(&encode_coord(self.center_lat));
buf
}
}
impl TileType {
pub fn from_code(code: u8) -> Option<Self> {
match code {
0 => Some(TileType::Unknown),
1 => Some(TileType::Mvt),
2 => Some(TileType::Png),
3 => Some(TileType::Jpeg),
4 => Some(TileType::Webp),
5 => Some(TileType::Avif),
_ => None,
}
}
}
impl Header {
pub fn from_bytes(bytes: &[u8]) -> Result<Header> {
let err = |msg: String| Error::PMTilesRead(msg);
if bytes.len() < 127 {
return Err(err(format!(
"file too short for PMTiles header: {} bytes (need 127)",
bytes.len()
)));
}
if &bytes[0..7] != PMTILES_MAGIC {
return Err(err("bad magic: not a PMTiles archive".to_string()));
}
if bytes[7] != PMTILES_VERSION {
return Err(err(format!(
"unsupported PMTiles version {} (only v3 is supported)",
bytes[7]
)));
}
let read_u64 =
|at: usize| u64::from_le_bytes(bytes[at..at + 8].try_into().expect("8-byte slice"));
let read_coord = |at: usize| {
f64::from(i32::from_le_bytes(
bytes[at..at + 4].try_into().expect("4-byte slice"),
)) / 10_000_000.0
};
let internal_compression = Compression::from_code(bytes[97])
.ok_or_else(|| err(format!("invalid internal compression code {}", bytes[97])))?;
let tile_compression = Compression::from_code(bytes[98])
.ok_or_else(|| err(format!("invalid tile compression code {}", bytes[98])))?;
let tile_type = TileType::from_code(bytes[99])
.ok_or_else(|| err(format!("invalid tile type code {}", bytes[99])))?;
Ok(Header {
root_dir_offset: read_u64(8),
root_dir_length: read_u64(16),
json_metadata_offset: read_u64(24),
json_metadata_length: read_u64(32),
leaf_dirs_offset: read_u64(40),
leaf_dirs_length: read_u64(48),
tile_data_offset: read_u64(56),
tile_data_length: read_u64(64),
addressed_tiles_count: read_u64(72),
tile_entries_count: read_u64(80),
tile_contents_count: read_u64(88),
clustered: bytes[96] == 1,
internal_compression,
tile_compression,
tile_type,
min_zoom: bytes[100],
max_zoom: bytes[101],
min_lon: read_coord(102),
min_lat: read_coord(106),
max_lon: read_coord(110),
max_lat: read_coord(114),
center_zoom: bytes[118],
center_lon: read_coord(119),
center_lat: read_coord(123),
})
}
}
pub fn tile_id(z: u8, x: u32, y: u32) -> u64 {
if z == 0 {
return 0;
}
let base_id: u64 = (1..z as u64).map(|i| 4u64.pow(i as u32)).sum();
let hilbert_idx = xy_to_hilbert(z, x, y);
base_id + hilbert_idx + 1
}
fn xy_to_hilbert(z: u8, x: u32, y: u32) -> u64 {
let n = 1u32 << z;
let mut rx: u32;
let mut ry: u32;
let mut s: u32;
let mut d: u64 = 0;
let mut x = x;
let mut y = y;
s = n / 2;
while s > 0 {
rx = if (x & s) > 0 { 1 } else { 0 };
ry = if (y & s) > 0 { 1 } else { 0 };
d += (s as u64) * (s as u64) * ((3 * rx) ^ ry) as u64;
if ry == 0 {
if rx == 1 {
x = n - 1 - x;
y = n - 1 - y;
}
std::mem::swap(&mut x, &mut y);
}
s /= 2;
}
d
}
pub fn tile_id_to_zxy(id: u64) -> Result<(u8, u32, u32)> {
let mut acc: u64 = 0;
for z in 0u8..=31 {
let num = 1u64 << (2 * u64::from(z));
if id - acc < num {
let (x, y) = hilbert_d2xy(z, id - acc);
return Ok((z, x, y));
}
acc += num;
}
Err(Error::PMTilesRead(format!(
"tile id {id} exceeds the zoom 31 address space"
)))
}
fn hilbert_d2xy(z: u8, d: u64) -> (u32, u32) {
let n = 1u64 << z;
let (mut x, mut y) = (0u64, 0u64);
let mut t = d;
let mut s = 1u64;
while s < n {
let rx = 1 & (t / 2);
let ry = 1 & (t ^ rx);
if ry == 0 {
if rx == 1 {
x = s - 1 - x;
y = s - 1 - y;
}
std::mem::swap(&mut x, &mut y);
}
x += s * rx;
y += s * ry;
t /= 4;
s *= 2;
}
(x as u32, y as u32)
}
#[derive(Debug, Clone)]
pub struct DirEntry {
pub tile_id: u64,
pub offset: u64,
pub length: u32,
pub run_length: u32, }
pub fn encode_varint(mut value: u64, buf: &mut Vec<u8>) {
while value >= 0x80 {
buf.push((value as u8) | 0x80);
value >>= 7;
}
buf.push(value as u8);
}
pub fn decode_varint(data: &[u8]) -> Option<(u64, usize)> {
let mut result: u64 = 0;
let mut shift = 0;
for (i, &byte) in data.iter().enumerate() {
result |= ((byte & 0x7f) as u64) << shift;
if byte & 0x80 == 0 {
return Some((result, i + 1));
}
shift += 7;
if shift >= 64 {
return None; }
}
None
}
pub fn encode_directory(entries: &[DirEntry]) -> Vec<u8> {
let mut buf = Vec::new();
encode_varint(entries.len() as u64, &mut buf);
if entries.is_empty() {
return buf;
}
let mut last_id = 0u64;
for entry in entries {
encode_varint(entry.tile_id - last_id, &mut buf);
last_id = entry.tile_id;
}
for entry in entries {
encode_varint(entry.run_length as u64, &mut buf);
}
for entry in entries {
encode_varint(entry.length as u64, &mut buf);
}
let mut expected_offset = 0u64;
for (i, entry) in entries.iter().enumerate() {
let is_contiguous = i > 0 && entry.offset == expected_offset;
if is_contiguous {
encode_varint(0, &mut buf);
} else {
encode_varint(entry.offset + 1, &mut buf);
}
expected_offset = entry.offset + entry.length as u64;
}
buf
}
pub fn decode_directory(data: &[u8]) -> Option<Vec<DirEntry>> {
let mut offset = 0;
let (count, consumed) = decode_varint(&data[offset..])?;
offset += consumed;
let count = count as usize;
if count == 0 {
return Some(Vec::new());
}
let mut entries = Vec::with_capacity(count);
let mut last_id = 0u64;
for _ in 0..count {
let (delta, consumed) = decode_varint(&data[offset..])?;
offset += consumed;
last_id += delta;
entries.push(DirEntry {
tile_id: last_id,
offset: 0,
length: 0,
run_length: 0,
});
}
for entry in entries.iter_mut() {
let (run_length, consumed) = decode_varint(&data[offset..])?;
offset += consumed;
entry.run_length = run_length as u32;
}
for entry in entries.iter_mut() {
let (length, consumed) = decode_varint(&data[offset..])?;
offset += consumed;
entry.length = length as u32;
}
let mut expected_offset = 0u64;
for (i, entry) in entries.iter_mut().enumerate() {
let (encoded_offset, consumed) = decode_varint(&data[offset..])?;
offset += consumed;
if encoded_offset == 0 && i > 0 {
entry.offset = expected_offset;
} else {
entry.offset = encoded_offset.saturating_sub(1);
}
expected_offset = entry.offset.checked_add(u64::from(entry.length))?;
}
Some(entries)
}
const MAX_ROOT_DIR_BYTES: usize = 16384 - 127;
const INITIAL_LEAF_SIZE: usize = 4096;
#[derive(Debug)]
pub struct DirectoryLayout {
pub root_bytes: Vec<u8>,
pub leaves_bytes: Vec<u8>,
pub num_leaves: usize,
}
fn build_root_leaves(
entries: &[DirEntry],
leaf_size: usize,
compression: Compression,
) -> std::io::Result<DirectoryLayout> {
let mut root_entries = Vec::new();
let mut leaves_bytes = Vec::new();
let mut num_leaves = 0;
for chunk in entries.chunks(leaf_size) {
num_leaves += 1;
let leaf_encoded = encode_directory(chunk);
let leaf_compressed = compression::compress(&leaf_encoded, compression)?;
root_entries.push(DirEntry {
tile_id: chunk[0].tile_id,
offset: leaves_bytes.len() as u64,
length: leaf_compressed.len() as u32,
run_length: 0, });
leaves_bytes.extend(leaf_compressed);
}
let root_encoded = encode_directory(&root_entries);
let root_compressed = compression::compress(&root_encoded, compression)?;
Ok(DirectoryLayout {
root_bytes: root_compressed,
leaves_bytes,
num_leaves,
})
}
pub fn make_root_leaves(
entries: &[DirEntry],
compression: Compression,
) -> std::io::Result<DirectoryLayout> {
let single_encoded = encode_directory(entries);
let single_compressed = compression::compress(&single_encoded, compression)?;
if single_compressed.len() <= MAX_ROOT_DIR_BYTES {
return Ok(DirectoryLayout {
root_bytes: single_compressed,
leaves_bytes: Vec::new(),
num_leaves: 0,
});
}
let mut leaf_size = INITIAL_LEAF_SIZE;
loop {
let layout = build_root_leaves(entries, leaf_size, compression)?;
if layout.root_bytes.len() <= MAX_ROOT_DIR_BYTES {
return Ok(layout);
}
leaf_size *= 2;
if leaf_size > entries.len() * 2 {
return build_root_leaves(entries, entries.len(), compression);
}
}
}
pub fn gzip_compress(data: &[u8]) -> std::io::Result<Vec<u8>> {
compression::compress(data, Compression::Gzip)
}
fn fields_json(fields: &HashMap<String, String>) -> String {
if fields.is_empty() {
return "{}".to_string();
}
let mut field_pairs: Vec<_> = fields.iter().collect();
field_pairs.sort_by_key(|(k, _)| *k);
let field_strings: Vec<String> = field_pairs
.iter()
.map(|(name, type_str)| format!(r#""{}":"{}""#, name, type_str))
.collect();
format!("{{{}}}", field_strings.join(","))
}
fn tilestats_json(layer_name: &str, total_features: u64, field_count: usize) -> String {
if total_features == 0 {
return String::new();
}
format!(
r#""tilestats":{{"layerCount":1,"layers":[{{"layer":"{}","count":{},"attributeCount":{}}}]}},"#,
layer_name, total_features, field_count
)
}
#[derive(Debug, Clone)]
struct TileEntry {
data: Option<Vec<u8>>,
hash: u64,
}
pub struct PmtilesWriter {
tiles: BTreeMap<u64, TileEntry>,
min_zoom: u8,
max_zoom: u8,
bounds: TileBounds,
layer_name: String,
fields: HashMap<String, String>,
total_features: u64,
features_per_zoom: HashMap<u8, u64>,
tile_compression: Compression,
internal_compression: Compression,
dedup_enabled: bool,
dedup_cache: DeduplicationCache,
vector_layers_json: Option<String>,
}
impl PmtilesWriter {
pub fn new() -> Self {
Self {
tiles: BTreeMap::new(),
min_zoom: 255,
max_zoom: 0,
bounds: TileBounds::empty(),
layer_name: "layer".to_string(),
fields: HashMap::new(),
total_features: 0,
features_per_zoom: HashMap::new(),
tile_compression: Compression::Gzip,
internal_compression: Compression::Gzip,
dedup_enabled: false,
dedup_cache: DeduplicationCache::new(),
vector_layers_json: None,
}
}
pub fn with_compression(compression: Compression) -> Self {
Self {
tiles: BTreeMap::new(),
min_zoom: 255,
max_zoom: 0,
bounds: TileBounds::empty(),
layer_name: "layer".to_string(),
fields: HashMap::new(),
total_features: 0,
features_per_zoom: HashMap::new(),
tile_compression: compression,
internal_compression: compression,
dedup_enabled: false,
dedup_cache: DeduplicationCache::new(),
vector_layers_json: None,
}
}
pub fn enable_deduplication(&mut self, enabled: bool) {
self.dedup_enabled = enabled;
}
pub fn set_tile_compression(&mut self, compression: Compression) {
self.tile_compression = compression;
}
pub fn set_internal_compression(&mut self, compression: Compression) {
self.internal_compression = compression;
}
pub fn tile_compression(&self) -> Compression {
self.tile_compression
}
pub fn internal_compression(&self) -> Compression {
self.internal_compression
}
pub fn is_dedup_enabled(&self) -> bool {
self.dedup_enabled
}
pub fn dedup_stats(&self) -> &DeduplicationStats {
self.dedup_cache.stats()
}
pub fn set_layer_name(&mut self, name: &str) {
self.layer_name = name.to_string();
}
pub fn set_vector_layers_json(&mut self, json: String) {
self.vector_layers_json = Some(json);
}
pub fn set_fields(&mut self, fields: HashMap<String, String>) {
self.fields = fields;
}
fn build_fields_json(&self) -> String {
fields_json(&self.fields)
}
fn build_tilestats_json(&self) -> String {
tilestats_json(&self.layer_name, self.total_features, self.fields.len())
}
pub fn add_tile(&mut self, z: u8, x: u32, y: u32, data: &[u8]) -> std::io::Result<()> {
self.add_tile_with_count(z, x, y, data, 0)
}
pub fn add_tile_with_count(
&mut self,
z: u8,
x: u32,
y: u32,
data: &[u8],
feature_count: usize,
) -> std::io::Result<()> {
let id = tile_id(z, x, y);
let uncompressed_size = data.len() as u32;
self.min_zoom = self.min_zoom.min(z);
self.max_zoom = self.max_zoom.max(z);
self.total_features += feature_count as u64;
*self.features_per_zoom.entry(z).or_insert(0) += feature_count as u64;
if self.dedup_enabled {
let hash = TileHasher::hash(data);
if self.dedup_cache.check(hash).is_some() {
self.dedup_cache.record_duplicate(uncompressed_size);
self.tiles.insert(
id,
TileEntry {
data: None, hash,
},
);
} else {
let compressed = compression::compress(data, self.tile_compression)?;
let compressed_len = compressed.len() as u32;
self.dedup_cache
.record_new(hash, 0, compressed_len, uncompressed_size);
self.tiles.insert(
id,
TileEntry {
data: Some(compressed),
hash,
},
);
}
} else {
let compressed = compression::compress(data, self.tile_compression)?;
let hash = TileHasher::hash(data);
self.tiles.insert(
id,
TileEntry {
data: Some(compressed),
hash,
},
);
}
Ok(())
}
pub fn add_tile_compressed(
&mut self,
z: u8,
x: u32,
y: u32,
compressed_data: Vec<u8>,
) -> std::io::Result<()> {
let id = tile_id(z, x, y);
let hash = TileHasher::hash(&compressed_data);
self.tiles.insert(
id,
TileEntry {
data: Some(compressed_data),
hash,
},
);
self.min_zoom = self.min_zoom.min(z);
self.max_zoom = self.max_zoom.max(z);
Ok(())
}
pub fn set_bounds(&mut self, bounds: &TileBounds) {
self.bounds = TileBounds::new(
bounds.lng_min,
bounds.lat_min.clamp(-85.05, 85.05),
bounds.lng_max,
bounds.lat_max.clamp(-85.05, 85.05),
);
}
pub fn tile_count(&self) -> usize {
self.tiles.len()
}
pub fn write_to_file(&self, path: &Path) -> Result<()> {
let file = File::create(path)
.map_err(|e| Error::PMTilesWrite(format!("Failed to create file: {}", e)))?;
let mut writer = BufWriter::new(file);
let mut tile_data_buf = Vec::new();
let mut entries = Vec::new();
let mut hash_to_offset: HashMap<u64, (u64, u32)> = HashMap::new();
let mut unique_contents = 0u64;
if self.dedup_enabled {
for (&id, entry) in &self.tiles {
let (offset, length) = if let Some(ref data) = entry.data {
let offset = tile_data_buf.len() as u64;
let length = data.len() as u32;
tile_data_buf.extend_from_slice(data);
hash_to_offset.insert(entry.hash, (offset, length));
unique_contents += 1;
(offset, length)
} else {
*hash_to_offset.get(&entry.hash).expect("Hash must exist")
};
if let Some(last) = entries.last_mut() {
let last_entry: &mut DirEntry = last;
if last_entry.offset == offset
&& id == last_entry.tile_id + last_entry.run_length as u64
{
last_entry.run_length += 1;
continue;
}
}
entries.push(DirEntry {
tile_id: id,
offset,
length,
run_length: 1,
});
}
} else {
for (&id, entry) in &self.tiles {
let data = entry.data.as_ref().expect("Non-dedup tiles must have data");
entries.push(DirEntry {
tile_id: id,
offset: tile_data_buf.len() as u64,
length: data.len() as u32,
run_length: 1,
});
tile_data_buf.extend_from_slice(data);
unique_contents += 1;
}
}
let layout = make_root_leaves(&entries, self.internal_compression)
.map_err(|e| Error::PMTilesWrite(format!("Failed to build directories: {}", e)))?;
let compressed_dir = layout.root_bytes;
let leaves_bytes = layout.leaves_bytes;
let min_z = if self.min_zoom == 255 {
0
} else {
self.min_zoom
};
let max_z = if self.max_zoom == 0 && self.tiles.is_empty() {
0
} else {
self.max_zoom
};
let tilestats_json = self.build_tilestats_json();
let vector_layers = match &self.vector_layers_json {
Some(json) => json.clone(),
None => format!(
r#"[{{"id":"{}","minzoom":{},"maxzoom":{},"fields":{}}}]"#,
self.layer_name,
min_z,
max_z,
self.build_fields_json()
),
};
let metadata = format!(
r#"{{"vector_layers":{},{}"format":"pbf","generator":"tylertoo"}}"#,
vector_layers, tilestats_json
);
let compressed_metadata =
compression::compress(metadata.as_bytes(), self.internal_compression)
.map_err(|e| Error::PMTilesWrite(format!("Failed to compress metadata: {}", e)))?;
let root_dir_offset = 127u64;
let root_dir_length = compressed_dir.len() as u64;
let metadata_offset = root_dir_offset + root_dir_length;
let metadata_length = compressed_metadata.len() as u64;
let leaf_dirs_offset = metadata_offset + metadata_length;
let leaf_dirs_length = leaves_bytes.len() as u64;
let tile_data_offset = leaf_dirs_offset + leaf_dirs_length;
let tile_data_length = tile_data_buf.len() as u64;
let header = Header {
root_dir_offset,
root_dir_length,
json_metadata_offset: metadata_offset,
json_metadata_length: metadata_length,
leaf_dirs_offset,
leaf_dirs_length,
tile_data_offset,
tile_data_length,
addressed_tiles_count: self.tiles.len() as u64,
tile_entries_count: entries.len() as u64,
tile_contents_count: unique_contents,
clustered: true,
internal_compression: self.internal_compression,
tile_compression: self.tile_compression,
tile_type: TileType::Mvt,
min_zoom: if self.min_zoom == 255 {
0
} else {
self.min_zoom
},
max_zoom: if self.max_zoom == 0 && self.tiles.is_empty() {
0
} else {
self.max_zoom
},
min_lon: self.bounds.lng_min,
min_lat: self.bounds.lat_min,
max_lon: self.bounds.lng_max,
max_lat: self.bounds.lat_max,
center_zoom: if self.tiles.is_empty() {
0
} else {
(self.min_zoom + self.max_zoom) / 2
},
center_lon: (self.bounds.lng_min + self.bounds.lng_max) / 2.0,
center_lat: (self.bounds.lat_min + self.bounds.lat_max) / 2.0,
};
writer
.write_all(&header.to_bytes())
.map_err(|e| Error::PMTilesWrite(format!("Failed to write header: {}", e)))?;
writer
.write_all(&compressed_dir)
.map_err(|e| Error::PMTilesWrite(format!("Failed to write directory: {}", e)))?;
writer
.write_all(&compressed_metadata)
.map_err(|e| Error::PMTilesWrite(format!("Failed to write metadata: {}", e)))?;
writer
.write_all(&leaves_bytes)
.map_err(|e| Error::PMTilesWrite(format!("Failed to write leaf directories: {}", e)))?;
writer
.write_all(&tile_data_buf)
.map_err(|e| Error::PMTilesWrite(format!("Failed to write tile data: {}", e)))?;
writer
.flush()
.map_err(|e| Error::PMTilesWrite(format!("Failed to flush: {}", e)))?;
Ok(())
}
}
impl Default for PmtilesWriter {
fn default() -> Self {
Self::new()
}
}
use std::path::PathBuf;
#[derive(Debug, Clone)]
struct StreamingDirEntry {
tile_id: u64,
offset: u64,
length: u32,
}
#[derive(Debug, Clone, Default)]
pub struct StreamingWriteStats {
pub total_tiles: u64,
pub unique_tiles: u64,
pub bytes_written: u64,
pub bytes_saved_dedup: u64,
}
impl StreamingWriteStats {
pub fn estimated_memory_bytes(&self) -> u64 {
self.total_tiles * 24 + self.unique_tiles * 40
}
}
pub struct StreamingPmtilesWriter {
temp_file: Option<BufWriter<File>>,
temp_path: PathBuf,
entries: Vec<StreamingDirEntry>,
dedup_cache: HashMap<u64, (u64, u32)>,
current_offset: u64,
min_zoom: u8,
max_zoom: u8,
declared_min_zoom: Option<u8>,
bounds: TileBounds,
layer_name: String,
fields: HashMap<String, String>,
vector_layers_json: Option<String>,
tile_compression: Compression,
internal_compression: Compression,
stats: StreamingWriteStats,
total_features: u64,
finalized: bool,
}
impl StreamingPmtilesWriter {
pub fn new(compression: Compression) -> std::io::Result<Self> {
Self::with_temp_dir(compression, std::env::temp_dir())
}
pub fn with_temp_dir(compression: Compression, temp_dir: PathBuf) -> std::io::Result<Self> {
use std::time::{SystemTime, UNIX_EPOCH};
let timestamp = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_nanos();
let pid = std::process::id();
let tid = std::thread::current().id();
let temp_path = temp_dir.join(format!("tylertoo-{}-{}-{:?}.tmp", timestamp, pid, tid));
let file = File::create(&temp_path)?;
let temp_file = BufWriter::with_capacity(64 * 1024, file);
Ok(Self {
temp_file: Some(temp_file),
temp_path,
entries: Vec::new(),
dedup_cache: HashMap::new(),
current_offset: 0,
min_zoom: 255,
max_zoom: 0,
declared_min_zoom: None,
bounds: TileBounds::empty(),
layer_name: "layer".to_string(),
fields: HashMap::new(),
vector_layers_json: None,
tile_compression: compression,
internal_compression: compression,
stats: StreamingWriteStats::default(),
total_features: 0,
finalized: false,
})
}
pub fn temp_path(&self) -> &Path {
&self.temp_path
}
pub fn set_layer_name(&mut self, name: &str) {
self.layer_name = name.to_string();
}
pub fn set_declared_min_zoom(&mut self, zoom: u8) {
self.declared_min_zoom = Some(zoom);
}
fn header_min_zoom(&self) -> u8 {
if self.entries.is_empty() {
return 0;
}
match self.declared_min_zoom {
Some(d) => self.min_zoom.min(d),
None => self.min_zoom,
}
}
pub fn set_fields(&mut self, fields: HashMap<String, String>) {
self.fields = fields;
}
pub fn set_vector_layers_json(&mut self, json: String) {
self.vector_layers_json = Some(json);
}
pub fn set_bounds(&mut self, bounds: &TileBounds) {
self.bounds = TileBounds::new(
bounds.lng_min,
bounds.lat_min.clamp(-85.05, 85.05),
bounds.lng_max,
bounds.lat_max.clamp(-85.05, 85.05),
);
}
pub fn stats(&self) -> &StreamingWriteStats {
&self.stats
}
pub fn add_tile(&mut self, z: u8, x: u32, y: u32, data: &[u8]) -> std::io::Result<()> {
self.add_tile_with_count(z, x, y, data, 0)
}
pub fn add_tile_with_count(
&mut self,
z: u8,
x: u32,
y: u32,
data: &[u8],
feature_count: usize,
) -> std::io::Result<()> {
let temp_file = self
.temp_file
.as_mut()
.ok_or_else(|| std::io::Error::other("Writer already finalized"))?;
let id = tile_id(z, x, y);
self.stats.total_tiles += 1;
self.total_features += feature_count as u64;
self.min_zoom = self.min_zoom.min(z);
self.max_zoom = self.max_zoom.max(z);
let hash = crate::dedup::TileHasher::hash(data);
if let Some((offset, length)) = self.dedup_cache.get(&hash) {
self.entries.push(StreamingDirEntry {
tile_id: id,
offset: *offset,
length: *length,
});
self.stats.bytes_saved_dedup += data.len() as u64;
return Ok(());
}
let compressed = compression::compress(data, self.tile_compression)?;
let compressed_len = compressed.len() as u32;
temp_file.write_all(&compressed)?;
let offset = self.current_offset;
self.dedup_cache.insert(hash, (offset, compressed_len));
self.entries.push(StreamingDirEntry {
tile_id: id,
offset,
length: compressed_len,
});
self.current_offset += compressed_len as u64;
self.stats.unique_tiles += 1;
self.stats.bytes_written += compressed_len as u64;
Ok(())
}
#[allow(clippy::too_many_arguments)]
pub fn add_tile_precompressed(
&mut self,
z: u8,
x: u32,
y: u32,
hash: u64,
compressed: &[u8],
raw_len: usize,
feature_count: usize,
) -> std::io::Result<()> {
let temp_file = self
.temp_file
.as_mut()
.ok_or_else(|| std::io::Error::other("Writer already finalized"))?;
let id = tile_id(z, x, y);
self.stats.total_tiles += 1;
self.total_features += feature_count as u64;
self.min_zoom = self.min_zoom.min(z);
self.max_zoom = self.max_zoom.max(z);
if let Some((offset, length)) = self.dedup_cache.get(&hash) {
self.entries.push(StreamingDirEntry {
tile_id: id,
offset: *offset,
length: *length,
});
self.stats.bytes_saved_dedup += raw_len as u64;
return Ok(());
}
let compressed_len = compressed.len() as u32;
temp_file.write_all(compressed)?;
let offset = self.current_offset;
self.dedup_cache.insert(hash, (offset, compressed_len));
self.entries.push(StreamingDirEntry {
tile_id: id,
offset,
length: compressed_len,
});
self.current_offset += compressed_len as u64;
self.stats.unique_tiles += 1;
self.stats.bytes_written += compressed_len as u64;
Ok(())
}
pub fn finalize(mut self, output_path: &Path) -> Result<StreamingWriteStats> {
self.write_archive(output_path)?;
drop(self.temp_file.take());
let _ = std::fs::remove_file(&self.temp_path);
self.finalized = true;
Ok(self.stats.clone())
}
pub fn checkpoint(&mut self, output_path: &Path) -> Result<()> {
self.write_archive(output_path)
}
fn write_archive(&mut self, output_path: &Path) -> Result<()> {
match self.temp_file.as_mut() {
Some(tf) => tf
.flush()
.map_err(|e| Error::PMTilesWrite(format!("Failed to flush temp file: {}", e)))?,
None => return Err(Error::PMTilesWrite("Writer already finalized".to_string())),
}
self.entries.sort_by_key(|e| e.tile_id);
let dir_entries = self.build_directory_entries();
let dir_layout = make_root_leaves(&dir_entries, self.internal_compression)
.map_err(|e| Error::PMTilesWrite(format!("Failed to build directory: {}", e)))?;
let metadata = self.build_metadata_json();
let compressed_metadata =
compression::compress(metadata.as_bytes(), self.internal_compression)
.map_err(|e| Error::PMTilesWrite(format!("Failed to compress metadata: {}", e)))?;
let root_dir_offset = 127u64;
let root_dir_length = dir_layout.root_bytes.len() as u64;
let metadata_offset = root_dir_offset + root_dir_length;
let metadata_length = compressed_metadata.len() as u64;
let leaf_dirs_offset = metadata_offset + metadata_length;
let leaf_dirs_length = dir_layout.leaves_bytes.len() as u64;
let tile_data_offset = leaf_dirs_offset + leaf_dirs_length;
let tile_data_length = self.current_offset;
let header = Header {
root_dir_offset,
root_dir_length,
json_metadata_offset: metadata_offset,
json_metadata_length: metadata_length,
leaf_dirs_offset,
leaf_dirs_length,
tile_data_offset,
tile_data_length,
addressed_tiles_count: self.stats.total_tiles,
tile_entries_count: dir_entries.len() as u64,
tile_contents_count: self.stats.unique_tiles,
clustered: true,
internal_compression: self.internal_compression,
tile_compression: self.tile_compression,
tile_type: TileType::Mvt,
min_zoom: self.header_min_zoom(),
max_zoom: if self.max_zoom == 0 && self.entries.is_empty() {
0
} else {
self.max_zoom
},
min_lon: self.bounds.lng_min,
min_lat: self.bounds.lat_min,
max_lon: self.bounds.lng_max,
max_lat: self.bounds.lat_max,
center_zoom: if self.entries.is_empty() {
0
} else {
(self.header_min_zoom() + self.max_zoom) / 2
},
center_lon: (self.bounds.lng_min + self.bounds.lng_max) / 2.0,
center_lat: (self.bounds.lat_min + self.bounds.lat_max) / 2.0,
};
let partial_path = {
let mut os = output_path.as_os_str().to_owned();
os.push(".partial");
PathBuf::from(os)
};
let output_file = File::create(&partial_path)
.map_err(|e| Error::PMTilesWrite(format!("Failed to create output file: {}", e)))?;
let mut writer = BufWriter::new(output_file);
writer
.write_all(&header.to_bytes())
.map_err(|e| Error::PMTilesWrite(format!("Failed to write header: {}", e)))?;
writer
.write_all(&dir_layout.root_bytes)
.map_err(|e| Error::PMTilesWrite(format!("Failed to write root directory: {}", e)))?;
writer
.write_all(&compressed_metadata)
.map_err(|e| Error::PMTilesWrite(format!("Failed to write metadata: {}", e)))?;
if !dir_layout.leaves_bytes.is_empty() {
writer.write_all(&dir_layout.leaves_bytes).map_err(|e| {
Error::PMTilesWrite(format!("Failed to write leaf directories: {}", e))
})?;
}
let mut temp_reader = File::open(&self.temp_path)
.map_err(|e| Error::PMTilesWrite(format!("Failed to reopen temp file: {}", e)))?;
std::io::copy(&mut temp_reader, &mut writer)
.map_err(|e| Error::PMTilesWrite(format!("Failed to copy tile data: {}", e)))?;
writer
.flush()
.map_err(|e| Error::PMTilesWrite(format!("Failed to flush output: {}", e)))?;
drop(writer);
std::fs::rename(&partial_path, output_path)
.map_err(|e| Error::PMTilesWrite(format!("Failed to publish archive: {}", e)))?;
Ok(())
}
fn build_directory_entries(&self) -> Vec<DirEntry> {
let mut dir_entries = Vec::new();
for entry in &self.entries {
if let Some(last) = dir_entries.last_mut() {
let last_entry: &mut DirEntry = last;
if last_entry.offset == entry.offset
&& entry.tile_id == last_entry.tile_id + last_entry.run_length as u64
{
last_entry.run_length += 1;
continue;
}
}
dir_entries.push(DirEntry {
tile_id: entry.tile_id,
offset: entry.offset,
length: entry.length,
run_length: 1,
});
}
dir_entries
}
fn build_metadata_json(&self) -> String {
let min_z = self.header_min_zoom();
let max_z = if self.max_zoom == 0 && self.entries.is_empty() {
0
} else {
self.max_zoom
};
let tilestats_json = self.build_tilestats_json();
let vector_layers = match &self.vector_layers_json {
Some(json) => json.clone(),
None => format!(
r#"[{{"id":"{}","minzoom":{},"maxzoom":{},"fields":{}}}]"#,
self.layer_name,
min_z,
max_z,
self.build_fields_json()
),
};
format!(
r#"{{"vector_layers":{},{}"format":"pbf","generator":"tylertoo"}}"#,
vector_layers, tilestats_json
)
}
fn build_fields_json(&self) -> String {
fields_json(&self.fields)
}
fn build_tilestats_json(&self) -> String {
tilestats_json(&self.layer_name, self.total_features, self.fields.len())
}
}
impl Drop for StreamingPmtilesWriter {
fn drop(&mut self) {
if !self.finalized {
let _ = std::fs::remove_file(&self.temp_path);
}
}
}
#[cfg(test)]
mod tests {
#[test]
fn large_archive_spills_into_leaf_directories() {
let mut writer = PmtilesWriter::new();
let mut rng = 0x2545_F491_4F6C_DD1Du64;
for i in 0..40_000u32 {
let (x, y) = ((i % 200) * 20, (i / 200) * 20);
rng = rng
.wrapping_mul(6364136223846793005)
.wrapping_add(1442695040888963407);
let len = 64 + (rng >> 33) as usize % 4096;
writer
.add_tile_compressed(14, x, y, vec![(i % 251) as u8; len])
.unwrap();
}
let tmp = tempfile::NamedTempFile::new().unwrap();
writer.write_to_file(tmp.path()).unwrap();
let bytes = std::fs::read(tmp.path()).unwrap();
let header = Header::from_bytes(&bytes).unwrap();
assert!(
header.root_dir_length <= 16384,
"root directory must fit the spec budget, got {}",
header.root_dir_length
);
assert!(
header.leaf_dirs_length > 0,
"an archive too big for one root must have leaf directories"
);
assert_eq!(
header.leaf_dirs_offset,
header.json_metadata_offset + header.json_metadata_length
);
assert_eq!(
header.tile_data_offset,
header.leaf_dirs_offset + header.leaf_dirs_length
);
let read_dir = |raw: &[u8]| -> Vec<DirEntry> {
let plain = compression::decompress(raw, header.internal_compression).unwrap();
decode_directory(&plain).expect("directory must decode")
};
let root = read_dir(
&bytes[header.root_dir_offset as usize
..(header.root_dir_offset + header.root_dir_length) as usize],
);
assert!(
root.iter().any(|e| e.run_length == 0),
"a spilled archive's root must contain at least one leaf pointer"
);
let mut found: HashMap<u64, Vec<u8>> = HashMap::new();
for e in &root {
let leaves = if e.run_length == 0 {
let start = (header.leaf_dirs_offset + e.offset) as usize;
read_dir(&bytes[start..start + e.length as usize])
} else {
vec![e.clone()]
};
for le in leaves {
let start = (header.tile_data_offset + le.offset) as usize;
found.insert(
le.tile_id,
bytes[start..start + le.length as usize].to_vec(),
);
}
}
assert_eq!(found.len(), 40_000, "every tile must be addressable");
for i in [0u32, 1, 19_899, 39_998, 39_999] {
let (x, y) = ((i % 200) * 20, (i / 200) * 20);
let data = found
.get(&tile_id(14, x, y))
.unwrap_or_else(|| panic!("tile {i} (z14/{x}/{y}) not found"));
let want = (i % 251) as u8;
assert!(
data.iter().all(|&b| b == want),
"tile {i} (z14/{x}/{y}) resolved to the wrong bytes: \
expected all {want}, got {:?}..",
&data[..data.len().min(8)]
);
}
}
use super::*;
use std::fs;
#[test]
fn test_header_size_is_127_bytes() {
let header = Header::default();
let bytes = header.to_bytes();
assert_eq!(
bytes.len(),
127,
"PMTiles v3 header must be exactly 127 bytes"
);
}
#[test]
fn test_header_magic_and_version() {
let header = Header::default();
let bytes = header.to_bytes();
assert_eq!(&bytes[0..7], b"PMTiles", "Magic number must be 'PMTiles'");
assert_eq!(bytes[7], 3, "Version must be 3");
}
#[test]
fn test_header_default_offsets() {
let header = Header::default();
let bytes = header.to_bytes();
let root_offset = u64::from_le_bytes(bytes[8..16].try_into().unwrap());
assert_eq!(root_offset, 127);
}
#[test]
fn test_header_bounds_encoding() {
let header = Header {
min_lon: -122.4194, min_lat: 37.7749,
max_lon: -122.3894,
max_lat: 37.8049,
..Default::default()
};
let bytes = header.to_bytes();
let min_lon_encoded = i32::from_le_bytes(bytes[102..106].try_into().unwrap());
let min_lon_decoded = min_lon_encoded as f64 / 10_000_000.0;
assert!(
(min_lon_decoded - header.min_lon).abs() < 0.0001,
"Lon encoding should preserve precision to ~0.0001 degrees"
);
}
#[test]
fn test_tile_id_zoom_0() {
assert_eq!(tile_id(0, 0, 0), 0);
}
#[test]
fn test_tile_id_zoom_1_matches_spec() {
assert_eq!(tile_id(1, 0, 0), 1);
assert_eq!(tile_id(1, 0, 1), 2);
assert_eq!(tile_id(1, 1, 1), 3);
assert_eq!(tile_id(1, 1, 0), 4);
}
#[test]
fn test_tile_id_zoom_2_base() {
assert_eq!(tile_id(2, 0, 0), 5);
}
#[test]
fn test_tile_id_unique_at_each_zoom() {
for z in 0..=4u8 {
let mut ids = Vec::new();
let n = 1u32 << z;
for y in 0..n {
for x in 0..n {
ids.push(tile_id(z, x, y));
}
}
let original_len = ids.len();
ids.sort();
ids.dedup();
assert_eq!(
ids.len(),
original_len,
"All tile IDs at zoom {} should be unique",
z
);
}
}
#[test]
fn test_tile_id_to_zxy_matches_spec_examples() {
assert_eq!(tile_id_to_zxy(0).unwrap(), (0, 0, 0));
assert_eq!(tile_id_to_zxy(1).unwrap(), (1, 0, 0));
assert_eq!(tile_id_to_zxy(2).unwrap(), (1, 0, 1));
assert_eq!(tile_id_to_zxy(3).unwrap(), (1, 1, 1));
assert_eq!(tile_id_to_zxy(4).unwrap(), (1, 1, 0));
assert_eq!(tile_id_to_zxy(5).unwrap(), (2, 0, 0));
}
#[test]
fn test_tile_id_to_zxy_inverts_tile_id_exhaustively() {
for z in 0..=5u8 {
let n = 1u32 << z;
for y in 0..n {
for x in 0..n {
assert_eq!(
tile_id_to_zxy(tile_id(z, x, y)).unwrap(),
(z, x, y),
"round-trip z={z} x={x} y={y}"
);
}
}
}
for (z, x, y) in [
(14u8, 4823u32, 6160u32),
(14, 0, 0),
(14, (1 << 14) - 1, (1 << 14) - 1),
(20, 123_456, 654_321),
(31, (1u32 << 31) - 1, 0),
] {
assert_eq!(tile_id_to_zxy(tile_id(z, x, y)).unwrap(), (z, x, y));
}
}
#[test]
fn test_tile_id_to_zxy_rejects_out_of_range() {
let past_z31 = (0..=31u8).map(|z| 1u64 << (2 * u64::from(z))).sum::<u64>();
assert!(tile_id_to_zxy(past_z31).is_err());
assert!(tile_id_to_zxy(u64::MAX).is_err());
}
#[test]
fn test_header_from_bytes_roundtrips_to_bytes() {
let header = Header {
root_dir_offset: 127,
root_dir_length: 421,
json_metadata_offset: 548,
json_metadata_length: 33,
leaf_dirs_offset: 581,
leaf_dirs_length: 1290,
tile_data_offset: 1871,
tile_data_length: 999_999,
addressed_tiles_count: 42,
tile_entries_count: 40,
tile_contents_count: 39,
clustered: true,
internal_compression: Compression::Gzip,
tile_compression: Compression::Zstd,
tile_type: TileType::Mvt,
min_zoom: 3,
max_zoom: 14,
min_lon: -75.1652,
min_lat: -33.8688,
max_lon: 151.2093,
max_lat: 48.8566,
center_zoom: 8,
center_lon: 2.3522,
center_lat: 39.9526,
};
let parsed = Header::from_bytes(&header.to_bytes()).unwrap();
assert_eq!(parsed.root_dir_offset, header.root_dir_offset);
assert_eq!(parsed.root_dir_length, header.root_dir_length);
assert_eq!(parsed.json_metadata_offset, header.json_metadata_offset);
assert_eq!(parsed.json_metadata_length, header.json_metadata_length);
assert_eq!(parsed.leaf_dirs_offset, header.leaf_dirs_offset);
assert_eq!(parsed.leaf_dirs_length, header.leaf_dirs_length);
assert_eq!(parsed.tile_data_offset, header.tile_data_offset);
assert_eq!(parsed.tile_data_length, header.tile_data_length);
assert_eq!(parsed.addressed_tiles_count, header.addressed_tiles_count);
assert_eq!(parsed.tile_entries_count, header.tile_entries_count);
assert_eq!(parsed.tile_contents_count, header.tile_contents_count);
assert_eq!(parsed.clustered, header.clustered);
assert_eq!(parsed.internal_compression, header.internal_compression);
assert_eq!(parsed.tile_compression, header.tile_compression);
assert_eq!(parsed.tile_type, header.tile_type);
assert_eq!(parsed.min_zoom, header.min_zoom);
assert_eq!(parsed.max_zoom, header.max_zoom);
for (got, want) in [
(parsed.min_lon, header.min_lon),
(parsed.min_lat, header.min_lat),
(parsed.max_lon, header.max_lon),
(parsed.max_lat, header.max_lat),
(parsed.center_lon, header.center_lon),
(parsed.center_lat, header.center_lat),
] {
assert!((got - want).abs() < 1e-6, "{got} vs {want}");
}
assert_eq!(parsed.center_zoom, header.center_zoom);
}
#[test]
fn test_header_from_bytes_rejects_garbage() {
assert!(Header::from_bytes(&[0u8; 50]).is_err());
let mut bytes = Header::default().to_bytes();
bytes[0] = b'X';
assert!(Header::from_bytes(&bytes).is_err());
let mut bytes = Header::default().to_bytes();
bytes[7] = 2;
assert!(Header::from_bytes(&bytes).is_err());
let mut bytes = Header::default().to_bytes();
bytes[97] = 9;
assert!(Header::from_bytes(&bytes).is_err());
let mut bytes = Header::default().to_bytes();
bytes[99] = 9;
assert!(Header::from_bytes(&bytes).is_err());
}
#[test]
fn test_tile_id_increasing_with_zoom() {
for z in 0..4u8 {
let n = 1u32 << z;
let max_id_at_z = (0..n)
.flat_map(|y| (0..n).map(move |x| tile_id(z, x, y)))
.max()
.unwrap();
let min_id_at_z_plus_1 = tile_id(z + 1, 0, 0);
assert!(
max_id_at_z < min_id_at_z_plus_1,
"Max ID at zoom {} ({}) should be < min ID at zoom {} ({})",
z,
max_id_at_z,
z + 1,
min_id_at_z_plus_1
);
}
}
#[test]
fn test_encode_varint_small_values() {
let mut buf = Vec::new();
encode_varint(0, &mut buf);
assert_eq!(buf, vec![0]);
buf.clear();
encode_varint(1, &mut buf);
assert_eq!(buf, vec![1]);
buf.clear();
encode_varint(127, &mut buf);
assert_eq!(buf, vec![127]);
}
#[test]
fn test_encode_varint_128() {
let mut buf = Vec::new();
encode_varint(128, &mut buf);
assert_eq!(buf, vec![0x80, 0x01]);
}
#[test]
fn test_encode_varint_300() {
let mut buf = Vec::new();
encode_varint(300, &mut buf);
assert_eq!(buf, vec![0xAC, 0x02]);
}
#[test]
fn test_varint_roundtrip() {
let test_values = [0u64, 1, 127, 128, 255, 256, 300, 16383, 16384, u64::MAX];
for &value in &test_values {
let mut buf = Vec::new();
encode_varint(value, &mut buf);
let (decoded, bytes_consumed) = decode_varint(&buf).expect("Should decode");
assert_eq!(decoded, value, "Roundtrip failed for {}", value);
assert_eq!(bytes_consumed, buf.len());
}
}
#[test]
fn test_encode_directory_empty() {
let entries: Vec<DirEntry> = vec![];
let encoded = encode_directory(&entries);
assert_eq!(encoded, vec![0]);
}
#[test]
fn test_encode_directory_single_entry() {
let entries = vec![DirEntry {
tile_id: 1,
offset: 0,
length: 100,
run_length: 1,
}];
let encoded = encode_directory(&entries);
assert!(!encoded.is_empty());
assert_eq!(encoded[0], 1);
}
#[test]
fn test_encode_directory_multiple_entries() {
let entries = vec![
DirEntry {
tile_id: 5,
offset: 0,
length: 100,
run_length: 1,
},
DirEntry {
tile_id: 42,
offset: 100,
length: 200,
run_length: 1,
},
DirEntry {
tile_id: 69,
offset: 300,
length: 50,
run_length: 1,
},
];
let encoded = encode_directory(&entries);
assert_eq!(encoded[0], 3);
assert!(encoded.len() < entries.len() * 24);
}
fn tippecanoe_style_leaf_root() -> Vec<u8> {
let mut buf = Vec::new();
encode_varint(3, &mut buf);
for delta in [0, 1000, 1000] {
encode_varint(delta, &mut buf);
}
for _ in 0..3 {
encode_varint(0, &mut buf);
}
for len in [100, 200, 50] {
encode_varint(len, &mut buf);
}
for encoded_offset in [1, 0, 0] {
encode_varint(encoded_offset, &mut buf);
}
buf
}
fn dir_tuples(entries: &[DirEntry]) -> Vec<(u64, u64, u32, u32)> {
entries
.iter()
.map(|e| (e.tile_id, e.offset, e.length, e.run_length))
.collect()
}
#[test]
fn test_decode_directory_resolves_contiguous_leaf_offsets() {
let entries = decode_directory(&tippecanoe_style_leaf_root()).unwrap();
assert_eq!(
dir_tuples(&entries),
vec![(0, 0, 100, 0), (1000, 100, 200, 0), (2000, 300, 50, 0)]
);
}
#[test]
fn test_encode_directory_contiguous_leaf_entries_encode_as_zero() {
let entries = vec![
DirEntry {
tile_id: 0,
offset: 0,
length: 100,
run_length: 0,
},
DirEntry {
tile_id: 1000,
offset: 100,
length: 200,
run_length: 0,
},
DirEntry {
tile_id: 2000,
offset: 300,
length: 50,
run_length: 0,
},
];
let encoded = encode_directory(&entries);
assert_eq!(encoded, tippecanoe_style_leaf_root());
assert_eq!(
dir_tuples(&decode_directory(&encoded).unwrap()),
dir_tuples(&entries)
);
}
#[test]
fn test_decode_directory_rejects_offset_that_overflows_the_accumulator() {
let mut buf = Vec::new();
encode_varint(2, &mut buf); encode_varint(0, &mut buf); encode_varint(1, &mut buf); encode_varint(1, &mut buf); encode_varint(1, &mut buf);
encode_varint(5, &mut buf); encode_varint(5, &mut buf);
encode_varint(u64::MAX, &mut buf); encode_varint(0, &mut buf);
assert!(decode_directory(&buf).is_none());
}
#[test]
fn test_gzip_compress_roundtrip() {
use flate2::read::GzDecoder;
use std::io::Read;
let original = b"Hello, PMTiles! This is test data.";
let compressed = gzip_compress(original).expect("Should compress");
let mut decoder = GzDecoder::new(&compressed[..]);
let mut decompressed = Vec::new();
decoder
.read_to_end(&mut decompressed)
.expect("Should decompress");
assert_eq!(decompressed, original);
}
#[test]
fn test_writer_creation() {
let writer = PmtilesWriter::new();
assert_eq!(writer.tile_count(), 0);
}
#[test]
fn test_writer_add_single_tile() {
let mut writer = PmtilesWriter::new();
let mvt_data = vec![0x1a, 0x00];
writer.add_tile(0, 0, 0, &mvt_data).unwrap();
assert_eq!(writer.tile_count(), 1);
}
#[test]
fn test_writer_creates_valid_pmtiles_file() {
let mut writer = PmtilesWriter::new();
let mvt_data = vec![0x1a, 0x00];
writer.add_tile(0, 0, 0, &mvt_data).unwrap();
writer.set_bounds(&TileBounds::new(-180.0, -85.0, 180.0, 85.0));
let path = Path::new("/tmp/test-pmtiles-writer.pmtiles");
let _ = fs::remove_file(path);
writer.write_to_file(path).expect("Should write file");
assert!(path.exists(), "File should exist");
let data = fs::read(path).unwrap();
assert_eq!(&data[0..7], b"PMTiles");
assert_eq!(data[7], 3);
assert!(data.len() > 127);
let root_offset = u64::from_le_bytes(data[8..16].try_into().unwrap());
assert_eq!(root_offset, 127);
let _ = fs::remove_file(path);
}
#[test]
fn test_writer_multiple_tiles_multiple_zooms() {
let mut writer = PmtilesWriter::new();
for z in 0..3u8 {
let n = 1u32 << z;
for x in 0..n {
for y in 0..n {
let mvt_data = vec![0x1a, z, x as u8, y as u8];
writer.add_tile(z, x, y, &mvt_data).unwrap();
}
}
}
assert_eq!(writer.tile_count(), 21);
writer.set_bounds(&TileBounds::new(-180.0, -85.0, 180.0, 85.0));
let path = Path::new("/tmp/test-pmtiles-multi.pmtiles");
let _ = fs::remove_file(path);
writer.write_to_file(path).expect("Should write file");
let data = fs::read(path).unwrap();
assert_eq!(&data[0..7], b"PMTiles");
assert_eq!(data[7], 3);
let addressed_count = u64::from_le_bytes(data[72..80].try_into().unwrap());
assert_eq!(addressed_count, 21);
assert_eq!(data[100], 0); assert_eq!(data[101], 2);
let _ = fs::remove_file(path);
}
#[test]
fn test_writer_empty_tileset() {
let writer = PmtilesWriter::new();
let path = Path::new("/tmp/test-pmtiles-empty.pmtiles");
let _ = fs::remove_file(path);
writer.write_to_file(path).expect("Should write empty file");
let data = fs::read(path).unwrap();
assert_eq!(&data[0..7], b"PMTiles");
let _ = fs::remove_file(path);
}
#[test]
fn test_writer_tile_ordering() {
let mut writer = PmtilesWriter::new();
writer.add_tile(2, 3, 3, &[1, 2, 3]).unwrap();
writer.add_tile(0, 0, 0, &[4, 5, 6]).unwrap();
writer.add_tile(1, 1, 0, &[7, 8, 9]).unwrap();
assert_eq!(writer.tile_count(), 3);
let path = Path::new("/tmp/test-pmtiles-ordering.pmtiles");
let _ = fs::remove_file(path);
writer.write_to_file(path).expect("Should write file");
assert!(path.exists());
let _ = fs::remove_file(path);
}
#[test]
fn test_writer_bounds_preserved() {
let mut writer = PmtilesWriter::new();
writer.add_tile(0, 0, 0, &[1, 2, 3]).unwrap();
let bounds = TileBounds::new(-122.5, 37.7, -122.3, 37.9);
writer.set_bounds(&bounds);
let path = Path::new("/tmp/test-pmtiles-bounds.pmtiles");
let _ = fs::remove_file(path);
writer.write_to_file(path).expect("Should write file");
let data = fs::read(path).unwrap();
let decode_coord = |offset: usize| -> f64 {
let val = i32::from_le_bytes(data[offset..offset + 4].try_into().unwrap());
val as f64 / 10_000_000.0
};
let min_lon = decode_coord(102);
let min_lat = decode_coord(106);
let max_lon = decode_coord(110);
let max_lat = decode_coord(114);
assert!((min_lon - bounds.lng_min).abs() < 0.0001);
assert!((min_lat - bounds.lat_min).abs() < 0.0001);
assert!((max_lon - bounds.lng_max).abs() < 0.0001);
assert!((max_lat - bounds.lat_max).abs() < 0.0001);
let _ = fs::remove_file(path);
}
#[test]
fn test_build_fields_json_empty() {
let writer = PmtilesWriter::new();
assert_eq!(writer.build_fields_json(), "{}");
}
#[test]
fn test_build_fields_json_with_fields() {
let mut writer = PmtilesWriter::new();
let mut fields = HashMap::new();
fields.insert("name".to_string(), "String".to_string());
fields.insert("area".to_string(), "Number".to_string());
writer.set_fields(fields);
let json = writer.build_fields_json();
assert_eq!(json, r#"{"area":"Number","name":"String"}"#);
}
#[test]
fn test_writer_field_metadata_in_output() {
use flate2::read::GzDecoder;
use std::io::Read;
let mut writer = PmtilesWriter::new();
writer.add_tile(0, 0, 0, &[1, 2, 3]).unwrap();
writer.set_layer_name("buildings");
let mut fields = HashMap::new();
fields.insert("name".to_string(), "String".to_string());
fields.insert("height".to_string(), "Number".to_string());
writer.set_fields(fields);
let path = Path::new("/tmp/test-pmtiles-fields.pmtiles");
let _ = fs::remove_file(path);
writer.write_to_file(path).expect("Should write file");
let data = fs::read(path).unwrap();
let metadata_offset = u64::from_le_bytes(data[24..32].try_into().unwrap()) as usize;
let metadata_length = u64::from_le_bytes(data[32..40].try_into().unwrap()) as usize;
let compressed_metadata = &data[metadata_offset..metadata_offset + metadata_length];
let mut decoder = GzDecoder::new(compressed_metadata);
let mut metadata_json = String::new();
decoder
.read_to_string(&mut metadata_json)
.expect("Should decompress metadata");
assert!(metadata_json.contains(r#""height":"Number""#));
assert!(metadata_json.contains(r#""name":"String""#));
assert!(metadata_json.contains(r#""id":"buildings""#));
let _ = fs::remove_file(path);
}
#[test]
fn streaming_writer_declared_min_zoom_widens_header_and_layer_range() {
use flate2::read::GzDecoder;
use std::io::Read;
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("declared.pmtiles");
let mut writer = StreamingPmtilesWriter::new(Compression::Gzip).unwrap();
writer.set_layer_name("t");
writer.set_declared_min_zoom(0);
writer.add_tile(2, 1, 1, &[0x1a, 0x00]).unwrap();
writer.add_tile(4, 5, 5, &[0x1a, 0x01]).unwrap();
writer.finalize(&path).unwrap();
let data = fs::read(&path).unwrap();
let header = Header::from_bytes(&data[..127]).unwrap();
assert_eq!(header.min_zoom, 0, "declared minimum wins over observed z2");
assert_eq!(header.max_zoom, 4);
let start = header.json_metadata_offset as usize;
let end = start + header.json_metadata_length as usize;
let mut json = String::new();
GzDecoder::new(&data[start..end])
.read_to_string(&mut json)
.unwrap();
assert!(
json.contains(r#""minzoom":0"#),
"vector_layers must advertise the declared minimum: {json}"
);
}
#[test]
fn streaming_writer_declared_min_zoom_never_hides_written_tiles() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("declared-narrow.pmtiles");
let mut writer = StreamingPmtilesWriter::new(Compression::Gzip).unwrap();
writer.set_layer_name("t");
writer.set_declared_min_zoom(3);
writer.add_tile(2, 1, 1, &[0x1a, 0x00]).unwrap();
writer.finalize(&path).unwrap();
let data = fs::read(&path).unwrap();
assert_eq!(Header::from_bytes(&data[..127]).unwrap().min_zoom, 2);
}
#[test]
fn streaming_writer_declared_min_zoom_over_zero_tiles_stays_z0() {
use flate2::read::GzDecoder;
use std::io::Read;
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("declared-empty.pmtiles");
let mut writer = StreamingPmtilesWriter::new(Compression::Gzip).unwrap();
writer.set_layer_name("t");
writer.set_declared_min_zoom(3);
writer.finalize(&path).unwrap();
let data = fs::read(&path).unwrap();
let header = Header::from_bytes(&data[..127]).unwrap();
assert_eq!(header.min_zoom, 0, "empty archive is z0..z0, not z3..z0");
assert_eq!(header.max_zoom, 0);
assert_eq!(header.center_zoom, 0);
let start = header.json_metadata_offset as usize;
let end = start + header.json_metadata_length as usize;
let mut json = String::new();
GzDecoder::new(&data[start..end])
.read_to_string(&mut json)
.unwrap();
assert!(
json.contains(r#""minzoom":0"#),
"vector_layers of an empty archive stay at minzoom 0: {json}"
);
}
#[test]
fn test_writer_with_compression_constructor() {
let writer = PmtilesWriter::with_compression(Compression::Brotli);
assert_eq!(writer.tile_compression(), Compression::Brotli);
assert_eq!(writer.internal_compression(), Compression::Brotli);
}
#[test]
fn test_writer_set_compression() {
let mut writer = PmtilesWriter::new();
assert_eq!(writer.tile_compression(), Compression::Gzip);
writer.set_tile_compression(Compression::Zstd);
assert_eq!(writer.tile_compression(), Compression::Zstd);
writer.set_internal_compression(Compression::Brotli);
assert_eq!(writer.internal_compression(), Compression::Brotli);
}
#[test]
fn test_writer_brotli_compression() {
let mut writer = PmtilesWriter::with_compression(Compression::Brotli);
let mvt_data = vec![0x1a; 100];
writer.add_tile(0, 0, 0, &mvt_data).unwrap();
writer.set_bounds(&TileBounds::new(-180.0, -85.0, 180.0, 85.0));
let path = Path::new("/tmp/test-pmtiles-brotli.pmtiles");
let _ = fs::remove_file(path);
writer
.write_to_file(path)
.expect("Should write file with brotli");
let data = fs::read(path).unwrap();
assert_eq!(&data[0..7], b"PMTiles");
assert_eq!(data[7], 3);
assert_eq!(data[97], Compression::Brotli as u8);
assert_eq!(data[98], Compression::Brotli as u8);
let _ = fs::remove_file(path);
}
#[test]
fn test_writer_zstd_compression() {
let mut writer = PmtilesWriter::with_compression(Compression::Zstd);
let mvt_data = vec![0x1a; 100];
writer.add_tile(0, 0, 0, &mvt_data).unwrap();
writer.set_bounds(&TileBounds::new(-180.0, -85.0, 180.0, 85.0));
let path = Path::new("/tmp/test-pmtiles-zstd.pmtiles");
let _ = fs::remove_file(path);
writer
.write_to_file(path)
.expect("Should write file with zstd");
let data = fs::read(path).unwrap();
assert_eq!(data[97], Compression::Zstd as u8);
assert_eq!(data[98], Compression::Zstd as u8);
let _ = fs::remove_file(path);
}
#[test]
fn test_writer_no_compression() {
let mut writer = PmtilesWriter::with_compression(Compression::None);
let mvt_data = vec![0x1a, 0x00];
writer.add_tile(0, 0, 0, &mvt_data).unwrap();
writer.set_bounds(&TileBounds::new(-180.0, -85.0, 180.0, 85.0));
let path = Path::new("/tmp/test-pmtiles-none.pmtiles");
let _ = fs::remove_file(path);
writer
.write_to_file(path)
.expect("Should write file without compression");
let data = fs::read(path).unwrap();
assert_eq!(data[97], Compression::None as u8);
assert_eq!(data[98], Compression::None as u8);
let _ = fs::remove_file(path);
}
#[test]
fn test_writer_mixed_compression() {
let mut writer = PmtilesWriter::new();
writer.set_internal_compression(Compression::Gzip);
writer.set_tile_compression(Compression::Zstd);
let mvt_data = vec![0x1a; 100];
writer.add_tile(0, 0, 0, &mvt_data).unwrap();
writer.set_bounds(&TileBounds::new(-180.0, -85.0, 180.0, 85.0));
let path = Path::new("/tmp/test-pmtiles-mixed.pmtiles");
let _ = fs::remove_file(path);
writer
.write_to_file(path)
.expect("Should write file with mixed compression");
let data = fs::read(path).unwrap();
assert_eq!(data[97], Compression::Gzip as u8); assert_eq!(data[98], Compression::Zstd as u8);
let _ = fs::remove_file(path);
}
#[test]
fn test_writer_dedup_identical_tiles() {
let mut writer = PmtilesWriter::new();
writer.enable_deduplication(true);
let ocean_tile = vec![0x1a, 0x00]; writer.add_tile(1, 0, 0, &ocean_tile).unwrap();
writer.add_tile(1, 0, 1, &ocean_tile).unwrap();
writer.add_tile(1, 1, 1, &ocean_tile).unwrap();
assert_eq!(writer.tile_count(), 3);
let stats = writer.dedup_stats();
assert_eq!(stats.total_tiles, 3);
assert_eq!(stats.unique_tiles, 1);
assert_eq!(stats.duplicates_eliminated, 2);
}
#[test]
fn test_writer_dedup_mixed_tiles() {
let mut writer = PmtilesWriter::new();
writer.enable_deduplication(true);
let tile_a = vec![0x1a, 0x01];
let tile_b = vec![0x1a, 0x02];
writer.add_tile(0, 0, 0, &tile_a).unwrap();
writer.add_tile(1, 0, 0, &tile_a).unwrap(); writer.add_tile(1, 0, 1, &tile_b).unwrap();
writer.add_tile(1, 1, 1, &tile_a).unwrap(); writer.add_tile(1, 1, 0, &tile_b).unwrap(); writer.add_tile(2, 0, 0, &tile_b).unwrap();
let stats = writer.dedup_stats();
assert_eq!(stats.total_tiles, 6);
assert_eq!(stats.unique_tiles, 2);
assert_eq!(stats.duplicates_eliminated, 4);
}
#[test]
fn test_writer_dedup_disabled_by_default() {
let writer = PmtilesWriter::new();
assert!(!writer.is_dedup_enabled());
}
#[test]
fn test_writer_dedup_file_size_reduction() {
let mut writer_dedup = PmtilesWriter::new();
writer_dedup.enable_deduplication(true);
let ocean_tile = vec![0x1a, 0x00, 0x01, 0x02, 0x03, 0x04, 0x05];
for z in 0..3u8 {
let n = 1u32 << z;
for x in 0..n {
for y in 0..n {
writer_dedup.add_tile(z, x, y, &ocean_tile).unwrap();
}
}
}
let path_dedup = Path::new("/tmp/test-pmtiles-dedup-enabled.pmtiles");
let _ = fs::remove_file(path_dedup);
writer_dedup.write_to_file(path_dedup).unwrap();
let size_dedup = fs::metadata(path_dedup).unwrap().len();
let mut writer_no_dedup = PmtilesWriter::new();
for z in 0..3u8 {
let n = 1u32 << z;
for x in 0..n {
for y in 0..n {
writer_no_dedup.add_tile(z, x, y, &ocean_tile).unwrap();
}
}
}
let path_no_dedup = Path::new("/tmp/test-pmtiles-dedup-disabled.pmtiles");
let _ = fs::remove_file(path_no_dedup);
writer_no_dedup.write_to_file(path_no_dedup).unwrap();
let size_no_dedup = fs::metadata(path_no_dedup).unwrap().len();
assert!(
size_dedup < size_no_dedup,
"Deduplicated file ({} bytes) should be smaller than non-deduplicated ({} bytes)",
size_dedup,
size_no_dedup
);
let _ = fs::remove_file(path_dedup);
let _ = fs::remove_file(path_no_dedup);
}
#[test]
fn test_writer_dedup_run_length_consecutive() {
let mut writer = PmtilesWriter::new();
writer.enable_deduplication(true);
let tile = vec![0x1a, 0x10];
writer.add_tile(1, 0, 0, &tile).unwrap(); writer.add_tile(1, 0, 1, &tile).unwrap(); writer.add_tile(1, 1, 1, &tile).unwrap(); writer.add_tile(1, 1, 0, &tile).unwrap();
let path = Path::new("/tmp/test-pmtiles-runlength.pmtiles");
let _ = fs::remove_file(path);
writer.write_to_file(path).unwrap();
let data = fs::read(path).unwrap();
let addressed_count = u64::from_le_bytes(data[72..80].try_into().unwrap());
let entries_count = u64::from_le_bytes(data[80..88].try_into().unwrap());
let contents_count = u64::from_le_bytes(data[88..96].try_into().unwrap());
assert_eq!(addressed_count, 4);
assert_eq!(entries_count, 1);
assert_eq!(contents_count, 1);
let _ = fs::remove_file(path);
}
#[test]
fn test_writer_dedup_header_stats() {
let mut writer = PmtilesWriter::new();
writer.enable_deduplication(true);
let tile_a = vec![0x1a, 0x01];
let tile_b = vec![0x1a, 0x02];
let tile_c = vec![0x1a, 0x03];
for _ in 0..3 {
writer.add_tile(0, 0, 0, &tile_a).unwrap(); }
let mut writer2 = PmtilesWriter::new();
writer2.enable_deduplication(true);
writer2.add_tile(0, 0, 0, &tile_a).unwrap();
writer2.add_tile(1, 0, 0, &tile_a).unwrap(); writer2.add_tile(1, 0, 1, &tile_a).unwrap(); writer2.add_tile(1, 1, 1, &tile_b).unwrap();
writer2.add_tile(1, 1, 0, &tile_b).unwrap(); writer2.add_tile(2, 0, 0, &tile_c).unwrap();
writer2.add_tile(2, 0, 1, &tile_a).unwrap(); writer2.add_tile(2, 1, 0, &tile_b).unwrap(); writer2.add_tile(2, 1, 1, &tile_c).unwrap(); writer2.add_tile(2, 2, 0, &tile_c).unwrap();
let path = Path::new("/tmp/test-pmtiles-header-stats.pmtiles");
let _ = fs::remove_file(path);
writer2.write_to_file(path).unwrap();
let data = fs::read(path).unwrap();
let addressed_count = u64::from_le_bytes(data[72..80].try_into().unwrap());
let contents_count = u64::from_le_bytes(data[88..96].try_into().unwrap());
assert_eq!(addressed_count, 10, "Should address 10 tiles");
assert_eq!(contents_count, 3, "Should have 3 unique contents");
let _ = fs::remove_file(path);
}
#[test]
fn test_streaming_writer_creates_temp_file() {
let writer =
StreamingPmtilesWriter::new(Compression::Gzip).expect("Should create streaming writer");
assert!(
writer.temp_path().exists(),
"Temp file should be created at {:?}",
writer.temp_path()
);
let temp_path = writer.temp_path().to_path_buf();
drop(writer);
assert!(
!temp_path.exists(),
"Temp file should be cleaned up on drop"
);
}
#[test]
fn test_streaming_writer_add_tile_writes_to_temp() {
let mut writer =
StreamingPmtilesWriter::new(Compression::Gzip).expect("Should create streaming writer");
let mvt_data = vec![0x1a, 0x00, 0x01, 0x02];
writer.add_tile(0, 0, 0, &mvt_data).unwrap();
let stats = writer.stats();
assert_eq!(stats.total_tiles, 1);
assert_eq!(stats.unique_tiles, 1);
assert!(
stats.bytes_written > 0,
"Should have written bytes to temp file"
);
assert_eq!(
stats.bytes_written, writer.current_offset,
"bytes_written should match current_offset"
);
}
#[test]
fn test_streaming_writer_dedup_same_content() {
let mut writer =
StreamingPmtilesWriter::new(Compression::Gzip).expect("Should create streaming writer");
let ocean_tile = vec![0x1a, 0x00];
writer.add_tile(1, 0, 0, &ocean_tile).unwrap();
writer.add_tile(1, 0, 1, &ocean_tile).unwrap();
writer.add_tile(1, 1, 1, &ocean_tile).unwrap();
let stats = writer.stats();
assert_eq!(stats.total_tiles, 3, "Should track 3 total tiles");
assert_eq!(stats.unique_tiles, 1, "Should only have 1 unique tile");
assert!(
stats.bytes_saved_dedup > 0,
"Should have saved bytes via deduplication"
);
}
#[test]
fn test_streaming_writer_finalize_creates_valid_pmtiles() {
let mut writer =
StreamingPmtilesWriter::new(Compression::Gzip).expect("Should create streaming writer");
writer.add_tile(0, 0, 0, &[0x1a, 0x00]).unwrap();
writer.add_tile(1, 0, 0, &[0x1a, 0x01]).unwrap();
writer.add_tile(1, 0, 1, &[0x1a, 0x02]).unwrap();
writer.set_bounds(&TileBounds::new(-180.0, -85.0, 180.0, 85.0));
let output_path = Path::new("/tmp/test-streaming-pmtiles.pmtiles");
let _ = fs::remove_file(output_path);
let stats = writer.finalize(output_path).expect("Should finalize");
assert!(output_path.exists(), "Output file should exist");
let data = fs::read(output_path).unwrap();
assert_eq!(&data[0..7], b"PMTiles", "Should have PMTiles magic");
assert_eq!(data[7], 3, "Should be version 3");
let addressed_count = u64::from_le_bytes(data[72..80].try_into().unwrap());
assert_eq!(addressed_count, 3, "Should have 3 addressed tiles");
assert_eq!(stats.total_tiles, 3);
assert_eq!(stats.unique_tiles, 3);
let _ = fs::remove_file(output_path);
}
#[test]
fn test_streaming_writer_memory_bounded() {
let mut writer =
StreamingPmtilesWriter::new(Compression::Gzip).expect("Should create streaming writer");
let mut count = 0;
for z in 0..10u8 {
let max_coord = 1u32 << z; for x in 0..max_coord.min(10) {
for y in 0..max_coord.min(10) {
let data = vec![0x1a, z, (x & 0xFF) as u8, (y & 0xFF) as u8, count as u8];
writer.add_tile(z, x, y, &data).unwrap();
count += 1;
if count >= 1000 {
break;
}
}
if count >= 1000 {
break;
}
}
if count >= 1000 {
break;
}
}
let stats = writer.stats();
let estimated_mem = stats.estimated_memory_bytes();
assert!(
estimated_mem < 200_000, "Memory usage should be bounded, got {} bytes",
estimated_mem
);
}
#[test]
fn test_streaming_writer_matches_non_streaming_output() {
let tiles_data = vec![
(0, 0, 0, vec![0x1a, 0x00]),
(1, 0, 0, vec![0x1a, 0x01]),
(1, 0, 1, vec![0x1a, 0x02]),
(1, 1, 1, vec![0x1a, 0x00]), ];
let mut non_streaming = PmtilesWriter::with_compression(Compression::Gzip);
non_streaming.enable_deduplication(true);
non_streaming.set_layer_name("test");
non_streaming.set_bounds(&TileBounds::new(-180.0, -85.0, 180.0, 85.0));
for (z, x, y, data) in &tiles_data {
non_streaming.add_tile(*z, *x, *y, data).unwrap();
}
let non_streaming_path = Path::new("/tmp/test-compare-non-streaming.pmtiles");
let _ = fs::remove_file(non_streaming_path);
non_streaming.write_to_file(non_streaming_path).unwrap();
let mut streaming = StreamingPmtilesWriter::new(Compression::Gzip).unwrap();
streaming.set_layer_name("test");
streaming.set_bounds(&TileBounds::new(-180.0, -85.0, 180.0, 85.0));
for (z, x, y, data) in &tiles_data {
streaming.add_tile(*z, *x, *y, data).unwrap();
}
let streaming_path = Path::new("/tmp/test-compare-streaming.pmtiles");
let _ = fs::remove_file(streaming_path);
streaming.finalize(streaming_path).unwrap();
let ns_data = fs::read(non_streaming_path).unwrap();
let s_data = fs::read(streaming_path).unwrap();
assert_eq!(
&ns_data[0..8],
&s_data[0..8],
"Header magic/version should match"
);
let ns_addressed = u64::from_le_bytes(ns_data[72..80].try_into().unwrap());
let s_addressed = u64::from_le_bytes(s_data[72..80].try_into().unwrap());
assert_eq!(ns_addressed, s_addressed, "Addressed tiles should match");
let ns_contents = u64::from_le_bytes(ns_data[88..96].try_into().unwrap());
let s_contents = u64::from_le_bytes(s_data[88..96].try_into().unwrap());
assert_eq!(ns_contents, s_contents, "Unique contents should match");
let _ = fs::remove_file(non_streaming_path);
let _ = fs::remove_file(streaming_path);
}
#[test]
fn test_streaming_writer_with_feature_count() {
use flate2::read::GzDecoder;
use std::io::Read;
let mut writer = StreamingPmtilesWriter::new(Compression::Gzip).unwrap();
writer.set_layer_name("buildings");
writer
.add_tile_with_count(0, 0, 0, &[0x1a, 0x00], 100)
.unwrap();
writer
.add_tile_with_count(1, 0, 0, &[0x1a, 0x01], 50)
.unwrap();
writer
.add_tile_with_count(1, 0, 1, &[0x1a, 0x02], 75)
.unwrap();
let output_path = Path::new("/tmp/test-streaming-features.pmtiles");
let _ = fs::remove_file(output_path);
writer.finalize(output_path).unwrap();
let data = fs::read(output_path).unwrap();
let metadata_offset = u64::from_le_bytes(data[24..32].try_into().unwrap()) as usize;
let metadata_length = u64::from_le_bytes(data[32..40].try_into().unwrap()) as usize;
let compressed_metadata = &data[metadata_offset..metadata_offset + metadata_length];
let mut decoder = GzDecoder::new(compressed_metadata);
let mut metadata_json = String::new();
decoder.read_to_string(&mut metadata_json).unwrap();
assert!(
metadata_json.contains("\"count\":225"),
"Should have total feature count 225, got: {}",
metadata_json
);
let _ = fs::remove_file(output_path);
}
const INITIAL_FETCH_SIZE: usize = 16384;
const HEADER_SIZE: usize = 127;
const MAX_ROOT_DIR_SIZE: usize = INITIAL_FETCH_SIZE - HEADER_SIZE;
#[test]
fn test_large_archive_uses_leaf_directories() {
let num_tiles = 10_000;
let mut writer = StreamingPmtilesWriter::new(Compression::Gzip).unwrap();
writer.set_layer_name("test");
writer.set_bounds(&TileBounds::new(-180.0, -85.0, 180.0, 85.0));
for i in 0..num_tiles {
let x = i % 4096;
let y = i / 4096;
let data = vec![0x1a, (i & 0xff) as u8, ((i >> 8) & 0xff) as u8];
writer.add_tile(12, x as u32, y as u32, &data).unwrap();
}
let output_path = Path::new("/tmp/test-leaf-directories.pmtiles");
let _ = fs::remove_file(output_path);
writer.finalize(output_path).unwrap();
let data = fs::read(output_path).unwrap();
let root_dir_length = u64::from_le_bytes(data[16..24].try_into().unwrap()) as usize;
let leaf_dirs_offset = u64::from_le_bytes(data[40..48].try_into().unwrap());
let leaf_dirs_length = u64::from_le_bytes(data[48..56].try_into().unwrap());
eprintln!(
"Archive stats: root_dir_length={}, leaf_dirs_offset={}, leaf_dirs_length={}, MAX={}",
root_dir_length, leaf_dirs_offset, leaf_dirs_length, MAX_ROOT_DIR_SIZE
);
assert!(
root_dir_length <= MAX_ROOT_DIR_SIZE,
"Root directory ({} bytes) must fit in initial fetch ({} bytes)",
root_dir_length,
MAX_ROOT_DIR_SIZE
);
assert!(
leaf_dirs_offset > 0,
"Large archive should have leaf directories (offset={})",
leaf_dirs_offset
);
assert!(
leaf_dirs_length > 0,
"Large archive should have leaf directories (length={})",
leaf_dirs_length
);
let _ = fs::remove_file(output_path);
}
#[test]
fn test_small_archive_no_leaf_directories() {
let mut writer = StreamingPmtilesWriter::new(Compression::Gzip).unwrap();
writer.set_layer_name("test");
for i in 0..10 {
let data = vec![0x1a, i as u8];
writer.add_tile(0, 0, 0, &data).unwrap();
}
let output_path = Path::new("/tmp/test-no-leaf-directories.pmtiles");
let _ = fs::remove_file(output_path);
writer.finalize(output_path).unwrap();
let data = fs::read(output_path).unwrap();
let leaf_dirs_offset = u64::from_le_bytes(data[40..48].try_into().unwrap());
let leaf_dirs_length = u64::from_le_bytes(data[48..56].try_into().unwrap());
assert_eq!(
leaf_dirs_length, 0,
"Small archive should not have leaf directories"
);
let metadata_offset = u64::from_le_bytes(data[24..32].try_into().unwrap());
let metadata_length = u64::from_le_bytes(data[32..40].try_into().unwrap());
let tile_data_offset = u64::from_le_bytes(data[56..64].try_into().unwrap());
assert_ne!(
leaf_dirs_offset, 0,
"leaf_dirs_offset must never be 0 -- go-pmtiles rejects the archive"
);
assert_eq!(leaf_dirs_offset, metadata_offset + metadata_length);
assert_eq!(leaf_dirs_offset, tile_data_offset);
let _ = fs::remove_file(output_path);
}
#[test]
fn test_leaf_directory_entries_have_run_length_zero() {
let num_tiles = 3000;
let mut writer = StreamingPmtilesWriter::new(Compression::Gzip).unwrap();
writer.set_layer_name("test");
for i in 0..num_tiles {
let x = i % 1024;
let y = i / 1024;
let data = vec![0x1a, (i & 0xff) as u8];
writer.add_tile(10, x as u32, y as u32, &data).unwrap();
}
let output_path = Path::new("/tmp/test-leaf-run-length.pmtiles");
let _ = fs::remove_file(output_path);
writer.finalize(output_path).unwrap();
let data = fs::read(output_path).unwrap();
let root_dir_offset = u64::from_le_bytes(data[8..16].try_into().unwrap()) as usize;
let root_dir_length = u64::from_le_bytes(data[16..24].try_into().unwrap()) as usize;
let leaf_dirs_length = u64::from_le_bytes(data[48..56].try_into().unwrap());
if leaf_dirs_length > 0 {
let compressed_root = &data[root_dir_offset..root_dir_offset + root_dir_length];
use flate2::read::GzDecoder;
use std::io::Read;
let mut decoder = GzDecoder::new(compressed_root);
let mut decompressed = Vec::new();
decoder.read_to_end(&mut decompressed).unwrap();
let entries = decode_directory(&decompressed).unwrap();
for entry in &entries {
assert_eq!(
entry.run_length, 0,
"Root directory entries pointing to leaves must have run_length=0, got {}",
entry.run_length
);
}
}
let _ = fs::remove_file(output_path);
}
#[test]
fn add_tile_precompressed_matches_add_tile_with_count() {
let tiles: Vec<(u8, u32, u32, Vec<u8>)> = vec![
(0, 0, 0, vec![0x1a, 0x05, b'h', b'e', b'l', b'l', b'o']),
(1, 0, 0, vec![0x1a, 0x03, b'a', b'b', b'c']),
(1, 0, 1, vec![0x1a, 0x03, b'x', b'y', b'z']),
(1, 1, 1, vec![0x1a, 0x05, b'h', b'e', b'l', b'l', b'o']), ];
let dir = std::env::temp_dir();
let mut serial = StreamingPmtilesWriter::new(Compression::Gzip).unwrap();
serial.set_layer_name("test");
serial.set_bounds(&TileBounds::new(-180.0, -85.0, 180.0, 85.0));
for (z, x, y, data) in &tiles {
serial.add_tile_with_count(*z, *x, *y, data, 1).unwrap();
}
let serial_path = dir.join("gpq-227-serial.pmtiles");
let _ = fs::remove_file(&serial_path);
serial.finalize(&serial_path).unwrap();
let mut parallel = StreamingPmtilesWriter::new(Compression::Gzip).unwrap();
parallel.set_layer_name("test");
parallel.set_bounds(&TileBounds::new(-180.0, -85.0, 180.0, 85.0));
for (z, x, y, data) in &tiles {
let hash = TileHasher::hash(data);
let compressed = compression::compress(data, Compression::Gzip).unwrap();
parallel
.add_tile_precompressed(*z, *x, *y, hash, &compressed, data.len(), 1)
.unwrap();
}
let parallel_path = dir.join("gpq-227-parallel.pmtiles");
let _ = fs::remove_file(¶llel_path);
parallel.finalize(¶llel_path).unwrap();
let a = fs::read(&serial_path).unwrap();
let b = fs::read(¶llel_path).unwrap();
assert_eq!(
a, b,
"precompressed archive must be byte-identical to serial compression"
);
let _ = fs::remove_file(&serial_path);
let _ = fs::remove_file(¶llel_path);
}
#[test]
fn checkpoint_produces_valid_capped_pmtiles() {
let mut writer = StreamingPmtilesWriter::new(Compression::Gzip).unwrap();
writer.set_layer_name("test");
writer.set_bounds(&TileBounds::new(-180.0, -85.0, 180.0, 85.0));
writer.add_tile(0, 0, 0, &[0x1a, 0x00]).unwrap();
writer.add_tile(1, 0, 0, &[0x1a, 0x01]).unwrap();
writer.add_tile(1, 0, 1, &[0x1a, 0x02]).unwrap();
let ckpt_path = std::env::temp_dir().join("gpq-229-checkpoint.pmtiles");
let _ = fs::remove_file(&ckpt_path);
writer
.checkpoint(&ckpt_path)
.expect("checkpoint should succeed");
let ck = fs::read(&ckpt_path).unwrap();
assert_eq!(&ck[0..7], b"PMTiles", "checkpoint has PMTiles magic");
assert_eq!(ck[7], 3, "checkpoint is version 3");
assert_eq!(ck[100], 0, "checkpoint min_zoom == 0");
assert_eq!(
ck[101], 1,
"checkpoint max_zoom capped at last finished level"
);
let ck_addressed = u64::from_le_bytes(ck[72..80].try_into().unwrap());
assert_eq!(ck_addressed, 3, "checkpoint holds the 3 finished tiles");
writer.add_tile(2, 0, 0, &[0x1a, 0x03]).unwrap();
let final_path = std::env::temp_dir().join("gpq-229-checkpoint-final.pmtiles");
let _ = fs::remove_file(&final_path);
writer
.finalize(&final_path)
.expect("finalize after checkpoint");
let fin = fs::read(&final_path).unwrap();
assert_eq!(fin[101], 2, "final max_zoom includes the finer level");
let fin_addressed = u64::from_le_bytes(fin[72..80].try_into().unwrap());
assert_eq!(fin_addressed, 4, "final holds all 4 tiles");
let _ = fs::remove_file(&ckpt_path);
let _ = fs::remove_file(&final_path);
}
#[test]
fn checkpoint_then_finalize_byte_identical_to_no_checkpoint() {
let tiles: Vec<(u8, u32, u32, Vec<u8>)> = vec![
(0, 0, 0, vec![0x1a, 0x05, b'h', b'e', b'l', b'l', b'o']),
(1, 0, 0, vec![0x1a, 0x03, b'a', b'b', b'c']),
(1, 0, 1, vec![0x1a, 0x03, b'x', b'y', b'z']),
(1, 1, 1, vec![0x1a, 0x05, b'h', b'e', b'l', b'l', b'o']), ];
let dir = std::env::temp_dir();
let mut plain = StreamingPmtilesWriter::new(Compression::Gzip).unwrap();
plain.set_layer_name("test");
plain.set_bounds(&TileBounds::new(-180.0, -85.0, 180.0, 85.0));
for (z, x, y, data) in &tiles {
plain.add_tile(*z, *x, *y, data).unwrap();
}
let plain_path = dir.join("gpq-229-plain.pmtiles");
let _ = fs::remove_file(&plain_path);
plain.finalize(&plain_path).unwrap();
let mut ckpt = StreamingPmtilesWriter::new(Compression::Gzip).unwrap();
ckpt.set_layer_name("test");
ckpt.set_bounds(&TileBounds::new(-180.0, -85.0, 180.0, 85.0));
let ckpt_path = dir.join("gpq-229-intermediate.pmtiles");
let final_path = dir.join("gpq-229-checkpointed-final.pmtiles");
let _ = fs::remove_file(&final_path);
for (i, (z, x, y, data)) in tiles.iter().enumerate() {
ckpt.add_tile(*z, *x, *y, data).unwrap();
if i == 0 {
let _ = fs::remove_file(&ckpt_path);
ckpt.checkpoint(&ckpt_path).unwrap();
}
}
ckpt.finalize(&final_path).unwrap();
let a = fs::read(&plain_path).unwrap();
let b = fs::read(&final_path).unwrap();
assert_eq!(
a, b,
"checkpointed run must be byte-identical to the non-checkpointed finalize"
);
let _ = fs::remove_file(&plain_path);
let _ = fs::remove_file(&ckpt_path);
let _ = fs::remove_file(&final_path);
}
}