use std::fmt::Debug;
use async_trait::async_trait;
use crate::core::{RunId, Spend, StoreError, Timestamp};
#[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,
}
#[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 },
}
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 }
}
#[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}"),
}
}
#[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);
}
None
}
#[must_use]
pub fn covers(&self, agent: Option<&crate::journal::AgentIdentity>) -> 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),
}
}
}
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}"),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Halt {
pub scope: HaltScope,
pub reason: String,
}
#[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 spent {spent} of its {limit} {unit} for period {period} — \
this does not reset until the period does"
)]
SpentOut {
tenant: String,
period: String,
unit: &'static str,
spent: u64,
limit: u64,
},
#[error("{scope} is halted by an operator (tenant '{tenant}'): {reason}")]
Halted {
tenant: String,
scope: HaltScope,
reason: String,
},
#[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,
limit: Option<u32>,
at: Timestamp,
) -> Result<(), QuotaError>;
async fn release(&self, run: RunId) -> Result<(), StoreError>;
async fn set_halt(&self, scope: &HaltScope, reason: Option<&str>) -> Result<(), 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>;
}
pub fn check_spend(
tenant: &str,
period: &str,
quota: &TenantQuota,
spent: Spend,
) -> Result<(), QuotaError> {
if let Some(limit) = quota.max_tokens_per_period
&& spent.tokens >= limit
{
return Err(QuotaError::SpentOut {
tenant: tenant.to_owned(),
period: period.to_owned(),
unit: "tokens",
spent: spent.tokens,
limit,
});
}
if let Some(limit) = quota.max_minor_units_per_period
&& spent.minor_units >= limit
{
return Err(QuotaError::SpentOut {
tenant: tenant.to_owned(),
period: period.to_owned(),
unit: "minor units",
spent: spent.minor_units,
limit,
});
}
Ok(())
}