use axum::{
body::Body,
extract::{Query, State},
http::{header, StatusCode},
response::{IntoResponse, Response},
Json,
};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use std::collections::HashMap;
use tokio_stream::wrappers::ReceiverStream;
use crate::client::upgrade::Tier;
use crate::server::db::{CloudDb, UserRow};
use crate::server::{AppState, AuthUser};
#[derive(Serialize)]
pub struct AdminStatsResp {
pub total_users: i64,
pub users_by_tier: HashMap<String, i64>,
pub mrr_usd: f64,
pub total_tokens_saved_all_time: i64,
pub tokens_saved_last_30d: i64,
pub top_commands: Vec<TopCommand>,
pub new_users_last_7d: i64,
pub new_users_last_30d: i64,
}
#[derive(Serialize)]
pub struct TopCommand {
pub program: String,
pub tokens_saved: i64,
}
pub async fn stats(
State(state): State<AppState>,
AuthUser(user_id): AuthUser,
) -> impl IntoResponse {
let user = match state.db.get_user(&user_id) {
Some(u) => u,
None => {
return (
StatusCode::UNAUTHORIZED,
Json(serde_json::json!({"error": "not found"})),
)
.into_response();
}
};
if !user.is_admin {
return (
StatusCode::FORBIDDEN,
Json(serde_json::json!({"error": "forbidden"})),
)
.into_response();
}
let s = state.db.admin_stats();
let resp = AdminStatsResp {
total_users: s.total_users,
users_by_tier: s.users_by_tier.into_iter().collect(),
mrr_usd: s.mrr_usd,
total_tokens_saved_all_time: s.total_tokens_saved_all_time,
tokens_saved_last_30d: s.tokens_saved_last_30d,
top_commands: s
.top_commands
.into_iter()
.map(|(program, tokens_saved)| TopCommand {
program,
tokens_saved,
})
.collect(),
new_users_last_7d: s.new_users_last_7d,
new_users_last_30d: s.new_users_last_30d,
};
Json(resp).into_response()
}
#[derive(Serialize)]
pub struct TierChange {
pub old_tier: Option<String>,
pub new_tier: String,
pub reason: String,
pub at: String,
}
#[derive(Serialize)]
pub struct MonthUsage {
pub month: String,
pub tokens_used: i64,
pub tokens_saved: i64,
pub limit_reached_at: Option<String>,
}
fn active_tier_at<'a>(timeline: &'a [TierChange], ts: &str) -> &'a str {
let mut tier = "free";
for e in timeline {
if e.at.as_str() <= ts {
tier = &e.new_tier;
} else {
break;
}
}
tier
}
fn months_from_events(events: &[(String, i64, i64)], timeline: &[TierChange]) -> Vec<MonthUsage> {
use std::collections::BTreeMap;
let mut acc: BTreeMap<String, (i64, i64, Option<String>)> = BTreeMap::new();
for (ts, sent, saved) in events {
let month = ts.get(0..7).unwrap_or_default().to_string();
let entry = acc.entry(month).or_insert((0, 0, None));
entry.0 += sent;
entry.1 += saved;
if entry.2.is_none() {
if let Some(quota) = Tier::parse_tier(active_tier_at(timeline, ts)).cloud_token_quota()
{
if entry.0 as u64 >= quota {
entry.2 = Some(ts.get(0..10).unwrap_or(ts).to_string());
}
}
}
}
acc.into_iter()
.map(|(month, (used, saved, limit_reached_at))| MonthUsage {
month,
tokens_used: used,
tokens_saved: saved,
limit_reached_at,
})
.collect()
}
#[derive(Serialize)]
pub struct SubscriberExport {
pub subscriber: String,
pub joined_at: String,
pub current_tier: String,
pub current_token_quota: Option<u64>,
pub tier_timeline: Vec<TierChange>,
pub monthly_usage: Vec<MonthUsage>,
}
fn mask_id(user_id: &str) -> String {
let digest = Sha256::digest(user_id.as_bytes());
let hex: String = digest.iter().take(6).map(|b| format!("{b:02x}")).collect();
format!("sub_{hex}")
}
fn deny_if_not_admin(state: &AppState, user_id: &str) -> Option<Response> {
match state.db.get_user(user_id) {
Some(u) if u.is_admin => None,
Some(_) => Some(
(
StatusCode::FORBIDDEN,
Json(serde_json::json!({"error": "forbidden"})),
)
.into_response(),
),
None => Some(
(
StatusCode::UNAUTHORIZED,
Json(serde_json::json!({"error": "not found"})),
)
.into_response(),
),
}
}
fn build_subscriber(db: &CloudDb, u: UserRow) -> SubscriberExport {
let quota = Tier::parse_tier(&u.tier).cloud_token_quota();
let tier_timeline: Vec<TierChange> = db
.tier_events_for(&u.id)
.into_iter()
.map(|e| TierChange {
old_tier: e.old_tier,
new_tier: e.new_tier,
reason: e.reason,
at: e.occurred_at,
})
.collect();
let monthly_usage = months_from_events(&db.usage_events(&u.id), &tier_timeline);
SubscriberExport {
subscriber: mask_id(&u.id),
joined_at: u.created_at,
current_tier: u.tier,
current_token_quota: quota,
tier_timeline,
monthly_usage,
}
}
#[derive(Deserialize)]
pub struct ExportQuery {
#[serde(default)]
pub offset: i64,
pub limit: Option<i64>,
}
#[derive(Serialize)]
pub struct PagedExport {
pub total: i64,
pub limit: i64,
pub offset: i64,
pub subscribers: Vec<SubscriberExport>,
}
pub async fn export(
State(state): State<AppState>,
AuthUser(user_id): AuthUser,
Query(q): Query<ExportQuery>,
) -> Response {
if let Some(resp) = deny_if_not_admin(&state, &user_id) {
return resp;
}
let limit = q.limit.unwrap_or(50).clamp(1, 500);
let offset = q.offset.max(0);
let total = state.db.count_users();
let subscribers = state
.db
.users_page(limit, offset)
.into_iter()
.map(|u| build_subscriber(&state.db, u))
.collect();
Json(PagedExport {
total,
limit,
offset,
subscribers,
})
.into_response()
}
pub async fn export_csv(State(state): State<AppState>, AuthUser(user_id): AuthUser) -> Response {
if let Some(resp) = deny_if_not_admin(&state, &user_id) {
return resp;
}
let db = state.db.clone();
let (tx, rx) = tokio::sync::mpsc::channel::<Result<String, std::io::Error>>(8);
tokio::task::spawn_blocking(move || {
let header = "subscriber,joined_at,current_tier,token_quota,cycle_month,tokens_used,tokens_saved,pct_of_quota,maxed_on\n";
if tx.blocking_send(Ok(header.to_string())).is_err() {
return;
}
let total = db.count_users();
let page = 500;
let mut offset = 0;
while offset < total {
let users = db.users_page(page, offset);
if users.is_empty() {
break;
}
let mut buf = String::new();
for u in users {
append_csv_rows(&mut buf, &build_subscriber(&db, u));
}
if tx.blocking_send(Ok(buf)).is_err() {
return;
}
offset += page;
}
});
let body = Body::from_stream(ReceiverStream::new(rx));
(
StatusCode::OK,
[
(header::CONTENT_TYPE, "text/csv; charset=utf-8"),
(
header::CONTENT_DISPOSITION,
"attachment; filename=\"bctx-subscribers.csv\"",
),
],
body,
)
.into_response()
}
fn append_csv_rows(buf: &mut String, s: &SubscriberExport) {
use std::fmt::Write;
let quota = s
.current_token_quota
.map(|q| q.to_string())
.unwrap_or_default();
if s.monthly_usage.is_empty() {
let _ = writeln!(
buf,
"{},{},{},{},,0,0,,",
s.subscriber, s.joined_at, s.current_tier, quota
);
return;
}
for m in &s.monthly_usage {
let pct = s
.current_token_quota
.map(|q| format!("{:.1}", m.tokens_used as f64 / q as f64 * 100.0))
.unwrap_or_default();
let _ = writeln!(
buf,
"{},{},{},{},{},{},{},{},{}",
s.subscriber,
s.joined_at,
s.current_tier,
quota,
m.month,
m.tokens_used,
m.tokens_saved,
pct,
m.limit_reached_at.as_deref().unwrap_or("")
);
}
}
#[derive(Serialize)]
pub struct CohortRow {
pub month: String,
pub signups: i64,
pub converted: i64,
}
#[derive(Serialize)]
pub struct QuotaRow {
pub subscriber: String,
pub tier: String,
pub tokens_used: i64,
pub quota: Option<u64>,
pub pct: Option<f64>,
}
#[derive(Serialize)]
pub struct ConversionResp {
pub signups: i64,
pub ever_paid: i64,
pub paid_now: i64,
pub churned: i64,
pub conversion_rate: f64,
pub median_days_to_convert: Option<i64>,
pub cohorts: Vec<CohortRow>,
pub quota_usage: Vec<QuotaRow>,
}
pub async fn conversion(State(state): State<AppState>, AuthUser(user_id): AuthUser) -> Response {
if let Some(resp) = deny_if_not_admin(&state, &user_id) {
return resp;
}
let (signups, ever_paid, paid_now) = state.db.conversion_counts();
let churned = (ever_paid - paid_now).max(0);
let conversion_rate = if signups > 0 {
ever_paid as f64 / signups as f64
} else {
0.0
};
let mut days = state.db.days_to_convert();
days.sort_unstable();
let median_days_to_convert = (!days.is_empty()).then(|| days[days.len() / 2]);
let cohorts = state
.db
.cohorts()
.into_iter()
.map(|(month, signups, converted)| CohortRow {
month,
signups,
converted,
})
.collect();
let month = chrono::Utc::now().format("%Y-%m").to_string();
let quota_usage = state
.db
.quota_usage_top(&month, 50)
.into_iter()
.map(|(uid, tier, tokens_used)| {
let quota = Tier::parse_tier(&tier).cloud_token_quota();
let pct = quota.map(|q| tokens_used as f64 / q as f64 * 100.0);
QuotaRow {
subscriber: mask_id(&uid),
tier,
tokens_used,
quota,
pct,
}
})
.collect();
Json(ConversionResp {
signups,
ever_paid,
paid_now,
churned,
conversion_rate,
median_days_to_convert,
cohorts,
quota_usage,
})
.into_response()
}
#[derive(Deserialize)]
pub struct PromoteReq {
pub email: String,
}
async fn set_admin(
state: AppState,
user_id: String,
email: String,
is_admin: bool,
) -> impl IntoResponse {
let caller = match state.db.get_user(&user_id) {
Some(u) => u,
None => {
return (
StatusCode::UNAUTHORIZED,
Json(serde_json::json!({"error": "not found"})),
)
.into_response();
}
};
if !caller.is_admin {
return (
StatusCode::FORBIDDEN,
Json(serde_json::json!({"error": "forbidden"})),
)
.into_response();
}
match state.db.set_admin_by_email(&email, is_admin) {
Ok(true) => Json(serde_json::json!({"ok": true, "email": email})).into_response(),
Ok(false) => (
StatusCode::NOT_FOUND,
Json(serde_json::json!({"error": "user not found"})),
)
.into_response(),
Err(e) => (
StatusCode::INTERNAL_SERVER_ERROR,
Json(serde_json::json!({"error": e.to_string()})),
)
.into_response(),
}
}
pub async fn promote(
State(state): State<AppState>,
AuthUser(user_id): AuthUser,
Json(req): Json<PromoteReq>,
) -> impl IntoResponse {
set_admin(state, user_id, req.email, true).await
}
pub async fn demote(
State(state): State<AppState>,
AuthUser(user_id): AuthUser,
Json(req): Json<PromoteReq>,
) -> impl IntoResponse {
set_admin(state, user_id, req.email, false).await
}
#[cfg(test)]
mod tests {
use super::*;
fn change(new_tier: &str, at: &str) -> TierChange {
TierChange {
old_tier: None,
new_tier: new_tier.into(),
reason: "test".into(),
at: at.into(),
}
}
#[test]
fn limit_reached_accounts_for_midmonth_upgrade() {
let timeline = vec![
change("free", "2026-05-01 00:00:00"),
change("beacon", "2026-05-05 00:00:00"),
];
let events = vec![
("2026-05-02 09:00:00".to_string(), 10_000, 4_000),
("2026-05-06 09:00:00".to_string(), 60_000, 20_000),
("2026-05-08 09:00:00".to_string(), 50_000, 18_000),
("2026-06-03 09:00:00".to_string(), 5_000, 2_000),
];
let months = months_from_events(&events, &timeline);
assert_eq!(months.len(), 2);
assert_eq!(months[0].month, "2026-05");
assert_eq!(months[0].tokens_used, 120_000);
assert_eq!(months[0].limit_reached_at.as_deref(), Some("2026-05-08"));
assert_eq!(months[1].month, "2026-06");
assert_eq!(months[1].limit_reached_at, None);
}
#[test]
fn free_tier_never_reaches_limit() {
let timeline = vec![change("free", "2026-05-01 00:00:00")];
let events = vec![("2026-05-02 09:00:00".to_string(), 9_999_999, 1)];
let months = months_from_events(&events, &timeline);
assert_eq!(months[0].limit_reached_at, None);
}
#[test]
fn mask_id_is_stable_and_carries_no_email() {
let a = mask_id("user-uuid-123");
assert_eq!(a, mask_id("user-uuid-123"));
assert_ne!(a, mask_id("user-uuid-124"));
assert!(a.starts_with("sub_"));
}
#[test]
fn csv_rows_cover_quota_maxed_and_empty() {
let sub = SubscriberExport {
subscriber: "sub_abc".into(),
joined_at: "2026-05-03 10:00:00".into(),
current_tier: "beacon".into(),
current_token_quota: Some(100_000),
tier_timeline: vec![],
monthly_usage: vec![
MonthUsage {
month: "2026-05".into(),
tokens_used: 110_000,
tokens_saved: 38_000,
limit_reached_at: Some("2026-05-12".into()),
},
MonthUsage {
month: "2026-06".into(),
tokens_used: 40_000,
tokens_saved: 9_000,
limit_reached_at: None,
},
],
};
let mut buf = String::new();
append_csv_rows(&mut buf, &sub);
let lines: Vec<&str> = buf.lines().collect();
assert_eq!(lines.len(), 2);
assert_eq!(
lines[0],
"sub_abc,2026-05-03 10:00:00,beacon,100000,2026-05,110000,38000,110.0,2026-05-12"
);
assert_eq!(
lines[1],
"sub_abc,2026-05-03 10:00:00,beacon,100000,2026-06,40000,9000,40.0,"
);
let mut empty = String::new();
append_csv_rows(
&mut empty,
&SubscriberExport {
subscriber: "sub_def".into(),
joined_at: "2026-06-01 00:00:00".into(),
current_tier: "free".into(),
current_token_quota: None,
tier_timeline: vec![],
monthly_usage: vec![],
},
);
assert_eq!(empty.trim_end(), "sub_def,2026-06-01 00:00:00,free,,,0,0,,");
}
}