use std::fs::File;
use std::io::{Read as _, Seek as _, SeekFrom};
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::mpsc::{Sender, channel};
use oxigis_render::ByteRange;
use oxigis_ui::{RangeJob, RangeSink, RangeTransport, TileError};
pub const FILE_RANGE_WORKER_THREADS: usize = 4;
pub const MAX_FILE_RANGE_BYTES: u64 = 32 * 1024 * 1024;
struct Job {
path: String,
range: ByteRange,
job: RangeJob,
sink: RangeSink,
}
pub struct FileRangeTransport {
queues: Vec<Sender<Job>>,
next: AtomicUsize,
}
impl FileRangeTransport {
pub fn new() -> Result<Self, std::io::Error> {
let pins = Arc::new(FileValidatorPins::default());
let mut queues = Vec::with_capacity(FILE_RANGE_WORKER_THREADS);
for index in 0..FILE_RANGE_WORKER_THREADS {
let (tx, rx) = channel::<Job>();
let pins = Arc::clone(&pins);
std::thread::Builder::new()
.name(format!("oxigis-archive-{index}"))
.spawn(move || {
let mut open: Option<(String, File)> = None;
for queued in rx {
let result = read_range(&mut open, &pins, &queued.path, queued.range);
queued.sink.deliver(queued.job, result);
}
})?;
queues.push(tx);
}
Ok(Self {
queues,
next: AtomicUsize::new(0),
})
}
}
impl core::fmt::Debug for FileRangeTransport {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("FileRangeTransport")
.field("workers", &self.queues.len())
.finish()
}
}
impl RangeTransport for FileRangeTransport {
fn request_range(&self, url: String, range: ByteRange, job: RangeJob, sink: RangeSink) {
if self.queues.is_empty() {
sink.deliver(job, Err(TileError::permanent("no archive worker threads")));
return;
}
let index = self.next.fetch_add(1, Ordering::Relaxed) % self.queues.len();
let Some(queue) = self.queues.get(index) else {
sink.deliver(
job,
Err(TileError::permanent("archive worker queue vanished")),
);
return;
};
let queued = Job {
path: url,
range,
job,
sink: sink.clone(),
};
if let Err(error) = queue.send(queued) {
sink.deliver(
job,
Err(TileError::permanent(format!(
"archive worker is gone: {error}"
))),
);
}
}
}
const DRIFT_ADVICE: &str = "the archive changed on disk; remove and re-add the layer";
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
struct FileValidator {
len: u64,
modified: Option<std::time::SystemTime>,
platform_id: Option<u64>,
}
impl FileValidator {
fn observe(metadata: &std::fs::Metadata) -> Self {
Self {
len: metadata.len(),
modified: metadata.modified().ok(),
platform_id: platform_id_of(metadata),
}
}
}
#[cfg(unix)]
fn platform_id_of(metadata: &std::fs::Metadata) -> Option<u64> {
use std::os::unix::fs::MetadataExt as _;
Some(metadata.ino())
}
#[cfg(not(unix))]
fn platform_id_of(_metadata: &std::fs::Metadata) -> Option<u64> {
None
}
#[derive(Debug, Default)]
struct FileValidatorPins {
entries: std::sync::Mutex<Vec<(String, FileValidator)>>,
}
impl FileValidatorPins {
fn check(&self, path: &str, observed: FileValidator) -> Result<(), TileError> {
let mut entries = self.lock();
if let Some((_, pinned)) = entries.iter_mut().find(|(held, _)| held == path) {
if pinned.len != observed.len {
return Err(TileError::permanent(format!(
"{path}: {DRIFT_ADVICE} (its length went from {} to {} bytes)",
pinned.len, observed.len
)));
}
if let (Some(was), Some(now)) = (pinned.modified, observed.modified)
&& was != now
{
return Err(TileError::permanent(format!(
"{path}: {DRIFT_ADVICE} (its modified time changed)"
)));
}
if let (Some(was), Some(now)) = (pinned.platform_id, observed.platform_id)
&& was != now
{
return Err(TileError::permanent(format!(
"{path}: {DRIFT_ADVICE} (it is no longer the same file on disk)"
)));
}
if pinned.modified.is_none() {
pinned.modified = observed.modified;
}
if pinned.platform_id.is_none() {
pinned.platform_id = observed.platform_id;
}
return Ok(());
}
while entries.len() >= crate::range_http::MAX_PINNED_VALIDATORS && !entries.is_empty() {
entries.remove(0);
}
entries.push((path.to_owned(), observed));
Ok(())
}
fn lock(&self) -> std::sync::MutexGuard<'_, Vec<(String, FileValidator)>> {
self.entries
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
}
fn read_range(
open: &mut Option<(String, File)>,
pins: &FileValidatorPins,
path: &str,
range: ByteRange,
) -> Result<Vec<u8>, TileError> {
let matches = open.as_ref().is_some_and(|(held, _)| held.as_str() == path);
if !matches {
match File::open(path) {
Ok(file) => *open = Some((path.to_owned(), file)),
Err(error) => {
*open = None;
return Err(classify_io(&error, path));
}
}
}
let Some((_, file)) = open.as_mut() else {
return Err(TileError::transient(format!("{path} could not be opened")));
};
let metadata = match file.metadata() {
Ok(metadata) => metadata,
Err(error) => {
*open = None;
return Err(classify_io(&error, path));
}
};
if let Err(error) = pins.check(path, FileValidator::observe(&metadata)) {
*open = None;
return Err(error);
}
read_from(file, path, range).inspect_err(|error| {
if !error.retryable() {
*open = None;
}
})
}
fn read_from(file: &mut File, path: &str, range: ByteRange) -> Result<Vec<u8>, TileError> {
let length = range.end.saturating_sub(range.start);
if length > MAX_FILE_RANGE_BYTES {
return Err(TileError::permanent(format!(
"a {length}-byte read of {path} exceeds the {MAX_FILE_RANGE_BYTES}-byte cap"
)));
}
let Ok(capacity) = usize::try_from(length) else {
return Err(TileError::permanent(format!(
"a {length}-byte read of {path} does not fit this machine's address space"
)));
};
if let Err(error) = file.seek(SeekFrom::Start(range.start)) {
return Err(classify_io(&error, path));
}
let mut buffer = vec![0u8; capacity];
let mut filled = 0usize;
while filled < capacity {
let Some(rest) = buffer.get_mut(filled..) else {
break;
};
match file.read(rest) {
Ok(0) => break,
Ok(read) => filled = filled.saturating_add(read),
Err(ref error) if error.kind() == std::io::ErrorKind::Interrupted => {}
Err(error) => return Err(classify_io(&error, path)),
}
}
if filled == 0 {
return Err(TileError::permanent(format!(
"byte {} is past the end of {path}",
range.start
)));
}
buffer.truncate(filled);
Ok(buffer)
}
fn classify_io(error: &std::io::Error, path: &str) -> TileError {
let message = format!("{path}: {error}");
match error.kind() {
std::io::ErrorKind::NotFound | std::io::ErrorKind::PermissionDenied => {
TileError::permanent(message)
}
_ => TileError::transient(message),
}
}
#[cfg(test)]
mod tests {
use super::{
DRIFT_ADVICE, FILE_RANGE_WORKER_THREADS, FileRangeTransport, FileValidator,
FileValidatorPins, MAX_FILE_RANGE_BYTES, read_range,
};
use oxigis_render::ByteRange;
fn temp_archive(name: &str, bytes: &[u8]) -> std::path::PathBuf {
let stamp = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |elapsed| elapsed.as_nanos());
let path = std::env::temp_dir().join(format!("oxigis-{name}-{stamp}.bin"));
std::fs::write(&path, bytes).expect("the fixture must be writable");
path
}
fn range(start: u64, end: u64) -> ByteRange {
ByteRange::new(start, end).expect("a non-empty range")
}
#[test]
fn the_transport_starts_a_worker_pool() {
let transport = FileRangeTransport::new().expect("worker threads must start");
assert!(format!("{transport:?}").contains(&FILE_RANGE_WORKER_THREADS.to_string()));
}
#[test]
fn a_range_inside_the_file_reads_exactly_those_bytes() {
let path = temp_archive("inside", &[0, 1, 2, 3, 4, 5, 6, 7]);
let mut open = None;
let pins = FileValidatorPins::default();
let bytes = read_range(&mut open, &pins, &path.display().to_string(), range(2, 5))
.expect("an in-bounds read");
assert_eq!(bytes, vec![2, 3, 4]);
assert!(open.is_some());
let _removed = std::fs::remove_file(&path);
}
#[test]
fn a_range_running_past_the_end_comes_back_short_rather_than_failing() {
let path = temp_archive("short", &[9, 8, 7]);
let mut open = None;
let pins = FileValidatorPins::default();
let bytes = read_range(
&mut open,
&pins,
&path.display().to_string(),
range(1, 16_384),
)
.expect("a short read at EOF is legitimate");
assert_eq!(bytes, vec![8, 7]);
let _removed = std::fs::remove_file(&path);
}
#[test]
fn a_range_wider_than_the_cap_is_refused_rather_than_truncated() {
let path = temp_archive("oversize", &[0u8; 8]);
let mut open = None;
let pins = FileValidatorPins::default();
let error = read_range(
&mut open,
&pins,
&path.display().to_string(),
range(0, MAX_FILE_RANGE_BYTES + 1),
)
.expect_err("a range past the cap must be refused, not silently shrunk");
assert!(!error.retryable(), "{error}");
assert!(error.message().contains("exceeds"), "{error}");
let _removed = std::fs::remove_file(&path);
}
#[test]
fn a_range_starting_past_the_end_is_a_permanent_failure() {
let path = temp_archive("past", &[1, 2, 3]);
let mut open = None;
let pins = FileValidatorPins::default();
let error = read_range(
&mut open,
&pins,
&path.display().to_string(),
range(99, 128),
)
.expect_err("a start past the end is not a short read");
assert!(!error.retryable());
assert!(error.message().contains("past the end"), "{error}");
let _removed = std::fs::remove_file(&path);
}
#[test]
fn a_missing_file_is_permanent_not_retried_for_ever() {
let mut open = None;
let pins = FileValidatorPins::default();
let error = read_range(
&mut open,
&pins,
"I:\\this\\path\\does\\not\\exist\\at\\all.pmtiles",
range(0, 16),
)
.expect_err("a missing file must fail");
assert!(!error.retryable(), "{error}");
assert!(open.is_none());
}
#[test]
fn a_file_rewritten_in_place_is_refused_on_the_next_read() {
let path = temp_archive("rewritten", &[0u8; 8]);
let path_str = path.display().to_string();
let mut open = None;
let pins = FileValidatorPins::default();
let _first = read_range(&mut open, &pins, &path_str, range(0, 4)).expect("first read");
std::fs::write(&path, [1u8; 16]).expect("the fixture must be rewritable");
let error = read_range(&mut open, &pins, &path_str, range(0, 4))
.expect_err("a length change on an already-pinned path must be refused");
assert!(!error.retryable(), "{error}");
assert!(error.message().contains(DRIFT_ADVICE), "{error}");
let again = read_range(&mut open, &pins, &path_str, range(0, 4))
.expect_err("the refusal must persist for the life of this transport");
assert!(again.message().contains(DRIFT_ADVICE), "{again}");
let _removed = std::fs::remove_file(&path);
}
#[test]
fn an_unchanged_file_reads_repeatedly_without_tripping_drift() {
let path = temp_archive("stable", &[7u8; 32]);
let path_str = path.display().to_string();
let mut open = None;
let pins = FileValidatorPins::default();
for attempt in 0..5 {
read_range(&mut open, &pins, &path_str, range(0, 4))
.unwrap_or_else(|error| panic!("read {attempt} of an unchanged file: {error}"));
}
let _removed = std::fs::remove_file(&path);
}
#[test]
fn a_changed_platform_id_at_the_same_length_is_drift() {
let pins = FileValidatorPins::default();
let first = FileValidator {
len: 100,
modified: None,
platform_id: Some(11),
};
assert_eq!(pins.check("p", first), Ok(()));
let second = FileValidator {
platform_id: Some(22),
..first
};
let error = pins
.check("p", second)
.expect_err("a changed platform id is drift");
assert!(!error.retryable(), "{error}");
assert!(error.message().contains(DRIFT_ADVICE), "{error}");
}
#[test]
fn a_field_unknown_on_one_side_does_not_manufacture_drift() {
let pins = FileValidatorPins::default();
let unknown = FileValidator {
len: 10,
modified: None,
platform_id: None,
};
assert_eq!(pins.check("p", unknown), Ok(()));
let now_known = FileValidator {
modified: Some(std::time::SystemTime::now()),
..unknown
};
assert_eq!(pins.check("p", now_known), Ok(()));
let different_time = FileValidator {
modified: Some(std::time::SystemTime::now() + std::time::Duration::from_secs(3600)),
..now_known
};
let error = pins
.check("p", different_time)
.expect_err("the modified time pinned by the second call must now guard");
assert!(error.message().contains(DRIFT_ADVICE), "{error}");
}
#[test]
fn the_pin_store_is_bounded() {
let pins = FileValidatorPins::default();
let cap = crate::range_http::MAX_PINNED_VALIDATORS;
for index in 0..(cap * 2) {
let validator = FileValidator {
len: index as u64,
modified: None,
platform_id: None,
};
assert_eq!(pins.check(&format!("path-{index}"), validator), Ok(()));
}
assert_eq!(pins.lock().len(), cap);
}
#[test]
fn a_real_archive_streams_through_the_provider_that_reads_a_remote_one() {
use oxigis_ui::{ArchiveTileProvider, TileProvider as _};
let path = temp_archive("pmtiles", &oxigis_render::pmtiles::sample_pmtiles_raster());
let transport = FileRangeTransport::new().expect("worker threads must start");
let provider = ArchiveTileProvider::pmtiles(
path.display().to_string(),
&egui::Context::default(),
Box::new(transport),
)
.expect("the provider must build");
let tile = oxigis_render::TileId::new(0, 0, 0).expect("0/0/0");
let mut decoded = None;
for _ in 0..200 {
if let Some(pixels) = provider.tile(tile) {
decoded = Some(pixels);
break;
}
if let Some(failure) = provider.failure() {
panic!("the local archive failed to open: {failure}");
}
std::thread::sleep(std::time::Duration::from_millis(25));
}
let pixels = decoded.expect("a tile must arrive within 5 s");
assert_eq!(pixels.width(), 2);
assert_eq!(&pixels.rgba()[..3], &[220, 40, 40]);
let _removed = std::fs::remove_file(&path);
}
#[test]
fn a_local_mbtiles_path_pages_all_the_way_to_a_drawn_tile() {
use oxigis_core::ArchiveFormat;
use oxigis_ui::{ArchiveContent, ArchiveProbe, ArchiveTileProvider, TileProvider as _};
let path = temp_archive("mbtiles", &oxigis_ui::sample_mbtiles_raster());
let location = path.display().to_string();
let probe = ArchiveProbe::start(
location.clone(),
ArchiveFormat::MbTiles,
&egui::Context::default(),
Box::new(FileRangeTransport::new().expect("worker threads must start")),
);
let mut surveyed = None;
for _ in 0..200 {
if let Some(answer) = probe.take() {
surveyed = Some(answer.expect("the local archive must survey cleanly"));
break;
}
std::thread::sleep(std::time::Duration::from_millis(25));
}
let opened = surveyed.expect("the survey must finish within 5 s");
assert_eq!(opened.info.content, ArchiveContent::Raster);
assert_eq!(opened.info.min_zoom, 0);
assert_eq!(opened.info.max_zoom, 2);
assert_eq!(opened.location, location);
let declared_total = std::fs::metadata(&path).map(|metadata| metadata.len()).ok();
assert!(declared_total.is_some(), "the fixture was just written");
let provider = ArchiveTileProvider::paged_mbtiles(
location,
&egui::Context::default(),
Box::new(FileRangeTransport::new().expect("worker threads must start")),
declared_total,
)
.expect("the provider must build");
let root = oxigis_render::TileId::new(0, 0, 0).expect("0/0/0");
let mut decoded = None;
for _ in 0..200 {
if let Some(pixels) = provider.tile(root) {
decoded = Some(pixels);
break;
}
if let Some(failure) = provider.failure() {
panic!("the local archive failed to page: {failure}");
}
std::thread::sleep(std::time::Duration::from_millis(25));
}
let pixels = decoded.unwrap_or_else(|| {
panic!(
"0/0/0 must arrive within 5 s; stats {:?}, failure {:?}",
provider.stats(),
provider.failure()
)
});
assert_eq!(pixels.width(), 2);
assert_eq!(
&pixels.rgba()[..3],
&[10, 120, 200],
"the MBTiles fixture's colour — a wrong archive would say so"
);
let flipped = oxigis_render::TileId::new(1, 0, 1).expect("1/0/1");
let mut second = None;
for _ in 0..200 {
if let Some(pixels) = provider.tile(flipped) {
second = Some(pixels);
break;
}
if let Some(failure) = provider.failure() {
panic!("1/0/1 failed to page: {failure}");
}
std::thread::sleep(std::time::Duration::from_millis(25));
}
let second = second.unwrap_or_else(|| {
panic!(
"1/0/1 must arrive within 5 s; stats {:?}, failure {:?}",
provider.stats(),
provider.failure()
)
});
assert_eq!(second.width(), 2);
assert!(provider.is_open());
assert_eq!(provider.stats().failed, 0);
let _removed = std::fs::remove_file(&path);
}
}