diff --git a/Cargo.lock b/Cargo.lock index 5e926e5..14798c5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2,6 +2,29 @@ # It is not intended for manual editing. version = 4 +[[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.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ddd31a130427c27518df266943a5308ed92d4b226cc639f5a8f1002816174301" +dependencies = [ + "memchr", +] + [[package]] name = "android_system_properties" version = "0.1.5" @@ -38,6 +61,183 @@ version = "0.7.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7c02d123df017efcdfbd739ef81735b36c5ba83ec3c59c80a9d7ecc718f92e50" +[[package]] +name = "arrow" +version = "55.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f3f15b4c6b148206ff3a2b35002e08929c2462467b62b9c02036d9c34f9ef994" +dependencies = [ + "arrow-arith", + "arrow-array", + "arrow-buffer", + "arrow-cast", + "arrow-data", + "arrow-ipc", + "arrow-ord", + "arrow-row", + "arrow-schema", + "arrow-select", + "arrow-string", +] + +[[package]] +name = "arrow-arith" +version = "55.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "30feb679425110209ae35c3fbf82404a39a4c0436bb3ec36164d8bffed2a4ce4" +dependencies = [ + "arrow-array", + "arrow-buffer", + "arrow-data", + "arrow-schema", + "chrono", + "num", +] + +[[package]] +name = "arrow-array" +version = "55.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "70732f04d285d49054a48b72c54f791bb3424abae92d27aafdf776c98af161c8" +dependencies = [ + "ahash", + "arrow-buffer", + "arrow-data", + "arrow-schema", + "chrono", + "half", + "hashbrown 0.15.5", + "num", +] + +[[package]] +name = "arrow-buffer" +version = "55.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "169b1d5d6cb390dd92ce582b06b23815c7953e9dfaaea75556e89d890d19993d" +dependencies = [ + "bytes", + "half", + "num", +] + +[[package]] +name = "arrow-cast" +version = "55.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e4f12eccc3e1c05a766cafb31f6a60a46c2f8efec9b74c6e0648766d30686af8" +dependencies = [ + "arrow-array", + "arrow-buffer", + "arrow-data", + "arrow-schema", + "arrow-select", + "atoi", + "base64", + "chrono", + "half", + "lexical-core", + "num", + "ryu", +] + +[[package]] +name = "arrow-data" +version = "55.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8de1ce212d803199684b658fc4ba55fb2d7e87b213de5af415308d2fee3619c2" +dependencies = [ + "arrow-buffer", + "arrow-schema", + "half", + "num", +] + +[[package]] +name = "arrow-ipc" +version = "55.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d9ea5967e8b2af39aff5d9de2197df16e305f47f404781d3230b2dc672da5d92" +dependencies = [ + "arrow-array", + "arrow-buffer", + "arrow-data", + "arrow-schema", + "flatbuffers", +] + +[[package]] +name = "arrow-ord" +version = "55.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6506e3a059e3be23023f587f79c82ef0bcf6d293587e3272d20f2d30b969b5a7" +dependencies = [ + "arrow-array", + "arrow-buffer", + "arrow-data", + "arrow-schema", + "arrow-select", +] + +[[package]] +name = "arrow-row" +version = "55.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "52bf7393166beaf79b4bed9bfdf19e97472af32ce5b6b48169d321518a08cae2" +dependencies = [ + "arrow-array", + "arrow-buffer", + "arrow-data", + "arrow-schema", + "half", +] + +[[package]] +name = "arrow-schema" +version = "55.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "af7686986a3bf2254c9fb130c623cdcb2f8e1f15763e7c71c310f0834da3d292" + +[[package]] +name = "arrow-select" +version = "55.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dd2b45757d6a2373faa3352d02ff5b54b098f5e21dccebc45a21806bc34501e5" +dependencies = [ + "ahash", + "arrow-array", + "arrow-buffer", + "arrow-data", + "arrow-schema", + "num", +] + +[[package]] +name = "arrow-string" +version = "55.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0377d532850babb4d927a06294314b316e23311503ed580ec6ce6a0158f49d40" +dependencies = [ + "arrow-array", + "arrow-buffer", + "arrow-data", + "arrow-schema", + "arrow-select", + "memchr", + "num", + "regex", + "regex-syntax", +] + +[[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" @@ -102,6 +302,12 @@ dependencies = [ "tracing", ] +[[package]] +name = "base64" +version = "0.22.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "72b3254f16251a8381aa12e40e3c4d2f0199f8c6508fbecb9d91f575e0fbb8c6" + [[package]] name = "bincode" version = "1.3.3" @@ -196,6 +402,26 @@ dependencies = [ "windows-link", ] +[[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 = "constant_time_eq" version = "0.4.2" @@ -217,6 +443,12 @@ dependencies = [ "libc", ] +[[package]] +name = "crunchy" +version = "0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "460fbee9c2c2f33933d720630a6a0bac33ba7053db5344fac858d4b8952d77d5" + [[package]] name = "equivalent" version = "1.0.2" @@ -245,6 +477,16 @@ version = "0.1.9" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5baebc0774151f905a1a2cc41989300b1e6fbb29aff0ceffa1064fdd3088d582" +[[package]] +name = "flatbuffers" +version = "25.12.19" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "35f6839d7b3b98adde531effaf34f0c2badc6f4735d26fe74709d8e513a96ef3" +dependencies = [ + "bitflags", + "rustc_version", +] + [[package]] name = "fnv" version = "1.0.7" @@ -299,6 +541,17 @@ dependencies = [ "slab", ] +[[package]] +name = "getrandom" +version = "0.2.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ff2abc00be7fca6ebc474524697ae276ad847ad0a6b3faa4bcb027e9a4614ad0" +dependencies = [ + "cfg-if", + "libc", + "wasi", +] + [[package]] name = "getrandom" version = "0.3.4" @@ -324,6 +577,18 @@ dependencies = [ "wasip3", ] +[[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" @@ -467,6 +732,12 @@ dependencies = [ "serde_core", ] +[[package]] +name = "integer-encoding" +version = "3.0.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8bb03732005da905c88227371639bf1ad885cc712789c011c31c5fb3ab3ccf02" + [[package]] name = "itoa" version = "1.0.18" @@ -507,12 +778,75 @@ version = "0.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "09edd9e8b54e49e587e4f6295a7d29c3ea94d469cb40ab8ca70b288248a81db2" +[[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.186" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "68ab91017fe16c622486840e4c83c9a37afeff978bd239b5293d61ece587de66" +[[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.12.1" @@ -590,6 +924,70 @@ dependencies = [ "windows-sys", ] +[[package]] +name = "num" +version = "0.4.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "35bd024e8b2ff75562e5f34e7f4905839deb4b22955ef5e73d2fea1b9813cb23" +dependencies = [ + "num-bigint", + "num-complex", + "num-integer", + "num-iter", + "num-rational", + "num-traits", +] + +[[package]] +name = "num-bigint" +version = "0.4.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a5e44f723f1133c9deac646763579fdb3ac745e418f2a7af9cd0c431da1f20b9" +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.46" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7969661fd2958a5cb096e56c8e1ad0444ac2bbcd0061bd28660485a44879858f" +dependencies = [ + "num-traits", +] + +[[package]] +name = "num-iter" +version = "0.1.45" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1429034a0490724d0075ebb2bc9e875d6503c3cf69e235a8941aa757d83ef5bf" +dependencies = [ + "autocfg", + "num-integer", + "num-traits", +] + +[[package]] +name = "num-rational" +version = "0.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f83d14da390562dca69fc84082e73e548e1ad308d24accdedd2720017cb37824" +dependencies = [ + "num-bigint", + "num-integer", + "num-traits", +] + [[package]] name = "num-traits" version = "0.2.19" @@ -597,6 +995,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "071dfc062690e90b734c0b2273ce72ad0ffa95f0c74596bc250dcfd960262841" dependencies = [ "autocfg", + "libm", ] [[package]] @@ -614,6 +1013,15 @@ version = "1.21.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9f7c3e4beb33f85d45ae3e3a1792185706c8e16d043238c593331cc7cd313b50" +[[package]] +name = "ordered-float" +version = "2.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "68f19d67e5a2795c94e73e0bb1cc1a7edeb2e28efd39e2e1c9b7a40c1108b11c" +dependencies = [ + "num-traits", +] + [[package]] name = "parking_lot" version = "0.12.5" @@ -637,6 +1045,40 @@ dependencies = [ "windows-link", ] +[[package]] +name = "parquet" +version = "55.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b17da4150748086bd43352bc77372efa9b6e3dbd06a04831d2a98c041c225cfa" +dependencies = [ + "ahash", + "arrow-array", + "arrow-buffer", + "arrow-cast", + "arrow-data", + "arrow-ipc", + "arrow-schema", + "arrow-select", + "base64", + "bytes", + "chrono", + "half", + "hashbrown 0.15.5", + "num", + "num-bigint", + "paste", + "seq-macro", + "thrift", + "twox-hash", + "zstd", +] + +[[package]] +name = "paste" +version = "1.0.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "57c0d7b74b563b49d38dae00a0c37d4d6de9b432382b2892f0574ddcae73fd0a" + [[package]] name = "percent-encoding" version = "2.3.2" @@ -806,12 +1248,44 @@ dependencies = [ "bitflags", ] +[[package]] +name = "regex" +version = "1.12.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e10754a14b9137dd7b1e3e5b0493cc9171fdd105e0ab477f51b72e7f3ac0e276" +dependencies = [ + "aho-corasick", + "memchr", + "regex-automata", + "regex-syntax", +] + +[[package]] +name = "regex-automata" +version = "0.4.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6e1dd4122fc1595e8162618945476892eefca7b88c52820e74af6262213cae8f" +dependencies = [ + "aho-corasick", + "memchr", + "regex-syntax", +] + [[package]] name = "regex-syntax" version = "0.8.10" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "dc897dd8d9e8bd1ed8cdad82b5966c3e0ecae09fb1907d58efaa013543185d0a" +[[package]] +name = "rustc_version" +version = "0.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cfcb3a22ef46e85b45de6ee7e79d063319ebb6594faafcf1c225ea92ab6e9b92" +dependencies = [ + "semver", +] + [[package]] name = "rustix" version = "1.1.4" @@ -861,6 +1335,12 @@ version = "1.0.28" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8a7852d02fc848982e0c167ef163aaff9cd91dc640ba85e263cb1ce46fae51cd" +[[package]] +name = "seq-macro" +version = "0.3.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1bc711410fbe7399f390ca1c3b60ad0f53f80e95c5eb935e52268a0e2cd49acc" + [[package]] name = "serde" version = "1.0.228" @@ -1056,6 +1536,26 @@ dependencies = [ "cfg-if", ] +[[package]] +name = "thrift" +version = "0.17.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7e54bc85fc7faa8bc175c4bab5b92ba8d9a3ce893d0e9f42cc455c8ab16a9e09" +dependencies = [ + "byteorder", + "integer-encoding", + "ordered-float", +] + +[[package]] +name = "tiny-keccak" +version = "2.0.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2c9d3793400a45f954c52e73d068316d76b6f4e36977e3fcebb13a2721e80237" +dependencies = [ + "crunchy", +] + [[package]] name = "tokio" version = "1.52.3" @@ -1199,6 +1699,7 @@ name = "usagedb" version = "0.3.0" dependencies = [ "anyhow", + "arrow", "axum", "bincode", "blake3", @@ -1209,6 +1710,7 @@ dependencies = [ "hyper", "lz4_flex", "memmap2", + "parquet", "proptest", "serde", "serde_json", @@ -1240,6 +1742,12 @@ version = "0.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65" +[[package]] +name = "version_check" +version = "0.9.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a" + [[package]] name = "wait-timeout" version = "0.2.1" diff --git a/Cargo.toml b/Cargo.toml index 3ac4046..78d6aa2 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -27,6 +27,8 @@ tracing = "0.1.44" tracing-subscriber = "0.3.23" uuid = { version = "1.23.1", features = ["v4"] } zstd = "0.13.3" +arrow = { version = "55", default-features = false, features = ["ipc"] } +parquet = { version = "55", default-features = false, features = ["arrow", "zstd"] } [dev-dependencies] proptest = "1.11.0" diff --git a/src/export/mod.rs b/src/export/mod.rs new file mode 100644 index 0000000..7e91c23 --- /dev/null +++ b/src/export/mod.rs @@ -0,0 +1,7 @@ +//! Export raw usage events to open analytical formats (Phase D). +//! +//! Currently supports Apache Parquet (`parquet` submodule). The reviewer's +//! framing was "internal format optimized for billing; external format +//! Parquet for warehouse/BI/debug" — this module is the bridge. + +pub mod parquet; diff --git a/src/export/parquet.rs b/src/export/parquet.rs new file mode 100644 index 0000000..d08c870 --- /dev/null +++ b/src/export/parquet.rs @@ -0,0 +1,197 @@ +//! Apache Parquet export of raw usage events. +//! +//! Schema is intentionally flat — no nested types — so any Parquet +//! consumer can read it without struct/map support. `correction_ref` +//! and `dimensions` are flattened/serialized: +//! +//! | Column | Arrow type | Notes | +//! | --- | --- | --- | +//! | `event_id` | Utf8 | | +//! | `kind` | Utf8 | "Usage" / "Correction" / "Retraction" | +//! | `correction_original_event_id` | Utf8 (nullable) | flattened from correction_ref | +//! | `correction_reason` | Utf8 (nullable) | flattened from correction_ref | +//! | `account_id` | Utf8 | | +//! | `subscription_id` | Utf8 (nullable) | | +//! | `product_id` | Utf8 | | +//! | `meter_id` | Utf8 | | +//! | `model_id` | Utf8 (nullable) | | +//! | `timestamp_ms` | Int64 | | +//! | `quantity` | Decimal128(38, 0) | i128 fits exactly | +//! | `unit` | Utf8 | | +//! | `source` | Utf8 | | +//! | `dimensions_canonical` | Utf8 | serde_json of `SmallDimensions` (BTreeMap → canonical order) | +//! | `ingested_at_ms` | Int64 | | +//! +//! Compression: zstd. Row group size: default (~64K rows). + +use std::fs::File; +use std::path::Path; +use std::sync::Arc; + +use arrow::array::{ArrayRef, Decimal128Array, Int64Array, RecordBatch, StringArray}; +use arrow::datatypes::{DataType, Field, Schema}; +use parquet::arrow::ArrowWriter; +use parquet::basic::{Compression, ZstdLevel}; +use parquet::file::properties::WriterProperties; + +use crate::model::event::{EventKind, UsageEvent}; +use crate::runtime::state::AppState; +use crate::storage::segment_reader::RawSegmentReader; + +/// Stats returned by an export run. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct ExportStats { + pub events_exported: u64, + pub segments_read: usize, +} + +/// Export every raw segment in the manifest to a single Parquet file. +/// +/// This is the simplest possible export — no partitioning, no filtering, +/// no chunking. Sufficient for "dump-and-load into the warehouse" use +/// cases at MVP scale. Later phases can add per-day partitioning and +/// streaming row groups for arbitrarily large exports. +pub async fn export_raw_segments( + state: &AppState, + output_path: &Path, +) -> anyhow::Result { + // Snapshot raw segment paths under the manifest read lock. + let segment_paths: Vec = { + let manifest = state.manifest.read().await; + manifest + .raw_segments + .iter() + .map(|s| state.config.db_root.join(format!("{}.seg", s.segment_id))) + .collect() + }; + + let mut all_events: Vec = Vec::new(); + let mut segments_read = 0usize; + for path in &segment_paths { + if !path.exists() { + tracing::warn!("export: manifest references missing segment {:?}", path); + continue; + } + let mut reader = RawSegmentReader::new(path.clone())?; + while let Some(event) = reader.read_next()? { + all_events.push(event); + } + segments_read += 1; + } + + write_parquet(&all_events, output_path)?; + + Ok(ExportStats { + events_exported: all_events.len() as u64, + segments_read, + }) +} + +/// Build an Arrow RecordBatch + write it to `output_path` as a Parquet file. +/// Public for testability — callers usually go through `export_raw_segments`. +pub fn write_parquet(events: &[UsageEvent], output_path: &Path) -> anyhow::Result<()> { + let schema = Arc::new(parquet_schema()); + + // Build per-column arrays. `Vec::with_capacity` up front avoids + // reallocations on the large segments. + let n = events.len(); + let mut event_id = Vec::with_capacity(n); + let mut kind = Vec::with_capacity(n); + let mut correction_original = Vec::with_capacity(n); + let mut correction_reason = Vec::with_capacity(n); + let mut account_id = Vec::with_capacity(n); + let mut subscription_id = Vec::with_capacity(n); + let mut product_id = Vec::with_capacity(n); + let mut meter_id = Vec::with_capacity(n); + let mut model_id = Vec::with_capacity(n); + let mut timestamp_ms = Vec::with_capacity(n); + let mut quantity = Vec::with_capacity(n); + let mut unit = Vec::with_capacity(n); + let mut source = Vec::with_capacity(n); + let mut dimensions_canonical = Vec::with_capacity(n); + let mut ingested_at_ms = Vec::with_capacity(n); + + for e in events { + event_id.push(e.event_id.0.clone()); + kind.push(match e.kind { + EventKind::Usage => "Usage", + EventKind::Correction => "Correction", + EventKind::Retraction => "Retraction", + }.to_string()); + correction_original.push(e.correction_ref.as_ref().map(|c| c.original_event_id.0.clone())); + correction_reason.push(e.correction_ref.as_ref().map(|c| c.reason.clone())); + account_id.push(e.account_id.0.clone()); + subscription_id.push(e.subscription_id.as_ref().map(|s| s.0.clone())); + product_id.push(e.product_id.0.clone()); + meter_id.push(e.meter_id.0.clone()); + model_id.push(e.model_id.as_ref().map(|m| m.0.clone())); + timestamp_ms.push(e.timestamp_ms); + quantity.push(e.quantity); + unit.push(e.unit.0.clone()); + source.push(e.source.0.clone()); + // BTreeMap → JSON is already canonical-by-key. + dimensions_canonical.push(serde_json::to_string(&e.dimensions).unwrap_or_default()); + ingested_at_ms.push(e.ingested_at_ms); + } + + let arrays: Vec = vec![ + Arc::new(StringArray::from(event_id)), + Arc::new(StringArray::from(kind)), + Arc::new(StringArray::from(correction_original)), + Arc::new(StringArray::from(correction_reason)), + Arc::new(StringArray::from(account_id)), + Arc::new(StringArray::from(subscription_id)), + Arc::new(StringArray::from(product_id)), + Arc::new(StringArray::from(meter_id)), + Arc::new(StringArray::from(model_id)), + Arc::new(Int64Array::from(timestamp_ms)), + Arc::new( + Decimal128Array::from(quantity) + .with_precision_and_scale(38, 0) + .map_err(|e| anyhow::anyhow!("decimal precision: {}", e))?, + ), + Arc::new(StringArray::from(unit)), + Arc::new(StringArray::from(source)), + Arc::new(StringArray::from(dimensions_canonical)), + Arc::new(Int64Array::from(ingested_at_ms)), + ]; + + let batch = RecordBatch::try_new(schema.clone(), arrays) + .map_err(|e| anyhow::anyhow!("RecordBatch::try_new: {}", e))?; + + let file = File::create(output_path)?; + let props = WriterProperties::builder() + .set_compression(Compression::ZSTD(ZstdLevel::default())) + .build(); + let mut writer = ArrowWriter::try_new(file, schema, Some(props)) + .map_err(|e| anyhow::anyhow!("ArrowWriter::try_new: {}", e))?; + writer + .write(&batch) + .map_err(|e| anyhow::anyhow!("ArrowWriter::write: {}", e))?; + writer + .close() + .map_err(|e| anyhow::anyhow!("ArrowWriter::close: {}", e))?; + Ok(()) +} + +/// The Parquet schema. Public so tests can compare round-trip output +/// against the canonical shape. +pub fn parquet_schema() -> Schema { + Schema::new(vec![ + Field::new("event_id", DataType::Utf8, false), + Field::new("kind", DataType::Utf8, false), + Field::new("correction_original_event_id", DataType::Utf8, true), + Field::new("correction_reason", DataType::Utf8, true), + Field::new("account_id", DataType::Utf8, false), + Field::new("subscription_id", DataType::Utf8, true), + Field::new("product_id", DataType::Utf8, false), + Field::new("meter_id", DataType::Utf8, false), + Field::new("model_id", DataType::Utf8, true), + Field::new("timestamp_ms", DataType::Int64, false), + Field::new("quantity", DataType::Decimal128(38, 0), false), + Field::new("unit", DataType::Utf8, false), + Field::new("source", DataType::Utf8, false), + Field::new("dimensions_canonical", DataType::Utf8, false), + Field::new("ingested_at_ms", DataType::Int64, false), + ]) +} diff --git a/src/lib.rs b/src/lib.rs index 7c7b620..0df175f 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,5 +1,6 @@ pub mod api; pub mod compact; +pub mod export; pub mod ingest; pub mod model; pub mod query; diff --git a/tests/parquet_export.rs b/tests/parquet_export.rs new file mode 100644 index 0000000..1780630 --- /dev/null +++ b/tests/parquet_export.rs @@ -0,0 +1,248 @@ +//! Integration tests for Parquet export (Phase D). +//! +//! Verifies: +//! - Round-trip: every event written via `export_raw_segments` reads +//! back from the Parquet file with the same field values +//! - Empty manifest produces a valid zero-row file +//! - The schema matches the canonical shape (column names + types) + +use std::path::PathBuf; +use std::sync::Arc; + +use arrow::datatypes::DataType; +use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder; +use tokio::sync::{Mutex, RwLock}; + +use usagedb::export::parquet::{export_raw_segments, parquet_schema, write_parquet}; +use usagedb::ingest::dedupe::HotDedupe; +use usagedb::ingest::flusher::build_segment_meta; +use usagedb::ingest::memtable::Memtable; +use usagedb::ingest::wal::Wal; +use usagedb::model::dimensions::SmallDimensions; +use usagedb::model::event::{CorrectionRef, EventKind, UsageEvent}; +use usagedb::model::ids::{ + AccountId, EventId, MeterId, ModelId, ProductId, SourceId, SubscriptionId, Unit, +}; +use usagedb::runtime::config::Config; +use usagedb::runtime::state::{AppState, AppStateInner}; +use usagedb::storage::manifest::Manifest; +use usagedb::storage::segment_writer::RawSegmentWriter; + +fn tmp_root() -> PathBuf { + let dir = tempfile::tempdir().expect("tempdir"); + let p = dir.path().to_path_buf(); + std::mem::forget(dir); + p +} + +fn tmp_file(name: &str) -> PathBuf { + let dir = tempfile::tempdir().expect("tempdir"); + let p = dir.path().join(name); + std::mem::forget(dir); + p +} + +fn rich_event(id: &str, account: &str, ts: i64, qty: i128) -> UsageEvent { + let mut dims = SmallDimensions::default(); + dims.inner.insert("provider".into(), "anthropic".into()); + UsageEvent { + event_id: EventId(id.to_string()), + kind: EventKind::Usage, + correction_ref: None, + account_id: AccountId(account.to_string()), + subscription_id: Some(SubscriptionId("sub_1".into())), + product_id: ProductId("ai_gateway".into()), + meter_id: MeterId("tokens.input".into()), + timestamp_ms: ts, + quantity: qty, + unit: Unit("token".into()), + source: SourceId("agentcore".into()), + model_id: Some(ModelId("claude-sonnet-4".into())), + dimensions: dims, + ingested_at_ms: ts + 5, + } +} + +fn build_state(db_root: PathBuf) -> AppState { + let config = Config { + db_root: db_root.clone(), + default_bucket_count: 2, + ..Config::default() + }; + std::fs::create_dir_all(&config.db_root).unwrap(); + let wal = Wal::open(db_root.join("wal"), 0).unwrap(); + let manifest = Manifest { bucket_count: 2, ..Manifest::default() }; + let (flush_sender, _r) = tokio::sync::mpsc::channel(4); + Arc::new(AppStateInner { + config, + dedupe: Mutex::new(HotDedupe::new(1000)), + wal: Mutex::new(wal), + memtable: Mutex::new(Memtable::new()), + manifest: RwLock::new(manifest), + flush_sender, + }) +} + +async fn commit_segment(state: &AppState, events: &[UsageEvent], bucket: u32) { + let segment_id = format!("raw_{}", uuid::Uuid::new_v4().simple()); + let path = state.config.db_root.join(format!("{}.seg", segment_id)); + let mut writer = RawSegmentWriter::new(path).unwrap(); + for e in events { writer.write_event(e).unwrap(); } + let (_rows, checksum) = writer.finish().unwrap(); + let meta = build_segment_meta(&segment_id, events, bucket, checksum); + let mut manifest = state.manifest.write().await; + manifest.raw_segments.push(meta); + manifest.save(&state.config.db_root).unwrap(); +} + +/// Read all events back out of a Parquet file. Returns rows as +/// (event_id, kind, account_id, quantity_i128). +fn read_parquet_minimal(path: &PathBuf) -> Vec<(String, String, String, i128)> { + use arrow::array::{Array, Decimal128Array, StringArray}; + let file = std::fs::File::open(path).unwrap(); + let reader = ParquetRecordBatchReaderBuilder::try_new(file) + .unwrap() + .build() + .unwrap(); + let mut out = Vec::new(); + for batch in reader { + let batch = batch.unwrap(); + let event_id = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let kind = batch + .column(1) + .as_any() + .downcast_ref::() + .unwrap(); + let account_id = batch + .column(4) + .as_any() + .downcast_ref::() + .unwrap(); + let quantity = batch + .column(10) + .as_any() + .downcast_ref::() + .unwrap(); + for i in 0..batch.num_rows() { + out.push(( + event_id.value(i).to_string(), + kind.value(i).to_string(), + account_id.value(i).to_string(), + quantity.value(i), + )); + } + } + out +} + +#[test] +fn write_parquet_round_trips_event_fields() { + let out = tmp_file("rt.parquet"); + let events = vec![ + rich_event("evt_a", "acc_1", 1_000, 100), + rich_event("evt_b", "acc_2", 2_000, i128::MAX / 4), + UsageEvent { + kind: EventKind::Correction, + correction_ref: Some(CorrectionRef { + original_event_id: EventId("evt_a".into()), + reason: "overcount".into(), + }), + ..rich_event("evt_c", "acc_1", 3_000, -100) + }, + ]; + write_parquet(&events, &out).unwrap(); + + let rows = read_parquet_minimal(&out); + assert_eq!(rows.len(), 3); + + assert_eq!(rows[0].0, "evt_a"); + assert_eq!(rows[0].1, "Usage"); + assert_eq!(rows[0].2, "acc_1"); + assert_eq!(rows[0].3, 100); + + assert_eq!(rows[1].0, "evt_b"); + assert_eq!(rows[1].3, i128::MAX / 4); + + assert_eq!(rows[2].0, "evt_c"); + assert_eq!(rows[2].1, "Correction"); + assert_eq!(rows[2].3, -100); +} + +#[test] +fn schema_matches_canonical_shape() { + let schema = parquet_schema(); + let expected_columns: &[(&str, bool)] = &[ + ("event_id", false), + ("kind", false), + ("correction_original_event_id", true), + ("correction_reason", true), + ("account_id", false), + ("subscription_id", true), + ("product_id", false), + ("meter_id", false), + ("model_id", true), + ("timestamp_ms", false), + ("quantity", false), + ("unit", false), + ("source", false), + ("dimensions_canonical", false), + ("ingested_at_ms", false), + ]; + assert_eq!(schema.fields().len(), expected_columns.len()); + for (i, (name, nullable)) in expected_columns.iter().enumerate() { + let f = schema.field(i); + assert_eq!(f.name(), name, "column {} name", i); + assert_eq!(f.is_nullable(), *nullable, "column {} nullability", name); + } + // quantity must be Decimal128(38, 0) so i128 fits exactly. + let quantity_field = schema.field(10); + assert_eq!(quantity_field.data_type(), &DataType::Decimal128(38, 0)); +} + +#[test] +fn empty_input_writes_valid_zero_row_parquet() { + let out = tmp_file("empty.parquet"); + write_parquet(&[], &out).unwrap(); + let rows = read_parquet_minimal(&out); + assert!(rows.is_empty()); + // Sanity: a real file was created. + assert!(std::fs::metadata(&out).unwrap().len() > 0); +} + +#[tokio::test] +async fn export_raw_segments_round_trips_through_manifest() { + let root = tmp_root(); + let state = build_state(root.clone()); + + let events_a = vec![ + rich_event("a1", "acc_x", 1_000, 10), + rich_event("a2", "acc_x", 1_500, 20), + ]; + let events_b = vec![rich_event("b1", "acc_y", 2_000, 30)]; + commit_segment(&state, &events_a, 0).await; + commit_segment(&state, &events_b, 1).await; + + let out = tmp_file("e2e.parquet"); + let stats = export_raw_segments(&state, &out).await.unwrap(); + assert_eq!(stats.events_exported, 3); + assert_eq!(stats.segments_read, 2); + + let rows = read_parquet_minimal(&out); + let mut ids: Vec = rows.iter().map(|r| r.0.clone()).collect(); + ids.sort(); + assert_eq!(ids, vec!["a1".to_string(), "a2".to_string(), "b1".to_string()]); +} + +#[tokio::test] +async fn export_empty_manifest_produces_zero_row_file() { + let root = tmp_root(); + let state = build_state(root.clone()); + let out = tmp_file("empty_e2e.parquet"); + let stats = export_raw_segments(&state, &out).await.unwrap(); + assert_eq!(stats.events_exported, 0); + assert_eq!(stats.segments_read, 0); +}