use crate::db; use crate::types::{ AvailableFiltersDto, Config, LandingDto, LookupDto, MonthlyAggregate, MonthlySpeciesBreakdown, PriceTrend, SpeciesDto, SummaryDto, TopSpecies, }; use axum::{ Json, Router, body::Body, extract::{Query, State}, http::{HeaderValue, Method, StatusCode, header}, response::{IntoResponse, Response}, routing::get, }; use rust_embed::Embed; use serde::Deserialize; use std::sync::Arc; use tokio::sync::Mutex; use tower_http::cors::{AllowOrigin, 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 month_from: Option, pub month_to: 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 } const MAX_LIMIT: u32 = 10000; const LEAF_ONLY_FILTER: &str = "species_code != 'TOTAL' AND gear_code != 'TOTAL' AND zone_code != 'TOTAL'"; #[derive(Debug, Clone, Deserialize, Default)] pub struct SummaryQuery { pub species: Option, pub species_multi: Option, pub month: Option, pub month_from: Option, pub month_to: Option, pub zone: Option, pub zone_multi: Option, pub gear: Option, pub gear_multi: Option, } impl SummaryQuery { fn parse_multi(value: &Option) -> Vec { match value { Some(v) if !v.is_empty() => { v.split(',').map(|s| s.trim().to_string()).filter(|s| !s.is_empty()).collect() } _ => Vec::new(), } } fn build_conditions(&self) -> (Vec, Vec>) { let mut conditions: Vec = Vec::new(); let mut args: Vec> = Vec::new(); conditions.push(format!("({LEAF_ONLY_FILTER})")); if let Some(ref month) = self.month { conditions.push("month = ?".to_string()); args.push(Box::new(month.clone())); } if let Some(ref month_from) = self.month_from { conditions.push("month >= ?".to_string()); args.push(Box::new(month_from.clone())); } if let Some(ref month_to) = self.month_to { conditions.push("month <= ?".to_string()); args.push(Box::new(month_to.clone())); } let species_list = Self::parse_multi(&self.species_multi); if !species_list.is_empty() { let placeholders: Vec = (0..species_list.len()) .map(|_| "?".to_string()) .collect(); conditions.push(format!("species_code IN ({})", placeholders.join(", "))); for s in &species_list { args.push(Box::new(s.clone())); } } else if let Some(ref species) = self.species { conditions.push("species_code = ?".to_string()); args.push(Box::new(species.clone())); } let zone_list = Self::parse_multi(&self.zone_multi); if !zone_list.is_empty() { let placeholders: Vec = (0..zone_list.len()) .map(|_| "?".to_string()) .collect(); conditions.push(format!("zone_code IN ({})", placeholders.join(", "))); for z in &zone_list { args.push(Box::new(z.clone())); } } else if let Some(ref zone) = self.zone { conditions.push("zone_code = ?".to_string()); args.push(Box::new(zone.clone())); } let gear_list = Self::parse_multi(&self.gear_multi); if !gear_list.is_empty() { let placeholders: Vec = (0..gear_list.len()) .map(|_| "?".to_string()) .collect(); conditions.push(format!("gear_code IN ({})", placeholders.join(", "))); for g in &gear_list { args.push(Box::new(g.clone())); } } else if let Some(ref gear) = self.gear { conditions.push("gear_code = ?".to_string()); args.push(Box::new(gear.clone())); } (conditions, args) } fn build_where_and(&self) -> (String, Vec>) { let (conditions, args) = self.build_conditions(); let clause = if conditions.is_empty() { String::new() } else { format!(" AND {}", conditions.join(" AND ")) }; (clause, args) } fn build_where_prefix(&self) -> (String, Vec>) { let (conditions, args) = self.build_conditions(); let clause = if conditions.is_empty() { String::new() } else { format!(" WHERE {}", conditions.join(" AND ")) }; (clause, args) } } #[derive(Debug, Clone, Deserialize)] pub struct AvailableFiltersQuery { pub month_from: Option, pub month_to: Option, pub species_multi: Option, pub zone_multi: Option, pub gear_multi: Option, } impl AvailableFiltersQuery { fn parse_multi(value: &Option) -> Vec { match value { Some(v) if !v.is_empty() => { v.split(',').map(|s| s.trim().to_string()).filter(|s| !s.is_empty()).collect() } _ => Vec::new(), } } fn build_conditions_excluding(&self, exclude: &str) -> (Vec, Vec>) { let mut conditions: Vec = Vec::new(); let mut args: Vec> = Vec::new(); conditions.push(format!("({LEAF_ONLY_FILTER})")); if let Some(ref month_from) = self.month_from { conditions.push("month >= ?".to_string()); args.push(Box::new(month_from.clone())); } if let Some(ref month_to) = self.month_to { conditions.push("month <= ?".to_string()); args.push(Box::new(month_to.clone())); } if exclude != "species" { let species_list = Self::parse_multi(&self.species_multi); if !species_list.is_empty() { let placeholders: Vec = species_list.iter().map(|_| "?".to_string()).collect(); conditions.push(format!("species_code IN ({})", placeholders.join(", "))); for s in &species_list { args.push(Box::new(s.clone())); } } } if exclude != "zone" { let zone_list = Self::parse_multi(&self.zone_multi); if !zone_list.is_empty() { let placeholders: Vec = zone_list.iter().map(|_| "?".to_string()).collect(); conditions.push(format!("zone_code IN ({})", placeholders.join(", "))); for z in &zone_list { args.push(Box::new(z.clone())); } } } if exclude != "gear" { let gear_list = Self::parse_multi(&self.gear_multi); if !gear_list.is_empty() { let placeholders: Vec = gear_list.iter().map(|_| "?".to_string()).collect(); conditions.push(format!("gear_code IN ({})", placeholders.join(", "))); for g in &gear_list { args.push(Box::new(g.clone())); } } } (conditions, args) } fn build_where_prefix_excluding(&self, exclude: &str) -> (String, Vec>) { let (conditions, args) = self.build_conditions_excluding(exclude); let clause = if conditions.is_empty() { String::new() } else { format!(" WHERE {}", conditions.join(" AND ")) }; (clause, args) } } #[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 { let cors = if state.config.allowed_origins.is_empty() { tracing::warn!("CORS is permissive — no allowed_origins configured"); CorsLayer::permissive() } else { let origins: Vec = state .config .allowed_origins .iter() .filter_map(|s| s.parse().ok()) .collect(); if origins.len() < state.config.allowed_origins.len() { tracing::warn!( "{} origin(s) failed to parse (non-ASCII?): {:?}", state.config.allowed_origins.len() - origins.len(), state.config.allowed_origins ); } if origins.is_empty() { tracing::warn!("No valid origins configured — falling back to permissive CORS"); CorsLayer::permissive() } else { CorsLayer::new() .allow_methods([Method::GET]) .allow_headers([header::CONTENT_TYPE]) .allow_origin(AllowOrigin::list(origins)) } }; Router::new() .route("/healthz", get(healthz)) .route("/api/species", get(get_species)) .route("/api/zones", get(get_zones)) .route("/api/gear", get(get_gear)) .route("/api/available-filters", get(get_available_filters)) .route("/api/landings", get(get_landings)) .route("/api/summary", get(get_summary)) .route("/api/summary/monthly-breakdown", get(get_monthly_breakdown)) .route("/api/export.parquet", get(export_parquet)) .fallback(static_handler) .layer(cors) .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 WHERE code != 'TOTAL' 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 capped_limit = std::cmp::min(params.limit, MAX_LIMIT); 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 \ AND species_code != 'TOTAL' AND gear_code != 'TOTAL' AND zone_code != 'TOTAL'", ); 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 month_from) = params.month_from { sql.push_str(&format!(" AND month >= ${idx}")); args.push(Box::new(month_from.clone())); idx += 1; } if let Some(ref month_to) = params.month_to { sql.push_str(&format!(" AND month <= ${idx}")); args.push(Box::new(month_to.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(capped_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 result = tokio::task::spawn_blocking(move || -> db::Result { let conn = conn.blocking_lock(); let (where_clause, args) = params.build_where_and(); let arg_refs: Vec<&dyn duckdb::ToSql> = args.iter().map(|b| b.as_ref()).collect(); 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 1=1{where_clause} \ GROUP BY month ORDER BY month" ); let mut stmt = conn.prepare(&monthly_sql)?; let monthly: Vec = stmt .query_map(arg_refs.as_slice(), |row| { Ok(MonthlyAggregate { month: row.get(0)?, total_mass: row.get(1)?, total_value: row.get(2)?, }) })? .collect::, duckdb::Error>>()?; 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 1=1{where_clause} \ GROUP BY species_code, species_label \ ORDER BY total_value DESC NULLS LAST \ LIMIT 10" ); let mut stmt = conn.prepare(&top_sql)?; let top_species: Vec = stmt .query_map(arg_refs.as_slice(), |row| { Ok(TopSpecies { species_code: row.get(0)?, species_label: row.get(1)?, total_value: row.get(2)?, total_mass: row.get(3)?, }) })? .collect::, duckdb::Error>>()?; 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 mut price_args: Vec<&dyn duckdb::ToSql> = Vec::with_capacity(args.len() * 2); for a in &args { price_args.push(a.as_ref()); } for a in &args { price_args.push(a.as_ref()); } let mut stmt = conn.prepare(&price_sql)?; let price_trend: Vec = stmt .query_map(price_args.as_slice(), |row| { Ok(PriceTrend { month: row.get(0)?, price_per_kg: row.get(1)?, }) })? .collect::, duckdb::Error>>()?; 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 get_monthly_breakdown( State(state): State, Query(params): Query, ) -> ApiResult>> { let conn = state.conn.clone(); let result = tokio::task::spawn_blocking(move || -> db::Result> { let conn = conn.blocking_lock(); let (where_clause, args) = params.build_where_prefix(); let sql = format!( "SELECT month, species_code, MAX(species_label) 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 month, species_code \ ORDER BY month ASC, total_value DESC NULLS LAST" ); 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(MonthlySpeciesBreakdown { month: row.get(0)?, species_code: row.get(1)?, species_label: row.get(2)?, total_value: row.get(3)?, total_mass: row.get(4)?, }) })? .collect::, duckdb::Error>>()?; Ok(rows) }) .await .map_err(|e| ApiError { status: StatusCode::INTERNAL_SERVER_ERROR, message: format!("Task join error: {e}"), })??; Ok(Json(result)) } async fn get_available_filters( State(state): State, Query(params): Query, ) -> ApiResult> { let conn = state.conn.clone(); let result = tokio::task::spawn_blocking(move || -> db::Result { let conn = conn.blocking_lock(); let (species_clause, species_args) = params.build_where_prefix_excluding("species"); let species_arg_refs: Vec<&dyn duckdb::ToSql> = species_args.iter().map(|b| b.as_ref()).collect(); let species_sql = if species_clause.is_empty() { "SELECT DISTINCT species_code FROM landings WHERE species_code != 'TOTAL' ORDER BY species_code" } else { "SELECT DISTINCT species_code FROM landings WHERE species_code != 'TOTAL'" }; let species_sql = if species_clause.is_empty() { species_sql.to_string() } else { format!("SELECT DISTINCT species_code FROM landings{species_clause} AND species_code != 'TOTAL' ORDER BY species_code") }; let mut stmt = conn.prepare(&species_sql)?; let species: Vec = stmt .query_map(species_arg_refs.as_slice(), |row| row.get(0))? .collect::, duckdb::Error>>()?; let (zone_clause, zone_args) = params.build_where_prefix_excluding("zone"); let zone_arg_refs: Vec<&dyn duckdb::ToSql> = zone_args.iter().map(|b| b.as_ref()).collect(); let zone_sql = if zone_clause.is_empty() { "SELECT DISTINCT zone_code FROM landings ORDER BY zone_code".to_string() } else { format!("SELECT DISTINCT zone_code FROM landings{zone_clause} ORDER BY zone_code") }; let mut stmt = conn.prepare(&zone_sql)?; let zones: Vec = stmt .query_map(zone_arg_refs.as_slice(), |row| row.get(0))? .collect::, duckdb::Error>>()?; let (gear_clause, gear_args) = params.build_where_prefix_excluding("gear"); let gear_arg_refs: Vec<&dyn duckdb::ToSql> = gear_args.iter().map(|b| b.as_ref()).collect(); let gear_sql = if gear_clause.is_empty() { "SELECT DISTINCT gear_code FROM landings ORDER BY gear_code".to_string() } else { format!("SELECT DISTINCT gear_code FROM landings{gear_clause} ORDER BY gear_code") }; let mut stmt = conn.prepare(&gear_sql)?; let gear: Vec = stmt .query_map(gear_arg_refs.as_slice(), |row| row.get(0))? .collect::, duckdb::Error>>()?; Ok(AvailableFiltersDto { species, zones, gear, }) }) .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 = tokio::task::spawn_blocking(move || -> ApiResult { let conn = conn.blocking_lock(); let tmp = tempfile::Builder::new() .prefix("hagfish-") .suffix(".parquet") .tempfile() .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)?; drop(conn); Ok(path_str) }) .await .map_err(|e| ApiError { status: StatusCode::INTERNAL_SERVER_ERROR, message: format!("Task join error: {e}"), })??; let file = tokio::fs::File::open(&file_path) .await .map_err(|e| ApiError { status: StatusCode::INTERNAL_SERVER_ERROR, message: format!("Failed to open exported file: {e}"), })?; let stream = tokio_util::io::ReaderStream::new(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); } }); Response::builder() .status(StatusCode::OK) .header(header::CONTENT_TYPE, "application/octet-stream") .header( header::CONTENT_DISPOSITION, "attachment; filename=\"landings.parquet\"", ) .body(body) .map_err(|e| ApiError { status: StatusCode::INTERNAL_SERVER_ERROR, message: format!("Failed to build response: {e}"), }) } async fn static_handler(uri: axum::http::Uri) -> Response { let path = uri.path().trim_start_matches('/'); let asset_path = if let Some(stripped) = path.strip_prefix("static/") { stripped } else { path }; let asset_path = if asset_path.is_empty() { "index.html" } else { asset_path }; let asset = StaticAssets::get(asset_path).or_else(|| StaticAssets::get("index.html")); match asset { Some(file) => { let mime = mime_guess::from_path(asset_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_or_else(|_| { Response::builder() .status(StatusCode::INTERNAL_SERVER_ERROR) .body(Body::from("Internal server error")) .unwrap() }) } None => Response::builder() .status(StatusCode::NOT_FOUND) .body(Body::from("Not found")) .unwrap(), } } async fn get_zones(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 zone WHERE code != 'TOTAL' ORDER BY code")?; let rows = stmt.query_map([], |row| { Ok(LookupDto { 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_gear(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 gear WHERE code != 'TOTAL' ORDER BY code")?; let rows = stmt.query_map([], |row| { Ok(LookupDto { 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)) } #[cfg(test)] mod tests { use super::*; use crate::types::Landing; use duckdb::Connection; use std::collections::HashMap; use std::collections::HashSet; 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: "TR1".to_string(), zone_code: "FO".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: "TR1".to_string(), zone_code: "FO".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: "TR1".to_string(), zone_code: "FO".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: "TR1".to_string(), zone_code: "FO".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: "TR1".to_string(), zone_code: "FO".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: "TR1".to_string(), zone_code: "FO".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 { allowed_origins: vec!["https://hagfisk.poc.xn--fl-5ka.fo".to_string()], ..Default::default() }), } } fn empty_test_state() -> AppState { let conn = Connection::open_in_memory().unwrap(); db::init_schema(&conn).unwrap(); AppState { conn: Arc::new(Mutex::new(conn)), config: Arc::new(Config { allowed_origins: vec![], ..Default::default() }), } } async fn spawn_test_server(state: AppState) -> String { 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(); }); format!("http://{addr}") } #[tokio::test] async fn test_healthz_returns_ok() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/healthz")).await.unwrap(); assert_eq!(resp.status(), StatusCode::OK); } #[tokio::test] async fn test_get_species_returns_all() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/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 base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/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 base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/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 base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/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_landings_range_filter() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!( "{base}/api/landings?month_from=2024M01&month_to=2024M01" )) .await .unwrap(); assert_eq!(resp.status(), StatusCode::OK); let body: Vec = resp.json().await.unwrap(); assert!(body.iter().all(|l| l.month == "2024M01")); assert_eq!(body.len(), 4); let resp = reqwest::get(format!( "{base}/api/landings?month_from=2024M01&month_to=2024M02" )) .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_summary_monthly_aggregates() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/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 base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/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 base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/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_summary_with_species_filter() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/api/summary?species=COD")) .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() - 1000.0).abs() < f64::EPSILON); assert!((jan.total_value.unwrap() - 5000.0).abs() < f64::EPSILON); assert!(!body.top_species.is_empty()); for ts in &body.top_species { assert_eq!(ts.species_code, "COD"); } assert_eq!(body.price_trend.len(), 2); for pt in &body.price_trend { assert!(pt.price_per_kg.is_some()); } } #[tokio::test] async fn test_get_summary_with_multi_species_filter() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/api/summary?species_multi=COD,HER")) .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); } #[tokio::test] async fn test_get_summary_with_zone_filter() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/api/summary?zone=FO")) .await .unwrap(); assert_eq!(resp.status(), StatusCode::OK); let body: SummaryDto = resp.json().await.unwrap(); assert!(!body.monthly.is_empty()); } #[tokio::test] async fn test_get_summary_with_gear_filter() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/api/summary?gear=TR1")) .await .unwrap(); assert_eq!(resp.status(), StatusCode::OK); let body: SummaryDto = resp.json().await.unwrap(); assert!(!body.monthly.is_empty()); } #[tokio::test] async fn test_get_summary_with_month_range_filter() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!( "{base}/api/summary?month_from=2024M01&month_to=2024M01" )) .await .unwrap(); assert_eq!(resp.status(), StatusCode::OK); let body: SummaryDto = resp.json().await.unwrap(); assert_eq!(body.monthly.len(), 1); assert_eq!(body.monthly[0].month, "2024M01"); } #[tokio::test] async fn test_get_summary_with_combined_filters() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!( "{base}/api/summary?month_from=2024M01&month_to=2024M01&species=COD" )) .await .unwrap(); assert_eq!(resp.status(), StatusCode::OK); let body: SummaryDto = resp.json().await.unwrap(); assert_eq!(body.monthly.len(), 1); assert_eq!(body.monthly[0].month, "2024M01"); assert!((body.monthly[0].total_mass.unwrap() - 1000.0).abs() < f64::EPSILON); } #[tokio::test] async fn test_get_landings_faroese_label_preserved() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/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 base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/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 base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/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 base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/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_export_parquet_empty_db() { let state = empty_test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/api/export.parquet")) .await .unwrap(); assert_eq!(resp.status(), StatusCode::OK); let bytes = resp.bytes().await.unwrap(); assert_eq!(&bytes[..4], b"PAR1"); assert!(bytes.len() > 4); } #[tokio::test] async fn test_concurrent_export_requests() { let state = test_state(); let base = spawn_test_server(state).await; let url = format!("{base}/api/export.parquet"); let (resp1, resp2) = tokio::join!(reqwest::get(&url), reqwest::get(&url)); let resp1 = resp1.unwrap(); let resp2 = resp2.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!(!bytes1.is_empty()); assert!(!bytes2.is_empty()); } #[tokio::test] async fn test_static_handler_falls_back_to_index_html() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/nonexistent-path-xyz")) .await .unwrap(); assert_eq!(resp.status(), StatusCode::OK); let body = resp.text().await.unwrap(); assert!(body.contains("Fisheries Dashboard")); } #[tokio::test] async fn test_sql_injection_safe() { let state = test_state(); let base = spawn_test_server(state.clone()).await; let malicious_params = "?species='; DROP TABLE landings; --"; let url = format!("{base}/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.conn.blocking_lock(); conn.query_row("SELECT COUNT(*) FROM landings", [], |row| row.get(0)) .unwrap() }) .await .unwrap(); assert_eq!(count, 6); } #[tokio::test] async fn test_sql_injection_all_params_safe() { let state = test_state(); let base = spawn_test_server(state.clone()).await; let malicious_params = "?month_from=2024M01'; DROP TABLE landings;--&month_to=2024M12&species=COD"; let url = format!("{base}/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.conn.blocking_lock(); conn.query_row("SELECT COUNT(*) FROM landings", [], |row| row.get(0)) .unwrap() }) .await .unwrap(); assert_eq!(count, 6); } #[tokio::test] async fn test_get_monthly_breakdown_returns_species_by_month() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/api/summary/monthly-breakdown")) .await .unwrap(); assert_eq!(resp.status(), StatusCode::OK); let body: Vec = resp.json().await.unwrap(); assert!(!body.is_empty(), "Should have breakdown rows"); let months: HashSet<_> = body.iter().map(|r| r.month.as_str()).collect(); let species: HashSet<_> = body.iter().map(|r| r.species_code.as_str()).collect(); assert!(!months.is_empty()); assert!(!species.is_empty()); for row in &body { assert!(!row.species_code.is_empty()); assert!(!row.species_label.is_empty()); assert!(row.total_value.is_some() || row.total_mass.is_some()); } } #[tokio::test] async fn test_get_monthly_breakdown_filtered_by_month_range() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!( "{base}/api/summary/monthly-breakdown?month_from=2024M01&month_to=2024M01" )) .await .unwrap(); assert_eq!(resp.status(), StatusCode::OK); let body: Vec = resp.json().await.unwrap(); for row in &body { assert_eq!(row.month, "2024M01"); } } #[tokio::test] async fn test_get_monthly_breakdown_filtered_by_species() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/api/summary/monthly-breakdown?species=COD")) .await .unwrap(); assert_eq!(resp.status(), StatusCode::OK); let body: Vec = resp.json().await.unwrap(); for row in &body { assert_eq!(row.species_code, "COD"); assert_eq!(row.species_label, "Toskur"); } } #[tokio::test] async fn test_get_available_filters_no_constraints() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/api/available-filters")) .await .unwrap(); assert_eq!(resp.status(), StatusCode::OK); let body: AvailableFiltersDto = resp.json().await.unwrap(); assert!(body.species.contains(&"COD".to_string())); assert!(body.zones.contains(&"FO".to_string())); assert!(body.gear.contains(&"TR1".to_string())); } #[tokio::test] async fn test_get_available_filters_with_species_constraint() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/api/available-filters?species_multi=COD")) .await .unwrap(); assert_eq!(resp.status(), StatusCode::OK); let body: AvailableFiltersDto = resp.json().await.unwrap(); assert!(body.species.contains(&"COD".to_string())); assert!(body.zones.contains(&"FO".to_string())); assert!(body.gear.contains(&"TR1".to_string())); } #[tokio::test] async fn test_get_available_filters_excludes_unrelated_species() { let state = test_state(); let base = spawn_test_server(state).await; let resp = reqwest::get(format!("{base}/api/available-filters?zone_multi=FO")) .await .unwrap(); assert_eq!(resp.status(), StatusCode::OK); let body: AvailableFiltersDto = resp.json().await.unwrap(); assert!(body.species.contains(&"COD".to_string())); assert!(body.species.contains(&"HER".to_string())); } }