From 7a64f00014c4e08198ddc933eac88b4f40886536 Mon Sep 17 00:00:00 2001 From: Bartal Laearsson Date: Sun, 16 Aug 2026 23:53:30 +0100 Subject: [PATCH] phase 2 ready for qa review --- src/db.rs | 67 +++++++++++++++++++++++++++---------------------------- 1 file changed, 33 insertions(+), 34 deletions(-) diff --git a/src/db.rs b/src/db.rs index 768561d..f6937b5 100644 --- a/src/db.rs +++ b/src/db.rs @@ -119,33 +119,30 @@ pub fn update_lookups(conn: &Connection, lookup_maps: &LookupMap) -> Result<()> /// 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: &mut Connection, rows: &[Landing]) -> Result { +pub fn upsert_landings(conn: &Connection, rows: &[Landing]) -> Result { if rows.is_empty() { return Ok(0); } let months: HashSet<&str> = rows.iter().map(|r| r.month.as_str()).collect(); - let tx = conn.transaction()?; + let tx = conn.unchecked_transaction()?; for month in &months { tx.execute("DELETE FROM landings WHERE month = ?", params![month])?; } { - let mut stmt = tx.prepare( - "INSERT INTO landings ( - month, species_code, species_label, gear_code, zone_code, - processing_code, preservation_code, shipsize_code, - measure_code, value - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", - )?; - + let mut app = tx.appender("landings")?; for row in rows { - stmt.execute(params![ + app.append_row(params![ row.month, row.species_code, row.species_label, @@ -158,6 +155,8 @@ pub fn upsert_landings(conn: &mut Connection, rows: &[Landing]) -> Result row.value, ])?; } + // Flush explicitly — Drop discards errors per DuckDB docs. + app.flush()?; } tx.commit()?; @@ -288,10 +287,10 @@ mod tests { #[test] fn test_upsert_insert_and_idempotency() { - let mut conn = test_conn(); + let conn = test_conn(); let rows = sample_landings(); - let inserted = upsert_landings(&mut conn, &rows).unwrap(); + let inserted = upsert_landings(&conn, &rows).unwrap(); assert_eq!(inserted, 3); let count: i64 = conn @@ -299,7 +298,7 @@ mod tests { .unwrap(); assert_eq!(count, 3); - let inserted = upsert_landings(&mut conn, &rows).unwrap(); + let inserted = upsert_landings(&conn, &rows).unwrap(); assert_eq!(inserted, 3); let count: i64 = conn @@ -310,8 +309,8 @@ mod tests { #[test] fn test_upsert_replaces_month_data() { - let mut conn = test_conn(); - upsert_landings(&mut conn, &sample_landings()).unwrap(); + let conn = test_conn(); + upsert_landings(&conn, &sample_landings()).unwrap(); let modified = vec![Landing { month: "2024M01".to_string(), @@ -325,7 +324,7 @@ mod tests { measure_code: "MASS".to_string(), value: Some(9999.9), }]; - upsert_landings(&mut conn, &modified).unwrap(); + upsert_landings(&conn, &modified).unwrap(); // 2024M01 had 2 rows, now has 1. 2024M02 untouched. let count: i64 = conn @@ -346,16 +345,16 @@ mod tests { #[test] fn test_upsert_empty_input() { - let mut conn = test_conn(); - let result = upsert_landings(&mut conn, &[]); + 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 mut conn = test_conn(); - upsert_landings(&mut conn, &sample_landings()).unwrap(); + let conn = test_conn(); + upsert_landings(&conn, &sample_landings()).unwrap(); // Upsert both months in one call — both should be replaced. let rows = vec![ @@ -384,7 +383,7 @@ mod tests { value: Some(222.0), }, ]; - upsert_landings(&mut conn, &rows).unwrap(); + upsert_landings(&conn, &rows).unwrap(); let jan_val: f64 = conn .query_row( @@ -415,7 +414,7 @@ mod tests { #[test] fn test_null_preserved_on_insert() { - let mut conn = test_conn(); + let conn = test_conn(); let rows = vec![Landing { month: "2024M01".to_string(), species_code: "COD".to_string(), @@ -428,7 +427,7 @@ mod tests { measure_code: "VALUE".to_string(), value: None, }]; - upsert_landings(&mut conn, &rows).unwrap(); + upsert_landings(&conn, &rows).unwrap(); let val: Option = conn .query_row( @@ -442,8 +441,8 @@ mod tests { #[test] fn test_null_preserved_after_upsert() { - let mut conn = test_conn(); - upsert_landings(&mut conn, &sample_landings()).unwrap(); + let conn = test_conn(); + upsert_landings(&conn, &sample_landings()).unwrap(); let val: Option = conn .query_row( @@ -467,8 +466,8 @@ mod tests { #[test] fn test_get_last_month_populated() { - let mut conn = test_conn(); - upsert_landings(&mut conn, &sample_landings()).unwrap(); + let conn = test_conn(); + upsert_landings(&conn, &sample_landings()).unwrap(); let last = get_last_month(&conn).unwrap(); assert_eq!(last.as_deref(), Some("2024M02")); @@ -476,7 +475,7 @@ mod tests { #[test] fn test_get_last_month_single_month() { - let mut conn = test_conn(); + let conn = test_conn(); let rows = vec![Landing { month: "2015M03".to_string(), species_code: "COD".to_string(), @@ -489,7 +488,7 @@ mod tests { measure_code: "MASS".to_string(), value: Some(100.0), }]; - upsert_landings(&mut conn, &rows).unwrap(); + upsert_landings(&conn, &rows).unwrap(); let last = get_last_month(&conn).unwrap(); assert_eq!(last.as_deref(), Some("2015M03")); @@ -587,8 +586,8 @@ mod tests { #[test] fn test_export_parquet() { - let mut conn = test_conn(); - upsert_landings(&mut conn, &sample_landings()).unwrap(); + 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"; @@ -610,7 +609,7 @@ mod tests { #[test] fn test_faroese_labels_round_trip() { - let mut conn = test_conn(); + let conn = test_conn(); let rows = vec![Landing { month: "2024M01".to_string(), @@ -624,7 +623,7 @@ mod tests { measure_code: "MASS".to_string(), value: Some(100.0), }]; - upsert_landings(&mut conn, &rows).unwrap(); + upsert_landings(&conn, &rows).unwrap(); let label: String = conn .query_row("SELECT species_label FROM landings LIMIT 1", [], |row| {