use crate::core::{EffectKey, Phase, RunId, Spend, StepId, Timestamp};
use crate::quota::{
QuotaError, QuotaSettlement, QuotaStore, RateCeiling, RateReservation, SpendHold, TenantQuota,
};
use super::conformance::Report;
pub async fn check(store: &dyn QuotaStore, report: &mut Report) {
let at = Timestamp::from_unix_timestamp(1_760_000_000).expect("a valid test instant");
zero_admits_nothing(store, at, report).await;
ceiling_refuses_and_frees(store, at, report).await;
reserving_twice_takes_one_slot(store, at, report).await;
spend_accrues_per_period(store, report).await;
spend_holds(store, at, report).await;
halt(store, report).await;
rate(store, report).await;
}
fn slots(n: u32) -> TenantQuota {
TenantQuota {
max_concurrent_runs: Some(n),
..TenantQuota::default()
}
}
fn tokens_per_period(limit: u64) -> TenantQuota {
TenantQuota {
max_tokens_per_period: Some(limit),
..TenantQuota::default()
}
}
fn hold(period: &str, tokens: u64) -> SpendHold {
SpendHold {
period: period.to_owned(),
amount: Spend::tokens(tokens),
}
}
#[allow(clippy::too_many_lines, clippy::many_single_char_names)]
async fn spend_holds(store: &dyn QuotaStore, at: Timestamp, report: &mut Report) {
let period = "2998-01";
let quota = tokens_per_period(1_000);
let (a, b, c, d) = (
RunId::generate(),
RunId::generate(),
RunId::generate(),
RunId::generate(),
);
report.checked += 1;
for (run, amount) in [(a, 600), (c, 400)] {
if let Err(e) = store
.reserve(run, "a, Some(&hold(period, amount)), at)
.await
{
report.record(
"holds that fit the ceiling are admitted",
format!(
"holding {amount} of a 1000-token period failed with `{e}` — a period \
filled exactly is one its admitted work can reach and not pass"
),
);
return;
}
}
report.checked += 1;
if let Err(e) = store.reserve(a, "a, Some(&hold(period, 600)), at).await {
report.record(
"holding one run twice is idempotent",
format!("a retried admission was refused against its own hold: {e}"),
);
}
match store.reserved(period).await {
Ok(s) if s.tokens == 1_000 => {}
Ok(s) => report.record(
"a hold is what the period counts",
format!(
"two holds of 600 and 400 read back as {} reserved",
s.tokens
),
),
Err(e) => report.record("reading what is reserved", format!("{e}")),
}
report.checked += 1;
match store.reserve(b, "a, Some(&hold(period, 1)), at).await {
Err(QuotaError::SpentOut {
settled: 0,
reserved: 1_000,
requested: 1,
limit: 1_000,
..
}) => {}
Err(e) => report.record(
"a full period refuses and names what is reserved",
format!("refused with `{e}` rather than naming settled and reserved apart"),
),
Ok(()) => {
report.record(
"outstanding holds count against the ceiling",
"a period whose holds already reach its ceiling admitted one more, so \
the ceiling counts settled spend only and every run admitted before any \
settles may spend its whole budget",
);
let _ = store.release(b).await;
}
}
report.checked += 1;
if let Ok(held) = store.running_runs(100).await
&& held.contains(&b)
{
report.record(
"a refused hold takes no slot",
"the slot was inserted although the spend hold was refused",
);
}
report.checked += 1;
let pass = |epoch, tokens, concludes| QuotaSettlement {
run: a,
epoch,
period: Some(period.to_owned()),
spend: Spend::tokens(tokens),
release_slot: epoch == 1,
concludes,
};
if let Err(e) = store.settle(&pass(1, 200, false)).await {
report.record("settling a pass under a hold", format!("{e}"));
return;
}
match (store.spent(period).await, store.reserved(period).await) {
(Ok(spent), Ok(held)) if spent.tokens == 200 && held.tokens == 800 => {}
(Ok(spent), Ok(held)) => report.record(
"a pass settlement moves its spend out of the hold",
format!(
"after a 200-token pass the period reads {} settled and {} reserved — a \
suspended run is charged twice unless its settled spend leaves its hold",
spent.tokens, held.tokens
),
),
(Err(e), _) | (_, Err(e)) => report.record("reading the period", format!("{e}")),
}
report.checked += 1;
if let Err(e) = store.settle(&pass(2, 100, true)).await {
report.record("settling the concluding pass", format!("{e}"));
return;
}
match (store.spent(period).await, store.reserved(period).await) {
(Ok(spent), Ok(held)) if spent.tokens == 300 && held.tokens == 400 => {}
(Ok(spent), Ok(held)) => report.record(
"the concluding settlement releases what the run did not spend",
format!(
"after the run concluded having spent 300 of its 600 the period reads \
{} settled and {} reserved",
spent.tokens, held.tokens
),
),
(Err(e), _) | (_, Err(e)) => report.record("reading the period", format!("{e}")),
}
report.checked += 1;
if store.settle(&pass(2, 100, false)).await.is_ok() {
report.record(
"one pass key names one conclusion",
"the same run/epoch was accepted as concluding and as not",
);
}
report.checked += 1;
if let Err(e) = store.reserve(d, "a, Some(&hold(period, 300)), at).await {
report.record(
"a released hold makes room",
format!("300 tokens freed by a conclusion could not be held again: {e}"),
);
}
report.checked += 1;
let later = "2998-02";
if let Err(e) = store.carry(c, later).await {
report.record("carrying a hold", format!("{e}"));
}
match (store.reserved(period).await, store.reserved(later).await) {
(Ok(old), Ok(new)) if old.tokens == 300 && new.tokens == 400 => {}
(Ok(old), Ok(new)) => report.record(
"a carried hold moves whole",
format!(
"after carrying a 400-token hold the old period holds {} and the new {}",
old.tokens, new.tokens
),
),
(Err(e), _) | (_, Err(e)) => report.record("reading carried holds", format!("{e}")),
}
report.checked += 1;
match store.reservations(100).await {
Ok(held) => {
let mut runs: Vec<RunId> = held.iter().map(|h| h.run).collect();
runs.sort();
let mut want = vec![c, d];
want.sort();
if runs != want {
report.record(
"the runs holding spend are listed",
format!("expected {want:?} to hold spend, the listing says {held:?}"),
);
}
}
Err(e) => report.record("listing holds", format!("{e}")),
}
report.checked += 1;
for run in [c, d] {
if let Err(e) = store.release(run).await {
report.record("releasing a hold", format!("{e}"));
}
}
match (store.reserved(period).await, store.reserved(later).await) {
(Ok(old), Ok(new)) if old.is_free() && new.is_free() => {}
(Ok(old), Ok(new)) => report.record(
"a release drops the hold with the slot",
format!(
"after releasing every holder, {old} and {new} are still held — a crash \
before the admission record strands a reservation over no run"
),
),
(Err(e), _) | (_, Err(e)) => report.record("reading released holds", format!("{e}")),
}
}
pub async fn check_race(first: &dyn QuotaStore, second: &dyn QuotaStore, report: &mut Report) {
let at = Timestamp::from_unix_timestamp(1_760_000_000).expect("a valid test instant");
let quota = tokens_per_period(1_000);
let period = "2997-01";
let runs: Vec<RunId> = (0..20).map(|_| RunId::generate()).collect();
let spend = hold(period, 100);
let attempts = runs.iter().enumerate().map(|(i, run)| {
let store = if i % 2 == 0 { first } else { second };
let (quota, spend) = ("a, &spend);
async move { store.reserve(*run, quota, Some(spend), at).await }
});
let outcomes = futures_util::future::join_all(attempts).await;
report.checked += 1;
let admitted = outcomes.iter().filter(|o| o.is_ok()).count();
let refused = outcomes
.iter()
.filter(|o| matches!(o, Err(QuotaError::SpentOut { .. })))
.count();
if admitted != 10 || refused != 10 {
report.record(
"concurrent admissions at the spend ceiling admit only what fits",
format!(
"{admitted} holds of 100 landed in a 1000-token period ({refused} refused, \
{} failed otherwise) — the check and the insert are not one decision",
20 - admitted - refused
),
);
}
match first.reserved(period).await {
Ok(s) if s.tokens == 1_000 => {}
Ok(s) => report.record(
"concurrent admissions at the spend ceiling admit only what fits",
format!("the period holds {} tokens after the race", s.tokens),
),
Err(e) => report.record("reading the raced period", format!("{e}")),
}
for run in runs {
let _ = first.release(run).await;
}
}
async fn zero_admits_nothing(store: &dyn QuotaStore, at: Timestamp, report: &mut Report) {
report.checked += 1;
let run = RunId::generate();
match store.reserve(run, &slots(0), None, at).await {
Err(QuotaError::TooManyRuns { .. }) => {}
Err(e) => report.record(
"a ceiling of zero admits nothing",
format!("reserving under a zero ceiling failed with `{e}` rather than a refusal"),
),
Ok(()) => {
report.record(
"a ceiling of zero admits nothing",
"a run was admitted under a ceiling of zero — the value an \
operator sets to stop a tenant dead, and the one a limit \
compared inside its counting loop never sees",
);
let _ = store.release(run).await;
}
}
}
async fn ceiling_refuses_and_frees(store: &dyn QuotaStore, at: Timestamp, report: &mut Report) {
let first = RunId::generate();
let second = RunId::generate();
report.checked += 1;
if let Err(e) = store.reserve(first, &slots(1), None, at).await {
report.record(
"a run fits under a ceiling of one",
format!("the first reservation failed: {e}"),
);
return;
}
report.checked += 1;
match store.reserve(second, &slots(1), None, at).await {
Err(QuotaError::TooManyRuns { running, .. }) => {
report.checked += 1;
if running == 0 {
report.record(
"a refusal reports how many runs are executing",
"the refusal said zero runs are executing, which tells an \
operator asking why they are throttled precisely nothing",
);
}
}
Err(e) => report.record(
"a tenant at its ceiling is refused",
format!("failed with `{e}` rather than reporting the ceiling"),
),
Ok(()) => report.record(
"a tenant at its ceiling is refused",
"a second run was admitted past a ceiling of one, so the ceiling \
bounds nothing",
),
}
report.checked += 1;
if let Err(e) = store
.settle(&QuotaSettlement {
run: first,
epoch: 1,
period: None,
spend: Spend::default(),
release_slot: true,
concludes: true,
})
.await
{
report.record("settling and releasing a slot", format!("{e}"));
return;
}
report.checked += 1;
match store.reserve(second, &slots(1), None, at).await {
Ok(()) => {
let _ = store.release(second).await;
}
Err(e) => report.record(
"releasing makes room",
format!(
"the slot freed by a finished run could not be reused: {e}. A \
ceiling is back-pressure, and one that never frees is a tenant \
permanently stopped by its first burst"
),
),
}
}
async fn reserving_twice_takes_one_slot(
store: &dyn QuotaStore,
at: Timestamp,
report: &mut Report,
) {
let run = RunId::generate();
report.checked += 1;
if let Err(e) = store.reserve(run, &slots(1), None, at).await {
report.record("reserving a run", format!("{e}"));
return;
}
report.checked += 1;
match store.reserve(run, &slots(1), None, at).await {
Ok(()) => {}
Err(e) => report.record(
"reserving one run twice is idempotent",
format!(
"a retried admission was refused against its own slot ({e}), so \
a transient error during admission costs the tenant capacity \
until something releases a run it never really started"
),
),
}
report.checked += 1;
match store.running().await {
Ok(1) => {}
Ok(n) => report.record(
"reserving one run twice takes one slot",
format!(
"{n} slots are held for one run, so every retry permanently shrinks the ceiling"
),
),
Err(e) => report.record("counting running runs", format!("{e}")),
}
report.checked += 1;
match store.running_runs(100).await {
Ok(held) if held == vec![run] => {}
Ok(held) => report.record(
"naming the runs that hold slots",
format!(
"one run holds a slot and the listing says {held:?} — an operator \
told only how many cannot tell a live run from a slot a dead \
instance stranded, which is the case the accounting exists for"
),
),
Err(e) => report.record("naming the runs that hold slots", format!("{e}")),
}
let _ = store.release(run).await;
report.checked += 1;
match store.running_runs(100).await {
Ok(held) if held.is_empty() => {}
Ok(held) => report.record(
"releasing a slot takes the run off the listing",
format!("the released run is still listed as holding a slot: {held:?}"),
),
Err(e) => report.record(
"releasing a slot takes the run off the listing",
format!("{e}"),
),
}
}
async fn spend_accrues_per_period(store: &dyn QuotaStore, report: &mut Report) {
let (this, next) = ("2999-01", "2999-02");
let run = RunId::generate();
let first = QuotaSettlement {
run,
epoch: 1,
period: Some(this.to_owned()),
spend: Spend::tokens(400),
release_slot: false,
concludes: false,
};
let second = QuotaSettlement {
run,
epoch: 2,
period: Some(this.to_owned()),
spend: Spend::tokens(600),
release_slot: false,
concludes: false,
};
report.checked += 1;
for settlement in [&first, &second] {
if let Err(e) = store.settle(settlement).await {
report.record("settling spend", format!("{e}"));
return;
}
}
report.checked += 1;
match store.spent(this).await {
Ok(s) if s.tokens == 1_000 => {}
Ok(s) => report.record(
"accruals sum",
format!(
"two accruals of 400 and 600 totalled {} rather than 1000. \
Reading a total, adding to it and writing it back loses one of \
two concurrent updates — and what it loses is spend a tenant \
has already incurred, so the ceiling drifts upward under load",
s.tokens
),
),
Err(e) => report.record("reading spend", format!("{e}")),
}
report.checked += 1;
if let Err(e) = store.settle(&first).await {
report.record("retrying an identical settlement", format!("{e}"));
}
match store.spent(this).await {
Ok(s) if s.tokens == 1_000 => {}
Ok(s) => report.record(
"an identical settlement accrues once",
format!("retrying one pass changed the total to {} tokens", s.tokens),
),
Err(e) => report.record("reading spend after a settlement retry", format!("{e}")),
}
report.checked += 1;
let changed = QuotaSettlement {
spend: Spend::tokens(401),
..first.clone()
};
if store.settle(&changed).await.is_ok() {
report.record(
"one pass key names one exact settlement",
"the same run/epoch accepted a different spend, so a retry can rewrite the bill",
);
}
report.checked += 1;
match store.spent(next).await {
Ok(s) if s.tokens == 0 => {}
Ok(s) => report.record(
"periods are independent",
format!(
"an untouched period already reports {} tokens, so a ceiling \
would never reset and a tenant is billed forever for one month",
s.tokens
),
),
Err(e) => report.record("reading an untouched period", format!("{e}")),
}
}
#[allow(clippy::too_many_lines)]
async fn halt(store: &dyn QuotaStore, report: &mut Report) {
use crate::quota::HaltScope;
let tenant = HaltScope::Tenant;
let agent = HaltScope::agent("payments-clerk");
let revision = HaltScope::revision(crate::core::Digest::of(b"a manifest revision"));
let standing = |halts: &[crate::quota::Halt], scope: &HaltScope| -> Option<String> {
halts
.iter()
.find(|h| &h.scope == scope)
.map(|h| h.reason.clone())
};
report.checked += 1;
match store.halts().await {
Ok(halts) if halts.is_empty() => {}
Ok(halts) => report.record(
"a fresh tenant is not halted",
format!(
"an untouched tenant reports {halts:?}, so a plane would refuse \
every run it was never told to refuse"
),
),
Err(e) => report.record("reading the halts", format!("{e}")),
}
report.checked += 1;
let thrower =
|actor: &str| crate::core::Operator::asserted(actor).expect("a battery names its operator");
let at = crate::core::Timestamp::from_unix_timestamp(1_700_000_000).expect("a fixed instant");
if let Err(e) = store
.set_halt(&tenant, &thrower("ops-alice"), at, "incident 42")
.await
{
report.record("setting the halt", format!("{e}"));
}
match store.halts().await {
Ok(halts) if standing(&halts, &tenant).as_deref() == Some("incident 42") => {}
Ok(other) => report.record(
"the halt survives being written",
format!(
"after halting, the store reports {other:?} — a switch that does \
not read back is one an operator believes they threw"
),
),
Err(e) => report.record("reading the halt back", format!("{e}")),
}
report.checked += 1;
if let Err(e) = store
.set_halt(&tenant, &thrower("ops-alice"), at, "incident 43")
.await
{
report.record("re-halting", format!("{e}"));
}
match store.halts().await {
Ok(halts) if standing(&halts, &tenant).as_deref() == Some("incident 43") => {}
Ok(other) => report.record(
"re-halting replaces the reason",
format!("expected the newer reason, got {other:?}"),
),
Err(e) => report.record("re-reading the halt", format!("{e}")),
}
report.checked += 1;
if let Err(e) = store
.set_halt(&agent, &thrower("ops-bob"), at, "agent 12 is looping")
.await
{
report.record("halting one agent", format!("{e}"));
}
if let Err(e) = store
.set_halt(&revision, &thrower("ops-bob"), at, "bad deploy")
.await
{
report.record("halting one revision", format!("{e}"));
}
match store.halts().await {
Ok(halts)
if standing(&halts, &tenant).as_deref() == Some("incident 43")
&& standing(&halts, &agent).as_deref() == Some("agent 12 is looping")
&& standing(&halts, &revision).as_deref() == Some("bad deploy") => {}
Ok(other) => report.record(
"scopes are independent",
format!(
"a narrow halt overwrote a broader one, or was not kept: {other:?} — \
an incident that widens must not un-stop what was already stopped"
),
),
Err(e) => report.record("reading several standing halts", format!("{e}")),
}
report.checked += 1;
if let Err(e) = store.lift_halt(&agent).await {
report.record("lifting one scope", format!("{e}"));
}
match store.halts().await {
Ok(halts)
if standing(&halts, &agent).is_none()
&& standing(&halts, &tenant).as_deref() == Some("incident 43") => {}
Ok(other) => report.record(
"lifting one scope leaves the others",
format!(
"after lifting the agent halt the store reports {other:?} — lifting \
a narrow stop must not lift the broad one it sits under"
),
),
Err(e) => report.record("reading a partly lifted halt", format!("{e}")),
}
report.checked += 1;
for scope in [&tenant, &revision] {
if let Err(e) = store.lift_halt(scope).await {
report.record("lifting the halt", format!("{e}"));
}
}
match store.halts().await {
Ok(halts) if halts.is_empty() => {}
Ok(other) => report.record(
"a lifted halt stays lifted",
format!(
"the tenant is still halted by {other:?} after the stop was \
lifted, so an incident that is over never ends"
),
),
Err(e) => report.record("reading a lifted halt", format!("{e}")),
}
report.checked += 1;
match store.lift_halt(&tenant).await {
Ok(false) => {}
Ok(true) => report.record(
"lifting an unset halt",
"the store reported that a halt was standing when none was".to_owned(),
),
Err(e) => report.record("lifting an unset halt", format!("{e}")),
}
report.checked += 1;
let by = crate::core::Operator::authenticated("ops-carol").expect("a name");
if let Err(e) = store.set_halt(&tenant, &by, at, "incident 44").await {
report.record("halting with an authenticated operator", format!("{e}"));
}
match store.halts().await {
Ok(halts) => match halts.iter().find(|h| h.scope == tenant) {
Some(h) if h.by == by && h.at == at => {}
Some(h) => report.record(
"a halt keeps who threw it",
format!(
"the store read back {:?} at {:?} rather than {by:?} at {at:?} — an \
emergency stop nobody is named on cannot be asked about afterwards",
h.by, h.at
),
),
None => report.record(
"a halt keeps who threw it",
"the halt did not read back at all".to_owned(),
),
},
Err(e) => report.record("reading an attributed halt", format!("{e}")),
}
if let Err(e) = store.lift_halt(&tenant).await {
report.record("clearing the attributed halt", format!("{e}"));
}
}
fn dispatch(n: u32) -> EffectKey {
EffectKey::derive(StepId(n), Phase::Forward, 0, 1, "tool.call", b"{}")
}
fn rated(
grant: &str,
run: RunId,
key: EffectKey,
ceiling: RateCeiling,
at: Timestamp,
) -> RateReservation {
RateReservation {
grant: grant.to_owned(),
run,
dispatch: key,
ceilings: vec![ceiling],
at,
exempt: false,
}
}
fn instant(seconds: i64) -> Timestamp {
Timestamp::from_unix_timestamp(seconds).expect("a valid test instant")
}
#[allow(clippy::too_many_lines)]
async fn rate(store: &dyn QuotaStore, report: &mut Report) {
let three = RateCeiling {
count: 3,
window_seconds: 60,
};
let at = instant(1_760_000_000);
let grant = "tool://conformance/refund";
let run = RunId::generate();
report.checked += 1;
for n in 0..3 {
if let Err(e) = store
.reserve_rate(&rated(grant, run, dispatch(n), three, at))
.await
{
report.record(
"a rate ceiling admits its count",
format!("dispatch {} of 3 was refused: {e}", n + 1),
);
return;
}
}
report.checked += 1;
match store
.reserve_rate(&rated(grant, run, dispatch(3), three, at))
.await
{
Err(QuotaError::RateLimited { reached: 3, .. }) => {}
Err(e) => report.record(
"a full window refuses and says what it counted",
format!("the fourth dispatch failed with `{e}`"),
),
Ok(()) => report.record(
"a full window refuses",
"a fourth dispatch was admitted under a ceiling of three",
),
}
report.checked += 1;
if let Err(e) = store
.reserve_rate(&rated(grant, run, dispatch(0), three, at))
.await
{
report.record(
"re-reserving a present dispatch spends nothing and is not judged",
format!(
"a dispatch already counted was refused against its own row: {e} — every retry of a call the window admitted would be refused"
),
);
}
report.checked += 1;
match store.rate_room(grant, &[three], at).await {
Err(QuotaError::RateLimited { reached: 3, .. }) => {}
Err(QuotaError::RateLimited { reached, .. }) => report.record(
"re-reserving a present dispatch spends nothing",
format!("three dispatches and a retry read back as {reached}"),
),
other => report.record(
"asking for room takes none and answers a full window",
format!("a full window answered {other:?}"),
),
}
let same = "tool://conformance/same-call";
let two = RateCeiling {
count: 2,
window_seconds: 60,
};
report.checked += 1;
for _ in 0..2 {
if let Err(e) = store
.reserve_rate(&rated(same, RunId::generate(), dispatch(0), two, at))
.await
{
report.record(
"two runs making the same call each count",
format!("the second run's identical call was refused: {e}"),
);
return;
}
}
report.checked += 1;
if store
.reserve_rate(&rated(same, RunId::generate(), dispatch(0), two, at))
.await
.is_ok()
{
report.record(
"two runs making the same call each count",
"a third run's identical call was admitted under a ceiling of two, so \
identical calls from different runs share one row and the ceiling \
admits every run making the same call",
);
}
let edge = "tool://conformance/boundary";
let hourly = RateCeiling {
count: 20,
window_seconds: 3_600,
};
let boundary: i64 = 1_760_004_000;
let edge_run = RunId::generate();
report.checked += 1;
for n in 0..20 {
if let Err(e) = store
.reserve_rate(&rated(
edge,
edge_run,
dispatch(n),
hourly,
instant(boundary - 60),
))
.await
{
report.record(
"a rate ceiling admits its count",
format!(
"dispatch {} of 20 before the boundary was refused: {e}",
n + 1
),
);
return;
}
}
let past = (20..40)
.map(|n| rated(edge, edge_run, dispatch(n), hourly, instant(boundary + 60)))
.collect::<Vec<_>>();
let mut admitted = 0;
for r in &past {
if store.reserve_rate(r).await.is_ok() {
admitted += 1;
}
}
report.checked += 1;
if admitted > 0 {
report.record(
"a window boundary does not double the ceiling",
format!(
"twenty dispatches a minute before an hour boundary and {admitted} more a \
minute after it were admitted under twenty an hour — the window is a \
fixed bucket, not a sliding one"
),
);
}
report.checked += 1;
if let Err(e) = store
.reserve_rate(&rated(
edge,
edge_run,
dispatch(99),
hourly,
instant(boundary - 60 + 3_600 + 1),
))
.await
{
report.record(
"a window that has passed makes room",
format!("a dispatch an hour after the burst was refused: {e}"),
);
}
let undo = "tool://conformance/undo";
let one = RateCeiling {
count: 1,
window_seconds: 60,
};
let undo_run = RunId::generate();
report.checked += 1;
let _ = store
.reserve_rate(&rated(undo, undo_run, dispatch(0), one, at))
.await;
let mut exempt = rated(undo, undo_run, dispatch(1), one, at);
exempt.exempt = true;
if let Err(e) = store.reserve_rate(&exempt).await {
report.record(
"an undo is counted and never refused",
format!("an exempt reservation was refused: {e}"),
);
}
report.checked += 1;
match store
.rate_room(
undo,
&[RateCeiling {
count: 2,
window_seconds: 60,
}],
at,
)
.await
{
Err(QuotaError::RateLimited { reached: 2, .. }) => {}
other => report.record(
"an undo is counted",
format!("a dispatch and an undo under a ceiling of two left room: {other:?}"),
),
}
}
pub async fn check_rate_tenants(
first: &dyn QuotaStore,
other: &dyn QuotaStore,
report: &mut Report,
) {
let at = instant(1_760_000_000);
let grant = "tool://conformance/tenants";
let one = RateCeiling {
count: 1,
window_seconds: 60,
};
report.checked += 1;
if let Err(e) = first
.reserve_rate(&rated(grant, RunId::generate(), dispatch(0), one, at))
.await
{
report.record("a rate ceiling admits its count", format!("{e}"));
return;
}
if let Err(e) = other
.reserve_rate(&rated(grant, RunId::generate(), dispatch(0), one, at))
.await
{
report.record(
"one tenant's rate count does not throttle another",
format!(
"tenant '{}' was refused by tenant '{}''s count: {e}",
other.tenant(),
first.tenant()
),
);
}
}
pub async fn check_rate_race(first: &dyn QuotaStore, second: &dyn QuotaStore, report: &mut Report) {
let at = instant(1_760_000_000);
let grant = "tool://conformance/race";
let ten = RateCeiling {
count: 10,
window_seconds: 3_600,
};
let attempts = (0..40u32).map(|n| {
let store = if n % 2 == 0 { first } else { second };
let reservation = rated(grant, RunId::generate(), dispatch(n), ten, at);
async move { store.reserve_rate(&reservation).await }
});
let outcomes = futures_util::future::join_all(attempts).await;
report.checked += 1;
let admitted = outcomes.iter().filter(|o| o.is_ok()).count();
let refused = outcomes
.iter()
.filter(|o| matches!(o, Err(QuotaError::RateLimited { .. })))
.count();
if admitted != 10 || refused != 30 {
report.record(
"concurrent dispatches at a rate ceiling admit exactly the ceiling",
format!(
"{admitted} of 40 dispatches landed under a ceiling of ten ({refused} \
refused, {} failed otherwise) — the count and the insert are not one \
decision",
40 - admitted - refused
),
);
}
}