lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

commit 60661379104182ecbe693597b859cf9bfc177867
parent 91e037fc4ddd3fa82248139e594ce913164123fd
Author: triesap <tyson@radroots.org>
Date:   Mon, 20 Jul 2026 03:45:31 +0000

geocoder: bound asset downloads

- resolve hosts asynchronously with bounded runtime teardown
- stream verified responses into atomic same-directory temp files
- expose stable timeout and transport failure phases
- prove overflow, stalls, redirects, cleanup, and target preservation

Diffstat:
MCHANGELOG.md | 7+++++++
MCargo.lock | 125++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcontracts/releases/1.0.0-alpha.1.toml | 12++++++++++++
Mcrates/geocoder/Cargo.toml | 3++-
Mcrates/geocoder/src/asset.rs | 753+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----
Mcrates/geocoder/src/error.rs | 52+++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcrates/geocoder/src/lib.rs | 2+-
7 files changed, 915 insertions(+), 39 deletions(-)

diff --git a/CHANGELOG.md b/CHANGELOG.md @@ -20,6 +20,13 @@ publish policy both pass for the same source revision. identifiers as opaque strings. SQLite integer and text values normalize to the same lossless public representation, and mixed-storage candidate order remains deterministic. +- Default GeoNames asset installation now uses cancellable asynchronous DNS, + explicit connect, response, read, and total deadlines, denied redirects, and + bounded runtime shutdown. HTTP bodies stream through incremental length and + SHA-256 verification into a same-directory tempfile that is synced, checked + for SQLite integrity and schema, and atomically persisted. Public download + errors expose stable typed phases and detail fields instead of + `reqwest::Error`; trusted injected fetchers retain a bounded byte adapter. - Trusted event-contract admission now has one signature-verified entry point. Profile, root Post, Reply, Comment, DeletionRequest, and FoodAvailability retain typed admitted values; other registered events require full contract diff --git a/Cargo.lock b/Cargo.lock @@ -1736,6 +1736,18 @@ dependencies = [ ] [[package]] +name = "enum-as-inner" +version = "0.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a1e6a265c649f3f5979b601d26f1d05ada116434c87741c9493cb56218f76cbc" +dependencies = [ + "heck", + "proc-macro2", + "quote", + "syn 2.0.117", +] + +[[package]] name = "enum-map" version = "2.7.3" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2300,6 +2312,52 @@ dependencies = [ ] [[package]] +name = "hickory-proto" +version = "0.25.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f8a6fe56c0038198998a6f217ca4e7ef3a5e51f46163bd6dd60b5c71ca6c6502" +dependencies = [ + "async-trait", + "cfg-if", + "data-encoding", + "enum-as-inner", + "futures-channel", + "futures-io", + "futures-util", + "idna", + "ipnet", + "once_cell", + "rand 0.9.2", + "ring", + "thiserror 2.0.18", + "tinyvec", + "tokio", + "tracing", + "url", +] + +[[package]] +name = "hickory-resolver" +version = "0.25.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dc62a9a99b0bfb44d2ab95a7208ac952d31060efc16241c87eaf36406fecf87a" +dependencies = [ + "cfg-if", + "futures-util", + "hickory-proto", + "ipconfig", + "moka", + "once_cell", + "parking_lot", + "rand 0.9.2", + "resolv-conf", + "smallvec", + "thiserror 2.0.18", + "tokio", + "tracing", +] + +[[package]] name = "hkdf" version = "0.12.4" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2665,6 +2723,19 @@ dependencies = [ ] [[package]] +name = "ipconfig" +version = "0.3.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4d40460c0ce33d6ce4b0630ad68ff63d6661961c48b6dba35e5a4d81cfb48222" +dependencies = [ + "socket2 0.6.3", + "widestring", + "windows-registry", + "windows-result", + "windows-sys 0.61.2", +] + +[[package]] name = "ipnet" version = "2.12.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -3106,6 +3177,23 @@ dependencies = [ ] [[package]] +name = "moka" +version = "0.12.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "957228ad12042ee839f93c8f257b62b4c0ab5eaae1d4fa60de53b27c9d7c5046" +dependencies = [ + "crossbeam-channel", + "crossbeam-epoch", + "crossbeam-utils", + "equivalent", + "parking_lot", + "portable-atomic", + "smallvec", + "tagptr", + "uuid", +] + +[[package]] name = "mti" version = "1.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -3420,6 +3508,10 @@ name = "once_cell" version = "1.21.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9f7c3e4beb33f85d45ae3e3a1792185706c8e16d043238c593331cc7cd313b50" +dependencies = [ + "critical-section", + "portable-atomic", +] [[package]] name = "once_cell_polyfill" @@ -4383,6 +4475,7 @@ dependencies = [ "sqlx", "tempfile", "thiserror 1.0.69", + "tokio", "url", ] @@ -5169,9 +5262,9 @@ checksum = "eddd3ca559203180a307f12d114c268abf583f59b03cb906fd0b3ff8646c1147" dependencies = [ "base64 0.22.1", "bytes", - "futures-channel", "futures-core", "futures-util", + "hickory-resolver", "http", "http-body", "http-body-util", @@ -5180,6 +5273,7 @@ dependencies = [ "hyper-util", "js-sys", "log", + "once_cell", "percent-encoding", "pin-project-lite", "quinn", @@ -5204,6 +5298,12 @@ dependencies = [ ] [[package]] +name = "resolv-conf" +version = "0.7.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1e061d1b48cb8d38042de4ae0a7a6401009d6143dc80d2e2d6f31f0bdd6470c7" + +[[package]] name = "rfc6979" version = "0.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -7213,6 +7313,12 @@ dependencies = [ ] [[package]] +name = "tagptr" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7b2093cf4c8eb1e67749a6762251bc9cd836b6fc171623bd0a9d324d37af2417" + +[[package]] name = "tap" version = "1.0.1" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -8153,6 +8259,12 @@ dependencies = [ ] [[package]] +name = "widestring" +version = "1.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "72069c3113ab32ab29e5584db3c6ec55d416895e60715417b5b883a357c3e471" + +[[package]] name = "winapi" version = "0.3.9" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -8235,6 +8347,17 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f0805222e57f7521d6a62e36fa9163bc891acd422f971defe97d64e70d0a4fe5" [[package]] +name = "windows-registry" +version = "0.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "02752bf7fbdcce7f2a27a742f798510f3e5ad88dbe84871e5168e2120c3d5720" +dependencies = [ + "windows-link", + "windows-result", + "windows-strings", +] + +[[package]] name = "windows-result" version = "0.4.1" source = "registry+https://github.com/rust-lang/crates.io-index" diff --git a/contracts/releases/1.0.0-alpha.1.toml b/contracts/releases/1.0.0-alpha.1.toml @@ -344,3 +344,15 @@ semver_impacts = [ "change_exported_algorithm_behavior", ] summary = "Replace raw event-store migration SQL and unrestricted destructive rollback with a transactional, checksummed schema authority, exact legacy adoption, descriptor-owned catalog validation and deltas, tamper-evident fail-closed history, a no-write current-schema fast path, SQLite and FTS5 integrity checks, read-only schema inspection, and terminal pool-closing rollback with a governed version floor." + +[[changes]] +id = "geonames-asset-download-boundary" +classification = "breaking" +semver_impacts = [ + "add_exported_type", + "add_exported_function", + "change_exported_field_type", + "change_exported_enum_variant", + "change_exported_algorithm_behavior", +] +summary = "Replace dependency-specific GeoNames download errors and whole-body default installation with stable typed failure phases, cancellable DNS and bounded deadlines, and incrementally verified atomic streaming installation." diff --git a/crates/geocoder/Cargo.toml b/crates/geocoder/Cargo.toml @@ -19,12 +19,13 @@ test-fixture-geonames-asset = [] hex = { workspace = true } futures-executor = { workspace = true } radroots_runtime_paths = { workspace = true } -reqwest = { workspace = true, features = ["blocking", "rustls-tls"] } +reqwest = { workspace = true, features = ["hickory-dns", "rustls-tls"] } serde = { workspace = true, features = ["derive"] } sha2 = { workspace = true } sqlx = { workspace = true, features = ["derive", "sqlite-bundled"] } tempfile = { workspace = true } thiserror = { workspace = true } +tokio = { workspace = true, features = ["rt", "time"] } url = { workspace = true } [lints.rust] diff --git a/crates/geocoder/src/asset.rs b/crates/geocoder/src/asset.rs @@ -1,6 +1,8 @@ use std::fs::{self, File, OpenOptions}; use std::io::{Read, Write}; use std::path::{Path, PathBuf}; +use std::thread; +use std::time::Duration; use radroots_runtime_paths::default_shared_geonames_database_path_from_cache_root; use sha2::{Digest, Sha256}; @@ -8,7 +10,7 @@ use sqlx::Connection; use sqlx::sqlite::{SqliteConnectOptions, SqliteConnection}; use url::Url; -use crate::GeocoderError; +use crate::{GeoNamesAssetDownloadError, GeoNamesAssetDownloadPhase, GeocoderError}; pub const GEONAMES_ASSET_VERSION: &str = "1.0"; pub const GEONAMES_ASSET_FILE_NAME: &str = "geonames-1.0.db"; @@ -34,6 +36,13 @@ pub const GEONAMES_1_0_ASSET: GeoNamesAssetSpec = GeoNamesAssetSpec { sha256: GEONAMES_ASSET_SHA256, }; +const GEONAMES_HTTP_CONNECT_TIMEOUT: Duration = Duration::from_secs(10); +const GEONAMES_HTTP_RESPONSE_TIMEOUT: Duration = Duration::from_secs(15); +const GEONAMES_HTTP_READ_TIMEOUT: Duration = Duration::from_secs(15); +const GEONAMES_HTTP_TOTAL_TIMEOUT: Duration = Duration::from_secs(120); +const GEONAMES_HTTP_RUNTIME_SHUTDOWN_TIMEOUT: Duration = Duration::from_millis(250); +const GEONAMES_HTTP_INITIAL_CAPACITY_MAX: usize = 16 * 1024; + #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub struct GeoNamesAssetSpec { pub version: &'static str, @@ -64,7 +73,44 @@ pub struct GeoNamesAssetStatus { } pub trait GeoNamesAssetFetcher { + /// Returns a complete asset from a trusted or injected source. + /// + /// The default writer adapter bounds this result before installation. + /// Network fetchers should override [`Self::fetch_to_writer`] to avoid + /// buffering the complete asset. fn fetch(&self, url: &str) -> Result<Vec<u8>, GeocoderError>; + + /// Returns a complete asset after enforcing its maximum logical size. + fn fetch_with_max_bytes( + &self, + url: &str, + maximum_bytes: u64, + ) -> Result<Vec<u8>, GeocoderError> { + let bytes = self.fetch(url)?; + let actual = u64::try_from(bytes.len()).unwrap_or(u64::MAX); + if actual > maximum_bytes { + return Err(asset_download_error( + url, + GeoNamesAssetDownloadError::ResponseTooLarge { + maximum: maximum_bytes, + observed_at_least: actual, + }, + )); + } + Ok(bytes) + } + + /// Writes a logically bounded asset to an installation destination. + fn fetch_to_writer( + &self, + url: &str, + maximum_bytes: u64, + destination: &mut (dyn Write + Send), + ) -> Result<(), GeocoderError> { + let bytes = self.fetch_with_max_bytes(url, maximum_bytes)?; + destination.write_all(&bytes)?; + Ok(()) + } } #[derive(Clone, Debug, Default)] @@ -73,25 +119,246 @@ pub struct GeoNamesBlockingHttpFetcher; impl GeoNamesAssetFetcher for GeoNamesBlockingHttpFetcher { #[cfg_attr(coverage_nightly, coverage(off))] fn fetch(&self, url: &str) -> Result<Vec<u8>, GeocoderError> { - let response = - reqwest::blocking::get(url).map_err(|source| GeocoderError::AssetDownload { - url: url.to_owned(), - source, - })?; - let response = - response - .error_for_status() - .map_err(|source| GeocoderError::AssetDownload { - url: url.to_owned(), - source, - })?; - response - .bytes() - .map(|bytes| bytes.to_vec()) - .map_err(|source| GeocoderError::AssetDownload { - url: url.to_owned(), - source, + self.fetch_with_max_bytes(url, GEONAMES_ASSET_BYTE_SIZE) + } + + #[cfg_attr(coverage_nightly, coverage(off))] + fn fetch_with_max_bytes( + &self, + url: &str, + maximum_bytes: u64, + ) -> Result<Vec<u8>, GeocoderError> { + let initial_capacity = usize::try_from(maximum_bytes) + .unwrap_or(usize::MAX) + .min(GEONAMES_HTTP_INITIAL_CAPACITY_MAX); + let mut bytes = Vec::with_capacity(initial_capacity); + fetch_http_asset_to_writer_with_policy( + url, + maximum_bytes, + &mut bytes, + GeoNamesHttpFetchPolicy::production(), + )?; + Ok(bytes) + } + + #[cfg_attr(coverage_nightly, coverage(off))] + fn fetch_to_writer( + &self, + url: &str, + maximum_bytes: u64, + destination: &mut (dyn Write + Send), + ) -> Result<(), GeocoderError> { + fetch_http_asset_to_writer_with_policy( + url, + maximum_bytes, + destination, + GeoNamesHttpFetchPolicy::production(), + ) + } +} + +#[derive(Clone, Copy)] +struct GeoNamesHttpFetchPolicy { + connect_timeout: Duration, + response_timeout: Duration, + read_timeout: Duration, + total_timeout: Duration, + runtime_shutdown_timeout: Duration, +} + +impl GeoNamesHttpFetchPolicy { + const fn production() -> Self { + Self { + connect_timeout: GEONAMES_HTTP_CONNECT_TIMEOUT, + response_timeout: GEONAMES_HTTP_RESPONSE_TIMEOUT, + read_timeout: GEONAMES_HTTP_READ_TIMEOUT, + total_timeout: GEONAMES_HTTP_TOTAL_TIMEOUT, + runtime_shutdown_timeout: GEONAMES_HTTP_RUNTIME_SHUTDOWN_TIMEOUT, + } + } +} + +#[cfg_attr(coverage_nightly, coverage(off))] +fn fetch_http_asset_to_writer_with_policy( + url: &str, + maximum_bytes: u64, + destination: &mut (dyn Write + Send), + policy: GeoNamesHttpFetchPolicy, +) -> Result<(), GeocoderError> { + thread::scope(|scope| { + let worker_url = url.to_owned(); + let worker = thread::Builder::new() + .name("radroots-geonames-http".to_owned()) + .spawn_scoped(scope, move || { + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .map_err(|source| { + asset_download_error( + &worker_url, + GeoNamesAssetDownloadError::Runtime { + detail: source.to_string(), + }, + ) + })?; + let result = runtime.block_on(async { + tokio::time::timeout( + policy.total_timeout, + fetch_http_asset_async(&worker_url, maximum_bytes, destination, policy), + ) + .await + .map_err(|_| { + timeout_error( + &worker_url, + GeoNamesAssetDownloadPhase::Total, + policy.total_timeout, + ) + })? + }); + runtime.shutdown_timeout(policy.runtime_shutdown_timeout); + result }) + .map_err(|source| { + asset_download_error( + url, + GeoNamesAssetDownloadError::Runtime { + detail: source.to_string(), + }, + ) + })?; + worker + .join() + .map_err(|_| asset_download_error(url, GeoNamesAssetDownloadError::WorkerTerminated))? + }) +} + +async fn fetch_http_asset_async( + url: &str, + maximum_bytes: u64, + destination: &mut (dyn Write + Send), + policy: GeoNamesHttpFetchPolicy, +) -> Result<(), GeocoderError> { + let client = reqwest::Client::builder() + .connect_timeout(policy.connect_timeout) + .hickory_dns(true) + // Source validation grants authority to one exact host. + .redirect(reqwest::redirect::Policy::none()) + .build() + .map_err(|source| request_error(url, GeoNamesAssetDownloadPhase::Setup, source))?; + let mut response = tokio::time::timeout(policy.response_timeout, client.get(url).send()) + .await + .map_err(|_| { + timeout_error( + url, + GeoNamesAssetDownloadPhase::Response, + policy.response_timeout, + ) + })? + .map_err(|source| { + if source.is_timeout() && source.is_connect() { + timeout_error( + url, + GeoNamesAssetDownloadPhase::Connect, + policy.connect_timeout, + ) + } else { + let phase = if source.is_connect() { + GeoNamesAssetDownloadPhase::Connect + } else { + GeoNamesAssetDownloadPhase::Response + }; + request_error(url, phase, source) + } + })?; + + if !response.status().is_success() { + return Err(asset_download_error( + url, + GeoNamesAssetDownloadError::HttpStatus { + status: response.status().as_u16(), + }, + )); + } + if let Some(content_length) = response.content_length() + && content_length > maximum_bytes + { + return Err(asset_download_error( + url, + GeoNamesAssetDownloadError::ResponseTooLarge { + maximum: maximum_bytes, + observed_at_least: content_length, + }, + )); + } + + let mut downloaded = 0_u64; + loop { + let chunk = tokio::time::timeout(policy.read_timeout, response.chunk()) + .await + .map_err(|_| timeout_error(url, GeoNamesAssetDownloadPhase::Read, policy.read_timeout))? + .map_err(|source| response_read_error(url, source))?; + let Some(chunk) = chunk else { + break; + }; + let chunk_len = u64::try_from(chunk.len()).unwrap_or(u64::MAX); + let observed_at_least = downloaded.saturating_add(chunk_len); + if observed_at_least > maximum_bytes { + return Err(asset_download_error( + url, + GeoNamesAssetDownloadError::ResponseTooLarge { + maximum: maximum_bytes, + observed_at_least, + }, + )); + } + destination.write_all(&chunk)?; + downloaded = observed_at_least; + } + destination.flush()?; + Ok(()) +} + +fn request_error( + url: &str, + phase: GeoNamesAssetDownloadPhase, + source: reqwest::Error, +) -> GeocoderError { + asset_download_error( + url, + GeoNamesAssetDownloadError::Request { + phase, + detail: source.to_string(), + }, + ) +} + +fn response_read_error(url: &str, source: reqwest::Error) -> GeocoderError { + asset_download_error( + url, + GeoNamesAssetDownloadError::Read { + detail: source.to_string(), + }, + ) +} + +fn timeout_error(url: &str, phase: GeoNamesAssetDownloadPhase, timeout: Duration) -> GeocoderError { + asset_download_error( + url, + GeoNamesAssetDownloadError::Timeout { + phase, + timeout_ms: duration_millis(timeout), + }, + ) +} + +fn duration_millis(duration: Duration) -> u64 { + u64::try_from(duration.as_millis()).unwrap_or(u64::MAX) +} + +fn asset_download_error(url: &str, source: GeoNamesAssetDownloadError) -> GeocoderError { + GeocoderError::AssetDownload { + url: url.to_owned(), + source, } } @@ -148,8 +415,7 @@ where if inspection.state == GeoNamesAssetState::Available { return Ok(inspection); } - let bytes = fetcher.fetch(spec.url)?; - install_geonames_asset_bytes(path, spec, &bytes)?; + install_geonames_asset_with_fetcher(path, spec, fetcher)?; let mut status = validate_geonames_asset_file(path, spec)?; status.state = GeoNamesAssetState::Refreshed; Ok(status) @@ -245,33 +511,113 @@ pub fn validate_geonames_asset_spec_source(spec: &GeoNamesAssetSpec) -> Result<( } #[cfg_attr(coverage_nightly, coverage(off))] -fn install_geonames_asset_bytes( +fn install_geonames_asset_with_fetcher<F>( path: &Path, spec: &GeoNamesAssetSpec, - bytes: &[u8], -) -> Result<(), GeocoderError> { - if bytes.len() as u64 != spec.byte_size { - return Err(GeocoderError::InvalidAssetLength { - path: path.to_path_buf(), - expected: spec.byte_size, - actual: bytes.len() as u64, - }); - } + fetcher: &F, +) -> Result<(), GeocoderError> +where + F: GeoNamesAssetFetcher, +{ let parent = path.parent().unwrap_or_else(|| Path::new(".")); fs::create_dir_all(parent)?; let mut tempfile = tempfile::Builder::new() .prefix(&format!(".{}.", spec.file_name)) .suffix(".tmp") .tempfile_in(parent)?; - tempfile.as_file_mut().write_all(bytes)?; + let identity = { + let mut writer = GeoNamesAssetIdentityWriter::new(tempfile.as_file_mut(), spec.byte_size); + fetcher.fetch_to_writer(spec.url, spec.byte_size, &mut writer)?; + writer.flush()?; + writer.finish() + }; + validate_downloaded_asset_identity(path, spec, &identity)?; tempfile.as_file_mut().sync_all()?; - validate_geonames_asset_file(tempfile.path(), spec)?; + validate_sqlite_integrity_and_schema(tempfile.path())?; tempfile .persist(path) .map(|_| ()) .map_err(|error| GeocoderError::Io(error.error)) } +struct GeoNamesAssetIdentity { + byte_size: u64, + sha256: String, +} + +struct GeoNamesAssetIdentityWriter<'a> { + destination: &'a mut File, + maximum_bytes: u64, + byte_size: u64, + hasher: Sha256, +} + +impl<'a> GeoNamesAssetIdentityWriter<'a> { + fn new(destination: &'a mut File, maximum_bytes: u64) -> Self { + Self { + destination, + maximum_bytes, + byte_size: 0, + hasher: Sha256::new(), + } + } + + fn finish(self) -> GeoNamesAssetIdentity { + GeoNamesAssetIdentity { + byte_size: self.byte_size, + sha256: hex::encode(self.hasher.finalize()), + } + } +} + +impl Write for GeoNamesAssetIdentityWriter<'_> { + fn write(&mut self, bytes: &[u8]) -> std::io::Result<usize> { + let requested = u64::try_from(bytes.len()).unwrap_or(u64::MAX); + let observed_at_least = self.byte_size.saturating_add(requested); + if observed_at_least > self.maximum_bytes { + return Err(std::io::Error::new( + std::io::ErrorKind::FileTooLarge, + format!( + "GeoNames asset exceeds maximum {} bytes; observed at least {observed_at_least}", + self.maximum_bytes + ), + )); + } + let written = self.destination.write(bytes)?; + self.byte_size = self + .byte_size + .saturating_add(u64::try_from(written).unwrap_or(u64::MAX)); + self.hasher.update(&bytes[..written]); + Ok(written) + } + + fn flush(&mut self) -> std::io::Result<()> { + self.destination.flush() + } +} + +fn validate_downloaded_asset_identity( + path: &Path, + spec: &GeoNamesAssetSpec, + identity: &GeoNamesAssetIdentity, +) -> Result<(), GeocoderError> { + if identity.byte_size != spec.byte_size { + return Err(GeocoderError::InvalidAssetLength { + path: path.to_path_buf(), + expected: spec.byte_size, + actual: identity.byte_size, + }); + } + if identity.sha256 != spec.sha256 { + return Err(GeocoderError::InvalidAssetSha256 { + path: path.to_path_buf(), + expected: spec.sha256.to_owned(), + actual: identity.sha256.clone(), + }); + } + Ok(()) +} + #[cfg_attr(coverage_nightly, coverage(off))] fn validate_sqlite_integrity_and_schema(path: &Path) -> Result<(), GeocoderError> { let mut conn = futures_executor::block_on(SqliteConnection::connect_with( @@ -379,7 +725,11 @@ impl Drop for GeoNamesAssetLock { mod tests { use std::cell::Cell; use std::fs; + use std::io::{Read, Write}; + use std::net::{TcpListener, TcpStream}; use std::path::{Path, PathBuf}; + use std::thread; + use std::time::{Duration, Instant}; use sha2::Digest; use sqlx::Connection; @@ -387,11 +737,12 @@ mod tests { use super::{ GEONAMES_ASSET_HOST, GeoNamesAssetFetcher, GeoNamesAssetSpec, GeoNamesAssetState, - ensure_geonames_asset_path_with_fetcher, inspect_geonames_asset_path, + GeoNamesHttpFetchPolicy, ensure_geonames_asset_path_with_fetcher, + fetch_http_asset_to_writer_with_policy, inspect_geonames_asset_path, is_invalid_asset_error, lock_path_for_asset, validate_geonames_asset_file, validate_geonames_asset_spec_source, }; - use crate::GeocoderError; + use crate::{GeoNamesAssetDownloadError, GeoNamesAssetDownloadPhase, GeocoderError}; const TEST_URL: &str = "https://assets.radroots.io/data/geonames/geonames-test.db"; @@ -407,6 +758,218 @@ mod tests { } } + struct WriterOnlyFetcher { + bytes: Vec<u8>, + calls: Cell<usize>, + } + + impl GeoNamesAssetFetcher for WriterOnlyFetcher { + fn fetch(&self, _url: &str) -> Result<Vec<u8>, GeocoderError> { + panic!("writer-oriented install must not call the Vec adapter") + } + + fn fetch_to_writer( + &self, + _url: &str, + maximum_bytes: u64, + destination: &mut (dyn Write + Send), + ) -> Result<(), GeocoderError> { + self.calls.set(self.calls.get() + 1); + assert!(self.bytes.len() as u64 <= maximum_bytes); + for chunk in self.bytes.chunks(257) { + destination.write_all(chunk)?; + } + Ok(()) + } + } + + #[test] + fn blocking_http_fetch_streams_a_bounded_success_response() { + let server = LoopbackHttpServer::spawn(|mut stream| { + read_request(&mut stream); + stream + .write_all( + b"HTTP/1.1 200 OK\r\nContent-Length: 5\r\nConnection: close\r\n\r\nhello", + ) + .expect("success response"); + }); + + let bytes = fetch_http_bytes(&server.url, 5, test_http_policy()).expect("bounded download"); + assert_eq!(bytes, b"hello"); + } + + #[test] + fn blocking_http_fetch_allows_progress_to_exceed_read_timeout_cumulatively() { + let server = LoopbackHttpServer::spawn(|mut stream| { + read_request(&mut stream); + stream + .write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 5\r\nConnection: close\r\n\r\n") + .expect("response headers"); + for byte in b"hello" { + stream.write_all(&[*byte]).expect("response byte"); + stream.flush().expect("response flush"); + thread::sleep(Duration::from_millis(75)); + } + }); + let policy = GeoNamesHttpFetchPolicy { + connect_timeout: Duration::from_millis(500), + response_timeout: Duration::from_secs(1), + read_timeout: Duration::from_millis(200), + total_timeout: Duration::from_secs(2), + runtime_shutdown_timeout: Duration::from_millis(100), + }; + + let bytes = fetch_http_bytes(&server.url, 5, policy).expect("progressing download"); + assert_eq!(bytes, b"hello"); + } + + #[test] + fn blocking_http_fetch_rejects_oversized_chunked_response_without_buffering_it() { + let server = LoopbackHttpServer::spawn(|mut stream| { + read_request(&mut stream); + stream + .write_all( + b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\nConnection: close\r\n\r\n4\r\nabcd\r\n4\r\nefgh\r\n0\r\n\r\n", + ) + .expect("chunked response"); + }); + + assert!(matches!( + fetch_http_bytes(&server.url, 5, test_http_policy()), + Err(GeocoderError::AssetDownload { + source: GeoNamesAssetDownloadError::ResponseTooLarge { + maximum: 5, + observed_at_least, + }, + .. + }) if observed_at_least >= 6 + )); + } + + #[test] + fn blocking_http_fetch_times_out_a_stalled_response_body() { + let server = LoopbackHttpServer::spawn(|mut stream| { + read_request(&mut stream); + stream + .write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 1\r\nConnection: close\r\n\r\n") + .expect("response headers"); + stream.flush().expect("response flush"); + thread::sleep(Duration::from_millis(600)); + }); + let policy = GeoNamesHttpFetchPolicy { + connect_timeout: Duration::from_millis(500), + response_timeout: Duration::from_secs(1), + read_timeout: Duration::from_millis(200), + total_timeout: Duration::from_secs(2), + runtime_shutdown_timeout: Duration::from_millis(100), + }; + + let started_at = Instant::now(); + let result = fetch_http_bytes(&server.url, 1, policy); + let elapsed = started_at.elapsed(); + assert!( + matches!( + result, + Err(GeocoderError::AssetDownload { + source: GeoNamesAssetDownloadError::Timeout { + phase: GeoNamesAssetDownloadPhase::Read, + timeout_ms: 200, + }, + .. + }) + ), + "{result:?}" + ); + assert!(elapsed >= Duration::from_millis(100), "{elapsed:?}"); + assert!(elapsed < Duration::from_secs(3), "{elapsed:?}"); + } + + #[test] + fn blocking_http_fetch_reports_response_timeout_without_elapsed_heuristics() { + let server = LoopbackHttpServer::spawn(|mut stream| { + read_request(&mut stream); + thread::sleep(Duration::from_millis(600)); + }); + let policy = GeoNamesHttpFetchPolicy { + connect_timeout: Duration::from_millis(500), + response_timeout: Duration::from_millis(200), + read_timeout: Duration::from_secs(1), + total_timeout: Duration::from_secs(2), + runtime_shutdown_timeout: Duration::from_millis(100), + }; + + assert!(matches!( + fetch_http_bytes(&server.url, 1, policy), + Err(GeocoderError::AssetDownload { + source: GeoNamesAssetDownloadError::Timeout { + phase: GeoNamesAssetDownloadPhase::Response, + timeout_ms: 200, + }, + .. + }) + )); + } + + #[test] + fn blocking_http_fetch_enforces_total_deadline_across_progressing_reads() { + let server = LoopbackHttpServer::spawn(|mut stream| { + read_request(&mut stream); + stream + .write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 4\r\nConnection: close\r\n\r\n") + .expect("response headers"); + for byte in b"slow" { + if stream.write_all(&[*byte]).is_err() { + break; + } + let _ = stream.flush(); + thread::sleep(Duration::from_millis(150)); + } + }); + let policy = GeoNamesHttpFetchPolicy { + connect_timeout: Duration::from_millis(500), + response_timeout: Duration::from_secs(1), + read_timeout: Duration::from_millis(500), + total_timeout: Duration::from_millis(250), + runtime_shutdown_timeout: Duration::from_millis(100), + }; + + let started_at = Instant::now(); + let result = fetch_http_bytes(&server.url, 4, policy); + let elapsed = started_at.elapsed(); + assert!(matches!( + result, + Err(GeocoderError::AssetDownload { + source: GeoNamesAssetDownloadError::Timeout { + phase: GeoNamesAssetDownloadPhase::Total, + timeout_ms: 250, + }, + .. + }) + )); + assert!(elapsed >= Duration::from_millis(125), "{elapsed:?}"); + assert!(elapsed < Duration::from_secs(3), "{elapsed:?}"); + } + + #[test] + fn blocking_http_fetch_does_not_follow_redirects_outside_the_validated_source() { + let server = LoopbackHttpServer::spawn(|mut stream| { + read_request(&mut stream); + stream + .write_all( + b"HTTP/1.1 302 Found\r\nLocation: https://example.com/geonames.db\r\nContent-Length: 0\r\nConnection: close\r\n\r\n", + ) + .expect("redirect response"); + }); + + assert!(matches!( + fetch_http_bytes(&server.url, 5, test_http_policy()), + Err(GeocoderError::AssetDownload { + source: GeoNamesAssetDownloadError::HttpStatus { status: 302 }, + .. + }) + )); + } + #[test] fn geonames_asset_missing_available_invalid_and_refreshed_states_are_reported() { let tempdir = tempfile::tempdir().expect("tempdir"); @@ -445,6 +1008,49 @@ mod tests { } #[test] + fn geonames_asset_install_uses_streaming_writer_and_atomic_replacement() { + let tempdir = tempfile::tempdir().expect("tempdir"); + let bytes = fixture_database_bytes(); + let spec = fixture_spec(&bytes, TEST_URL); + let target = tempdir.path().join("geonames-test.db"); + let fetcher = WriterOnlyFetcher { + bytes, + calls: Cell::new(0), + }; + + let status = + ensure_geonames_asset_path_with_fetcher(&target, &spec, &fetcher).expect("install"); + assert_eq!(status.state, GeoNamesAssetState::Refreshed); + assert_eq!(fetcher.calls.get(), 1); + assert_no_install_tempfiles(tempdir.path()); + } + + #[test] + fn geonames_asset_failed_identity_validation_preserves_existing_target() { + let tempdir = tempfile::tempdir().expect("tempdir"); + let bytes = fixture_database_bytes(); + let spec = fixture_spec(&bytes, TEST_URL); + let target = tempdir.path().join("geonames-test.db"); + let previous = b"previous asset"; + fs::write(&target, previous).expect("previous target"); + let mut wrong_bytes = bytes; + let last = wrong_bytes.last_mut().expect("nonempty fixture"); + *last ^= 0x01; + let fetcher = WriterOnlyFetcher { + bytes: wrong_bytes, + calls: Cell::new(0), + }; + + assert!(matches!( + ensure_geonames_asset_path_with_fetcher(&target, &spec, &fetcher), + Err(GeocoderError::InvalidAssetSha256 { .. }) + )); + assert_eq!(fs::read(&target).expect("preserved target"), previous); + assert_eq!(fetcher.calls.get(), 1); + assert_no_install_tempfiles(tempdir.path()); + } + + #[test] fn geonames_asset_rejects_wrong_host_length_hash_sqlite_and_schema() { let tempdir = tempfile::tempdir().expect("tempdir"); let bytes = fixture_database_bytes(); @@ -631,6 +1237,83 @@ mod tests { .expect("execute fixture sql batch"); } + fn fetch_http_bytes( + url: &str, + maximum_bytes: u64, + policy: GeoNamesHttpFetchPolicy, + ) -> Result<Vec<u8>, GeocoderError> { + let mut bytes = Vec::new(); + fetch_http_asset_to_writer_with_policy(url, maximum_bytes, &mut bytes, policy)?; + Ok(bytes) + } + + fn test_http_policy() -> GeoNamesHttpFetchPolicy { + GeoNamesHttpFetchPolicy { + connect_timeout: Duration::from_secs(1), + response_timeout: Duration::from_secs(1), + read_timeout: Duration::from_secs(1), + total_timeout: Duration::from_secs(2), + runtime_shutdown_timeout: Duration::from_millis(100), + } + } + + struct LoopbackHttpServer { + url: String, + thread: Option<thread::JoinHandle<()>>, + } + + impl LoopbackHttpServer { + fn spawn(handler: impl FnOnce(TcpStream) + Send + 'static) -> Self { + let listener = TcpListener::bind("127.0.0.1:0").expect("loopback bind"); + let address = listener.local_addr().expect("loopback address"); + let thread = thread::spawn(move || { + let (stream, _) = listener.accept().expect("loopback accept"); + handler(stream); + }); + Self { + url: format!("http://{address}/geonames.db"), + thread: Some(thread), + } + } + } + + impl Drop for LoopbackHttpServer { + fn drop(&mut self) { + if let Some(thread) = self.thread.take() { + let _ = thread.join(); + } + } + } + + fn read_request(stream: &mut TcpStream) { + stream + .set_read_timeout(Some(Duration::from_secs(1))) + .expect("request read timeout"); + let mut request = Vec::with_capacity(2048); + let mut buffer = [0_u8; 256]; + loop { + let read = stream.read(&mut buffer).expect("request"); + assert_ne!(read, 0, "connection closed before complete request headers"); + request.extend_from_slice(&buffer[..read]); + assert!(request.len() <= 32 * 1024, "request headers are bounded"); + if request.windows(4).any(|window| window == b"\r\n\r\n") { + return; + } + } + } + + fn assert_no_install_tempfiles(parent: &Path) { + let entries = fs::read_dir(parent).expect("asset parent"); + for entry in entries { + let file_name = entry.expect("asset parent entry").file_name(); + let file_name = file_name.to_string_lossy(); + assert!( + !(file_name.starts_with(".geonames-test.db.") && file_name.ends_with(".tmp")), + "temporary install file leaked: {file_name}" + ); + } + } + const FIXTURE_SCHEMA: &str = r#" CREATE TABLE countries(id TEXT, name TEXT); CREATE TABLE admin1(country_id TEXT, id INTEGER, name TEXT); diff --git a/crates/geocoder/src/error.rs b/crates/geocoder/src/error.rs @@ -1,5 +1,55 @@ use thiserror::Error; +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +#[non_exhaustive] +pub enum GeoNamesAssetDownloadPhase { + Setup, + Connect, + Response, + Read, + Total, +} + +impl std::fmt::Display for GeoNamesAssetDownloadPhase { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter.write_str(match self { + Self::Setup => "setup", + Self::Connect => "connect", + Self::Response => "response", + Self::Read => "read", + Self::Total => "total", + }) + } +} + +#[derive(Debug, Error)] +#[non_exhaustive] +pub enum GeoNamesAssetDownloadError { + #[error("HTTP worker runtime failed: {detail}")] + Runtime { detail: String }, + #[error("HTTP worker terminated unexpectedly")] + WorkerTerminated, + #[error("{phase} request failed: {detail}")] + Request { + phase: GeoNamesAssetDownloadPhase, + detail: String, + }, + #[error("HTTP status {status}")] + HttpStatus { status: u16 }, + #[error("{phase} timeout after {timeout_ms} ms")] + Timeout { + phase: GeoNamesAssetDownloadPhase, + timeout_ms: u64, + }, + #[error("response read failed: {detail}")] + Read { detail: String }, + #[error("response body exceeds maximum {maximum} bytes; observed at least {observed_at_least}")] + ResponseTooLarge { + maximum: u64, + observed_at_least: u64, + }, +} + #[derive(Debug, Error)] pub enum GeocoderError { #[error("sqlite error: {0}")] @@ -49,7 +99,7 @@ pub enum GeocoderError { AssetDownload { url: String, #[source] - source: reqwest::Error, + source: GeoNamesAssetDownloadError, }, #[error("country center not found for {country_id}")] CountryCenterNotFound { country_id: String }, diff --git a/crates/geocoder/src/lib.rs b/crates/geocoder/src/lib.rs @@ -15,7 +15,7 @@ pub use asset::{ inspect_default_geonames_asset_in_cache_root, inspect_geonames_asset_path, validate_geonames_asset_file, validate_geonames_asset_spec_source, }; -pub use error::GeocoderError; +pub use error::{GeoNamesAssetDownloadError, GeoNamesAssetDownloadPhase, GeocoderError}; pub use geocoder::Geocoder; pub use model::{ GeocoderCountryListResult, GeocoderLocalityCandidate, GeocoderLocalityInput,