add todo
This commit is contained in:
+509
@@ -0,0 +1,509 @@
|
||||
//! Data ingestion from PX-Web API.
|
||||
//!
|
||||
//! Handles metadata fetching, query construction, data retrieval, and
|
||||
//! parsing of JSON-stat2 responses into structured rows.
|
||||
|
||||
use crate::types::*;
|
||||
use reqwest::Client;
|
||||
use tracing::{debug, info, warn};
|
||||
use std::collections::HashMap;
|
||||
|
||||
const PX_WEB_LANGUAGE: &str = "fo"; // Faroese labels
|
||||
|
||||
/// Fetches metadata from PX-Web API endpoint.
|
||||
///
|
||||
/// Returns a HashMap mapping variable_id → code → label mappings.
|
||||
/// This lookup is used to decode opaque codes in data responses.
|
||||
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() {
|
||||
return Err(IngestError::HttpError(reqwest::Error::from(
|
||||
std::io::Error::new(
|
||||
std::io::ErrorKind::Other,
|
||||
format!("API returned status: {}", resp.status()),
|
||||
),
|
||||
)));
|
||||
}
|
||||
|
||||
let meta: MetadataResponse = resp.json().await?;
|
||||
debug!("Received {} variables from metadata", meta.variables.len());
|
||||
|
||||
// Build lookup map from metadata
|
||||
let mut lookup_map: LookupMap = HashMap::new();
|
||||
|
||||
for var in &meta.variables {
|
||||
let codes = meta.values.get(&var.id);
|
||||
if let Some(codes) = codes {
|
||||
let labels: HashMap<String, String> = codes
|
||||
.iter()
|
||||
.map(|vm| (vm.code.clone(), vm.text.clone()))
|
||||
.collect();
|
||||
lookup_map.insert(var.id.clone(), labels);
|
||||
info!("Loaded {} codes for variable '{}'", labels.len(), var.id);
|
||||
} else {
|
||||
warn!("No codes found for variable '{}'", var.id);
|
||||
}
|
||||
}
|
||||
|
||||
Ok(lookup_map)
|
||||
}
|
||||
|
||||
/// Constructs a query body for fetching all landing data.
|
||||
///
|
||||
/// Sets all categorical dimensions to "*" (all values) and measures to MASS+VALUE.
|
||||
/// Used for full backfill ingestion.
|
||||
pub fn build_query(all_months: &[String]) -> Query {
|
||||
Query {
|
||||
query: vec![
|
||||
QueryItem {
|
||||
id: "Tid".to_string(), // Month
|
||||
values: all_months.to_vec(),
|
||||
},
|
||||
QueryItem {
|
||||
id: "Art".to_string(), // Species
|
||||
values: vec!["*".to_string()],
|
||||
},
|
||||
QueryItem {
|
||||
id: "Redskab".to_string(), // Gear
|
||||
values: vec!["*".to_string()],
|
||||
},
|
||||
QueryItem {
|
||||
id: "Økonomisk zone".to_string(),
|
||||
values: vec!["*".to_string()],
|
||||
},
|
||||
QueryItem {
|
||||
id: "Tilstand".to_string(), // Processing
|
||||
values: vec!["TOTAL".to_string()],
|
||||
},
|
||||
QueryItem {
|
||||
id: "Konservering".to_string(), // Preservation
|
||||
values: vec!["TOTAL".to_string()],
|
||||
},
|
||||
QueryItem {
|
||||
id: "Skibsstørrelse".to_string(), // Vessel size
|
||||
values: vec!["TOTAL".to_string()],
|
||||
},
|
||||
QueryItem {
|
||||
id: "Måleenhed".to_string(), // Measure
|
||||
values: vec!["MASS".to_string(), "VALUE".to_string()],
|
||||
},
|
||||
],
|
||||
language: PX_WEB_LANGUAGE.to_string(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Fetches data from PX-Web API using the provided query.
|
||||
pub async fn fetch_data(client: &Client, url: &str, query: &Query) -> Result<DataResponse> {
|
||||
info!("Sending data query for {} month(s)", query.query[0].values.len());
|
||||
|
||||
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 status: {}", resp.status()),
|
||||
),
|
||||
)));
|
||||
}
|
||||
|
||||
let data: DataResponse = resp.json().await?;
|
||||
|
||||
if data.dataset.value.is_empty() {
|
||||
return Err(IngestError::EmptyDataset);
|
||||
}
|
||||
|
||||
info!(
|
||||
"Retrieved {} data points from API",
|
||||
data.dataset.value.len()
|
||||
);
|
||||
|
||||
Ok(data)
|
||||
}
|
||||
|
||||
/// Extracts dimension codes from a flat row index in JSON-stat2 format.
|
||||
///
|
||||
/// JSON-stat2 uses a flattened array where each element corresponds to a unique
|
||||
/// combination of dimension values. The dimension.keys array tells us how many
|
||||
/// values exist per dimension. We decode the flat index into per-dimension indices.
|
||||
fn decode_key_indices(flat_index: usize, key_counts: &[usize]) -> Vec<usize> {
|
||||
let mut indices = Vec::with_capacity(key_counts.len());
|
||||
let mut remaining = flat_index;
|
||||
|
||||
// Process dimensions in reverse (last dimension varies fastest)
|
||||
for count in key_counts.iter().rev() {
|
||||
indices.push((remaining % count) as usize);
|
||||
remaining /= count;
|
||||
}
|
||||
|
||||
indices.reverse(); // Restore original order
|
||||
indices
|
||||
}
|
||||
|
||||
/// Parses a single row from JSON-stat2 format into a DataRow.
|
||||
///
|
||||
/// # Arguments
|
||||
/// * `row_index` - Index into the dataset's value array
|
||||
/// * `dataset` - The complete dataset response
|
||||
/// * `lookup_maps` - Code → label mappings for each dimension
|
||||
///
|
||||
/// # Errors
|
||||
/// Returns an error if the row_index exceeds bounds or required lookups are missing.
|
||||
pub fn parse_row(
|
||||
row_index: usize,
|
||||
dataset: &DataResponse,
|
||||
lookup_maps: &LookupMap,
|
||||
) -> Result<DataRow> {
|
||||
let dim_info = &dataset.dataset.dimension;
|
||||
let categories = &dim_info.category;
|
||||
let keys = &dim_info.keys;
|
||||
let value = dataset.dataset.value.get(row_index).copied().flatten();
|
||||
|
||||
// Validate bounds
|
||||
if row_index >= dataset.dataset.value.len() {
|
||||
return Err(IngestError::InvalidValueCode(format!(
|
||||
"Row index {} out of bounds (max {})",
|
||||
row_index,
|
||||
dataset.dataset.value.len()
|
||||
)));
|
||||
}
|
||||
|
||||
// Get category label lists for each dimension (in key order)
|
||||
let category_lists: Vec<Vec<&str>> = keys
|
||||
.iter()
|
||||
.filter_map(|k| {
|
||||
categories.label.get(k).map(|labels| {
|
||||
labels.iter().map(|s| s.as_str()).collect::<Vec<_>>()
|
||||
})
|
||||
})
|
||||
.collect();
|
||||
|
||||
if category_lists.len() != keys.len() {
|
||||
return Err(IngestError::MissingDimension(
|
||||
format!(
|
||||
"Category count mismatch: {} keys vs {} category lists",
|
||||
keys.len(),
|
||||
category_lists.len()
|
||||
)
|
||||
));
|
||||
}
|
||||
|
||||
// Decode flat index into per-dimension indices
|
||||
let key_counts: Vec<usize> = category_lists.iter().map(|c| c.len()).collect();
|
||||
let indices = decode_key_indices(row_index, &key_counts);
|
||||
|
||||
// Map indices back to actual dimension values
|
||||
let dimension_values: Vec<&str> = indices
|
||||
.into_iter()
|
||||
.zip(category_lists.iter())
|
||||
.map(|(idx, list)| list[idx])
|
||||
.collect();
|
||||
|
||||
// Expected 8 dimensions: [month, species, gear, zone, processing, preservation, shipsize, measure]
|
||||
const EXPECTED_DIMS: usize = 8;
|
||||
if dimension_values.len() != EXPECTED_DIMS {
|
||||
return Err(IngestError::MissingDimension(format!(
|
||||
"Expected {} dimensions, got {}",
|
||||
EXPECTED_DIMS,
|
||||
dimension_values.len()
|
||||
)));
|
||||
}
|
||||
|
||||
let [month_code, species_code, gear_code, zone_code, processing_code, preservation_code, shipsize_code, measure_code]: [&str; 8] =
|
||||
dimension_values.try_into().unwrap();
|
||||
|
||||
// Helper to safely get label from lookup map
|
||||
fn get_label<'a>(maps: &'a LookupMap, var_name: &str, code: &str) -> &'a str {
|
||||
maps.get(var_name)
|
||||
.and_then(|m| m.get(code))
|
||||
.map(|s| s.as_str())
|
||||
.unwrap_or(code)
|
||||
}
|
||||
|
||||
let data_row = DataRow {
|
||||
month: month_code.to_string(),
|
||||
species_code: species_code.to_string(),
|
||||
species_label: get_label(lookup_maps, "Art", species_code).to_string(),
|
||||
gear_code: gear_code.to_string(),
|
||||
gear_label: get_label(lookup_maps, "Redskab", gear_code).to_string(),
|
||||
zone_code: zone_code.to_string(),
|
||||
zone_label: get_label(lookup_maps, "Økonomisk zone", zone_code).to_string(),
|
||||
processing_code: processing_code.to_string(),
|
||||
processing_label: get_label(lookup_maps, "Tilstand", processing_code).to_string(),
|
||||
preservation_code: preservation_code.to_string(),
|
||||
preservation_label: get_label(lookup_maps, "Konservering", preservation_code).to_string(),
|
||||
shipsize_code: shipsize_code.to_string(),
|
||||
shipsize_label: get_label(lookup_maps, "Skibsstørrelse", shipsize_code).to_string(),
|
||||
measure_code: measure_code.to_string(),
|
||||
measure_label: get_label(lookup_maps, "Måleenhed", measure_code).to_string(),
|
||||
value, // Already handled None case above
|
||||
};
|
||||
|
||||
Ok(data_row)
|
||||
}
|
||||
|
||||
/// Converts DataRow to Landing struct for database insertion.
|
||||
/// Removes redundant fields that aren't needed in the fact table.
|
||||
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,
|
||||
}
|
||||
}
|
||||
|
||||
/// Collects all months available in the metadata response.
|
||||
/// Used for building incremental ingestion queries.
|
||||
pub fn extract_available_months(meta: &MetadataResponse) -> Vec<String> {
|
||||
meta.values
|
||||
.get("Tid")
|
||||
.map(|codes| codes.iter().map(|c| c.code.clone()).collect())
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use std::sync::OnceLock;
|
||||
|
||||
static CLIENT: OnceLock<Client> = OnceLock::new();
|
||||
|
||||
fn client() -> &'static Client {
|
||||
CLIENT.get_or_init(Client::new)
|
||||
}
|
||||
|
||||
/// Test fixture: mock JSON-stat2 response with Faroese Unicode characters.
|
||||
/// Simulates a response with 2 months × 2 species × 2 measures = 8 values.
|
||||
fn mock_dataset_response() -> DataResponse {
|
||||
DataResponse {
|
||||
dataset: Dataset {
|
||||
// Keys specify cardinality per dimension (2 months, 2 species, ..., 2 measures)
|
||||
keys: vec![
|
||||
"2".to_string(), // Time (2 months)
|
||||
"2".to_string(), // Art (2 species)
|
||||
"1".to_string(), // Redskab (1 gear - TOTAL)
|
||||
"1".to_string(), // Zone (1 zone - TOTAL)
|
||||
"1".to_string(), // Tilstand (1 - TOTAL)
|
||||
"1".to_string(), // Konservering (1 - TOTAL)
|
||||
"1".to_string(), // Skibsstørrelse (1 - TOTAL)
|
||||
"2".to_string(), // Måleenhed (2 - MASS, VALUE)
|
||||
],
|
||||
category: CategoryInfo {
|
||||
label: HashMap::from([
|
||||
(
|
||||
"0".to_string(),
|
||||
vec!["2015M01".to_string(), "2015M02".to_string()],
|
||||
),
|
||||
(
|
||||
"1".to_string(),
|
||||
vec!["Sild".to_string(), "Þorskur".to_string()],
|
||||
),
|
||||
("2".to_string(), vec!["TOTAL".to_string()]),
|
||||
("3".to_string(), vec!["TOTAL".to_string()]),
|
||||
("4".to_string(), vec!["TOTAL".to_string()]),
|
||||
("5".to_string(), vec!["TOTAL".to_string()]),
|
||||
("6".to_string(), vec!["TOTAL".to_string()]),
|
||||
(
|
||||
"7".to_string(),
|
||||
vec!["MASS".to_string(), "VALUE".to_string()],
|
||||
),
|
||||
]),
|
||||
index: None,
|
||||
},
|
||||
value: vec![
|
||||
Some(1234.5), // 2015M01, Sild, TOTAL..., MASS
|
||||
Some(2345.6), // 2015M01, Sild, TOTAL..., VALUE
|
||||
Some(-1.0), // 2015M01, Þorskur, TOTAL..., MASS (marker)
|
||||
Some(3456.7), // 2015M01, Þorskur, TOTAL..., VALUE
|
||||
Some(4567.8), // 2015M02, Sild, TOTAL..., MASS
|
||||
Some(5678.9), // 2015M02, Sild, TOTAL..., VALUE
|
||||
None, // 2015M02, Þorskur, TOTAL..., MASS (missing)
|
||||
Some(6789.0), // 2015M02, Þorskur, TOTAL..., VALUE
|
||||
],
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
/// Test fixture: mock lookup maps with Faroese labels.
|
||||
fn mock_lookup_maps() -> LookupMap {
|
||||
let mut maps = LookupMap::new();
|
||||
maps.insert(
|
||||
"0".to_string(),
|
||||
HashMap::from([
|
||||
("2015M01".to_string(), "Januar 2015".to_string()),
|
||||
("2015M02".to_string(), "Februar 2015".to_string()),
|
||||
]),
|
||||
);
|
||||
maps.insert(
|
||||
"1".to_string(),
|
||||
HashMap::from([
|
||||
("Sild".to_string(), "Sild".to_string()),
|
||||
("Þorskur".to_string(), "Þorskur".to_string()),
|
||||
]),
|
||||
);
|
||||
maps.insert(
|
||||
"7".to_string(),
|
||||
HashMap::from([
|
||||
("MASS".to_string(), "Kilo".to_string()),
|
||||
("VALUE".to_string(), "Krónur".to_string()),
|
||||
]),
|
||||
);
|
||||
maps
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_build_query_all_wildcards() {
|
||||
let months = vec!["2015M01".to_string(), "2015M02".to_string()];
|
||||
let query = build_query(&months);
|
||||
|
||||
assert_eq!(query.language, PX_WEB_LANGUAGE);
|
||||
assert_eq!(query.query.len(), 8); // All 8 dimensions
|
||||
|
||||
// Check that species is wildcard
|
||||
let species_query = query.query.iter().find(|q| q.id == "Art").unwrap();
|
||||
assert_eq!(species_query.values, vec!["*".to_string()]);
|
||||
|
||||
// Check that measures includes both MASS and VALUE
|
||||
let measure_query = query.query.iter().find(|q| q.id == "Måleenhed").unwrap();
|
||||
assert!(measure_query.values.contains(&"MASS".to_string()));
|
||||
assert!(measure_query.values.contains(&"VALUE".to_string()));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_decode_key_indices_basic() {
|
||||
// 2 × 2 × 2 = 8 combinations
|
||||
let key_counts = vec![2, 2, 2];
|
||||
|
||||
// Index 0 should give [0, 0, 0]
|
||||
let indices = decode_key_indices(0, &key_counts);
|
||||
assert_eq!(indices, vec![0, 0, 0]);
|
||||
|
||||
// Index 7 should give [1, 1, 1]
|
||||
let indices = decode_key_indices(7, &key_counts);
|
||||
assert_eq!(indices, vec![1, 1, 1]);
|
||||
|
||||
// Index 3 should give [0, 1, 1] (middle combination)
|
||||
let indices = decode_key_indices(3, &key_counts);
|
||||
assert_eq!(indices, vec![0, 1, 1]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_parse_row_faroese_unicode() {
|
||||
let dataset = mock_dataset_response();
|
||||
let lookup_maps = mock_lookup_maps();
|
||||
|
||||
// Parse first row (index 0)
|
||||
let row = parse_row(0, &dataset, &lookup_maps).expect("Failed to parse row");
|
||||
|
||||
// Verify Faroese characters survive intact
|
||||
assert_eq!(row.month, "2015M01");
|
||||
assert_eq!(row.species_code, "Sild");
|
||||
assert!(row.species_label.contains("Sild"));
|
||||
assert_eq!(row.measure_code, "MASS");
|
||||
|
||||
// Verify value is preserved
|
||||
assert!(row.value.is_some());
|
||||
assert_eq!(row.value.unwrap(), 1234.5);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_parse_row_second_species() {
|
||||
let dataset = mock_dataset_response();
|
||||
let lookup_maps = mock_lookup_maps();
|
||||
|
||||
// Row index 2 should be Þorskur (second species, first time, MASS)
|
||||
let row = parse_row(2, &dataset, &lookup_maps).expect("Failed to parse row");
|
||||
|
||||
// Verify the special character survives
|
||||
assert_eq!(row.species_code, "Þorskur");
|
||||
assert!(row.species_label.contains("Þ"));
|
||||
assert!(row.species_label.contains("orskur"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_parse_row_missing_value_handling() {
|
||||
let dataset = mock_dataset_response();
|
||||
let lookup_maps = mock_lookup_maps();
|
||||
|
||||
// Row index 6 has None value (missing data for 2015M02, Þorskur, MASS)
|
||||
let row = parse_row(6, &dataset, &lookup_maps).expect("Failed to parse row");
|
||||
|
||||
// Missing values should become None, not panic
|
||||
assert!(row.value.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_parse_row_out_of_bounds() {
|
||||
let dataset = mock_dataset_response();
|
||||
let lookup_maps = mock_lookup_maps();
|
||||
|
||||
// Requesting out-of-bounds index should return error
|
||||
let result = parse_row(100, &dataset, &lookup_maps);
|
||||
assert!(result.is_err());
|
||||
|
||||
if let Err(IngestError::InvalidValueCode(msg)) = result {
|
||||
assert!(msg.contains("out of bounds"));
|
||||
} else {
|
||||
panic!("Expected InvalidValueCode error");
|
||||
}
|
||||
}
|
||||
|
||||
#[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("Failed to parse row");
|
||||
let landing = data_row_to_landing(&row);
|
||||
|
||||
// Verify core fields are copied
|
||||
assert_eq!(landing.month, row.month);
|
||||
assert_eq!(landing.species_code, row.species_code);
|
||||
assert_eq!(landing.measure_code, row.measure_code);
|
||||
assert_eq!(landing.value, row.value);
|
||||
|
||||
// Verify lookup fields are preserved
|
||||
assert_eq!(landing.species_label, row.species_label);
|
||||
assert_eq!(landing.gear_code, row.gear_code);
|
||||
assert_eq!(landing.zone_code, row.zone_code);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_extract_available_months() {
|
||||
let meta = MetadataResponse {
|
||||
variables: vec![],
|
||||
values: HashMap::from([(
|
||||
"Tid".to_string(),
|
||||
vec![
|
||||
ValueMeta { code: "2015M01".to_string(), text: "Jan 2015".to_string() },
|
||||
ValueMeta { code: "2015M02".to_string(), text: "Feb 2015".to_string() },
|
||||
],
|
||||
)]),
|
||||
};
|
||||
|
||||
let months = extract_available_months(&meta);
|
||||
assert_eq!(months.len(), 2);
|
||||
assert!(months.contains(&"2015M01".to_string()));
|
||||
assert!(months.contains(&"2015M02".to_string()));
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user