Files
hagfish/src/ingest.rs
T
2026-08-16 22:22:59 +01:00

742 lines
24 KiB
Rust

//! 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, warn};
/// 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)";
const DIM_ZONE: &str = "Economic Zone (GEONOM2023)";
const DIM_PROCESSING: &str = "Processing (EUMOFAPresentation)";
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; 2] = [-1.0, f64::NAN];
/// 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);
let resp = client
.get(url)
.header("Accept", "application/json")
.send()
.await?;
if !resp.status().is_success() {
let body = resp.text().await.unwrap_or_default();
return Err(IngestError::HttpError(reqwest::Error::from(
std::io::Error::new(
std::io::ErrorKind::Other,
format!("API returned error: {}", body),
),
)));
}
let meta: MetadataResponse = resp.json().await?;
debug!("Received {} variables from metadata", meta.variables.len());
for var in &meta.variables {
info!(
"Variable: code='{}' text='{}' role={:?} values={}",
var.code,
var.text,
var.role,
var.values.len()
);
}
let mut lookup_map: LookupMap = HashMap::with_capacity(meta.variables.len());
for var in &meta.variables {
if var.values.is_empty() {
debug!("Variable '{}' has no values", var.code);
continue;
}
let lookup = var.lookup_map();
info!("Loaded {} codes for '{}'", lookup.len(), var.code);
lookup_map.insert(var.code.clone(), lookup);
}
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()
.find(|v| v.code == DIM_MONTH)
.map(|v| v.values.clone())
.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(),
"build_query requires at least one month"
);
Query {
query: vec![
QueryItem {
code: DIM_MONTH.to_string(),
selection: Selection {
filter: "item".to_string(),
values: all_months.to_vec(),
},
},
QueryItem {
code: DIM_SPECIES.to_string(),
selection: Selection {
filter: "all".to_string(),
values: vec!["*".to_string()],
},
},
QueryItem {
code: DIM_GEAR.to_string(),
selection: Selection {
filter: "all".to_string(),
values: vec!["*".to_string()],
},
},
QueryItem {
code: DIM_ZONE.to_string(),
selection: Selection {
filter: "all".to_string(),
values: vec!["*".to_string()],
},
},
QueryItem {
code: DIM_PROCESSING.to_string(),
selection: Selection {
filter: "all".to_string(),
values: vec!["*".to_string()],
},
},
QueryItem {
code: DIM_PRESERVATION.to_string(),
selection: Selection {
filter: "all".to_string(),
values: vec!["*".to_string()],
},
},
QueryItem {
code: DIM_SHIPSIZE.to_string(),
selection: Selection {
filter: "all".to_string(),
values: vec!["*".to_string()],
},
},
QueryItem {
code: DIM_MEASURE.to_string(),
selection: Selection {
filter: "item".to_string(),
values: vec!["MASS".to_string(), "VALUE".to_string()],
},
},
],
response: QueryResponse::default(),
}
}
/// 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
.iter()
.find(|q| q.code == DIM_MONTH)
.map(|q| q.selection.values.len())
.unwrap_or(0);
info!("Sending data query for {} month(s)", month_count);
let resp = client.post(url).json(query).send().await?;
if !resp.status().is_success() {
let body = resp.text().await.unwrap_or_default();
warn!("API error response: {}", body);
return Err(IngestError::HttpError(reqwest::Error::from(
std::io::Error::new(
std::io::ErrorKind::Other,
format!("API returned error: {}", body),
),
)));
}
let data: DataResponse = resp.json().await?;
if data.dataset.value.is_empty() {
return Err(IngestError::EmptyDataset);
}
info!("Retrieved {} data points", data.dataset.value.len());
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;
for &size in key_sizes.iter().rev() {
indices.push(remaining % size);
remaining /= size;
}
indices.reverse();
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,
lookup_maps: &LookupMap,
) -> Result<DataRow> {
let dim_info = &dataset.dataset.dimension;
let dim_order = &dim_info.id;
if dim_order.len() != 8 {
return Err(IngestError::MissingDimension(format!(
"Expected 8 dimensions, got {}. Dimensions: {:?}",
dim_order.len(),
dim_order
)));
}
if row_index >= dataset.dataset.value.len() {
return Err(IngestError::InvalidValueCode(format!(
"Row index {} out of bounds (max {})",
row_index,
dataset.dataset.value.len()
)));
}
// 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| {
let dim = dim_info.dimensions.get(dim_code).ok_or_else(|| {
IngestError::MissingDimension(format!(
"Dimension '{}' not found in response",
dim_code
))
})?;
let mut entries: Vec<(String, usize)> = dim
.category
.index
.iter()
.map(|(code, &pos)| (code.clone(), pos))
.collect();
entries.sort_by_key(|(_, pos)| *pos);
let list: Vec<(String, String)> = entries
.into_iter()
.map(|(code, _)| {
let label = dim.category.label.get(&code).cloned().unwrap_or_else(|| {
lookup_maps
.get(dim_code)
.and_then(|m| m.get(&code))
.cloned()
.unwrap_or_else(|| code.clone())
});
(code, label)
})
.collect();
Ok(list)
})
.collect::<Result<Vec<_>>>()?;
let key_sizes: Vec<usize> = category_lists.iter().map(|c| c.len()).collect();
let indices = decode_key_indices(row_index, &key_sizes);
let dimension_values: Vec<(String, String)> = indices
.into_iter()
.zip(category_lists.iter())
.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());
}
let get = |code: &str| -> (String, String) {
dim_map
.get(code)
.cloned()
.unwrap_or_else(|| ("".to_string(), "".to_string()))
};
let (month_code, month_label) = get(DIM_MONTH);
let (species_code, species_label) = get(DIM_SPECIES);
let (gear_code, gear_label) = get(DIM_GEAR);
let (zone_code, zone_label) = get(DIM_ZONE);
let (processing_code, processing_label) = get(DIM_PROCESSING);
let (preservation_code, preservation_label) = get(DIM_PRESERVATION);
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,
other => other,
};
debug!(
"Row {}: month={} species={} ({}) measure={} value={:?}",
row_index, month_code, species_code, species_label, measure_code, value
);
Ok(DataRow {
month: month_code,
species_code,
species_label,
gear_code,
gear_label,
zone_code,
zone_label,
processing_code,
processing_label,
preservation_code,
preservation_label,
shipsize_code,
shipsize_label,
measure_code,
measure_label,
value,
})
}
/// Converts a `DataRow` to a `Landing` struct for database insertion.
pub fn data_row_to_landing(row: &DataRow) -> Landing {
Landing {
month: row.month.clone(),
species_code: row.species_code.clone(),
species_label: row.species_label.clone(),
gear_code: row.gear_code.clone(),
zone_code: row.zone_code.clone(),
processing_code: row.processing_code.clone(),
preservation_code: row.preservation_code.clone(),
shipsize_code: row.shipsize_code.clone(),
measure_code: row.measure_code.clone(),
value: row.value,
}
}
#[cfg(test)]
mod tests {
use super::*;
fn mock_dataset_response() -> DataResponse {
let mut dimensions = HashMap::new();
dimensions.insert(
DIM_MONTH.to_string(),
Dimension {
label: "Mánaður".to_string(),
category: CategoryInfo {
index: HashMap::from([("2015M01".to_string(), 0), ("2015M02".to_string(), 1)]),
label: HashMap::from([
("2015M01".to_string(), "Januar 2015".to_string()),
("2015M02".to_string(), "Februar 2015".to_string()),
]),
},
},
);
dimensions.insert(
DIM_SPECIES.to_string(),
Dimension {
label: "Fiskaslag".to_string(),
category: CategoryInfo {
index: HashMap::from([
("148XXXXXXX00000".to_string(), 0),
("183XXXXXXX00000".to_string(), 1),
]),
label: HashMap::from([
("148XXXXXXX00000".to_string(), "Sild".to_string()),
("183XXXXXXX00000".to_string(), "Þorskur".to_string()),
]),
},
},
);
for dim_code in [
DIM_GEAR,
DIM_ZONE,
DIM_PROCESSING,
DIM_PRESERVATION,
DIM_SHIPSIZE,
] {
dimensions.insert(
dim_code.to_string(),
Dimension {
label: dim_code.to_string(),
category: CategoryInfo {
index: HashMap::from([("TOTAL".to_string(), 0)]),
label: HashMap::from([("TOTAL".to_string(), "Total".to_string())]),
},
},
);
}
dimensions.insert(
DIM_MEASURE.to_string(),
Dimension {
label: "Mát".to_string(),
category: CategoryInfo {
index: HashMap::from([("MASS".to_string(), 0), ("VALUE".to_string(), 1)]),
label: HashMap::from([
("MASS".to_string(), "Nøgd".to_string()),
("VALUE".to_string(), "Virði".to_string()),
]),
},
},
);
// Dimension order: measure, species, gear, zone, processing,
// preservation, shipsize, month
// Sizes: 2, 2, 1, 1, 1, 1, 1, 2 = 8 values
DataResponse {
dataset: Dataset {
dimension: DimInfo {
id: vec![
DIM_MEASURE.to_string(),
DIM_SPECIES.to_string(),
DIM_GEAR.to_string(),
DIM_ZONE.to_string(),
DIM_PROCESSING.to_string(),
DIM_PRESERVATION.to_string(),
DIM_SHIPSIZE.to_string(),
DIM_MONTH.to_string(),
],
size: vec![2, 2, 1, 1, 1, 1, 1, 2],
dimensions,
},
value: vec![
Some(1234.5), // MASS, Sild, ..., 2015M01
Some(2345.6), // MASS, Sild, ..., 2015M02
Some(-1.0), // MASS, Þorskur, ..., 2015M01 (sentinel)
Some(3456.7), // MASS, Þorskur, ..., 2015M02
Some(4567.8), // VALUE, Sild, ..., 2015M01
Some(5678.9), // VALUE, Sild, ..., 2015M02
None, // VALUE, Þorskur, ..., 2015M01 (null)
Some(6789.0), // VALUE, Þorskur, ..., 2015M02
],
status: vec![],
},
}
}
fn mock_lookup_maps() -> LookupMap {
let mut maps = LookupMap::new();
maps.insert(
DIM_MONTH.to_string(),
HashMap::from([
("2015M01".to_string(), "Januar 2015".to_string()),
("2015M02".to_string(), "Februar 2015".to_string()),
]),
);
maps.insert(
DIM_SPECIES.to_string(),
HashMap::from([
("148XXXXXXX00000".to_string(), "Sild".to_string()),
("183XXXXXXX00000".to_string(), "Þorskur".to_string()),
]),
);
maps.insert(
DIM_MEASURE.to_string(),
HashMap::from([
("MASS".to_string(), "Nøgd".to_string()),
("VALUE".to_string(), "Virði".to_string()),
]),
);
maps
}
#[test]
fn test_build_query_uses_correct_dimension_codes() {
let months = vec!["2015M01".to_string(), "2015M02".to_string()];
let query = build_query(&months);
assert_eq!(query.response.format, "json-stat2");
assert_eq!(query.query.len(), 8);
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[1].code, DIM_SPECIES);
assert_eq!(query.query[1].selection.filter, "all");
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_decode_key_indices_basic() {
let key_sizes = vec![2, 2, 2];
assert_eq!(decode_key_indices(0, &key_sizes), vec![0, 0, 0]);
assert_eq!(decode_key_indices(7, &key_sizes), vec![1, 1, 1]);
assert_eq!(decode_key_indices(3, &key_sizes), vec![0, 1, 1]);
}
#[test]
fn test_decode_key_indices_single_dim() {
let key_sizes = vec![5];
assert_eq!(decode_key_indices(0, &key_sizes), vec![0]);
assert_eq!(decode_key_indices(4, &key_sizes), vec![4]);
}
#[test]
fn test_parse_row_faroese_unicode() {
let dataset = mock_dataset_response();
let lookup_maps = mock_lookup_maps();
let row = parse_row(0, &dataset, &lookup_maps).expect("parse failed");
assert_eq!(row.month, "2015M01");
assert_eq!(row.species_code, "148XXXXXXX00000");
assert_eq!(row.species_label, "Sild");
assert_eq!(row.measure_code, "MASS");
assert_eq!(row.measure_label, "Nøgd");
assert_eq!(row.value, Some(1234.5));
}
#[test]
fn test_parse_row_thorn_character() {
let dataset = mock_dataset_response();
let lookup_maps = mock_lookup_maps();
// Row 2: MASS, Þorskur, ..., 2015M01
let row = parse_row(2, &dataset, &lookup_maps).expect("parse failed");
assert_eq!(row.species_code, "183XXXXXXX00000");
assert_eq!(row.species_label, "Þorskur");
assert!(row.species_label.contains('Þ'));
}
#[test]
fn test_parse_row_sentinel_coerced_to_none() {
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!(
row.value.is_none(),
"Sentinel -1.0 must be None, got {:?}",
row.value
);
}
#[test]
fn test_parse_row_explicit_null() {
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());
}
#[test]
fn test_parse_row_out_of_bounds() {
let dataset = mock_dataset_response();
let lookup_maps = mock_lookup_maps();
let result = parse_row(100, &dataset, &lookup_maps);
assert!(result.is_err());
match result {
Err(IngestError::InvalidValueCode(msg)) => assert!(msg.contains("out of bounds")),
_ => panic!("Expected InvalidValueCode"),
}
}
#[test]
fn test_parse_row_wrong_dimension_count() {
let mut dataset = mock_dataset_response();
dataset.dataset.dimension.id.pop();
dataset.dataset.dimension.size.pop();
let lookup_maps = mock_lookup_maps();
let result = parse_row(0, &dataset, &lookup_maps);
assert!(result.is_err());
}
#[test]
fn test_data_row_to_landing() {
let dataset = mock_dataset_response();
let lookup_maps = mock_lookup_maps();
let row = parse_row(0, &dataset, &lookup_maps).expect("parse failed");
let landing = data_row_to_landing(&row);
assert_eq!(landing.month, row.month);
assert_eq!(landing.species_code, row.species_code);
assert_eq!(landing.species_label, row.species_label);
assert_eq!(landing.measure_code, row.measure_code);
assert_eq!(landing.value, row.value);
}
#[test]
fn test_faroese_unicode_all_rows() {
let dataset = mock_dataset_response();
let lookup_maps = mock_lookup_maps();
for i in 0..8 {
let row = parse_row(i, &dataset, &lookup_maps);
assert!(row.is_ok(), "Row {} failed", i);
if let Ok(r) = row {
assert!(
!r.species_label.contains('\u{FFFD}'),
"Replacement char in row {}",
i
);
}
}
}
#[test]
fn test_is_sentinel() {
assert!(is_sentinel(-1.0));
assert!(is_sentinel(f64::NAN));
assert!(!is_sentinel(0.0));
assert!(!is_sentinel(1234.5));
assert!(!is_sentinel(0.001));
}
#[test]
fn test_extract_available_months() {
let meta = MetadataResponse {
title: Some("test".to_string()),
variables: vec![VariableMeta {
code: DIM_MONTH.to_string(),
text: "Mánaður".to_string(),
role: None,
values: vec!["2015M01".to_string(), "2015M02".to_string()],
value_texts: vec!["Jan 2015".to_string(), "Feb 2015".to_string()],
extra: HashMap::new(),
}],
};
let months = extract_available_months(&meta);
assert_eq!(months.len(), 2);
assert!(months.contains(&"2015M01".to_string()));
assert!(months.contains(&"2015M02".to_string()));
}
#[test]
fn test_extract_available_months_no_month_var() {
let meta = MetadataResponse {
title: None,
variables: vec![VariableMeta {
code: DIM_SPECIES.to_string(),
text: "Fiskaslag".to_string(),
role: None,
values: vec!["TOTAL".to_string()],
value_texts: vec!["Tils.".to_string()],
extra: HashMap::new(),
}],
};
let months = extract_available_months(&meta);
assert!(months.is_empty());
}
#[test]
#[should_panic(expected = "requires at least one month")]
fn test_build_query_empty_panics() {
let _ = build_query(&[]);
}
#[test]
fn test_variable_meta_lookup_map() {
let var = VariableMeta {
code: "measure".to_string(),
text: "mát".to_string(),
role: None,
values: vec!["MASS".to_string(), "VALUE".to_string()],
value_texts: vec!["Nøgd".to_string(), "Virði".to_string()],
extra: HashMap::new(),
};
let map = var.lookup_map();
assert_eq!(map.get("MASS"), Some(&"Nøgd".to_string()));
assert_eq!(map.get("VALUE"), Some(&"Virði".to_string()));
assert_eq!(map.len(), 2);
}
}