use std::collections::{BTreeMap, VecDeque};
use std::fmt;
use std::num::NonZeroU32;
use crate::script::Tenant;
pub const MAX_QUEUE_PER_TENANT: usize = 65_536;
pub const MAX_TENANT_WEIGHT: u32 = 1_000;
pub const MAX_TOTAL_QUEUE: usize = 65_536;
pub const MAX_WEIGHTED_TENANTS: usize = 5_000;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum TenancyError {
TooManyWeightedTenants {
limit: usize,
},
WeightTooHigh {
weight: u32,
limit: u32,
},
}
impl fmt::Display for TenancyError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match *self {
Self::TooManyWeightedTenants { limit } => {
write!(
formatter,
"the policy names weights for more than {limit} tenants"
)
}
Self::WeightTooHigh { weight, limit } => {
write!(
formatter,
"a tenant weight of {weight} exceeds the bound of {limit}"
)
}
}
}
}
impl std::error::Error for TenancyError {}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct TenancyPolicy {
per_tenant_limit: usize,
queue_per_tenant: usize,
queue_total: usize,
weights: BTreeMap<Tenant, NonZeroU32>,
}
impl TenancyPolicy {
#[must_use]
pub fn new(per_tenant_limit: usize, queue_per_tenant: usize) -> Self {
Self {
per_tenant_limit: per_tenant_limit.max(1),
queue_per_tenant: queue_per_tenant.min(MAX_QUEUE_PER_TENANT),
queue_total: MAX_TOTAL_QUEUE,
weights: BTreeMap::new(),
}
}
#[must_use]
pub fn with_queue_total(mut self, queue_total: usize) -> Self {
self.queue_total = queue_total.min(MAX_TOTAL_QUEUE);
self
}
pub fn with_weight(
mut self,
tenant: &Tenant,
weight: NonZeroU32,
) -> Result<Self, TenancyError> {
let declared = weight.get();
if declared > MAX_TENANT_WEIGHT {
let refusal = Err(TenancyError::WeightTooHigh {
weight: declared,
limit: MAX_TENANT_WEIGHT,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "with_weight: returning an error to the caller");
return refusal;
}
if !self.weights.contains_key(tenant) && self.weights.len() >= MAX_WEIGHTED_TENANTS {
let refusal = Err(TenancyError::TooManyWeightedTenants {
limit: MAX_WEIGHTED_TENANTS,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "with_weight: returning an error to the caller");
return refusal;
}
self.weights.insert(tenant.clone(), weight);
Ok(self)
}
#[must_use]
pub const fn per_tenant_limit(&self) -> usize {
self.per_tenant_limit
}
#[must_use]
pub const fn queue_per_tenant(&self) -> usize {
self.queue_per_tenant
}
#[must_use]
pub const fn queue_total(&self) -> usize {
self.queue_total
}
#[must_use]
pub fn weight_of(&self, tenant: &Tenant) -> NonZeroU32 {
match self.weights.get(tenant).copied() {
Some(weight) => weight,
None => NonZeroU32::MIN,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum SpawnRefused {
TenantAtCapacity {
tenant: Tenant,
limit: usize,
},
SupervisorQueueFull {
limit: usize,
},
Cancelled,
}
impl fmt::Display for SpawnRefused {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match *self {
Self::TenantAtCapacity { ref tenant, limit } => write!(
formatter,
"tenant {tenant} already holds {limit} waiting admissions"
),
Self::SupervisorQueueFull { limit } => write!(
formatter,
"the supervisor already holds {limit} waiting admissions across its tenants"
),
Self::Cancelled => {
formatter.write_str("the supervisor is cancelled and admits no new work")
}
}
}
}
impl std::error::Error for SpawnRefused {}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum TryArrival<P> {
Immediate(P),
Contended,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum Arrival<P> {
Immediate(P),
Queued,
Refused {
limit: usize,
},
SupervisorFull {
limit: usize,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct Grant<W, P> {
pub tenant: Tenant,
pub waiter: W,
pub permit: P,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum GrantOutcome<W, P> {
Granted(Grant<W, P>),
Idle(P),
}
#[derive(Debug)]
struct Entry<W> {
queue: VecDeque<W>,
abandoned: usize,
in_flight: usize,
deficit: u32,
turn_open: bool,
on_ring: bool,
}
impl<W> Default for Entry<W> {
fn default() -> Self {
Self {
queue: VecDeque::new(),
abandoned: 0,
in_flight: 0,
deficit: 0,
turn_open: false,
on_ring: false,
}
}
}
impl<W> Entry<W> {
fn live(&self) -> usize {
self.queue.len().saturating_sub(self.abandoned)
}
fn is_idle(&self) -> bool {
self.queue.is_empty() && self.in_flight == 0 && !self.on_ring
}
}
#[derive(Debug)]
pub struct DeficitRoundRobin<W> {
policy: TenancyPolicy,
entries: BTreeMap<Tenant, Entry<W>>,
ring: VecDeque<Tenant>,
retained: usize,
}
impl<W> DeficitRoundRobin<W> {
#[must_use]
pub fn new(policy: TenancyPolicy) -> Self {
Self {
policy,
entries: BTreeMap::new(),
ring: VecDeque::new(),
retained: 0,
}
}
#[must_use]
pub fn policy(&self) -> &TenancyPolicy {
&self.policy
}
#[must_use]
pub fn has_eligible(&self) -> bool {
!self.ring.is_empty()
}
#[must_use]
pub fn in_flight_of(&self, tenant: &Tenant) -> usize {
self.entries.get(tenant).map_or(0, |entry| entry.in_flight)
}
#[must_use]
pub fn queued_of(&self, tenant: &Tenant) -> usize {
self.entries.get(tenant).map_or(0, Entry::live)
}
#[must_use]
pub fn retained_of(&self, tenant: &Tenant) -> usize {
self.entries
.get(tenant)
.map_or(0, |entry| entry.queue.len())
}
#[must_use]
pub fn retained_total(&self) -> usize {
self.retained
}
#[must_use]
pub fn tenants_with_queues(&self) -> usize {
self.entries
.values()
.filter(|entry| !entry.queue.is_empty())
.count()
}
pub fn arrive<P, F>(&mut self, tenant: &Tenant, waiter: W, acquire: F) -> Arrival<P>
where
F: FnOnce() -> Option<P>,
{
let ceiling = self.policy.per_tenant_limit;
let queue_limit = self.policy.queue_per_tenant;
let entry = self.entries.entry(tenant.clone()).or_default();
if self.ring.is_empty()
&& entry.in_flight < ceiling
&& let Some(permit) = acquire()
{
entry.in_flight = entry.in_flight.saturating_add(1);
return Arrival::Immediate(permit);
}
if entry.live() >= queue_limit {
if entry.is_idle() {
self.entries.remove(tenant);
}
return Arrival::Refused { limit: queue_limit };
}
if self.retained >= self.policy.queue_total {
if entry.is_idle() {
self.entries.remove(tenant);
}
return Arrival::SupervisorFull {
limit: self.policy.queue_total,
};
}
entry.queue.push_back(waiter);
self.retained = self.retained.saturating_add(1);
if !entry.on_ring && entry.in_flight < ceiling {
entry.on_ring = true;
self.ring.push_back(tenant.clone());
}
Arrival::Queued
}
pub fn try_arrive<P, F>(&mut self, tenant: &Tenant, acquire: F) -> TryArrival<P>
where
F: FnOnce() -> Option<P>,
{
let ceiling = self.policy.per_tenant_limit;
if !self.ring.is_empty() || self.entries.contains_key(tenant) {
return TryArrival::Contended;
}
let Some(permit) = acquire() else {
return TryArrival::Contended;
};
let entry = self.entries.entry(tenant.clone()).or_default();
if entry.in_flight >= ceiling {
return TryArrival::Contended;
}
entry.in_flight = entry.in_flight.saturating_add(1);
TryArrival::Immediate(permit)
}
pub fn note_release(&mut self, tenant: &Tenant) {
let Some(entry) = self.entries.get_mut(tenant) else {
return;
};
entry.in_flight = entry.in_flight.saturating_sub(1);
if !entry.queue.is_empty()
&& !entry.on_ring
&& entry.in_flight < self.policy.per_tenant_limit
{
entry.on_ring = true;
self.ring.push_back(tenant.clone());
}
if entry.is_idle() {
self.entries.remove(tenant);
}
}
pub fn note_abandoned<F>(&mut self, tenant: &Tenant, is_live: &mut F)
where
F: FnMut(&W) -> bool,
{
let Some(entry) = self.entries.get_mut(tenant) else {
return;
};
entry.abandoned = entry.abandoned.saturating_add(1).min(entry.queue.len());
if entry.abandoned > entry.live() {
let before = entry.queue.len();
entry.queue.retain(|waiter| is_live(waiter));
self.retained = self
.retained
.saturating_sub(before.saturating_sub(entry.queue.len()));
entry.abandoned = 0;
}
}
pub fn grant<P, F>(&mut self, permit: P, is_live: &mut F) -> GrantOutcome<W, P>
where
F: FnMut(&W) -> bool,
{
loop {
let Some(front) = self.ring.front().cloned() else {
return GrantOutcome::Idle(permit);
};
let Some(entry) = self.entries.get_mut(&front) else {
self.ring.pop_front();
continue;
};
while entry.queue.front().is_some_and(|waiter| !is_live(waiter)) {
entry.queue.pop_front();
entry.abandoned = entry.abandoned.saturating_sub(1);
self.retained = self.retained.saturating_sub(1);
}
if entry.queue.is_empty() {
self.ring.pop_front();
entry.on_ring = false;
entry.deficit = 0;
entry.turn_open = false;
if entry.is_idle() {
self.entries.remove(&front);
}
continue;
}
if !entry.turn_open {
let weight = self.policy.weight_of(&front).get();
entry.deficit = entry.deficit.saturating_add(weight);
entry.turn_open = true;
}
let Some(waiter) = entry.queue.pop_front() else {
continue;
};
self.retained = self.retained.saturating_sub(1);
entry.deficit = entry.deficit.saturating_sub(1);
entry.in_flight = entry.in_flight.saturating_add(1);
let drained = entry.queue.is_empty();
let at_ceiling = entry.in_flight >= self.policy.per_tenant_limit;
if drained || at_ceiling || entry.deficit == 0 {
self.ring.pop_front();
entry.turn_open = false;
if drained || at_ceiling {
entry.on_ring = false;
entry.deficit = 0;
if entry.is_idle() {
self.entries.remove(&front);
}
} else {
self.ring.push_back(front.clone());
}
}
return GrantOutcome::Granted(Grant {
tenant: front,
waiter,
permit,
});
}
}
}
#[cfg(test)]
mod tests {
use super::{
Arrival, DeficitRoundRobin, GrantOutcome, MAX_QUEUE_PER_TENANT, MAX_TENANT_WEIGHT,
TenancyError, TenancyPolicy,
};
use crate::script::Tenant;
use std::error::Error;
use std::num::NonZeroU32;
type Core = DeficitRoundRobin<u32>;
fn tenant(name: &str) -> Result<Tenant, Box<dyn Error>> {
Ok(Tenant::new(name)?)
}
fn all_live(_waiter: &u32) -> bool {
true
}
fn parked<P>(arrival: Arrival<P>) -> bool {
matches!(arrival, Arrival::Queued)
}
#[test]
fn a_fresh_core_admits_immediately_from_the_pool() -> Result<(), Box<dyn Error>> {
let mut core: Core = DeficitRoundRobin::new(TenancyPolicy::new(2, 4));
let t0 = tenant("t-0")?;
assert_eq!(
core.arrive(&t0, 100, || Some(7_u32)),
Arrival::Immediate(7),
"an empty ring and an under-ceiling tenant admit at once"
);
assert_eq!(
core.in_flight_of(&t0),
1,
"the immediate admission is charged to the tenant"
);
let queued: Arrival<u32> = Arrival::Queued;
assert_eq!(
core.arrive(&t0, 101, || None),
queued,
"with a ring member present, the next arrival queues rather than jumping it"
);
Ok(())
}
#[test]
fn a_refusal_names_the_queue_bound_and_retains_nothing() -> Result<(), Box<dyn Error>> {
let mut core: Core = DeficitRoundRobin::new(TenancyPolicy::new(1, 2));
let t0 = tenant("t-0")?;
let cell = std::cell::Cell::new(1_u32);
let take = || {
let left = cell.get();
if left == 0 {
return None;
}
cell.set(left.saturating_sub(1));
Some(left.saturating_sub(1))
};
assert_eq!(
core.arrive(&t0, 1, take),
Arrival::Immediate(0),
"the pool's one permit is taken at once"
);
let queued: Arrival<u32> = Arrival::Queued;
assert_eq!(core.arrive(&t0, 2, || None), queued, "the second waits");
assert_eq!(core.arrive(&t0, 3, || None), queued, "the third waits");
let spent = || None::<u32>;
assert_eq!(
core.arrive(&t0, 4, spent),
Arrival::Refused { limit: 2 },
"the fourth is past the queue bound and the refusal names it"
);
assert_eq!(
core.queued_of(&t0),
2,
"only the two live waiters are counted"
);
Ok(())
}
#[test]
fn two_backlogged_tenants_alternate_permits() -> Result<(), Box<dyn Error>> {
let mut core: Core = DeficitRoundRobin::new(TenancyPolicy::new(4, 8));
let first = tenant("a")?;
let second = tenant("b")?;
let first_arrival = core.arrive(&first, 1, || None::<u32>);
assert!(parked(first_arrival), "the first tenant queues");
let second_arrival = core.arrive(&second, 2, || None::<u32>);
assert!(parked(second_arrival), "the second tenant queues");
for (tenant, waiter) in [(&first, 3_u32), (&second, 4_u32)] {
let arrived = core.arrive(tenant, waiter, || None::<u32>);
assert!(parked(arrived), "the backlog grows on both tenants");
}
let first_third = core.arrive(&first, 5, || None::<u32>);
assert!(parked(first_third), "the first tenant queues a third");
let second_third = core.arrive(&second, 6, || None::<u32>);
assert!(parked(second_third), "the second tenant queues a third");
let mut served: Vec<u32> = Vec::new();
for permit in 0..4_u32 {
let mut live = all_live;
let outcome = core.grant(permit, &mut live);
let GrantOutcome::Granted(grant) = outcome else {
return Err("a queued tenant must take the permit".into());
};
served.push(grant.waiter);
core.note_release(&grant.tenant);
}
assert_eq!(
served,
vec![1, 2, 3, 4],
"equal weights alternate strictly: first, second, first, second"
);
Ok(())
}
#[test]
fn a_double_weight_takes_two_permits_per_turn() -> Result<(), Box<dyn Error>> {
let heavy = tenant("heavy")?;
let light = tenant("light")?;
let two: NonZeroU32 = 2_u32.try_into()?;
let policy = TenancyPolicy::new(8, 16).with_weight(&heavy, two)?;
let mut core: Core = DeficitRoundRobin::new(policy);
let queued: Arrival<u32> = Arrival::Queued;
assert_eq!(core.arrive(&heavy, 1, || None), queued, "heavy queues");
assert_eq!(core.arrive(&light, 2, || None), queued, "light queues");
assert_eq!(
core.arrive(&heavy, 3, || None),
queued,
"heavy queues again"
);
assert_eq!(
core.arrive(&light, 4, || None),
queued,
"light queues again"
);
for waiter in 5..11_u32 {
let _heavy_arrival = core.arrive(&heavy, waiter, || None::<u32>);
}
for waiter in 11..17_u32 {
let _light_arrival = core.arrive(&light, waiter, || None::<u32>);
}
let mut heavy_served: u32 = 0;
let mut light_served: u32 = 0;
for permit in 0..9_u32 {
let mut live = all_live;
let outcome = core.grant(permit, &mut live);
let GrantOutcome::Granted(grant) = outcome else {
return Err("a backlogged tenant must take the permit".into());
};
if grant.tenant == heavy {
heavy_served = heavy_served.saturating_add(1);
} else {
light_served = light_served.saturating_add(1);
}
core.note_release(&grant.tenant);
}
assert_eq!(
(heavy_served, light_served),
(6, 3),
"weight 2 takes two permits for every one weight 1 takes"
);
Ok(())
}
#[test]
fn a_tenant_at_its_ceiling_is_parked_until_it_releases() -> Result<(), Box<dyn Error>> {
let mut core: Core = DeficitRoundRobin::new(TenancyPolicy::new(1, 8));
let loud = tenant("loud")?;
let quiet = tenant("quiet")?;
assert_eq!(
core.arrive(&loud, 1, || Some(10_u32)),
Arrival::Immediate(10),
"the loud tenant takes the pool's one permit"
);
let queued: Arrival<u32> = Arrival::Queued;
assert_eq!(core.arrive(&loud, 2, || None), queued, "its next waits");
assert_eq!(
core.arrive(&quiet, 3, || None),
queued,
"the quiet tenant waits"
);
core.note_release(&loud);
let mut live = all_live;
let first = core.grant(20, &mut live);
let GrantOutcome::Granted(grant) = first else {
return Err("the quiet tenant is eligible and must be served".into());
};
assert_eq!(
grant.waiter, 3,
"the permit went to the tenant under its ceiling"
);
core.note_release(&grant.tenant);
let second = core.grant(21, &mut live);
let GrantOutcome::Granted(next) = second else {
return Err("the loud tenant is eligible again and must be served".into());
};
assert_eq!(
next.waiter, 2,
"the loud tenant's waiter runs as soon as it has room"
);
Ok(())
}
#[test]
fn an_abandoned_waiter_is_skipped_and_stops_spending_the_bound() -> Result<(), Box<dyn Error>> {
let mut core: Core = DeficitRoundRobin::new(TenancyPolicy::new(1, 1));
let t0 = tenant("t-0")?;
assert_eq!(
core.arrive(&t0, 7, || Some(1_u32)),
Arrival::Immediate(1),
"one in flight"
);
let queued: Arrival<u32> = Arrival::Queued;
assert_eq!(
core.arrive(&t0, 8, || None),
queued,
"the queue bound has room"
);
let spent = || None::<u32>;
assert_eq!(
core.arrive(&t0, 9, spent),
Arrival::Refused { limit: 1 },
"the bound is reached while the waiter lives"
);
let mut gone = |waiter: &u32| *waiter != 8;
core.note_abandoned(&t0, &mut gone);
assert_eq!(
core.arrive(&t0, 10, || None),
queued,
"an abandoned waiter no longer spends the tenant's queue room"
);
core.note_release(&t0);
let mut live = |waiter: &u32| *waiter != 8;
let outcome = core.grant(2, &mut live);
let GrantOutcome::Granted(grant) = outcome else {
return Err("the live waiter must be served".into());
};
assert_eq!(
grant.waiter, 10,
"the abandoned head was skipped, not served"
);
assert_eq!(
core.queued_of(&t0),
0,
"the skipped waiter is gone from the queue and the count with it"
);
Ok(())
}
#[test]
fn an_idle_tenant_costs_nothing_after_it_drains() -> Result<(), Box<dyn Error>> {
let mut core: Core = DeficitRoundRobin::new(TenancyPolicy::new(2, 2));
let t0 = tenant("t-0")?;
assert_eq!(
core.arrive(&t0, 1, || Some(1_u32)),
Arrival::Immediate(1),
"admitted"
);
core.note_release(&t0);
assert_eq!(
core.tenants_with_queues(),
0,
"a drained tenant holds no queue and no entry"
);
Ok(())
}
#[test]
fn an_idle_pool_answers_idle_rather_than_spinning() -> Result<(), Box<dyn Error>> {
let mut core: Core = DeficitRoundRobin::new(TenancyPolicy::new(4, 8));
let t0 = tenant("t-0")?;
assert_eq!(
core.arrive(&t0, 1, || Some(1)),
Arrival::Immediate(1),
"one in flight"
);
core.note_release(&t0);
assert!(
!core.has_eligible(),
"with nothing waiting and the tenant drained, no tenant is eligible"
);
let mut live = all_live;
let outcome = core.grant(9, &mut live);
let GrantOutcome::Idle(permit) = outcome else {
return Err("a permit nobody can take is handed back".into());
};
assert_eq!(
permit, 9,
"the idle permit is returned to the caller's pool"
);
Ok(())
}
#[test]
fn the_declared_bounds_are_clamped_to_their_ceilings() {
let policy = TenancyPolicy::new(0, MAX_QUEUE_PER_TENANT.saturating_add(1));
assert_eq!(
policy.per_tenant_limit(),
1,
"a zero ceiling reads as one, as Supervisor::new reads its own bound"
);
assert_eq!(
policy.queue_per_tenant(),
MAX_QUEUE_PER_TENANT,
"a queue bound past the ceiling is clamped to it"
);
}
#[test]
fn a_weight_past_its_bound_is_refused_naming_both_numbers() -> Result<(), Box<dyn Error>> {
let heavy = tenant("heavy")?;
let past = MAX_TENANT_WEIGHT.saturating_add(1);
let Some(too_heavy) = NonZeroU32::new(past) else {
return Err("a weight one past its bound is never zero".into());
};
let outcome = TenancyPolicy::new(1, 1).with_weight(&heavy, too_heavy);
assert_eq!(
outcome.err(),
Some(TenancyError::WeightTooHigh {
weight: past,
limit: MAX_TENANT_WEIGHT,
}),
"a weight above the ceiling is refused with the weight and the bound"
);
Ok(())
}
}