Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
95 changes: 52 additions & 43 deletions src/api/http_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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(),
})))
}

Expand Down
64 changes: 29 additions & 35 deletions src/compact/worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<String> = 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));
}
};

Expand Down Expand Up @@ -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<String>)> = 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));
}
Expand Down
11 changes: 6 additions & 5 deletions src/ingest/flusher.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
43 changes: 24 additions & 19 deletions src/rollup/worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -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<usize> {
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)
}
Expand Down
6 changes: 5 additions & 1 deletion src/runtime/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
43 changes: 43 additions & 0 deletions src/runtime/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,4 +24,47 @@ pub struct AppStateInner {
pub flush_sender: tokio::sync::mpsc::Sender<FlushMessage>,
}

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<F, T>(&self, op: F) -> std::io::Result<T>
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<F, T>(&self, op: F) -> std::io::Result<Option<T>>
where
F: FnOnce(&mut Manifest) -> Option<T>,
{
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<AppStateInner>;
21 changes: 20 additions & 1 deletion src/storage/compression.rs
Original file line number Diff line number Diff line change
@@ -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};

Expand All @@ -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<i32> = OnceLock::new();
*LEVEL.get_or_init(|| {
std::env::var("USAGEDB_ZSTD_LEVEL")
.ok()
.and_then(|v| v.trim().parse::<i32>().ok())
.filter(|l| (1..=22).contains(l))
.unwrap_or(19)
})
}

pub fn compress(data: &[u8], codec: CompressionCodec) -> IoResult<Vec<u8>> {
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()
}
Expand Down
Loading
Loading