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 },
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:<delegation subject>",
];
#[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 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>,
subject: Option<&str>,
) -> 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 { id } => subject.is_some_and(|s| s == id),
}
}
}
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 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,
by: &crate::core::Operator,
at: crate::core::Timestamp,
reason: &str,
) -> Result<(), StoreError>;
async fn lift_halt(&self, scope: &HaltScope) -> 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>;
}
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(())
}