use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use anyhow::{Context, Result};
use serde::{Deserialize, Serialize};
use tracing::{info, warn};
use crate::limiter::{gcra_admit, Gcra};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Bucket {
Orders,
RegisteredDomain,
IdentifierSet,
}
impl Bucket {
pub fn label(&self) -> &'static str {
match self {
Bucket::Orders => "orders",
Bucket::RegisteredDomain => "registered_domain",
Bucket::IdentifierSet => "identifier_set",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Decision {
Allow,
Defer {
bucket: Bucket,
retry_at_unix: i64,
},
}
#[derive(Debug, Clone)]
pub struct CaProfile {
pub name: &'static str,
pub orders: Gcra,
pub registered_domain: Gcra,
pub identifier_set: Gcra,
}
impl CaProfile {
fn lets_encrypt() -> Result<CaProfile> {
Ok(CaProfile {
name: "letsencrypt",
orders: Gcra::from_parts(300, Duration::from_secs(3 * 3_600), 300)?,
registered_domain: Gcra::from_parts(50, Duration::from_secs(7 * 86_400), 50)?,
identifier_set: Gcra::from_parts(5, Duration::from_secs(7 * 86_400), 5)?,
})
}
pub fn for_directory(url: &str) -> Option<CaProfile> {
let host = url
.split_once("://")
.map(|(_, r)| r)
.unwrap_or(url)
.split('/')
.next()
.unwrap_or("");
if host.ends_with("letsencrypt.org") {
CaProfile::lets_encrypt().ok()
} else {
None
}
}
fn gcra(&self, b: Bucket) -> &Gcra {
match b {
Bucket::Orders => &self.orders,
Bucket::RegisteredDomain => &self.registered_domain,
Bucket::IdentifierSet => &self.identifier_set,
}
}
}
pub fn group_key(host: &str) -> String {
let h = host.trim().trim_end_matches('.').to_ascii_lowercase();
if h.is_empty() || h.parse::<std::net::IpAddr>().is_ok() {
return h;
}
let labels: Vec<&str> = h.split('.').collect();
if labels.len() <= 2 {
return h;
}
labels[labels.len() - 2..].join(".")
}
pub fn identifier_set_key(domains: &[String]) -> String {
let mut d: Vec<String> = domains
.iter()
.map(|s| s.trim().trim_end_matches('.').to_ascii_lowercase())
.filter(|s| !s.is_empty())
.collect();
d.sort();
d.dedup();
d.join(",")
}
#[derive(Debug, Default, Serialize, Deserialize)]
struct LedgerFile {
#[serde(default)]
tats: HashMap<String, u64>,
}
pub struct IssuanceBudget {
profile: CaProfile,
path: PathBuf,
file: LedgerFile,
}
impl IssuanceBudget {
pub fn load(directory_url: &str, cache_dir: &str) -> Option<IssuanceBudget> {
let profile = CaProfile::for_directory(directory_url)?;
let path = Path::new(cache_dir).join("issuance-budget.json");
let file = std::fs::read_to_string(&path)
.ok()
.and_then(|s| serde_json::from_str::<LedgerFile>(&s).ok())
.unwrap_or_default();
info!(
ca = profile.name,
ledger = %path.display(),
entries = file.tats.len(),
"ACME issuance budget active"
);
Some(IssuanceBudget {
profile,
path,
file,
})
}
fn key(bucket: Bucket, k: &str) -> String {
format!("{}:{}", bucket.label(), k)
}
pub fn check(&self, domains: &[String], now_unix: i64) -> Decision {
let now_us = (now_unix.max(0) as u64).saturating_mul(1_000_000);
for (bucket, k) in self.keys_for(domains) {
let g = self.profile.gcra(bucket);
let stored = self.file.tats.get(&Self::key(bucket, &k)).copied();
if gcra_admit(stored, now_us, g).is_none() {
return Decision::Defer {
bucket,
retry_at_unix: (g.next_admit_at(stored, now_us) / 1_000_000) as i64,
};
}
}
Decision::Allow
}
pub fn debit(&mut self, domains: &[String], now_unix: i64) -> Result<()> {
let now_us = (now_unix.max(0) as u64).saturating_mul(1_000_000);
for (bucket, k) in self.keys_for(domains) {
let g = self.profile.gcra(bucket);
let key = Self::key(bucket, &k);
let stored = self.file.tats.get(&key).copied();
if let Some(new_tat) = gcra_admit(stored, now_us, g) {
self.file.tats.insert(key, new_tat);
}
}
self.persist()
}
fn keys_for(&self, domains: &[String]) -> Vec<(Bucket, String)> {
let mut out = vec![
(Bucket::Orders, "account".to_string()),
(Bucket::IdentifierSet, identifier_set_key(domains)),
];
let mut groups: Vec<String> = domains.iter().map(|d| group_key(d)).collect();
groups.sort();
groups.dedup();
for g in groups {
out.push((Bucket::RegisteredDomain, g));
}
out
}
fn persist(&self) -> Result<()> {
if let Some(dir) = self.path.parent() {
std::fs::create_dir_all(dir).with_context(|| format!("creating {}", dir.display()))?;
}
let tmp = self.path.with_extension("json.tmp");
let body = serde_json::to_string(&self.file).context("serialising the issuance ledger")?;
std::fs::write(&tmp, body).with_context(|| format!("writing {}", tmp.display()))?;
std::fs::rename(&tmp, &self.path)
.with_context(|| format!("renaming into {}", self.path.display()))?;
Ok(())
}
pub fn remaining(&self, domains: &[String], now_unix: i64) -> Vec<(Bucket, String, u64)> {
let now_us = (now_unix.max(0) as u64).saturating_mul(1_000_000);
self.keys_for(domains)
.into_iter()
.map(|(bucket, k)| {
let stored = self.file.tats.get(&Self::key(bucket, &k)).copied();
let left = self.profile.gcra(bucket).remaining(stored, now_us);
(bucket, k, left)
})
.collect()
}
pub fn ca_name(&self) -> &'static str {
self.profile.name
}
}
pub fn now_unix() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs() as i64)
.unwrap_or(0)
}
pub fn render_metrics(budget: &IssuanceBudget, domains: &[String], now_unix: i64) -> String {
let mut out = String::new();
out.push_str(
"# HELP edgeguard_acme_budget_remaining Issuance THIS EDGE's own ledger still allows, per bucket. Not a fleet total: the registered_domain limit is shared across every edge under that domain, and the control plane's GET /v3/acme/budget is the fleet figure.\n",
);
out.push_str("# TYPE edgeguard_acme_budget_remaining gauge\n");
for (bucket, key, left) in budget.remaining(domains, now_unix) {
let label = match bucket {
Bucket::IdentifierSet => short_hash(&key),
_ => key.clone(),
};
out.push_str(&format!(
"edgeguard_acme_budget_remaining{{ca=\"{}\",bucket=\"{}\",key=\"{}\"}} {left}\n",
budget.ca_name(),
bucket.label(),
escape_label(&label),
));
}
out
}
fn short_hash(s: &str) -> String {
let mut h: u64 = 0xcbf2_9ce4_8422_2325;
for b in s.as_bytes() {
h ^= *b as u64;
h = h.wrapping_mul(0x1000_0000_01b3);
}
format!("{h:016x}")
}
fn escape_label(s: &str) -> String {
s.replace('\\', "\\\\").replace('"', "\\\"")
}
pub fn warn_deferred(bucket: Bucket, retry_at_unix: i64, domains: &[String]) {
warn!(
bucket = bucket.label(),
retry_at_unix,
domains = ?domains,
"ACME issuance deferred: the CA's rate limit for this bucket is exhausted. \
The existing certificate (if any) keeps serving; no self-signed certificate is \
substituted on a public name."
);
}
#[cfg(test)]
mod tests {
use super::*;
fn budget(dir: &Path) -> IssuanceBudget {
IssuanceBudget::load(
"https://acme-v02.api.letsencrypt.org/directory",
dir.to_str().unwrap(),
)
.expect("letsencrypt profile")
}
fn tmpdir(tag: &str) -> PathBuf {
let d = std::env::temp_dir().join(format!("eg-acme-budget-{tag}-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&d);
std::fs::create_dir_all(&d).unwrap();
d
}
#[test]
fn the_identifier_set_bucket_stops_a_restart_loop_at_five() {
let d = tmpdir("restart");
let mut b = budget(&d);
let domains = vec!["edge.example.com".to_string()];
let now = 1_800_000_000;
for i in 0..5 {
assert_eq!(
b.check(&domains, now),
Decision::Allow,
"order {i} must be allowed"
);
b.debit(&domains, now).unwrap();
}
match b.check(&domains, now) {
Decision::Defer {
bucket,
retry_at_unix,
} => {
assert_eq!(bucket, Bucket::IdentifierSet);
assert!(retry_at_unix > now, "a deferral must say when to retry");
}
Decision::Allow => panic!("the sixth order must be deferred"),
}
}
#[test]
fn the_budget_survives_a_restart() {
let d = tmpdir("persist");
let domains = vec!["edge.example.com".to_string()];
let now = 1_800_000_000;
{
let mut b = budget(&d);
for _ in 0..5 {
b.debit(&domains, now).unwrap();
}
}
let b2 = budget(&d);
assert!(
matches!(b2.check(&domains, now), Decision::Defer { .. }),
"the ledger must be read back from disk"
);
}
#[test]
fn the_bucket_refills_over_time_rather_than_resetting_at_a_boundary() {
let d = tmpdir("refill");
let mut b = budget(&d);
let domains = vec!["edge.example.com".to_string()];
let now = 1_800_000_000;
for _ in 0..5 {
b.debit(&domains, now).unwrap();
}
assert!(matches!(b.check(&domains, now), Decision::Defer { .. }));
const EI: i64 = 7 * 86_400 / 5;
assert!(matches!(
b.check(&domains, now + EI - 60),
Decision::Defer { .. }
));
assert_eq!(b.check(&domains, now + EI + 60), Decision::Allow);
b.debit(&domains, now + EI + 60).unwrap();
assert!(matches!(
b.check(&domains, now + EI + 60),
Decision::Defer { .. }
));
}
#[test]
fn a_different_identifier_set_has_its_own_bucket() {
let d = tmpdir("sets");
let mut b = budget(&d);
let now = 1_800_000_000;
let a = vec!["a.example.com".to_string()];
for _ in 0..5 {
b.debit(&a, now).unwrap();
}
assert!(matches!(b.check(&a, now), Decision::Defer { .. }));
let c = vec!["b.example.com".to_string()];
assert_eq!(b.check(&c, now), Decision::Allow);
}
#[test]
fn the_registered_domain_bucket_binds_across_different_hostnames() {
let d = tmpdir("domain");
let mut b = budget(&d);
let now = 1_800_000_000;
for i in 0..50 {
let h = vec![format!("h{i}.example.com")];
assert_eq!(b.check(&h, now), Decision::Allow, "host {i}");
b.debit(&h, now).unwrap();
}
let next = vec!["h50.example.com".to_string()];
match b.check(&next, now) {
Decision::Defer { bucket, .. } => assert_eq!(bucket, Bucket::RegisteredDomain),
Decision::Allow => panic!("the 51st registered-domain order must be deferred"),
}
assert_eq!(b.check(&["x.other.com".to_string()], now), Decision::Allow);
}
#[test]
fn identifier_set_key_is_order_and_case_insensitive() {
let a = identifier_set_key(&["b.example.com".into(), "A.example.com".into()]);
let b = identifier_set_key(&["a.example.com".into(), "B.EXAMPLE.COM.".into()]);
assert_eq!(a, b);
assert_ne!(a, identifier_set_key(&["a.example.com".into()]));
}
#[test]
fn group_key_errs_toward_merging_never_splitting() {
assert_eq!(group_key("a.b.example.com"), "example.com");
assert_eq!(group_key("example.com"), "example.com");
assert_eq!(group_key("EXAMPLE.COM."), "example.com");
assert_eq!(group_key("a.example.co.uk"), "co.uk");
assert_eq!(group_key("b.other.co.uk"), "co.uk");
assert_eq!(group_key("127.0.0.1"), "127.0.0.1");
assert_eq!(group_key("localhost"), "localhost");
}
#[test]
fn an_unrecognised_ca_gets_no_budget() {
let d = tmpdir("unknown");
assert!(
IssuanceBudget::load("https://acme.internal/directory", d.to_str().unwrap()).is_none()
);
assert!(CaProfile::for_directory("https://ca.example.com/dir").is_none());
assert!(
CaProfile::for_directory("https://acme-staging-v02.api.letsencrypt.org/directory")
.is_some()
);
}
#[test]
fn metrics_never_publish_the_customers_hostname_set() {
let d = tmpdir("metrics");
let b = budget(&d);
let domains = vec!["secret-customer.example.com".to_string()];
let text = render_metrics(&b, &domains, 1_800_000_000);
assert!(
!text.contains("secret-customer"),
"the hostname set leaked into /metrics:\n{text}"
);
assert!(
text.contains("bucket=\"registered_domain\",key=\"example.com\""),
"{text}"
);
assert!(text.contains("edgeguard_acme_budget_remaining"), "{text}");
}
#[test]
fn remaining_counts_down_and_is_reported_per_bucket() {
let d = tmpdir("remaining");
let mut b = budget(&d);
let domains = vec!["edge.example.com".to_string()];
let now = 1_800_000_000;
let before = b.remaining(&domains, now);
let set_before = before
.iter()
.find(|(bk, _, _)| *bk == Bucket::IdentifierSet)
.unwrap()
.2;
assert_eq!(set_before, 5);
b.debit(&domains, now).unwrap();
let after = b.remaining(&domains, now);
let set_after = after
.iter()
.find(|(bk, _, _)| *bk == Bucket::IdentifierSet)
.unwrap()
.2;
assert_eq!(
set_after, 4,
"an operator must see the budget draining before it is gone"
);
}
}