use std::collections::BTreeSet;
use std::path::{Path, PathBuf};
use std::time::{Duration, Instant};
use axum::{Json, extract::State, http::StatusCode, response::IntoResponse};
use serde::Serialize;
use serde_json::json;
use super::team::TeamAppState;
pub(super) const STORAGE_CACHE_TTL: Duration = Duration::from_mins(1);
pub(super) const QUOTA_ENV: &str = "LEANCTX_TEAM_STORAGE_QUOTA_BYTES";
pub(super) const DEFAULT_TEAM_STORAGE_QUOTA_BYTES: u64 = 5 * 1024 * 1024 * 1024;
#[derive(Debug, Clone, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct StorageComponent {
pub id: String,
pub bytes: u64,
}
#[derive(Debug, Clone)]
pub struct StorageReport {
pub used_bytes: u64,
pub components: Vec<StorageComponent>,
}
#[derive(Default)]
pub struct StorageCache(Option<(Instant, StorageReport)>);
#[derive(Clone)]
pub struct StorageRoots {
pub data_root: PathBuf,
pub workspaces: Vec<(String, PathBuf)>,
pub quota_bytes: u64,
}
impl StorageRoots {
fn measure(&self) -> StorageReport {
let mut components = Vec::new();
let data_bytes = dir_allocated_bytes(&self.data_root);
components.push(StorageComponent {
id: "server-data".into(),
bytes: data_bytes,
});
for (id, state_dir) in &self.workspaces {
if state_dir.starts_with(&self.data_root) {
continue; }
components.push(StorageComponent {
id: format!("workspace:{id}"),
bytes: dir_allocated_bytes(state_dir),
});
}
let used_bytes = components
.iter()
.map(|c| c.bytes)
.fold(0u64, u64::saturating_add);
StorageReport {
used_bytes,
components,
}
}
}
pub async fn v1_storage(State(state): State<TeamAppState>) -> impl IntoResponse {
let (report, age) = cached_report(&state).await;
let body = json!({
"schemaVersion": 1,
"measuredAt": chrono::Utc::now().to_rfc3339(),
"usedBytes": report.used_bytes,
"quotaBytes": state.team.storage_roots.quota_bytes,
"components": report.components,
"cacheAgeSeconds": age.as_secs(),
});
(StatusCode::OK, Json(body))
}
pub async fn v1_usage(State(state): State<TeamAppState>) -> impl IntoResponse {
let dir = state.team.savings_store_dir.lock().await.clone();
let summary = tokio::task::spawn_blocking(move || super::savings_summary::aggregate(&dir))
.await
.unwrap_or_default();
let (report, _) = cached_report(&state).await;
let storage = json!({
"used_bytes": report.used_bytes,
"quota_bytes": state.team.storage_roots.quota_bytes,
});
let body = json!({
"schemaVersion": 1,
"generatedAt": chrono::Utc::now().to_rfc3339(),
"savings": {
"memberCount": summary.member_count,
"savedTokens": summary.totals.saved_tokens,
"netSavedTokens": summary.totals.net_saved_tokens,
"savedUsd": summary.totals.saved_usd,
},
"toolCalls": summary.totals.total_events,
"storage": storage,
});
(StatusCode::OK, Json(body))
}
async fn cached_report(state: &TeamAppState) -> (StorageReport, Duration) {
{
let cache = state.team.storage_cache.lock().await;
if let Some((at, report)) = cache.0.as_ref() {
let age = at.elapsed();
if age < STORAGE_CACHE_TTL {
return (report.clone(), age);
}
}
}
let roots = state.team.storage_roots.clone();
let report = tokio::task::spawn_blocking(move || roots.measure())
.await
.unwrap_or(StorageReport {
used_bytes: 0,
components: Vec::new(),
});
let mut cache = state.team.storage_cache.lock().await;
cache.0 = Some((Instant::now(), report.clone()));
(report, Duration::ZERO)
}
fn quota_bytes_from_env() -> Option<u64> {
std::env::var(QUOTA_ENV).ok()?.trim().parse::<u64>().ok()
}
pub(super) fn resolve_quota_bytes(env_override: Option<u64>, config_quota: Option<u64>) -> u64 {
env_override
.or(config_quota)
.unwrap_or(DEFAULT_TEAM_STORAGE_QUOTA_BYTES)
}
pub(super) fn storage_roots_from_config(
audit_log_path: &Path,
workspaces: &[(String, PathBuf)],
config_quota: Option<u64>,
) -> StorageRoots {
let data_root = audit_log_path
.parent()
.unwrap_or_else(|| Path::new("."))
.to_path_buf();
let workspaces = workspaces
.iter()
.map(|(id, root)| (id.clone(), root.join(".lean-ctx")))
.collect();
StorageRoots {
data_root,
workspaces,
quota_bytes: resolve_quota_bytes(quota_bytes_from_env(), config_quota),
}
}
fn dir_allocated_bytes(path: &Path) -> u64 {
let mut seen: BTreeSet<(u64, u64)> = BTreeSet::new();
let mut total = 0u64;
let mut stack = vec![path.to_path_buf()];
while let Some(p) = stack.pop() {
let Ok(meta) = std::fs::symlink_metadata(&p) else {
continue;
};
if meta.is_symlink() {
continue;
}
if meta.is_dir() {
if let Ok(entries) = std::fs::read_dir(&p) {
stack.extend(entries.flatten().map(|e| e.path()));
}
continue;
}
#[cfg(unix)]
{
use std::os::unix::fs::MetadataExt;
if !seen.insert((meta.dev(), meta.ino())) {
continue; }
total = total.saturating_add(meta.blocks().saturating_mul(512));
}
#[cfg(not(unix))]
{
let _ = &mut seen;
total = total.saturating_add(meta.len());
}
}
total
}
#[cfg(test)]
mod tests {
use super::*;
fn temp_dir(tag: &str) -> PathBuf {
let d =
std::env::temp_dir().join(format!("leanctx_team_billing_{tag}_{}", std::process::id()));
let _ = std::fs::remove_dir_all(&d);
std::fs::create_dir_all(&d).unwrap();
d
}
#[test]
fn missing_dir_measures_zero() {
let missing = std::env::temp_dir().join("leanctx_team_billing_does_not_exist_xyz");
let _ = std::fs::remove_dir_all(&missing);
assert_eq!(dir_allocated_bytes(&missing), 0);
}
#[test]
fn quota_resolution_precedence() {
assert_eq!(resolve_quota_bytes(Some(7), Some(9)), 7);
assert_eq!(resolve_quota_bytes(None, Some(9)), 9);
assert_eq!(
resolve_quota_bytes(None, None),
DEFAULT_TEAM_STORAGE_QUOTA_BYTES
);
assert_eq!(DEFAULT_TEAM_STORAGE_QUOTA_BYTES, 5_368_709_120);
}
#[test]
fn allocated_bytes_cover_written_content() {
let d = temp_dir("alloc");
std::fs::write(d.join("a.jsonl"), vec![b'x'; 10_000]).unwrap();
std::fs::create_dir_all(d.join("nested")).unwrap();
std::fs::write(d.join("nested/b.bin"), vec![b'y'; 5_000]).unwrap();
let measured = dir_allocated_bytes(&d);
assert!(measured >= 15_000, "measured {measured} < logical 15000");
assert!(
measured < 15_000 * 64,
"measured {measured} implausibly large"
);
let _ = std::fs::remove_dir_all(&d);
}
#[cfg(unix)]
#[test]
fn hard_links_count_once_and_symlinks_do_not_escape() {
let d = temp_dir("links");
std::fs::write(d.join("real.bin"), vec![b'z'; 8_192]).unwrap();
std::fs::hard_link(d.join("real.bin"), d.join("hard.bin")).unwrap();
let outside = temp_dir("links_outside");
std::fs::write(outside.join("big.bin"), vec![b'w'; 100_000]).unwrap();
std::os::unix::fs::symlink(outside.join("big.bin"), d.join("escape.bin")).unwrap();
let measured = dir_allocated_bytes(&d);
let single = dir_allocated_bytes(&outside); assert!(measured < single, "symlink target must not be billed");
assert!(
(8_192..16_384).contains(&measured),
"hard link double-counted: {measured}"
);
let _ = std::fs::remove_dir_all(&d);
let _ = std::fs::remove_dir_all(&outside);
}
#[test]
fn workspace_state_dirs_under_data_root_are_not_double_counted() {
let d = temp_dir("dedupe");
let audit = d.join("audit.jsonl");
std::fs::write(&audit, "x").unwrap();
let ws_root = d.join("ws1");
std::fs::create_dir_all(ws_root.join(".lean-ctx")).unwrap();
std::fs::write(ws_root.join(".lean-ctx/events.jsonl"), vec![b'e'; 4_096]).unwrap();
let ext = temp_dir("dedupe_ext");
std::fs::create_dir_all(ext.join(".lean-ctx")).unwrap();
std::fs::write(ext.join(".lean-ctx/k.jsonl"), vec![b'k'; 4_096]).unwrap();
let roots = storage_roots_from_config(
&audit,
&[
("inside".into(), ws_root.clone()),
("outside".into(), ext.clone()),
],
None,
);
let report = roots.measure();
let ids: Vec<&str> = report.components.iter().map(|c| c.id.as_str()).collect();
assert!(ids.contains(&"server-data"));
assert!(
!ids.contains(&"workspace:inside"),
"nested workspace double-counted"
);
assert!(ids.contains(&"workspace:outside"));
let sum: u64 = report.components.iter().map(|c| c.bytes).sum();
assert_eq!(report.used_bytes, sum);
let _ = std::fs::remove_dir_all(&d);
let _ = std::fs::remove_dir_all(&ext);
}
}