From bd63b7e78846a811297e02762d45ddf47f7e6d81 Mon Sep 17 00:00:00 2001 From: Bartal Laearsson Date: Sun, 16 Aug 2026 23:42:55 +0100 Subject: [PATCH] starting on duckdb database work --- Cargo.lock | 814 +++++++++++++++++++++++++++++++++++++++++++++++++- Cargo.toml | 2 + src/db.rs | 637 +++++++++++++++++++++++++++++++++++++++ src/ingest.rs | 16 +- src/main.rs | 1 + 5 files changed, 1457 insertions(+), 13 deletions(-) create mode 100644 src/db.rs diff --git a/Cargo.lock b/Cargo.lock index fabe43a..d8809ee 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2,6 +2,26 @@ # It is not intended for manual editing. version = 4 +[[package]] +name = "adler2" +version = "2.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "320119579fcad9c21884f5c4861d16174d0e06250625266f50fe6898340abefa" + +[[package]] +name = "ahash" +version = "0.8.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a15f179cd60c4584b8a8c596927aadc462e27f2ca70c04e0071964a73ba7a75" +dependencies = [ + "cfg-if", + "const-random", + "getrandom 0.3.4", + "once_cell", + "version_check", + "zerocopy", +] + [[package]] name = "aho-corasick" version = "1.1.5" @@ -11,12 +31,193 @@ dependencies = [ "memchr", ] +[[package]] +name = "android_system_properties" +version = "0.1.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ae221649c9976a6f6c56ae1facf410f3ddb33cc661c4b7b61020a912d4237fbc" +dependencies = [ + "libc", +] + [[package]] name = "anyhow" version = "1.0.104" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "330a5ed07fa54e4702c9d6c4174f74427fc0ef6e214bbd677ae50a5099946470" +[[package]] +name = "arbitrary" +version = "1.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c3d036a3c4ab069c7b410a2ce876bd74808d2d0888a82667669f8e783a898bf1" +dependencies = [ + "derive_arbitrary", +] + +[[package]] +name = "arrow" +version = "58.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6cfdd0833e32a9874d2b55089333ad310c0be208aafa277385ce2461dec90be3" +dependencies = [ + "arrow-arith", + "arrow-array", + "arrow-buffer", + "arrow-cast", + "arrow-data", + "arrow-ord", + "arrow-row", + "arrow-schema", + "arrow-select", + "arrow-string", +] + +[[package]] +name = "arrow-arith" +version = "58.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0a41203398f0eaa6f7ec8e62c0da742a21abf282c148fc157f6c35c90e29981a" +dependencies = [ + "arrow-array", + "arrow-buffer", + "arrow-data", + "arrow-schema", + "chrono", + "num-traits", +] + +[[package]] +name = "arrow-array" +version = "58.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ae33dad492b7df00a217563a7b0ef2874df68a0deea1b1a3acf628152f7f7a69" +dependencies = [ + "ahash", + "arrow-buffer", + "arrow-data", + "arrow-schema", + "chrono", + "half", + "hashbrown 0.17.1", + "num-complex", + "num-integer", + "num-traits", +] + +[[package]] +name = "arrow-buffer" +version = "58.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b9552f96391c005e6ab449fa941420935e7e062489b12b8b1b08879b2163f5b5" +dependencies = [ + "bytes", + "half", + "num-bigint", + "num-traits", +] + +[[package]] +name = "arrow-cast" +version = "58.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3a8a327c9649f30d8406995f27642b68df354713cca3baaaf100f076f18d5f34" +dependencies = [ + "arrow-array", + "arrow-buffer", + "arrow-data", + "arrow-ord", + "arrow-schema", + "arrow-select", + "atoi", + "base64 0.22.1", + "chrono", + "comfy-table", + "half", + "lexical-core", + "num-traits", + "ryu", +] + +[[package]] +name = "arrow-data" +version = "58.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2b24852db04738907e06c04ea61e42fe7fda962a34513022dc0d0e754fb7976b" +dependencies = [ + "arrow-buffer", + "arrow-schema", + "half", + "num-integer", + "num-traits", +] + +[[package]] +name = "arrow-ord" +version = "58.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "63a083ec750f5c043f02946b4baf05fcdbb55f4560a3277055caca5cc99f3eb0" +dependencies = [ + "arrow-array", + "arrow-buffer", + "arrow-data", + "arrow-schema", + "arrow-select", +] + +[[package]] +name = "arrow-row" +version = "58.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "514ba0ef0d4c5896202dae736251ce415abb43a950bed570fb7981b8716c0e4c" +dependencies = [ + "arrow-array", + "arrow-buffer", + "arrow-data", + "arrow-schema", + "half", +] + +[[package]] +name = "arrow-schema" +version = "58.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "21ca356ad6425cecb6eb7b28e4f659f1ee7880fbb1a16127de7dd62901efee9e" +dependencies = [ + "bitflags", +] + +[[package]] +name = "arrow-select" +version = "58.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c58da39eb3d8350ad4a549e5c2bc49284dac554016c69829310350f1731b0aad" +dependencies = [ + "ahash", + "arrow-array", + "arrow-buffer", + "arrow-data", + "arrow-schema", + "num-traits", +] + +[[package]] +name = "arrow-string" +version = "58.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b6789b388467525e3271326b6b4915666ecfdf5142aef09779445c954b67543c" +dependencies = [ + "arrow-array", + "arrow-buffer", + "arrow-data", + "arrow-schema", + "arrow-select", + "memchr", + "num-traits", + "regex", + "regex-syntax", +] + [[package]] name = "assert-json-diff" version = "2.0.2" @@ -27,18 +228,39 @@ dependencies = [ "serde_json", ] +[[package]] +name = "atoi" +version = "2.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f28d99ec8bfea296261ca1af174f24225171fea9664ba9003cbebee704810528" +dependencies = [ + "num-traits", +] + [[package]] name = "atomic-waker" version = "1.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1505bd5d3d116872e7271a6d4e16d81d0c8570876c8de68093a09ac269d8aac0" +[[package]] +name = "autocfg" +version = "1.5.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f2032f911046de80f0a198e0901378627c33f59ea0ac00e363d481118bd70a53" + [[package]] name = "base64" version = "0.22.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "72b3254f16251a8381aa12e40e3c4d2f0199f8c6508fbecb9d91f575e0fbb8c6" +[[package]] +name = "base64" +version = "0.23.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ac07cdecf99051d9a5238b80f35af32cdeba5b336e55d957b318b50137e18da5" + [[package]] name = "bitflags" version = "2.13.1" @@ -57,6 +279,12 @@ version = "1.12.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "fc652a48c352aef3ea3aed32080501cf3ef6ed5da78602a020c991775b0aff04" +[[package]] +name = "cast" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "37b2a672a2cb129a2e41c10b1224bb368f9f37a2b16b612598138befd7b37eb5" + [[package]] name = "cc" version = "1.4.3" @@ -64,6 +292,8 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "509591b7bcd67f4ef775afad7662703b4935daaa6ec0e5605cfb1090b32a2b6d" dependencies = [ "find-msvc-tools", + "jobserver", + "libc", "shlex", ] @@ -73,6 +303,17 @@ version = "1.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9330f8b2ff13f34540b44e946ef35111825727b38d33286ef986142615121801" +[[package]] +name = "chrono" +version = "0.4.45" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1aa79e62e7697b8e29b513a68abacf485adcd1fe8284a4316c5ae868e6633327" +dependencies = [ + "iana-time-zone", + "num-traits", + "windows-link", +] + [[package]] name = "colored" version = "3.1.1" @@ -82,6 +323,37 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "comfy-table" +version = "7.1.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4a65ebfec4fb190b6f90e944a817d60499ee0744e582530e2c9900a22e591d9a" +dependencies = [ + "crossterm", + "unicode-segmentation", + "unicode-width", +] + +[[package]] +name = "const-random" +version = "0.1.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "87e00182fe74b066627d63b85fd550ac2998d4b0bd86bfed477a0ae4c7c71359" +dependencies = [ + "const-random-macro", +] + +[[package]] +name = "const-random-macro" +version = "0.1.16" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f9d839f2a20b0aee515dc581a6172f2321f96cab76c1a38a4c584a194955390e" +dependencies = [ + "getrandom 0.2.17", + "once_cell", + "tiny-keccak", +] + [[package]] name = "core-foundation" version = "0.9.4" @@ -108,6 +380,54 @@ version = "0.8.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "773648b94d0e5d620f64f280777445740e61fe701025087ec8b57f45c791888b" +[[package]] +name = "crc32fast" +version = "1.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9481c1c90cbf2ac953f07c8d4a58aa3945c425b7185c9154d67a65e4230da511" +dependencies = [ + "cfg-if", +] + +[[package]] +name = "crossterm" +version = "0.28.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "829d955a0bb380ef178a640b91779e3987da38c9aea133b20614cfed8cdea9c6" +dependencies = [ + "bitflags", + "crossterm_winapi", + "parking_lot", + "rustix 0.38.44", + "winapi", +] + +[[package]] +name = "crossterm_winapi" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "acdd7c62a3665c7f6830a51635d9ac9b23ed385797f70a83bb8bafe9c572ab2b" +dependencies = [ + "winapi", +] + +[[package]] +name = "crunchy" +version = "0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "460fbee9c2c2f33933d720630a6a0bac33ba7053db5344fac858d4b8952d77d5" + +[[package]] +name = "derive_arbitrary" +version = "1.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1e567bd82dcff979e4b03460c307b3cdc9e96fde3d73bed1496d2bc75d9dd62a" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", +] + [[package]] name = "displaydoc" version = "0.2.7" @@ -119,6 +439,23 @@ dependencies = [ "syn 3.0.3", ] +[[package]] +name = "duckdb" +version = "1.10505.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "970e05eedd3f55c435194d9104f90a9b4a79a80d6e73251bc9ff43e178130c4e" +dependencies = [ + "arrow", + "cast", + "comfy-table", + "fallible-iterator", + "fallible-streaming-iterator", + "hashlink", + "libduckdb-sys", + "num-integer", + "strum", +] + [[package]] name = "encoding_rs" version = "0.8.35" @@ -144,24 +481,63 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "fallible-iterator" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2acce4a10f12dc2fb14a218589d4f1f62ef011b2d0cc4b3cb1bba8e94da14649" + +[[package]] +name = "fallible-streaming-iterator" +version = "0.1.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7360491ce676a36bf9bb3c56c1aa791658183a54d2744120f27285738d90465a" + [[package]] name = "fastrand" version = "2.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "da7c62ceae207dd37ea5b845da6a0696c799f85e97da1ab5b7910be3c1c80223" +[[package]] +name = "filetime" +version = "0.2.29" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5c287a33c7f0a620c38e641e7f60827713987b3c0f26e8ddc9462cc69cf75759" +dependencies = [ + "cfg-if", + "libc", +] + [[package]] name = "find-msvc-tools" version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d45db016d36b838f563236e9193d0ee6ce38f3f68b6c94e914b4929c96bbb890" +[[package]] +name = "flate2" +version = "1.1.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "843fba2746e448b37e26a819579957415c8cef339bf08564fe8b7ddbd959573c" +dependencies = [ + "crc32fast", + "miniz_oxide", + "zlib-rs", +] + [[package]] name = "fnv" version = "1.0.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3f9eec918d3f24069decb9af1554cad7c880e2da24a9afd88aca000531ab82c1" +[[package]] +name = "foldhash" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d9c4f5dac5e15c24eb999c26181a6ca40b39fe946cbe4c263c7209467bc83af2" + [[package]] name = "foreign-types" version = "0.3.2" @@ -283,10 +659,12 @@ name = "hagfish" version = "0.1.0" dependencies = [ "anyhow", + "duckdb", "mockito", "reqwest", "serde", "serde_json", + "tempfile", "thiserror", "tokio", "tokio-test", @@ -294,12 +672,48 @@ dependencies = [ "tracing-subscriber", ] +[[package]] +name = "half" +version = "2.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6ea2d84b969582b4b1864a92dc5d27cd2b77b622a8d79306834f1be5ba20d84b" +dependencies = [ + "cfg-if", + "crunchy", + "num-traits", + "zerocopy", +] + +[[package]] +name = "hashbrown" +version = "0.15.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9229cfe53dfd69f0609a49f65461bd93001ea1ef889cd5529dd176593f5338a1" +dependencies = [ + "foldhash", +] + [[package]] name = "hashbrown" version = "0.17.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ed5909b6e89a2db4456e54cd5f673791d7eca6732202bbf2a9cc504fe2f9b84a" +[[package]] +name = "hashlink" +version = "0.10.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7382cf6263419f2d8df38c55d7da83da5c18aef87fc7a7fc1fb1e344edfe14c1" +dependencies = [ + "hashbrown 0.15.5", +] + +[[package]] +name = "heck" +version = "0.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea" + [[package]] name = "http" version = "1.5.0" @@ -404,7 +818,7 @@ version = "0.1.20" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "96547c2556ec9d12fb1578c4eaf448b04993e7fb79cbaad930a656880a6bdfa0" dependencies = [ - "base64", + "base64 0.22.1", "bytes", "futures-channel", "futures-util", @@ -423,6 +837,30 @@ dependencies = [ "windows-registry", ] +[[package]] +name = "iana-time-zone" +version = "0.1.65" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e31bc9ad994ba00e440a8aa5c9ef0ec67d5cb5e5cb0cc7f8b744a35b389cc470" +dependencies = [ + "android_system_properties", + "core-foundation-sys", + "iana-time-zone-haiku", + "js-sys", + "log", + "wasm-bindgen", + "windows-core", +] + +[[package]] +name = "iana-time-zone-haiku" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f31827a206f56af32e590ba56d5d2d085f558508192593743f16b2306495269f" +dependencies = [ + "cc", +] + [[package]] name = "icu_collections" version = "2.3.0" @@ -534,7 +972,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d466e9454f08e4a911e14806c24e16fba1b4c121d1ea474396f396069cf949d9" dependencies = [ "equivalent", - "hashbrown", + "hashbrown 0.17.1", ] [[package]] @@ -549,6 +987,16 @@ version = "1.0.18" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8f42a60cbdf9a97f5d2305f08a87dc4e09308d1276d28c869c684d7777685682" +[[package]] +name = "jobserver" +version = "0.1.35" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1c00acbd29eabad4a2392fa0e921c874934dbbf4194312ad20f04a0ed67a3cb3" +dependencies = [ + "getrandom 0.4.3", + "libc", +] + [[package]] name = "js-sys" version = "0.3.104" @@ -566,12 +1014,98 @@ version = "1.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "bbd2bcb4c963f2ddae06a2efc7e9f3591312473c50c6685e1f298068316e66fe" +[[package]] +name = "lexical-core" +version = "1.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7d8d125a277f807e55a77304455eb7b1cb52f2b18c143b60e766c120bd64a594" +dependencies = [ + "lexical-parse-float", + "lexical-parse-integer", + "lexical-util", + "lexical-write-float", + "lexical-write-integer", +] + +[[package]] +name = "lexical-parse-float" +version = "1.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "52a9f232fbd6f550bc0137dcb5f99ab674071ac2d690ac69704593cb4abbea56" +dependencies = [ + "lexical-parse-integer", + "lexical-util", +] + +[[package]] +name = "lexical-parse-integer" +version = "1.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9a7a039f8fb9c19c996cd7b2fcce303c1b2874fe1aca544edc85c4a5f8489b34" +dependencies = [ + "lexical-util", +] + +[[package]] +name = "lexical-util" +version = "1.0.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2604dd126bb14f13fb5d1bd6a66155079cb9fa655b37f875b3a742c705dbed17" + +[[package]] +name = "lexical-write-float" +version = "1.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "50c438c87c013188d415fbabbb1dceb44249ab81664efbd31b14ae55dabb6361" +dependencies = [ + "lexical-util", + "lexical-write-integer", +] + +[[package]] +name = "lexical-write-integer" +version = "1.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "409851a618475d2d5796377cad353802345cba92c867d9fbcde9cf4eac4e14df" +dependencies = [ + "lexical-util", +] + [[package]] name = "libc" version = "0.2.189" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3eaf3ede3fee6db1a4c2ee091bf8a8b4dccdc6d17f656fb07896ee72867612f2" +[[package]] +name = "libduckdb-sys" +version = "1.10505.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6cb514dab5e271e849235c1cb98bd65a2ae107fbd619a6740219319c54a71d95" +dependencies = [ + "cc", + "flate2", + "pkg-config", + "serde", + "serde_json", + "tar", + "ureq", + "vcpkg", + "zip", +] + +[[package]] +name = "libm" +version = "0.2.16" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b6d2cec3eae94f9f509c767b45932f1ada8350c4bdb85af2fcab4a3c14807981" + +[[package]] +name = "linux-raw-sys" +version = "0.4.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d26c52dbd32dccf2d10cac7725f8eae5296885fb5703b261f7d0a0739ec807ab" + [[package]] name = "linux-raw-sys" version = "0.12.1" @@ -620,6 +1154,16 @@ version = "0.3.17" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6877bb514081ee2a7ff5ef9de3281f14a4dd4bceac4c09388074a6b5df8a139a" +[[package]] +name = "miniz_oxide" +version = "0.8.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1fa76a2c86f704bdb222d66965fb3d63269ce38518b83cb0575fca855ebb6316" +dependencies = [ + "adler2", + "simd-adler32", +] + [[package]] name = "mio" version = "1.2.2" @@ -682,6 +1226,44 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "num-bigint" +version = "0.4.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c89e69e7e0f03bea5ef08013795c25018e101932225a656383bd384495ecc367" +dependencies = [ + "num-integer", + "num-traits", +] + +[[package]] +name = "num-complex" +version = "0.4.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "73f88a1307638156682bada9d7604135552957b7818057dcef22705b4d509495" +dependencies = [ + "num-traits", +] + +[[package]] +name = "num-integer" +version = "0.1.47" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7ce2d95d4b3734dc35aa2f45e1aa22cd416814592a4f9d9205e11affd5b8e10b" +dependencies = [ + "num-traits", +] + +[[package]] +name = "num-traits" +version = "0.2.19" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "071dfc062690e90b734c0b2273ce72ad0ffa95f0c74596bc250dcfd960262841" +dependencies = [ + "autocfg", + "libm", +] + [[package]] name = "once_cell" version = "1.21.4" @@ -893,7 +1475,7 @@ version = "0.12.28" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "eddd3ca559203180a307f12d114c268abf583f59b03cb906fd0b3ff8646c1147" dependencies = [ - "base64", + "base64 0.22.1", "bytes", "encoding_rs", "futures-core", @@ -941,6 +1523,19 @@ dependencies = [ "windows-sys 0.52.0", ] +[[package]] +name = "rustix" +version = "0.38.44" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fdb5bc1ae2baa591800df16c9ca78619bf65c0488b41b96ccec5d11220d8c154" +dependencies = [ + "bitflags", + "errno", + "libc", + "linux-raw-sys 0.4.15", + "windows-sys 0.52.0", +] + [[package]] name = "rustix" version = "1.1.4" @@ -950,7 +1545,7 @@ dependencies = [ "bitflags", "errno", "libc", - "linux-raw-sys", + "linux-raw-sys 0.12.1", "windows-sys 0.61.2", ] @@ -960,7 +1555,9 @@ version = "0.23.43" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0283386ce02abc0151e1761d08802dfe86c173b0b494af5cbc086574e453da06" dependencies = [ + "log", "once_cell", + "ring", "rustls-pki-types", "rustls-webpki", "subtle", @@ -1117,6 +1714,12 @@ dependencies = [ "libc", ] +[[package]] +name = "simd-adler32" +version = "0.3.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3a219298ac11a56ea9a6d2120044824d6f01aeb034955e7af7bc16858527deea" + [[package]] name = "similar" version = "2.7.0" @@ -1151,6 +1754,27 @@ version = "1.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6ce2be8dc25455e1f91df71bfa12ad37d7af1092ae736f3a6cd0e37bc7810596" +[[package]] +name = "strum" +version = "0.27.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "af23d6f6c1a224baef9d3f61e287d2761385a5b88fdab4eb4c6f11aeb54c4bcf" +dependencies = [ + "strum_macros", +] + +[[package]] +name = "strum_macros" +version = "0.27.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7695ce3845ea4b33927c055a39dc438a45b059f7c1b3d91d38d10355fb8cbca7" +dependencies = [ + "heck", + "proc-macro2", + "quote", + "syn 2.0.119", +] + [[package]] name = "subtle" version = "2.6.1" @@ -1220,6 +1844,17 @@ dependencies = [ "libc", ] +[[package]] +name = "tar" +version = "0.4.46" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3f6221d9a6003c78398e3b239969f352578258df48c8eb051caadae0015bc840" +dependencies = [ + "filetime", + "libc", + "xattr", +] + [[package]] name = "tempfile" version = "3.27.0" @@ -1229,7 +1864,7 @@ dependencies = [ "fastrand", "getrandom 0.4.3", "once_cell", - "rustix", + "rustix 1.1.4", "windows-sys 0.61.2", ] @@ -1262,6 +1897,15 @@ dependencies = [ "cfg-if", ] +[[package]] +name = "tiny-keccak" +version = "2.0.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2c9d3793400a45f954c52e73d068316d76b6f4e36977e3fcebb13a2721e80237" +dependencies = [ + "crunchy", +] + [[package]] name = "tinystr" version = "0.8.4" @@ -1474,12 +2118,52 @@ version = "1.0.24" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75" +[[package]] +name = "unicode-segmentation" +version = "1.13.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c6f5d3c3b1bf09027a88a6bc961fc00497d651009560b5463668dc81b0fa87a8" + +[[package]] +name = "unicode-width" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b4ac048d71ede7ee76d585517add45da530660ef4390e49b098733c6e897f254" + [[package]] name = "untrusted" version = "0.9.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8ecb6da28b8a351d773b68d5825ac39017e680750f980f3a1a85cd8dd28a47c1" +[[package]] +name = "ureq" +version = "3.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "972d7902c8735f2695410b8aed7df6ed12a47394aa1c8d7af49f0497b731a94d" +dependencies = [ + "base64 0.23.1", + "log", + "percent-encoding", + "rustls", + "rustls-pki-types", + "ureq-proto", + "utf8-zero", + "webpki-roots", +] + +[[package]] +name = "ureq-proto" +version = "0.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "da5f78b09e6941e1a0f2e30e695e4b120377b54d5e0aec11b594bb57b3971613" +dependencies = [ + "base64 0.23.1", + "http", + "httparse", + "log", +] + [[package]] name = "url" version = "2.5.8" @@ -1492,6 +2176,12 @@ dependencies = [ "serde", ] +[[package]] +name = "utf8-zero" +version = "0.8.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b8c0a043c9540bae7c578c88f91dda8bd82e59ae27c21baca69c8b191aaf5a6e" + [[package]] name = "utf8_iter" version = "1.0.4" @@ -1510,6 +2200,12 @@ version = "0.2.15" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "accd4ea62f7bb7a82fe23066fb0957d48ef677f6eeb8215f372f52e48bb32426" +[[package]] +name = "version_check" +version = "0.9.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a" + [[package]] name = "want" version = "0.3.1" @@ -1599,6 +2295,72 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "webpki-roots" +version = "1.0.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7dcd9d09a39985f5344844e66b0c530a33843579125f23e21e9f0f220850f22a" +dependencies = [ + "rustls-pki-types", +] + +[[package]] +name = "winapi" +version = "0.3.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5c839a674fcd7a98952e593242ea400abe93992746761e38641405d28b00f419" +dependencies = [ + "winapi-i686-pc-windows-gnu", + "winapi-x86_64-pc-windows-gnu", +] + +[[package]] +name = "winapi-i686-pc-windows-gnu" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ac3b87c63620426dd9b991e5ce0329eff545bccbbb34f3be09ff6fb6ab51b7b6" + +[[package]] +name = "winapi-x86_64-pc-windows-gnu" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "712e227841d057c1ee1cd2fb22fa7e5a5461ae8e48fa2ca79ec42cfc1931183f" + +[[package]] +name = "windows-core" +version = "0.62.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b8e83a14d34d0623b51dce9581199302a221863196a1dde71a7663a4c2be9deb" +dependencies = [ + "windows-implement", + "windows-interface", + "windows-link", + "windows-result", + "windows-strings", +] + +[[package]] +name = "windows-implement" +version = "0.60.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "053e2e040ab57b9dc951b72c264860db7eb3b0200ba345b4e4c3b14f67855ddf" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", +] + +[[package]] +name = "windows-interface" +version = "0.59.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3f316c4a2570ba26bbec722032c4099d8c8bc095efccdc15688708623367e358" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", +] + [[package]] name = "windows-link" version = "0.2.1" @@ -1728,6 +2490,16 @@ version = "0.6.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3ad82d2a33cdc9674dc7465672f271e096168fcdbe0f799d9e6db8c5892679dc" +[[package]] +name = "xattr" +version = "1.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "32e45ad4206f6d2479085147f02bc2ef834ac85886624a23575ae137c8aa8156" +dependencies = [ + "libc", + "rustix 1.1.4", +] + [[package]] name = "yoke" version = "0.8.3" @@ -1831,8 +2603,40 @@ dependencies = [ "syn 3.0.3", ] +[[package]] +name = "zip" +version = "6.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "eb2a05c7c36fde6c09b08576c9f7fb4cda705990f73b58fe011abf7dfb24168b" +dependencies = [ + "arbitrary", + "crc32fast", + "flate2", + "indexmap", + "memchr", + "zopfli", +] + +[[package]] +name = "zlib-rs" +version = "0.6.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "34b31d188d9d685a4f9c7b46d6e36631b07058d2cfe190267adce54dc230bf12" + [[package]] name = "zmij" version = "1.0.23" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "29666d0abbfad1e3dc4dcf6144730dd3a3ab225bbbdac83319345b1b44ccfc1b" + +[[package]] +name = "zopfli" +version = "0.8.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f05cd8797d63865425ff89b5c4a48804f35ba0ce8d125800027ad6017d2b5249" +dependencies = [ + "bumpalo", + "crc32fast", + "log", + "simd-adler32", +] diff --git a/Cargo.toml b/Cargo.toml index c45f563..5904dcf 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -14,7 +14,9 @@ thiserror = "2" tracing = "0.1" tracing-subscriber = { version = "0.3", features = ["env-filter"] } anyhow = "1" +duckdb = { version = "1.10505", features = ["bundled"] } [dev-dependencies] mockito = "1" tokio-test = "0.4" +tempfile = "3" diff --git a/src/db.rs b/src/db.rs new file mode 100644 index 0000000..768561d --- /dev/null +++ b/src/db.rs @@ -0,0 +1,637 @@ +//! 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 = std::result::Result; + +/// Opens a DuckDB connection at the given path and initializes the schema. +pub fn init(path: &str) -> Result { + 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. +/// +/// 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 { + if rows.is_empty() { + return Ok(0); + } + + let months: HashSet<&str> = rows.iter().map(|r| r.month.as_str()).collect(); + + let tx = conn.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 (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", + )?; + + for row in rows { + stmt.execute(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, + ])?; + } + } + + 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> { + let result = conn.query_row("SELECT MAX(month) FROM landings", [], |row| { + let val: Option = 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 { + 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 mut conn = test_conn(); + let rows = sample_landings(); + + let inserted = upsert_landings(&mut 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(&mut 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 mut conn = test_conn(); + upsert_landings(&mut 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(&mut 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 mut conn = test_conn(); + let result = upsert_landings(&mut 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(); + + // 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(&mut 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 mut 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(&mut conn, &rows).unwrap(); + + let val: Option = 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 mut conn = test_conn(); + upsert_landings(&mut conn, &sample_landings()).unwrap(); + + let val: Option = 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 mut conn = test_conn(); + upsert_landings(&mut 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 mut 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(&mut 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 mut conn = test_conn(); + upsert_landings(&mut 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 mut 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(&mut 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}')); + } +} diff --git a/src/ingest.rs b/src/ingest.rs index 0238b4c..f639e93 100644 --- a/src/ingest.rs +++ b/src/ingest.rs @@ -16,7 +16,7 @@ use crate::types::*; use reqwest::Client; use std::collections::HashMap; -use tracing::{debug, info, warn}; +use tracing::{debug, info}; /// Language code for PX-Web queries. const PX_WEB_LANGUAGE: &str = "fo"; @@ -479,24 +479,24 @@ mod tests { fn mock_lookup_maps() -> LookupMap { let mut maps = LookupMap::new(); maps.insert( - "0".to_string(), // Tid / Month + DIM_MONTH.to_string(), HashMap::from([ ("2015M01".to_string(), "Januar 2015".to_string()), ("2015M02".to_string(), "Februar 2015".to_string()), ]), ); maps.insert( - "1".to_string(), // Art / Species + DIM_SPECIES.to_string(), HashMap::from([ - ("Sild".to_string(), "Sild".to_string()), - ("Þorskur".to_string(), "Þorskur".to_string()), + ("148XXXXXXX00000".to_string(), "Sild".to_string()), + ("183XXXXXXX00000".to_string(), "Þorskur".to_string()), ]), ); maps.insert( - "7".to_string(), // Måleenhed / Measure + DIM_MEASURE.to_string(), HashMap::from([ - ("MASS".to_string(), "Kilo".to_string()), - ("VALUE".to_string(), "Krónur".to_string()), + ("MASS".to_string(), "Nøgd".to_string()), + ("VALUE".to_string(), "Virði".to_string()), ]), ); maps diff --git a/src/main.rs b/src/main.rs index 0c6638b..6881768 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,3 +1,4 @@ +mod db; mod ingest; mod types;