Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_geom/src/tile/pmtiles.rs

19.7 KiB, 1 run

created by r1870400018:60373, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! PMTiles version 3: a whole tile pyramid in one file, read by byte range.
2//!
3//! An archive is a 127-byte header, a root directory, optional JSON metadata, leaf directories
4//! and the tile data, and every tile is found by at most four ranged reads with no index
5//! server: the header and root, up to three levels of leaf directory, the tile. That is what
6//! lets a street map be served from a plain file by any web server that answers `Range`, or
7//! read straight from disk. The format is Protomaps' specification at
8//! <https://github.com/protomaps/PMTiles/blob/main/spec/v3/spec.md>.
9//!
10//! The reading is split so that nothing here does input or output of its own. [`Header`],
11//! [`decode_directory`] and [`next`] are pure: a caller fetching ranges asynchronously -- over
12//! HTTP, from a browser, from object storage -- drives them itself, one read per step. An
13//! [`Archive`] drives them over any synchronous [`RangeSource`], of which a local file
14//! ([`FileSource`]) and a buffer in memory are provided here; an HTTP one lives with the HTTP
15//! client, in `fe2o3_net`, so that this crate carries no network code.
16//!
17//! Tile ids number the tiles of all zooms along one Hilbert curve per zoom, so that tiles near
18//! each other on the map are near each other in the file and runs of identical tiles -- open
19//! ocean -- collapse into one directory entry.
20
21use oxedyne_fe2o3_core::prelude::*;
22
23use std::io::Read;
24#[cfg(not(unix))]
25use std::io::{
26 Seek,
27 SeekFrom,
28};
29
30pub const HEADER_LEN: usize = 127;
31pub const MAGIC: &[u8; 7] = b"PMTiles";
32pub const VERSION: u8 = 3;
33pub const MAX_DEPTH: usize = 4; // the root and three levels of leaves
34pub const MAX_ZOOM: u8 = 31;
35
36/// How the directories or the tiles of an archive are compressed.
37#[derive(Clone, Copy, Debug, PartialEq, Eq)]
38pub enum Compression {
39 Unknown,
40 None,
41 Gzip,
42 Brotli,
43 Zstd,
44}
45
46impl Compression {
47 fn from_code(c: u8) -> Outcome<Self> {
48 match c {
49 0 => Ok(Self::Unknown),
50 1 => Ok(Self::None),
51 2 => Ok(Self::Gzip),
52 3 => Ok(Self::Brotli),
53 4 => Ok(Self::Zstd),
54 _ => Err(err!("Compression {} is none the specification names.", c; Invalid, Input)),
55 }
56 }
57}
58
59/// What the tiles of an archive are.
60#[derive(Clone, Copy, Debug, PartialEq, Eq)]
61pub enum TileType {
62 Unknown,
63 Mvt,
64 Png,
65 Jpeg,
66 Webp,
67 Avif,
68 Mlt,
69}
70
71impl TileType {
72 fn from_code(c: u8) -> Outcome<Self> {
73 match c {
74 0 => Ok(Self::Unknown),
75 1 => Ok(Self::Mvt),
76 2 => Ok(Self::Png),
77 3 => Ok(Self::Jpeg),
78 4 => Ok(Self::Webp),
79 5 => Ok(Self::Avif),
80 6 => Ok(Self::Mlt),
81 _ => Err(err!("Tile type {} is none the specification names.", c; Invalid, Input)),
82 }
83 }
84}
85
86/// The fixed 127 bytes an archive begins with.
87#[derive(Clone, Debug, PartialEq)]
88pub struct Header {
89 pub root_offset: u64,
90 pub root_length: u64,
91 pub metadata_offset: u64,
92 pub metadata_length: u64,
93 pub leaf_offset: u64,
94 pub leaf_length: u64,
95 pub data_offset: u64,
96 pub data_length: u64,
97 pub addressed_tiles: u64,
98 pub tile_entries: u64,
99 pub tile_contents: u64,
100 pub clustered: bool,
101 pub internal_compression: Compression,
102 pub tile_compression: Compression,
103 pub tile_type: TileType,
104 pub min_zoom: u8,
105 pub max_zoom: u8,
106 pub min_lon_e7: i32, // degrees times ten million, as the file holds them
107 pub min_lat_e7: i32,
108 pub max_lon_e7: i32,
109 pub max_lat_e7: i32,
110 pub centre_zoom: u8,
111 pub centre_lon_e7: i32,
112 pub centre_lat_e7: i32,
113}
114
115fn u64_at(b: &[u8], at: usize) -> u64 {
116 let mut a = [0u8; 8];
117 a.copy_from_slice(&b[at..at + 8]);
118 u64::from_le_bytes(a)
119}
120
121fn e7_at(b: &[u8], at: usize) -> i32 {
122 i32::from_le_bytes([b[at], b[at + 1], b[at + 2], b[at + 3]])
123}
124
125impl Header {
126 /// Reads the header from the first 127 bytes (or more) of an archive.
127 pub fn parse(b: &[u8]) -> Outcome<Self> {
128 if b.len() < HEADER_LEN {
129 return Err(err!("A PMTiles header is {} bytes and {} were given.", HEADER_LEN, b.len();
130 Invalid, Input, TooSmall));
131 }
132 if &b[..7] != MAGIC {
133 return Err(err!("The archive does not begin with the PMTiles magic."; Invalid, Input));
134 }
135 if b[7] != VERSION {
136 return Err(err!("The archive is PMTiles version {}, and this reads version {}.",
137 b[7], VERSION; Invalid, Input, Version));
138 }
139 let h = Self {
140 root_offset: u64_at(b, 8),
141 root_length: u64_at(b, 16),
142 metadata_offset: u64_at(b, 24),
143 metadata_length: u64_at(b, 32),
144 leaf_offset: u64_at(b, 40),
145 leaf_length: u64_at(b, 48),
146 data_offset: u64_at(b, 56),
147 data_length: u64_at(b, 64),
148 addressed_tiles: u64_at(b, 72),
149 tile_entries: u64_at(b, 80),
150 tile_contents: u64_at(b, 88),
151 clustered: b[96] == 1,
152 internal_compression: res!(Compression::from_code(b[97])),
153 tile_compression: res!(Compression::from_code(b[98])),
154 tile_type: res!(TileType::from_code(b[99])),
155 min_zoom: b[100],
156 max_zoom: b[101],
157 min_lon_e7: e7_at(b, 102),
158 min_lat_e7: e7_at(b, 106),
159 max_lon_e7: e7_at(b, 110),
160 max_lat_e7: e7_at(b, 114),
161 centre_zoom: b[118],
162 centre_lon_e7: e7_at(b, 119),
163 centre_lat_e7: e7_at(b, 123),
164 };
165 if h.min_zoom > h.max_zoom || h.max_zoom > MAX_ZOOM {
166 return Err(err!("The archive claims zooms {} to {}.", h.min_zoom, h.max_zoom;
167 Invalid, Input, Range));
168 }
169 if h.root_length > 16_384 * 1024 {
170 return Err(err!("The archive claims a root directory of {} bytes.", h.root_length;
171 Invalid, Input, Excessive));
172 }
173 Ok(h)
174 }
175
176 /// The bounds as `[min lon, min lat, max lon, max lat]` in degrees times ten million.
177 pub fn bounds_e7(&self) -> [i32; 4] {
178 [self.min_lon_e7, self.min_lat_e7, self.max_lon_e7, self.max_lat_e7]
179 }
180
181 /// The bounds as `[min lon, min lat, max lon, max lat]` in degrees.
182 pub fn bounds(&self) -> [f64; 4] {
183 let b = self.bounds_e7();
184 [b[0] as f64 / 1.0e7, b[1] as f64 / 1.0e7, b[2] as f64 / 1.0e7, b[3] as f64 / 1.0e7]
185 }
186}
187
188/// One directory entry: a run of tiles sharing one stored content, or, with a run length of
189/// nought, a pointer to a leaf directory.
190#[derive(Clone, Copy, Debug, PartialEq, Eq)]
191pub struct Entry {
192 pub tile_id: u64,
193 pub offset: u64, // from the start of the tile data, or of the leaf directories
194 pub length: u32,
195 pub run_length: u32,
196}
197
198/// The id of tile `z/x/y`: the tiles of every shallower zoom, then the tile's position along
199/// the Hilbert curve through its own zoom.
200pub fn zxy_to_id(z: u8, x: u32, y: u32) -> Outcome<u64> {
201 if z > MAX_ZOOM {
202 return Err(err!("Zoom {} is past the deepest, {}.", z, MAX_ZOOM; Invalid, Input, Range));
203 }
204 let n = 1u64 << z;
205 let (mut x, mut y) = (x as u64, y as u64);
206 if x >= n || y >= n {
207 return Err(err!("Tile {}/{}/{} is off a grid {} tiles across.", z, x, y, n; Invalid, Input, Range));
208 }
209 let base = ((1u128 << (2 * z as u32)) - 1) / 3;
210 let mut d: u64 = 0;
211 let mut s = n / 2;
212 while s > 0 {
213 let rx = (x & s) != 0;
214 let ry = (y & s) != 0;
215 d += s * s * ((if rx { 3 } else { 0 }) ^ (if ry { 1 } else { 0 }));
216 // Turn the quadrant so the curve runs on through the next level.
217 if !ry {
218 if rx {
219 x = n - 1 - x;
220 y = n - 1 - y;
221 }
222 std::mem::swap(&mut x, &mut y);
223 }
224 s /= 2;
225 }
226 Ok(base as u64 + d)
227}
228
229/// The tile a tile id names.
230pub fn id_to_zxy(id: u64) -> Outcome<(u8, u32, u32)> {
231 let mut base: u128 = 0;
232 for z in 0..=MAX_ZOOM {
233 let count = 1u128 << (2 * z as u32);
234 if (id as u128) < base + count {
235 let n = 1u64 << z;
236 let mut t = (id as u128 - base) as u64;
237 let (mut x, mut y) = (0u64, 0u64);
238 let mut s = 1u64;
239 while s < n {
240 let rx = 1 & (t / 2);
241 let ry = 1 & (t ^ rx);
242 if ry == 0 {
243 if rx == 1 {
244 x = s - 1 - x;
245 y = s - 1 - y;
246 }
247 std::mem::swap(&mut x, &mut y);
248 }
249 x += s * rx;
250 y += s * ry;
251 t /= 4;
252 s *= 2;
253 }
254 return Ok((z, x as u32, y as u32));
255 }
256 base += count;
257 }
258 Err(err!("Tile id {} is past zoom {}.", id, MAX_ZOOM; Invalid, Input, Range))
259}
260
261fn varint(b: &[u8], pos: &mut usize) -> Outcome<u64> {
262 let mut v: u64 = 0;
263 let mut shift = 0u32;
264 loop {
265 let byte = match b.get(*pos) {
266 Some(x) => *x,
267 None => return Err(err!("A directory ends inside a varint at byte {}.", *pos;
268 Invalid, Input, Decode)),
269 };
270 *pos += 1;
271 if shift >= 64 {
272 return Err(err!("A directory varint at byte {} is longer than ten bytes.", *pos;
273 Invalid, Input, Decode));
274 }
275 v |= ((byte & 0x7f) as u64) << shift;
276 if byte < 0x80 {
277 return Ok(v);
278 }
279 shift += 7;
280 }
281}
282
283/// Decodes a directory from its bytes, already decompressed: an entry count, then the tile
284/// ids as deltas, the run lengths, the lengths and the offsets, each column in turn, all
285/// varints. An offset of nought after the first entry means "straight after the last".
286pub fn decode_directory(b: &[u8]) -> Outcome<Vec<Entry>> {
287 let mut pos = 0usize;
288 let n = res!(varint(b, &mut pos)) as usize;
289 // Every entry costs at least four bytes, so a count past that is a lie and sizes nothing.
290 if n > b.len() / 4 + 1 {
291 return Err(err!("A directory of {} bytes claims {} entries.", b.len(), n; Invalid, Input, Decode));
292 }
293 let mut out = vec![Entry { tile_id: 0, offset: 0, length: 0, run_length: 0 }; n];
294 let mut id: u64 = 0;
295 for e in out.iter_mut() {
296 id = match id.checked_add(res!(varint(b, &mut pos))) {
297 Some(v) => v,
298 None => return Err(err!("A directory's tile ids overflow."; Invalid, Input, Decode, Overflow)),
299 };
300 e.tile_id = id;
301 }
302 for e in out.iter_mut() {
303 e.run_length = res!(u32_of(res!(varint(b, &mut pos)), "A run length"));
304 }
305 for e in out.iter_mut() {
306 e.length = res!(u32_of(res!(varint(b, &mut pos)), "A length"));
307 }
308 for i in 0..n {
309 let v = res!(varint(b, &mut pos));
310 out[i].offset = if v == 0 && i > 0 {
311 out[i - 1].offset + out[i - 1].length as u64
312 } else if v == 0 {
313 return Err(err!("A directory's first offset is written as 'after the last'.";
314 Invalid, Input, Decode));
315 } else {
316 v - 1
317 };
318 }
319 if pos != b.len() {
320 return Err(err!("A directory has {} bytes after its last entry.", b.len() - pos;
321 Invalid, Input, Decode));
322 }
323 Ok(out)
324}
325
326fn u32_of(v: u64, what: &str) -> Outcome<u32> {
327 if v > u32::MAX as u64 {
328 return Err(err!("{} of {} is more than 32 bits.", what, v; Invalid, Input, Decode));
329 }
330 Ok(v as u32)
331}
332
333/// The entry covering a tile id: the last entry at or before it, if that is a leaf pointer or
334/// a run that reaches it.
335pub fn find(entries: &[Entry], id: u64) -> Option<&Entry> {
336 let at = entries.partition_point(|e| e.tile_id <= id);
337 if at == 0 {
338 return None;
339 }
340 let e = &entries[at - 1];
341 if e.run_length == 0 || id - e.tile_id < e.run_length as u64 {
342 Some(e)
343 } else {
344 None
345 }
346}
347
348/// What to read next to find a tile, as absolute byte ranges of the archive.
349#[derive(Clone, Copy, Debug, PartialEq, Eq)]
350pub enum Next {
351 Tile { offset: u64, length: u64 }, // the tile, as stored
352 Leaf { offset: u64, length: u64 }, // a leaf directory to decode and search in turn
353 Missing, // no such tile
354}
355
356/// One step of finding a tile in a directory already read.
357pub fn next(header: &Header, dir: &[Entry], id: u64) -> Next {
358 match find(dir, id) {
359 None => Next::Missing,
360 Some(e) if e.run_length == 0 => Next::Leaf {
361 offset: header.leaf_offset + e.offset, length: e.length as u64 },
362 Some(e) => Next::Tile { offset: header.data_offset + e.offset, length: e.length as u64 },
363 }
364}
365
366/// Undoes an archive's compression of its directories or tiles.
367pub fn decompress(bytes: &[u8], c: Compression) -> Outcome<Vec<u8>> {
368 match c {
369 Compression::None | Compression::Unknown => Ok(bytes.to_vec()),
370 Compression::Gzip => {
371 let mut out = Vec::new();
372 let mut dec = flate2::read::GzDecoder::new(bytes);
373 res!(dec.read_to_end(&mut out), Decode, IO);
374 Ok(out)
375 },
376 other => Err(err!("{:?} compression is not read here; only gzip and none are.", other;
377 Invalid, Input, Unimplemented)),
378 }
379}
380
381// ---------------------------------------------------------------------------------------------
382// Synchronous reading
383// ---------------------------------------------------------------------------------------------
384
385/// Anything that can hand over a byte range: a file, a buffer, a remote file by HTTP `Range`.
386///
387/// A read takes `&self`, so one source serves many threads at once, as a tile server's
388/// blocking workers do.
389pub trait RangeSource {
390 fn read(&self, offset: u64, len: u64) -> Outcome<Vec<u8>>;
391}
392
393/// A local archive file, read by position so that concurrent reads need no lock.
394pub struct FileSource {
395 #[cfg(unix)]
396 file: std::fs::File,
397 #[cfg(not(unix))]
398 file: std::sync::Mutex<std::fs::File>,
399}
400
401impl FileSource {
402 pub fn open<P: AsRef<std::path::Path>>(path: P) -> Outcome<Self> {
403 let file = res!(std::fs::File::open(path.as_ref()), File, Read);
404 #[cfg(unix)]
405 return Ok(Self { file });
406 #[cfg(not(unix))]
407 return Ok(Self { file: std::sync::Mutex::new(file) });
408 }
409}
410
411impl RangeSource for FileSource {
412 fn read(&self, offset: u64, len: u64) -> Outcome<Vec<u8>> {
413 if len > u32::MAX as u64 {
414 return Err(err!("A read of {} bytes is more than one tile or directory needs.", len;
415 Invalid, Input, Excessive));
416 }
417 let mut buf = vec![0u8; len as usize];
418 #[cfg(unix)]
419 {
420 use std::os::unix::fs::FileExt;
421 res!(self.file.read_exact_at(&mut buf, offset), File, Read);
422 }
423 #[cfg(not(unix))]
424 {
425 let mut f = lock_mutex!(self.file);
426 res!(f.seek(SeekFrom::Start(offset)), File, Seek);
427 res!(f.read_exact(&mut buf), File, Read);
428 }
429 Ok(buf)
430 }
431}
432
433impl RangeSource for &[u8] {
434 fn read(&self, offset: u64, len: u64) -> Outcome<Vec<u8>> {
435 let end = offset.checked_add(len);
436 match end {
437 Some(e) if e <= self.len() as u64 => Ok(self[offset as usize..e as usize].to_vec()),
438 _ => Err(err!("Bytes {}+{} are past the {} there are.", offset, len, self.len();
439 Invalid, Input, Range)),
440 }
441 }
442}
443
444/// An archive over a synchronous source, with its root directory and the leaf directories it
445/// has read lately kept in memory. Reads take `&self`, so an archive shared between threads
446/// serves them all; the leaf cache is the only thing they share, behind a lock held for a
447/// lookup and never across a read.
448pub struct Archive<S: RangeSource> {
449 src: S,
450 header: Header,
451 root: Vec<Entry>,
452 leaves: std::sync::Mutex<Vec<(u64, std::sync::Arc<Vec<Entry>>)>>, // most recent last
453}
454
455const LEAF_CACHE: usize = 32;
456
457impl<S: RangeSource> Archive<S> {
458 pub fn open(src: S) -> Outcome<Self> {
459 let head = res!(src.read(0, HEADER_LEN as u64));
460 let header = res!(Header::parse(&head));
461 let raw = res!(src.read(header.root_offset, header.root_length));
462 let root = res!(decode_directory(&res!(decompress(&raw, header.internal_compression))));
463 Ok(Self { src, header, root, leaves: std::sync::Mutex::new(Vec::new()) })
464 }
465
466 pub fn header(&self) -> &Header { &self.header }
467
468 pub fn root(&self) -> &[Entry] { &self.root }
469
470 /// The archive's JSON metadata, decompressed.
471 pub fn metadata(&self) -> Outcome<String> {
472 let raw = res!(self.src.read(self.header.metadata_offset, self.header.metadata_length));
473 let bytes = res!(decompress(&raw, self.header.internal_compression));
474 match String::from_utf8(bytes) {
475 Ok(s) => Ok(s),
476 Err(_) => Err(err!("The archive's metadata is not UTF-8."; Invalid, Input, UTF8)),
477 }
478 }
479
480 /// Tile `z/x/y` as stored, still compressed as [`Header::tile_compression`] says, or
481 /// `None` where the archive has no such tile. Stored bytes are what a server passes on
482 /// with a `Content-Encoding`; [`Archive::tile_decoded`] undoes it.
483 pub fn tile(&self, z: u8, x: u32, y: u32) -> Outcome<Option<Vec<u8>>> {
484 if z < self.header.min_zoom || z > self.header.max_zoom {
485 return Ok(None);
486 }
487 let id = res!(zxy_to_id(z, x, y));
488 let mut step = next(&self.header, &self.root, id);
489 for _ in 0..MAX_DEPTH {
490 match step {
491 Next::Missing => return Ok(None),
492 Next::Tile { offset, length } => return Ok(Some(res!(self.src.read(offset, length)))),
493 Next::Leaf { offset, length } => {
494 let dir = res!(self.leaf(offset, length));
495 step = next(&self.header, &dir, id);
496 },
497 }
498 }
499 Err(err!("Tile {}/{}/{} is nested deeper than {} directories.", z, x, y, MAX_DEPTH;
500 Invalid, Input, Decode))
501 }
502
503 /// A leaf directory, from the cache or read and decoded. Two threads missing the same
504 /// leaf at once both read it, which costs a read and keeps the lock off the disk.
505 fn leaf(&self, offset: u64, length: u64) -> Outcome<std::sync::Arc<Vec<Entry>>> {
506 {
507 let mut cache = lock_mutex!(self.leaves);
508 if let Some(i) = cache.iter().position(|(o, _)| *o == offset) {
509 let hit = cache.remove(i);
510 let dir = hit.1.clone();
511 cache.push(hit);
512 return Ok(dir);
513 }
514 }
515 let raw = res!(self.src.read(offset, length));
516 let dir = res!(decode_directory(&res!(decompress(&raw, self.header.internal_compression))));
517 if dir.is_empty() {
518 return Err(err!("The leaf directory at byte {} is empty.", offset; Invalid, Input, Decode));
519 }
520 let dir = std::sync::Arc::new(dir);
521 let mut cache = lock_mutex!(self.leaves);
522 cache.push((offset, dir.clone()));
523 if cache.len() > LEAF_CACHE {
524 cache.remove(0);
525 }
526 Ok(dir)
527 }
528
529 /// Tile `z/x/y` with the archive's tile compression undone.
530 pub fn tile_decoded(&self, z: u8, x: u32, y: u32) -> Outcome<Option<Vec<u8>>> {
531 let c = self.header.tile_compression;
532 match res!(self.tile(z, x, y)) {
533 Some(raw) => Ok(Some(res!(decompress(&raw, c)))),
534 None => Ok(None),
535 }
536 }
537}