diff --git a/Cargo.lock b/Cargo.lock index 7053aeb..e25cacb 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2398,6 +2398,16 @@ dependencies = [ "tracing-core", ] +[[package]] +name = "tracing-serde" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "704b1aeb7be0d0a84fc9828cae51dab5970fee5088f83d1dd7ee6f6246fc6ff1" +dependencies = [ + "serde", + "tracing-core", +] + [[package]] name = "tracing-subscriber" version = "0.3.23" @@ -2408,12 +2418,15 @@ dependencies = [ "nu-ansi-term", "once_cell", "regex-automata", + "serde", + "serde_json", "sharded-slab", "smallvec", "thread_local", "tracing", "tracing-core", "tracing-log", + "tracing-serde", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 38439a9..75fb93c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -12,7 +12,7 @@ serde = { version = "1", features = ["derive"] } serde_json = "1" thiserror = "2" tracing = "0.1" -tracing-subscriber = { version = "0.3", features = ["env-filter"] } +tracing-subscriber = { version = "0.3", features = ["env-filter", "json"] } anyhow = "1" duckdb = { version = "1.10505", features = ["bundled"] } axum = "0.8" diff --git a/src/api.rs b/src/api.rs index 517815c..1265698 100644 --- a/src/api.rs +++ b/src/api.rs @@ -263,7 +263,7 @@ async fn get_summary( let conn = conn.blocking_lock(); let where_clause = if species_filter.is_some() { - " WHERE species_code = ?".to_string() + " AND species_code = ?".to_string() } else { String::new() }; @@ -272,7 +272,7 @@ async fn get_summary( "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} \ + FROM landings WHERE 1=1{where_clause} \ GROUP BY month ORDER BY month" ); @@ -307,7 +307,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_clause} \ + FROM landings WHERE 1=1{where_clause} \ GROUP BY species_code, species_label \ ORDER BY total_value DESC NULLS LAST \ LIMIT 10" diff --git a/src/ingest.rs b/src/ingest.rs index 82b53cd..f41fbe7 100644 --- a/src/ingest.rs +++ b/src/ingest.rs @@ -3,8 +3,6 @@ use reqwest::Client; use std::collections::HashMap; use tracing::{debug, info}; -const PX_WEB_LANGUAGE: &str = "fo"; - const DIM_MONTH: &str = "month"; const DIM_SPECIES: &str = "Species (ASFIS2022)"; const DIM_GEAR: &str = "Fishing Gear (ISSCFG2016)"; @@ -74,6 +72,17 @@ pub fn extract_available_months(meta: &MetadataResponse) -> Vec { .unwrap_or_default() } +pub fn chunk_months(months: &[String], batch_size: usize) -> Vec> { + if months.is_empty() { + return vec![]; + } + + months + .chunks(batch_size) + .map(|chunk| chunk.to_vec()) + .collect() +} + pub fn build_query(all_months: &[String]) -> Query { assert!( !all_months.is_empty(), @@ -113,22 +122,22 @@ pub fn build_query(all_months: &[String]) -> Query { QueryItem { code: DIM_PROCESSING.to_string(), selection: Selection { - filter: "all".to_string(), - values: vec!["*".to_string()], + filter: "item".to_string(), + values: vec!["TOTAL".to_string()], }, }, QueryItem { code: DIM_PRESERVATION.to_string(), selection: Selection { - filter: "all".to_string(), - values: vec!["*".to_string()], + filter: "item".to_string(), + values: vec!["TOTAL".to_string()], }, }, QueryItem { code: DIM_SHIPSIZE.to_string(), selection: Selection { - filter: "all".to_string(), - values: vec!["*".to_string()], + filter: "item".to_string(), + values: vec!["TOTAL".to_string()], }, }, QueryItem { @@ -161,7 +170,20 @@ pub async fn fetch_data(client: &Client, url: &str, query: &Query) -> Result(); + + let data: DataResponse = match serde_json::from_str(&body) { + Ok(d) => d, + Err(e) => { + tracing::error!( + error = %e, + body_preview = %preview, + "failed to deserialize JSON-stat2 response" + ); + return Err(IngestError::JsonError(e)); + } + }; if data.dataset.value.is_empty() { return Err(IngestError::EmptyDataset); @@ -185,13 +207,8 @@ fn decode_key_indices(flat_index: usize, key_sizes: &[usize]) -> Vec { indices } -pub fn parse_row( - row_index: usize, - dataset: &DataResponse, - lookup_maps: &LookupMap, -) -> Result { - let dim_info = &dataset.dataset.dimension; - let dim_order = &dim_info.id; +pub fn parse_row(row_index: usize, dataset: &Dataset, lookup_maps: &LookupMap) -> Result { + let dim_order = &dataset.id; if dim_order.len() != 8 { return Err(IngestError::MissingDimension(format!( @@ -201,18 +218,18 @@ pub fn parse_row( ))); } - if row_index >= dataset.dataset.value.len() { + if row_index >= dataset.value.len() { return Err(IngestError::InvalidValueCode(format!( "Row index {} out of bounds (max {})", row_index, - dataset.dataset.value.len() + dataset.value.len() ))); } let category_lists: Vec> = dim_order .iter() .map(|dim_code| { - let dim = dim_info.dimensions.get(dim_code).ok_or_else(|| { + let dim = dataset.dimension.get(dim_code).ok_or_else(|| { IngestError::MissingDimension(format!( "Dimension '{}' not found in response", dim_code @@ -275,7 +292,7 @@ pub fn parse_row( let (shipsize_code, shipsize_label) = get(DIM_SHIPSIZE); let (measure_code, measure_label) = get(DIM_MEASURE); - let raw_value = dataset.dataset.value[row_index]; + let raw_value = dataset.value[row_index]; let value = match raw_value { Some(v) if is_sentinel(v) => None, other => other, @@ -394,20 +411,18 @@ mod tests { DataResponse { dataset: Dataset { - dimension: DimInfo { - id: vec![ - DIM_MONTH.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_MEASURE.to_string(), - ], - size: vec![2, 2, 1, 1, 1, 1, 1, 2], - dimensions, - }, + id: vec![ + DIM_MONTH.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_MEASURE.to_string(), + ], + size: vec![2, 2, 1, 1, 1, 1, 1, 2], + dimension: dimensions, value: vec![ Some(1234.5), Some(2345.6), @@ -418,7 +433,8 @@ mod tests { None, Some(6789.0), ], - status: vec![], + extra_fields: HashMap::new(), + status_array: vec![], }, } } @@ -464,6 +480,18 @@ mod tests { assert_eq!(query.query[1].code, DIM_SPECIES); assert_eq!(query.query[1].selection.filter, "all"); + 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[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] @@ -499,7 +527,7 @@ mod tests { let dataset = mock_dataset_response(); let lookup_maps = mock_lookup_maps(); - let row = parse_row(0, &dataset, &lookup_maps).expect("parse failed"); + let row = parse_row(0, &dataset.dataset, &lookup_maps).expect("parse failed"); assert_eq!(row.month, "2015M01"); assert_eq!(row.species_code, "148XXXXXXX00000"); @@ -514,7 +542,7 @@ mod tests { let dataset = mock_dataset_response(); let lookup_maps = mock_lookup_maps(); - let row = parse_row(2, &dataset, &lookup_maps).expect("parse failed"); + let row = parse_row(2, &dataset.dataset, &lookup_maps).expect("parse failed"); assert_eq!(row.species_code, "183XXXXXXX00000"); assert_eq!(row.species_label, "Toskur"); @@ -526,7 +554,7 @@ mod tests { let dataset = mock_dataset_response(); let lookup_maps = mock_lookup_maps(); - let row = parse_row(2, &dataset, &lookup_maps).expect("parse failed"); + let row = parse_row(2, &dataset.dataset, &lookup_maps).expect("parse failed"); assert!( row.value.is_none(), @@ -540,7 +568,7 @@ mod tests { let dataset = mock_dataset_response(); let lookup_maps = mock_lookup_maps(); - let row = parse_row(6, &dataset, &lookup_maps).expect("parse failed"); + let row = parse_row(6, &dataset.dataset, &lookup_maps).expect("parse failed"); assert!(row.value.is_none()); } @@ -550,7 +578,7 @@ mod tests { let dataset = mock_dataset_response(); let lookup_maps = mock_lookup_maps(); - let result = parse_row(100, &dataset, &lookup_maps); + let result = parse_row(100, &dataset.dataset, &lookup_maps); assert!(result.is_err()); match result { @@ -562,11 +590,11 @@ mod tests { #[test] fn test_parse_row_wrong_dimension_count() { let mut dataset = mock_dataset_response(); - dataset.dataset.dimension.id.pop(); - dataset.dataset.dimension.size.pop(); + dataset.dataset.id.pop(); + dataset.dataset.size.pop(); let lookup_maps = mock_lookup_maps(); - let result = parse_row(0, &dataset, &lookup_maps); + let result = parse_row(0, &dataset.dataset, &lookup_maps); assert!(result.is_err()); } @@ -575,7 +603,7 @@ mod tests { let dataset = mock_dataset_response(); let lookup_maps = mock_lookup_maps(); - let row = parse_row(0, &dataset, &lookup_maps).expect("parse failed"); + let row = parse_row(0, &dataset.dataset, &lookup_maps).expect("parse failed"); let landing = data_row_to_landing(&row); assert_eq!(landing.month, row.month); @@ -591,7 +619,7 @@ mod tests { let lookup_maps = mock_lookup_maps(); for i in 0..8 { - let row = parse_row(i, &dataset, &lookup_maps); + let row = parse_row(i, &dataset.dataset, &lookup_maps); assert!(row.is_ok(), "Row {} failed", i); if let Ok(r) = row { @@ -672,4 +700,68 @@ mod tests { assert_eq!(map.get("VALUE"), Some(&"Virði".to_string())); assert_eq!(map.len(), 2); } + + #[test] + fn test_chunk_months_basic() { + let months = vec![ + "2024M01".to_string(), + "2024M02".to_string(), + "2024M03".to_string(), + "2024M04".to_string(), + "2024M05".to_string(), + ]; + + let batches = chunk_months(&months, 3); + + assert_eq!(batches.len(), 2); + assert_eq!(batches[0], vec!["2024M01", "2024M02", "2024M03"]); + assert_eq!(batches[1], vec!["2024M04", "2024M05"]); + } + + #[test] + fn test_chunk_months_exact_division() { + let months = vec![ + "2024M01".to_string(), + "2024M02".to_string(), + "2024M03".to_string(), + "2024M04".to_string(), + "2024M05".to_string(), + "2024M06".to_string(), + ]; + + let batches = chunk_months(&months, 3); + + assert_eq!(batches.len(), 2); + assert_eq!(batches[0].len(), 3); + assert_eq!(batches[1].len(), 3); + } + + #[test] + fn test_chunk_months_empty() { + let months: Vec = vec![]; + let batches = chunk_months(&months, 12); + assert!(batches.is_empty()); + } + + #[test] + fn test_chunk_months_single_item() { + let months = vec!["2024M01".to_string()]; + let batches = chunk_months(&months, 12); + assert_eq!(batches.len(), 1); + assert_eq!(batches[0], vec!["2024M01"]); + } + + #[test] + fn test_chunk_months_large_batch() { + let months: Vec = (1..=37) + .map(|m| format!("2024M{:02}", ((m - 1) % 12) + 1)) + .collect(); + let batches = chunk_months(&months, 12); + + assert_eq!(batches.len(), 4); + assert_eq!(batches[0].len(), 12); + assert_eq!(batches[1].len(), 12); + assert_eq!(batches[2].len(), 12); + assert_eq!(batches[3].len(), 1); + } } diff --git a/src/main.rs b/src/main.rs index 67e644a..15687f4 100644 --- a/src/main.rs +++ b/src/main.rs @@ -6,12 +6,26 @@ mod types; use std::path::Path; use std::sync::Arc; +use std::time::Instant; use axum::serve; use clap::Parser; use tokio::net::TcpListener; use tracing_subscriber::EnvFilter; +fn init_logging() { + let filter = EnvFilter::from_default_env(); + + if std::env::var("HAGFISH_LOG_FORMAT").as_deref() == Ok("json") { + tracing_subscriber::fmt() + .with_env_filter(filter) + .json() + .init(); + } else { + tracing_subscriber::fmt().with_env_filter(filter).init(); + } +} + fn cleanup_stale_parquet_files() { let tmp_dir = std::env::temp_dir(); match std::fs::read_dir(&tmp_dir) { @@ -30,15 +44,13 @@ fn cleanup_stale_parquet_files() { } } } - Err(e) => tracing::warn!("Failed to list temp dir: {e}"), + Err(e) => tracing::warn!(error = %e, "Failed to list temp dir"), } } #[tokio::main] async fn main() -> anyhow::Result<()> { - tracing_subscriber::fmt() - .with_env_filter(EnvFilter::from_default_env()) - .init(); + init_logging(); let cli = cli::Cli::parse(); let config = types::Config::load(&cli.config)?; @@ -63,7 +75,7 @@ async fn run_serve(config: &types::Config) -> anyhow::Result<()> { let app = api::build_router(state); let listener = TcpListener::bind(&config.bind_address).await?; - tracing::info!("hagfish listening on http://{}", config.bind_address); + tracing::info!(address = %config.bind_address, "hagfish server started"); serve(listener, app).await?; @@ -71,10 +83,21 @@ async fn run_serve(config: &types::Config) -> anyhow::Result<()> { } async fn run_ingest(config: &types::Config, full: bool) -> anyhow::Result<()> { + let run_start = Instant::now(); + + tracing::info!( + mode = if full { "full" } else { "incremental" }, + "starting ingestion run" + ); + let conn = db::init(&config.duckdb_path)?; let client = reqwest::Client::new(); - let (lookup_maps, metadata) = ingest::fetch_metadata(&client, &config.data_source_url).await?; + let (lookup_maps, metadata) = ingest::fetch_metadata(&client, &config.data_source_url) + .await + .inspect_err(|e| { + tracing::error!(error = %e, stage = "fetch_metadata", "ingestion failed"); + })?; db::update_lookups(&conn, &lookup_maps)?; @@ -94,36 +117,145 @@ async fn run_ingest(config: &types::Config, full: bool) -> anyhow::Result<()> { }; if months_to_fetch.is_empty() { - tracing::info!("No new months to ingest"); + let elapsed = run_start.elapsed(); + tracing::info!( + mode = if full { "full" } else { "incremental" }, + months_requested = 0, + total_rows = 0, + batches = 0, + errors = 0, + elapsed_ms = elapsed.as_millis() as u64, + "ingestion complete — no new months" + ); return Ok(()); } - tracing::info!("Ingesting {} month(s)", months_to_fetch.len()); + const BATCH_SIZE: usize = 12; + let batches = ingest::chunk_months(&months_to_fetch, BATCH_SIZE); - // TODO: 5.2 — batch into chunks of <=12 months when count exceeds 12 - let query = ingest::build_query(&months_to_fetch); - let data = ingest::fetch_data(&client, &config.data_source_url, &query).await?; + tracing::info!( + mode = if full { "full" } else { "incremental" }, + months_requested = months_to_fetch.len(), + batches = batches.len(), + "beginning batched ingestion" + ); - let mut landings = Vec::with_capacity(data.dataset.value.len()); - for i in 0..data.dataset.value.len() { - let row = ingest::parse_row(i, &data, &lookup_maps)?; - landings.push(ingest::data_row_to_landing(&row)); + let mut total_rows = 0; + let mut total_skipped = 0; + let mut errors = 0; + + for (i, batch) in batches.iter().enumerate() { + let batch_start = Instant::now(); + + tracing::info!( + batch = i + 1, + total_batches = batches.len(), + months_in_batch = batch.len(), + months = %batch.join(", "), + "fetching batch" + ); + + let query = ingest::build_query(batch); + + let data = match ingest::fetch_data(&client, &config.data_source_url, &query).await { + Ok(d) => d, + Err(e) => { + tracing::error!( + batch = i + 1, + error = %e, + stage = "fetch_data", + "batch failed, continuing to next batch" + ); + errors += 1; + continue; + } + }; + + let mut landings = Vec::with_capacity(data.dataset.value.len()); + let mut parse_errors = 0; + + for j in 0..data.dataset.value.len() { + match ingest::parse_row(j, &data.dataset, &lookup_maps) { + Ok(row) => { + if row.value.is_some() { + landings.push(ingest::data_row_to_landing(&row)); + } else { + total_skipped += 1; + } + } + Err(e) => { + tracing::warn!( + batch = i + 1, + row_index = j, + error = %e, + "failed to parse row" + ); + parse_errors += 1; + errors += 1; + } + } + } + + let count = match db::upsert_landings(&conn, &landings) { + Ok(c) => c, + Err(e) => { + tracing::error!( + batch = i + 1, + error = %e, + stage = "upsert_landings", + "batch upsert failed" + ); + errors += 1; + continue; + } + }; + + total_rows += count; + let batch_elapsed = batch_start.elapsed(); + + tracing::info!( + batch = i + 1, + total_batches = batches.len(), + rows_inserted = count, + null_rows_skipped = total_skipped, + parse_errors = parse_errors, + elapsed_ms = batch_elapsed.as_millis() as u64, + "batch complete" + ); } - let count = db::upsert_landings(&conn, &landings)?; - tracing::info!("Ingested {} rows across {} month(s)", count, months_to_fetch.len()); + let elapsed = run_start.elapsed(); + + tracing::info!( + mode = if full { "full" } else { "incremental" }, + months_requested = months_to_fetch.len(), + batches = batches.len(), + total_rows = total_rows, + null_rows_skipped = total_skipped, + errors = errors, + elapsed_ms = elapsed.as_millis() as u64, + "ingestion run complete" + ); Ok(()) } fn run_export(config: &types::Config, out: &Path) -> anyhow::Result<()> { + let run_start = Instant::now(); + let conn = db::init(&config.duckdb_path)?; let path_str = out .to_str() .ok_or_else(|| anyhow::anyhow!("Output path contains invalid UTF-8"))?; db::export_parquet(&conn, path_str)?; - tracing::info!("Exported landings to {}", path_str); + + let elapsed = run_start.elapsed(); + tracing::info!( + path = path_str, + elapsed_ms = elapsed.as_millis() as u64, + "export complete" + ); Ok(()) } diff --git a/src/types.rs b/src/types.rs index 57044b0..fad5978 100644 --- a/src/types.rs +++ b/src/types.rs @@ -108,25 +108,20 @@ impl Default for QueryResponse { #[derive(Debug, Clone, Deserialize, Serialize)] pub struct DataResponse { + #[serde(flatten)] pub dataset: Dataset, } #[derive(Debug, Clone, Deserialize, Serialize)] pub struct Dataset { - pub dimension: DimInfo, - pub value: Vec>, - #[serde(default, skip_serializing_if = "Vec::is_empty")] - pub status: Vec, -} - -#[derive(Debug, Clone, Deserialize, Serialize)] -pub struct DimInfo { #[serde(default)] pub id: Vec, #[serde(default)] pub size: Vec, #[serde(flatten)] - pub dimensions: HashMap, + pub extra_fields: HashMap, + pub dimension: HashMap, + pub value: Vec>, } #[derive(Debug, Clone, Deserialize, Serialize)] @@ -143,6 +138,7 @@ pub struct CategoryInfo { pub label: HashMap, } +#[allow(dead_code)] #[derive(Debug, Clone)] pub struct DataRow { pub month: String, @@ -261,9 +257,6 @@ pub enum IngestError { #[error("API returned empty data set")] EmptyDataset, - - #[error("Unicode decode error: {0}")] - UnicodeError(String), } pub type LookupMap = HashMap>;