use std::collections::BTreeMap;
use std::collections::HashMap;
use std::collections::HashSet;
use std::fs::File;
use std::io::BufWriter;
use std::path::Path;
use std::sync::Arc;
use std::time::{Duration, Instant};
use arrow_array::cast::AsArray;
use arrow_array::types::{
Decimal128Type, Float32Type, Float64Type, Int16Type, Int32Type, Int64Type, Int8Type,
UInt16Type, UInt32Type, UInt64Type, UInt8Type,
};
use arrow_array::{Array, RecordBatch};
use arrow_cast::display::{ArrayFormatter, FormatOptions};
use arrow_schema::{DataType, Schema};
use crossbeam_channel::{Receiver, Sender};
use geo::{BoundingRect, CoordsIter, Geometry, MapCoords};
use geoarrow::array::from_arrow_array;
use geoarrow_array::GeoArrowArray;
use prost::Message;
use rayon::prelude::*;
use serde::Serialize;
use tempfile::NamedTempFile;
use crate::batch_processor::extract_geometries_from_array;
use crate::clip::{clip_geometry_simple, geometry_is_simple};
use crate::compression::{self, Compression};
use crate::dedup::TileHasher;
use crate::mvt::{LayerBuilder, PropertyValue, TileBuilder};
use crate::pmtiles_writer::StreamingPmtilesWriter;
use crate::tile::{tile_ranges_for_bbox, BboxTileRanges, TileBounds, TileCoord};
use super::coalesce::COALESCED_COUNT_COLUMN;
use super::level::{zoom_for_gsd, Crs, Mode, OverviewsMeta};
use super::pipe::scoped_pipe;
use super::pipeline::{auto_backing, available_memory_bytes, SinkBacking};
use super::properties::PropertySelection;
use super::reader::{OverviewReader, ReaderError};
use super::writer::LEVEL_COLUMN;
const DEFAULT_EXTENT: u32 = 4096;
const DEFAULT_TILE_BUFFER_PX: u32 = 8;
const NOMINAL_TILE_PIXELS: f64 = 256.0;
#[inline]
fn buffer_fraction(opts: &ExportOptions) -> f64 {
f64::from(opts.tile_buffer) / NOMINAL_TILE_PIXELS
}
pub const DEFAULT_TILE_SIZE_LIMIT: usize = 500 * 1024;
const WEBMERC_HALF_M: f64 = 20_037_508.342_789_244;
#[derive(Debug, Clone)]
pub struct ExportOptions {
pub layer_name: String,
pub tile_buffer: u32,
pub extent: u32,
pub tile_size_limit: Option<usize>,
pub simple_clip_fastpath: bool,
pub partition_wave: usize,
pub feature_order: FeatureOrder,
pub min_zoom: Option<u8>,
pub properties: PropertySelection,
}
impl Default for ExportOptions {
fn default() -> Self {
Self {
layer_name: "overview".to_string(),
tile_buffer: DEFAULT_TILE_BUFFER_PX,
extent: DEFAULT_EXTENT,
tile_size_limit: Some(DEFAULT_TILE_SIZE_LIMIT),
simple_clip_fastpath: true,
partition_wave: PARTITION_WAVE_AUTO,
feature_order: FeatureOrder::default(),
min_zoom: None,
properties: PropertySelection::default(),
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize)]
pub struct ZoomReport {
pub zoom: u8,
pub level: usize,
pub level_feature_count: usize,
pub tile_count: usize,
pub tile_feature_count: usize,
pub oversized_tiles: usize,
}
#[derive(Debug, Clone, PartialEq, Serialize)]
pub struct ExportReport {
pub mode: String,
pub min_zoom: u8,
pub max_zoom: u8,
pub zooms: Vec<ZoomReport>,
pub total_tiles: usize,
pub total_tile_features: usize,
pub oversized_tiles: usize,
pub duration_secs: f64,
}
#[derive(Debug, thiserror::Error)]
pub enum ExportError {
#[error("overview reader error: {0}")]
Reader(#[from] ReaderError),
#[error("io error: {0}")]
Io(#[from] std::io::Error),
#[error("arrow error: {0}")]
Arrow(#[from] arrow_schema::ArrowError),
#[error("{0}")]
Core(#[from] crate::Error),
#[error("unsupported input CRS {crs:?}: export requires EPSG:4326 or EPSG:3857")]
UnsupportedCrs {
crs: String,
},
#[error("overview file has no geometry column")]
NoGeometryColumn,
#[error(
"declared min_zoom {declared} is finer than the coarsest level present (zoom \
{coarsest}): the header can widen the zoom range over empty zooms, not narrow \
it over real tiles (#380)"
)]
DeclaredMinZoomTooFine { declared: u8, coarsest: u8 },
#[error(
"included property {name:?} is not a property this overview file exports \
(exportable: {available})"
)]
UnknownProperty { name: String, available: String },
#[error(
"property {name:?} is excluded but {knob} reads it; keep it in the selection or \
drop the knob"
)]
PropertyRequiredByKnob { name: String, knob: String },
}
#[cfg(test)]
struct Feature {
geom: Geometry<f64>,
props: Vec<(String, PropertyValue)>,
}
struct EncodedTile {
x: u32,
y: u32,
data: Vec<u8>,
hash: u64,
raw_len: usize,
feature_count: usize,
oversized: bool,
}
pub fn zoom_for_level(meta: &OverviewsMeta, level_idx: usize) -> u8 {
let level = &meta.levels[level_idx];
if let Some(z) = level.zoom {
return z;
}
let z = zoom_for_gsd(level.gsd).round();
z.clamp(0.0, 255.0) as u8
}
const DEFAULT_PARTITION_TARGET: usize = 32_768;
pub const PARTITION_WAVE_AUTO: usize = 0;
pub const PARTITION_WAVE_MIN: usize = 6;
pub const PARTITION_WAVE_FALLBACK_MAX: usize = 16;
const EXPORT_WAVE_RAM_FRACTION: f64 = 0.5;
const PARTITION_SLOT_TRANSIENT_BYTES: u64 = 64 * 1024 * 1024;
const MEMBER_MEMORY_INFLATION: u64 = 2;
fn memory_safe_wave(available_ram_bytes: Option<u64>) -> usize {
match available_ram_bytes {
Some(ram) => {
let budget = ((ram as f64) * EXPORT_WAVE_RAM_FRACTION) as u64;
(budget / PARTITION_SLOT_TRANSIENT_BYTES) as usize
}
None => PARTITION_WAVE_FALLBACK_MAX,
}
}
fn auto_partition_wave(cores: usize, available_ram_bytes: Option<u64>) -> usize {
cores
.min(memory_safe_wave(available_ram_bytes))
.max(PARTITION_WAVE_MIN)
}
fn memory_safe_level_wave(
ceiling: usize,
partitions: &[Partition],
mean_member_bytes: Option<u64>,
available_ram_bytes: Option<u64>,
) -> usize {
let (Some(mean), Some(ram)) = (mean_member_bytes.filter(|&b| b > 0), available_ram_bytes)
else {
return ceiling;
};
let max_members = partitions.iter().map(|p| p.members).max().unwrap_or(0);
if max_members == 0 {
return ceiling;
}
let per_slot = (max_members as u64)
.saturating_mul(mean)
.saturating_mul(MEMBER_MEMORY_INFLATION)
.max(1);
let budget = ((ram as f64) * EXPORT_WAVE_RAM_FRACTION) as u64;
let safe = (budget / per_slot).max(1) as usize;
ceiling.min(safe)
}
pub fn resolve_partition_wave(requested: usize) -> usize {
if requested == PARTITION_WAVE_AUTO {
let cores = std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(PARTITION_WAVE_MIN);
auto_partition_wave(cores, available_memory_bytes())
} else {
requested
}
}
fn resolve_and_log_partition_wave(requested: usize) -> usize {
let wave = resolve_partition_wave(requested);
let cores = std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(0);
if requested == PARTITION_WAVE_AUTO {
let available = available_memory_bytes();
let avail_str = available.map_or_else(
|| "unknown".to_string(),
|b| format!("{} MiB", b / (1024 * 1024)),
);
log::info!(
"[export] partition wave: {wave} partition(s) per band read \
(auto: {cores} core(s), memory-safe cap {}, avail {avail_str})",
memory_safe_wave(available),
);
} else {
log::info!(
"[export] partition wave: {wave} partition(s) per band read \
(explicit; {cores} core(s) detected)"
);
}
wave
}
const EXPORT_BATCH_SIZE: usize = 8_192;
const CHECKPOINT_INTERVAL: Duration = Duration::from_secs(60);
const WAVE_LOG_INTERVAL: Duration = Duration::from_secs(30);
pub fn export_pmtiles(
input_path: impl AsRef<Path>,
output_path: impl AsRef<Path>,
options: &ExportOptions,
) -> Result<ExportReport, ExportError> {
export_pmtiles_with_partition_target(input_path, output_path, options, DEFAULT_PARTITION_TARGET)
}
fn export_pmtiles_with_partition_target(
input_path: impl AsRef<Path>,
output_path: impl AsRef<Path>,
options: &ExportOptions,
partition_target: usize,
) -> Result<ExportReport, ExportError> {
export_pmtiles_impl(
input_path,
output_path,
options,
partition_target,
false,
None,
)
}
#[allow(clippy::too_many_arguments)]
fn plan_levels(
scans: &[LevelScan],
meta: &OverviewsMeta,
num_levels: usize,
partition_target: usize,
ceiling_wave: usize,
auto_wave: bool,
mean_member_bytes: Option<u64>,
available_ram: Option<u64>,
) -> Vec<LevelPlan> {
scans
.iter()
.enumerate()
.map(|(level_idx, scan)| {
let zoom = zoom_for_level(meta, level_idx);
let partitions = plan_partitions(&scan.tile_counts, zoom, partition_target);
let wave = if auto_wave {
let w = memory_safe_level_wave(
ceiling_wave,
&partitions,
mean_member_bytes,
available_ram,
);
if w < ceiling_wave {
let densest = partitions.iter().map(|p| p.members).max().unwrap_or(0);
let budget_mib = available_ram
.map(|r| ((r as f64 * EXPORT_WAVE_RAM_FRACTION) as u64) / (1024 * 1024))
.unwrap_or(0);
log::info!(
"[export] level {}/{num_levels} z{zoom}: wave {ceiling_wave} → {w} \
(densest partition {densest} members × ~{} B/member × {MEMBER_MEMORY_INFLATION}, \
{budget_mib} MiB budget) — #311 density guard",
level_idx + 1,
mean_member_bytes.unwrap_or(0),
);
}
w
} else {
ceiling_wave
};
LevelPlan {
zoom,
partitions,
wave,
}
})
.collect()
}
fn export_pmtiles_impl(
input_path: impl AsRef<Path>,
output_path: impl AsRef<Path>,
options: &ExportOptions,
partition_target: usize,
force_legacy_pass2: bool,
backing_override: Option<SinkBacking>,
) -> Result<ExportReport, ExportError> {
let start = Instant::now();
let input_path = input_path.as_ref();
let reader = OverviewReader::open(input_path)?;
let meta = reader.meta().clone();
let crs = detect_crs(input_path)?;
let num_levels = reader.num_levels();
let coarsest_zoom = zoom_for_level(&meta, 0);
let max_zoom = zoom_for_level(&meta, num_levels - 1);
let min_zoom = match options.min_zoom {
Some(declared) if declared > coarsest_zoom => {
return Err(ExportError::DeclaredMinZoomTooFine {
declared,
coarsest: coarsest_zoom,
});
}
Some(declared) if declared < coarsest_zoom => {
log::info!(
"[export] widening the declared zoom range to z{declared}..z{max_zoom}; \
the file's coarsest level is z{coarsest_zoom}, so z{declared}..z{} hold \
no tiles",
coarsest_zoom - 1
);
declared
}
_ => coarsest_zoom,
};
let ceiling_wave = resolve_and_log_partition_wave(options.partition_wave);
let auto_wave = options.partition_wave == PARTITION_WAVE_AUTO;
let mean_member_bytes = reader.finest_level_mean_row_bytes();
let available_ram = available_memory_bytes();
let mut writer = StreamingPmtilesWriter::new(Compression::Gzip)?;
writer.set_layer_name(&options.layer_name);
writer.set_declared_min_zoom(min_zoom);
let published = PublishedNames::from_reader(&reader).with_selection(
reader.schema(),
geometry_index(reader.schema()).ok_or(ExportError::NoGeometryColumn)?,
&options.properties,
&options.feature_order,
)?;
warn_if_order_column_is_absent(reader.schema(), &published, &options.feature_order);
writer.set_fields(field_metadata(
reader.schema(),
geometry_index(reader.schema()),
&published,
));
let t_scan = Instant::now();
let scans = scan_all_levels(&reader, crs, &meta, options)?;
log::info!(
"[export] scan complete: {num_levels} levels, single read, {:.2}s",
t_scan.elapsed().as_secs_f64()
);
let mut overall_bounds: Option<TileBounds> = None;
for scan in &scans {
if let Some(b) = &scan.bounds {
match &mut overall_bounds {
Some(acc) => acc.expand(b),
None => overall_bounds = Some(*b),
}
}
}
if let Some(b) = &overall_bounds {
writer.set_bounds(b);
}
let mut last_checkpoint = Instant::now();
let mut zooms: Vec<ZoomReport> = Vec::with_capacity(num_levels);
let plans = plan_levels(
&scans,
&meta,
num_levels,
partition_target,
ceiling_wave,
auto_wave,
mean_member_bytes,
available_ram,
);
let single_read = matches!(reader.mode(), Mode::Partitioning) && !force_legacy_pass2;
let mut store = if single_read {
let buffered_rows: usize = scans.iter().map(|s| s.feature_count).sum();
let backing = backing_override.unwrap_or_else(|| {
let b = auto_backing(
Mode::Partitioning,
buffered_rows,
available_ram,
mean_member_bytes,
);
log::info!(
"[export] pass2 single-read fan-out (#235): ~{buffered_rows} buffered \
member row(s) across {num_levels} level(s) → {b:?}"
);
b
});
Some(fill_member_store(
&reader,
&FanoutCtx {
crs,
plans: &plans,
opts: options,
published: &published,
},
backing,
)?)
} else {
None
};
for (level_idx, plan) in plans.iter().enumerate() {
zooms.push(export_level(
&mut writer,
store.as_mut(),
&ExportLevelCtx {
reader: &reader,
scan: &scans[level_idx],
plan,
level_idx,
num_levels,
crs,
published: &published,
options,
start,
},
)?);
if level_idx + 1 < num_levels && last_checkpoint.elapsed() >= CHECKPOINT_INTERVAL {
let t_ckpt = Instant::now();
writer.checkpoint(output_path.as_ref())?;
last_checkpoint = Instant::now();
log::info!(
"[export] checkpoint written: zooms {}..={} salvageable ({:.2}s)",
zoom_for_level(&meta, 0),
plan.zoom,
t_ckpt.elapsed().as_secs_f64(),
);
}
}
let t_finalize = Instant::now();
writer.finalize(output_path.as_ref())?;
log::info!(
"[export] finalize complete: {:.2}s (total {:.0}s)",
t_finalize.elapsed().as_secs_f64(),
start.elapsed().as_secs_f64(),
);
let total_tiles = zooms.iter().map(|z| z.tile_count).sum();
let total_tile_features = zooms.iter().map(|z| z.tile_feature_count).sum();
let oversized_tiles = zooms.iter().map(|z| z.oversized_tiles).sum();
Ok(ExportReport {
mode: format!("{:?}", reader.mode()).to_lowercase(),
min_zoom,
max_zoom,
zooms,
total_tiles,
total_tile_features,
oversized_tiles,
duration_secs: start.elapsed().as_secs_f64(),
})
}
struct ExportLevelCtx<'a> {
reader: &'a OverviewReader,
scan: &'a LevelScan,
plan: &'a LevelPlan,
level_idx: usize,
num_levels: usize,
crs: Crs,
published: &'a PublishedNames,
options: &'a ExportOptions,
start: Instant,
}
fn export_level(
writer: &mut StreamingPmtilesWriter,
mut store: Option<&mut MemberStore>,
level: &ExportLevelCtx<'_>,
) -> Result<ZoomReport, ExportError> {
let ExportLevelCtx {
reader,
scan,
plan,
level_idx,
num_levels,
crs,
published,
options,
start,
} = *level;
let zoom = plan.zoom;
let partitions = &plan.partitions;
let partition_wave = plan.wave;
let t_tiles = Instant::now();
let ctx = LevelCtx {
reader,
level_idx,
crs,
zoom,
opts: options,
published,
};
let mut tile_count = 0usize;
let mut tile_feature_count = 0usize;
let mut oversized = 0usize;
let mut write_secs = 0f64;
let total_waves = partitions.len().div_ceil(partition_wave);
let mut last_wave_log = Instant::now();
for (wave_idx, wave) in partitions.chunks(partition_wave).enumerate() {
let results: Vec<Vec<EncodedTile>> = match store.as_deref_mut() {
Some(s) => encode_wave_from_store(s, level_idx, wave_idx, wave, zoom, options)?,
None => process_wave(&ctx, wave)?,
};
let t_write = Instant::now();
for tiles in &results {
for t in tiles {
tile_feature_count += t.feature_count;
if t.oversized {
oversized += 1;
}
writer.add_tile_precompressed(
zoom,
t.x,
t.y,
t.hash,
&t.data,
t.raw_len,
t.feature_count,
)?;
}
tile_count += tiles.len();
}
write_secs += t_write.elapsed().as_secs_f64();
if last_wave_log.elapsed() >= WAVE_LOG_INTERVAL {
log::info!(
"[export] level {}/{num_levels} z{zoom}: wave {}/{total_waves}, \
{tile_count} tiles, {:.0}s",
level_idx + 1,
wave_idx + 1,
start.elapsed().as_secs_f64(),
);
last_wave_log = Instant::now();
}
}
log::info!(
"[export] level {}/{num_levels} z{zoom} done: {} feats, {tile_count} tiles, \
{} partitions, {:.1}s (total {:.0}s)",
level_idx + 1,
scan.feature_count,
partitions.len(),
t_tiles.elapsed().as_secs_f64(),
start.elapsed().as_secs_f64(),
);
log::debug!(
"[profile] z{zoom} (level {level_idx}, {} feats, {tile_count} tiles, {} partitions): \
clip+encode+gzip={:.2}s write={write_secs:.2}s",
scan.feature_count,
partitions.len(),
t_tiles.elapsed().as_secs_f64() - write_secs,
);
Ok(ZoomReport {
zoom,
level: level_idx,
level_feature_count: scan.feature_count,
tile_count,
tile_feature_count,
oversized_tiles: oversized,
})
}
#[inline]
fn tile_key(x: u32, y: u32) -> u64 {
((x as u64) << 32) | y as u64
}
struct Member {
key: u64,
seq: u64,
geom: Geometry<f64>,
props: Arc<Vec<(String, PropertyValue)>>,
}
#[cfg_attr(test, derive(Debug, PartialEq))]
struct LevelScan {
feature_count: usize,
bounds: Option<TileBounds>,
tile_counts: BTreeMap<u64, usize>,
}
struct Partition {
key_lo: u64,
key_hi: u64,
bbox: TileBounds,
members: usize,
}
struct LevelCtx<'a> {
reader: &'a OverviewReader,
level_idx: usize,
crs: Crs,
zoom: u8,
opts: &'a ExportOptions,
published: &'a PublishedNames,
}
#[derive(Clone, Copy)]
struct FanoutCtx<'a> {
crs: Crs,
plans: &'a [LevelPlan],
opts: &'a ExportOptions,
published: &'a PublishedNames,
}
struct LevelPlan {
zoom: u8,
partitions: Vec<Partition>,
wave: usize,
}
const MEMBER_STORE_RAM_BUDGET: usize = 128 * 1024 * 1024;
const SINGLE_READ_IN_FLIGHT: usize = 4;
fn put_u32(buf: &mut Vec<u8>, v: u32) {
buf.extend_from_slice(&v.to_le_bytes());
}
fn put_u64(buf: &mut Vec<u8>, v: u64) {
buf.extend_from_slice(&v.to_le_bytes());
}
fn put_coord(buf: &mut Vec<u8>, c: geo::Coord<f64>) {
put_u64(buf, c.x.to_bits());
put_u64(buf, c.y.to_bits());
}
fn put_line_string(buf: &mut Vec<u8>, ls: &geo::LineString<f64>) {
put_u32(buf, ls.0.len() as u32);
for c in &ls.0 {
put_coord(buf, *c);
}
}
fn put_polygon(buf: &mut Vec<u8>, p: &geo::Polygon<f64>) {
put_u32(buf, 1 + p.interiors().len() as u32);
put_line_string(buf, p.exterior());
for r in p.interiors() {
put_line_string(buf, r);
}
}
fn encode_geometry(buf: &mut Vec<u8>, g: &Geometry<f64>) {
match g {
Geometry::Point(p) => {
buf.push(0);
put_coord(buf, p.0);
}
Geometry::Line(l) => {
buf.push(1);
put_coord(buf, l.start);
put_coord(buf, l.end);
}
Geometry::LineString(ls) => {
buf.push(2);
put_line_string(buf, ls);
}
Geometry::Polygon(p) => {
buf.push(3);
put_polygon(buf, p);
}
Geometry::MultiPoint(mp) => {
buf.push(4);
put_u32(buf, mp.0.len() as u32);
for p in &mp.0 {
put_coord(buf, p.0);
}
}
Geometry::MultiLineString(mls) => {
buf.push(5);
put_u32(buf, mls.0.len() as u32);
for ls in &mls.0 {
put_line_string(buf, ls);
}
}
Geometry::MultiPolygon(mp) => {
buf.push(6);
put_u32(buf, mp.0.len() as u32);
for p in &mp.0 {
put_polygon(buf, p);
}
}
Geometry::GeometryCollection(gc) => {
buf.push(7);
put_u32(buf, gc.0.len() as u32);
for g in &gc.0 {
encode_geometry(buf, g);
}
}
Geometry::Rect(r) => {
buf.push(8);
put_coord(buf, r.min());
put_coord(buf, r.max());
}
Geometry::Triangle(t) => {
buf.push(9);
put_coord(buf, t.v1());
put_coord(buf, t.v2());
put_coord(buf, t.v3());
}
}
}
fn encode_member(buf: &mut Vec<u8>, m: &Member) {
put_u64(buf, m.key);
put_u64(buf, m.seq);
put_u32(buf, m.props.len() as u32);
for (name, v) in m.props.iter() {
put_u32(buf, name.len() as u32);
buf.extend_from_slice(name.as_bytes());
match v {
PropertyValue::String(s) => {
buf.push(0);
put_u32(buf, s.len() as u32);
buf.extend_from_slice(s.as_bytes());
}
PropertyValue::Float(f) => {
buf.push(1);
put_u32(buf, f.to_bits());
}
PropertyValue::Double(d) => {
buf.push(2);
put_u64(buf, d.to_bits());
}
PropertyValue::Int(i) => {
buf.push(3);
put_u64(buf, *i as u64);
}
PropertyValue::UInt(u) => {
buf.push(4);
put_u64(buf, *u);
}
PropertyValue::Bool(b) => {
buf.push(5);
buf.push(*b as u8);
}
}
}
encode_geometry(buf, &m.geom);
}
fn spill_corrupt() -> ExportError {
ExportError::Io(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"member spill decode: truncated or corrupt record",
))
}
struct SpillCursor<'a> {
buf: &'a [u8],
pos: usize,
}
impl<'a> SpillCursor<'a> {
fn new(buf: &'a [u8]) -> Self {
SpillCursor { buf, pos: 0 }
}
fn is_empty(&self) -> bool {
self.pos == self.buf.len()
}
fn take(&mut self, n: usize) -> Result<&'a [u8], ExportError> {
let end = self.pos.checked_add(n).filter(|&e| e <= self.buf.len());
let Some(end) = end else {
return Err(spill_corrupt());
};
let out = &self.buf[self.pos..end];
self.pos = end;
Ok(out)
}
fn u8(&mut self) -> Result<u8, ExportError> {
Ok(self.take(1)?[0])
}
fn u32(&mut self) -> Result<u32, ExportError> {
Ok(u32::from_le_bytes(self.take(4)?.try_into().unwrap()))
}
fn u64(&mut self) -> Result<u64, ExportError> {
Ok(u64::from_le_bytes(self.take(8)?.try_into().unwrap()))
}
fn coord(&mut self) -> Result<geo::Coord<f64>, ExportError> {
let x = f64::from_bits(self.u64()?);
let y = f64::from_bits(self.u64()?);
Ok(geo::coord! { x: x, y: y })
}
fn line_string(&mut self) -> Result<geo::LineString<f64>, ExportError> {
let n = self.u32()? as usize;
let mut coords = Vec::with_capacity(n);
for _ in 0..n {
coords.push(self.coord()?);
}
Ok(geo::LineString(coords))
}
fn polygon(&mut self) -> Result<geo::Polygon<f64>, ExportError> {
let rings = self.u32()? as usize;
if rings == 0 {
return Err(spill_corrupt());
}
let exterior = self.line_string()?;
let mut interiors = Vec::with_capacity(rings - 1);
for _ in 1..rings {
interiors.push(self.line_string()?);
}
Ok(geo::Polygon::new(exterior, interiors))
}
fn string(&mut self) -> Result<String, ExportError> {
let n = self.u32()? as usize;
String::from_utf8(self.take(n)?.to_vec()).map_err(|_| spill_corrupt())
}
}
fn decode_geometry(cur: &mut SpillCursor<'_>) -> Result<Geometry<f64>, ExportError> {
Ok(match cur.u8()? {
0 => Geometry::Point(geo::Point(cur.coord()?)),
1 => Geometry::Line(geo::Line::new(cur.coord()?, cur.coord()?)),
2 => Geometry::LineString(cur.line_string()?),
3 => Geometry::Polygon(cur.polygon()?),
4 => {
let n = cur.u32()? as usize;
let mut pts = Vec::with_capacity(n);
for _ in 0..n {
pts.push(geo::Point(cur.coord()?));
}
Geometry::MultiPoint(geo::MultiPoint(pts))
}
5 => {
let n = cur.u32()? as usize;
let mut lines = Vec::with_capacity(n);
for _ in 0..n {
lines.push(cur.line_string()?);
}
Geometry::MultiLineString(geo::MultiLineString(lines))
}
6 => {
let n = cur.u32()? as usize;
let mut polys = Vec::with_capacity(n);
for _ in 0..n {
polys.push(cur.polygon()?);
}
Geometry::MultiPolygon(geo::MultiPolygon(polys))
}
7 => {
let n = cur.u32()? as usize;
let mut geoms = Vec::with_capacity(n);
for _ in 0..n {
geoms.push(decode_geometry(cur)?);
}
Geometry::GeometryCollection(geo::GeometryCollection(geoms))
}
8 => Geometry::Rect(geo::Rect::new(cur.coord()?, cur.coord()?)),
9 => Geometry::Triangle(geo::Triangle::new(cur.coord()?, cur.coord()?, cur.coord()?)),
_ => return Err(spill_corrupt()),
})
}
fn decode_member(cur: &mut SpillCursor<'_>) -> Result<Member, ExportError> {
let key = cur.u64()?;
let seq = cur.u64()?;
let n_props = cur.u32()? as usize;
let mut props = Vec::with_capacity(n_props);
for _ in 0..n_props {
let name = cur.string()?;
let value = match cur.u8()? {
0 => PropertyValue::String(cur.string()?),
1 => PropertyValue::Float(f32::from_bits(cur.u32()?)),
2 => PropertyValue::Double(f64::from_bits(cur.u64()?)),
3 => PropertyValue::Int(cur.u64()? as i64),
4 => PropertyValue::UInt(cur.u64()?),
5 => PropertyValue::Bool(cur.u8()? != 0),
_ => return Err(spill_corrupt()),
};
props.push((name, value));
}
let geom = decode_geometry(cur)?;
Ok(Member {
key,
seq,
geom,
props: Arc::new(props),
})
}
fn member_bytes_estimate(m: &Member) -> usize {
48 + m.geom.coords_count() * 16
+ m.props
.iter()
.map(|(n, v)| {
24 + n.len()
+ match v {
PropertyValue::String(s) => s.len(),
_ => 8,
}
})
.sum::<usize>()
}
struct MemberSpill {
writer: BufWriter<File>,
read: File,
_temp: NamedTempFile,
offset: u64,
segments: Vec<Vec<Vec<(u64, u64)>>>,
}
struct MemberStore {
buckets: Vec<Vec<Vec<Member>>>,
bucket_bytes: Vec<Vec<usize>>,
buffered_bytes: usize,
spill: Option<MemberSpill>,
}
impl MemberStore {
fn new(plans: &[LevelPlan], backing: SinkBacking) -> Result<Self, ExportError> {
let shape: Vec<usize> = plans
.iter()
.map(|p| p.partitions.len().div_ceil(p.wave.max(1)))
.collect();
let buckets = shape
.iter()
.map(|&w| (0..w).map(|_| Vec::new()).collect())
.collect();
let bucket_bytes = shape.iter().map(|&w| vec![0usize; w]).collect();
let spill = match backing {
SinkBacking::Ram => None,
SinkBacking::Spill => {
let temp = NamedTempFile::new()?;
let writer = BufWriter::new(temp.reopen()?);
let read = temp.reopen()?;
Some(MemberSpill {
writer,
read,
_temp: temp,
offset: 0,
segments: shape
.iter()
.map(|&w| (0..w).map(|_| Vec::new()).collect())
.collect(),
})
}
};
Ok(MemberStore {
buckets,
bucket_bytes,
buffered_bytes: 0,
spill,
})
}
fn push(&mut self, level: usize, wave: usize, m: Member) -> Result<(), ExportError> {
let est = member_bytes_estimate(&m);
self.buckets[level][wave].push(m);
self.bucket_bytes[level][wave] += est;
self.buffered_bytes += est;
if self.spill.is_some() && self.buffered_bytes > MEMBER_STORE_RAM_BUDGET {
self.flush_largest_until(MEMBER_STORE_RAM_BUDGET / 2)?;
}
Ok(())
}
fn flush_largest_until(&mut self, floor: usize) -> Result<(), ExportError> {
while self.buffered_bytes > floor {
let mut best = (0usize, 0usize, 0usize);
for (li, level) in self.bucket_bytes.iter().enumerate() {
for (wi, &b) in level.iter().enumerate() {
if b > best.2 {
best = (li, wi, b);
}
}
}
if best.2 == 0 {
break;
}
self.flush_bucket(best.0, best.1)?;
}
Ok(())
}
fn flush_bucket(&mut self, level: usize, wave: usize) -> Result<(), ExportError> {
use std::io::Write;
let members = std::mem::take(&mut self.buckets[level][wave]);
let bytes = std::mem::take(&mut self.bucket_bytes[level][wave]);
self.buffered_bytes -= bytes;
if members.is_empty() {
return Ok(());
}
let spill = self
.spill
.as_mut()
.expect("flush_bucket requires spill backing");
let mut buf = Vec::with_capacity(bytes + 16);
put_u64(&mut buf, members.len() as u64);
for m in &members {
encode_member(&mut buf, m);
}
spill.writer.write_all(&buf)?;
spill.segments[level][wave].push((spill.offset, buf.len() as u64));
spill.offset += buf.len() as u64;
Ok(())
}
fn finish_fill(&mut self) -> Result<(), ExportError> {
use std::io::Write;
if self.spill.is_some() {
for li in 0..self.buckets.len() {
for wi in 0..self.buckets[li].len() {
self.flush_bucket(li, wi)?;
}
}
self.spill
.as_mut()
.expect("spill checked above")
.writer
.flush()?;
}
Ok(())
}
fn take_wave(&mut self, level: usize, wave: usize) -> Result<Vec<Member>, ExportError> {
use std::io::{Read, Seek, SeekFrom};
let mut out = Vec::new();
if let Some(spill) = self.spill.as_mut() {
for &(off, len) in &spill.segments[level][wave] {
spill.read.seek(SeekFrom::Start(off))?;
let mut buf = vec![0u8; len as usize];
spill.read.read_exact(&mut buf)?;
let mut cur = SpillCursor::new(&buf);
let n = cur.u64()? as usize;
out.reserve(n);
for _ in 0..n {
out.push(decode_member(&mut cur)?);
}
if !cur.is_empty() {
return Err(spill_corrupt());
}
}
spill.segments[level][wave] = Vec::new();
}
let bytes = std::mem::take(&mut self.bucket_bytes[level][wave]);
self.buffered_bytes -= bytes;
out.append(&mut self.buckets[level][wave]);
Ok(out)
}
}
fn fill_member_store(
reader: &OverviewReader,
ctx: &FanoutCtx<'_>,
backing: SinkBacking,
) -> Result<MemberStore, ExportError> {
let plans = ctx.plans;
let num_levels = plans.len();
let mut store = MemberStore::new(plans, backing)?;
let t_fill = Instant::now();
let mut seq = 0u64;
let store_ref = &mut store;
let seq_ref = &mut seq;
scoped_pipe(
SINGLE_READ_IN_FLIGHT,
|tx: &Sender<(usize, RecordBatch)>| -> Result<(), ExportError> {
for band in 0..num_levels {
let band_reader = reader.read_band_with_batch_size(band, EXPORT_BATCH_SIZE)?;
for batch in band_reader {
if tx.send((band, batch?)).is_err() {
return Ok(()); }
}
}
Ok(())
},
|rx: Receiver<(usize, RecordBatch)>| -> Result<(), ExportError> {
for (band, batch) in rx.iter() {
fanout_batch_members(&batch, band, ctx, seq_ref, store_ref)?;
}
Ok(())
},
)?;
store.finish_fill()?;
log::info!(
"[export] pass2 fill (#235): {seq} row(s) read once and fanned across \
{num_levels} level(s) in {:.2}s",
t_fill.elapsed().as_secs_f64(),
);
Ok(store)
}
fn fanout_batch_members(
batch: &RecordBatch,
band: usize,
ctx: &FanoutCtx<'_>,
seq: &mut u64,
store: &mut MemberStore,
) -> Result<(), ExportError> {
let FanoutCtx {
crs,
plans,
opts,
published,
} = *ctx;
let schema = batch.schema();
let geom_idx = geometry_index(&schema).ok_or(ExportError::NoGeometryColumn)?;
let geom_field = schema.field(geom_idx).clone();
let garr: Arc<dyn GeoArrowArray> =
from_arrow_array(batch.column(geom_idx).as_ref(), &geom_field)
.map_err(|e| crate::Error::GeoParquetRead(format!("geometry decode: {e}")))?;
let mut geoms: Vec<Geometry<f64>> = Vec::with_capacity(batch.num_rows());
extract_geometries_from_array(garr.as_ref(), &mut geoms)?;
if matches!(crs, Crs::Epsg3857) {
geoms = geoms.par_iter().map(reproject_3857_to_4326).collect();
}
let targets: Vec<usize> = (band..plans.len())
.filter(|&k| !plans[k].partitions.is_empty())
.collect();
type RowMembers = Vec<Vec<(u64, Geometry<f64>)>>;
let mut per_level: Vec<RowMembers> = targets
.iter()
.map(|&k| {
geoms
.par_iter()
.map(|g| feature_tile_members(g, plans[k].zoom, opts, 0, u64::MAX))
.collect()
})
.collect();
drop(geoms);
if per_level.iter().all(|rows| rows.iter().all(Vec::is_empty)) {
*seq += batch.num_rows() as u64;
return Ok(());
}
let prop_cols = property_columns(&schema, geom_idx, published);
let mut extracted: Vec<(String, Vec<Option<PropertyValue>>)> =
Vec::with_capacity(prop_cols.len());
for &(idx, ref name) in &prop_cols {
extracted.push((name.clone(), extract_property_column(batch.column(idx))));
}
for row in 0..batch.num_rows() {
if per_level.iter().all(|rows| rows[row].is_empty()) {
*seq += 1;
continue;
}
let mut props = Vec::with_capacity(extracted.len());
for (name, col) in &extracted {
if let Some(v) = &col[row] {
props.push((name.clone(), v.clone()));
}
}
let props = Arc::new(props);
for (ti, &k) in targets.iter().enumerate() {
let items = std::mem::take(&mut per_level[ti][row]);
let plan = &plans[k];
for (key, geom) in items {
let pidx = route_partition(&plan.partitions, key);
let wave = pidx / plan.wave.max(1);
store.push(
k,
wave,
Member {
key,
seq: *seq,
geom,
props: Arc::clone(&props),
},
)?;
}
}
*seq += 1;
}
Ok(())
}
fn encode_wave_from_store(
store: &mut MemberStore,
level_idx: usize,
wave_idx: usize,
wave: &[Partition],
zoom: u8,
opts: &ExportOptions,
) -> Result<Vec<Vec<EncodedTile>>, ExportError> {
let members = store.take_wave(level_idx, wave_idx)?;
let mut buckets: Vec<Vec<Member>> = (0..wave.len()).map(|_| Vec::new()).collect();
for m in members {
buckets[route_partition(wave, m.key)].push(m);
}
buckets
.into_par_iter()
.map(|members| encode_members(members, zoom, opts))
.collect::<Result<_, _>>()
}
fn scan_all_levels(
reader: &OverviewReader,
crs: Crs,
meta: &OverviewsMeta,
opts: &ExportOptions,
) -> Result<Vec<LevelScan>, ExportError> {
let num_levels = reader.num_levels();
let partitioning = matches!(reader.mode(), Mode::Partitioning);
let zooms: Vec<u8> = (0..num_levels).map(|k| zoom_for_level(meta, k)).collect();
let mut scans: Vec<LevelScan> = (0..num_levels)
.map(|_| LevelScan {
feature_count: 0,
bounds: None,
tile_counts: BTreeMap::new(),
})
.collect();
for j in 0..num_levels {
let last = if partitioning { num_levels - 1 } else { j };
let band_reader = reader.read_band_with_batch_size(j, EXPORT_BATCH_SIZE)?;
for batch in band_reader {
let bboxes = decode_batch_bboxes(&batch?, crs)?;
for k in j..=last {
let scan = &mut scans[k];
scan.feature_count += bboxes.len();
let (batch_bounds, batch_counts) = size_bboxes(&bboxes, zooms[k], opts);
if let Some(b) = batch_bounds {
match &mut scan.bounds {
Some(acc) => acc.expand(&b),
None => scan.bounds = Some(b),
}
}
for (key, v) in batch_counts {
*scan.tile_counts.entry(key).or_insert(0) += v;
}
}
}
}
Ok(scans)
}
fn decode_batch_bboxes(
batch: &RecordBatch,
crs: Crs,
) -> Result<Vec<Option<TileBounds>>, ExportError> {
let schema = batch.schema();
let geom_idx = geometry_index(&schema).ok_or(ExportError::NoGeometryColumn)?;
let geom_field = schema.field(geom_idx).clone();
let garr: Arc<dyn GeoArrowArray> =
from_arrow_array(batch.column(geom_idx).as_ref(), &geom_field)
.map_err(|e| crate::Error::GeoParquetRead(format!("geometry decode: {e}")))?;
let mut geoms: Vec<Geometry<f64>> = Vec::with_capacity(batch.num_rows());
extract_geometries_from_array(garr.as_ref(), &mut geoms)?;
let bboxes = geoms
.par_iter()
.map(|geom| {
geom.bounding_rect().map(|rect| {
let bbox = TileBounds::new(rect.min().x, rect.min().y, rect.max().x, rect.max().y);
if matches!(crs, Crs::Epsg3857) {
let (lng_min, lat_min) = webmerc_to_lnglat(bbox.lng_min, bbox.lat_min);
let (lng_max, lat_max) = webmerc_to_lnglat(bbox.lng_max, bbox.lat_max);
TileBounds::new(lng_min, lat_min, lng_max, lat_max)
} else {
bbox
}
})
})
.collect();
Ok(bboxes)
}
fn size_bboxes(
bboxes: &[Option<TileBounds>],
zoom: u8,
opts: &ExportOptions,
) -> (Option<TileBounds>, HashMap<u64, usize>) {
bboxes
.par_iter()
.fold(
|| (None::<TileBounds>, HashMap::<u64, usize>::new()),
|(mut bounds, mut counts), bbox| {
if let Some(bbox) = bbox {
match &mut bounds {
Some(acc) => acc.expand(bbox),
None => bounds = Some(*bbox),
}
let ranges = member_ranges(bbox, zoom, opts);
for tc in tiles_in_ranges(&ranges, zoom) {
*counts.entry(tile_key(tc.x, tc.y)).or_insert(0) += 1;
}
}
(bounds, counts)
},
)
.reduce(
|| (None::<TileBounds>, HashMap::<u64, usize>::new()),
|(b_a, c_a), (b_b, c_b)| {
let bounds = match (b_a, b_b) {
(Some(mut a), Some(b)) => {
a.expand(&b);
Some(a)
}
(Some(a), None) | (None, Some(a)) => Some(a),
(None, None) => None,
};
let (mut dst, src) = if c_a.len() >= c_b.len() {
(c_a, c_b)
} else {
(c_b, c_a)
};
for (k, v) in src {
*dst.entry(k).or_insert(0) += v;
}
(bounds, dst)
},
)
}
#[cfg(test)]
fn scan_level(
reader: &OverviewReader,
level_idx: usize,
crs: Crs,
zoom: u8,
opts: &ExportOptions,
) -> Result<LevelScan, ExportError> {
let batch_reader = reader.read_level_with_batch_size(level_idx, None, EXPORT_BATCH_SIZE)?;
let mut scan = LevelScan {
feature_count: 0,
bounds: None,
tile_counts: BTreeMap::new(),
};
for batch in batch_reader {
let bboxes = decode_batch_bboxes(&batch?, crs)?;
scan.feature_count += bboxes.len();
let (batch_bounds, batch_counts) = size_bboxes(&bboxes, zoom, opts);
if let Some(b) = batch_bounds {
match &mut scan.bounds {
Some(acc) => acc.expand(&b),
None => scan.bounds = Some(b),
}
}
for (k, v) in batch_counts {
*scan.tile_counts.entry(k).or_insert(0) += v;
}
}
Ok(scan)
}
fn plan_partitions(
tile_counts: &BTreeMap<u64, usize>,
zoom: u8,
partition_target: usize,
) -> Vec<Partition> {
let target = partition_target.max(1);
let mut out: Vec<Partition> = Vec::new();
let mut cur: Option<Partition> = None;
for (&key, &count) in tile_counts {
let tb = TileCoord::new((key >> 32) as u32, key as u32, zoom).bounds();
match cur.as_mut() {
Some(p) => {
p.key_hi = key;
p.bbox.expand(&tb);
p.members += count;
}
None => {
cur = Some(Partition {
key_lo: key,
key_hi: key,
bbox: tb,
members: count,
});
}
}
if cur.as_ref().is_some_and(|p| p.members >= target) {
out.push(cur.take().unwrap());
}
}
out.extend(cur);
let margin = 360.0 / 2f64.powi(zoom as i32);
for p in &mut out {
p.bbox = TileBounds::new(
p.bbox.lng_min - margin,
p.bbox.lat_min - margin,
p.bbox.lng_max + margin,
p.bbox.lat_max + margin,
);
}
out
}
fn process_wave(
ctx: &LevelCtx<'_>,
wave: &[Partition],
) -> Result<Vec<Vec<EncodedTile>>, ExportError> {
debug_assert!(!wave.is_empty(), "wave must be non-empty");
let bbox = match ctx.crs {
Crs::Epsg4326 => {
let mut b = wave[0].bbox;
for p in &wave[1..] {
b.expand(&p.bbox);
}
Some([b.lng_min, b.lat_min, b.lng_max, b.lat_max])
}
Crs::Epsg3857 => None,
};
let batch_reader =
ctx.reader
.read_level_with_batch_size(ctx.level_idx, bbox, EXPORT_BATCH_SIZE)?;
let t_collect = Instant::now();
let mut buckets: Vec<Vec<Member>> = (0..wave.len()).map(|_| Vec::new()).collect();
let mut seq = 0u64;
for batch in batch_reader {
let batch = batch?;
collect_wave_members(ctx, wave, &batch, &mut seq, &mut buckets)?;
}
let collect_secs = t_collect.elapsed().as_secs_f64();
let n_members: usize = buckets.iter().map(Vec::len).sum();
let t_encode = Instant::now();
let tiles: Vec<Vec<EncodedTile>> = buckets
.into_par_iter()
.map(|members| encode_members(members, ctx.zoom, ctx.opts))
.collect::<Result<_, _>>()?;
let n_tiles: usize = tiles.iter().map(Vec::len).sum();
log::debug!(
"[profile] z{} wave [{:x}..{:x}] ({} partitions): rows_read={seq} \
members={n_members} tiles={n_tiles} collect={collect_secs:.2}s \
sort+encode={:.2}s",
ctx.zoom,
wave[0].key_lo,
wave[wave.len() - 1].key_hi,
wave.len(),
t_encode.elapsed().as_secs_f64(),
);
Ok(tiles)
}
#[inline]
fn route_partition(wave: &[Partition], key: u64) -> usize {
let idx = wave.partition_point(|p| p.key_lo <= key).saturating_sub(1);
debug_assert!(
key >= wave[idx].key_lo && key <= wave[idx].key_hi,
"member key {key:x} falls outside its wave partition [{:x}..{:x}]",
wave[idx].key_lo,
wave[idx].key_hi,
);
idx
}
fn collect_wave_members(
ctx: &LevelCtx<'_>,
wave: &[Partition],
batch: &RecordBatch,
seq: &mut u64,
buckets: &mut [Vec<Member>],
) -> Result<(), ExportError> {
debug_assert_eq!(buckets.len(), wave.len());
let key_lo = wave[0].key_lo;
let key_hi = wave[wave.len() - 1].key_hi;
let schema = batch.schema();
let geom_idx = geometry_index(&schema).ok_or(ExportError::NoGeometryColumn)?;
let geom_field = schema.field(geom_idx).clone();
let garr: Arc<dyn GeoArrowArray> =
from_arrow_array(batch.column(geom_idx).as_ref(), &geom_field)
.map_err(|e| crate::Error::GeoParquetRead(format!("geometry decode: {e}")))?;
let mut geoms: Vec<Geometry<f64>> = Vec::with_capacity(batch.num_rows());
extract_geometries_from_array(garr.as_ref(), &mut geoms)?;
if matches!(ctx.crs, Crs::Epsg3857) {
geoms = geoms.par_iter().map(reproject_3857_to_4326).collect();
}
let row_members: Vec<Vec<(u64, Geometry<f64>)>> = geoms
.par_iter()
.map(|g| feature_tile_members(g, ctx.zoom, ctx.opts, key_lo, key_hi))
.collect();
drop(geoms);
if row_members.iter().all(|v| v.is_empty()) {
*seq += row_members.len() as u64;
return Ok(());
}
let prop_cols = property_columns(&schema, geom_idx, ctx.published);
let mut extracted: Vec<(String, Vec<Option<PropertyValue>>)> =
Vec::with_capacity(prop_cols.len());
for &(idx, ref name) in &prop_cols {
extracted.push((name.clone(), extract_property_column(batch.column(idx))));
}
for (row, items) in row_members.into_iter().enumerate() {
if items.is_empty() {
*seq += 1;
continue;
}
let mut props = Vec::with_capacity(extracted.len());
for (name, col) in &extracted {
if let Some(v) = &col[row] {
props.push((name.clone(), v.clone()));
}
}
let props = Arc::new(props);
for (key, geom) in items {
buckets[route_partition(wave, key)].push(Member {
key,
seq: *seq,
geom,
props: Arc::clone(&props),
});
}
*seq += 1;
}
Ok(())
}
fn feature_tile_members(
geom: &Geometry<f64>,
zoom: u8,
opts: &ExportOptions,
key_lo: u64,
key_hi: u64,
) -> Vec<(u64, Geometry<f64>)> {
let Some(rect) = geom.bounding_rect() else {
return Vec::new();
};
let bbox = TileBounds::new(rect.min().x, rect.min().y, rect.max().x, rect.max().y);
let ranges = member_ranges(&bbox, zoom, opts);
let mut out = Vec::new();
let assume_simple = geometry_is_simple(geom);
let direct_cost = tile_span(&ranges).saturating_mul(geom.coords_count() as u64);
if direct_cost <= DIRECT_CLIP_BUDGET {
feature_tile_members_direct(
geom,
&bbox,
zoom,
opts,
key_lo,
key_hi,
&ranges,
assume_simple,
&mut out,
);
} else {
let root = covering_tile(&ranges, zoom);
split_feature_into_tiles(
root,
geom,
&bbox,
zoom,
opts,
key_lo,
key_hi,
&ranges,
assume_simple,
&mut out,
);
}
out
}
#[inline]
fn buffer_deg_at_zoom(zoom: u8, opts: &ExportOptions) -> f64 {
let tile_width = 360.0 / f64::from(1u32 << zoom);
tile_width * buffer_fraction(opts)
}
#[inline]
fn expand_bbox(bbox: &TileBounds, buffer: f64) -> TileBounds {
let was_wrapped = bbox.lng_min > bbox.lng_max;
let mut lng_min = (bbox.lng_min - buffer).clamp(-180.0, 180.0);
let lng_max = (bbox.lng_max + buffer).clamp(-180.0, 180.0);
if !was_wrapped && lng_min > lng_max {
lng_min = lng_max;
}
TileBounds::new(
lng_min,
(bbox.lat_min - buffer).clamp(-90.0, 90.0),
lng_max,
(bbox.lat_max + buffer).clamp(-90.0, 90.0),
)
}
#[inline]
fn member_ranges(bbox: &TileBounds, zoom: u8, opts: &ExportOptions) -> BboxTileRanges {
tile_ranges_for_bbox(&expand_bbox(bbox, buffer_deg_at_zoom(zoom, opts)), zoom)
}
const DIRECT_CLIP_BUDGET: u64 = 100_000;
const CASCADE_PARALLEL_DEPTH: u32 = 3;
#[inline]
fn tile_span(ranges: &BboxTileRanges) -> u64 {
let x1 = (ranges.x.1 - ranges.x.0 + 1) as u64;
let x2 = ranges.x2.map_or(0, |(a, b)| (b - a + 1) as u64);
let y = (ranges.y.1 - ranges.y.0 + 1) as u64;
(x1 + x2) * y
}
#[allow(clippy::too_many_arguments)]
fn feature_tile_members_direct(
geom: &Geometry<f64>,
bbox: &TileBounds,
zoom: u8,
opts: &ExportOptions,
key_lo: u64,
key_hi: u64,
ranges: &BboxTileRanges,
assume_simple: bool,
out: &mut Vec<(u64, Geometry<f64>)>,
) {
for tc in tiles_in_ranges(ranges, zoom) {
let key = tile_key(tc.x, tc.y);
if key < key_lo || key > key_hi {
continue;
}
let tb = tc.bounds();
let buffer_deg = tb.width() * buffer_fraction(opts);
if bbox_within_buffered(bbox, &tb, buffer_deg) {
out.push((key, geom.clone()));
} else if let Some(clipped) = clip_geometry_simple(
geom,
&tb,
buffer_deg,
assume_simple,
opts.simple_clip_fastpath,
) {
out.push((key, clipped));
}
}
}
fn tiles_in_ranges(ranges: &BboxTileRanges, zoom: u8) -> impl Iterator<Item = TileCoord> + '_ {
let bands = std::iter::once(ranges.x).chain(ranges.x2);
bands.flat_map(move |(x0, x1)| {
(x0..=x1)
.flat_map(move |x| (ranges.y.0..=ranges.y.1).map(move |y| TileCoord::new(x, y, zoom)))
})
}
#[inline]
fn covering_tile(ranges: &BboxTileRanges, zoom: u8) -> TileCoord {
if ranges.x2.is_some() {
return TileCoord::new(0, 0, 0);
}
let (x_lo, x_hi) = ranges.x;
let (y_lo, y_hi) = ranges.y;
let mut s = 0u8;
while s < zoom && !((x_lo >> s) == (x_hi >> s) && (y_lo >> s) == (y_hi >> s)) {
s += 1;
}
TileCoord::new(x_lo >> s, y_lo >> s, zoom - s)
}
#[inline]
fn node_overlaps_ranges(node: TileCoord, zoom: u8, ranges: &BboxTileRanges) -> bool {
let shift = (zoom - node.z) as u32;
let x_lo = (node.x as u64) << shift;
let x_hi = (((node.x as u64) + 1) << shift) - 1;
let y_lo = (node.y as u64) << shift;
let y_hi = (((node.y as u64) + 1) << shift) - 1;
let overlaps = |a: (u32, u32), b: (u64, u64)| a.0 as u64 <= b.1 && b.0 <= a.1 as u64;
if !overlaps(ranges.y, (y_lo, y_hi)) {
return false;
}
overlaps(ranges.x, (x_lo, x_hi)) || ranges.x2.is_some_and(|x2| overlaps(x2, (x_lo, x_hi)))
}
#[inline]
fn node_key_overlaps(node: TileCoord, zoom: u8, key_lo: u64, key_hi: u64) -> bool {
let shift = (zoom - node.z) as u32;
let x_lo = (node.x as u64) << shift;
let x_hi = (((node.x as u64) + 1) << shift) - 1;
let y_lo = (node.y as u64) << shift;
let y_hi = (((node.y as u64) + 1) << shift) - 1;
let min_key = (x_lo << 32) | y_lo;
let max_key = (x_hi << 32) | y_hi;
max_key >= key_lo && min_key <= key_hi
}
#[allow(clippy::too_many_arguments)]
fn split_feature_into_tiles(
node: TileCoord,
cur: &Geometry<f64>,
cur_bbox: &TileBounds,
zoom: u8,
opts: &ExportOptions,
key_lo: u64,
key_hi: u64,
ranges: &BboxTileRanges,
assume_simple: bool,
out: &mut Vec<(u64, Geometry<f64>)>,
) {
if !node_overlaps_ranges(node, zoom, ranges) || !node_key_overlaps(node, zoom, key_lo, key_hi) {
return;
}
let tb = node.bounds();
let buffer_deg = tb.width() * buffer_fraction(opts);
if node.z == zoom {
let key = tile_key(node.x, node.y);
if key < key_lo || key > key_hi {
return;
}
if bbox_within_buffered(cur_bbox, &tb, buffer_deg) {
out.push((key, cur.clone()));
} else if let Some(clipped) = clip_geometry_simple(
cur,
&tb,
buffer_deg,
assume_simple,
opts.simple_clip_fastpath,
) {
out.push((key, clipped));
}
return;
}
let children = node.children().expect("non-leaf node has children");
if bbox_within_buffered(cur_bbox, &tb, buffer_deg) {
fanout_children(
children,
node.z,
cur,
cur_bbox,
zoom,
opts,
key_lo,
key_hi,
ranges,
assume_simple,
out,
);
} else if let Some(clipped) = clip_geometry_simple(
cur,
&tb,
buffer_deg,
assume_simple,
opts.simple_clip_fastpath,
) {
let cbbox = match clipped.bounding_rect() {
Some(r) => TileBounds::new(r.min().x, r.min().y, r.max().x, r.max().y),
None => return,
};
fanout_children(
children,
node.z,
&clipped,
&cbbox,
zoom,
opts,
key_lo,
key_hi,
ranges,
assume_simple,
out,
);
}
}
#[allow(clippy::too_many_arguments)]
fn fanout_children(
children: [TileCoord; 4],
node_z: u8,
child_geom: &Geometry<f64>,
child_bbox: &TileBounds,
zoom: u8,
opts: &ExportOptions,
key_lo: u64,
key_hi: u64,
ranges: &BboxTileRanges,
assume_simple: bool,
out: &mut Vec<(u64, Geometry<f64>)>,
) {
if (zoom - node_z) as u32 > CASCADE_PARALLEL_DEPTH {
let mut left = Vec::new();
let mut right = Vec::new();
rayon::join(
|| {
for &child in &children[..2] {
split_feature_into_tiles(
child,
child_geom,
child_bbox,
zoom,
opts,
key_lo,
key_hi,
ranges,
assume_simple,
&mut left,
);
}
},
|| {
for &child in &children[2..] {
split_feature_into_tiles(
child,
child_geom,
child_bbox,
zoom,
opts,
key_lo,
key_hi,
ranges,
assume_simple,
&mut right,
);
}
},
);
out.append(&mut left);
out.append(&mut right);
} else {
for child in children {
split_feature_into_tiles(
child,
child_geom,
child_bbox,
zoom,
opts,
key_lo,
key_hi,
ranges,
assume_simple,
out,
);
}
}
}
#[cfg(test)]
fn members_direct_vec(
geom: &Geometry<f64>,
zoom: u8,
opts: &ExportOptions,
key_lo: u64,
key_hi: u64,
) -> Vec<(u64, Geometry<f64>)> {
let Some(rect) = geom.bounding_rect() else {
return Vec::new();
};
let bbox = TileBounds::new(rect.min().x, rect.min().y, rect.max().x, rect.max().y);
let ranges = member_ranges(&bbox, zoom, opts);
let assume_simple = geometry_is_simple(geom);
let mut out = Vec::new();
feature_tile_members_direct(
geom,
&bbox,
zoom,
opts,
key_lo,
key_hi,
&ranges,
assume_simple,
&mut out,
);
out
}
#[cfg(test)]
fn members_recursive_vec(
geom: &Geometry<f64>,
zoom: u8,
opts: &ExportOptions,
key_lo: u64,
key_hi: u64,
) -> Vec<(u64, Geometry<f64>)> {
let Some(rect) = geom.bounding_rect() else {
return Vec::new();
};
let bbox = TileBounds::new(rect.min().x, rect.min().y, rect.max().x, rect.max().y);
let ranges = member_ranges(&bbox, zoom, opts);
let root = covering_tile(&ranges, zoom);
let assume_simple = geometry_is_simple(geom);
let mut out = Vec::new();
split_feature_into_tiles(
root,
geom,
&bbox,
zoom,
opts,
key_lo,
key_hi,
&ranges,
assume_simple,
&mut out,
);
out
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub enum FeatureOrder {
#[default]
Input,
Column {
name: String,
descending: bool,
},
}
impl std::str::FromStr for FeatureOrder {
type Err = String;
fn from_str(s: &str) -> Result<Self, Self::Err> {
if s == "input" {
return Ok(Self::Input);
}
let (name, descending) = match s.rsplit_once(':') {
Some((n, "asc")) => (n, false),
Some((n, "desc")) => (n, true),
Some((_, other)) => {
return Err(format!(
"unknown sort direction {other:?}: expected `asc` or `desc` \
(as in `--feature-order level:desc`)"
))
}
None => (s, false),
};
if name.is_empty() {
return Err("expected `input` or a column name, e.g. `level:desc`".to_string());
}
Ok(Self::Column {
name: name.to_string(),
descending,
})
}
}
enum OrderKey<'a> {
Missing,
Number(f64),
Text(&'a str),
}
impl<'a> OrderKey<'a> {
fn of(m: &'a Member, name: &str) -> Self {
let Some((_, v)) = m.props.iter().find(|(k, _)| k == name) else {
return Self::Missing;
};
match v {
PropertyValue::String(s) => Self::Text(s),
PropertyValue::Float(f) => Self::number(*f as f64),
PropertyValue::Double(d) => Self::number(*d),
PropertyValue::Int(i) => Self::number(*i as f64),
PropertyValue::UInt(u) => Self::number(*u as f64),
PropertyValue::Bool(b) => Self::number(*b as u8 as f64),
}
}
fn number(v: f64) -> Self {
if v.is_nan() {
Self::Missing
} else {
Self::Number(v)
}
}
fn cmp(&self, other: &Self) -> std::cmp::Ordering {
use std::cmp::Ordering;
match (self, other) {
(Self::Missing, Self::Missing) => Ordering::Equal,
(Self::Missing, _) => Ordering::Less,
(_, Self::Missing) => Ordering::Greater,
(Self::Number(a), Self::Number(b)) => a.total_cmp(b),
(Self::Number(_), Self::Text(_)) => Ordering::Less,
(Self::Text(_), Self::Number(_)) => Ordering::Greater,
(Self::Text(a), Self::Text(b)) => a.cmp(b),
}
}
}
fn warn_if_order_column_is_absent(
schema: &Schema,
published: &PublishedNames,
order: &FeatureOrder,
) {
let FeatureOrder::Column { name, .. } = order else {
return;
};
if !order_column_is_absent(schema, published, order) {
return;
}
log::warn!(
"--feature-order names {name:?}, which this layer does not publish, so \
the order is unchanged (input row order). Available: {}",
published_property_names(schema, published).join(", ")
);
}
fn order_column_is_absent(
schema: &Schema,
published: &PublishedNames,
order: &FeatureOrder,
) -> bool {
let FeatureOrder::Column { name, .. } = order else {
return false;
};
!published_property_names(schema, published)
.iter()
.any(|n| n == name)
}
fn published_property_names(schema: &Schema, published: &PublishedNames) -> Vec<String> {
let Some(geom_idx) = geometry_index(schema) else {
return Vec::new(); };
property_columns(schema, geom_idx, published)
.into_iter()
.map(|(_, n)| n)
.collect()
}
fn compare_by_feature_order(
a: &Member,
b: &Member,
name: &str,
descending: bool,
) -> std::cmp::Ordering {
let ord = OrderKey::of(a, name).cmp(&OrderKey::of(b, name));
let ord = if descending { ord.reverse() } else { ord };
ord.then_with(|| a.seq.cmp(&b.seq))
}
fn sort_members_for_encode(members: &mut [Member], order: &FeatureOrder) {
members.par_sort_unstable_by_key(|m| (m.key, m.seq));
let FeatureOrder::Column { name, descending } = order else {
return;
};
for tile in members.chunk_by_mut(|a, b| a.key == b.key) {
tile.sort_by(|a, b| compare_by_feature_order(a, b, name, *descending));
}
}
fn encode_members(
mut members: Vec<Member>,
zoom: u8,
opts: &ExportOptions,
) -> Result<Vec<EncodedTile>, ExportError> {
sort_members_for_encode(&mut members, &opts.feature_order);
let groups: Vec<&[Member]> = members.chunk_by(|a, b| a.key == b.key).collect();
groups
.into_par_iter()
.filter_map(|g| {
let (x, y) = ((g[0].key >> 32) as u32, g[0].key as u32);
let tb = TileCoord::new(x, y, zoom).bounds();
let (data, count, oversized) = encode_tile(g, &tb, opts);
if count == 0 {
return None;
}
let hash = TileHasher::hash(&data);
let raw_len = data.len();
let compressed = match compression::compress(&data, Compression::Gzip) {
Ok(c) => c,
Err(e) => return Some(Err(ExportError::from(e))),
};
Some(Ok(EncodedTile {
x,
y,
data: compressed,
hash,
raw_len,
feature_count: count,
oversized,
}))
})
.collect()
}
#[inline]
fn bbox_within_buffered(bbox: &TileBounds, tb: &TileBounds, buffer: f64) -> bool {
bbox.lng_min >= tb.lng_min - buffer
&& bbox.lat_min >= tb.lat_min - buffer
&& bbox.lng_max <= tb.lng_max + buffer
&& bbox.lat_max <= tb.lat_max + buffer
}
fn encode_tile(
members: &[Member],
tb: &TileBounds,
opts: &ExportOptions,
) -> (Vec<u8>, usize, bool) {
let data = build_mvt(members.iter(), tb, opts);
match opts.tile_size_limit {
Some(limit) if limit > 0 && data.len() > limit && members.len() > 1 => {
let keep_frac = limit as f64 / data.len() as f64;
let keep = ((members.len() as f64 * keep_frac).floor() as usize).max(1);
let kept = shed_to_fit(members, keep, opts);
let keep = kept.len();
let data = build_mvt(kept, tb, opts);
log::warn!(
"oversized tile ({} bytes > {limit} limit): dropped {} of {} features (one pass)",
data.len(),
members.len() - keep,
members.len()
);
(data, keep, true)
}
_ => {
let count = members.len();
(data, count, false)
}
}
}
fn shed_to_fit<'a>(members: &'a [Member], keep: usize, opts: &ExportOptions) -> Vec<&'a Member> {
let mut kept = select_kept_members(members, keep);
if let FeatureOrder::Column { name, descending } = &opts.feature_order {
kept.sort_by(|a, b| compare_by_feature_order(a, b, name, *descending));
}
kept
}
fn select_kept_members(members: &[Member], keep: usize) -> Vec<&Member> {
let n = members.len();
let keep = keep.clamp(1, n);
let point_like = members.iter().all(|m| m.geom.coords_count() <= 1);
if point_like {
(0..keep).map(|i| &members[i * n / keep]).collect()
} else {
let mut ranked: Vec<&Member> = members.iter().collect();
ranked.sort_by_key(|m| std::cmp::Reverse(m.geom.coords_count()));
ranked.truncate(keep);
ranked
}
}
fn build_mvt<'a>(
members: impl IntoIterator<Item = &'a Member>,
tb: &TileBounds,
opts: &ExportOptions,
) -> Vec<u8> {
let mut layer = LayerBuilder::new(opts.layer_name.clone()).with_extent(opts.extent);
for (i, m) in members.into_iter().enumerate() {
layer.add_feature(Some(i as u64), &m.geom, &m.props, tb);
}
let mut tb_builder = TileBuilder::new();
tb_builder.add_layer(layer.build());
tb_builder.build().encode_to_vec()
}
#[cfg(test)]
fn encode_level_tiles(features: &[Feature], zoom: u8, opts: &ExportOptions) -> Vec<EncodedTile> {
let mut members = Vec::new();
for (fi, f) in features.iter().enumerate() {
let props = Arc::new(f.props.clone());
for (key, geom) in feature_tile_members(&f.geom, zoom, opts, 0, u64::MAX) {
members.push(Member {
key,
seq: fi as u64,
geom,
props: Arc::clone(&props),
});
}
}
let mut tiles = encode_members(members, zoom, opts).expect("in-memory gzip is infallible");
for t in &mut tiles {
t.data = crate::compression::decompress(&t.data, Compression::Gzip)
.expect("gzip roundtrip of just-compressed tile");
}
tiles
}
#[cfg(test)]
fn read_level_features(
reader: &OverviewReader,
level_idx: usize,
crs: Crs,
) -> Result<Vec<Feature>, ExportError> {
let batch_reader = reader.read_level(level_idx, None)?;
let mut out = Vec::new();
for batch in batch_reader {
let batch = batch?;
decode_batch(&batch, crs, &mut out)?;
}
Ok(out)
}
#[cfg(test)]
fn decode_batch(batch: &RecordBatch, crs: Crs, out: &mut Vec<Feature>) -> Result<(), ExportError> {
let schema = batch.schema();
let geom_idx = geometry_index(&schema).ok_or(ExportError::NoGeometryColumn)?;
let geom_field = schema.field(geom_idx).clone();
let garr: Arc<dyn GeoArrowArray> =
from_arrow_array(batch.column(geom_idx).as_ref(), &geom_field)
.map_err(|e| crate::Error::GeoParquetRead(format!("geometry decode: {e}")))?;
let mut geoms: Vec<Geometry<f64>> = Vec::with_capacity(batch.num_rows());
extract_geometries_from_array(garr.as_ref(), &mut geoms)?;
let prop_cols = property_columns(&schema, geom_idx, &PublishedNames::identity());
let mut extracted: Vec<(String, Vec<Option<PropertyValue>>)> =
Vec::with_capacity(prop_cols.len());
for &(idx, ref name) in &prop_cols {
extracted.push((name.clone(), extract_property_column(batch.column(idx))));
}
for row in 0..batch.num_rows() {
let mut geom = geoms[row].clone();
if matches!(crs, Crs::Epsg3857) {
geom = reproject_3857_to_4326(&geom);
}
let mut props = Vec::with_capacity(extracted.len());
for (name, col) in &extracted {
if let Some(v) = &col[row] {
props.push((name.clone(), v.clone()));
}
}
out.push(Feature { geom, props });
}
Ok(())
}
fn geometry_index(schema: &Schema) -> Option<usize> {
schema
.fields()
.iter()
.position(|f| f.name() == "geometry")
.or_else(|| {
schema
.fields()
.iter()
.position(|f| f.name().contains("geom"))
})
}
#[derive(Debug, Clone, Default)]
struct PublishedNames {
restored: HashMap<String, String>,
suppressed: HashSet<String>,
}
impl PublishedNames {
fn identity() -> Self {
Self::default()
}
fn from_renames(
renames: &BTreeMap<String, String>,
schema: &Schema,
suppressed: HashSet<String>,
) -> Self {
let occupied: HashSet<String> = schema
.fields()
.iter()
.filter(|f| {
!f.name().eq_ignore_ascii_case(LEVEL_COLUMN) && !suppressed.contains(f.name())
})
.map(|f| f.name().to_ascii_lowercase())
.collect();
let mut occupied = occupied;
let mut restored: HashMap<String, String> = HashMap::new();
for (renamed, source) in renames {
if schema.column_with_name(renamed).is_none() {
continue;
}
if !occupied.insert(source.to_ascii_lowercase()) {
continue;
}
restored.insert(renamed.clone(), source.clone());
}
for (renamed, source) in &restored {
log::info!(
"[export] publishing overview column {renamed:?} under its \
source name {source:?} (it was renamed on convert to clear \
the reserved overview column, #288/#359)"
);
}
Self {
restored,
suppressed,
}
}
fn from_meta(meta: &OverviewsMeta, schema: &Schema, suppressed: HashSet<String>) -> Self {
match meta
.generalization
.as_ref()
.and_then(|g| g.renamed_columns.as_ref())
{
Some(renames) => Self::from_renames(renames, schema, suppressed),
None => Self {
suppressed,
..Self::identity()
},
}
}
fn from_reader(reader: &OverviewReader) -> Self {
let meta = reader.meta();
let mut suppressed = HashSet::new();
let counter = meta
.generalization
.as_ref()
.and_then(|g| g.coalescing.as_ref())
.map(|c| c.coalesced_count_column.as_str())
.filter(|&c| c == COALESCED_COUNT_COLUMN);
if let Some(column) = counter {
if reader.int_column_max(column) == Some(1) {
log::info!(
"[export] not exporting {column:?}: it is 1 on every row \
(no lines were merged), so it says nothing about the data; \
the overview file keeps it (#379)"
);
suppressed.insert(column.to_string());
}
}
Self::from_meta(meta, reader.schema(), suppressed)
}
fn with_selection(
mut self,
schema: &Schema,
geom_idx: usize,
selection: &PropertySelection,
feature_order: &FeatureOrder,
) -> Result<Self, ExportError> {
if selection.is_identity() {
return Ok(self);
}
if let Some(include) = &selection.include {
let exported: Vec<String> = property_columns(schema, geom_idx, &self)
.into_iter()
.map(|(_, n)| n)
.collect();
for name in include {
if exported.contains(name) {
continue;
}
let withheld = schema
.fields()
.iter()
.map(|f| f.name())
.find(|n| self.is_suppressed(n) && self.publish(n) == name)
.cloned();
if let Some(schema_name) = withheld {
log::info!(
"[export] exporting {schema_name:?} after all: it was withheld \
(#379) but --include-property names it"
);
self.suppressed.remove(&schema_name);
}
}
}
let exportable: Vec<(usize, String)> = property_columns(schema, geom_idx, &self);
let published: Vec<&str> = exportable.iter().map(|(_, n)| n.as_str()).collect();
if let Some(include) = &selection.include {
for name in include {
if !published.contains(&name.as_str()) {
return Err(ExportError::UnknownProperty {
name: name.clone(),
available: published
.iter()
.map(|n| format!("{n:?}"))
.collect::<Vec<_>>()
.join(", "),
});
}
}
}
for name in &selection.exclude {
if !published.contains(&name.as_str()) {
log::warn!(
"[export] excluded property {name:?} is not exported by this file anyway"
);
}
}
if let FeatureOrder::Column { name, .. } = feature_order {
if published.contains(&name.as_str()) && !selection.keeps(name) {
return Err(ExportError::PropertyRequiredByKnob {
name: name.clone(),
knob: "--feature-order".to_string(),
});
}
}
let mut dropped = 0usize;
for (idx, name) in &exportable {
if !selection.keeps(name) {
self.suppressed.insert(schema.field(*idx).name().clone());
dropped += 1;
}
}
let kept: Vec<&str> = published
.iter()
.copied()
.filter(|n| selection.keeps(n))
.collect();
log::info!(
"[export] property selection: keeping {} of {} ({}), dropping {dropped}",
kept.len(),
published.len(),
if kept.is_empty() {
"none — geometry only".to_string()
} else {
kept.iter()
.map(|n| format!("{n:?}"))
.collect::<Vec<_>>()
.join(", ")
},
);
Ok(self)
}
fn is_suppressed(&self, schema_name: &str) -> bool {
self.suppressed.contains(schema_name)
}
fn publish<'a>(&'a self, schema_name: &'a str) -> &'a str {
self.restored
.get(schema_name)
.map_or(schema_name, String::as_str)
}
}
fn property_columns(
schema: &Schema,
geom_idx: usize,
published: &PublishedNames,
) -> Vec<(usize, String)> {
schema
.fields()
.iter()
.enumerate()
.filter(|&(i, f)| {
i != geom_idx
&& !f.name().eq_ignore_ascii_case(LEVEL_COLUMN)
&& !published.is_suppressed(f.name())
&& is_supported_scalar(f.data_type())
})
.map(|(i, f)| (i, published.publish(f.name()).to_string()))
.collect()
}
fn is_supported_scalar(dt: &DataType) -> bool {
matches!(
dt,
DataType::Utf8
| DataType::LargeUtf8
| DataType::Boolean
| DataType::Int8
| DataType::Int16
| DataType::Int32
| DataType::Int64
| DataType::UInt8
| DataType::UInt16
| DataType::UInt32
| DataType::UInt64
| DataType::Float32
| DataType::Float64
) || is_temporal_scalar(dt)
|| matches!(dt, DataType::Decimal128(_, _) | DataType::Decimal256(_, _))
}
fn is_temporal_scalar(dt: &DataType) -> bool {
match dt {
DataType::Date32 | DataType::Date64 | DataType::Time32(_) | DataType::Time64(_) => true,
DataType::Timestamp(_, tz) => tz.as_ref().is_none_or(|tz| utc_equivalent_tz(tz)),
_ => false,
}
}
fn utc_equivalent_tz(tz: &str) -> bool {
matches!(
tz.trim(),
"UTC" | "utc" | "Utc" | "Z" | "z" | "+00:00" | "-00:00" | "00:00"
)
}
fn formattable_temporal_type(dt: &DataType) -> DataType {
match dt {
DataType::Timestamp(unit, Some(tz)) if utc_equivalent_tz(tz) => {
DataType::Timestamp(*unit, Some("+00:00".into()))
}
other => other.clone(),
}
}
fn field_metadata(
schema: &Schema,
geom_idx: Option<usize>,
published: &PublishedNames,
) -> HashMap<String, String> {
let geom_idx = geom_idx.unwrap_or(usize::MAX);
let mut out = HashMap::new();
for (i, f) in schema.fields().iter().enumerate() {
if i == geom_idx
|| f.name().eq_ignore_ascii_case(LEVEL_COLUMN)
|| published.is_suppressed(f.name())
{
continue;
}
let ty = match f.data_type() {
DataType::Utf8 | DataType::LargeUtf8 => "String",
DataType::Boolean => "Boolean",
dt if is_temporal_scalar(dt) => "String",
dt if is_supported_scalar(dt) => "Number",
_ => continue,
};
out.insert(published.publish(f.name()).to_string(), ty.to_string());
}
out
}
fn extract_property_column(col: &dyn arrow_array::Array) -> Vec<Option<PropertyValue>> {
let n = col.len();
macro_rules! prim {
($ty:ty, $variant:ident, $cast:ty) => {{
let a = col.as_primitive::<$ty>();
(0..n)
.map(|i| {
if a.is_null(i) {
None
} else {
Some(PropertyValue::$variant(a.value(i) as $cast))
}
})
.collect()
}};
}
macro_rules! strcol {
($off:ty) => {{
let a = col.as_string::<$off>();
(0..n)
.map(|i| {
if a.is_null(i) {
None
} else {
Some(PropertyValue::String(a.value(i).to_string()))
}
})
.collect()
}};
}
match col.data_type() {
DataType::Utf8 => strcol!(i32),
DataType::LargeUtf8 => strcol!(i64),
DataType::Boolean => {
let a = col.as_boolean();
(0..n)
.map(|i| {
if a.is_null(i) {
None
} else {
Some(PropertyValue::Bool(a.value(i)))
}
})
.collect()
}
DataType::Int8 => prim!(Int8Type, Int, i64),
DataType::Int16 => prim!(Int16Type, Int, i64),
DataType::Int32 => prim!(Int32Type, Int, i64),
DataType::Int64 => prim!(Int64Type, Int, i64),
DataType::UInt8 => prim!(UInt8Type, UInt, u64),
DataType::UInt16 => prim!(UInt16Type, UInt, u64),
DataType::UInt32 => prim!(UInt32Type, UInt, u64),
DataType::UInt64 => prim!(UInt64Type, UInt, u64),
DataType::Float32 => prim!(Float32Type, Float, f32),
DataType::Float64 => prim!(Float64Type, Double, f64),
DataType::Decimal128(_, scale) => {
let a = col.as_primitive::<Decimal128Type>();
let scale = i32::from(*scale);
let f = 10f64.powi(scale);
(0..n)
.map(|i| {
(!a.is_null(i)).then(|| {
let raw = a.value(i);
let exact = (scale <= 0)
.then(|| {
10i128
.checked_pow((-scale) as u32)
.and_then(|m| raw.checked_mul(m))
})
.flatten()
.and_then(|v| i64::try_from(v).ok());
match exact {
Some(v) => PropertyValue::Int(v),
None => PropertyValue::Double(raw as f64 / f),
}
})
})
.collect()
}
DataType::Decimal256(_, _) => {
let opts = FormatOptions::default();
match ArrayFormatter::try_new(col, &opts) {
Ok(fmt) => (0..n)
.map(|i| {
(!col.is_null(i))
.then(|| fmt.value(i).to_string().parse().ok())
.flatten()
.map(PropertyValue::Double)
})
.collect(),
Err(_) => vec![None; n],
}
}
dt if is_temporal_scalar(dt) => {
let opts = FormatOptions::default();
let retyped: Option<arrow_array::ArrayRef> = match formattable_temporal_type(dt) {
t if t == *dt => None,
t => col
.to_data()
.into_builder()
.data_type(t)
.build()
.ok()
.map(arrow_array::make_array),
};
let target: &dyn Array = retyped.as_deref().unwrap_or(col);
let out = match ArrayFormatter::try_new(target, &opts) {
Ok(fmt) => (0..n)
.map(|i| {
(!col.is_null(i)).then(|| PropertyValue::String(fmt.value(i).to_string()))
})
.collect(),
Err(e) => {
log::warn!("cannot format {dt:?} column as text, dropping it: {e}");
vec![None; n]
}
};
out
}
_ => vec![None; n],
}
}
fn detect_crs(path: &Path) -> Result<Crs, ExportError> {
let info = crate::quality::extract_crs(path).map_err(ExportError::Core)?;
if info.is_wgs84 {
return Ok(Crs::Epsg4326);
}
if let Some(id) = &info.identifier {
let up = id.to_uppercase();
if up.contains("3857") || up.contains("900913") {
return Ok(Crs::Epsg3857);
}
}
if info.identifier.is_none() && info.name.is_none() {
return Ok(Crs::Epsg4326);
}
Err(ExportError::UnsupportedCrs {
crs: info
.identifier
.or(info.name)
.unwrap_or_else(|| "unknown".to_string()),
})
}
#[inline]
fn webmerc_to_lnglat(x: f64, y: f64) -> (f64, f64) {
let lng = x / WEBMERC_HALF_M * 180.0;
let lat = (2.0 * (y / WEBMERC_HALF_M * std::f64::consts::PI).exp().atan()
- std::f64::consts::FRAC_PI_2)
.to_degrees();
(lng, lat)
}
fn reproject_3857_to_4326(g: &Geometry<f64>) -> Geometry<f64> {
g.map_coords(|c| {
let (lng, lat) = webmerc_to_lnglat(c.x, c.y);
geo::coord! { x: lng, y: lat }
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::mvt::{command_decode, zigzag_decode};
use crate::overview::level::{
gsd, ClusteringProvenance, CoalescingProvenance, Generalization, Level, Mode,
};
use crate::overview::writer::{
LevelSpec, LevelWriteOutcome, OverviewWriter, OverviewWriterOptions,
};
use crate::vector_tile::tile::GeomType;
use crate::vector_tile::Tile;
use arrow_array::{ArrayRef, Int64Array, RecordBatch, StringArray};
use arrow_schema::{Field, Schema};
use geo::{Geometry, LineString, Point};
use geoarrow::array::GeometryBuilder;
use geoarrow::datatypes::GeometryType;
fn build_geometry_array(geoms: &[Geometry<f64>]) -> geoarrow::array::GeometryArray {
let typ = GeometryType::new(Default::default());
let mut b = GeometryBuilder::new(typ).with_prefer_multi(false);
b.extend_from_iter(geoms.iter().map(Some));
b.finish()
}
fn geometry_field() -> Field {
build_geometry_array(&[Geometry::Point(Point::new(0.0, 0.0))])
.data_type()
.to_field("geometry", true)
}
fn source_schema() -> Schema {
Schema::new(vec![
Field::new("id", DataType::Int64, false),
Field::new("name", DataType::Utf8, false),
geometry_field(),
])
}
fn batch(schema: &Arc<Schema>, ids: &[i64], geoms: &[Geometry<f64>]) -> RecordBatch {
let id = Int64Array::from(ids.to_vec());
let name = StringArray::from(ids.iter().map(|i| format!("f{i}")).collect::<Vec<_>>());
let geom = build_geometry_array(geoms);
RecordBatch::try_new(
schema.clone(),
vec![Arc::new(id), Arc::new(name), Arc::new(geom.to_array_ref())],
)
.unwrap()
}
fn write_mode_fixture(
path: &Path,
level_geoms: &[(Vec<i64>, Vec<Geometry<f64>>)],
mode: Mode,
max_row_group_size: usize,
) -> OverviewsMeta {
let schema = Arc::new(source_schema());
let specs: Vec<LevelSpec> = (0..level_geoms.len())
.map(|k| LevelSpec::new(gsd((2 + 2 * k) as u8), Some((2 + 2 * k) as u8)))
.collect();
let mut opts = OverviewWriterOptions::new(mode, specs);
opts.max_row_group_size = max_row_group_size;
let mut writer = OverviewWriter::create(path, &schema, opts).unwrap();
for (k, (ids, geoms)) in level_geoms.iter().enumerate() {
assert_eq!(
writer
.write_level(
k,
Some(ids.len()),
std::iter::once(batch(&schema, ids, geoms)),
)
.unwrap(),
LevelWriteOutcome::Written
);
}
writer.finish().unwrap()
}
fn write_fixture(path: &Path, level_geoms: &[(Vec<i64>, Vec<Geometry<f64>>)]) -> OverviewsMeta {
write_mode_fixture(path, level_geoms, Mode::Duplicating, 10_000)
}
fn decode_tile(bytes: &[u8]) -> Tile {
Tile::decode(bytes).unwrap()
}
fn decode_coords(geom: &[u32]) -> Vec<(i32, i32)> {
let mut coords = Vec::new();
let (mut cx, mut cy) = (0i32, 0i32);
let mut i = 0;
while i < geom.len() {
let (cmd, count) = command_decode(geom[i]);
i += 1;
if cmd == 7 {
continue;
}
for _ in 0..count {
cx += zigzag_decode(geom[i]);
cy += zigzag_decode(geom[i + 1]);
coords.push((cx, cy));
i += 2;
}
}
coords
}
fn equivalence_fixture(path: &Path) {
let pa = Geometry::Point(Point::new(-120.0, 40.0));
let pb = Geometry::Point(Point::new(120.0, -40.0));
let wide = Geometry::LineString(LineString::from(vec![
(-100.0, 10.0),
(-80.0, 12.0),
(-60.0, 8.0),
(-40.0, 11.0),
]));
let concave = Geometry::Polygon(geo::Polygon::new(
LineString::from(vec![
(-100.0, 25.0),
(-98.0, 25.0),
(-98.0, 27.0),
(-99.0, 25.5),
(-100.0, 27.0),
(-100.0, 25.0),
]),
vec![],
));
let mut coords = Vec::new();
for k in 0..200 {
coords.push((30.0 + k as f64 * 0.3, -20.0 + (k as f64 * 0.1).sin()));
}
let wiggly = Geometry::LineString(LineString::from(coords));
write_fixture(
path,
&[
(vec![0, 1], vec![pa.clone(), pb.clone()]),
(vec![0, 1, 2, 3, 4], vec![pa, pb, wide, concave, wiggly]),
],
);
}
#[test]
fn export_archive_matches_pre_refactor_reference() {
let tin = tempfile::NamedTempFile::new().unwrap();
equivalence_fixture(tin.path());
let tout = tempfile::NamedTempFile::new().unwrap();
let opts = ExportOptions {
layer_name: "ref".to_string(),
..Default::default()
};
let report = export_pmtiles(tin.path(), tout.path(), &opts).unwrap();
assert_eq!(report.total_tiles, 13);
let bytes = std::fs::read(tout.path()).unwrap();
assert_eq!(
format!("{:016x}", crate::dedup::TileHasher::hash(&bytes)),
"0dfdf8a0bb941a0c",
"archive bytes diverged from the pre-refactor reference"
);
use crate::compression;
use crate::pmtiles_writer::{decode_directory, tile_id_to_zxy, Header};
let header = Header::from_bytes(&bytes).unwrap();
let root = compression::decompress(
&bytes[header.root_dir_offset as usize
..(header.root_dir_offset + header.root_dir_length) as usize],
header.internal_compression,
)
.unwrap();
let entries = decode_directory(&root).expect("root directory must decode");
assert!(
entries.iter().all(|e| e.run_length > 0),
"13 tiles must fit the root directory; this walk does not follow leaves"
);
let mut seen = 0usize;
for e in &entries {
let start = (header.tile_data_offset + e.offset) as usize;
let raw = &bytes[start..start + e.length as usize];
let plain = compression::decompress(raw, header.tile_compression).unwrap();
let decoded = decode_tile(&plain);
let features: usize = decoded.layers.iter().map(|l| l.features.len()).sum();
let zxy = tile_id_to_zxy(e.tile_id).unwrap();
assert!(
features > 0,
"tile {zxy:?} was emitted with no features -- membership is too wide"
);
seen += e.run_length as usize;
}
assert_eq!(
seen, 13,
"every counted tile must be addressable in the archive"
);
}
#[test]
fn partitioned_export_is_partition_invariant() {
let tin = tempfile::NamedTempFile::new().unwrap();
equivalence_fixture(tin.path());
let opts = ExportOptions {
layer_name: "ref".to_string(),
..Default::default()
};
let t_many = tempfile::NamedTempFile::new().unwrap();
let t_one = tempfile::NamedTempFile::new().unwrap();
let r_many =
export_pmtiles_with_partition_target(tin.path(), t_many.path(), &opts, 1).unwrap();
let r_one = export_pmtiles(tin.path(), t_one.path(), &opts).unwrap();
let b_many = std::fs::read(t_many.path()).unwrap();
let b_one = std::fs::read(t_one.path()).unwrap();
assert_eq!(
b_many, b_one,
"partitioned archive bytes diverge from single-partition archive"
);
assert_eq!(r_many.zooms, r_one.zooms);
assert_eq!(r_many.total_tiles, r_one.total_tiles);
assert_eq!(r_many.total_tile_features, r_one.total_tile_features);
assert_eq!(r_many.oversized_tiles, r_one.oversized_tiles);
}
#[test]
fn auto_partition_wave_memory_preflight() {
const MIB: u64 = 1024 * 1024;
const GIB: u64 = 1024 * MIB;
assert_eq!(auto_partition_wave(64, Some(256 * GIB)), 64);
assert_eq!(auto_partition_wave(24, Some(256 * GIB)), 24);
assert_eq!(auto_partition_wave(16, Some(54 * GIB)), 16);
assert_eq!(auto_partition_wave(64, Some(4 * GIB)), 32);
assert_eq!(auto_partition_wave(64, Some(2 * GIB)), 16);
assert_eq!(auto_partition_wave(64, Some(GIB)), 8);
assert_eq!(auto_partition_wave(16, Some(256 * MIB)), PARTITION_WAVE_MIN);
assert_eq!(auto_partition_wave(16, Some(0)), PARTITION_WAVE_MIN);
assert_eq!(auto_partition_wave(2, Some(256 * GIB)), PARTITION_WAVE_MIN);
assert_eq!(auto_partition_wave(64, None), PARTITION_WAVE_FALLBACK_MAX);
assert_eq!(auto_partition_wave(8, None), 8);
assert_eq!(auto_partition_wave(2, None), PARTITION_WAVE_MIN);
}
fn partitions_with_members(counts: &[usize]) -> Vec<Partition> {
counts
.iter()
.enumerate()
.map(|(i, &m)| Partition {
key_lo: i as u64,
key_hi: i as u64,
bbox: TileBounds::new(0.0, 0.0, 0.0, 0.0),
members: m,
})
.collect()
}
#[test]
fn memory_safe_level_wave_uses_densest_partition() {
const MIB: u64 = 1024 * 1024;
const GIB: u64 = 1024 * MIB;
let sparse = partitions_with_members(&[32_768, 30_000, 31_000]);
assert_eq!(
memory_safe_level_wave(16, &sparse, Some(512), Some(54 * GIB)),
16,
"a sparse level must not narrow below the ceiling"
);
let dense = partitions_with_members(&[6_500_000, 100, 200]);
assert_eq!(
memory_safe_level_wave(16, &dense, Some(600), Some(54 * GIB)),
3,
);
let extreme = partitions_with_members(&[60_000_000, 100]);
assert_eq!(
memory_safe_level_wave(16, &extreme, Some(600), Some(54 * GIB)),
1,
);
assert_eq!(
memory_safe_level_wave(4, &sparse, Some(512), Some(256 * GIB)),
4,
);
assert_eq!(
memory_safe_level_wave(16, &dense, None, Some(54 * GIB)),
16,
"unknown member size falls back to the ceiling"
);
assert_eq!(
memory_safe_level_wave(16, &dense, Some(0), Some(54 * GIB)),
16,
"a degenerate zero member size falls back to the ceiling"
);
assert_eq!(
memory_safe_level_wave(16, &dense, Some(600), None),
16,
"unprobeable RAM falls back to the ceiling"
);
assert_eq!(
memory_safe_level_wave(16, &[], Some(600), Some(54 * GIB)),
16,
"an empty level falls back to the ceiling"
);
}
#[test]
fn resolve_partition_wave_auto_and_explicit() {
let auto = resolve_partition_wave(PARTITION_WAVE_AUTO);
assert!(
auto >= PARTITION_WAVE_MIN,
"auto wave {auto} below floor {PARTITION_WAVE_MIN}"
);
let cores = std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(PARTITION_WAVE_MIN);
assert_eq!(
auto,
auto_partition_wave(cores, available_memory_bytes()),
"auto must equal the preflight decision on this box's cores + RAM"
);
assert_eq!(resolve_partition_wave(1), 1);
assert_eq!(resolve_partition_wave(6), 6);
assert_eq!(
resolve_partition_wave(PARTITION_WAVE_FALLBACK_MAX + 100),
PARTITION_WAVE_FALLBACK_MAX + 100
);
}
#[test]
fn export_wave_width_is_byte_invariant() {
let tin = tempfile::NamedTempFile::new().unwrap();
equivalence_fixture(tin.path());
let t_narrow = tempfile::NamedTempFile::new().unwrap();
let t_wide = tempfile::NamedTempFile::new().unwrap();
let narrow = ExportOptions {
layer_name: "ref".to_string(),
partition_wave: 1,
..Default::default()
};
let wide = ExportOptions {
layer_name: "ref".to_string(),
partition_wave: 16,
..Default::default()
};
export_pmtiles_with_partition_target(tin.path(), t_narrow.path(), &narrow, 1).unwrap();
export_pmtiles_with_partition_target(tin.path(), t_wide.path(), &wide, 1).unwrap();
assert_eq!(
std::fs::read(t_narrow.path()).unwrap(),
std::fs::read(t_wide.path()).unwrap(),
"archive bytes diverge between wave widths 1 and 16"
);
}
fn members_by_key(v: Vec<(u64, Geometry<f64>)>) -> HashMap<u64, Geometry<f64>> {
let mut m = HashMap::new();
for (k, g) in v {
assert!(m.insert(k, g).is_none(), "duplicate member for one feature");
}
m
}
fn tile_mvt(key: u64, geom: &Geometry<f64>, zoom: u8, opts: &ExportOptions) -> Vec<u8> {
let (x, y) = ((key >> 32) as u32, key as u32);
let tb = TileCoord::new(x, y, zoom).bounds();
let m = Member {
key,
seq: 0,
geom: geom.clone(),
props: Arc::new(Vec::new()),
};
build_mvt(std::iter::once(&m), &tb, opts)
}
fn assert_recursive_matches_direct(geom: &Geometry<f64>, zoom: u8, opts: &ExportOptions) {
let rec = members_by_key(members_recursive_vec(geom, zoom, opts, 0, u64::MAX));
let dir = members_by_key(members_direct_vec(geom, zoom, opts, 0, u64::MAX));
let mut rk: Vec<u64> = rec.keys().copied().collect();
rk.sort_unstable();
let mut dk: Vec<u64> = dir.keys().copied().collect();
dk.sort_unstable();
assert_eq!(rk, dk, "tile key set diverges at z{zoom}");
for k in rk {
assert_eq!(
tile_mvt(k, &rec[&k], zoom, opts),
tile_mvt(k, &dir[&k], zoom, opts),
"tile ({}, {}) MVT diverges at z{zoom}",
(k >> 32) as u32,
k as u32,
);
}
}
fn equivalence_corpus() -> Vec<(&'static str, Geometry<f64>, Vec<u8>)> {
let seam_poly = Geometry::Polygon(geo::Polygon::new(
LineString::from(vec![
(-73.985, 40.700),
(-73.955, 40.700),
(-73.955, 40.730),
(-73.985, 40.730),
(-73.985, 40.700),
]),
vec![],
));
let concave = Geometry::Polygon(geo::Polygon::new(
LineString::from(vec![
(-100.0, 25.0),
(-98.0, 25.0),
(-98.0, 27.0),
(-99.0, 25.5),
(-100.0, 27.0),
(-100.0, 25.0),
]),
vec![],
));
let big_box = Geometry::Polygon(geo::Polygon::new(
LineString::from(vec![
(-20.0, -20.0),
(20.0, -20.0),
(20.0, 20.0),
(-20.0, 20.0),
(-20.0, -20.0),
]),
vec![],
));
let mut wig = Vec::new();
for k in 0..200 {
wig.push((30.0 + k as f64 * 0.3, -20.0 + (k as f64 * 0.1).sin()));
}
let wiggly = Geometry::LineString(LineString::from(wig));
let antimeridian = Geometry::LineString(LineString::from(vec![
(170.0, 5.0),
(178.0, 8.0),
(-176.0, 6.0),
(-170.0, 9.0),
]));
vec![
(
"interior_point",
Geometry::Point(Point::new(-73.97, 40.71)),
vec![4, 8, 14],
),
("seam_poly", seam_poly, vec![10, 12, 14]),
("concave", concave, vec![4, 6, 8]),
("big_box", big_box, vec![3, 5, 6]),
("wiggly", wiggly, vec![4, 6]),
("antimeridian", antimeridian, vec![3, 5]),
]
}
#[test]
fn recursive_split_matches_direct_oracle() {
let opts = ExportOptions {
simple_clip_fastpath: false,
..Default::default()
};
for (name, geom, zooms) in equivalence_corpus() {
for z in zooms {
std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
assert_recursive_matches_direct(&geom, z, &opts);
}))
.unwrap_or_else(|_| panic!("equivalence failed for feature '{name}' at z{z}"));
}
}
}
#[test]
fn tile_buffer_is_tile_pixels_not_extent_units() {
let opts = ExportOptions::default();
assert_eq!(opts.tile_buffer, 8);
assert_eq!(opts.extent, 4096);
let in_mvt_units = buffer_fraction(&opts) * f64::from(opts.extent);
assert!(
(in_mvt_units - 128.0).abs() < 1e-9,
"8 tile pixels must be 128 MVT units at extent 4096, got {in_mvt_units}"
);
let finer = ExportOptions {
extent: 8192,
..Default::default()
};
assert!(
(buffer_fraction(&finer) - buffer_fraction(&opts)).abs() < 1e-12,
"extent must not affect the buffer width"
);
}
#[test]
fn point_in_neighbour_buffer_is_a_member_of_that_tile() {
let opts = ExportOptions::default();
let zoom = 1;
let buffer_deg = buffer_deg_at_zoom(zoom, &opts);
let west = tile_key(0, 0);
let east = tile_key(1, 0);
let inside = 0.9 * buffer_deg;
let keys = members_by_key(feature_tile_members(
&Geometry::Point(geo::Point::new(inside, 10.0)),
zoom,
&opts,
0,
u64::MAX,
));
assert!(keys.contains_key(&east), "must be in its own tile (1,0)");
assert!(
keys.contains_key(&west),
"point {inside}° is within tile (0,0)'s {buffer_deg}° buffer but was \
not made a member of it; emitted tiles: {:?}",
keys.keys().collect::<Vec<_>>()
);
let outside = 1.5 * buffer_deg;
let keys = members_by_key(feature_tile_members(
&Geometry::Point(geo::Point::new(outside, 10.0)),
zoom,
&opts,
0,
u64::MAX,
));
assert!(keys.contains_key(&east), "must be in its own tile (1,0)");
assert!(
!keys.contains_key(&west),
"point {outside}° is {buffer_deg}° beyond tile (0,0) and must not be \
a member of it; emitted tiles: {:?}",
keys.keys().collect::<Vec<_>>()
);
}
#[test]
fn out_of_domain_bbox_does_not_claim_every_tile_column() {
for tile_buffer in [0, 8, 128] {
let opts = ExportOptions {
tile_buffer,
..Default::default()
};
for zoom in [1u8, 10, 14] {
for lng in [180.5, -180.5] {
let bbox = TileBounds::new(lng, -17.0, lng, -17.0);
let widened = expand_bbox(&bbox, buffer_deg_at_zoom(zoom, &opts));
assert!(
widened.lng_min <= widened.lng_max,
"expand_bbox manufactured a wrapped bbox at lng {lng}, z{zoom}, \
buffer {tile_buffer}: {widened:?}"
);
let ranges = member_ranges(&bbox, zoom, &opts);
assert!(
ranges.x2.is_none(),
"out-of-domain bbox at lng {lng}, z{zoom}, buffer {tile_buffer} \
took the antimeridian branch: {ranges:?}"
);
let span = u64::from(ranges.x.1 - ranges.x.0 + 1);
assert!(
span <= 2,
"out-of-domain bbox at lng {lng}, z{zoom}, buffer {tile_buffer} \
claims {span} tile columns"
);
}
}
}
}
#[test]
fn genuinely_wrapped_bbox_survives_expansion() {
let opts = ExportOptions::default();
let bbox = TileBounds::new(170.0, 0.0, -170.0, 10.0);
let widened = expand_bbox(&bbox, buffer_deg_at_zoom(4, &opts));
assert!(
widened.lng_min > widened.lng_max,
"a genuinely wrapped bbox must not be collapsed: {widened:?}"
);
}
#[test]
fn scan_counts_cover_every_emitted_member_key() {
use std::collections::HashSet;
for tile_buffer in [0, 8, 64] {
let opts = ExportOptions {
tile_buffer,
..Default::default()
};
for zoom in [1u8, 4, 9] {
let geoms = [
Geometry::Point(geo::Point::new(0.05, 10.0)),
Geometry::Point(geo::Point::new(-179.9, -80.0)),
Geometry::Point(geo::Point::new(179.9, 84.0)),
Geometry::LineString(geo::LineString::from(vec![(-20.0, -10.0), (25.0, 30.0)])),
];
let bboxes: Vec<Option<TileBounds>> = geoms
.iter()
.map(|g| {
g.bounding_rect()
.map(|r| TileBounds::new(r.min().x, r.min().y, r.max().x, r.max().y))
})
.collect();
let (_, counts) = size_bboxes(&bboxes, zoom, &opts);
let planned: HashSet<u64> = counts.keys().copied().collect();
for g in &geoms {
for (key, _) in feature_tile_members(g, zoom, &opts, 0, u64::MAX) {
assert!(
planned.contains(&key),
"emitted key {key} at z{zoom} buffer {tile_buffer} was not \
planned by the scan pass"
);
}
}
}
}
}
#[test]
fn recursive_split_is_partition_range_invariant() {
let opts = ExportOptions::default();
for (name, geom, zooms) in equivalence_corpus() {
for z in zooms {
let full = members_by_key(feature_tile_members(&geom, z, &opts, 0, u64::MAX));
if full.is_empty() {
continue;
}
let mut keys: Vec<u64> = full.keys().copied().collect();
keys.sort_unstable();
let split = keys[keys.len() / 2];
let lo = feature_tile_members(&geom, z, &opts, 0, split);
let hi = feature_tile_members(&geom, z, &opts, split + 1, u64::MAX);
let mut union = HashMap::new();
for (k, g) in lo.into_iter().chain(hi) {
assert!(
union.insert(k, g).is_none(),
"tile {k:x} emitted by both partitions ('{name}' z{z})"
);
}
let mut uk: Vec<u64> = union.keys().copied().collect();
uk.sort_unstable();
assert_eq!(uk, keys, "partitioned tile set diverges ('{name}' z{z})");
for k in keys {
assert_eq!(
tile_mvt(k, &union[&k], z, &opts),
tile_mvt(k, &full[&k], z, &opts),
"partitioned tile {k:x} bytes diverge ('{name}' z{z})"
);
}
}
}
}
#[test]
fn zoom_mapping_uses_explicit_zoom() {
let meta = OverviewsMeta {
version: "0.1.0".to_string(),
mode: Some(Mode::Duplicating),
canonical_level: Some(1),
levels: vec![
Level {
row_group_end: 0,
gsd: gsd(4),
zoom: Some(4),
},
Level {
row_group_end: 1,
gsd: gsd(7),
zoom: Some(7),
},
],
generalization: None,
};
assert_eq!(zoom_for_level(&meta, 0), 4);
assert_eq!(zoom_for_level(&meta, 1), 7);
}
#[test]
fn zoom_mapping_derives_from_gsd_when_absent() {
let meta = OverviewsMeta {
version: "0.1.0".to_string(),
mode: Some(Mode::Duplicating),
canonical_level: Some(1),
levels: vec![
Level {
row_group_end: 0,
gsd: gsd(3),
zoom: None,
},
Level {
row_group_end: 1,
gsd: gsd(8),
zoom: None,
},
],
generalization: None,
};
assert_eq!(zoom_for_level(&meta, 0), 3);
assert_eq!(zoom_for_level(&meta, 1), 8);
}
#[test]
fn partitioning_prefix_semantics_via_reader() {
let schema = Arc::new(source_schema());
let tmp = tempfile::NamedTempFile::new().unwrap();
let specs = vec![
LevelSpec::new(gsd(4), Some(4)),
LevelSpec::new(gsd(6), Some(6)),
];
let mut opts = OverviewWriterOptions::new(Mode::Partitioning, specs);
opts.max_row_group_size = 10_000;
let mut w = OverviewWriter::create(tmp.path(), &schema, opts).unwrap();
assert_eq!(
w.write_level(
0,
Some(1),
std::iter::once(batch(
&schema,
&[0],
&[Geometry::Point(Point::new(1.0, 1.0))],
)),
)
.unwrap(),
LevelWriteOutcome::Written
);
assert_eq!(
w.write_level(
1,
Some(1),
std::iter::once(batch(
&schema,
&[1],
&[Geometry::Point(Point::new(2.0, 2.0))],
)),
)
.unwrap(),
LevelWriteOutcome::Written
);
w.finish().unwrap();
let reader = OverviewReader::open(tmp.path()).unwrap();
let l0 = read_level_features(&reader, 0, Crs::Epsg4326).unwrap();
let l1 = read_level_features(&reader, 1, Crs::Epsg4326).unwrap();
assert_eq!(l0.len(), 1);
assert_eq!(l1.len(), 2);
}
fn write_partitioning_fixture(
path: &Path,
level_geoms: &[(Vec<i64>, Vec<Geometry<f64>>)],
) -> OverviewsMeta {
write_mode_fixture(path, level_geoms, Mode::Partitioning, 2)
}
#[test]
fn fanout_scan_matches_per_level_prefix_scan_partitioning() {
let tmp = tempfile::NamedTempFile::new().unwrap();
let poly = |x: f64, y: f64| {
Geometry::Polygon(geo::Polygon::new(
LineString::from(vec![
(x, y),
(x + 8.0, y),
(x + 8.0, y + 6.0),
(x, y + 6.0),
(x, y),
]),
vec![],
))
};
let line = Geometry::LineString(LineString::from(vec![
(-70.0, -10.0),
(-40.0, 5.0),
(-10.0, -5.0),
(20.0, 8.0),
]));
let meta = write_partitioning_fixture(
tmp.path(),
&[
(
vec![0, 1, 2],
vec![
Geometry::Point(Point::new(-120.0, 40.0)),
Geometry::Point(Point::new(100.0, -30.0)),
poly(-100.0, 20.0),
],
),
(
vec![3, 4, 5, 6],
vec![
line,
Geometry::Point(Point::new(0.0, 0.0)),
Geometry::Point(Point::new(-119.5, 39.5)),
poly(10.0, -40.0),
],
),
(
vec![7, 8, 9, 10, 11],
vec![
Geometry::Point(Point::new(-118.0, 34.0)),
Geometry::Point(Point::new(-117.5, 33.8)),
Geometry::Point(Point::new(-117.0, 34.2)),
Geometry::Point(Point::new(2.0, 48.0)),
poly(30.0, 30.0),
],
),
],
);
let reader = OverviewReader::open(tmp.path()).unwrap();
assert!(matches!(reader.mode(), Mode::Partitioning));
let num_levels = reader.num_levels();
assert_eq!(num_levels, 3);
let opts = ExportOptions::default();
let scans = scan_all_levels(&reader, Crs::Epsg4326, &meta, &opts).unwrap();
assert_eq!(scans.len(), num_levels);
assert_eq!(scans[0].feature_count, 3);
assert_eq!(scans[1].feature_count, 7);
assert_eq!(scans[2].feature_count, 12);
for (level_idx, scan) in scans.iter().enumerate() {
let zoom = zoom_for_level(&meta, level_idx);
let oracle = scan_level(&reader, level_idx, Crs::Epsg4326, zoom, &opts).unwrap();
assert_eq!(
*scan, oracle,
"fan-out scan differs from per-level prefix scan at level {level_idx}"
);
assert!(!scan.tile_counts.is_empty());
assert!(scan.bounds.is_some());
}
}
#[test]
fn temporal_and_decimal_columns_survive_as_properties() {
use arrow_array::{Date32Array, Decimal128Array, TimestampMicrosecondArray};
let dates = Date32Array::from(vec![Some(20710), None]);
let got = extract_property_column(&dates);
assert_eq!(
got[0],
Some(PropertyValue::String("2026-09-14".to_string())),
"DATE must render as an ISO date"
);
assert_eq!(got[1], None, "null stays null");
let ts = TimestampMicrosecondArray::from(vec![Some(1_789_389_296_000_000), None]);
let got = extract_property_column(&ts);
let rendered = match &got[0] {
Some(PropertyValue::String(s)) => s.clone(),
other => panic!("TIMESTAMP must render as a string, got {other:?}"),
};
assert!(
rendered.starts_with("2026-09-14T12:34:56"),
"TIMESTAMP must render as ISO 8601, got {rendered:?}"
);
assert_eq!(got[1], None);
let dec = Decimal128Array::from(vec![Some(150_i128), None])
.with_precision_and_scale(4, 2)
.unwrap();
let got = extract_property_column(&dec);
match got[0] {
Some(PropertyValue::Double(v)) => assert!(
(v - 1.5).abs() < 1e-9,
"DECIMAL(4,2) 150 must be 1.5, got {v}"
),
ref other => panic!("DECIMAL must render as a double, got {other:?}"),
}
assert_eq!(got[1], None);
}
#[test]
fn utc_stamped_timestamps_encode_rather_than_vanishing() {
use arrow_array::TimestampMicrosecondArray;
use arrow_schema::{Field, TimeUnit};
let dt = DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into()));
let ts = TimestampMicrosecondArray::from(vec![Some(1_789_389_296_000_000), None])
.with_data_type(dt.clone());
let got = extract_property_column(&ts);
let rendered = match &got[0] {
Some(PropertyValue::String(s)) => s.clone(),
other => panic!("UTC-stamped TIMESTAMP must render as a string, got {other:?}"),
};
assert!(
rendered.starts_with("2026-09-14T12:34:56"),
"expected an ISO 8601 rendering, got {rendered:?}"
);
assert!(
rendered.ends_with('Z') || rendered.contains('+'),
"a tz-stamped value must carry its offset, got {rendered:?}"
);
assert_eq!(got[1], None, "null stays null");
let schema = Schema::new(vec![Field::new("acq", dt, true)]);
assert_eq!(
field_metadata(&schema, None, &PublishedNames::identity())
.get("acq")
.map(String::as_str),
Some("String")
);
assert_eq!(
property_columns(&schema, usize::MAX, &PublishedNames::identity()).len(),
1
);
}
#[test]
fn renamed_source_column_is_published_under_its_source_name() {
use arrow_schema::Field;
let schema = Schema::new(vec![
Field::new("id", DataType::Int64, false),
Field::new("level_", DataType::Float64, true),
Field::new(LEVEL_COLUMN, DataType::UInt8, false),
Field::new("geometry", DataType::Binary, false),
]);
let published = PublishedNames::from_renames(
&BTreeMap::from([("level_".to_string(), "level".to_string())]),
&schema,
HashSet::new(),
);
let cols = property_columns(&schema, 3, &published);
let names: Vec<&str> = cols.iter().map(|(_, n)| n.as_str()).collect();
assert_eq!(
names,
vec!["id", "level"],
"the renamed column must publish as `level`, and tylertoo's own \
`level` must still be excluded"
);
assert_eq!(cols[1].0, 1);
let fields = field_metadata(&schema, Some(3), &published);
assert_eq!(fields.get("level").map(String::as_str), Some("Number"));
assert!(
!fields.contains_key("level_"),
"the internal name must not be advertised: {fields:?}"
);
}
#[test]
fn rename_is_not_restored_when_the_source_name_is_still_occupied() {
use arrow_schema::Field;
let schema = Schema::new(vec![
Field::new("point_count_", DataType::Int64, true),
Field::new("point_count", DataType::Int64, false),
Field::new("geometry", DataType::Binary, false),
]);
let published = PublishedNames::from_renames(
&BTreeMap::from([("point_count_".to_string(), "point_count".to_string())]),
&schema,
HashSet::new(),
);
let names: Vec<String> = property_columns(&schema, 2, &published)
.into_iter()
.map(|(_, n)| n)
.collect();
assert_eq!(
names,
vec!["point_count_", "point_count"],
"restoring onto a live exported column would merge two columns"
);
}
#[test]
fn two_renames_onto_one_source_name_restore_at_most_one() {
use arrow_schema::Field;
let schema = Schema::new(vec![
Field::new("level_", DataType::Float64, true),
Field::new("level__", DataType::Float64, true),
Field::new("level", DataType::Int32, false),
Field::new("geometry", DataType::Binary, false),
]);
let published = PublishedNames::from_renames(
&BTreeMap::from([
("level_".to_string(), "level".to_string()),
("level__".to_string(), "level".to_string()),
]),
&schema,
HashSet::new(),
);
let names: Vec<String> = property_columns(&schema, 3, &published)
.into_iter()
.map(|(_, n)| n)
.collect();
assert_eq!(
names.iter().filter(|n| *n == "level").count(),
1,
"the source name may be published once, not twice: {names:?}"
);
assert_eq!(names, vec!["level", "level__"]);
}
#[test]
fn absent_rename_provenance_publishes_schema_names() {
use arrow_schema::Field;
let schema = Schema::new(vec![
Field::new("level_", DataType::Float64, true),
Field::new("geometry", DataType::Binary, false),
]);
let names: Vec<String> = property_columns(&schema, 1, &PublishedNames::identity())
.into_iter()
.map(|(_, n)| n)
.collect();
assert_eq!(names, vec!["level_"]);
}
#[test]
fn unrenderable_timezone_is_neither_selected_nor_advertised() {
use arrow_schema::{Field, TimeUnit};
let dt = DataType::Timestamp(TimeUnit::Microsecond, Some("America/New_York".into()));
assert!(
!is_supported_scalar(&dt),
"a named zone must not be claimed"
);
let schema = Schema::new(vec![Field::new("local", dt, true)]);
assert!(!field_metadata(&schema, None, &PublishedNames::identity()).contains_key("local"));
assert!(property_columns(&schema, usize::MAX, &PublishedNames::identity()).is_empty());
}
#[test]
fn unscaled_decimals_stay_exact() {
use arrow_array::Decimal128Array;
let id = 9_007_199_254_740_993_i128; let dec = Decimal128Array::from(vec![Some(id), None])
.with_precision_and_scale(38, 0)
.unwrap();
let got = extract_property_column(&dec);
assert_eq!(
got[0],
Some(PropertyValue::Int(9_007_199_254_740_993)),
"an unscaled decimal must survive exactly, not round through f64"
);
assert_eq!(got[1], None);
let neg = Decimal128Array::from(vec![Some(150_i128)])
.with_precision_and_scale(10, -3)
.unwrap();
assert_eq!(
extract_property_column(&neg)[0],
Some(PropertyValue::Int(150_000))
);
let huge = Decimal128Array::from(vec![Some(i128::MAX)])
.with_precision_and_scale(38, 0)
.unwrap();
assert!(matches!(
extract_property_column(&huge)[0],
Some(PropertyValue::Double(_))
));
}
#[test]
fn decimal256_encodes_as_a_double() {
use arrow_array::Decimal256Array;
use arrow_buffer::i256;
let dec = Decimal256Array::from(vec![Some(i256::from_i128(150)), None])
.with_precision_and_scale(40, 2)
.unwrap();
let got = extract_property_column(&dec);
match got[0] {
Some(PropertyValue::Double(v)) => {
assert!(
(v - 1.5).abs() < 1e-9,
"DECIMAL256(40,2) 150 must be 1.5, got {v}"
);
}
ref other => panic!("DECIMAL256 must render as a double, got {other:?}"),
}
assert_eq!(got[1], None);
}
#[test]
fn date64_and_time_columns_render_as_iso_text() {
use arrow_array::{Date64Array, Time32SecondArray, Time64MicrosecondArray};
let d64 = Date64Array::from(vec![Some(1_789_344_000_000)]);
assert_eq!(
extract_property_column(&d64)[0],
Some(PropertyValue::String("2026-09-14T00:00:00".to_string())),
"Date64 keeps arrow's datetime rendering -- truncating would discard \
a time component the type may legitimately carry"
);
let t32 = Time32SecondArray::from(vec![Some(45_296)]);
assert_eq!(
extract_property_column(&t32)[0],
Some(PropertyValue::String("12:34:56".to_string()))
);
let t64 = Time64MicrosecondArray::from(vec![Some(45_296_000_000)]);
match &extract_property_column(&t64)[0] {
Some(PropertyValue::String(s)) => assert!(s.starts_with("12:34:56"), "got {s:?}"),
other => panic!("Time64 must render as a string, got {other:?}"),
}
}
#[test]
fn temporal_properties_reach_a_decoded_tile() {
use arrow_array::{Date32Array, TimestampMicrosecondArray};
use arrow_schema::TimeUnit;
let tz = DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into()));
let schema = Arc::new(Schema::new(vec![
Field::new("acq_date", DataType::Date32, false),
Field::new("acq_ts", tz.clone(), false),
geometry_field(),
]));
let geoms = vec![Geometry::Point(Point::new(-120.0, 40.0))];
let tmp = tempfile::NamedTempFile::new().unwrap();
let mut opts =
OverviewWriterOptions::new(Mode::Duplicating, vec![LevelSpec::new(gsd(2), Some(2))]);
opts.max_row_group_size = 10_000;
let mut writer = OverviewWriter::create(tmp.path(), &schema, opts).unwrap();
let rb = RecordBatch::try_new(
schema.clone(),
vec![
Arc::new(Date32Array::from(vec![20710])),
Arc::new(
TimestampMicrosecondArray::from(vec![1_789_389_296_000_000]).with_data_type(tz),
),
Arc::new(build_geometry_array(&geoms).to_array_ref()),
],
)
.unwrap();
assert_eq!(
writer.write_level(0, Some(1), std::iter::once(rb)).unwrap(),
LevelWriteOutcome::Written
);
writer.finish().unwrap();
let reader = OverviewReader::open(tmp.path()).unwrap();
let feats = read_level_features(&reader, 0, Crs::Epsg4326).unwrap();
let tiles = encode_level_tiles(&feats, 2, &ExportOptions::default());
assert_eq!(tiles.len(), 1);
let decoded = decode_tile(&tiles[0].data);
let layer = &decoded.layers[0];
for want in ["acq_date", "acq_ts"] {
assert!(
layer.keys.iter().any(|k| k == want),
"{want} never reached the tile; keys were {:?}",
layer.keys
);
}
let strings: Vec<&str> = layer
.values
.iter()
.filter_map(|v| v.string_value.as_deref())
.collect();
assert!(
strings.contains(&"2026-09-14"),
"DATE value missing from the tile: {strings:?}"
);
assert!(
strings.iter().any(|v| v.starts_with("2026-09-14T12:34:56")),
"TIMESTAMP value missing from the tile: {strings:?}"
);
}
#[test]
fn temporal_and_decimal_columns_are_advertised_in_field_metadata() {
use arrow_schema::{Field, TimeUnit};
let schema = Schema::new(vec![
Field::new("d", DataType::Date32, true),
Field::new("ts", DataType::Timestamp(TimeUnit::Microsecond, None), true),
Field::new("dec", DataType::Decimal128(10, 3), true),
]);
let meta = field_metadata(&schema, None, &PublishedNames::identity());
assert_eq!(meta.get("d").map(String::as_str), Some("String"));
assert_eq!(meta.get("ts").map(String::as_str), Some("String"));
assert_eq!(meta.get("dec").map(String::as_str), Some("Number"));
}
#[test]
fn export_feature_counts_and_level_absent_from_props() {
let a = Geometry::Point(Point::new(-120.0, 40.0));
let b = Geometry::Point(Point::new(120.0, -40.0));
let tmp = tempfile::NamedTempFile::new().unwrap();
write_fixture(
tmp.path(),
&[
(vec![0, 1], vec![a.clone(), b.clone()]),
(vec![0, 1], vec![a.clone(), b.clone()]),
],
);
let reader = OverviewReader::open(tmp.path()).unwrap();
let feats = read_level_features(&reader, 0, Crs::Epsg4326).unwrap();
assert_eq!(feats.len(), 2);
let tiles = encode_level_tiles(&feats, 2, &ExportOptions::default());
assert_eq!(tiles.len(), 2, "two far-apart points => two tiles");
for t in &tiles {
assert_eq!(t.feature_count, 1);
let decoded = decode_tile(&t.data);
let layer = &decoded.layers[0];
assert_eq!(layer.features.len(), 1);
assert!(
!layer.keys.iter().any(|k| k.eq_ignore_ascii_case("level")),
"level column leaked into MVT properties: {:?}",
layer.keys
);
assert!(layer.keys.iter().any(|k| k == "id"));
assert!(layer.keys.iter().any(|k| k == "name"));
assert_eq!(
decoded.layers[0].features[0].r#type,
Some(GeomType::Point as i32)
);
}
}
#[test]
fn clipped_features_within_tile_plus_buffer() {
let line = Geometry::LineString(LineString::from(vec![
(-100.0, 10.0),
(-80.0, 12.0),
(-60.0, 8.0),
(-40.0, 11.0),
]));
let tmp = tempfile::NamedTempFile::new().unwrap();
write_fixture(
tmp.path(),
&[(vec![0], vec![line.clone()]), (vec![0], vec![line.clone()])],
);
let reader = OverviewReader::open(tmp.path()).unwrap();
let feats = read_level_features(&reader, 1, Crs::Epsg4326).unwrap();
let opts = ExportOptions::default();
let tiles = encode_level_tiles(&feats, 4, &opts);
assert!(tiles.len() >= 2, "wide line must span multiple tiles");
let extent = opts.extent as i32;
let slack = extent * opts.tile_buffer as i32 / 256 + 4;
for t in &tiles {
let decoded = decode_tile(&t.data);
for f in &decoded.layers[0].features {
for (x, y) in decode_coords(&f.geometry) {
assert!(
x >= -slack - extent && x <= extent + slack + extent,
"x {x} outside buffered tile bounds"
);
assert!(
y >= -slack - extent && y <= extent + slack + extent,
"y {y} outside buffered tile bounds"
);
}
}
}
}
#[test]
fn bbox_contained_feature_bypasses_clip_with_identical_output() {
let poly = Geometry::Polygon(geo::Polygon::new(
LineString::from(vec![
(-100.0, 25.0),
(-98.0, 25.0),
(-98.0, 27.0),
(-99.0, 25.5), (-100.0, 27.0),
(-100.0, 25.0),
]),
vec![],
));
let feats = vec![Feature {
geom: poly.clone(),
props: vec![],
}];
let opts = ExportOptions::default();
let tiles = encode_level_tiles(&feats, 4, &opts);
assert_eq!(tiles.len(), 1, "fully-contained feature => exactly 1 tile");
assert_eq!(tiles[0].feature_count, 1);
let tc = TileCoord::new(tiles[0].x, tiles[0].y, 4);
let member = Member {
key: tile_key(tiles[0].x, tiles[0].y),
seq: 0,
geom: poly.clone(),
props: Arc::new(vec![]),
};
let expected = build_mvt([&member], &tc.bounds(), &opts);
assert_eq!(
tiles[0].data, expected,
"contained feature must bypass the clip (geometry emitted as-is)"
);
}
#[test]
fn seam_crossing_feature_still_clips() {
let line = Geometry::LineString(LineString::from(vec![(-100.0, 10.0), (-40.0, 11.0)]));
let feats = vec![Feature {
geom: line,
props: vec![],
}];
let opts = ExportOptions::default();
let tiles = encode_level_tiles(&feats, 4, &opts);
assert!(tiles.len() >= 2, "seam-crossing line must span tiles");
let extent = opts.extent as i32;
let slack = opts.tile_buffer as i32 * extent / 256 + 4;
for t in &tiles {
let decoded = decode_tile(&t.data);
for f in &decoded.layers[0].features {
for (x, y) in decode_coords(&f.geometry) {
assert!(
x >= -slack && x <= extent + slack && y >= -slack && y <= extent + slack,
"({x},{y}) outside buffered tile: geometry was not clipped"
);
}
}
}
}
#[test]
fn oversized_valve_fires_with_tiny_limit() {
let mut coords = Vec::new();
for k in 0..200 {
coords.push((-100.0 + k as f64 * 0.001, 40.0 + (k as f64 * 0.01).sin()));
}
let line = Geometry::LineString(LineString::from(coords));
let tmp = tempfile::NamedTempFile::new().unwrap();
let geoms: Vec<Geometry<f64>> = (0..20).map(|_| line.clone()).collect();
let ids: Vec<i64> = (0..20).collect();
write_fixture(
tmp.path(),
&[(ids.clone(), geoms.clone()), (ids.clone(), geoms.clone())],
);
let reader = OverviewReader::open(tmp.path()).unwrap();
let feats = read_level_features(&reader, 1, Crs::Epsg4326).unwrap();
let no_limit = ExportOptions {
tile_size_limit: None,
..Default::default()
};
let none = encode_level_tiles(&feats, 6, &no_limit);
assert!(none.iter().all(|t| !t.oversized));
let full_count: usize = none.iter().map(|t| t.feature_count).sum();
let opts = ExportOptions {
tile_size_limit: Some(64),
..Default::default()
};
let limited = encode_level_tiles(&feats, 6, &opts);
assert!(
limited.iter().any(|t| t.oversized),
"tiny --tile-size-limit must trip the valve"
);
let limited_count: usize = limited.iter().map(|t| t.feature_count).sum();
assert!(
limited_count < full_count,
"valve must drop features ({limited_count} !< {full_count})"
);
}
#[test]
fn default_export_caps_tiles_at_500_kib() {
assert_eq!(DEFAULT_TILE_SIZE_LIMIT, 500 * 1024);
assert_eq!(
ExportOptions::default().tile_size_limit,
Some(DEFAULT_TILE_SIZE_LIMIT)
);
}
fn point_member(x: f64, y: f64, seq: u64) -> Member {
Member {
key: 0,
seq,
geom: Geometry::Point(Point::new(x, y)),
props: Arc::new(Vec::new()),
}
}
fn ordered_member(seq: u64, level: f64) -> Member {
Member {
key: 0,
seq,
geom: Geometry::Point(Point::new(seq as f64 * 0.001, 0.0)),
props: Arc::new(vec![("level".to_string(), PropertyValue::Double(level))]),
}
}
fn paint_order(members: Vec<Member>, order: &FeatureOrder) -> Vec<u64> {
let mut m = members;
sort_members_for_encode(&mut m, order);
m.iter().map(|x| x.seq).collect()
}
#[test]
fn default_feature_order_is_input_row_order() {
let members = vec![
ordered_member(7, 0.5),
ordered_member(2, 0.2),
ordered_member(5, 0.35),
];
assert_eq!(paint_order(members, &FeatureOrder::Input), vec![2, 5, 7]);
}
#[test]
fn column_feature_order_sorts_ascending_within_the_tile() {
let members = vec![
ordered_member(0, 0.5),
ordered_member(1, 0.2),
ordered_member(2, 0.425),
ordered_member(3, 0.275),
];
let order = FeatureOrder::Column {
name: "level".to_string(),
descending: false,
};
assert_eq!(paint_order(members, &order), vec![1, 3, 2, 0]);
}
#[test]
fn column_feature_order_sorts_descending_when_asked() {
let members = vec![
ordered_member(0, 0.5),
ordered_member(1, 0.2),
ordered_member(2, 0.425),
];
let order = FeatureOrder::Column {
name: "level".to_string(),
descending: true,
};
assert_eq!(paint_order(members, &order), vec![0, 2, 1]);
}
#[test]
fn column_feature_order_breaks_ties_by_input_order() {
let members = vec![
ordered_member(9, 0.2),
ordered_member(4, 0.2),
ordered_member(6, 0.2),
];
let order = FeatureOrder::Column {
name: "level".to_string(),
descending: false,
};
assert_eq!(paint_order(members, &order), vec![4, 6, 9]);
}
#[test]
fn features_missing_the_order_column_sort_first() {
let mut bare = ordered_member(0, 0.0);
bare.props = Arc::new(Vec::new());
let members = vec![ordered_member(1, 0.5), bare, ordered_member(2, 0.2)];
let order = FeatureOrder::Column {
name: "level".to_string(),
descending: false,
};
assert_eq!(paint_order(members, &order), vec![0, 2, 1]);
}
#[test]
fn nan_in_the_order_column_does_not_abort_the_sort() {
let members: Vec<Member> = (0..64)
.map(|i| {
ordered_member(
i,
if i.is_multiple_of(3) {
f64::NAN
} else {
(64 - i) as f64
},
)
})
.collect();
let order = FeatureOrder::Column {
name: "level".to_string(),
descending: false,
};
let painted = paint_order(members, &order);
assert_eq!(painted.len(), 64, "no feature may be lost");
let level_of = |seq: u64| {
if seq.is_multiple_of(3) {
f64::NAN
} else {
(64 - seq) as f64
}
};
let ranked: Vec<f64> = painted
.iter()
.map(|s| level_of(*s))
.filter(|v| !v.is_nan())
.collect();
assert!(
ranked.windows(2).all(|w| w[0] <= w[1]),
"the comparable values must be ascending, got {ranked:?}"
);
let first_number = painted.iter().position(|s| !level_of(*s).is_nan()).unwrap();
assert!(
painted[..first_number]
.iter()
.all(|s| level_of(*s).is_nan()),
"NaN sorts with the unrankable features, underneath"
);
}
#[test]
fn column_feature_order_never_reorders_across_tiles() {
let mk = |key: u64, seq: u64, level: f64| Member {
key,
seq,
geom: Geometry::Point(Point::new(0.0, 0.0)),
props: Arc::new(vec![("level".to_string(), PropertyValue::Double(level))]),
};
let mut members = vec![mk(1, 0, 0.1), mk(0, 1, 0.9), mk(0, 2, 0.2)];
sort_members_for_encode(
&mut members,
&FeatureOrder::Column {
name: "level".to_string(),
descending: false,
},
);
let keys: Vec<u64> = members.iter().map(|m| m.key).collect();
assert_eq!(keys, vec![0, 0, 1], "tiles must stay grouped and ascending");
assert_eq!(
members.iter().map(|m| m.seq).collect::<Vec<_>>(),
vec![2, 1, 0],
"sorting applies within each tile, not across the whole level"
);
}
#[test]
fn numeric_order_keys_compare_across_property_variants() {
let mk = |seq: u64, v: PropertyValue| Member {
key: 0,
seq,
geom: Geometry::Point(Point::new(0.0, 0.0)),
props: Arc::new(vec![("n".to_string(), v)]),
};
let members = vec![
mk(0, PropertyValue::Int(10)),
mk(1, PropertyValue::Double(2.5)),
mk(2, PropertyValue::UInt(7)),
];
let order = FeatureOrder::Column {
name: "n".to_string(),
descending: false,
};
assert_eq!(paint_order(members, &order), vec![1, 2, 0]);
}
#[test]
fn feature_order_parses_the_cli_spellings() {
assert_eq!(
"input".parse::<FeatureOrder>().unwrap(),
FeatureOrder::Input
);
assert_eq!(
"level".parse::<FeatureOrder>().unwrap(),
FeatureOrder::Column {
name: "level".to_string(),
descending: false
}
);
assert_eq!(
"level:desc".parse::<FeatureOrder>().unwrap(),
FeatureOrder::Column {
name: "level".to_string(),
descending: true
}
);
assert_eq!(
"level:asc".parse::<FeatureOrder>().unwrap(),
FeatureOrder::Column {
name: "level".to_string(),
descending: false
}
);
assert!("level:dsc".parse::<FeatureOrder>().is_err());
assert!("".parse::<FeatureOrder>().is_err());
assert!(":asc".parse::<FeatureOrder>().is_err());
}
#[test]
fn drop_valve_point_tile_keeps_uniform_stride() {
let members: Vec<Member> = (0..10).map(|i| point_member(i as f64, 0.0, i)).collect();
let kept = select_kept_members(&members, 5);
let xs: Vec<f64> = kept
.iter()
.map(|m| match &m.geom {
Geometry::Point(p) => p.x(),
_ => unreachable!("built points"),
})
.collect();
assert_eq!(xs, vec![0.0, 2.0, 4.0, 6.0, 8.0]);
}
#[test]
fn drop_valve_sized_tile_keeps_largest() {
let members: Vec<Member> = (1..=5)
.map(|n| {
let coords: Vec<(f64, f64)> = (0..n).map(|k| (k as f64, 0.0)).collect();
Member {
key: 0,
seq: n as u64,
geom: Geometry::LineString(LineString::from(coords)),
props: Arc::new(Vec::new()),
}
})
.collect();
let kept = select_kept_members(&members, 2);
let counts: Vec<usize> = kept.iter().map(|m| m.geom.coords_count()).collect();
assert_eq!(counts, vec![5, 4]);
}
#[test]
fn absent_order_column_is_reported() {
use arrow_schema::Field;
let schema = Schema::new(vec![
Field::new("elevation", DataType::Float64, true),
Field::new("level", DataType::Int32, false),
Field::new("geometry", DataType::Binary, false),
]);
let column = |n: &str| FeatureOrder::Column {
name: n.to_string(),
descending: false,
};
let plain = PublishedNames::identity();
assert!(!order_column_is_absent(
&schema,
&plain,
&column("elevation")
));
assert!(
order_column_is_absent(&schema, &plain, &column("elevatoin")),
"a misspelled column must be detectable"
);
assert!(!order_column_is_absent(
&schema,
&plain,
&FeatureOrder::Input
));
assert!(order_column_is_absent(&schema, &plain, &column("level")));
}
#[test]
fn order_column_resolves_against_the_restored_source_name() {
use arrow_schema::Field;
let schema = Schema::new(vec![
Field::new("level_", DataType::Float64, true),
Field::new("level", DataType::Int32, false),
Field::new("geometry", DataType::Binary, false),
]);
let published = PublishedNames::from_renames(
&BTreeMap::from([("level_".to_string(), "level".to_string())]),
&schema,
HashSet::new(),
);
let column = |n: &str| FeatureOrder::Column {
name: n.to_string(),
descending: false,
};
assert!(
!order_column_is_absent(&schema, &published, &column("level")),
"the restored source name is what the tile carries"
);
assert!(
order_column_is_absent(&schema, &published, &column("level_")),
"the internal name is not published, so sorting by it would no-op"
);
}
#[test]
fn oversized_valve_keeps_the_requested_feature_order() {
let members: Vec<Member> = (1..=5u64)
.map(|i| {
let coords: Vec<(f64, f64)> = (0..i * 10)
.map(|v| (v as f64 * 1e-6, v as f64 * 1e-6))
.collect();
Member {
key: 0,
seq: i,
geom: Geometry::LineString(coords.into_iter().collect()),
props: Arc::new(vec![("level".to_string(), PropertyValue::Double(i as f64))]),
}
})
.collect();
let ordered = ExportOptions {
feature_order: FeatureOrder::Column {
name: "level".to_string(),
descending: false,
},
..ExportOptions::default()
};
let kept = shed_to_fit(&members, 3, &ordered);
let levels: Vec<u64> = kept.iter().map(|m| m.seq).collect();
assert_eq!(
levels,
vec![3, 4, 5],
"the three largest survive, but they must be painted in ascending \
`level` order, not in the valve's largest-first ranking"
);
let desc = ExportOptions {
feature_order: FeatureOrder::Column {
name: "level".to_string(),
descending: true,
},
..ExportOptions::default()
};
let kept: Vec<u64> = shed_to_fit(&members, 3, &desc)
.iter()
.map(|m| m.seq)
.collect();
assert_eq!(kept, vec![5, 4, 3]);
let kept: Vec<u64> = shed_to_fit(&members, 3, &ExportOptions::default())
.iter()
.map(|m| m.seq)
.collect();
assert_eq!(kept, vec![5, 4, 3], "default order must not change");
}
#[test]
fn size_limit_zero_disables_valve() {
let members: Vec<Member> = (0..64)
.map(|i| point_member(-100.0 + i as f64 * 0.001, 40.0, i))
.collect();
let tb = TileCoord::new(0, 0, 0).bounds();
let opts = ExportOptions {
tile_size_limit: Some(0),
..Default::default()
};
let (_data, count, oversized) = encode_tile(&members, &tb, &opts);
assert!(!oversized, "Some(0) must disable the cap");
assert_eq!(count, members.len());
}
#[test]
fn member_spill_codec_roundtrip_is_exact() {
use geo::{
coord, GeometryCollection, Line, MultiLineString, MultiPoint, MultiPolygon, Rect,
Triangle,
};
let poly = geo::Polygon::new(
LineString::from(vec![
(0.0, 0.0),
(4.0, 0.0),
(4.0, 4.0),
(0.0, 4.0),
(0.0, 0.0),
]),
vec![LineString::from(vec![
(1.0, 1.0),
(2.0, 1.0),
(2.0, 2.0),
(1.0, 2.0),
(1.0, 1.0),
])],
);
let geoms: Vec<Geometry<f64>> = vec![
Geometry::Point(Point::new(-0.0, std::f64::consts::PI)),
Geometry::Line(Line::new(
coord! { x: -1.5, y: 2.5 },
coord! { x: 3.25, y: -4.75 },
)),
Geometry::LineString(LineString::from(vec![(1.1, 2.2), (3.3, 4.4)])),
Geometry::Polygon(poly.clone()),
Geometry::MultiPoint(MultiPoint::from(vec![
Point::new(1.0, 2.0),
Point::new(-3.0, -4.0),
])),
Geometry::MultiLineString(MultiLineString::new(vec![LineString::from(vec![
(0.0, 1.0),
(2.0, 3.0),
])])),
Geometry::MultiPolygon(MultiPolygon::new(vec![poly.clone()])),
Geometry::GeometryCollection(GeometryCollection::from(vec![
Geometry::Point(Point::new(9.0, 9.0)),
Geometry::Polygon(poly),
])),
Geometry::Rect(Rect::new(
coord! { x: 0.0, y: 0.0 },
coord! { x: 1.0, y: 1.0 },
)),
Geometry::Triangle(Triangle::new(
coord! { x: 0.0, y: 0.0 },
coord! { x: 1.0, y: 0.0 },
coord! { x: 0.0, y: 1.0 },
)),
];
let props: Vec<(String, PropertyValue)> = vec![
("s".to_string(), PropertyValue::String("héllo".to_string())),
("f".to_string(), PropertyValue::Float(-0.0f32)),
("d".to_string(), PropertyValue::Double(f64::MIN_POSITIVE)),
("i".to_string(), PropertyValue::Int(-42)),
("u".to_string(), PropertyValue::UInt(u64::MAX)),
("b".to_string(), PropertyValue::Bool(true)),
];
for (i, g) in geoms.into_iter().enumerate() {
let m = Member {
key: tile_key(7, 11) + i as u64,
seq: u64::MAX - i as u64,
geom: g,
props: Arc::new(props.clone()),
};
let mut buf = Vec::new();
encode_member(&mut buf, &m);
let mut cur = SpillCursor::new(&buf);
let back = decode_member(&mut cur).unwrap();
assert!(cur.is_empty(), "trailing bytes after decode (geometry {i})");
assert_eq!(back.key, m.key);
assert_eq!(back.seq, m.seq);
assert_eq!(*back.props, *m.props);
let mut buf2 = Vec::new();
encode_member(&mut buf2, &back);
assert_eq!(buf, buf2, "re-encoded member bytes differ (geometry {i})");
}
}
#[test]
fn member_store_spill_segments_roundtrip() {
let plans = vec![LevelPlan {
zoom: 4,
partitions: partitions_with_members(&[10, 10]),
wave: 1,
}];
let mut store = MemberStore::new(&plans, SinkBacking::Spill).unwrap();
let mk = |key: u64, seq: u64| Member {
key,
seq,
geom: Geometry::Point(Point::new(seq as f64, -1.0)),
props: Arc::new(vec![("id".to_string(), PropertyValue::Int(seq as i64))]),
};
store.push(0, 0, mk(0, 0)).unwrap();
store.push(0, 1, mk(1, 1)).unwrap();
store.flush_bucket(0, 0).unwrap();
store.push(0, 0, mk(0, 2)).unwrap();
store.finish_fill().unwrap();
let w0 = store.take_wave(0, 0).unwrap();
let w1 = store.take_wave(0, 1).unwrap();
assert_eq!(w0.iter().map(|m| m.seq).collect::<Vec<_>>(), vec![0, 2]);
assert_eq!(w1.iter().map(|m| m.seq).collect::<Vec<_>>(), vec![1]);
assert!(w0.iter().all(|m| m.key == 0));
assert!(w1.iter().all(|m| m.key == 1));
assert_eq!(w0[0].props[0].1, PropertyValue::Int(0));
}
fn partitioning_equivalence_fixture(path: &Path) {
let wide = Geometry::LineString(LineString::from(vec![
(-100.0, 10.0),
(-80.0, 12.0),
(-60.0, 8.0),
(-40.0, 11.0),
]));
let concave = Geometry::Polygon(geo::Polygon::new(
LineString::from(vec![
(-100.0, 25.0),
(-98.0, 25.0),
(-98.0, 27.0),
(-99.0, 25.5),
(-100.0, 27.0),
(-100.0, 25.0),
]),
vec![],
));
let big_box = Geometry::Polygon(geo::Polygon::new(
LineString::from(vec![
(-20.0, -20.0),
(20.0, -20.0),
(20.0, 20.0),
(-20.0, 20.0),
(-20.0, -20.0),
]),
vec![],
));
let mut coords = Vec::new();
for k in 0..200 {
coords.push((30.0 + k as f64 * 0.3, -20.0 + (k as f64 * 0.1).sin()));
}
let wiggly = Geometry::LineString(LineString::from(coords));
write_partitioning_fixture(
path,
&[
(
vec![0, 1, 2],
vec![
Geometry::Point(Point::new(-120.0, 40.0)),
Geometry::Point(Point::new(120.0, -40.0)),
big_box,
],
),
(
vec![3, 4, 5, 6],
vec![
wide,
Geometry::Point(Point::new(0.0, 0.0)),
Geometry::Point(Point::new(-119.5, 39.5)),
concave,
],
),
(
vec![7, 8, 9, 10],
vec![
wiggly,
Geometry::Point(Point::new(-118.0, 34.0)),
Geometry::Point(Point::new(2.0, 48.0)),
Geometry::Point(Point::new(-117.5, 33.8)),
],
),
],
);
}
#[test]
fn partitioning_single_read_pass2_matches_legacy() {
let tin = tempfile::NamedTempFile::new().unwrap();
partitioning_equivalence_fixture(tin.path());
for (target, wave) in [(1usize, 1usize), (1, 3), (DEFAULT_PARTITION_TARGET, 2)] {
let opts = ExportOptions {
layer_name: "ref".to_string(),
partition_wave: wave,
..Default::default()
};
let t_legacy = tempfile::NamedTempFile::new().unwrap();
let r_legacy =
export_pmtiles_impl(tin.path(), t_legacy.path(), &opts, target, true, None)
.unwrap();
let legacy_bytes = std::fs::read(t_legacy.path()).unwrap();
for backing in [SinkBacking::Ram, SinkBacking::Spill] {
let t_new = tempfile::NamedTempFile::new().unwrap();
let r_new = export_pmtiles_impl(
tin.path(),
t_new.path(),
&opts,
target,
false,
Some(backing),
)
.unwrap();
assert_eq!(
std::fs::read(t_new.path()).unwrap(),
legacy_bytes,
"single-read ({backing:?}, target {target}, wave {wave}) \
archive diverges from legacy pass 2"
);
assert_eq!(r_new.zooms, r_legacy.zooms);
assert_eq!(r_new.total_tiles, r_legacy.total_tiles);
assert_eq!(r_new.total_tile_features, r_legacy.total_tile_features);
assert_eq!(r_new.oversized_tiles, r_legacy.oversized_tiles);
}
}
}
#[test]
fn partitioning_export_is_partition_invariant() {
let tin = tempfile::NamedTempFile::new().unwrap();
partitioning_equivalence_fixture(tin.path());
let opts = ExportOptions {
layer_name: "ref".to_string(),
..Default::default()
};
let t_many = tempfile::NamedTempFile::new().unwrap();
let t_one = tempfile::NamedTempFile::new().unwrap();
let r_many =
export_pmtiles_with_partition_target(tin.path(), t_many.path(), &opts, 1).unwrap();
let r_one = export_pmtiles(tin.path(), t_one.path(), &opts).unwrap();
assert_eq!(
std::fs::read(t_many.path()).unwrap(),
std::fs::read(t_one.path()).unwrap(),
"partitioning archives diverge across partition targets"
);
assert_eq!(r_many.zooms, r_one.zooms);
assert_eq!(r_many.total_tiles, r_one.total_tiles);
}
#[test]
fn export_pmtiles_writes_nonempty_archive() {
let a = Geometry::Point(Point::new(-120.0, 40.0));
let b = Geometry::Point(Point::new(120.0, -40.0));
let tin = tempfile::NamedTempFile::new().unwrap();
write_fixture(
tin.path(),
&[
(vec![0, 1], vec![a.clone(), b.clone()]),
(vec![0, 1], vec![a.clone(), b.clone()]),
],
);
let tout = tempfile::NamedTempFile::new().unwrap();
let report = export_pmtiles(tin.path(), tout.path(), &ExportOptions::default()).unwrap();
assert_eq!(report.min_zoom, 2);
assert_eq!(report.max_zoom, 4);
assert_eq!(report.zooms.len(), 2);
assert!(report.total_tiles >= 2);
let meta = std::fs::metadata(tout.path()).unwrap();
assert!(meta.len() > 127, "archive must be larger than the header");
let mut magic = [0u8; 7];
use std::io::Read;
std::fs::File::open(tout.path())
.unwrap()
.read_exact(&mut magic)
.unwrap();
assert_eq!(&magic, b"PMTiles");
}
fn coalescing_generalization(renames: Option<BTreeMap<String, String>>) -> Generalization {
Generalization {
engine: "tylertoo test".to_string(),
gsd_base: None,
cascade: None,
collapse: None,
representation: None,
levels: vec![],
ranking: None,
density_drop: None,
clustering: None,
coalescing: Some(CoalescingProvenance {
enabled: true,
snap_tolerance_gsd_factor: 1.0,
junction_angle: Some(0.0),
max_level_rows: Some(2_000_000),
coalesced_count_column: "coalesced_count".to_string(),
}),
renamed_columns: renames,
}
}
fn clustering_generalization() -> Generalization {
Generalization {
clustering: Some(ClusteringProvenance {
enabled: true,
point_count_column: "point_count".to_string(),
accumulated: vec![],
}),
coalescing: None,
..coalescing_generalization(None)
}
}
fn write_counter_fixture(
path: &Path,
generalization: Generalization,
counter: &str,
counts: &[Vec<i32>],
user_column: Option<&str>,
) {
use arrow_array::Int32Array;
let a = Geometry::Point(Point::new(-120.0, 40.0));
let b = Geometry::Point(Point::new(120.0, -40.0));
let mut fields = vec![Field::new("id", DataType::Int64, false)];
if let Some(name) = user_column {
fields.push(Field::new(name, DataType::Int32, false));
}
fields.push(geometry_field());
fields.push(Field::new(counter, DataType::Int32, false));
let schema = Arc::new(Schema::new(fields));
let specs: Vec<LevelSpec> = (0..counts.len())
.map(|k| LevelSpec::new(gsd((2 + 2 * k) as u8), Some((2 + 2 * k) as u8)))
.collect();
let mut opts = OverviewWriterOptions::new(Mode::Duplicating, specs);
opts.generalization = Some(generalization);
let mut writer = OverviewWriter::create(path, &schema, opts).unwrap();
for (k, level_counts) in counts.iter().enumerate() {
let n = level_counts.len();
let geoms: Vec<Geometry<f64>> = (0..n)
.map(|i| if i % 2 == 0 { a.clone() } else { b.clone() })
.collect();
let mut columns: Vec<ArrayRef> = vec![Arc::new(Int64Array::from(
(0..n as i64).collect::<Vec<_>>(),
))];
if user_column.is_some() {
columns.push(Arc::new(Int32Array::from(vec![7; n])));
}
columns.push(Arc::new(build_geometry_array(&geoms).to_array_ref()));
columns.push(Arc::new(Int32Array::from(level_counts.clone())));
let batch = RecordBatch::try_new(schema.clone(), columns).unwrap();
assert_eq!(
writer
.write_level(k, Some(n), std::iter::once(batch))
.unwrap(),
LevelWriteOutcome::Written
);
}
writer.finish().unwrap();
}
fn write_coalesced_fixture(path: &Path, counts: &[Vec<i32>]) {
write_counter_fixture(
path,
coalescing_generalization(None),
"coalesced_count",
counts,
None,
);
}
fn exported_property_names(input: &Path) -> Vec<String> {
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
let tout = tempfile::NamedTempFile::new().unwrap();
export_pmtiles(input, tout.path(), &ExportOptions::default()).unwrap();
let decoded = tempfile::NamedTempFile::new().unwrap();
crate::decode::decode_pmtiles(
tout.path(),
decoded.path(),
&crate::decode::DecodeOptions::default(),
)
.unwrap();
let file = std::fs::File::open(decoded.path()).unwrap();
ParquetRecordBatchReaderBuilder::try_new(file)
.unwrap()
.schema()
.fields()
.iter()
.map(|f| f.name().clone())
.filter(|n| {
!matches!(
n.as_str(),
"zoom" | "layer" | "mvt_id" | "geometry" | "bbox"
)
})
.collect()
}
#[test]
fn coalesced_count_is_not_exported_when_it_is_1_everywhere() {
let tin = tempfile::NamedTempFile::new().unwrap();
write_coalesced_fixture(tin.path(), &[vec![1, 1], vec![1, 1]]);
let names = exported_property_names(tin.path());
assert_eq!(names, vec!["id"], "coalesced_count leaked into the tiles");
let reader = OverviewReader::open(tin.path()).unwrap();
let published = PublishedNames::from_reader(&reader);
let fields = field_metadata(reader.schema(), geometry_index(reader.schema()), &published);
assert!(
!fields.contains_key("coalesced_count"),
"vector_layers fields must not advertise it: {fields:?}"
);
}
#[test]
fn coalesced_count_explicit_include_overrides_suppression() {
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
let tin = tempfile::NamedTempFile::new().unwrap();
write_coalesced_fixture(tin.path(), &[vec![1, 1], vec![1, 1]]);
let opts = ExportOptions {
properties: PropertySelection {
include: Some(vec!["coalesced_count".to_string()]),
..Default::default()
},
..ExportOptions::default()
};
let tout = tempfile::NamedTempFile::new().unwrap();
export_pmtiles(tin.path(), tout.path(), &opts).expect("explicit include exports");
let decoded = tempfile::NamedTempFile::new().unwrap();
crate::decode::decode_pmtiles(
tout.path(),
decoded.path(),
&crate::decode::DecodeOptions::default(),
)
.unwrap();
let file = std::fs::File::open(decoded.path()).unwrap();
let names: Vec<String> = ParquetRecordBatchReaderBuilder::try_new(file)
.unwrap()
.schema()
.fields()
.iter()
.map(|f| f.name().clone())
.filter(|n| {
!matches!(
n.as_str(),
"zoom" | "layer" | "mvt_id" | "geometry" | "bbox"
)
})
.collect();
assert_eq!(names, vec!["coalesced_count"]);
let reader = OverviewReader::open(tin.path()).unwrap();
let geom_idx = geometry_index(reader.schema()).unwrap();
let published = PublishedNames::from_reader(&reader)
.with_selection(
reader.schema(),
geom_idx,
&opts.properties,
&opts.feature_order,
)
.unwrap();
let fields = field_metadata(reader.schema(), Some(geom_idx), &published);
assert_eq!(
fields.keys().collect::<Vec<_>>(),
vec!["coalesced_count"],
"{fields:?}"
);
}
#[test]
fn coalesced_count_is_exported_when_coalescing_merged_something() {
let tin = tempfile::NamedTempFile::new().unwrap();
write_coalesced_fixture(tin.path(), &[vec![2, 1], vec![1, 1]]);
let names = exported_property_names(tin.path());
assert_eq!(names, vec!["coalesced_count", "id"]);
}
#[test]
fn suppressed_counter_frees_its_name_for_a_renamed_source_column() {
let tin = tempfile::NamedTempFile::new().unwrap();
write_counter_fixture(
tin.path(),
coalescing_generalization(Some(BTreeMap::from([(
"coalesced_count_".to_string(),
"coalesced_count".to_string(),
)]))),
"coalesced_count",
&[vec![1, 1], vec![1, 1]],
Some("coalesced_count_"),
);
let reader = OverviewReader::open(tin.path()).unwrap();
let published = PublishedNames::from_reader(&reader);
let schema = reader.schema();
let names: Vec<String> =
property_columns(schema, geometry_index(schema).unwrap(), &published)
.into_iter()
.map(|(_, n)| n)
.collect();
assert_eq!(
names,
vec!["id", "coalesced_count"],
"the withheld counter no longer occupies its name"
);
let fields = field_metadata(schema, geometry_index(schema), &published);
assert_eq!(
fields.get("coalesced_count").map(String::as_str),
Some("Number")
);
assert!(!fields.contains_key("coalesced_count_"), "{fields:?}");
let mut exported = exported_property_names(tin.path());
exported.sort();
assert_eq!(exported, vec!["coalesced_count", "id"]);
}
#[test]
fn point_count_is_exported_even_when_it_is_1_everywhere() {
let tin = tempfile::NamedTempFile::new().unwrap();
write_counter_fixture(
tin.path(),
clustering_generalization(),
"point_count",
&[vec![1, 1], vec![1, 1]],
None,
);
let names = exported_property_names(tin.path());
assert_eq!(names, vec!["id", "point_count"], "point_count must stay");
let reader = OverviewReader::open(tin.path()).unwrap();
let published = PublishedNames::from_reader(&reader);
let fields = field_metadata(reader.schema(), geometry_index(reader.schema()), &published);
assert_eq!(
fields.get("point_count").map(String::as_str),
Some("Number")
);
}
#[test]
fn footer_cannot_redirect_suppression_onto_a_data_column() {
let tin = tempfile::NamedTempFile::new().unwrap();
let mut generalization = coalescing_generalization(None);
generalization
.coalescing
.as_mut()
.unwrap()
.coalesced_count_column = "id".to_string();
write_counter_fixture(
tin.path(),
generalization,
"coalesced_count",
&[vec![1, 1], vec![1, 1]],
None,
);
let names = exported_property_names(tin.path());
assert!(
names.contains(&"id".to_string()),
"id was withheld: {names:?}"
);
}
#[test]
fn export_declared_min_zoom_is_written_to_header_and_report() {
let a = Geometry::Point(Point::new(-120.0, 40.0));
let b = Geometry::Point(Point::new(120.0, -40.0));
let tin = tempfile::NamedTempFile::new().unwrap();
write_fixture(
tin.path(),
&[
(vec![0, 1], vec![a.clone(), b.clone()]),
(vec![0, 1], vec![a.clone(), b.clone()]),
],
);
let tout = tempfile::NamedTempFile::new().unwrap();
let opts = ExportOptions {
min_zoom: Some(0),
..ExportOptions::default()
};
let report = export_pmtiles(tin.path(), tout.path(), &opts).unwrap();
assert_eq!(report.min_zoom, 0, "report must cover the declared range");
assert_eq!(report.max_zoom, 4);
assert_eq!(
report.zooms.len(),
2,
"no phantom per-zoom rows for empty zooms"
);
let data = std::fs::read(tout.path()).unwrap();
let header = crate::pmtiles_writer::Header::from_bytes(&data[..127]).unwrap();
assert_eq!(header.min_zoom, 0);
assert_eq!(header.max_zoom, 4);
}
#[test]
fn export_declared_min_zoom_finer_than_coarsest_level_is_rejected() {
let a = Geometry::Point(Point::new(-120.0, 40.0));
let tin = tempfile::NamedTempFile::new().unwrap();
write_fixture(
tin.path(),
&[(vec![0], vec![a.clone()]), (vec![0], vec![a.clone()])],
);
let tout = tempfile::NamedTempFile::new().unwrap();
let opts = ExportOptions {
min_zoom: Some(3),
..ExportOptions::default()
};
let err = export_pmtiles(tin.path(), tout.path(), &opts).unwrap_err();
let msg = err.to_string();
assert!(
msg.contains("min_zoom") && msg.contains('3') && msg.contains('2'),
"error must name the declared and actual minimum: {msg}"
);
}
fn exported_names_with(selection: PropertySelection) -> Result<Vec<String>, ExportError> {
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
let a = Geometry::Point(Point::new(-120.0, 40.0));
let b = Geometry::Point(Point::new(120.0, -40.0));
let tin = tempfile::NamedTempFile::new().unwrap();
write_fixture(
tin.path(),
&[
(vec![0, 1], vec![a.clone(), b.clone()]),
(vec![0, 1], vec![a.clone(), b.clone()]),
],
);
let tout = tempfile::NamedTempFile::new().unwrap();
let opts = ExportOptions {
properties: selection,
..ExportOptions::default()
};
export_pmtiles(tin.path(), tout.path(), &opts)?;
let decoded = tempfile::NamedTempFile::new().unwrap();
crate::decode::decode_pmtiles(
tout.path(),
decoded.path(),
&crate::decode::DecodeOptions::default(),
)
.unwrap();
let file = std::fs::File::open(decoded.path()).unwrap();
Ok(ParquetRecordBatchReaderBuilder::try_new(file)
.unwrap()
.schema()
.fields()
.iter()
.map(|f| f.name().clone())
.filter(|n| {
!matches!(
n.as_str(),
"zoom" | "layer" | "mvt_id" | "geometry" | "bbox"
)
})
.collect())
}
#[test]
fn export_property_selection_include_exclude_and_none() {
let include = PropertySelection {
include: Some(vec!["name".to_string()]),
..Default::default()
};
assert_eq!(exported_names_with(include).unwrap(), vec!["name"]);
let exclude = PropertySelection {
exclude: vec!["name".to_string()],
..Default::default()
};
assert_eq!(exported_names_with(exclude).unwrap(), vec!["id"]);
let none = PropertySelection {
exclude_all: true,
..Default::default()
};
assert!(exported_names_with(none).unwrap().is_empty());
}
#[test]
fn export_property_selection_rejects_unknown_include() {
let typo = PropertySelection {
include: Some(vec!["nmae".to_string()]),
..Default::default()
};
let msg = exported_names_with(typo).unwrap_err().to_string();
assert!(
msg.contains("\"nmae\"") && msg.contains("\"name\""),
"{msg}"
);
}
#[test]
fn export_property_selection_cannot_drop_the_feature_order_column() {
let a = Geometry::Point(Point::new(-120.0, 40.0));
let tin = tempfile::NamedTempFile::new().unwrap();
write_fixture(
tin.path(),
&[(vec![0], vec![a.clone()]), (vec![0], vec![a])],
);
let tout = tempfile::NamedTempFile::new().unwrap();
let opts = ExportOptions {
feature_order: FeatureOrder::Column {
name: "id".to_string(),
descending: false,
},
properties: PropertySelection {
exclude: vec!["id".to_string()],
..Default::default()
},
..ExportOptions::default()
};
let msg = export_pmtiles(tin.path(), tout.path(), &opts)
.unwrap_err()
.to_string();
assert!(
msg.contains("--feature-order") && msg.contains("\"id\""),
"{msg}"
);
}
}