diff --git a/src/api/http_server.rs b/src/api/http_server.rs index f6060bb..decffb4 100644 --- a/src/api/http_server.rs +++ b/src/api/http_server.rs @@ -737,42 +737,49 @@ async fn handle_close_period( let (frozen_quantity, frozen_event_count) = snapshot_period_totals(&state, &account_id, from_ms, to_ms).await; - let mut manifest = state.manifest.write().await; - // Re-check inside the write lock — another caller might have raced us. - if manifest.closed_periods.iter().any(|p| { - p.account_id == account_id && p.year == year && p.month == month - }) { - return Ok(Json(serde_json::json!({ + // `commit_manifest_if` re-checks inside the write lock for race + // safety; returning `None` from the closure (already closed) skips + // the save and leaves the manifest unchanged (review P0 #2). + let closed_at_ms = now_ms(); + let committed = state + .commit_manifest_if(|manifest| { + if manifest.closed_periods.iter().any(|p| { + p.account_id == account_id && p.year == year && p.month == month + }) { + return None; + } + let watermark_at_close_ms = manifest.watermarks.hourly_rollup_ms; + let entry = ClosedPeriod { + account_id: account_id.clone(), + year, + month, + closed_at_ms, + frozen_quantity: Some(frozen_quantity), + frozen_event_count: Some(frozen_event_count), + watermark_at_close_ms: Some(watermark_at_close_ms), + }; + manifest.closed_periods.push(entry.clone()); + Some((entry, watermark_at_close_ms)) + }) + .await + .map_err(|e| AppError(anyhow::anyhow!("manifest save: {}", e)))?; + + match committed { + Some((entry, watermark_at_close_ms)) => Ok(Json(serde_json::json!({ + "account_id": account_id, + "period": period, + "state": "Closed", + "closed_at_ms": closed_at_ms, + "watermark_at_close_ms": watermark_at_close_ms, + "frozen": frozen_json(&entry), + }))), + None => Ok(Json(serde_json::json!({ "account_id": account_id, "period": period, "state": "Closed", "already_closed": true, - }))); + }))), } - let closed_at_ms = now_ms(); - let watermark_at_close_ms = manifest.watermarks.hourly_rollup_ms; - let entry = ClosedPeriod { - account_id: account_id.clone(), - year, - month, - closed_at_ms, - frozen_quantity: Some(frozen_quantity), - frozen_event_count: Some(frozen_event_count), - watermark_at_close_ms: Some(watermark_at_close_ms), - }; - manifest.closed_periods.push(entry.clone()); - manifest - .save(&state.config.db_root) - .map_err(|e| AppError(anyhow::anyhow!("manifest save: {}", e)))?; - - Ok(Json(serde_json::json!({ - "account_id": account_id, - "period": period, - "state": "Closed", - "closed_at_ms": closed_at_ms, - "watermark_at_close_ms": watermark_at_close_ms, - "frozen": frozen_json(&entry), - }))) } /// Compute period bounds `[from_ms, to_ms)` from (year, month) using UTC. @@ -859,23 +866,25 @@ async fn handle_reopen_period( let (year, month) = parse_period(&period) .map_err(|e| AppError(anyhow::anyhow!("invalid period: {}", e)))?; - let mut manifest = state.manifest.write().await; - let before = manifest.closed_periods.len(); - manifest.closed_periods.retain(|p| { - !(p.account_id == account_id && p.year == year && p.month == month) - }); - let removed = before - manifest.closed_periods.len(); - if removed > 0 { - manifest - .save(&state.config.db_root) - .map_err(|e| AppError(anyhow::anyhow!("manifest save: {}", e)))?; - } + // `commit_manifest_if` only saves when something was removed, + // matching the prior `if removed > 0` save guard (review P0 #2). + let removed = state + .commit_manifest_if(|manifest| { + let before = manifest.closed_periods.len(); + manifest.closed_periods.retain(|p| { + !(p.account_id == account_id && p.year == year && p.month == month) + }); + let n = before - manifest.closed_periods.len(); + if n > 0 { Some(n) } else { None } + }) + .await + .map_err(|e| AppError(anyhow::anyhow!("manifest save: {}", e)))?; Ok(Json(serde_json::json!({ "account_id": account_id, "period": period, "state": "Open", - "removed": removed > 0, + "removed": removed.is_some(), }))) } diff --git a/src/compact/worker.rs b/src/compact/worker.rs index 1a84810..ad75b85 100644 --- a/src/compact/worker.rs +++ b/src/compact/worker.rs @@ -167,49 +167,43 @@ impl CompactionWorker { let _ = std::fs::remove_file(&idx); } - // Atomically swap old → new in the manifest. - let committed = { - let mut manifest = self.state.manifest.write().await; - - // Defensive: every input must still be in raw_segments. If a - // concurrent worker (none today, but future-proofing) already - // removed any, abort this plan and clean up the output. - let input_set: HashSet<&String> = plan.segment_ids.iter().collect(); + // Atomically swap old → new in the manifest. `commit_manifest_if` + // mutates a clone and only publishes after the on-disk save + // succeeds; the closure returns `None` if a concurrent worker + // already removed any of our planned inputs, leaving the + // manifest untouched (review P0 #2). + let input_set_owned: HashSet = plan.segment_ids.iter().cloned().collect(); + let committed = match self.state.commit_manifest_if(|manifest| { + let input_set: HashSet<&String> = input_set_owned.iter().collect(); let still_present = manifest .raw_segments .iter() .filter(|s| input_set.contains(&s.segment_id)) .count(); if still_present != plan.segment_ids.len() { - warn!( - "Compaction plan inputs no longer all present in manifest ({}/{}); aborting", - still_present, plan.segment_ids.len() - ); - drop(manifest); - let _ = std::fs::remove_file(&output_path); - return Ok(false); + return None; } - manifest.raw_segments.retain(|s| !input_set.contains(&s.segment_id)); - manifest.raw_segments.push(new_meta); + manifest.raw_segments.push(new_meta.clone()); manifest.compacted_replacements.push(ReplacementRecord { old_segments: plan.segment_ids.clone(), new_segments: vec![output_id.clone()], committed_at_ms: now_ms, }); - - match manifest.save(&self.state.config.db_root) { - Ok(()) => true, - Err(e) => { - error!("Compaction manifest save failed: {} — cleaning up output", e); - // Roll back the in-memory mutation by re-reading from disk - // would be ideal; for now we surface the error and let - // the operator deal with it. The output file we just - // wrote is unreferenced; remove it. - drop(manifest); - let _ = std::fs::remove_file(&output_path); - return Err(anyhow::anyhow!("manifest save failed: {}", e)); - } + Some(still_present) + }).await { + Ok(Some(_)) => true, + Ok(None) => { + warn!( + "Compaction plan inputs no longer all present in manifest; aborting" + ); + let _ = std::fs::remove_file(&output_path); + return Ok(false); + } + Err(e) => { + error!("Compaction manifest save failed: {} — cleaning up output", e); + let _ = std::fs::remove_file(&output_path); + return Err(anyhow::anyhow!("manifest save failed: {}", e)); } }; @@ -259,15 +253,15 @@ impl CompactionWorker { // Update the manifest to drop the finalized records. We compare by // (committed_at_ms, old_segments) — adequately unique since each // tick generates a distinct ms timestamp per record. - let mut manifest = self.state.manifest.write().await; let finalized_keys: HashSet<(i64, Vec)> = to_finalize .iter() .map(|r| (r.committed_at_ms, r.old_segments.clone())) .collect(); - manifest - .compacted_replacements - .retain(|r| !finalized_keys.contains(&(r.committed_at_ms, r.old_segments.clone()))); - if let Err(e) = manifest.save(&self.state.config.db_root) { + if let Err(e) = self.state.commit_manifest(|manifest| { + manifest + .compacted_replacements + .retain(|r| !finalized_keys.contains(&(r.committed_at_ms, r.old_segments.clone()))); + }).await { error!("Failed to persist replacement cleanup: {}", e); return Err(anyhow::anyhow!("manifest save failed: {}", e)); } diff --git a/src/ingest/flusher.rs b/src/ingest/flusher.rs index 07d25fa..89594e0 100644 --- a/src/ingest/flusher.rs +++ b/src/ingest/flusher.rs @@ -104,17 +104,18 @@ impl FlusherWorker { } // All bucket segments are durable. Commit the manifest atomically - // with the WAL seal pointer. - let save_result = { - let mut manifest = self.state.manifest.write().await; + // with the WAL seal pointer. `commit_manifest` mutates a clone + // and only publishes it after the on-disk save succeeds, so a + // save failure leaves the in-memory manifest unchanged + // (review P0 #2). + let save_result = self.state.commit_manifest(|manifest| { for meta in &new_metas { manifest.raw_segments.push(meta.clone()); } if sealed_wal_id > manifest.last_sealed_wal_id { manifest.last_sealed_wal_id = sealed_wal_id; } - manifest.save(&self.state.config.db_root) - }; + }).await; if let Err(e) = save_result { // Manifest save failed. The segments we wrote are orphaned; diff --git a/src/rollup/worker.rs b/src/rollup/worker.rs index 7791bb2..82ef32e 100644 --- a/src/rollup/worker.rs +++ b/src/rollup/worker.rs @@ -258,14 +258,19 @@ impl RollupWorker { } // Atomically commit: append rollup segments + advance watermark. + // `commit_manifest` only publishes the in-memory mutation after + // the on-disk save succeeds, so a failure here leaves the + // manifest unchanged and we can safely unlink the orphaned + // segment files we just wrote (review P0 #2). let segments_written = new_segments.len(); if segments_written > 0 || target_hour > current_watermark { - let mut manifest = self.state.manifest.write().await; - for (meta, _path) in &new_segments { - manifest.rollup_segments.push(meta.clone()); - } - manifest.watermarks.hourly_rollup_ms = target_hour; - if let Err(e) = manifest.save(&self.state.config.db_root) { + let metas_to_commit: Vec<_> = new_segments.iter().map(|(m, _)| m.clone()).collect(); + if let Err(e) = self.state.commit_manifest(|manifest| { + for meta in &metas_to_commit { + manifest.rollup_segments.push(meta.clone()); + } + manifest.watermarks.hourly_rollup_ms = target_hour; + }).await { error!("rollup: manifest save failed: {} — cleaning up {} segment files", e, new_segments.len()); for (_meta, path) in &new_segments { let _ = std::fs::remove_file(path); @@ -299,21 +304,21 @@ impl RollupWorker { /// query that snapshotted the manifest can still open them. Cleanup /// of orphaned rollup files is a future operability task. pub async fn rebuild_rollups(&self, from_ms: i64, to_ms: i64) -> anyhow::Result { - let mut manifest = self.state.manifest.write().await; - let before = manifest.rollup_segments.len(); - manifest.rollup_segments.retain(|s| { - // Keep segments that don't overlap [from_ms, to_ms). - !(s.min_timestamp_ms < to_ms && s.max_timestamp_ms >= from_ms) - }); - let dropped = before - manifest.rollup_segments.len(); - - if manifest.watermarks.hourly_rollup_ms > from_ms { - manifest.watermarks.hourly_rollup_ms = from_ms; - } - manifest.save(&self.state.config.db_root)?; + let (dropped, new_watermark) = self.state.commit_manifest(|manifest| { + let before = manifest.rollup_segments.len(); + manifest.rollup_segments.retain(|s| { + // Keep segments that don't overlap [from_ms, to_ms). + !(s.min_timestamp_ms < to_ms && s.max_timestamp_ms >= from_ms) + }); + let dropped = before - manifest.rollup_segments.len(); + if manifest.watermarks.hourly_rollup_ms > from_ms { + manifest.watermarks.hourly_rollup_ms = from_ms; + } + (dropped, manifest.watermarks.hourly_rollup_ms) + }).await?; info!( "Rebuild scheduled: dropped {} rollup segment(s) in [{}, {}); watermark rewound to {}", - dropped, from_ms, to_ms, manifest.watermarks.hourly_rollup_ms + dropped, from_ms, to_ms, new_watermark ); Ok(dropped) } diff --git a/src/runtime/config.rs b/src/runtime/config.rs index 078147e..a9ebabb 100644 --- a/src/runtime/config.rs +++ b/src/runtime/config.rs @@ -60,7 +60,11 @@ impl Default for Config { Self { db_root: PathBuf::from("./data"), max_memtable_size_bytes: 64 * 1024 * 1024, // 64 MB - http_bind_address: "127.0.0.1:8080".to_string(), + // Bind loopback by default; override with USAGEDB_HTTP_BIND_ADDRESS + // (e.g. 0.0.0.0:8080 when running in a container reached over a + // private network). + http_bind_address: std::env::var("USAGEDB_HTTP_BIND_ADDRESS") + .unwrap_or_else(|_| "127.0.0.1:8080".to_string()), dedupe_capacity: 1_000_000, default_bucket_count: 64, rollup_tick_interval_secs: 30, diff --git a/src/runtime/state.rs b/src/runtime/state.rs index 3ad0f52..08902fa 100644 --- a/src/runtime/state.rs +++ b/src/runtime/state.rs @@ -24,4 +24,47 @@ pub struct AppStateInner { pub flush_sender: tokio::sync::mpsc::Sender, } +impl AppStateInner { + /// Atomic manifest commit: clone → mutate the clone → save to disk + /// → publish in-memory. If the save fails, the in-memory manifest + /// is unchanged and the error is surfaced (review P0 #2). + /// + /// The previous pattern — take the write lock, mutate in place, + /// then call `save` — left the in-memory state ahead of disk on + /// save failure. Subsequent reads would see writes that hadn't + /// reached the on-disk manifest, and any cleanup of segment files + /// (which several callers do on save failure) would leave the + /// in-memory manifest pointing at files that no longer existed. + pub async fn commit_manifest(&self, op: F) -> std::io::Result + where + F: FnOnce(&mut Manifest) -> T, + { + let mut guard = self.manifest.write().await; + let mut next = guard.clone(); + let value = op(&mut next); + next.save(&self.config.db_root)?; + *guard = next; + Ok(value) + } + + /// Like `commit_manifest`, but only writes to disk when `op` + /// returns `Some`. Use this when the closure decides — under the + /// write lock, with race-safety — whether the change is actually + /// needed (e.g. close-period rechecks whether the period was + /// already closed by a racing caller). + pub async fn commit_manifest_if(&self, op: F) -> std::io::Result> + where + F: FnOnce(&mut Manifest) -> Option, + { + let mut guard = self.manifest.write().await; + let mut next = guard.clone(); + let value = op(&mut next); + if value.is_some() { + next.save(&self.config.db_root)?; + *guard = next; + } + Ok(value) + } +} + pub type AppState = Arc; diff --git a/src/storage/compression.rs b/src/storage/compression.rs index 7bf4af7..ab5cb29 100644 --- a/src/storage/compression.rs +++ b/src/storage/compression.rs @@ -1,4 +1,5 @@ use std::io::{Read, Write, Result as IoResult}; +use std::sync::OnceLock; use zstd::stream::{read::Decoder as ZstdDecoder, write::Encoder as ZstdEncoder}; use lz4_flex::frame::{FrameDecoder, FrameEncoder}; @@ -9,10 +10,28 @@ pub enum CompressionCodec { None, } +/// zstd compression level for raw segments. Optimized for space, not write +/// speed: segments are immutable and written once by the background flusher / +/// compaction workers (off the ingest ack path, which only fsyncs the WAL), +/// so a higher level trades background CPU for a smaller on-disk footprint +/// that pays off on every read. Override with `USAGEDB_ZSTD_LEVEL` (1-22); +/// defaults to 19. Decompression is level-independent, so raising the level +/// never breaks older segments. +fn zstd_level() -> i32 { + static LEVEL: OnceLock = OnceLock::new(); + *LEVEL.get_or_init(|| { + std::env::var("USAGEDB_ZSTD_LEVEL") + .ok() + .and_then(|v| v.trim().parse::().ok()) + .filter(|l| (1..=22).contains(l)) + .unwrap_or(19) + }) +} + pub fn compress(data: &[u8], codec: CompressionCodec) -> IoResult> { match codec { CompressionCodec::Zstd => { - let mut encoder = ZstdEncoder::new(Vec::new(), 3)?; + let mut encoder = ZstdEncoder::new(Vec::new(), zstd_level())?; encoder.write_all(data)?; encoder.finish() } diff --git a/tests/p0_manifest_atomic.rs b/tests/p0_manifest_atomic.rs new file mode 100644 index 0000000..9bdfacb --- /dev/null +++ b/tests/p0_manifest_atomic.rs @@ -0,0 +1,151 @@ +//! Regression test for the P0 manifest-atomicity finding from external review. +//! +//! Before the fix, several call sites took the manifest write lock, +//! mutated the manifest in place, then called `save`. If `save` failed, +//! the in-memory state still held the mutation but the on-disk manifest +//! did not. Queries would see writes that hadn't reached disk, and +//! crash-recovery would re-derive state from an older manifest while +//! the running process believed it had committed newer state. +//! +//! The fix routes every mutation through `AppStateInner::commit_manifest` +//! (or `commit_manifest_if`), which clones the manifest, mutates the +//! clone, saves the clone, and only publishes the clone in-memory after +//! the save succeeds. A save failure leaves both disk and memory +//! untouched. +//! +//! This test forces a save failure by replacing the `manifest/` +//! directory with a regular file after the initial save, then exercises +//! the helper and asserts the in-memory manifest is unchanged. + +use std::path::PathBuf; +use std::sync::Arc; + +use tokio::sync::{Mutex, RwLock, mpsc}; +use usagedb::ingest::dedupe::HotDedupe; +use usagedb::ingest::memtable::Memtable; +use usagedb::ingest::wal::Wal; +use usagedb::runtime::config::Config; +use usagedb::runtime::state::{AppState, AppStateInner}; +use usagedb::storage::manifest::Manifest; + +fn tmp_root() -> PathBuf { + let dir = tempfile::tempdir().expect("tempdir"); + let path = dir.path().to_path_buf(); + std::mem::forget(dir); + path +} + +fn build_state(db_root: PathBuf) -> AppState { + let config = Config { + db_root: db_root.clone(), + default_bucket_count: 4, + ..Config::default() + }; + std::fs::create_dir_all(&config.db_root).unwrap(); + let wal = Wal::open(db_root.join("wal"), 0).unwrap(); + let mut manifest = Manifest { bucket_count: 4, ..Manifest::default() }; + // Persist an initial generation so the manifest_dir exists on disk + // and the in-memory generation counter is non-zero. + manifest.save(&db_root).expect("initial manifest save"); + let (flush_sender, _flush_receiver) = 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, + }) +} + +/// Sabotage manifest writes by replacing the `manifest/` directory +/// with a regular file. `create_dir_all` then fails on the next save +/// because the path exists but is not a directory — without touching +/// permissions, which keeps the test cross-platform. +fn break_manifest_writes(db_root: &PathBuf) { + let manifest_dir = db_root.join("manifest"); + std::fs::remove_dir_all(&manifest_dir).expect("remove manifest dir"); + std::fs::write(&manifest_dir, b"not a directory").expect("write sentinel file"); +} + +#[tokio::test(flavor = "current_thread")] +async fn commit_manifest_leaves_state_unchanged_on_save_failure() { + let root = tmp_root(); + let state = build_state(root.clone()); + + let (gen_before, bucket_before) = { + let m = state.manifest.read().await; + (m.current_generation, m.bucket_count) + }; + assert!(gen_before >= 1); + + break_manifest_writes(&root); + + let result = state + .commit_manifest(|m| { + // Mutate something visible — bucket_count is a single field + // and easy to inspect. + m.bucket_count = 999; + }) + .await; + assert!(result.is_err(), "save should have failed"); + + let m = state.manifest.read().await; + assert_eq!( + m.current_generation, gen_before, + "generation must not advance when save fails" + ); + assert_eq!( + m.bucket_count, bucket_before, + "in-memory mutation must be discarded when save fails" + ); +} + +#[tokio::test(flavor = "current_thread")] +async fn commit_manifest_if_skips_save_when_closure_returns_none() { + let root = tmp_root(); + let state = build_state(root.clone()); + + let gen_before = state.manifest.read().await.current_generation; + + let returned: Option<()> = state + .commit_manifest_if(|m| { + // Mutate, then decide not to commit. The mutation lives only + // in the closure's local clone — it must not be published. + m.bucket_count = 999; + None + }) + .await + .expect("no save attempted, no error"); + assert!(returned.is_none()); + + let m = state.manifest.read().await; + assert_eq!( + m.current_generation, gen_before, + "generation must not advance when closure returns None" + ); + assert_eq!(m.bucket_count, 4, "mutation in skipped closure must not stick"); +} + +#[tokio::test(flavor = "current_thread")] +async fn commit_manifest_publishes_only_after_save_succeeds() { + let root = tmp_root(); + let state = build_state(root.clone()); + + let gen_before = state.manifest.read().await.current_generation; + + state + .commit_manifest(|m| { + m.bucket_count = 8; + }) + .await + .expect("save should succeed"); + + let m = state.manifest.read().await; + assert_eq!( + m.current_generation, + gen_before + 1, + "generation must advance on successful save" + ); + assert_eq!(m.bucket_count, 8, "mutation must be visible after commit"); +}