From f9cfb99916225352013073beba0e52113f898aa9 Mon Sep 17 00:00:00 2001 From: Bartal Laearsson Date: Tue, 25 Aug 2026 23:33:54 +0100 Subject: [PATCH] add data validation --- justfile | 64 +++++++++++------------- src/api.rs | 56 ++++++++++++--------- src/ingest.rs | 134 +++++++++++++++++++++++++++++++------------------- src/main.rs | 29 +++++++++-- 4 files changed, 170 insertions(+), 113 deletions(-) diff --git a/justfile b/justfile index 411e124..04cd5cf 100644 --- a/justfile +++ b/justfile @@ -6,7 +6,19 @@ build: echo "Building hagfisk binary..." cargo build --release -# Sync binary to server +# Stop app service (must stop before DB operations) +stop-app: + ssh {{host}} "systemctl --user stop hagfisk.service" + +# Start app service +start-app: + ssh {{host}} "systemctl --user start hagfisk.service" + +# Restart app service +restart-app: + ssh {{host}} "systemctl --user restart hagfisk.service" + +# Deploy and restart (safe for regular updates) deploy: build ssh -q {{host}} "rm -f {{hagfisk_dir}}/hagfisk" echo "Syncing binary to remote..." @@ -16,50 +28,34 @@ deploy: build ssh -q {{host}} "systemctl --user daemon-reload" ssh -q {{host}} "systemctl --user restart hagfisk.service" -# Enable and start service -enable-now-app: - ssh {{host}} "systemctl --user enable --now hagfisk.service" - -# Start service on remote -start-app: +# Deploy with full re-ingest (stops service first) +deploy-reingest: build stop-app + echo "Syncing binary to remote..." + rsync -avz --progress ./target/release/hagfisk {{host}}:{{hagfisk_dir}}/hagfisk + echo "Restoring SELinux contexts..." + ssh -t -q {{host}} "sudo restorecon -RF {{hagfisk_dir}}" + ssh -q {{host}} "systemctl --user daemon-reload" + echo "Running full re-ingest..." + ssh {{host}} "cd {{hagfisk_dir}} && ./hagfisk ingest --full" + echo "Starting service..." ssh {{host}} "systemctl --user start hagfisk.service" -# Stop service on remote -stop-app: - ssh {{host}} "systemctl --user stop hagfisk.service" +# Manual full ingest (assumes service is stopped) +ingest-full: + ssh {{host}} "cd {{hagfisk_dir}} && ./hagfisk ingest --full" -# Restart service on remote -restart-app: - ssh {{host}} "systemctl --user restart hagfisk.service" - -# Stream logs from remote +# Stream logs from remote app logs-app: ssh {{host}} "journalctl --user -u hagfisk -f" -# Show status of hagfisk +# Show status of app ps-app: ssh {{host}} "systemctl --user status hagfisk.service" -# Enable and start ingest service and timer -enable-now-ingest: - ssh {{host}} "systemctl --user enable --now hagfisk-ingest.service hagfisk-ingest.timer" - -# Start ingest on remote -start-ingest: - ssh {{host}} "systemctl --user start hagfisk-ingest.service hagfisk-ingest.timer" - -# Stop ingest on remote -stop-ingest: - ssh {{host}} "systemctl --user stop hagfisk-ingest.service hagfisk-ingest.timer" - -# Restart ingest on remote -restart-ingest: - ssh {{host}} "systemctl --user restart hagfisk-ingest.service hagfisk-ingest.timer" - -# Stream ingest logs from remote +# Stream logs from remote ingest (manual run) logs-ingest: ssh {{host}} "journalctl --user -u hagfisk-ingest -f" -# Show status of hagfisk ingest +# Show status of ingest ps-ingest: ssh {{host}} "systemctl --user status hagfisk-ingest.service" diff --git a/src/api.rs b/src/api.rs index ca5cab9..3c4049b 100644 --- a/src/api.rs +++ b/src/api.rs @@ -48,6 +48,9 @@ fn default_limit() -> u32 { 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, @@ -75,6 +78,8 @@ impl SummaryQuery { 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())); @@ -177,6 +182,8 @@ impl AvailableFiltersQuery { 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())); @@ -323,7 +330,7 @@ async fn get_species(State(state): State) -> ApiResult db::Result> { let conn = conn.blocking_lock(); - let mut stmt = conn.prepare("SELECT code, label FROM species ORDER BY code")?; + 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)?, @@ -354,7 +361,8 @@ async fn get_landings( 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", + 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; @@ -468,7 +476,7 @@ async fn get_summary( 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} AND species_code != 'TOTAL' \ + FROM landings WHERE 1=1{where_clause} \ GROUP BY species_code, species_label \ ORDER BY total_value DESC NULLS LAST \ LIMIT 10" @@ -774,7 +782,7 @@ async fn get_zones(State(state): State) -> ApiResult db::Result> { let conn = conn.blocking_lock(); - let mut stmt = conn.prepare("SELECT code, label FROM zone ORDER BY code")?; + 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)?, @@ -796,7 +804,7 @@ async fn get_gear(State(state): State) -> ApiResult db::Result> { let conn = conn.blocking_lock(); - let mut stmt = conn.prepare("SELECT code, label FROM gear ORDER BY code")?; + 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)?, @@ -831,8 +839,8 @@ mod tests { month: "2024M01".to_string(), species_code: "COD".to_string(), species_label: "Toskur".to_string(), - gear_code: "TOTAL".to_string(), - zone_code: "TOTAL".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(), @@ -843,8 +851,8 @@ mod tests { month: "2024M01".to_string(), species_code: "COD".to_string(), species_label: "Toskur".to_string(), - gear_code: "TOTAL".to_string(), - zone_code: "TOTAL".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(), @@ -855,8 +863,8 @@ mod tests { month: "2024M01".to_string(), species_code: "HER".to_string(), species_label: "Sild".to_string(), - gear_code: "TOTAL".to_string(), - zone_code: "TOTAL".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(), @@ -867,8 +875,8 @@ mod tests { month: "2024M01".to_string(), species_code: "HER".to_string(), species_label: "Sild".to_string(), - gear_code: "TOTAL".to_string(), - zone_code: "TOTAL".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(), @@ -879,8 +887,8 @@ mod tests { month: "2024M02".to_string(), species_code: "COD".to_string(), species_label: "Toskur".to_string(), - gear_code: "TOTAL".to_string(), - zone_code: "TOTAL".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(), @@ -891,8 +899,8 @@ mod tests { month: "2024M02".to_string(), species_code: "COD".to_string(), species_label: "Toskur".to_string(), - gear_code: "TOTAL".to_string(), - zone_code: "TOTAL".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(), @@ -1153,7 +1161,7 @@ mod tests { let state = test_state(); let base = spawn_test_server(state).await; - let resp = reqwest::get(format!("{base}/api/summary?zone=TOTAL")) + let resp = reqwest::get(format!("{base}/api/summary?zone=FO")) .await .unwrap(); assert_eq!(resp.status(), StatusCode::OK); @@ -1167,7 +1175,7 @@ mod tests { let state = test_state(); let base = spawn_test_server(state).await; - let resp = reqwest::get(format!("{base}/api/summary?gear=TOTAL")) + let resp = reqwest::get(format!("{base}/api/summary?gear=TR1")) .await .unwrap(); assert_eq!(resp.status(), StatusCode::OK); @@ -1448,8 +1456,8 @@ mod tests { let body: AvailableFiltersDto = resp.json().await.unwrap(); assert!(body.species.contains(&"COD".to_string())); - assert!(body.zones.contains(&"TOTAL".to_string())); - assert!(body.gear.contains(&"TOTAL".to_string())); + assert!(body.zones.contains(&"FO".to_string())); + assert!(body.gear.contains(&"TR1".to_string())); } #[tokio::test] @@ -1464,8 +1472,8 @@ mod tests { let body: AvailableFiltersDto = resp.json().await.unwrap(); assert!(body.species.contains(&"COD".to_string())); - assert!(body.zones.contains(&"TOTAL".to_string())); - assert!(body.gear.contains(&"TOTAL".to_string())); + assert!(body.zones.contains(&"FO".to_string())); + assert!(body.gear.contains(&"TR1".to_string())); } #[tokio::test] @@ -1473,7 +1481,7 @@ mod tests { let state = test_state(); let base = spawn_test_server(state).await; - let resp = reqwest::get(format!("{base}/api/available-filters?zone_multi=TOTAL")) + let resp = reqwest::get(format!("{base}/api/available-filters?zone_multi=FO")) .await .unwrap(); assert_eq!(resp.status(), StatusCode::OK); diff --git a/src/ingest.rs b/src/ingest.rs index 4abf640..f4d8635 100644 --- a/src/ingest.rs +++ b/src/ingest.rs @@ -4,9 +4,9 @@ use std::collections::HashMap; use tracing::{debug, info}; const DIM_MONTH: &str = "month"; -const DIM_SPECIES: &str = "Species (ASFIS2022)"; -const DIM_GEAR: &str = "Fishing Gear (ISSCFG2016)"; -const DIM_ZONE: &str = "Economic Zone (GEONOM2023)"; +pub const DIM_SPECIES: &str = "Species (ASFIS2022)"; +pub const DIM_GEAR: &str = "Fishing Gear (ISSCFG2016)"; +pub const DIM_ZONE: &str = "Economic Zone (GEONOM2023)"; const DIM_PROCESSING: &str = "Processing (EUMOFAPresentation)"; const DIM_PRESERVATION: &str = "Preservation (EUMOFAPreservation)"; const DIM_SHIPSIZE: &str = "Shipsize"; @@ -70,6 +70,20 @@ pub fn extract_available_months(meta: &MetadataResponse) -> Vec { .unwrap_or_default() } +pub fn extract_non_total_values(meta: &MetadataResponse, dim_code: &str) -> Vec { + meta.variables + .iter() + .find(|v| v.code == dim_code) + .map(|v| { + v.values + .iter() + .filter(|val| **val != "TOTAL") + .cloned() + .collect() + }) + .unwrap_or_default() +} + pub fn chunk_months(months: &[String], batch_size: usize) -> Vec> { if months.is_empty() { return vec![]; @@ -81,8 +95,13 @@ pub fn chunk_months(months: &[String], batch_size: usize) -> Vec> { .collect() } -pub fn build_query(all_months: &[String]) -> Result { - if all_months.is_empty() { +pub fn build_query( + months: &[String], + species_codes: &[String], + gear_codes: &[String], + zone_codes: &[String], +) -> Result { + if months.is_empty() { return Err(IngestError::InvalidValueCode( "build_query requires at least one month".to_string(), )); @@ -94,28 +113,28 @@ pub fn build_query(all_months: &[String]) -> Result { code: DIM_MONTH.to_string(), selection: Selection { filter: "item".to_string(), - values: all_months.to_vec(), + values: months.to_vec(), }, }, QueryItem { code: DIM_SPECIES.to_string(), selection: Selection { - filter: "all".to_string(), - values: vec!["*".to_string()], + filter: "item".to_string(), + values: species_codes.to_vec(), }, }, QueryItem { code: DIM_GEAR.to_string(), selection: Selection { - filter: "all".to_string(), - values: vec!["*".to_string()], + filter: "item".to_string(), + values: gear_codes.to_vec(), }, }, QueryItem { code: DIM_ZONE.to_string(), selection: Selection { - filter: "all".to_string(), - values: vec!["*".to_string()], + filter: "item".to_string(), + values: zone_codes.to_vec(), }, }, QueryItem { @@ -463,53 +482,66 @@ mod tests { maps } - #[test] - fn test_build_query_uses_correct_dimension_codes() { - let months = vec!["2015M01".to_string(), "2015M02".to_string()]; - let query = build_query(&months).unwrap(); +#[test] +fn test_build_query_uses_correct_dimension_codes() { + let months = vec!["2015M01".to_string(), "2015M02".to_string()]; + let species = vec!["COD".to_string(), "HER".to_string()]; + let gear = vec!["TR1".to_string(), "GN".to_string()]; + let zones = vec!["FO".to_string(), "INTL".to_string()]; - assert_eq!(query.response.format, "json-stat2"); - assert_eq!(query.query.len(), 8); + let query = build_query(&months, &species, &gear, &zones).unwrap(); - assert_eq!(query.query[0].code, DIM_MONTH); - assert_eq!(query.query[0].selection.filter, "item"); - assert_eq!(query.query[0].selection.values.len(), 2); + assert_eq!(query.response.format, "json-stat2"); + assert_eq!(query.query.len(), 8); - assert_eq!(query.query[1].code, DIM_SPECIES); - assert_eq!(query.query[1].selection.filter, "all"); + assert_eq!(query.query[0].code, DIM_MONTH); + assert_eq!(query.query[0].selection.filter, "item"); + assert_eq!(query.query[0].selection.values.len(), 2); - assert_eq!(query.query[4].code, DIM_PROCESSING); - assert_eq!(query.query[4].selection.filter, "item"); - assert_eq!(query.query[4].selection.values, vec!["TOTAL"]); + assert_eq!(query.query[1].code, DIM_SPECIES); + assert_eq!(query.query[1].selection.filter, "item"); + assert_eq!(query.query[1].selection.values, vec!["COD", "HER"]); - assert_eq!(query.query[5].code, DIM_PRESERVATION); - assert_eq!(query.query[5].selection.filter, "item"); - assert_eq!(query.query[5].selection.values, vec!["TOTAL"]); + assert_eq!(query.query[2].code, DIM_GEAR); + assert_eq!(query.query[2].selection.filter, "item"); + assert_eq!(query.query[2].selection.values, vec!["TR1", "GN"]); - assert_eq!(query.query[6].code, DIM_SHIPSIZE); - assert_eq!(query.query[6].selection.filter, "item"); - assert_eq!(query.query[6].selection.values, vec!["TOTAL"]); + assert_eq!(query.query[3].code, DIM_ZONE); + assert_eq!(query.query[3].selection.filter, "item"); + assert_eq!(query.query[3].selection.values, vec!["FO", "INTL"]); - assert_eq!(query.query[7].code, DIM_MEASURE); - assert!( - query.query[7] - .selection - .values - .contains(&"MASS".to_string()) - ); - assert!( - query.query[7] - .selection - .values - .contains(&"VALUE".to_string()) - ); - } + assert_eq!(query.query[4].code, DIM_PROCESSING); + assert_eq!(query.query[4].selection.filter, "item"); + assert_eq!(query.query[4].selection.values, vec!["TOTAL"]); - #[test] - fn test_build_query_returns_error_on_empty() { - let result = build_query(&[]); - assert!(result.is_err()); - } + assert_eq!(query.query[5].code, DIM_PRESERVATION); + assert_eq!(query.query[5].selection.filter, "item"); + assert_eq!(query.query[5].selection.values, vec!["TOTAL"]); + + assert_eq!(query.query[6].code, DIM_SHIPSIZE); + assert_eq!(query.query[6].selection.filter, "item"); + assert_eq!(query.query[6].selection.values, vec!["TOTAL"]); + + assert_eq!(query.query[7].code, DIM_MEASURE); + assert!( + query.query[7] + .selection + .values + .contains(&"MASS".to_string()) + ); + assert!( + query.query[7] + .selection + .values + .contains(&"VALUE".to_string()) + ); +} + +#[test] +fn test_build_query_returns_error_on_empty() { + let result = build_query(&[], &[], &[], &[]); + assert!(result.is_err()); +} #[test] fn test_decode_key_indices_basic() { diff --git a/src/main.rs b/src/main.rs index 3ad87ee..90fac75 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,4 +1,3 @@ -// src/main.rs mod api; mod cli; mod db; @@ -150,8 +149,24 @@ async fn run_ingest(config: &types::Config, full: bool) -> anyhow::Result<()> { db::update_lookups(&conn, &lookup_maps)?; + let species_codes = ingest::extract_non_total_values(&metadata, ingest::DIM_SPECIES); + let gear_codes = ingest::extract_non_total_values(&metadata, ingest::DIM_GEAR); + let zone_codes = ingest::extract_non_total_values(&metadata, ingest::DIM_ZONE); + + tracing::info!( + species_count = species_codes.len(), + gear_count = gear_codes.len(), + zone_count = zone_codes.len(), + "Extracted non-TOTAL dimension values for leaf-only query" + ); + let all_months = ingest::extract_available_months(&metadata); + if full { + conn.execute("DELETE FROM landings", [])?; + tracing::info!("Cleared landings table for full re-ingest"); + } + let months_to_fetch: Vec = if full { all_months } else { @@ -203,7 +218,7 @@ async fn run_ingest(config: &types::Config, full: bool) -> anyhow::Result<()> { "fetching batch" ); - let query = ingest::build_query(batch)?; + let query = ingest::build_query(batch, &species_codes, &gear_codes, &zone_codes)?; let data = match ingest::fetch_data(&client, &config.data_source_url, &query).await { Ok(d) => d, @@ -225,6 +240,12 @@ async fn run_ingest(config: &types::Config, full: bool) -> anyhow::Result<()> { for j in 0..data.dataset.value.len() { match ingest::parse_row(j, &data.dataset, &lookup_maps) { Ok(row) => { + if row.species_code == "TOTAL" + || row.gear_code == "TOTAL" + || row.zone_code == "TOTAL" + { + continue; + } landings.push(ingest::data_row_to_landing(row)); } Err(e) => { @@ -331,8 +352,8 @@ mod tests { month: "2024M01".to_string(), species_code: "COD".to_string(), species_label: "Toskur".to_string(), - gear_code: "TOTAL".to_string(), - zone_code: "TOTAL".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(),