use std::fmt::Debug;
use async_trait::async_trait;
use crate::core::{Budget, EffectKey, RunId, Spend, StoreError, Timestamp};
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct RateCeiling {
pub count: u32,
pub window_seconds: u64,
}
impl std::fmt::Display for RateCeiling {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"at most {} per {} second(s)",
self.count, self.window_seconds
)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RateReservation {
pub grant: String,
pub run: RunId,
pub dispatch: EffectKey,
pub ceilings: Vec<RateCeiling>,
pub at: Timestamp,
pub exempt: bool,
}
#[must_use]
pub(crate) fn rate_window_start(at: i64, window_seconds: u64) -> i64 {
at.saturating_sub(i64::try_from(window_seconds).unwrap_or(i64::MAX))
}
pub fn check_rate(
tenant: &str,
grant: &str,
ceilings: &[RateCeiling],
instants: &[i64],
at: Timestamp,
) -> Result<(), QuotaError> {
let at = at.unix_timestamp();
if let Some(ceiling) = ceilings
.iter()
.find(|c| c.window_seconds == 0 || c.window_seconds > MAX_RATE_WINDOW_SECONDS)
{
return Err(QuotaError::UncountableRate {
grant: grant.to_owned(),
ceiling: *ceiling,
});
}
for ceiling in ceilings {
let start = rate_window_start(at, ceiling.window_seconds);
let reached = instants.iter().filter(|&&t| t > start).count();
if reached >= ceiling.count as usize {
return Err(QuotaError::RateLimited {
tenant: tenant.to_owned(),
grant: grant.to_owned(),
ceiling: *ceiling,
reached: u64::try_from(reached).unwrap_or(u64::MAX),
});
}
}
Ok(())
}
pub const MAX_RATE_WINDOW_SECONDS: u64 = 31 * 86_400;
#[must_use]
pub fn rate_prune_floor(at: Timestamp) -> i64 {
rate_window_start(at.unix_timestamp(), MAX_RATE_WINDOW_SECONDS)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct QuotaSettlement {
pub run: RunId,
pub epoch: u64,
pub period: Option<String>,
pub spend: Spend,
pub release_slot: bool,
pub concludes: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SpendHold {
pub period: String,
pub amount: Spend,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Held {
pub run: RunId,
pub period: String,
pub remaining: Spend,
}
pub fn reservation(
quota: &TenantQuota,
budget: &Budget,
width: u64,
) -> Result<Option<Spend>, &'static str> {
if !quota.bounds_spend() {
return Ok(None);
}
let hold = |bounded: bool,
ceiling: Option<u64>,
call: Option<u64>,
names: (&'static str, &'static str)|
-> Result<u64, &'static str> {
if !bounded {
return Ok(0);
}
let ceiling = ceiling.ok_or(names.0)?;
let call = call.ok_or(names.1)?;
Ok(ceiling.saturating_add(width.saturating_mul(call)))
};
Ok(Some(Spend {
tokens: hold(
quota.max_tokens_per_period.is_some(),
budget.max_tokens,
budget.max_call_tokens,
("max_tokens", "max_call_tokens"),
)?,
minor_units: hold(
quota.max_minor_units_per_period.is_some(),
budget.max_minor_units,
budget.max_call_minor_units,
("max_minor_units", "max_call_minor_units"),
)?,
}))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct TenantQuota {
pub max_concurrent_runs: Option<u32>,
pub max_tokens_per_period: Option<u64>,
pub max_minor_units_per_period: Option<u64>,
pub period: Period,
}
impl TenantQuota {
#[must_use]
pub const fn is_unlimited(&self) -> bool {
self.max_concurrent_runs.is_none()
&& self.max_tokens_per_period.is_none()
&& self.max_minor_units_per_period.is_none()
}
#[must_use]
pub const fn bounds_spend(&self) -> bool {
self.max_tokens_per_period.is_some() || self.max_minor_units_per_period.is_some()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum Period {
#[default]
Monthly,
Daily,
}
impl Period {
#[must_use]
pub fn key_for(self, at: Timestamp) -> String {
let d = at.date();
match self {
Self::Monthly => format!("{:04}-{:02}", d.year(), u8::from(d.month())),
Self::Daily => format!("{:04}-{:02}-{:02}", d.year(), u8::from(d.month()), d.day()),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub enum HaltScope {
Tenant,
Agent { name: String },
Revision { digest: crate::core::Digest },
Subject { id: String },
}
impl HaltScope {
pub fn agent(name: impl Into<String>) -> Self {
Self::Agent { name: name.into() }
}
#[must_use]
pub const fn revision(digest: crate::core::Digest) -> Self {
Self::Revision { digest }
}
pub fn subject(id: impl Into<String>) -> Self {
Self::Subject { id: id.into() }
}
#[must_use]
pub fn key(&self) -> String {
match self {
Self::Tenant => "tenant".to_owned(),
Self::Agent { name } => format!("agent:{name}"),
Self::Revision { digest } => format!("revision:{digest}"),
Self::Subject { id } => format!("subject:{id}"),
}
}
pub const FORMS: [&'static str; 4] = [
"tenant",
"agent:<metadata.name>",
"revision:<manifest digest>",
"subject:<principal on the chain>",
];
#[must_use]
pub fn forms() -> String {
Self::FORMS
.iter()
.map(|f| format!("'{f}'"))
.collect::<Vec<_>>()
.join(", ")
}
#[must_use]
pub fn parse(key: &str) -> Option<Self> {
if key == "tenant" {
return Some(Self::Tenant);
}
if let Some(name) = key.strip_prefix("agent:")
&& !name.is_empty()
{
return Some(Self::agent(name));
}
if let Some(hex) = key.strip_prefix("revision:") {
return crate::core::Digest::from_hex(hex).ok().map(Self::revision);
}
if let Some(id) = key.strip_prefix("subject:")
&& !id.is_empty()
{
return Some(Self::subject(id));
}
None
}
#[must_use]
pub fn withdraws_from(&self, chain: &crate::core::Delegation) -> Option<&str> {
self.withdrawn_subject()
.filter(|id| chain.links().any(|link| link.id == *id))
}
#[must_use]
pub fn withdrawn_subject(&self) -> Option<&str> {
match self {
Self::Subject { id } => Some(id.as_str()),
Self::Tenant | Self::Agent { .. } | Self::Revision { .. } => None,
}
}
#[must_use]
pub fn covers(
&self,
agent: Option<&crate::journal::AgentIdentity>,
chain: Option<&crate::core::Delegation>,
) -> bool {
match self {
Self::Tenant => true,
Self::Agent { name } => agent.is_some_and(|a| &a.name == name),
Self::Revision { digest } => agent.is_some_and(|a| &a.digest == digest),
Self::Subject { .. } => chain.is_some_and(|c| self.withdraws_from(c).is_some()),
}
}
}
impl std::fmt::Display for HaltScope {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Tenant => f.write_str("the whole tenant"),
Self::Agent { name } => write!(f, "agent '{name}'"),
Self::Revision { digest } => write!(f, "manifest revision {digest}"),
Self::Subject { id } => write!(f, "everything acting for '{id}'"),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Halt {
pub scope: HaltScope,
pub reason: String,
pub by: crate::core::Operator,
pub at: crate::core::Timestamp,
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum QuotaError {
#[error(
"tenant '{tenant}' already has {running} runs executing, which is its limit — \
this is back-pressure, not a fault: retry when one finishes"
)]
TooManyRuns { tenant: String, running: u32 },
#[error(
"tenant '{tenant}' has {settled} {unit} settled and {reserved} reserved by open \
runs against its {limit} for period {period}, and this run holds up to \
{requested} — settled spend resets with the period; reserved spend is released \
as its runs conclude, and `agentplane attention` names the stopped ones"
)]
SpentOut {
tenant: String,
period: String,
unit: &'static str,
settled: u64,
reserved: u64,
requested: u64,
limit: u64,
},
#[error(
"tenant '{tenant}' bounds {unit} per period, and this run declares no {field}, so \
what it can cost is unbounded and cannot be reserved — a manifest sets \
spec.budgets.max_tokens / max_minor_units, and bounds one call with \
max_input_tokens (and a price) on every model role; an embedder sets the same \
fields on its Budget"
)]
Unbounded {
tenant: String,
unit: &'static str,
field: &'static str,
},
#[error("{scope} is halted by an operator (tenant '{tenant}'): {reason}")]
Halted {
tenant: String,
scope: HaltScope,
reason: String,
},
#[error(
"'{grant}' is limited to {ceiling} for tenant '{tenant}', and {reached} dispatch(es) \
already fall in the window ending now — wait for the window to pass"
)]
RateLimited {
tenant: String,
grant: String,
ceiling: RateCeiling,
reached: u64,
},
#[error(
"'{grant}' states {ceiling}, which no store can count — a window is at least \
one second and at most {MAX_RATE_WINDOW_SECONDS} seconds"
)]
UncountableRate { grant: String, ceiling: RateCeiling },
#[error(
"the quota store could not be reached, and a ceiling that yields under load is not a ceiling: {0}"
)]
Unavailable(String),
}
impl From<StoreError> for QuotaError {
fn from(e: StoreError) -> Self {
Self::Unavailable(e.to_string())
}
}
#[async_trait]
pub trait QuotaStore: Send + Sync + Debug {
fn tenant(&self) -> &str;
async fn reserve(
&self,
run: RunId,
quota: &TenantQuota,
hold: Option<&SpendHold>,
at: Timestamp,
) -> Result<(), QuotaError>;
async fn release(&self, run: RunId) -> Result<(), StoreError>;
async fn carry(&self, run: RunId, period: &str) -> Result<(), StoreError>;
async fn reservations(&self, limit: usize) -> Result<Vec<Held>, StoreError>;
async fn reserved(&self, period: &str) -> Result<Spend, StoreError>;
async fn set_halt(
&self,
scope: &HaltScope,
by: &crate::core::Operator,
at: crate::core::Timestamp,
reason: &str,
) -> Result<(), StoreError>;
async fn lift_halt(&self, scope: &HaltScope) -> Result<bool, StoreError>;
async fn lift_halt_if(&self, standing: &Halt) -> Result<bool, StoreError>;
async fn halts(&self) -> Result<Vec<Halt>, StoreError>;
async fn settle(&self, settlement: &QuotaSettlement) -> Result<(), StoreError>;
async fn spent(&self, period: &str) -> Result<Spend, StoreError>;
async fn running(&self) -> Result<u32, StoreError>;
async fn running_runs(&self, limit: usize) -> Result<Vec<RunId>, StoreError>;
async fn reserve_rate(&self, reservation: &RateReservation) -> Result<(), QuotaError>;
async fn rate_room(
&self,
grant: &str,
ceilings: &[RateCeiling],
at: Timestamp,
) -> Result<(), QuotaError>;
}
pub fn check_spend(
tenant: &str,
period: &str,
quota: &TenantQuota,
settled: Spend,
reserved: Spend,
request: Spend,
) -> Result<(), QuotaError> {
for (unit, limit, settled, reserved, requested) in [
(
"tokens",
quota.max_tokens_per_period,
settled.tokens,
reserved.tokens,
request.tokens,
),
(
"minor units",
quota.max_minor_units_per_period,
settled.minor_units,
reserved.minor_units,
request.minor_units,
),
] {
if let Some(limit) = limit
&& settled.saturating_add(reserved).saturating_add(requested) > limit
{
return Err(QuotaError::SpentOut {
tenant: tenant.to_owned(),
period: period.to_owned(),
unit,
settled,
reserved,
requested,
limit,
});
}
}
Ok(())
}