From ca8367771f9ea6472708a5870bc80cb539179858 Mon Sep 17 00:00:00 2001 From: przemek Date: Sat, 16 May 2026 19:04:05 +0200 Subject: [PATCH] Admin CLI: check / rebuild-rollups / inspect-segment / verify-period / export-parquet MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The usagedb binary doubles as an admin tool. The existing HTTP server behavior moves under `usagedb serve` (still the default when no subcommand is given, so prior invocations work unchanged). Subcommands: check [--deep] Print manifest summary (generation, segment counts, watermark). --deep additionally opens every segment file and verifies checksum + structural validity. Fails non-zero if any segment is corrupt. rebuild-rollups --from --to Drops rollup segments overlapping the range + rewinds the watermark to `from`. Next server tick refills from raw segments. inspect-segment Prints SegmentMeta fields (bucket, time range, account/product/ meter/model sets, checksum, input_segment_ids for rollups) + first 5 rows as a sample. verify-period --account --from --to Same logic as the /verify HTTP endpoint. Prints raw_total, rollup_total, drift, and whether the period is sealed. export-parquet Calls export_raw_segments() to write every raw segment to one Parquet file with zstd compression. All commands accept --db-root (global, default ./data). They operate on the on-disk state without requiring a running server. Assumes the server is NOT running concurrently — file locking is on the backlog. Implementation: - New module src/admin.rs with cmd_* async functions taking AppState + args and returning Result. Each command formats its own output; the main binary just prints. - open_state_for_admin() loads the manifest directly via Manifest::load() — cheaper than full Recovery (no segment-scan dedupe rebuild, no WAL replay) for one-shot commands. - main.rs uses clap (derive macros) for arg parsing. The old server logic moves into run_server() — same code path, just behind a Command::Serve match arm. Tests (tests/admin.rs, 10 tests): - check_reports_manifest_summary - check_deep_passes_for_valid_segments - check_deep_fails_for_corrupt_segment (flips a byte mid-file) - check_on_fresh_db_errors_with_clear_message - inspect_segment_prints_meta_and_sample_rows - inspect_segment_errors_for_unknown_id - rebuild_rollups_drops_and_rewinds (tick → rebuild → verify state) - rebuild_rollups_rejects_invalid_dates - verify_period_reports_match_for_sealed_rollup - export_parquet_writes_file clap added as a regular dep (~1 MB). Total tests: 106 (was 96; +10). Clean under RUSTFLAGS=-D warnings. Co-Authored-By: Claude Opus 4.7 (1M context) --- Cargo.lock | 121 ++++++++++++++++++ Cargo.toml | 1 + README.md | 20 ++- src/admin.rs | 325 +++++++++++++++++++++++++++++++++++++++++++++++ src/lib.rs | 1 + src/main.rs | 120 ++++++++++++++++-- tests/admin.rs | 336 +++++++++++++++++++++++++++++++++++++++++++++++++ 7 files changed, 909 insertions(+), 15 deletions(-) create mode 100644 src/admin.rs create mode 100644 tests/admin.rs diff --git a/Cargo.lock b/Cargo.lock index 14798c5..b2219aa 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -34,6 +34,56 @@ dependencies = [ "libc", ] +[[package]] +name = "anstream" +version = "1.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "824a212faf96e9acacdbd09febd34438f8f711fb84e09a8916013cd7815ca28d" +dependencies = [ + "anstyle", + "anstyle-parse", + "anstyle-query", + "anstyle-wincon", + "colorchoice", + "is_terminal_polyfill", + "utf8parse", +] + +[[package]] +name = "anstyle" +version = "1.0.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "940b3a0ca603d1eade50a4846a2afffd5ef57a9feac2c0e2ec2e14f9ead76000" + +[[package]] +name = "anstyle-parse" +version = "1.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "52ce7f38b242319f7cabaa6813055467063ecdc9d355bbb4ce0c68908cd8130e" +dependencies = [ + "utf8parse", +] + +[[package]] +name = "anstyle-query" +version = "1.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "40c48f72fd53cd289104fc64099abca73db4166ad86ea0b4341abe65af83dadc" +dependencies = [ + "windows-sys", +] + +[[package]] +name = "anstyle-wincon" +version = "3.0.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "291e6a250ff86cd4a820112fb8898808a366d8f9f58ce16d1f538353ad55747d" +dependencies = [ + "anstyle", + "once_cell_polyfill", + "windows-sys", +] + [[package]] name = "anyhow" version = "1.0.102" @@ -402,6 +452,52 @@ dependencies = [ "windows-link", ] +[[package]] +name = "clap" +version = "4.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1ddb117e43bbf7dacf0a4190fef4d345b9bad68dfc649cb349e7d17d28428e51" +dependencies = [ + "clap_builder", + "clap_derive", +] + +[[package]] +name = "clap_builder" +version = "4.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "714a53001bf66416adb0e2ef5ac857140e7dc3a0c48fb28b2f10762fc4b5069f" +dependencies = [ + "anstream", + "anstyle", + "clap_lex", + "strsim", +] + +[[package]] +name = "clap_derive" +version = "4.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f2ce8604710f6733aa641a2b3731eaa1e8b3d9973d5e3565da11800813f997a9" +dependencies = [ + "heck", + "proc-macro2", + "quote", + "syn", +] + +[[package]] +name = "clap_lex" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c8d4a3bb8b1e0c1050499d1815f5ab16d04f0959b233085fb31653fbfc9d98f9" + +[[package]] +name = "colorchoice" +version = "1.0.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1d07550c9036bf2ae0c684c4297d503f838287c83c53686d05370d0e139ae570" + [[package]] name = "const-random" version = "0.1.18" @@ -738,6 +834,12 @@ version = "3.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8bb03732005da905c88227371639bf1ad885cc712789c011c31c5fb3ab3ccf02" +[[package]] +name = "is_terminal_polyfill" +version = "1.70.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a6cb138bb79a146c1bd460005623e142ef0181e3d0219cb493e02f7d08a35695" + [[package]] name = "itoa" version = "1.0.18" @@ -1013,6 +1115,12 @@ version = "1.21.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9f7c3e4beb33f85d45ae3e3a1792185706c8e16d043238c593331cc7cd313b50" +[[package]] +name = "once_cell_polyfill" +version = "1.70.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "384b8ab6d37215f3c5301a95a4accb5d64aa607f1fcb26a11b5303878451b4fe" + [[package]] name = "ordered-float" version = "2.10.1" @@ -1477,6 +1585,12 @@ dependencies = [ "windows-sys", ] +[[package]] +name = "strsim" +version = "0.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7da8b5736845d9f2fcb837ea5d9e2628564b3b043a70948a3f0b778838c5fb4f" + [[package]] name = "syn" version = "2.0.117" @@ -1706,6 +1820,7 @@ dependencies = [ "byteorder", "bytes", "chrono", + "clap", "http-body-util", "hyper", "lz4_flex", @@ -1725,6 +1840,12 @@ dependencies = [ "zstd", ] +[[package]] +name = "utf8parse" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821" + [[package]] name = "uuid" version = "1.23.1" diff --git a/Cargo.toml b/Cargo.toml index 78d6aa2..7c675c2 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -11,6 +11,7 @@ anyhow = "1.0.102" axum = "0.8.9" bincode = "1.3.3" blake3 = "1.5" +clap = { version = "4", features = ["derive"] } byteorder = "1.5.0" bytes = "1.11.1" chrono = { version = "0.4.44", features = ["serde"] } diff --git a/README.md b/README.md index bc324eb..9e379a6 100644 --- a/README.md +++ b/README.md @@ -73,10 +73,26 @@ Ingest response counts `accepted`, `duplicates` (same id + same payload), `confl ```bash cargo build cargo test -cargo run # HTTP server on 127.0.0.1:8080 +cargo run # HTTP server on 127.0.0.1:8080 (default) +cargo run -- --help # see admin subcommands ``` -Configuration is currently hardcoded in `Config::default()` (db_root `./data`, 64 MiB memtable, 1M dedupe entries). +Configuration is currently hardcoded in `Config::default()` (db_root `./data`, 64 MiB memtable, 1M dedupe entries). The `--db-root` flag overrides the data directory. + +## Admin CLI + +The `usagedb` binary doubles as an admin tool. Subcommands operate on the on-disk state without needing the HTTP server running: + +``` +usagedb serve # HTTP server (default) +usagedb check [--deep] # manifest summary; --deep verifies every segment +usagedb rebuild-rollups --from --to # drop rollups + rewind watermark +usagedb inspect-segment # metadata + sample rows +usagedb verify-period --account --from --to # raw vs rollup drift +usagedb export-parquet # dump every raw segment to Parquet +``` + +All commands accept `--db-root ` (default `./data`). Admin commands assume the server is **not** running concurrently — file locking to prevent that is on the backlog. ## Durability contract diff --git a/src/admin.rs b/src/admin.rs new file mode 100644 index 0000000..863bf25 --- /dev/null +++ b/src/admin.rs @@ -0,0 +1,325 @@ +//! Admin CLI commands. Each `cmd_*` function returns a formatted string +//! (or writes side effects + returns a short status). The binary entry +//! point in `src/main.rs` parses CLI args, calls one of these, and +//! prints the returned string. Splitting it this way keeps the +//! commands testable without spawning a subprocess. + +use std::path::{Path, PathBuf}; +use std::sync::Arc; +use std::time::Duration; + +use chrono::DateTime; +use tokio::sync::{Mutex, RwLock}; + +use crate::export::parquet::export_raw_segments; +use crate::ingest::dedupe::HotDedupe; +use crate::ingest::memtable::Memtable; +use crate::ingest::wal::Wal; +use crate::rollup::worker::RollupWorker; +use crate::runtime::config::Config; +use crate::runtime::state::{AppState, AppStateInner}; +use crate::storage::manifest::{Manifest, SegmentKind}; +use crate::storage::segment_reader::RawSegmentReader; + +/// Build a read-only AppState by loading the manifest directly. Does NOT +/// run the full Recovery flow (no segment-scan dedupe rebuild, no WAL +/// replay) — admin commands operate on the on-disk state as-is. Cheaper +/// than `run_startup_recovery` for one-shot commands. +pub fn open_state_for_admin(config: Config) -> anyhow::Result { + let manifest = Manifest::load(&config.db_root)?.ok_or_else(|| { + anyhow::anyhow!( + "no manifest found at {:?} — run the server at least once to initialize the DB", + config.db_root + ) + })?; + let wal_dir = config.db_root.join("wal"); + std::fs::create_dir_all(&wal_dir)?; + let wal = Wal::open(wal_dir, manifest.last_sealed_wal_id)?; + let (flush_sender, _r) = tokio::sync::mpsc::channel(4); + Ok(Arc::new(AppStateInner { + config, + // Minimal dedupe — admin commands don't ingest. + dedupe: Mutex::new(HotDedupe::new(1)), + wal: Mutex::new(wal), + memtable: Mutex::new(Memtable::new()), + manifest: RwLock::new(manifest), + flush_sender, + })) +} + +/// `usagedb check [--deep]` — print manifest summary; with --deep, also +/// open every segment to verify checksum + format. +pub async fn cmd_check(state: AppState, deep: bool) -> anyhow::Result { + let manifest = state.manifest.read().await; + let mut out = String::new(); + out.push_str(&format!("Database root: {:?}\n", state.config.db_root)); + out.push_str(&format!("Manifest generation: {}\n", manifest.current_generation)); + out.push_str(&format!("Bucket count: {}\n", manifest.bucket_count)); + out.push_str(&format!("Raw segments: {}\n", manifest.raw_segments.len())); + out.push_str(&format!("Rollup segments: {}\n", manifest.rollup_segments.len())); + out.push_str(&format!( + "Watermark: {} ms ({})\n", + manifest.watermarks.hourly_rollup_ms, + format_ms(manifest.watermarks.hourly_rollup_ms), + )); + out.push_str(&format!("Last sealed WAL: {}\n", manifest.last_sealed_wal_id)); + out.push_str(&format!( + "Pending replacements: {}\n", + manifest.compacted_replacements.len() + )); + + if deep { + out.push_str("\nVerifying segment files...\n"); + let mut errors = 0usize; + for meta in manifest.raw_segments.iter().chain(manifest.rollup_segments.iter()) { + let ext = match meta.kind { + SegmentKind::Raw => "seg", + SegmentKind::Rollup => "rseg", + }; + let path = state.config.db_root.join(format!("{}.{}", meta.segment_id, ext)); + let result = match meta.kind { + SegmentKind::Raw => verify_raw_segment(&path), + SegmentKind::Rollup => verify_rollup_segment(&path, meta.checksum), + }; + match result { + Ok(()) => out.push_str(&format!(" {} OK\n", meta.segment_id)), + Err(e) => { + errors += 1; + out.push_str(&format!(" {} ERROR: {}\n", meta.segment_id, e)); + } + } + } + if errors > 0 { + out.push_str(&format!("\n{} segment(s) failed verification.\n", errors)); + return Err(anyhow::anyhow!("{} segment(s) failed verification", errors)); + } + out.push_str("\nAll segments verified.\n"); + } + Ok(out) +} + +fn verify_raw_segment(path: &Path) -> anyhow::Result<()> { + // RawSegmentReader::new validates magic + checksum + end magic + version + // + decodes every column. If it returns Ok, the segment is structurally + // sound. + RawSegmentReader::new(path.to_path_buf())?; + Ok(()) +} + +fn verify_rollup_segment(path: &Path, expected_checksum: u64) -> anyhow::Result<()> { + use crate::rollup::reader::RollupSegmentReader; + RollupSegmentReader::open(path.to_path_buf(), expected_checksum)?; + Ok(()) +} + +/// `usagedb rebuild-rollups --from --to` — drop affected rollups + rewind +/// watermark; next server tick refills. +pub async fn cmd_rebuild_rollups( + state: AppState, + from: &str, + to: &str, +) -> anyhow::Result { + let from_ms = parse_rfc3339_ms(from, "from")?; + let to_ms = parse_rfc3339_ms(to, "to")?; + + let worker = RollupWorker::new( + state.clone(), + state.config.rollup_safety_lag_ms, + Duration::from_secs(60), + i64::MAX, // disable force-drain — admin command, no background ticks + ); + let dropped = worker.rebuild_rollups(from_ms, to_ms).await?; + Ok(format!( + "Dropped {} rollup segment(s) overlapping [{}, {}).\n\ + Watermark rewound to {}.\n\ + Run the server (or wait for its rollup worker) to refill the gap from raw segments.\n", + dropped, from, to, from + )) +} + +/// `usagedb inspect-segment ` — print segment metadata + a sample of rows. +pub async fn cmd_inspect_segment(state: AppState, segment_id: &str) -> anyhow::Result { + let manifest = state.manifest.read().await; + let raw_match = manifest.raw_segments.iter().find(|s| s.segment_id == segment_id); + let rollup_match = manifest.rollup_segments.iter().find(|s| s.segment_id == segment_id); + + let meta = raw_match.or(rollup_match).ok_or_else(|| { + anyhow::anyhow!("segment {} not found in manifest", segment_id) + })?; + let is_raw = matches!(meta.kind, SegmentKind::Raw); + + let mut out = String::new(); + out.push_str(&format!("Segment: {}\n", meta.segment_id)); + out.push_str(&format!(" Kind: {:?}\n", meta.kind)); + out.push_str(&format!(" Bucket: {}\n", meta.bucket)); + out.push_str(&format!(" Rows: {}\n", meta.row_count)); + out.push_str(&format!( + " Timestamp range: [{}, {}]\n", + format_ms(meta.min_timestamp_ms), + format_ms(meta.max_timestamp_ms), + )); + if let Some(min) = &meta.min_account_id { + out.push_str(&format!(" Account range: {}", min.0)); + if let Some(max) = &meta.max_account_id { + out.push_str(&format!(" .. {}", max.0)); + } + out.push('\n'); + } + out.push_str(&format!( + " Products: {}\n", + meta.product_ids + .iter() + .map(|p| p.0.as_str()) + .collect::>() + .join(", ") + )); + out.push_str(&format!( + " Meters: {}\n", + meta.meter_ids + .iter() + .map(|m| m.0.as_str()) + .collect::>() + .join(", ") + )); + out.push_str(&format!( + " Models: {}\n", + meta.model_ids + .iter() + .map(|m| m.0.as_str()) + .collect::>() + .join(", ") + )); + out.push_str(&format!( + " Quantity sum: {}\n", + meta.quantity_sum.map(|q| q.to_string()).unwrap_or_else(|| "(none)".into()) + )); + out.push_str(&format!(" Checksum: {:#018x}\n", meta.checksum)); + if !meta.input_segment_ids.is_empty() { + out.push_str(" Input segments (rollup provenance):\n"); + for id in &meta.input_segment_ids { + out.push_str(&format!(" - {}\n", id)); + } + } + + if is_raw { + let path = state.config.db_root.join(format!("{}.seg", segment_id)); + let mut reader = RawSegmentReader::new(path)?; + out.push_str("\n First few rows:\n"); + let mut shown = 0; + while let Some(e) = reader.read_next()? { + if shown >= 5 { + break; + } + out.push_str(&format!( + " {} {:?} acc={} product={} meter={} model={:?} ts={} qty={}\n", + e.event_id.0, + e.kind, + e.account_id.0, + e.product_id.0, + e.meter_id.0, + e.model_id.as_ref().map(|m| m.0.as_str()), + e.timestamp_ms, + e.quantity, + )); + shown += 1; + } + if shown == 0 { + out.push_str(" (empty segment)\n"); + } + } + Ok(out) +} + +/// `usagedb verify-period --account --from --to` — call the same logic +/// as the /verify endpoint and pretty-print the result. +pub async fn cmd_verify_period( + state: AppState, + account: &str, + from: &str, + to: &str, +) -> anyhow::Result { + use crate::query::executor::execute_plan; + use crate::query::plan::{AggregationFunction, QueryPlan, QuerySource}; + use std::collections::HashMap; + + let from_ms = parse_rfc3339_ms(from, "from")?; + let to_ms = parse_rfc3339_ms(to, "to")?; + + let mut metrics = HashMap::new(); + metrics.insert("quantity".to_string(), AggregationFunction::Sum); + let plan_raw = QueryPlan { + source: QuerySource::RawEvents, + account_id: Some(account.to_string()), + from_ms, + to_ms, + filters: vec![], + group_by: vec![], + metrics: metrics.clone(), + limit: None, + }; + let plan_rollup = QueryPlan { + source: QuerySource::RollupHourly, + ..plan_raw.clone() + }; + let raw_total = extract_sum(&execute_plan(&state, &plan_raw).await); + let rollup_total = extract_sum(&execute_plan(&state, &plan_rollup).await); + let drift = raw_total.saturating_sub(rollup_total); + let watermark_ms = state.manifest.read().await.watermarks.hourly_rollup_ms; + let period_sealed = to_ms <= watermark_ms; + + let mut out = String::new(); + out.push_str(&format!("Account: {}\n", account)); + out.push_str(&format!( + "Range: [{}, {})\n", + from, to + )); + out.push_str(&format!( + "Watermark: {} ({})\n", + watermark_ms, + format_ms(watermark_ms), + )); + out.push_str(&format!("Period sealed: {}\n", period_sealed)); + out.push_str(&format!("Raw total: {}\n", raw_total)); + out.push_str(&format!("Rollup total: {}\n", rollup_total)); + out.push_str(&format!( + "Drift: {} {}\n", + drift, + if drift == 0 { "(OK)" } else { "(MISMATCH)" } + )); + Ok(out) +} + +/// `usagedb export-parquet ` — dump every raw segment in the +/// manifest to one Parquet file. +pub async fn cmd_export_parquet(state: AppState, output: &Path) -> anyhow::Result { + let stats = export_raw_segments(&state, output).await?; + Ok(format!( + "Exported {} events from {} segment(s) to {:?}\n", + stats.events_exported, stats.segments_read, output + )) +} + +fn parse_rfc3339_ms(s: &str, field: &str) -> anyhow::Result { + DateTime::parse_from_rfc3339(s) + .map(|dt| dt.timestamp_millis()) + .map_err(|e| anyhow::anyhow!("invalid `{}` (not RFC 3339): {}", field, e)) +} + +fn extract_sum(result: &[serde_json::Value]) -> i128 { + result + .iter() + .filter_map(|v| v.get("quantity")) + .filter_map(|v| v.as_str()) + .filter_map(|s| s.parse().ok()) + .next() + .unwrap_or(0) +} + +fn format_ms(ms: i64) -> String { + DateTime::from_timestamp_millis(ms) + .map(|dt| dt.format("%Y-%m-%dT%H:%M:%SZ").to_string()) + .unwrap_or_else(|| "".into()) +} + +#[allow(dead_code)] +fn _force_pathbuf_use(_: &PathBuf) {} diff --git a/src/lib.rs b/src/lib.rs index 0df175f..d27e3eb 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,3 +1,4 @@ +pub mod admin; pub mod api; pub mod compact; pub mod export; diff --git a/src/main.rs b/src/main.rs index 072a417..8b00b9f 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,26 +1,120 @@ -use usagedb::runtime::config::Config; -use usagedb::runtime::state::{AppStateInner, AppState}; -use usagedb::ingest::wal::Wal; +use std::path::PathBuf; +use std::sync::Arc; +use std::time::Duration; + +use clap::{Parser, Subcommand}; +use tokio::sync::{Mutex, Notify, RwLock, mpsc}; +use tracing::{error, info}; + +use usagedb::admin::{ + cmd_check, cmd_export_parquet, cmd_inspect_segment, cmd_rebuild_rollups, cmd_verify_period, + open_state_for_admin, +}; +use usagedb::api::http_server::start_server; +use usagedb::compact::worker::CompactionWorker; use usagedb::ingest::flusher::FlusherWorker; +use usagedb::ingest::wal::Wal; use usagedb::rollup::worker::RollupWorker; -use usagedb::compact::worker::CompactionWorker; -use usagedb::api::http_server::start_server; +use usagedb::runtime::config::Config; use usagedb::runtime::recovery::Recovery; +use usagedb::runtime::state::{AppState, AppStateInner}; + +#[derive(Parser, Debug)] +#[command( + name = "usagedb", + about = "Embedded append-only usage database for AI billing", + version +)] +struct Cli { + /// Database root directory. + #[arg(long, global = true, default_value = "./data")] + db_root: PathBuf, + + #[command(subcommand)] + command: Option, +} -use tokio::sync::{RwLock, Mutex, Notify, mpsc}; -use std::sync::Arc; -use std::time::Duration; -use tracing::{info, error}; +#[derive(Subcommand, Debug)] +enum Command { + /// Run the HTTP server (default if no subcommand is given). + Serve, + /// Print manifest summary. With `--deep`, also open every segment and + /// verify checksums + structure. + Check { + #[arg(long, default_value_t = false)] + deep: bool, + }, + /// Drop rollup segments overlapping `[from, to)` and rewind the watermark + /// to `from`. The next server tick refills the gap from raw segments. + RebuildRollups { + #[arg(long)] + from: String, + #[arg(long)] + to: String, + }, + /// Read a specific segment and print its metadata + sample rows. + InspectSegment { + segment_id: String, + }, + /// Compute raw vs rollup totals + drift for an account over a period. + VerifyPeriod { + #[arg(long)] + account: String, + #[arg(long)] + from: String, + #[arg(long)] + to: String, + }, + /// Export every raw segment to a single Parquet file. + ExportParquet { + output: PathBuf, + }, +} #[tokio::main] async fn main() -> Result<(), Box> { tracing_subscriber::fmt::init(); - info!("Starting usageDb server..."); - let config = Config::default(); + let cli = Cli::parse(); + let mut config = Config::default(); + config.db_root = cli.db_root; + + match cli.command.unwrap_or(Command::Serve) { + Command::Serve => run_server(config).await?, + Command::Check { deep } => { + let state = open_state_for_admin(config)?; + let output = cmd_check(state, deep).await?; + print!("{}", output); + } + Command::RebuildRollups { from, to } => { + let state = open_state_for_admin(config)?; + let output = cmd_rebuild_rollups(state, &from, &to).await?; + print!("{}", output); + } + Command::InspectSegment { segment_id } => { + let state = open_state_for_admin(config)?; + let output = cmd_inspect_segment(state, &segment_id).await?; + print!("{}", output); + } + Command::VerifyPeriod { account, from, to } => { + let state = open_state_for_admin(config)?; + let output = cmd_verify_period(state, &account, &from, &to).await?; + print!("{}", output); + } + Command::ExportParquet { output } => { + let state = open_state_for_admin(config)?; + let output_msg = cmd_export_parquet(state, &output).await?; + print!("{}", output_msg); + } + } + Ok(()) +} + +async fn run_server(config: Config) -> Result<(), Box> { + info!("Starting usageDb server..."); std::fs::create_dir_all(&config.db_root)?; - // Run startup recovery: load manifest, clean tmp, replay WAL + // Run startup recovery: load manifest, clean tmp, replay WAL. let recovery = Recovery::new(config.db_root.clone()); let mut recovery_result = match recovery.run_startup_recovery(config.dedupe_capacity) { Ok(r) => r, @@ -89,7 +183,7 @@ async fn main() -> Result<(), Box> { start_server(state.clone()).await?; - // Shutdown flush (review P1 #6): drain the memtable + rotate the WAL + // Shutdown drain (review P1 #6): drain the memtable + rotate the WAL // so any events accumulated since the last size-based flush become // durable raw segments instead of staying stranded in WAL files. { diff --git a/tests/admin.rs b/tests/admin.rs new file mode 100644 index 0000000..7684a3b --- /dev/null +++ b/tests/admin.rs @@ -0,0 +1,336 @@ +//! Tests for the admin CLI commands. Drives the `cmd_*` functions +//! directly (the CLI binary just parses args and delegates), so we +//! don't have to spawn subprocesses. + +use std::path::PathBuf; +use std::sync::Arc; + +use tokio::sync::{Mutex, RwLock}; + +use usagedb::admin::{ + cmd_check, cmd_export_parquet, cmd_inspect_segment, cmd_rebuild_rollups, cmd_verify_period, + open_state_for_admin, +}; +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::{EventKind, UsageEvent}; +use usagedb::model::ids::{ + AccountId, EventId, MeterId, ModelId, ProductId, SourceId, SubscriptionId, Unit, + bucket_for_account, +}; +use usagedb::rollup::worker::RollupWorker; +use usagedb::runtime::config::Config; +use usagedb::runtime::state::{AppState, AppStateInner}; +use usagedb::storage::manifest::Manifest; +use usagedb::storage::segment_writer::RawSegmentWriter; + +const HOUR_MS: i64 = 3_600_000; + +fn tmp_root() -> PathBuf { + let dir = tempfile::tempdir().expect("tempdir"); + let p = dir.path().to_path_buf(); + std::mem::forget(dir); + p +} + +fn make_event(id: &str, account: &str, ts: i64, qty: i128) -> UsageEvent { + 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("test".into()), + model_id: Some(ModelId("m1".into())), + dimensions: SmallDimensions::default(), + ingested_at_ms: ts, + } +} + +/// Build a live AppState with a manifest already on disk so the +/// admin-side `open_state_for_admin` finds something to load. Returns +/// the live AppState too so the test can commit segments through the +/// same instance. +async fn setup_db_with_segments(events: Vec<(Vec, u32)>) -> (PathBuf, AppState) { + let root = tmp_root(); + let config = Config { + db_root: root.clone(), + default_bucket_count: 2, + ..Config::default() + }; + std::fs::create_dir_all(&config.db_root).unwrap(); + let wal = Wal::open(root.join("wal"), 0).unwrap(); + let manifest = Manifest { bucket_count: 2, ..Manifest::default() }; + let (flush_sender, _r) = tokio::sync::mpsc::channel(4); + let state: AppState = 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, + }); + + // Persist the manifest so `open_state_for_admin` finds something + // even when the test commits zero segments. + { + let mut m = state.manifest.write().await; + m.save(&state.config.db_root).unwrap(); + } + + for (evts, bucket) in events { + let segment_id = format!("raw_{}", uuid::Uuid::new_v4().simple()); + let path = state.config.db_root.join(format!("{}.seg", segment_id)); + let mut w = RawSegmentWriter::new(path).unwrap(); + for e in &evts { w.write_event(e).unwrap(); } + let (_rows, checksum) = w.finish().unwrap(); + let meta = build_segment_meta(&segment_id, &evts, bucket, checksum); + let mut m = state.manifest.write().await; + m.raw_segments.push(meta); + m.save(&state.config.db_root).unwrap(); + } + + (root, state) +} + +// ========================================================================= +// `check` +// ========================================================================= + +#[tokio::test] +async fn check_reports_manifest_summary() { + let (root, _live) = setup_db_with_segments(vec![ + (vec![make_event("a", "acc", 1000, 10)], 0), + (vec![make_event("b", "acc", 2000, 20)], 0), + ]) + .await; + drop(_live); // close the live state so admin can reopen cleanly + + let config = Config { db_root: root.clone(), ..Config::default() }; + let state = open_state_for_admin(config).unwrap(); + let out = cmd_check(state, false).await.unwrap(); + assert!(out.contains("Raw segments: 2"), "{}", out); + assert!(out.contains("Rollup segments: 0"), "{}", out); + assert!(out.contains("Manifest generation:"), "{}", out); +} + +#[tokio::test] +async fn check_deep_passes_for_valid_segments() { + let (root, _live) = setup_db_with_segments(vec![ + (vec![make_event("a", "acc", 1000, 10)], 0), + ]) + .await; + drop(_live); + + let config = Config { db_root: root.clone(), ..Config::default() }; + let state = open_state_for_admin(config).unwrap(); + let out = cmd_check(state, true).await.unwrap(); + assert!(out.contains("All segments verified"), "deep check should pass: {}", out); +} + +#[tokio::test] +async fn check_deep_fails_for_corrupt_segment() { + let (root, _live) = setup_db_with_segments(vec![ + (vec![make_event("a", "acc", 1000, 10)], 0), + ]) + .await; + drop(_live); + + // Corrupt the only segment file. + let seg = std::fs::read_dir(&root) + .unwrap() + .find_map(|e| { + let e = e.ok()?; + let n = e.file_name(); + let n = n.to_str()?; + if n.starts_with("raw_") && n.ends_with(".seg") { + Some(e.path()) + } else { + None + } + }) + .expect("segment file"); + let mut bytes = std::fs::read(&seg).unwrap(); + bytes[40] ^= 0xFF; + std::fs::write(&seg, bytes).unwrap(); + + let config = Config { db_root: root.clone(), ..Config::default() }; + let state = open_state_for_admin(config).unwrap(); + let err = cmd_check(state, true).await.unwrap_err(); + assert!(err.to_string().contains("failed verification"), "{}", err); +} + +#[tokio::test] +async fn check_on_fresh_db_errors_with_clear_message() { + let root = tmp_root(); + std::fs::create_dir_all(&root).unwrap(); + let config = Config { db_root: root, ..Config::default() }; + let err = match open_state_for_admin(config) { + Ok(_) => panic!("expected error on a fresh DB"), + Err(e) => e, + }; + assert!( + err.to_string().contains("no manifest found"), + "should give a clear error: {}", + err + ); +} + +// ========================================================================= +// `inspect-segment` +// ========================================================================= + +#[tokio::test] +async fn inspect_segment_prints_meta_and_sample_rows() { + let (root, live) = setup_db_with_segments(vec![ + (vec![make_event("evt_a", "acc_q", 1000, 100), make_event("evt_b", "acc_q", 2000, 200)], 0), + ]) + .await; + + let seg_id = { + let m = live.manifest.read().await; + m.raw_segments[0].segment_id.clone() + }; + drop(live); + + let config = Config { db_root: root.clone(), ..Config::default() }; + let state = open_state_for_admin(config).unwrap(); + let out = cmd_inspect_segment(state, &seg_id).await.unwrap(); + + assert!(out.contains(&seg_id), "{}", out); + assert!(out.contains("Rows: 2"), "{}", out); + assert!(out.contains("Bucket: 0"), "{}", out); + assert!(out.contains("evt_a"), "should show sample rows: {}", out); + assert!(out.contains("evt_b"), "{}", out); +} + +#[tokio::test] +async fn inspect_segment_errors_for_unknown_id() { + let (root, _live) = setup_db_with_segments(vec![]).await; + drop(_live); + + let config = Config { db_root: root.clone(), ..Config::default() }; + let state = open_state_for_admin(config).unwrap(); + let err = cmd_inspect_segment(state, "raw_bogus").await.unwrap_err(); + assert!(err.to_string().contains("not found in manifest"), "{}", err); +} + +// ========================================================================= +// `rebuild-rollups` +// ========================================================================= + +#[tokio::test] +async fn rebuild_rollups_drops_and_rewinds() { + let bucket = bucket_for_account(&AccountId("acc_r".into()), 2); + let h = 10 * HOUR_MS; + let (root, live) = setup_db_with_segments(vec![ + (vec![make_event("a", "acc_r", h + 1, 10), make_event("b", "acc_r", h + 2, 20)], bucket), + ]) + .await; + + // Seal hour 10 into a rollup via a real worker tick. + let worker = RollupWorker::new( + live.clone(), + 0, + std::time::Duration::from_secs(30), + i64::MAX, + ); + worker.tick(h + HOUR_MS + 1).await.unwrap(); + assert_eq!(live.manifest.read().await.rollup_segments.len(), 1); + drop(live); + + let config = Config { db_root: root.clone(), ..Config::default() }; + let state = open_state_for_admin(config).unwrap(); + let out = cmd_rebuild_rollups( + state.clone(), + "1970-01-01T00:00:00Z", + "2030-01-01T00:00:00Z", + ) + .await + .unwrap(); + assert!(out.contains("Dropped 1 rollup segment"), "{}", out); + + let m = state.manifest.read().await; + assert!(m.rollup_segments.is_empty()); + assert_eq!(m.watermarks.hourly_rollup_ms, 0); +} + +#[tokio::test] +async fn rebuild_rollups_rejects_invalid_dates() { + let (root, _live) = setup_db_with_segments(vec![]).await; + drop(_live); + + let config = Config { db_root: root, ..Config::default() }; + let state = open_state_for_admin(config).unwrap(); + let err = cmd_rebuild_rollups(state, "not-a-date", "2030-01-01T00:00:00Z") + .await + .unwrap_err(); + assert!(err.to_string().contains("invalid"), "{}", err); +} + +// ========================================================================= +// `verify-period` +// ========================================================================= + +#[tokio::test] +async fn verify_period_reports_match_for_sealed_rollup() { + let bucket = bucket_for_account(&AccountId("acc_v".into()), 2); + let h = 20 * HOUR_MS; + let (root, live) = setup_db_with_segments(vec![ + (vec![make_event("a", "acc_v", h + 1, 30), make_event("b", "acc_v", h + 2, 40)], bucket), + ]) + .await; + let worker = RollupWorker::new( + live.clone(), + 0, + std::time::Duration::from_secs(30), + i64::MAX, + ); + worker.tick(h + HOUR_MS + 1).await.unwrap(); + drop(live); + + let config = Config { db_root: root, ..Config::default() }; + let state = open_state_for_admin(config).unwrap(); + let out = cmd_verify_period( + state, + "acc_v", + "1970-01-01T00:00:00Z", + "2030-01-01T00:00:00Z", + ) + .await + .unwrap(); + assert!(out.contains("Raw total: 70"), "{}", out); + assert!(out.contains("Rollup total: 70"), "{}", out); + assert!(out.contains("Drift: 0 (OK)"), "{}", out); +} + +// ========================================================================= +// `export-parquet` +// ========================================================================= + +#[tokio::test] +async fn export_parquet_writes_file() { + let (root, _live) = setup_db_with_segments(vec![ + (vec![make_event("a", "acc", 1000, 5)], 0), + (vec![make_event("b", "acc", 2000, 7)], 0), + ]) + .await; + drop(_live); + + let out_path = root.join("export.parquet"); + let config = Config { db_root: root.clone(), ..Config::default() }; + let state = open_state_for_admin(config).unwrap(); + let msg = cmd_export_parquet(state, &out_path).await.unwrap(); + assert!(msg.contains("Exported 2 events"), "{}", msg); + assert!(msg.contains("2 segment(s)"), "{}", msg); + assert!(out_path.exists()); + assert!(std::fs::metadata(&out_path).unwrap().len() > 0); +}