#![forbid(unsafe_code)]
use chrono::{DateTime, Duration, NaiveDate, Utc};
use serde::{Deserialize, Serialize};
use serde_json::json;
use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
const RETENTION_DAYS: i64 = 90;
const TANTIVY_WALK_INTERVAL_SECS: i64 = 300;
const AVERAGE_WINDOW_DAYS: i64 = 30;
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
pub struct DayWrites {
pub bytes: u64,
pub samples: u32,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct LastSample {
ts: String,
lmdb_bytes: u64,
tantivy_bytes: u64,
}
#[derive(Debug, Default, Serialize, Deserialize)]
struct WriteBudgetState {
version: u8,
days: BTreeMap<String, DayWrites>,
last_sample: Option<LastSample>,
tracking_since: Option<String>,
}
impl WriteBudgetState {
const fn new() -> Self {
Self {
version: 1,
days: BTreeMap::new(),
last_sample: None,
tracking_since: None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct WriteBudgetDelta {
pub attributed_bytes: u64,
pub date: NaiveDate,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WriteBudgetReport {
pub today_bytes: u64,
pub yesterday_bytes: Option<u64>,
pub avg_30d_bytes: u64,
pub days_tracked: u32,
pub busiest_day: Option<(String, u64)>,
pub lmdb_bytes: u64,
pub tantivy_bytes: u64,
}
#[derive(Debug)]
pub struct WriteBudgetLedger {
path: PathBuf,
lmdb_data_mdb: PathBuf,
tantivy_dir: PathBuf,
state: WriteBudgetState,
last_tantivy_walk: Option<DateTime<Utc>>,
cached_tantivy_bytes: u64,
}
impl WriteBudgetLedger {
#[must_use]
pub fn paths(store_root: &Path) -> (PathBuf, PathBuf, PathBuf) {
(
store_root.join("write_budget.json"),
store_root.join("lmdb").join("data.mdb"),
store_root.join("lmdb").join("tantivy"),
)
}
#[must_use]
pub fn load(store_root: &Path) -> Self {
let (path, lmdb_data_mdb, tantivy_dir) = Self::paths(store_root);
let state = std::fs::read_to_string(&path)
.ok()
.and_then(|raw| {
let parsed = serde_json::from_str::<WriteBudgetState>(&raw).ok()?;
Some(parsed)
})
.unwrap_or_else(|| {
if path.exists() {
tracing::warn!(
path = %path.display(),
"write_budget.json unreadable — starting a fresh ledger"
);
}
WriteBudgetState::new()
});
Self {
path,
lmdb_data_mdb,
tantivy_dir,
state,
last_tantivy_walk: None,
cached_tantivy_bytes: 0,
}
}
pub fn observe(&mut self, now: DateTime<Utc>) -> WriteBudgetDelta {
let lmdb_bytes = std::fs::metadata(&self.lmdb_data_mdb).map_or(0, |m| m.len());
let walk_due = self
.last_tantivy_walk
.is_none_or(|last| (now - last).num_seconds() >= TANTIVY_WALK_INTERVAL_SECS);
if walk_due {
self.cached_tantivy_bytes = dir_size(&self.tantivy_dir);
self.last_tantivy_walk = Some(now);
}
let total_bytes = lmdb_bytes.saturating_add(self.cached_tantivy_bytes);
let (attributed, date) = match &self.state.last_sample {
Some(last) => {
let growth =
total_bytes.saturating_sub(last.lmdb_bytes.saturating_add(last.tantivy_bytes));
(growth, now.date_naive())
}
None => (0, now.date_naive()),
};
self.state.last_sample = Some(LastSample {
ts: now.to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
lmdb_bytes,
tantivy_bytes: self.cached_tantivy_bytes,
});
let day_key = date.to_string();
if attributed > 0 {
let entry = self.state.days.entry(day_key.clone()).or_default();
entry.bytes = entry.bytes.saturating_add(attributed);
entry.samples = entry.samples.saturating_add(1);
}
if self.state.tracking_since.is_none() {
self.state.tracking_since = Some(day_key);
}
self.prune(now.date_naive());
self.persist();
WriteBudgetDelta {
attributed_bytes: attributed,
date,
}
}
fn prune(&mut self, today: NaiveDate) {
let Some(cutoff) = today.checked_sub_signed(Duration::days(RETENTION_DAYS)) else {
return;
};
self.state.days.retain(|day, _| {
NaiveDate::parse_from_str(day, "%Y-%m-%d").is_ok_and(|d| d >= cutoff) });
}
fn persist(&self) {
match serde_json::to_string_pretty(&self.state) {
Ok(body) => {
let tmp = self
.path
.with_file_name(format!("write_budget.json.tmp.{}", std::process::id()));
if let Err(e) =
std::fs::write(&tmp, body).and_then(|()| std::fs::rename(&tmp, &self.path))
{
tracing::warn!(
path = %self.path.display(),
error = %e,
"failed to persist write budget ledger"
);
}
}
Err(e) => tracing::warn!(error = %e, "write budget serialization failed"),
}
}
pub fn fresh_report(&mut self) -> WriteBudgetReport {
let lmdb_bytes = std::fs::metadata(&self.lmdb_data_mdb).map_or(0, |m| m.len());
let walk_due = self
.last_tantivy_walk
.is_none_or(|last| (Utc::now() - last).num_seconds() >= TANTIVY_WALK_INTERVAL_SECS);
if walk_due {
self.cached_tantivy_bytes = dir_size(&self.tantivy_dir);
self.last_tantivy_walk = Some(Utc::now());
}
self.state.last_sample = Some(LastSample {
ts: Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
lmdb_bytes,
tantivy_bytes: self.cached_tantivy_bytes,
});
self.report()
}
#[must_use]
pub fn report(&self) -> WriteBudgetReport {
self.report_as_of(Utc::now())
}
#[must_use]
pub fn report_as_of(&self, now: DateTime<Utc>) -> WriteBudgetReport {
let today = now.date_naive();
let today_key = today.to_string();
let yesterday_key = (today - Duration::days(1)).to_string();
let today_bytes = self.state.days.get(&today_key).map_or(0, |d| d.bytes);
let yesterday_bytes = self.state.days.get(&yesterday_key).map(|d| d.bytes);
let mut sum = 0u64;
let mut counted = 0u32;
for offset in 0..AVERAGE_WINDOW_DAYS {
let key = (today - Duration::days(offset)).to_string();
if let Some(day) = self.state.days.get(&key) {
sum = sum.saturating_add(day.bytes);
counted += 1;
}
}
let avg = if counted > 0 {
sum / u64::from(counted)
} else {
0
};
let busiest_day = self
.state
.days
.iter()
.max_by_key(|(_, d)| d.bytes)
.map(|(day, d)| (day.clone(), d.bytes))
.filter(|(_, bytes)| *bytes > 0);
let (lmdb_bytes, tantivy_bytes) = match &self.state.last_sample {
Some(last) => (last.lmdb_bytes, last.tantivy_bytes),
None => (0, 0),
};
WriteBudgetReport {
today_bytes,
yesterday_bytes,
avg_30d_bytes: avg,
days_tracked: u32::try_from(self.state.days.len()).unwrap_or(u32::MAX),
busiest_day,
lmdb_bytes,
tantivy_bytes,
}
}
#[must_use]
pub fn report_json(&self) -> serde_json::Value {
let r = self.report();
json!({
"today_bytes": r.today_bytes,
"yesterday_bytes": r.yesterday_bytes,
"avg_30d_bytes": r.avg_30d_bytes,
"days_tracked": r.days_tracked,
"busiest_day": r.busiest_day,
"lmdb_bytes": r.lmdb_bytes,
"tantivy_bytes": r.tantivy_bytes,
})
}
#[must_use]
pub fn ledger_path(&self) -> &Path {
&self.path
}
#[cfg(test)]
const fn base_date() -> NaiveDate {
NaiveDate::from_ymd_opt(2026, 8, 28).expect("valid date")
}
}
fn dir_size(path: &Path) -> u64 {
let mut total = 0u64;
let Ok(entries) = std::fs::read_dir(path) else {
return 0;
};
for entry in entries.flatten() {
let Ok(meta) = entry.metadata() else {
continue;
};
if meta.is_dir() {
total = total.saturating_add(dir_size(&entry.path()));
} else {
total = total.saturating_add(meta.len());
}
}
total
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::{Datelike, TimeZone};
fn at(day: NaiveDate, hour: u32) -> DateTime<Utc> {
Utc.with_ymd_and_hms(day.year(), day.month(), day.day(), hour, 0, 0)
.single()
.expect("valid test time")
}
fn fake_store() -> (tempfile::TempDir, PathBuf) {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().to_path_buf();
std::fs::create_dir_all(root.join("lmdb/tantivy")).unwrap();
(dir, root)
}
fn write_sizes(root: &Path, lmdb: u64, tantivy: u64) {
let mdb = root.join("lmdb/data.mdb");
if mdb.exists() {
std::fs::remove_file(&mdb).unwrap();
}
std::fs::write(&mdb, vec![0u8; lmdb as usize]).unwrap();
let seg = root.join("lmdb/tantivy/seg.0");
if seg.exists() {
std::fs::remove_file(&seg).unwrap();
}
std::fs::write(&seg, vec![0u8; tantivy as usize]).unwrap();
}
#[test]
fn first_observation_baselines_without_attributing() {
let (_dir, root) = fake_store();
write_sizes(&root, 1_000, 500);
let mut ledger = WriteBudgetLedger::load(&root);
let delta = ledger.observe(at(WriteBudgetLedger::base_date(), 10));
assert_eq!(delta.attributed_bytes, 0, "first sample is a baseline");
assert_eq!(ledger.report().lmdb_bytes, 1_000);
assert_eq!(ledger.report().tantivy_bytes, 500);
}
#[test]
fn growth_attributed_to_current_day() {
let (_dir, root) = fake_store();
write_sizes(&root, 1_000, 500);
let mut ledger = WriteBudgetLedger::load(&root);
ledger.observe(at(WriteBudgetLedger::base_date(), 10));
write_sizes(&root, 1_000 + 4_200, 500 + 100);
let delta = ledger.observe(at(WriteBudgetLedger::base_date(), 11));
assert_eq!(delta.attributed_bytes, 4_300);
let report = ledger.report_as_of(at(WriteBudgetLedger::base_date(), 11));
assert_eq!(report.today_bytes, 4_300);
assert_eq!(report.yesterday_bytes, None);
assert_eq!(report.busiest_day, Some(("2026-08-28".into(), 4_300)));
}
#[test]
fn shrink_attributed_as_zero_not_negative() {
let (_dir, root) = fake_store();
write_sizes(&root, 10_000, 0);
let mut ledger = WriteBudgetLedger::load(&root);
ledger.observe(at(WriteBudgetLedger::base_date(), 10));
write_sizes(&root, 500, 0); let delta = ledger.observe(at(WriteBudgetLedger::base_date(), 11));
assert_eq!(delta.attributed_bytes, 0, "never credit deletes as writes");
assert_eq!(
ledger
.report_as_of(at(WriteBudgetLedger::base_date(), 11))
.today_bytes,
0
);
}
#[test]
fn utc_midnight_rollover_starts_new_day() {
let (_dir, root) = fake_store();
write_sizes(&root, 1_000, 0);
let mut ledger = WriteBudgetLedger::load(&root);
ledger.observe(at(WriteBudgetLedger::base_date(), 10));
write_sizes(&root, 3_000, 0);
ledger.observe(at(WriteBudgetLedger::base_date(), 22));
write_sizes(&root, 4_000, 0);
let next = WriteBudgetLedger::base_date().succ_opt().unwrap();
ledger.observe(at(next, 1));
let report = ledger.report_as_of(at(next, 1));
assert_eq!(report.today_bytes, 1_000);
assert_eq!(report.yesterday_bytes, Some(2_000));
}
#[test]
fn tantivy_walk_throttled_within_interval() {
let (_dir, root) = fake_store();
write_sizes(&root, 1_000, 100);
let mut ledger = WriteBudgetLedger::load(&root);
ledger.observe(at(WriteBudgetLedger::base_date(), 10));
write_sizes(&root, 1_000, 9_999);
let delta = ledger.observe(at(WriteBudgetLedger::base_date(), 10) + Duration::seconds(1));
assert_eq!(
delta.attributed_bytes, 0,
"stale tantivy cache must not over-attribute"
);
let delta = ledger.observe(at(WriteBudgetLedger::base_date(), 10) + Duration::seconds(301));
assert_eq!(delta.attributed_bytes, 9_899);
}
#[test]
fn ledger_roundtrips_and_recovers_from_corruption() {
let (_dir, root) = fake_store();
write_sizes(&root, 1_000, 0);
let mut ledger = WriteBudgetLedger::load(&root);
ledger.observe(at(WriteBudgetLedger::base_date(), 10));
write_sizes(&root, 5_000, 0);
ledger.observe(at(WriteBudgetLedger::base_date(), 12));
let reloaded = WriteBudgetLedger::load(&root);
assert_eq!(
reloaded
.report_as_of(at(WriteBudgetLedger::base_date(), 12))
.today_bytes,
4_000
);
std::fs::write(root.join("write_budget.json"), "{not json").unwrap();
let corrupted = WriteBudgetLedger::load(&root);
assert_eq!(corrupted.report().days_tracked, 0);
}
#[test]
fn retention_prunes_old_days() {
let (_dir, root) = fake_store();
write_sizes(&root, 1_000, 0);
let mut ledger = WriteBudgetLedger::load(&root);
ledger.observe(at(WriteBudgetLedger::base_date(), 10));
let ancient = (WriteBudgetLedger::base_date() - Duration::days(120)).to_string();
ledger.state.days.insert(
ancient.clone(),
DayWrites {
bytes: 7,
samples: 1,
},
);
ledger.state.days.insert(
"2026-08-28".into(),
DayWrites {
bytes: 5,
samples: 1,
},
);
ledger.prune(WriteBudgetLedger::base_date());
assert!(!ledger.state.days.contains_key(&ancient), "old days prune");
assert!(
ledger.state.days.contains_key("2026-08-28"),
"recent days survive"
);
}
}