use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use crate::memory::events::{MemoryEventSink, MemoryTimelineEvent};
pub const DEFAULT_RUNS_PER_HOUR: u32 = 12;
pub const DEFAULT_MAX_CONCURRENT: u32 = 1;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BackgroundBudgetConfig {
pub runs_per_window: u32,
pub window: Duration,
pub max_concurrent: u32,
}
impl Default for BackgroundBudgetConfig {
fn default() -> Self {
Self {
runs_per_window: DEFAULT_RUNS_PER_HOUR,
window: Duration::from_hours(1),
max_concurrent: DEFAULT_MAX_CONCURRENT,
}
}
}
#[derive(Debug, Default)]
struct RealmBudgetState {
starts: Vec<Instant>,
concurrent: u32,
}
struct BudgetInner {
config: BackgroundBudgetConfig,
realms: HashMap<String, RealmBudgetState>,
event_sink: Option<Arc<dyn MemoryEventSink>>,
}
#[derive(Clone)]
pub struct BackgroundBudget {
inner: Arc<Mutex<BudgetInner>>,
}
pub struct BudgetPermit {
inner: Arc<Mutex<BudgetInner>>,
realm: String,
}
impl std::fmt::Debug for BudgetPermit {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("BudgetPermit")
.field("realm", &self.realm)
.finish_non_exhaustive()
}
}
impl Drop for BudgetPermit {
fn drop(&mut self) {
let mut inner = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some(state) = inner.realms.get_mut(&self.realm) {
state.concurrent = state.concurrent.saturating_sub(1);
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum BudgetDenied {
WindowExhausted { used: u32, cap: u32 },
ConcurrencyCeiling { cap: u32 },
}
impl std::fmt::Display for BudgetDenied {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::WindowExhausted { used, cap } => {
write!(f, "window budget exhausted ({used}/{cap} runs)")
}
Self::ConcurrencyCeiling { cap } => {
write!(f, "concurrency ceiling reached ({cap} in flight)")
}
}
}
}
impl BackgroundBudget {
pub fn new(config: BackgroundBudgetConfig) -> Self {
Self {
inner: Arc::new(Mutex::new(BudgetInner {
config,
realms: HashMap::new(),
event_sink: None,
})),
}
}
pub fn set_event_sink(&self, sink: Arc<dyn MemoryEventSink>) {
self.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.event_sink = Some(sink);
}
pub fn try_acquire(&self, realm: &str, stage: &str) -> Result<BudgetPermit, BudgetDenied> {
let mut inner = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let window = inner.config.window;
let runs_cap = inner.config.runs_per_window;
let concurrent_cap = inner.config.max_concurrent;
let state = inner.realms.entry(realm.to_string()).or_default();
let now = Instant::now();
state
.starts
.retain(|start| now.duration_since(*start) < window);
let denied = if state.concurrent >= concurrent_cap {
Some(BudgetDenied::ConcurrencyCeiling {
cap: concurrent_cap,
})
} else if state.starts.len() as u32 >= runs_cap {
Some(BudgetDenied::WindowExhausted {
used: state.starts.len() as u32,
cap: runs_cap,
})
} else {
None
};
if let Some(denied) = denied {
tracing::warn!(
realm,
stage,
reason = %denied,
"agent memory background budget: run skipped"
);
if let Some(sink) = inner.event_sink.as_ref() {
sink.emit(MemoryTimelineEvent::BudgetDenied {
realm: realm.to_string(),
stage: stage.to_string(),
reason: denied.to_string(),
});
}
return Err(denied);
}
state.starts.push(now);
state.concurrent += 1;
Ok(BudgetPermit {
inner: self.inner.clone(),
realm: realm.to_string(),
})
}
}
#[cfg(test)]
#[allow(clippy::expect_used, clippy::unwrap_used)]
mod tests {
use super::*;
fn config(runs: u32, concurrent: u32) -> BackgroundBudgetConfig {
BackgroundBudgetConfig {
runs_per_window: runs,
window: Duration::from_hours(1),
max_concurrent: concurrent,
}
}
#[test]
fn window_cap_denies_after_budget_spent() {
let budget = BackgroundBudget::new(config(2, 10));
let p1 = budget.try_acquire("realm-a", "distiller").expect("run 1");
drop(p1);
let p2 = budget.try_acquire("realm-a", "distiller").expect("run 2");
drop(p2);
let denied = budget
.try_acquire("realm-a", "distiller")
.expect_err("third run in window must deny");
assert!(
matches!(denied, BudgetDenied::WindowExhausted { used: 2, cap: 2 }),
"{denied:?}"
);
budget
.try_acquire("realm-b", "distiller")
.expect("other realm has its own window");
}
#[test]
fn denial_emits_timeline_event_when_sink_wired() {
let budget = BackgroundBudget::new(config(1, 10));
let sink = Arc::new(crate::memory::events::CollectingEventSink::new());
budget.set_event_sink(sink.clone());
let _permit = budget.try_acquire("realm-a", "steward").expect("first");
let _ = budget
.try_acquire("realm-a", "steward")
.expect_err("window spent");
assert_eq!(sink.types(), vec!["memory.budget.denied"]);
let events = sink.events.lock().unwrap();
assert!(matches!(
&events[0],
MemoryTimelineEvent::BudgetDenied { realm, stage, .. }
if realm == "realm-a" && stage == "steward"
));
}
#[test]
fn concurrency_ceiling_releases_on_drop() {
let budget = BackgroundBudget::new(config(10, 1));
let permit = budget.try_acquire("realm-a", "distiller").expect("first");
let denied = budget
.try_acquire("realm-a", "distiller")
.expect_err("second concurrent run must deny");
assert!(
matches!(denied, BudgetDenied::ConcurrencyCeiling { cap: 1 }),
"{denied:?}"
);
drop(permit);
budget
.try_acquire("realm-a", "distiller")
.expect("slot released on drop");
}
}