phase 3 complete
This commit is contained in:
+866
@@ -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<Mutex<duckdb::Connection>>,
|
||||
#[allow(dead_code)]
|
||||
pub config: Arc<Config>,
|
||||
}
|
||||
|
||||
#[derive(Embed)]
|
||||
#[folder = "static/"]
|
||||
struct StaticAssets;
|
||||
|
||||
#[derive(Debug, Clone, Deserialize)]
|
||||
pub struct LandingsQuery {
|
||||
pub month: Option<String>,
|
||||
pub species: Option<String>,
|
||||
pub gear: Option<String>,
|
||||
pub zone: Option<String>,
|
||||
pub measure: Option<String>,
|
||||
#[serde(default = "default_limit")]
|
||||
pub limit: u32,
|
||||
}
|
||||
|
||||
fn default_limit() -> u32 {
|
||||
10000
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Deserialize)]
|
||||
pub struct SummaryQuery {
|
||||
pub species: Option<String>,
|
||||
}
|
||||
|
||||
#[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<duckdb::Error> 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<db::DbError> 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<T> = std::result::Result<T, ApiError>;
|
||||
|
||||
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<AppState>) -> impl IntoResponse {
|
||||
StatusCode::OK
|
||||
}
|
||||
|
||||
async fn get_species(State(state): State<AppState>) -> ApiResult<Json<Vec<SpeciesDto>>> {
|
||||
let conn = state.conn.clone();
|
||||
let rows = tokio::task::spawn_blocking(move || -> db::Result<Vec<SpeciesDto>> {
|
||||
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::<std::result::Result<Vec<_>, 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<AppState>,
|
||||
Query(params): Query<LandingsQuery>,
|
||||
) -> ApiResult<Json<Vec<LandingDto>>> {
|
||||
let conn = state.conn.clone();
|
||||
let rows = tokio::task::spawn_blocking(move || -> db::Result<Vec<LandingDto>> {
|
||||
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<Box<dyn duckdb::ToSql>> = 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::<std::result::Result<Vec<_>, 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<AppState>,
|
||||
Query(params): Query<SummaryQuery>,
|
||||
) -> ApiResult<Json<SummaryDto>> {
|
||||
let conn = state.conn.clone();
|
||||
let species_filter = params.species.clone();
|
||||
|
||||
let result = tokio::task::spawn_blocking(move || -> db::Result<SummaryDto> {
|
||||
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<MonthlyAggregate> = 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::<Vec<MonthlyAggregate>, duckdb::Error>(
|
||||
rows.collect::<std::result::Result<Vec<_>, _>>()?,
|
||||
)?
|
||||
} 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::<Vec<MonthlyAggregate>, duckdb::Error>(
|
||||
rows.collect::<std::result::Result<Vec<_>, _>>()?,
|
||||
)?
|
||||
};
|
||||
|
||||
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<TopSpecies> = 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::<Vec<TopSpecies>, duckdb::Error>(rows.collect::<std::result::Result<Vec<_>, _>>()?)?
|
||||
} 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::<Vec<TopSpecies>, duckdb::Error>(rows.collect::<std::result::Result<Vec<_>, _>>()?)?
|
||||
};
|
||||
|
||||
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<PriceTrend> = 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::<Vec<PriceTrend>, duckdb::Error>(rows.collect::<std::result::Result<Vec<_>, _>>()?)?
|
||||
} 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::<Vec<PriceTrend>, duckdb::Error>(rows.collect::<std::result::Result<Vec<_>, _>>()?)?
|
||||
};
|
||||
|
||||
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<AppState>) -> ApiResult<Response> {
|
||||
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<SpeciesDto> = 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<LandingDto> = 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<LandingDto> = 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<LandingDto> = 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<LandingDto> = 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<LandingDto> = 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<LandingDto> = 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);
|
||||
}
|
||||
}
|
||||
@@ -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<T> = std::result::Result<T, DbError>;
|
||||
|
||||
/// Opens a DuckDB connection at the given path and initializes the schema.
|
||||
pub fn init(path: &str) -> Result<Connection> {
|
||||
let conn = Connection::open(path)?;
|
||||
init_schema(&conn)?;
|
||||
@@ -39,7 +30,6 @@ pub fn init(path: &str) -> Result<Connection> {
|
||||
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<usize> {
|
||||
if rows.is_empty() {
|
||||
return Ok(0);
|
||||
@@ -154,7 +131,6 @@ pub fn upsert_landings(conn: &Connection, rows: &[Landing]) -> Result<usize> {
|
||||
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<usize> {
|
||||
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<Option<String>> {
|
||||
let result = conn.query_row("SELECT MAX(month) FROM landings", [], |row| {
|
||||
let val: Option<String> = row.get(0)?;
|
||||
@@ -182,12 +152,8 @@ pub fn get_last_month(conn: &Connection) -> Result<Option<String>> {
|
||||
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(),
|
||||
|
||||
+4
-60
@@ -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<LookupMap> {
|
||||
info!("Fetching metadata from {}", url);
|
||||
|
||||
@@ -89,9 +66,6 @@ pub async fn fetch_metadata(client: &Client, url: &str) -> Result<LookupMap> {
|
||||
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<String> {
|
||||
meta.variables
|
||||
.iter()
|
||||
@@ -100,13 +74,6 @@ pub fn extract_available_months(meta: &MetadataResponse) -> Vec<String> {
|
||||
.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<DataResponse> {
|
||||
let month_count = query
|
||||
.query
|
||||
@@ -206,9 +172,6 @@ pub async fn fetch_data(client: &Client, url: &str, query: &Query) -> Result<Dat
|
||||
Ok(data)
|
||||
}
|
||||
|
||||
/// Decodes a flat row-major index into per-dimension indices.
|
||||
///
|
||||
/// JSON-stat2 uses row-major order: the last dimension varies fastest.
|
||||
fn decode_key_indices(flat_index: usize, key_sizes: &[usize]) -> Vec<usize> {
|
||||
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<usize> {
|
||||
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<Vec<(String, String)>> = 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());
|
||||
|
||||
+53
-2
@@ -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(())
|
||||
}
|
||||
|
||||
+64
-93
@@ -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<VariableMeta>,
|
||||
}
|
||||
|
||||
/// 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<String>,
|
||||
/// Code values for this variable (parallel with valueTexts)
|
||||
#[serde(default)]
|
||||
pub values: Vec<String>,
|
||||
/// Display labels for each value (parallel with values)
|
||||
#[serde(default, rename = "valueTexts")]
|
||||
pub value_texts: Vec<String>,
|
||||
/// Catch-all for unknown fields
|
||||
#[serde(flatten)]
|
||||
pub extra: HashMap<String, serde_json::Value>,
|
||||
}
|
||||
|
||||
impl VariableMeta {
|
||||
/// Build a code→label lookup map from the parallel arrays.
|
||||
pub fn lookup_map(&self) -> HashMap<String, String> {
|
||||
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<QueryItem>,
|
||||
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<String>,
|
||||
}
|
||||
|
||||
/// 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<Option<f64>>,
|
||||
/// Status codes per value (optional)
|
||||
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
||||
pub status: Vec<String>,
|
||||
}
|
||||
|
||||
/// 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<String>,
|
||||
/// Size of each dimension
|
||||
#[serde(default)]
|
||||
pub size: Vec<usize>,
|
||||
/// One entry per dimension, keyed by dimension code
|
||||
#[serde(flatten)]
|
||||
pub dimensions: HashMap<String, Dimension>,
|
||||
}
|
||||
|
||||
/// 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<String, usize>,
|
||||
/// Maps category code → display label
|
||||
#[serde(default)]
|
||||
pub label: HashMap<String, String>,
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// 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<f64>,
|
||||
}
|
||||
|
||||
/// 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<f64>,
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// 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<f64>,
|
||||
}
|
||||
|
||||
impl From<Landing> 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<f64>,
|
||||
pub total_value: Option<f64>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct TopSpecies {
|
||||
pub species_code: String,
|
||||
pub species_label: String,
|
||||
pub total_value: Option<f64>,
|
||||
pub total_mass: Option<f64>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct PriceTrend {
|
||||
pub month: String,
|
||||
pub price_per_kg: Option<f64>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct SummaryDto {
|
||||
pub monthly: Vec<MonthlyAggregate>,
|
||||
pub top_species: Vec<TopSpecies>,
|
||||
pub price_trend: Vec<PriceTrend>,
|
||||
}
|
||||
|
||||
/// 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<String, HashMap<String, String>>;
|
||||
|
||||
/// Result alias using custom error type.
|
||||
pub type Result<T> = std::result::Result<T, IngestError>;
|
||||
|
||||
Reference in New Issue
Block a user