Files
hagfish/src/db.rs
T

637 lines
20 KiB
Rust

//! 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"),
("Economic Zone (GEONOM2023)", "zone"),
("Processing (EUMOFAPresentation)", "processing"),
("Preservation (EUMOFAPreservation)", "preservation"),
("Shipsize", "shipsize"),
];
#[derive(Debug, Error)]
pub enum DbError {
#[error("DuckDB error: {0}")]
Duckdb(#[from] duckdb::Error),
#[error("No data provided")]
EmptyInput,
}
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)?;
info!("Initialized DuckDB at {}", path);
Ok(conn)
}
/// Creates all tables and indexes if they don't exist.
pub fn init_schema(conn: &Connection) -> Result<()> {
conn.execute_batch(
"
CREATE TABLE IF NOT EXISTS species (
code TEXT PRIMARY KEY,
label TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS gear (
code TEXT PRIMARY KEY,
label TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS zone (
code TEXT PRIMARY KEY,
label TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS processing (
code TEXT PRIMARY KEY,
label TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS preservation (
code TEXT PRIMARY KEY,
label TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS shipsize (
code TEXT PRIMARY KEY,
label TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS landings (
month TEXT NOT NULL,
species_code TEXT NOT NULL,
species_label TEXT,
gear_code TEXT NOT NULL,
zone_code TEXT NOT NULL,
processing_code TEXT NOT NULL,
preservation_code TEXT NOT NULL,
shipsize_code TEXT NOT NULL,
measure_code TEXT NOT NULL,
value DOUBLE,
PRIMARY KEY (month, species_code, gear_code, zone_code,
processing_code, preservation_code, shipsize_code,
measure_code)
);
CREATE INDEX IF NOT EXISTS idx_landings_month
ON landings(month);
CREATE INDEX IF NOT EXISTS idx_landings_species
ON landings(species_code);
",
)?;
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) {
let sql = format!(
"INSERT INTO {} (code, label) VALUES (?, ?) \
ON CONFLICT(code) DO UPDATE SET label = excluded.label",
table_name
);
for (code, label) in codes {
conn.execute(&sql, params![code, label])?;
}
info!(
"Updated {} entries in '{}' lookup table",
codes.len(),
table_name
);
}
}
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);
}
let months: HashSet<&str> = rows.iter().map(|r| r.month.as_str()).collect();
let tx = conn.unchecked_transaction()?;
for month in &months {
tx.execute("DELETE FROM landings WHERE month = ?", params![month])?;
}
{
let mut app = tx.appender("landings")?;
for row in rows {
app.append_row(params![
row.month,
row.species_code,
row.species_label,
row.gear_code,
row.zone_code,
row.processing_code,
row.preservation_code,
row.shipsize_code,
row.measure_code,
row.value,
])?;
}
// Flush explicitly — Drop discards errors per DuckDB docs.
app.flush()?;
}
tx.commit()?;
info!(
"Upserted {} rows across {} month(s)",
rows.len(),
months.len()
);
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)?;
Ok(val)
})?;
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<()> {
conn.execute(
&format!(
"COPY (SELECT * FROM landings) TO '{}' (FORMAT PARQUET)",
path
),
[],
)?;
info!("Exported landings to Parquet at {}", path);
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::types::Landing;
use std::collections::HashMap;
fn test_conn() -> Connection {
let conn = Connection::open_in_memory().unwrap();
init_schema(&conn).unwrap();
conn
}
fn sample_landings() -> Vec<Landing> {
vec![
Landing {
month: "2024M01".to_string(),
species_code: "COD".to_string(),
species_label: "Þorskur".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(1234.5),
},
Landing {
month: "2024M01".to_string(),
species_code: "COD".to_string(),
species_label: "Þorskur".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: None,
},
Landing {
month: "2024M02".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(5678.9),
},
]
}
// --- Schema initialization ---
#[test]
fn test_init_creates_all_tables() {
let conn = test_conn();
for table in &[
"species",
"gear",
"zone",
"processing",
"preservation",
"shipsize",
"landings",
] {
let count: i64 = conn
.query_row(&format!("SELECT COUNT(*) FROM {}", table), [], |row| {
row.get(0)
})
.unwrap_or_else(|_| panic!("Failed to query {}", table));
assert_eq!(count, 0, "Table {} should be empty", table);
}
}
#[test]
fn test_init_is_idempotent() {
let conn = Connection::open_in_memory().unwrap();
init_schema(&conn).unwrap();
init_schema(&conn).unwrap();
}
// --- Upsert operations ---
#[test]
fn test_upsert_insert_and_idempotency() {
let conn = test_conn();
let rows = sample_landings();
let inserted = upsert_landings(&conn, &rows).unwrap();
assert_eq!(inserted, 3);
let count: i64 = conn
.query_row("SELECT COUNT(*) FROM landings", [], |row| row.get(0))
.unwrap();
assert_eq!(count, 3);
let inserted = upsert_landings(&conn, &rows).unwrap();
assert_eq!(inserted, 3);
let count: i64 = conn
.query_row("SELECT COUNT(*) FROM landings", [], |row| row.get(0))
.unwrap();
assert_eq!(count, 3);
}
#[test]
fn test_upsert_replaces_month_data() {
let conn = test_conn();
upsert_landings(&conn, &sample_landings()).unwrap();
let modified = vec![Landing {
month: "2024M01".to_string(),
species_code: "COD".to_string(),
species_label: "Þorskur".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(9999.9),
}];
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();
assert_eq!(count, 2);
let val: f64 = conn
.query_row(
"SELECT value FROM landings
WHERE month = '2024M01' AND measure_code = 'MASS'",
[],
|row| row.get(0),
)
.unwrap();
assert!((val - 9999.9).abs() < f64::EPSILON);
}
#[test]
fn test_upsert_empty_input() {
let conn = test_conn();
let result = upsert_landings(&conn, &[]);
assert!(result.is_ok());
assert_eq!(result.unwrap(), 0);
}
#[test]
fn test_upsert_multi_month_atomic() {
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(),
species_code: "COD".to_string(),
species_label: "Þorskur".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(111.0),
},
Landing {
month: "2024M02".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(222.0),
},
];
upsert_landings(&conn, &rows).unwrap();
let jan_val: f64 = conn
.query_row(
"SELECT value FROM landings WHERE month = '2024M01' AND measure_code = 'MASS'",
[],
|row| row.get(0),
)
.unwrap();
assert!((jan_val - 111.0).abs() < f64::EPSILON);
let feb_val: f64 = conn
.query_row(
"SELECT value FROM landings WHERE month = '2024M02' AND measure_code = 'MASS'",
[],
|row| row.get(0),
)
.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();
let rows = vec![Landing {
month: "2024M01".to_string(),
species_code: "COD".to_string(),
species_label: "Þorskur".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: None,
}];
upsert_landings(&conn, &rows).unwrap();
let val: Option<f64> = conn
.query_row(
"SELECT value FROM landings WHERE measure_code = 'VALUE'",
[],
|row| row.get(0),
)
.unwrap();
assert!(val.is_none());
}
#[test]
fn test_null_preserved_after_upsert() {
let conn = test_conn();
upsert_landings(&conn, &sample_landings()).unwrap();
let val: Option<f64> = conn
.query_row(
"SELECT value FROM landings
WHERE month = '2024M01' AND measure_code = 'VALUE'",
[],
|row| row.get(0),
)
.unwrap();
assert!(val.is_none());
}
// --- get_last_month ---
#[test]
fn test_get_last_month_empty_table() {
let conn = test_conn();
let result = get_last_month(&conn).unwrap();
assert!(result.is_none());
}
#[test]
fn test_get_last_month_populated() {
let conn = test_conn();
upsert_landings(&conn, &sample_landings()).unwrap();
let last = get_last_month(&conn).unwrap();
assert_eq!(last.as_deref(), Some("2024M02"));
}
#[test]
fn test_get_last_month_single_month() {
let conn = test_conn();
let rows = vec![Landing {
month: "2015M03".to_string(),
species_code: "COD".to_string(),
species_label: "Þorskur".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(100.0),
}];
upsert_landings(&conn, &rows).unwrap();
let last = get_last_month(&conn).unwrap();
assert_eq!(last.as_deref(), Some("2015M03"));
}
// --- Lookup tables ---
#[test]
fn test_update_lookups_populates_tables() {
let conn = test_conn();
let mut maps = LookupMap::new();
let mut species = HashMap::new();
species.insert("COD".to_string(), "Þorskur".to_string());
species.insert("HER".to_string(), "Sild".to_string());
maps.insert("Species (ASFIS2022)".to_string(), species);
let mut gear = HashMap::new();
gear.insert("TOTAL".to_string(), "Tilsamans".to_string());
maps.insert("Fishing Gear (ISSCFG2016)".to_string(), gear);
update_lookups(&conn, &maps).unwrap();
let label: String = conn
.query_row("SELECT label FROM species WHERE code = 'COD'", [], |row| {
row.get(0)
})
.unwrap();
assert_eq!(label, "Þorskur");
let label: String = conn
.query_row("SELECT label FROM gear WHERE code = 'TOTAL'", [], |row| {
row.get(0)
})
.unwrap();
assert_eq!(label, "Tilsamans");
let count: i64 = conn
.query_row("SELECT COUNT(*) FROM species", [], |row| row.get(0))
.unwrap();
assert_eq!(count, 2);
}
#[test]
fn test_update_lookups_upsert_existing() {
let conn = test_conn();
let mut maps = LookupMap::new();
let mut species = HashMap::new();
species.insert("COD".to_string(), "Old label".to_string());
maps.insert("Species (ASFIS2022)".to_string(), species);
update_lookups(&conn, &maps).unwrap();
let mut species2 = HashMap::new();
species2.insert("COD".to_string(), "Þorskur".to_string());
let mut maps2 = LookupMap::new();
maps2.insert("Species (ASFIS2022)".to_string(), species2);
update_lookups(&conn, &maps2).unwrap();
let label: String = conn
.query_row("SELECT label FROM species WHERE code = 'COD'", [], |row| {
row.get(0)
})
.unwrap();
assert_eq!(label, "Þorskur");
let count: i64 = conn
.query_row("SELECT COUNT(*) FROM species", [], |row| row.get(0))
.unwrap();
assert_eq!(count, 1);
}
#[test]
fn test_update_lookups_skips_missing_dims() {
let conn = test_conn();
let mut maps = LookupMap::new();
let mut species = HashMap::new();
species.insert("COD".to_string(), "Þorskur".to_string());
maps.insert("Species (ASFIS2022)".to_string(), species);
// No gear, zone, etc. — should not error
update_lookups(&conn, &maps).unwrap();
let count: i64 = conn
.query_row("SELECT COUNT(*) FROM gear", [], |row| row.get(0))
.unwrap();
assert_eq!(count, 0);
}
// --- Parquet export ---
#[test]
fn test_export_parquet() {
let conn = test_conn();
upsert_landings(&conn, &sample_landings()).unwrap();
let tmp = tempfile::NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap().to_string() + ".parquet";
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),
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(count, 3);
}
// --- Faroese Unicode round-trip ---
#[test]
fn test_faroese_labels_round_trip() {
let conn = test_conn();
let rows = vec![Landing {
month: "2024M01".to_string(),
species_code: "148XXXXXXX00000".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(100.0),
}];
upsert_landings(&conn, &rows).unwrap();
let label: String = conn
.query_row("SELECT species_label FROM landings LIMIT 1", [], |row| {
row.get(0)
})
.unwrap();
assert_eq!(label, "Sild");
assert!(!label.contains('\u{FFFD}'));
}
}