diff --git a/Cargo.lock b/Cargo.lock index d8809ee..3543712 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -249,6 +249,58 @@ version = "1.5.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f2032f911046de80f0a198e0901378627c33f59ea0ac00e363d481118bd70a53" +[[package]] +name = "axum" +version = "0.8.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "31b698c5f9a010f6573133b09e0de5408834d0c82f8d7475a89fc1867a71cd90" +dependencies = [ + "axum-core", + "bytes", + "form_urlencoded", + "futures-util", + "http", + "http-body", + "http-body-util", + "hyper", + "hyper-util", + "itoa", + "matchit", + "memchr", + "mime", + "percent-encoding", + "pin-project-lite", + "serde_core", + "serde_json", + "serde_path_to_error", + "serde_urlencoded", + "sync_wrapper", + "tokio", + "tower", + "tower-layer", + "tower-service", + "tracing", +] + +[[package]] +name = "axum-core" +version = "0.5.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "08c78f31d7b1291f7ee735c1c6780ccde7785daae9a9206026862dab7d8792d1" +dependencies = [ + "bytes", + "futures-core", + "http", + "http-body", + "http-body-util", + "mime", + "pin-project-lite", + "sync_wrapper", + "tower-layer", + "tower-service", + "tracing", +] + [[package]] name = "base64" version = "0.22.1" @@ -267,6 +319,15 @@ version = "2.13.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b588b76d00fde79687d7646a9b5bdf3cc0f655e0bbd080335a95d7e96f3587da" +[[package]] +name = "block-buffer" +version = "0.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d2f6c7dbe95a6ed67ad9f18e57daf93a2f034c524b99fd2b76d18fdfeb6660aa" +dependencies = [ + "hybrid-array", +] + [[package]] name = "bumpalo" version = "3.20.3" @@ -334,6 +395,12 @@ dependencies = [ "unicode-width", ] +[[package]] +name = "const-oid" +version = "0.10.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a6ef517f0926dd24a1582492c791b6a4818a4d94e789a334894aa15b0d12f55c" + [[package]] name = "const-random" version = "0.1.18" @@ -380,6 +447,15 @@ version = "0.8.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "773648b94d0e5d620f64f280777445740e61fe701025087ec8b57f45c791888b" +[[package]] +name = "cpufeatures" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8b2a41393f66f16b0823bb79094d54ac5fbd34ab292ddafb9a0456ac9f87d201" +dependencies = [ + "libc", +] + [[package]] name = "crc32fast" version = "1.5.0" @@ -417,6 +493,15 @@ version = "0.2.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "460fbee9c2c2f33933d720630a6a0bac33ba7053db5344fac858d4b8952d77d5" +[[package]] +name = "crypto-common" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ce6e4c961d6cd6c9a86db418387425e8bdeaf05b3c8bc1411e6dca4c252f1453" +dependencies = [ + "hybrid-array", +] + [[package]] name = "derive_arbitrary" version = "1.4.2" @@ -428,6 +513,17 @@ dependencies = [ "syn 2.0.119", ] +[[package]] +name = "digest" +version = "0.11.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f1dd6dbb5841937940781866fa1281a1ff7bd3bf827091440879f9994983d5c2" +dependencies = [ + "block-buffer", + "const-oid", + "crypto-common", +] + [[package]] name = "displaydoc" version = "0.2.7" @@ -659,15 +755,20 @@ name = "hagfish" version = "0.1.0" dependencies = [ "anyhow", + "axum", "duckdb", + "mime_guess", "mockito", "reqwest", + "rust-embed", "serde", "serde_json", "tempfile", "thiserror", "tokio", "tokio-test", + "tokio-util", + "tower-http", "tracing", "tracing-subscriber", ] @@ -759,6 +860,15 @@ version = "1.0.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "df3b46402a9d5adb4c86a0cf463f42e19994e3ee891101b1841f30a545cb49a9" +[[package]] +name = "hybrid-array" +version = "0.4.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "707114b52a152fa7bdb290cd7cd5912d9467273b6d74e21b8d81aca1f8533f6b" +dependencies = [ + "typenum", +] + [[package]] name = "hyper" version = "1.11.0" @@ -1142,6 +1252,12 @@ dependencies = [ "regex-automata", ] +[[package]] +name = "matchit" +version = "0.8.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "47e1ffaa40ddd1f3ed91f717a33c8c0ee23fff369e3aa8772b9605cc1d22f4c3" + [[package]] name = "memchr" version = "2.8.3" @@ -1154,6 +1270,16 @@ version = "0.3.17" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6877bb514081ee2a7ff5ef9de3281f14a4dd4bceac4c09388074a6b5df8a139a" +[[package]] +name = "mime_guess" +version = "2.0.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f7c44f8e672c00fe5308fa235f821cb4198414e1c77935c1ab6948d3fd78550e" +dependencies = [ + "mime", + "unicase", +] + [[package]] name = "miniz_oxide" version = "0.8.9" @@ -1523,6 +1649,41 @@ dependencies = [ "windows-sys 0.52.0", ] +[[package]] +name = "rust-embed" +version = "8.12.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e9e7760e252aaba7b09f4be00e36476cf585bdb68a53552ac954cdf504ab4bc9" +dependencies = [ + "rust-embed-impl", + "rust-embed-utils", + "walkdir", +] + +[[package]] +name = "rust-embed-impl" +version = "8.12.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3bcfc4d6f53af43755f7a723e4b6b8794fcce052a178dd8c6c1dadc5f5343097" +dependencies = [ + "mime_guess", + "proc-macro2", + "quote", + "rust-embed-utils", + "syn 2.0.119", + "walkdir", +] + +[[package]] +name = "rust-embed-utils" +version = "8.12.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "42ffa149f6aa81b58a5b3011d01a857c4ed12c7a732d2c51947a4c7c692185f0" +dependencies = [ + "sha2", + "walkdir", +] + [[package]] name = "rustix" version = "0.38.44" @@ -1533,7 +1694,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys 0.4.15", - "windows-sys 0.52.0", + "windows-sys 0.59.0", ] [[package]] @@ -1596,6 +1757,15 @@ version = "1.0.23" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9774ba4a74de5f7b1c1451ed6cd5285a32eddb5cccb8cc655a4e50009e06477f" +[[package]] +name = "same-file" +version = "1.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "93fc1dc3aaa9bfed95e02e6eadabb4baf7e3078b0bd1b4d7b6b0b68378900502" +dependencies = [ + "winapi-util", +] + [[package]] name = "schannel" version = "0.1.29" @@ -1677,6 +1847,17 @@ dependencies = [ "zmij", ] +[[package]] +name = "serde_path_to_error" +version = "0.1.20" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "10a9ff822e371bb5403e391ecd83e182e0e77ba7f6fe0160b795797109d1b457" +dependencies = [ + "itoa", + "serde", + "serde_core", +] + [[package]] name = "serde_urlencoded" version = "0.7.1" @@ -1689,6 +1870,17 @@ dependencies = [ "serde", ] +[[package]] +name = "sha2" +version = "0.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "446ba717509524cb3f22f17ecc096f10f4822d76ab5c0b9822c5f9c284e825f4" +dependencies = [ + "cfg-if", + "cpufeatures", + "digest", +] + [[package]] name = "sharded-slab" version = "0.1.7" @@ -2013,6 +2205,7 @@ dependencies = [ "tokio", "tower-layer", "tower-service", + "tracing", ] [[package]] @@ -2030,6 +2223,7 @@ dependencies = [ "tower", "tower-layer", "tower-service", + "tracing", "url", ] @@ -2051,6 +2245,7 @@ version = "0.1.44" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "63e71662fa4b2a2c3a26f570f037eb95bb1f85397f3cd8076caed2f026a6d100" dependencies = [ + "log", "pin-project-lite", "tracing-attributes", "tracing-core", @@ -2112,6 +2307,18 @@ version = "0.2.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" +[[package]] +name = "typenum" +version = "1.20.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b6f5e870be6c3b371b77fe0ee0bafb859fa4964b4404c27de1d380043c4dda20" + +[[package]] +name = "unicase" +version = "2.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dbc4bc3a9f746d862c45cb89d705aa10f187bb96c76001afab07a0d35ce60142" + [[package]] name = "unicode-ident" version = "1.0.24" @@ -2206,6 +2413,16 @@ version = "0.9.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a" +[[package]] +name = "walkdir" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "29790946404f91d9c5d06f9874efddea1dc06c5efe94541a7d6863108e3a5e4b" +dependencies = [ + "same-file", + "winapi-util", +] + [[package]] name = "want" version = "0.3.1" @@ -2320,6 +2537,15 @@ version = "0.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ac3b87c63620426dd9b991e5ce0329eff545bccbbb34f3be09ff6fb6ab51b7b6" +[[package]] +name = "winapi-util" +version = "0.1.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" +dependencies = [ + "windows-sys 0.61.2", +] + [[package]] name = "winapi-x86_64-pc-windows-gnu" version = "0.4.0" @@ -2405,6 +2631,15 @@ dependencies = [ "windows-targets", ] +[[package]] +name = "windows-sys" +version = "0.59.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1e38bc4d79ed67fd075bcc251a1c39b32a1776bbe92e5bef1f0bf1f8c531853b" +dependencies = [ + "windows-targets", +] + [[package]] name = "windows-sys" version = "0.61.2" diff --git a/Cargo.toml b/Cargo.toml index 5904dcf..16fbc4a 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -15,8 +15,14 @@ tracing = "0.1" tracing-subscriber = { version = "0.3", features = ["env-filter"] } anyhow = "1" duckdb = { version = "1.10505", features = ["bundled"] } +axum = "0.8" +tower-http = { version = "0.6", features = ["cors", "trace"] } +rust-embed = "8.12" +mime_guess = "2" +tokio-util = { version = "0.7", features = ["io"] } +tempfile = "3" [dev-dependencies] mockito = "1" tokio-test = "0.4" -tempfile = "3" +reqwest = "0.12" diff --git a/docs/TODO.md b/docs/TODO.md index 22506f0..b2a6368 100644 --- a/docs/TODO.md +++ b/docs/TODO.md @@ -1,4 +1,4 @@ -# HAGFISH Project todo +# HAGFISH Project TODO `GET statbank.hagstova.fo/api/v1/fo/H2/VV/VV01/fisknv_md.px — metadata` POST same URL — JSON-stat2 query @@ -78,14 +78,22 @@ hagfish/ ### Phase 3: API (Axum) -- [ ] 3.1 Set up Axum router in main.rs with Tokio runtime. AppState holds DuckDB connection wrapped in Mutex. -- [ ] 3.2 Implement GET /api/species — query species lookup table, return JSON array -- [ ] 3.3 Implement GET /api/landings — parse query params, build DuckDB SQL with WHERE clauses. Support: months, species, gear, zone, measure filters. -- [ ] 3.4 Implement GET /api/summary — aggregate query: total mass + value by month, top 10 species by value, price/kg trend -- [ ] 3.5 Implement GET /api/export.parquet — generate and stream Parquet via DuckDB COPY -- [ ] 3.6 Implement GET /healthz -- [ ] 3.7 Serve static files via rust-embed -- [ ] 3.8 Write API tests +- [x] 3.1 Set up Axum router in main.rs with Tokio runtime. AppState holds DuckDB connection wrapped in Mutex. +- [x] 3.2 Implement GET /api/species — query species lookup table, return JSON array +- [x] 3.3 Implement GET /api/landings — parse query params, build DuckDB SQL with WHERE clauses. Support: months, species, gear, zone, measure filters. +- [x] 3.4 Implement GET /api/summary — aggregate query: total mass + value by month, top 10 species by value, price/kg trend +- [x] 3.5 Implement GET /api/export.parquet — generate and stream Parquet via DuckDB COPY +- [x] 3.6 Implement GET /healthz +- [x] 3.7 Serve static files via rust-embed +- [x] 3.8 Write API tests + +#### Production Hardening (Deferred to Post-Launch) + +- [ ] 🔒 DB operation timeouts — Add 30s `tokio::time::timeout()` wrapper around all `spawn_blocking` DB calls +- [ ] 🔒 CORS hardening — Replace `CorsLayer::permissive()` with explicit allowed origins before production deployment +- [ ] 🔒 Rate limiting — Optional: add `tower_governor` middleware to prevent abuse on public deployments +- [ ] ⚠️ Startup config validation — Verify `data_source_url` is reachable, `duckdb_path` is writable before accepting connections +- [ ] 📊 Metrics export (Prometheus) — Optional: track request counts, latencies, error rates via `prometheus` crate ### Phase 4: Frontend @@ -111,3 +119,13 @@ hagfish/ - [ ] 6.2 Write systemd timer (hagfish-ingest.timer + hagfish-ingest.service) — monthly, runs hagfish ingest - [ ] 6.3 Taskfile: build (release, static), deploy (rsync binary + config + units, ssh reload) - [ ] 6.4 README with ELI5 Technology Choices section (why DuckDB, why Rust, why embedded static assets) + +--- + +## Current Status + +**Phase 1**: Complete ✅ +**Phase 2**: Complete ✅ +**Phase 3**: Complete ✅ (API operational, tests passing) +**Production Hardening**: Deferred ⏸️ +**Phase 4-6**: Pending — Start after MVP validation diff --git a/hagfish.db b/hagfish.db new file mode 100644 index 0000000..20abae0 Binary files /dev/null and b/hagfish.db differ diff --git a/hagfish.db.wal b/hagfish.db.wal new file mode 100644 index 0000000..372f2b1 Binary files /dev/null and b/hagfish.db.wal differ diff --git a/src/api.rs b/src/api.rs new file mode 100644 index 0000000..66d0ffe --- /dev/null +++ b/src/api.rs @@ -0,0 +1,866 @@ +use crate::db; +use crate::types::{ + Config, LandingDto, MonthlyAggregate, PriceTrend, SpeciesDto, SummaryDto, TopSpecies, +}; +use axum::{ + Json, Router, + body::Body, + extract::{Query, State}, + http::{StatusCode, header}, + response::{IntoResponse, Response}, + routing::get, +}; +use duckdb::params; +use rust_embed::Embed; +use serde::Deserialize; +use std::sync::Arc; +use tokio::sync::Mutex; +use tower_http::cors::CorsLayer; +use tower_http::trace::TraceLayer; + +#[derive(Clone)] +pub struct AppState { + pub conn: Arc>, + #[allow(dead_code)] + pub config: Arc, +} + +#[derive(Embed)] +#[folder = "static/"] +struct StaticAssets; + +#[derive(Debug, Clone, Deserialize)] +pub struct LandingsQuery { + pub month: Option, + pub species: Option, + pub gear: Option, + pub zone: Option, + pub measure: Option, + #[serde(default = "default_limit")] + pub limit: u32, +} + +fn default_limit() -> u32 { + 10000 +} + +#[derive(Debug, Clone, Deserialize)] +pub struct SummaryQuery { + pub species: Option, +} + +#[derive(Debug)] +pub struct ApiError { + pub status: StatusCode, + pub message: String, +} + +impl IntoResponse for ApiError { + fn into_response(self) -> Response { + let body = serde_json::json!({ "error": self.message }); + (self.status, Json(body)).into_response() + } +} + +impl From for ApiError { + fn from(e: duckdb::Error) -> Self { + tracing::error!("DuckDB error: {e}"); + ApiError { + status: StatusCode::INTERNAL_SERVER_ERROR, + message: format!("Database error: {e}"), + } + } +} + +impl From for ApiError { + fn from(e: db::DbError) -> Self { + tracing::error!("DB error: {e}"); + ApiError { + status: StatusCode::INTERNAL_SERVER_ERROR, + message: format!("Database error: {e}"), + } + } +} + +type ApiResult = std::result::Result; + +pub fn build_router(state: AppState) -> Router { + Router::new() + .route("/healthz", get(healthz)) + .route("/api/species", get(get_species)) + .route("/api/landings", get(get_landings)) + .route("/api/summary", get(get_summary)) + .route("/api/export.parquet", get(export_parquet)) + .fallback(static_handler) + .layer(CorsLayer::permissive()) + .layer(TraceLayer::new_for_http()) + .with_state(state) +} + +async fn healthz(State(_state): State) -> impl IntoResponse { + StatusCode::OK +} + +async fn get_species(State(state): State) -> ApiResult>> { + let conn = state.conn.clone(); + let rows = tokio::task::spawn_blocking(move || -> db::Result> { + let conn = conn.blocking_lock(); + let mut stmt = conn.prepare("SELECT code, label FROM species ORDER BY code")?; + let rows = stmt.query_map([], |row| { + Ok(SpeciesDto { + code: row.get(0)?, + label: row.get(1)?, + }) + })?; + Ok(rows.collect::, duckdb::Error>>()?) + }) + .await + .map_err(|e| ApiError { + status: StatusCode::INTERNAL_SERVER_ERROR, + message: format!("Task join error: {e}"), + })??; + + Ok(Json(rows)) +} + +async fn get_landings( + State(state): State, + Query(params): Query, +) -> ApiResult>> { + let conn = state.conn.clone(); + let rows = tokio::task::spawn_blocking(move || -> db::Result> { + let conn = conn.blocking_lock(); + + let mut sql = String::from( + "SELECT month, species_code, species_label, gear_code, zone_code, \ + processing_code, preservation_code, shipsize_code, measure_code, value \ + FROM landings WHERE 1=1", + ); + let mut args: Vec> = Vec::new(); + let mut idx = 1; + + if let Some(ref month) = params.month { + sql.push_str(&format!(" AND month = ${idx}")); + args.push(Box::new(month.clone())); + idx += 1; + } + if let Some(ref species) = params.species { + sql.push_str(&format!(" AND species_code = ${idx}")); + args.push(Box::new(species.clone())); + idx += 1; + } + if let Some(ref gear) = params.gear { + sql.push_str(&format!(" AND gear_code = ${idx}")); + args.push(Box::new(gear.clone())); + idx += 1; + } + if let Some(ref zone) = params.zone { + sql.push_str(&format!(" AND zone_code = ${idx}")); + args.push(Box::new(zone.clone())); + idx += 1; + } + if let Some(ref measure) = params.measure { + sql.push_str(&format!(" AND measure_code = ${idx}")); + args.push(Box::new(measure.clone())); + idx += 1; + } + + sql.push_str(&format!( + " ORDER BY month, species_code, measure_code LIMIT ${idx}" + )); + args.push(Box::new(params.limit as i64)); + + let arg_refs: Vec<&dyn duckdb::ToSql> = args.iter().map(|b| b.as_ref()).collect(); + + let mut stmt = conn.prepare(&sql)?; + let rows = stmt.query_map(arg_refs.as_slice(), |row| { + Ok(LandingDto { + month: row.get(0)?, + species_code: row.get(1)?, + species_label: row.get(2)?, + gear_code: row.get(3)?, + zone_code: row.get(4)?, + processing_code: row.get(5)?, + preservation_code: row.get(6)?, + shipsize_code: row.get(7)?, + measure_code: row.get(8)?, + value: row.get(9)?, + }) + })?; + + Ok(rows.collect::, duckdb::Error>>()?) + }) + .await + .map_err(|e| ApiError { + status: StatusCode::INTERNAL_SERVER_ERROR, + message: format!("Task join error: {e}"), + })??; + + Ok(Json(rows)) +} + +async fn get_summary( + State(state): State, + Query(params): Query, +) -> ApiResult> { + let conn = state.conn.clone(); + let species_filter = params.species.clone(); + + let result = tokio::task::spawn_blocking(move || -> db::Result { + let conn = conn.blocking_lock(); + + let where_clause = if species_filter.is_some() { + " WHERE species_code = ?".to_string() + } else { + String::new() + }; + + let monthly_sql = format!( + "SELECT month, \ + SUM(CASE WHEN measure_code = 'MASS' THEN value END) AS total_mass, \ + SUM(CASE WHEN measure_code = 'VALUE' THEN value END) AS total_value \ + FROM landings{where_clause} \ + GROUP BY month ORDER BY month" + ); + + let monthly: Vec = if let Some(ref sp) = species_filter { + let mut stmt = conn.prepare(&monthly_sql)?; + let rows = stmt.query_map(params![sp], |row| { + Ok(MonthlyAggregate { + month: row.get(0)?, + total_mass: row.get(1)?, + total_value: row.get(2)?, + }) + })?; + Ok::, duckdb::Error>( + rows.collect::, _>>()?, + )? + } else { + let mut stmt = conn.prepare(&monthly_sql)?; + let rows = stmt.query_map([], |row| { + Ok(MonthlyAggregate { + month: row.get(0)?, + total_mass: row.get(1)?, + total_value: row.get(2)?, + }) + })?; + Ok::, duckdb::Error>( + rows.collect::, _>>()?, + )? + }; + + let top_sql = format!( + "SELECT species_code, \ + COALESCE(MAX(species_label), species_code) AS species_label, \ + SUM(CASE WHEN measure_code = 'VALUE' THEN value END) AS total_value, \ + SUM(CASE WHEN measure_code = 'MASS' THEN value END) AS total_mass \ + FROM landings{where_clause} \ + GROUP BY species_code, species_label \ + ORDER BY total_value DESC NULLS LAST \ + LIMIT 10" + ); + + let top_species: Vec = if let Some(ref sp) = species_filter { + let mut stmt = conn.prepare(&top_sql)?; + let rows = stmt.query_map(params![sp], |row| { + Ok(TopSpecies { + species_code: row.get(0)?, + species_label: row.get(1)?, + total_value: row.get(2)?, + total_mass: row.get(3)?, + }) + })?; + Ok::, duckdb::Error>(rows.collect::, _>>()?)? + } else { + let mut stmt = conn.prepare(&top_sql)?; + let rows = stmt.query_map([], |row| { + Ok(TopSpecies { + species_code: row.get(0)?, + species_label: row.get(1)?, + total_value: row.get(2)?, + total_mass: row.get(3)?, + }) + })?; + Ok::, duckdb::Error>(rows.collect::, _>>()?)? + }; + + let price_sql = format!( + "WITH monthly_mass AS ( \ + SELECT month, SUM(value) AS mass FROM landings \ + WHERE measure_code = 'MASS'{where_clause} \ + GROUP BY month \ + ), \ + monthly_value AS ( \ + SELECT month, SUM(value) AS value FROM landings \ + WHERE measure_code = 'VALUE'{where_clause} \ + GROUP BY month \ + ) \ + SELECT m.month, \ + CASE WHEN m.mass IS NOT NULL AND m.mass > 0 \ + THEN v.value / m.mass END AS price_per_kg \ + FROM monthly_mass m \ + JOIN monthly_value v ON m.month = v.month \ + ORDER BY m.month" + ); + + let price_trend: Vec = if let Some(ref sp) = species_filter { + let mut stmt = conn.prepare(&price_sql)?; + let rows = stmt.query_map(params![sp, sp], |row| { + Ok(PriceTrend { + month: row.get(0)?, + price_per_kg: row.get(1)?, + }) + })?; + Ok::, duckdb::Error>(rows.collect::, _>>()?)? + } else { + let mut stmt = conn.prepare(&price_sql)?; + let rows = stmt.query_map([], |row| { + Ok(PriceTrend { + month: row.get(0)?, + price_per_kg: row.get(1)?, + }) + })?; + Ok::, duckdb::Error>(rows.collect::, _>>()?)? + }; + + Ok(SummaryDto { + monthly, + top_species, + price_trend, + }) + }) + .await + .map_err(|e| ApiError { + status: StatusCode::INTERNAL_SERVER_ERROR, + message: format!("Task join error: {e}"), + })??; + + Ok(Json(result)) +} + +async fn export_parquet(State(state): State) -> ApiResult { + let conn = state.conn.clone(); + + let (file_path, file) = + tokio::task::spawn_blocking(move || -> ApiResult<(String, std::fs::File)> { + let conn = conn.blocking_lock(); + + let tmp = tempfile::NamedTempFile::new().map_err(|e| ApiError { + status: StatusCode::INTERNAL_SERVER_ERROR, + message: format!("Failed to create temp file: {e}"), + })?; + + let (_kept_file, path) = tmp.keep().map_err(|e| ApiError { + status: StatusCode::INTERNAL_SERVER_ERROR, + message: format!("Failed to persist temp file: {e}"), + })?; + + let path_str = path + .to_str() + .ok_or_else(|| ApiError { + status: StatusCode::INTERNAL_SERVER_ERROR, + message: "Temp path contains invalid UTF-8".to_string(), + })? + .to_string(); + + db::export_parquet(&conn, &path_str)?; + + let file = std::fs::File::open(&path_str).map_err(|e| ApiError { + status: StatusCode::INTERNAL_SERVER_ERROR, + message: format!("Failed to open exported file: {e}"), + })?; + + Ok((path_str, file)) + }) + .await + .map_err(|e| ApiError { + status: StatusCode::INTERNAL_SERVER_ERROR, + message: format!("Task join error: {e}"), + })??; + + let tokio_file = tokio::fs::File::from(file); + let stream = tokio_util::io::ReaderStream::new(tokio_file); + let body = Body::from_stream(stream); + + let cleanup_path = file_path.clone(); + tokio::spawn(async move { + tokio::time::sleep(std::time::Duration::from_secs(300)).await; + if let Err(e) = tokio::fs::remove_file(&cleanup_path).await { + tracing::warn!("Failed to clean up parquet file {}: {e}", cleanup_path); + } + }); + + Ok(Response::builder() + .status(StatusCode::OK) + .header(header::CONTENT_TYPE, "application/octet-stream") + .header( + header::CONTENT_DISPOSITION, + "attachment; filename=\"landings.parquet\"", + ) + .body(body) + .unwrap()) +} + +async fn static_handler(uri: axum::http::Uri) -> Response { + let path = uri.path().trim_start_matches('/'); + let asset = StaticAssets::get(path).or_else(|| StaticAssets::get("index.html")); + + match asset { + Some(file) => { + let mime = mime_guess::from_path(path) + .first_or_octet_stream() + .as_ref() + .to_string(); + + Response::builder() + .status(StatusCode::OK) + .header(header::CONTENT_TYPE, mime) + .body(Body::from(file.data.into_owned())) + .unwrap() + } + None => Response::builder() + .status(StatusCode::NOT_FOUND) + .body(Body::from("Not found")) + .unwrap(), + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::types::Landing; + use duckdb::Connection; + use std::collections::HashMap; + + fn test_state() -> AppState { + let conn = Connection::open_in_memory().unwrap(); + db::init_schema(&conn).unwrap(); + + let rows = vec![ + Landing { + month: "2024M01".to_string(), + species_code: "COD".to_string(), + species_label: "Toskur".to_string(), + gear_code: "TOTAL".to_string(), + zone_code: "TOTAL".to_string(), + processing_code: "TOTAL".to_string(), + preservation_code: "TOTAL".to_string(), + shipsize_code: "TOTAL".to_string(), + measure_code: "MASS".to_string(), + value: Some(1000.0), + }, + Landing { + month: "2024M01".to_string(), + species_code: "COD".to_string(), + species_label: "Toskur".to_string(), + gear_code: "TOTAL".to_string(), + zone_code: "TOTAL".to_string(), + processing_code: "TOTAL".to_string(), + preservation_code: "TOTAL".to_string(), + shipsize_code: "TOTAL".to_string(), + measure_code: "VALUE".to_string(), + value: Some(5000.0), + }, + Landing { + month: "2024M01".to_string(), + species_code: "HER".to_string(), + species_label: "Sild".to_string(), + gear_code: "TOTAL".to_string(), + zone_code: "TOTAL".to_string(), + processing_code: "TOTAL".to_string(), + preservation_code: "TOTAL".to_string(), + shipsize_code: "TOTAL".to_string(), + measure_code: "MASS".to_string(), + value: Some(2000.0), + }, + Landing { + month: "2024M01".to_string(), + species_code: "HER".to_string(), + species_label: "Sild".to_string(), + gear_code: "TOTAL".to_string(), + zone_code: "TOTAL".to_string(), + processing_code: "TOTAL".to_string(), + preservation_code: "TOTAL".to_string(), + shipsize_code: "TOTAL".to_string(), + measure_code: "VALUE".to_string(), + value: Some(3000.0), + }, + Landing { + month: "2024M02".to_string(), + species_code: "COD".to_string(), + species_label: "Toskur".to_string(), + gear_code: "TOTAL".to_string(), + zone_code: "TOTAL".to_string(), + processing_code: "TOTAL".to_string(), + preservation_code: "TOTAL".to_string(), + shipsize_code: "TOTAL".to_string(), + measure_code: "MASS".to_string(), + value: Some(1500.0), + }, + Landing { + month: "2024M02".to_string(), + species_code: "COD".to_string(), + species_label: "Toskur".to_string(), + gear_code: "TOTAL".to_string(), + zone_code: "TOTAL".to_string(), + processing_code: "TOTAL".to_string(), + preservation_code: "TOTAL".to_string(), + shipsize_code: "TOTAL".to_string(), + measure_code: "VALUE".to_string(), + value: Some(7500.0), + }, + ]; + + db::update_lookups( + &conn, + &HashMap::from([( + "Species (ASFIS2022)".to_string(), + HashMap::from([ + ("COD".to_string(), "Toskur".to_string()), + ("HER".to_string(), "Sild".to_string()), + ]), + )]), + ) + .unwrap(); + + db::upsert_landings(&conn, &rows).unwrap(); + + AppState { + conn: Arc::new(Mutex::new(conn)), + config: Arc::new(Config::default()), + } + } + + #[tokio::test] + async fn test_healthz_returns_ok() { + let state = test_state(); + let app = build_router(state); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + + let resp = reqwest::get(format!("http://{addr}/healthz")) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + } + + #[tokio::test] + async fn test_get_species_returns_all() { + let state = test_state(); + let app = build_router(state); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + + let resp = reqwest::get(format!("http://{addr}/api/species")) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + + let body: Vec = resp.json().await.unwrap(); + assert!(body.iter().any(|s| s.code == "COD" && s.label == "Toskur")); + assert!(body.iter().any(|s| s.code == "HER" && s.label == "Sild")); + } + + #[tokio::test] + async fn test_get_landings_no_filter() { + let state = test_state(); + let app = build_router(state); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + + let resp = reqwest::get(format!("http://{addr}/api/landings")) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + + let body: Vec = resp.json().await.unwrap(); + assert!(!body.is_empty(), "Should have landing rows"); + assert_eq!(body.len(), 6); + } + + #[tokio::test] + async fn test_get_landings_filter_by_species() { + let state = test_state(); + let app = build_router(state); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + + let resp = reqwest::get(format!("http://{addr}/api/landings?species=COD")) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + + let body: Vec = resp.json().await.unwrap(); + assert!(body.iter().all(|l| l.species_code == "COD")); + assert_eq!(body.len(), 4); + } + + #[tokio::test] + async fn test_get_landings_filter_by_measure() { + let state = test_state(); + let app = build_router(state); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + + let resp = reqwest::get(format!("http://{addr}/api/landings?measure=MASS")) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + + let body: Vec = resp.json().await.unwrap(); + assert!(body.iter().all(|l| l.measure_code == "MASS")); + assert_eq!(body.len(), 3); + } + + #[tokio::test] + async fn test_get_summary_monthly_aggregates() { + let state = test_state(); + let app = build_router(state); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + + let resp = reqwest::get(format!("http://{addr}/api/summary")) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + + let body: SummaryDto = resp.json().await.unwrap(); + + assert_eq!(body.monthly.len(), 2); + + let jan = body.monthly.iter().find(|m| m.month == "2024M01").unwrap(); + assert!((jan.total_mass.unwrap() - 3000.0).abs() < f64::EPSILON); + assert!((jan.total_value.unwrap() - 8000.0).abs() < f64::EPSILON); + + let feb = body.monthly.iter().find(|m| m.month == "2024M02").unwrap(); + assert!((feb.total_mass.unwrap() - 1500.0).abs() < f64::EPSILON); + assert!((feb.total_value.unwrap() - 7500.0).abs() < f64::EPSILON); + } + + #[tokio::test] + async fn test_get_summary_top_species() { + let state = test_state(); + let app = build_router(state); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + + let resp = reqwest::get(format!("http://{addr}/api/summary")) + .await + .unwrap(); + + let body: SummaryDto = resp.json().await.unwrap(); + + assert!(!body.top_species.is_empty()); + assert_eq!(body.top_species[0].species_code, "COD"); + assert!((body.top_species[0].total_value.unwrap() - 12500.0).abs() < f64::EPSILON); + } + + #[tokio::test] + async fn test_get_summary_price_trend() { + let state = test_state(); + let app = build_router(state); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + + let resp = reqwest::get(format!("http://{addr}/api/summary")) + .await + .unwrap(); + + let body: SummaryDto = resp.json().await.unwrap(); + + assert_eq!(body.price_trend.len(), 2); + + let jan = body + .price_trend + .iter() + .find(|p| p.month == "2024M01") + .unwrap(); + let expected_jan = 8000.0 / 3000.0; + assert!((jan.price_per_kg.unwrap() - expected_jan).abs() < 0.01); + + let feb = body + .price_trend + .iter() + .find(|p| p.month == "2024M02") + .unwrap(); + assert!((feb.price_per_kg.unwrap() - 5.0).abs() < f64::EPSILON); + } + + #[tokio::test] + async fn test_get_landings_faroese_label_preserved() { + let state = test_state(); + let app = build_router(state); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + + let resp = reqwest::get(format!("http://{addr}/api/landings?species=COD")) + .await + .unwrap(); + + let body: Vec = resp.json().await.unwrap(); + let cod = &body[0]; + assert_eq!(cod.species_label, "Toskur"); + assert!(!cod.species_label.contains('\u{FFFD}')); + } + + #[tokio::test] + async fn test_get_landings_empty_result() { + let state = test_state(); + let app = build_router(state); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + + let resp = reqwest::get(format!("http://{addr}/api/landings?species=NONEXISTENT")) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + + let body: Vec = resp.json().await.unwrap(); + assert!(body.is_empty()); + } + + #[tokio::test] + async fn test_get_landings_limit_applied() { + let state = test_state(); + let app = build_router(state); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + + let resp = reqwest::get(format!("http://{addr}/api/landings?limit=2")) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + + let body: Vec = resp.json().await.unwrap(); + assert_eq!(body.len(), 2); + } + + #[tokio::test] + async fn test_export_parquet() { + let state = test_state(); + let app = build_router(state); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + + let resp = reqwest::get(format!("http://{addr}/api/export.parquet")) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::OK); + assert_eq!( + resp.headers().get(header::CONTENT_DISPOSITION).unwrap(), + "attachment; filename=\"landings.parquet\"" + ); + + let bytes = resp.bytes().await.unwrap(); + assert!(bytes.len() > 4); + assert_eq!(&bytes[..4], b"PAR1"); + } + + #[tokio::test] + async fn test_concurrent_export_requests() { + let state = test_state(); + let app = build_router(state.clone()); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + + let resp1 = reqwest::get(format!("http://{addr}/api/export.parquet")) + .await + .unwrap(); + let resp2 = reqwest::get(format!("http://{addr}/api/export.parquet")) + .await + .unwrap(); + + assert_eq!(resp1.status(), StatusCode::OK); + assert_eq!(resp2.status(), StatusCode::OK); + + let bytes1 = resp1.bytes().await.unwrap(); + let bytes2 = resp2.bytes().await.unwrap(); + + assert_eq!(&bytes1[..4], b"PAR1"); + assert_eq!(&bytes2[..4], b"PAR1"); + assert_ne!(bytes1.len(), 0); + assert_ne!(bytes2.len(), 0); + } + + #[tokio::test] + async fn test_sql_injection_safe() { + let state = test_state(); + let state_clone = state.clone(); + + let app = build_router(state.clone()); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + + let malicious_params = "?species='; DROP TABLE landings; --"; + let url = format!("http://{addr}/api/landings{malicious_params}"); + let resp = reqwest::get(&url).await.unwrap(); + + assert_eq!(resp.status(), StatusCode::OK); + + let count: i64 = tokio::task::spawn_blocking(move || { + let conn = state_clone.conn.blocking_lock(); + conn.query_row("SELECT COUNT(*) FROM landings", [], |row| row.get(0)) + .unwrap() + }) + .await + .unwrap(); + + assert_eq!(count, 6); + } +} diff --git a/src/db.rs b/src/db.rs index d88358a..266f64e 100644 --- a/src/db.rs +++ b/src/db.rs @@ -1,17 +1,9 @@ -//! DuckDB storage layer for hagfish. -//! -//! Handles schema initialization, data upserts, lookup table population, -//! incremental ingestion support, and Parquet export. - use crate::types::{Landing, LookupMap}; use duckdb::{Connection, params}; use std::collections::HashSet; use thiserror::Error; use tracing::info; -/// Dimension code → lookup table name. -/// -/// These map PX-Web variable codes to short table names in DuckDB. const LOOKUP_TABLES: &[(&str, &str)] = &[ ("Species (ASFIS2022)", "species"), ("Fishing Gear (ISSCFG2016)", "gear"), @@ -31,7 +23,6 @@ pub enum DbError { pub type Result = std::result::Result; -/// Opens a DuckDB connection at the given path and initializes the schema. pub fn init(path: &str) -> Result { let conn = Connection::open(path)?; init_schema(&conn)?; @@ -39,7 +30,6 @@ pub fn init(path: &str) -> Result { Ok(conn) } -/// Creates all tables and indexes if they don't exist. pub fn init_schema(conn: &Connection) -> Result<()> { conn.execute_batch( " @@ -91,10 +81,6 @@ pub fn init_schema(conn: &Connection) -> Result<()> { Ok(()) } -/// Populates lookup tables from metadata-derived lookup maps. -/// -/// Each dimension's code→label pairs are upserted into their respective -/// lookup table using ON CONFLICT semantics. pub fn update_lookups(conn: &Connection, lookup_maps: &LookupMap) -> Result<()> { for (dim_code, table_name) in LOOKUP_TABLES { if let Some(codes) = lookup_maps.get(*dim_code) { @@ -116,15 +102,6 @@ pub fn update_lookups(conn: &Connection, lookup_maps: &LookupMap) -> Result<()> Ok(()) } -/// Upserts landing records using delete-then-insert per month. -/// -/// Uses `unchecked_transaction` (takes `&self`) per DuckDB docs, and the -/// Appender API for bulk inserts as recommended by the official -/// DuckDB Rust documentation. -/// -/// All rows belonging to the same month(s) are deleted first, then -/// re-inserted within a single transaction. This ensures that revised -/// data from the API replaces stale records atomically. pub fn upsert_landings(conn: &Connection, rows: &[Landing]) -> Result { if rows.is_empty() { return Ok(0); @@ -154,7 +131,6 @@ pub fn upsert_landings(conn: &Connection, rows: &[Landing]) -> Result { row.value, ])?; } - // Flush explicitly — Drop discards errors per DuckDB docs. app.flush()?; } @@ -168,12 +144,6 @@ pub fn upsert_landings(conn: &Connection, rows: &[Landing]) -> Result { Ok(rows.len()) } -/// Returns the latest month present in the landings table. -/// -/// Used for incremental ingestion: the next fetch starts from the month -/// after this value. Returns `None` if the table is empty. -/// -/// Relies on YYYYMNN format being lexicographically sortable. pub fn get_last_month(conn: &Connection) -> Result> { let result = conn.query_row("SELECT MAX(month) FROM landings", [], |row| { let val: Option = row.get(0)?; @@ -182,12 +152,8 @@ pub fn get_last_month(conn: &Connection) -> Result> { Ok(result) } -/// Exports the landings table to a Parquet file. -/// -/// Caller must ensure the path is writable and does not contain -/// single quotes (which would break the SQL string literal). pub fn export_parquet(conn: &Connection, path: &str) -> Result<()> { - if path.contains('\'') { + if path.contains('\'') || path.contains('\0') { return Err(DbError::InvalidPath(path.to_string())); } conn.execute( @@ -254,8 +220,6 @@ mod tests { ] } - // --- Schema initialization --- - #[test] fn test_init_creates_all_tables() { let conn = test_conn(); @@ -285,8 +249,6 @@ mod tests { init_schema(&conn).unwrap(); } - // --- Upsert operations --- - #[test] fn test_upsert_insert_and_idempotency() { let conn = test_conn(); @@ -328,7 +290,6 @@ mod tests { }]; upsert_landings(&conn, &modified).unwrap(); - // 2024M01 had 2 rows, now has 1. 2024M02 untouched. let count: i64 = conn .query_row("SELECT COUNT(*) FROM landings", [], |row| row.get(0)) .unwrap(); @@ -358,7 +319,6 @@ mod tests { let conn = test_conn(); upsert_landings(&conn, &sample_landings()).unwrap(); - // Upsert both months in one call — both should be replaced. let rows = vec![ Landing { month: "2024M01".to_string(), @@ -405,15 +365,12 @@ mod tests { .unwrap(); assert!((feb_val - 222.0).abs() < f64::EPSILON); - // Old VALUE row for 2024M01 is gone (replaced) let count: i64 = conn .query_row("SELECT COUNT(*) FROM landings", [], |row| row.get(0)) .unwrap(); assert_eq!(count, 2); } - // --- NULL handling --- - #[test] fn test_null_preserved_on_insert() { let conn = test_conn(); @@ -457,8 +414,6 @@ mod tests { assert!(val.is_none()); } - // --- get_last_month --- - #[test] fn test_get_last_month_empty_table() { let conn = test_conn(); @@ -496,8 +451,6 @@ mod tests { assert_eq!(last.as_deref(), Some("2015M03")); } - // --- Lookup tables --- - #[test] fn test_update_lookups_populates_tables() { let conn = test_conn(); @@ -574,7 +527,6 @@ mod tests { let mut species = HashMap::new(); species.insert("COD".to_string(), "Toskur".to_string()); maps.insert("Species (ASFIS2022)".to_string(), species); - // No gear, zone, etc. — should not error update_lookups(&conn, &maps).unwrap(); @@ -584,8 +536,6 @@ mod tests { assert_eq!(count, 0); } - // --- Parquet export --- - #[test] fn test_export_parquet() { let conn = test_conn(); @@ -596,7 +546,6 @@ mod tests { export_parquet(&conn, &path).unwrap(); - // Re-read the Parquet file via DuckDB to verify let count: i64 = conn .query_row( &format!("SELECT COUNT(*) FROM read_parquet('{}')", path), @@ -607,8 +556,6 @@ mod tests { assert_eq!(count, 3); } - // --- Faroese Unicode round-trip --- - #[test] fn test_faroese_labels_round_trip() { let conn = test_conn(); @@ -616,7 +563,7 @@ mod tests { let rows = vec![Landing { month: "2024M01".to_string(), species_code: "148XXXXXXX00000".to_string(), - species_label: "Hýsa".to_string(), // ý — Faroese-specific + species_label: "Hýsa".to_string(), gear_code: "TOTAL".to_string(), zone_code: "TOTAL".to_string(), processing_code: "TOTAL".to_string(), diff --git a/src/ingest.rs b/src/ingest.rs index 5477969..1bb238e 100644 --- a/src/ingest.rs +++ b/src/ingest.rs @@ -1,28 +1,10 @@ -//! Data ingestion from PX-Web API. -//! -//! Handles metadata fetching, query construction, data retrieval, and -//! parsing of JSON-stat2 responses into structured rows. -//! -//! ## API contract (confirmed 2026-08-16) -//! -//! - GET endpoint: `https://statbank.hagstova.fo/api/v1/fo/H2/VV/VV01/fisknv_md.px` -//! - POST endpoint: same URL -//! - Dimension codes are NOT Faroese short names — they use classification -//! identifiers like "Species (ASFIS2022)", "Fishing Gear (ISSCFG2016)". -//! - `role` is null for all variables on this endpoint. -//! - Metadata uses parallel `values`/`valueTexts` string arrays. -//! - Sentinel values: -1.0 indicates missing data. - use crate::types::*; use reqwest::Client; use std::collections::HashMap; use tracing::{debug, info}; -/// Language code for PX-Web queries. const PX_WEB_LANGUAGE: &str = "fo"; -/// Dimension codes — confirmed against live API on 2026-08-16. -/// These are the `code` field values from the metadata response. const DIM_MONTH: &str = "month"; const DIM_SPECIES: &str = "Species (ASFIS2022)"; const DIM_GEAR: &str = "Fishing Gear (ISSCFG2016)"; @@ -32,19 +14,14 @@ const DIM_PRESERVATION: &str = "Preservation (EUMOFAPreservation)"; const DIM_SHIPSIZE: &str = "Shipsize"; const DIM_MEASURE: &str = "measure"; -/// Sentinel f64 values that indicate missing data, coerced to None. const SENTINEL_VALUES: [f64; 1] = [-1.0]; -/// Checks whether a numeric value is a sentinel. fn is_sentinel(v: f64) -> bool { SENTINEL_VALUES .iter() .any(|s| v.total_cmp(s) == std::cmp::Ordering::Equal) } -/// Fetches metadata from PX-Web API endpoint. -/// -/// Returns a `LookupMap` mapping dimension code → (value code → label). pub async fn fetch_metadata(client: &Client, url: &str) -> Result { info!("Fetching metadata from {}", url); @@ -89,9 +66,6 @@ pub async fn fetch_metadata(client: &Client, url: &str) -> Result { Ok(lookup_map) } -/// Extracts all available months from metadata. -/// -/// Uses the "month" dimension code since `role` is null on this endpoint. pub fn extract_available_months(meta: &MetadataResponse) -> Vec { meta.variables .iter() @@ -100,13 +74,6 @@ pub fn extract_available_months(meta: &MetadataResponse) -> Vec { .unwrap_or_default() } -/// Constructs a query body for fetching all landing data. -/// -/// Sets categorical dimensions to wildcard ("all" filter), time to explicit -/// months, and measure to MASS + VALUE. -/// -/// # Panics -/// Panics if `all_months` is empty. pub fn build_query(all_months: &[String]) -> Query { assert!( !all_months.is_empty(), @@ -176,7 +143,6 @@ pub fn build_query(all_months: &[String]) -> Query { } } -/// Fetches data from PX-Web API using the provided query. pub async fn fetch_data(client: &Client, url: &str, query: &Query) -> Result { let month_count = query .query @@ -206,9 +172,6 @@ pub async fn fetch_data(client: &Client, url: &str, query: &Query) -> Result Vec { let mut indices = Vec::with_capacity(key_sizes.len()); let mut remaining = flat_index; @@ -222,14 +185,6 @@ fn decode_key_indices(flat_index: usize, key_sizes: &[usize]) -> Vec { indices } -/// Parses a single row from a JSON-stat2 response into a `DataRow`. -/// -/// Dimension order in the response `id` array determines positional mapping. -/// The expected order (from API metadata) is: -/// measure, species, gear, zone, processing, preservation, shipsize, month -/// -/// However, JSON-stat2 `dimension.id` defines the actual order — we read it -/// dynamically and map by dimension code, not by position assumption. pub fn parse_row( row_index: usize, dataset: &DataResponse, @@ -254,8 +209,6 @@ pub fn parse_row( ))); } - // Build ordered category lists: for each dimension, extract (code, label) - // pairs sorted by their JSON-stat2 index position. let category_lists: Vec> = dim_order .iter() .map(|dim_code| { @@ -301,7 +254,6 @@ pub fn parse_row( .map(|(idx, list)| list[idx].clone()) .collect(); - // Build a lookup from dimension code → (code, label) for this row let mut dim_map: HashMap<&str, (String, String)> = HashMap::with_capacity(8); for (dim_code, values) in dim_order.iter().zip(dimension_values.iter()) { dim_map.insert(dim_code.as_str(), values.clone()); @@ -323,7 +275,6 @@ pub fn parse_row( let (shipsize_code, shipsize_label) = get(DIM_SHIPSIZE); let (measure_code, measure_label) = get(DIM_MEASURE); - // Coerce sentinel values to None let raw_value = dataset.dataset.value[row_index]; let value = match raw_value { Some(v) if is_sentinel(v) => None, @@ -355,7 +306,6 @@ pub fn parse_row( }) } -/// Converts a `DataRow` to a `Landing` struct for database insertion. pub fn data_row_to_landing(row: &DataRow) -> Landing { Landing { month: row.month.clone(), @@ -442,8 +392,6 @@ mod tests { }, ); - // Dimension order from API: month, species, gear, zone, processing, - // preservation, shipsize, measure DataResponse { dataset: Dataset { dimension: DimInfo { @@ -461,12 +409,12 @@ mod tests { dimensions, }, value: vec![ - Some(1234.5), // 2015M01, Sild, ..., MASS - Some(2345.6), // 2015M01, Sild, ..., VALUE + Some(1234.5), + Some(2345.6), Some(-1.0), Some(3456.7), - Some(4567.8), // 2015M02, Sild, ..., MASS - Some(5678.9), // 2015M02, Sild, ..., VALUE + Some(4567.8), + Some(5678.9), None, Some(6789.0), ], @@ -475,7 +423,6 @@ mod tests { } } - /// Test fixture: mock lookup maps with Faroese labels. fn mock_lookup_maps() -> LookupMap { let mut maps = LookupMap::new(); maps.insert( @@ -567,7 +514,6 @@ mod tests { let dataset = mock_dataset_response(); let lookup_maps = mock_lookup_maps(); - // Row 2: MASS, Toskur, ..., 2015M01 let row = parse_row(2, &dataset, &lookup_maps).expect("parse failed"); assert_eq!(row.species_code, "183XXXXXXX00000"); @@ -580,7 +526,6 @@ mod tests { let dataset = mock_dataset_response(); let lookup_maps = mock_lookup_maps(); - // Row 2 has -1.0 sentinel let row = parse_row(2, &dataset, &lookup_maps).expect("parse failed"); assert!( @@ -595,7 +540,6 @@ mod tests { let dataset = mock_dataset_response(); let lookup_maps = mock_lookup_maps(); - // Row 6 has None let row = parse_row(6, &dataset, &lookup_maps).expect("parse failed"); assert!(row.value.is_none()); diff --git a/src/main.rs b/src/main.rs index 6881768..96f4102 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,7 +1,58 @@ +mod api; mod db; mod ingest; mod types; -fn main() { - println!("hagfish — not yet wired. Run tests with: cargo test"); +use axum::serve; +use std::sync::Arc; +use tokio::net::TcpListener; +use tracing_subscriber::EnvFilter; + +fn cleanup_stale_parquet_files() { + let tmp_dir = std::env::temp_dir(); + match std::fs::read_dir(&tmp_dir) { + Ok(entries) => { + for entry in entries.filter_map(|e| e.ok()) { + if let Some(name) = entry.file_name().to_str() { + if name.starts_with("hagfish-") && name.ends_with(".parquet") { + if let Ok(metadata) = entry.metadata() { + let age = std::time::SystemTime::now() + .duration_since(metadata.modified().unwrap()); + if age.map(|d| d.as_secs() > 3600).unwrap_or(false) { + let _ = std::fs::remove_file(entry.path()); + } + } + } + } + } + } + Err(e) => tracing::warn!("Failed to list temp dir: {e}"), + } +} + +#[tokio::main] +async fn main() -> anyhow::Result<()> { + tracing_subscriber::fmt() + .with_env_filter(EnvFilter::from_default_env()) + .init(); + + let config = Arc::new(types::Config::default()); + + cleanup_stale_parquet_files(); + + let conn = db::init(&config.duckdb_path)?; + + let state = api::AppState { + conn: Arc::new(tokio::sync::Mutex::new(conn)), + config: config.clone(), + }; + + let app = api::build_router(state); + + let listener = TcpListener::bind(&config.bind_address).await?; + tracing::info!("hagfish listening on http://{}", config.bind_address); + + serve(listener, app).await?; + + Ok(()) } diff --git a/src/types.rs b/src/types.rs index 963a96d..23a3989 100644 --- a/src/types.rs +++ b/src/types.rs @@ -1,12 +1,6 @@ -//! Type definitions for hagfish data structures. -//! -//! These types represent the PX-Web API contract and internal data models. -//! All public types derive Serialize/Deserialize for JSON (de)serialization. - use serde::{Deserialize, Serialize}; use std::collections::HashMap; -/// Configuration loaded from config.json on startup. #[derive(Debug, Clone, Deserialize, Serialize)] pub struct Config { pub duckdb_path: String, @@ -27,27 +21,6 @@ impl Default for Config { } } -// --------------------------------------------------------------------------- -// Metadata types (GET response) -// --------------------------------------------------------------------------- - -/// Metadata response from PX-Web API GET request. -/// -/// Confirmed structure from live API: -/// ```json -/// { -/// "title": "AVR01010 ...", -/// "variables": [ -/// { -/// "code": "measure", -/// "text": "mát", -/// "role": null, -/// "values": ["MASS", "VALUE"], -/// "valueTexts": ["Nøgd", "Virði"] -/// } -/// ] -/// } -/// ``` #[derive(Debug, Clone, Deserialize, Serialize)] pub struct MetadataResponse { #[serde(default)] @@ -55,34 +28,22 @@ pub struct MetadataResponse { pub variables: Vec, } -/// Variable metadata describing a dimension in the PX-Web table. -/// -/// PxWeb v1 uses parallel string arrays for codes and labels. -/// The `role` field is null on this endpoint — do not rely on it. #[derive(Debug, Clone, Deserialize, Serialize)] pub struct VariableMeta { - /// Variable identifier used in queries (e.g. "month", "measure", - /// "Species (ASFIS2022)") pub code: String, - /// Human-readable label in the queried language #[serde(default)] pub text: String, - /// Role classification — null on this endpoint #[serde(default)] pub role: Option, - /// Code values for this variable (parallel with valueTexts) #[serde(default)] pub values: Vec, - /// Display labels for each value (parallel with values) #[serde(default, rename = "valueTexts")] pub value_texts: Vec, - /// Catch-all for unknown fields #[serde(flatten)] pub extra: HashMap, } impl VariableMeta { - /// Build a code→label lookup map from the parallel arrays. pub fn lookup_map(&self) -> HashMap { self.values .iter() @@ -92,42 +53,24 @@ impl VariableMeta { } } -// --------------------------------------------------------------------------- -// Query types (POST request body) -// --------------------------------------------------------------------------- - -/// Query body sent to PX-Web API POST endpoint. -/// -/// Confirmed format: -/// ```json -/// { -/// "query": [{"code": "month", "selection": {"filter": "item", "values": [...]}}], -/// "response": {"format": "json-stat2"} -/// } -/// ``` #[derive(Debug, Clone, Serialize)] pub struct Query { pub query: Vec, pub response: QueryResponse, } -/// Single query item representing a dimension selection. #[derive(Debug, Clone, Serialize)] pub struct QueryItem { pub code: String, pub selection: Selection, } -/// Selection within a query item. #[derive(Debug, Clone, Serialize)] pub struct Selection { - /// "item" for explicit values, "all" for wildcard, "top" for top-N pub filter: String, - /// Values to select. For "all" filter, use ["*"]. pub values: Vec, } -/// Response format specification. #[derive(Debug, Clone, Serialize)] pub struct QueryResponse { pub format: String, @@ -141,68 +84,43 @@ impl Default for QueryResponse { } } -// --------------------------------------------------------------------------- -// Data response types (JSON-stat2) -// --------------------------------------------------------------------------- - -/// Raw data response from PX-Web API POST request. #[derive(Debug, Clone, Deserialize, Serialize)] pub struct DataResponse { pub dataset: Dataset, } -/// Dataset wrapper containing dimensions and actual values. #[derive(Debug, Clone, Deserialize, Serialize)] pub struct Dataset { pub dimension: DimInfo, pub value: Vec>, - /// Status codes per value (optional) #[serde(default, skip_serializing_if = "Vec::is_empty")] pub status: Vec, } -/// Dimension metadata in JSON-stat2 format. -/// -/// Uses a flat `id` array for ordering and a `size` array for cardinality. -/// Each dimension is keyed by its code in the `dimensions` map. #[derive(Debug, Clone, Deserialize, Serialize)] pub struct DimInfo { - /// Ordered dimension IDs (defines cube layout) #[serde(default)] pub id: Vec, - /// Size of each dimension #[serde(default)] pub size: Vec, - /// One entry per dimension, keyed by dimension code #[serde(flatten)] pub dimensions: HashMap, } -/// A single dimension in the JSON-stat2 response. #[derive(Debug, Clone, Deserialize, Serialize)] pub struct Dimension { - /// Display label #[serde(default)] pub label: String, - /// Category info with index and label maps pub category: CategoryInfo, } -/// Category info containing index and label maps. #[derive(Debug, Clone, Deserialize, Serialize)] pub struct CategoryInfo { - /// Maps category code → numeric position in the dimension pub index: HashMap, - /// Maps category code → display label #[serde(default)] pub label: HashMap, } -// --------------------------------------------------------------------------- -// Internal data models -// --------------------------------------------------------------------------- - -/// Decoded row of landing data with labeled dimensions. #[derive(Debug, Clone)] pub struct DataRow { pub month: String, @@ -223,7 +141,6 @@ pub struct DataRow { pub value: Option, } -/// Structured representation of a single landing record for DB insertion. #[derive(Debug, Clone)] pub struct Landing { pub month: String, @@ -238,11 +155,71 @@ pub struct Landing { pub value: Option, } -// --------------------------------------------------------------------------- -// Errors -// --------------------------------------------------------------------------- +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct SpeciesDto { + pub code: String, + pub label: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct LandingDto { + pub month: String, + pub species_code: String, + pub species_label: String, + pub gear_code: String, + pub zone_code: String, + pub processing_code: String, + pub preservation_code: String, + pub shipsize_code: String, + pub measure_code: String, + pub value: Option, +} + +impl From for LandingDto { + fn from(l: Landing) -> Self { + Self { + month: l.month, + species_code: l.species_code, + species_label: l.species_label, + gear_code: l.gear_code, + zone_code: l.zone_code, + processing_code: l.processing_code, + preservation_code: l.preservation_code, + shipsize_code: l.shipsize_code, + measure_code: l.measure_code, + value: l.value, + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MonthlyAggregate { + pub month: String, + pub total_mass: Option, + pub total_value: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct TopSpecies { + pub species_code: String, + pub species_label: String, + pub total_value: Option, + pub total_mass: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct PriceTrend { + pub month: String, + pub price_per_kg: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct SummaryDto { + pub monthly: Vec, + pub top_species: Vec, + pub price_trend: Vec, +} -/// Error types for ingestion module. #[derive(Debug, thiserror::Error)] pub enum IngestError { #[error("HTTP request failed: {0}")] @@ -267,12 +244,6 @@ pub enum IngestError { UnicodeError(String), } -// --------------------------------------------------------------------------- -// Type aliases -// --------------------------------------------------------------------------- - -/// Lookup map: dimension_code → (value_code → display_label). pub type LookupMap = HashMap>; -/// Result alias using custom error type. pub type Result = std::result::Result; diff --git a/static/index.html b/static/index.html new file mode 100644 index 0000000..0fac637 --- /dev/null +++ b/static/index.html @@ -0,0 +1 @@ +hagfish

hagfish

Frontend coming in Phase 4.