use std::collections::{BTreeMap, BTreeSet};
use std::sync::Arc;
use std::time::Duration;
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
#[non_exhaustive]
pub enum ConstraintKey {
Host(String),
Origin(String),
Egress(String),
}
impl ConstraintKey {
#[must_use]
pub fn host(host: impl Into<String>) -> Self {
Self::Host(host.into())
}
#[must_use]
pub fn origin(origin: impl Into<String>) -> Self {
Self::Origin(origin.into())
}
#[must_use]
pub fn egress(egress: impl Into<String>) -> Self {
Self::Egress(egress.into())
}
}
impl std::fmt::Display for ConstraintKey {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match *self {
Self::Host(ref host) => write!(formatter, "host:{host}"),
Self::Origin(ref origin) => write!(formatter, "origin:{origin}"),
Self::Egress(ref egress) => write!(formatter, "egress:{egress}"),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum Resolved {
Origin(String),
Failed,
DoesNotExist,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
#[non_exhaustive]
pub enum RulesState {
#[default]
NotFetched,
Allowed,
Disallowed {
rule: String,
},
NoRules,
Unreachable {
since: Duration,
},
}
impl RulesState {
#[must_use]
fn effective(&self, at: Duration, grace: Duration) -> Self {
match *self {
Self::Unreachable { since } if at.saturating_sub(since) >= grace => Self::NoRules,
ref unchanged => unchanged.clone(),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
#[non_exhaustive]
pub enum DeferralKind {
RetryAfter,
RateLimited,
OriginPolicyUnreachable,
RulesNotFetched,
OriginUnmapped,
ResolutionFailed,
}
impl DeferralKind {
#[must_use]
pub const fn needs_prerequisite(&self) -> bool {
match *self {
Self::OriginUnmapped
| Self::ResolutionFailed
| Self::RulesNotFetched
| Self::OriginPolicyUnreachable => true,
Self::RetryAfter | Self::RateLimited => false,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum RejectKind {
Disallowed {
rule: String,
},
DoesNotExist,
Exhausted,
}
#[derive(Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum Admission {
Admit(InFlightPermit),
Defer {
resume_at: Duration,
kind: DeferralKind,
constraint: ConstraintKey,
},
Reject {
kind: RejectKind,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum Outcome {
Success,
CacheHit,
Failure,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum CompletionError {
ForeignPermit,
UnknownPermit,
}
impl std::fmt::Display for CompletionError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match *self {
Self::ForeignPermit => formatter.write_str(
"the permit was minted by a different frontier, so no reservation here matches it",
),
Self::UnknownPermit => formatter.write_str(
"no live reservation matches the permit: it was already completed, or never \
admitted here",
),
}
}
}
impl std::error::Error for CompletionError {}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum PolitenessError {
ZeroWindow,
InvertedWindow {
min_window: u32,
max_window: u32,
},
InvalidEgressWindow {
egress_window: u32,
max_window: u32,
},
InvalidDecreasePercent {
decrease_percent: u32,
},
}
impl std::fmt::Display for PolitenessError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match *self {
Self::ZeroWindow => formatter.write_str("a window of zero admits nothing"),
Self::InvertedWindow {
min_window,
max_window,
} => write!(
formatter,
"minimum window {min_window} exceeds maximum window {max_window}"
),
Self::InvalidEgressWindow {
egress_window,
max_window,
} => write!(
formatter,
"egress window {egress_window} is zero or above the maximum {max_window}"
),
Self::InvalidDecreasePercent { decrease_percent } => write!(
formatter,
"decrease factor {decrease_percent} is outside 1..=100"
),
}
}
}
impl std::error::Error for PolitenessError {}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub struct PolitenessPolicy {
min_window: u32,
max_window: u32,
egress_window: u32,
increase_after: u32,
decrease_percent: u32,
unreachable_grace: Duration,
retry_after_failure: Duration,
}
impl PolitenessPolicy {
pub const DEFAULT: Self = Self {
min_window: 1,
max_window: 4,
egress_window: 4,
increase_after: 8,
decrease_percent: 70,
unreachable_grace: Duration::from_secs(2_592_000),
retry_after_failure: Duration::from_secs(60),
};
pub fn new(
min_window: u32,
max_window: u32,
egress_window: u32,
increase_after: u32,
decrease_percent: u32,
) -> Result<Self, PolitenessError> {
if min_window == 0 {
let refusal = Err(PolitenessError::ZeroWindow);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "new: returning an error to the caller");
return refusal;
}
if min_window > max_window {
let refusal = Err(PolitenessError::InvertedWindow {
min_window,
max_window,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "new: returning an error to the caller");
return refusal;
}
if egress_window == 0 || egress_window > max_window {
let refusal = Err(PolitenessError::InvalidEgressWindow {
egress_window,
max_window,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "new: returning an error to the caller");
return refusal;
}
if decrease_percent == 0 || decrease_percent > 100 {
let refusal = Err(PolitenessError::InvalidDecreasePercent { decrease_percent });
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "new: returning an error to the caller");
return refusal;
}
Ok(Self {
min_window,
max_window,
egress_window,
increase_after,
decrease_percent,
..Self::DEFAULT
})
}
#[must_use]
pub const fn min_window(&self) -> u32 {
self.min_window
}
#[must_use]
pub const fn max_window(&self) -> u32 {
self.max_window
}
#[must_use]
pub const fn egress_window(&self) -> u32 {
self.egress_window
}
#[must_use]
pub const fn increase_after(&self) -> u32 {
self.increase_after
}
#[must_use]
pub const fn decrease_percent(&self) -> u32 {
self.decrease_percent
}
#[must_use]
pub const fn unreachable_grace(&self) -> Duration {
self.unreachable_grace
}
#[must_use]
pub const fn retry_after_failure(&self) -> Duration {
self.retry_after_failure
}
#[must_use]
pub fn decrease(&self, window: u32) -> u32 {
let scaled = window
.saturating_mul(self.decrease_percent)
.saturating_div(100);
scaled.clamp(self.min_window, self.max_window)
}
#[must_use]
pub fn increase(&self, window: u32) -> u32 {
window
.saturating_add(1)
.clamp(self.min_window, self.max_window)
}
}
impl Default for PolitenessPolicy {
fn default() -> Self {
Self::DEFAULT
}
}
#[derive(Debug)]
struct Issuer;
#[derive(Debug)]
#[must_use = "a permit dropped without `Frontier::complete` keeps its slots reserved; dropping \
the handle is not evidence that the request left flight"]
pub struct InFlightPermit {
issuer: Arc<Issuer>,
ordinal: u64,
host: String,
origin: Option<ConstraintKey>,
egress: ConstraintKey,
}
impl InFlightPermit {
#[must_use]
pub fn host(&self) -> &str {
&self.host
}
#[must_use]
pub fn holds(&self, key: &ConstraintKey) -> bool {
self.keys().iter().any(|held| held == key)
}
fn keys(&self) -> Vec<ConstraintKey> {
let mut keys = Vec::with_capacity(3);
keys.push(ConstraintKey::host(self.host.clone()));
keys.extend(self.origin.clone());
keys.push(self.egress.clone());
keys
}
fn server_keys(&self) -> Vec<ConstraintKey> {
let mut keys = Vec::with_capacity(2);
keys.push(ConstraintKey::host(self.host.clone()));
keys.extend(self.origin.clone());
keys
}
}
impl PartialEq for InFlightPermit {
fn eq(&self, other: &Self) -> bool {
self.ordinal == other.ordinal
}
}
impl Eq for InFlightPermit {}
#[derive(Debug)]
enum Verdict {
Admit,
Defer {
resume_at: Duration,
kind: DeferralKind,
constraint: ConstraintKey,
},
Reject {
kind: RejectKind,
},
}
impl Verdict {
fn refusal(self) -> Option<Admission> {
match self {
Self::Admit => None,
Self::Defer {
resume_at,
kind,
constraint,
} => Some(Admission::Defer {
resume_at,
kind,
constraint,
}),
Self::Reject { kind } => Some(Admission::Reject { kind }),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct KeyState {
window: u32,
in_flight: u32,
successes: u32,
blocked_until: Duration,
rules: RulesState,
}
fn later(current: Option<Verdict>, candidate: Verdict) -> Verdict {
let Some(current) = current else {
return candidate;
};
let take_candidate = match (deferral_order(¤t), deferral_order(&candidate)) {
(Some(held), Some(offered)) => offered > held,
(None, Some(_)) => true,
(Some(_), None) | (None, None) => false,
};
if take_candidate { candidate } else { current }
}
fn deferral_order(verdict: &Verdict) -> Option<(Duration, DeferralKind, &ConstraintKey)> {
match *verdict {
Verdict::Defer {
resume_at,
kind,
ref constraint,
} => Some((resume_at, kind, constraint)),
Verdict::Admit | Verdict::Reject { .. } => None,
}
}
#[derive(Debug)]
pub struct Frontier {
policy: PolitenessPolicy,
egress: ConstraintKey,
resolutions: BTreeMap<String, Resolved>,
constraints: BTreeMap<String, ConstraintKey>,
state: BTreeMap<ConstraintKey, KeyState>,
abandons: BTreeSet<String>,
issuer: Arc<Issuer>,
reservations: BTreeSet<u64>,
next_ordinal: u64,
}
impl Frontier {
#[must_use]
pub fn new(policy: PolitenessPolicy, egress: impl Into<String>) -> Self {
let egress = ConstraintKey::egress(egress);
let mut frontier = Self {
policy,
egress,
resolutions: BTreeMap::new(),
constraints: BTreeMap::new(),
state: BTreeMap::new(),
abandons: BTreeSet::new(),
issuer: Arc::new(Issuer),
reservations: BTreeSet::new(),
next_ordinal: 1,
};
let key = frontier.egress.clone();
let width = policy.egress_window();
frontier.entry(&key).window = width;
frontier
}
#[must_use]
pub fn policy(&self) -> &PolitenessPolicy {
&self.policy
}
#[must_use]
pub fn in_flight(&self, key: &ConstraintKey) -> u32 {
self.state.get(key).map_or(0, |entry| entry.in_flight)
}
#[must_use]
pub fn outstanding(&self) -> usize {
self.reservations.len()
}
#[must_use]
pub fn window(&self, key: &ConstraintKey) -> u32 {
self.state
.get(key)
.map_or(self.policy.min_window(), |entry| entry.window)
}
#[must_use]
pub fn rules(&self, key: &ConstraintKey) -> RulesState {
self.state
.get(key)
.map_or(RulesState::NotFetched, |entry| entry.rules.clone())
}
#[must_use]
pub fn is_abandoned(&self, host: &str) -> bool {
self.abandons.contains(host)
}
fn entry(&mut self, key: &ConstraintKey) -> &mut KeyState {
let floor = self.policy.min_window();
self.state.entry(key.clone()).or_insert_with(|| KeyState {
window: floor,
in_flight: 0,
successes: 0,
blocked_until: Duration::ZERO,
rules: RulesState::NotFetched,
})
}
pub fn observe_resolution(&mut self, host: &str, resolved: Resolved) {
match resolved {
Resolved::Origin(ref origin) => {
self.constraints
.insert(host.to_owned(), ConstraintKey::origin(origin.clone()));
}
Resolved::Failed | Resolved::DoesNotExist => {
self.constraints.remove(host);
}
}
self.resolutions.insert(host.to_owned(), resolved);
}
pub fn observe_rules(&mut self, key: &ConstraintKey, rules: RulesState) {
self.entry(key).rules = rules;
}
pub fn observe_retry_after(
&mut self,
permit: &InFlightPermit,
at: Duration,
delay: Duration,
) -> Result<(), CompletionError> {
self.verify(permit)?;
let until = at.saturating_add(delay);
for key in permit.server_keys() {
let entry = self.entry(&key);
if entry.blocked_until < until {
entry.blocked_until = until;
}
}
Ok(())
}
fn keys_for(&self, host: &str) -> Vec<ConstraintKey> {
let mut keys = Vec::with_capacity(3);
keys.push(ConstraintKey::host(host));
keys.extend(self.constraints.get(host).cloned());
keys.push(self.egress.clone());
keys
}
fn effective_rules(&self, key: &ConstraintKey, at: Duration) -> RulesState {
let grace = self.policy.unreachable_grace();
self.state.get(key).map_or(RulesState::NotFetched, |entry| {
entry.rules.effective(at, grace)
})
}
fn promote_rules(&mut self, key: &ConstraintKey, at: Duration) {
let grace = self.policy.unreachable_grace();
let entry = self.entry(key);
let promoted = entry.rules.effective(at, grace);
if promoted != entry.rules {
entry.rules = promoted;
}
}
fn decide(&self, host: &str, at: Duration) -> Verdict {
if self.abandons.contains(host) {
return Verdict::Reject {
kind: RejectKind::DoesNotExist,
};
}
match self.resolutions.get(host) {
None => {
return Verdict::Defer {
resume_at: at,
kind: DeferralKind::OriginUnmapped,
constraint: ConstraintKey::host(host),
};
}
Some(&Resolved::Failed) => {
return Verdict::Defer {
resume_at: at.saturating_add(self.policy.retry_after_failure()),
kind: DeferralKind::ResolutionFailed,
constraint: ConstraintKey::host(host),
};
}
Some(&Resolved::DoesNotExist) => {
return Verdict::Reject {
kind: RejectKind::DoesNotExist,
};
}
Some(&Resolved::Origin(_)) => {}
}
let host_key = ConstraintKey::host(host);
match self.effective_rules(&host_key, at) {
RulesState::Allowed | RulesState::NoRules => {}
RulesState::Disallowed { rule } => {
return Verdict::Reject {
kind: RejectKind::Disallowed { rule },
};
}
RulesState::NotFetched => {
return Verdict::Defer {
resume_at: at,
kind: DeferralKind::RulesNotFetched,
constraint: host_key,
};
}
RulesState::Unreachable { .. } => {
return Verdict::Defer {
resume_at: at.saturating_add(self.policy.retry_after_failure()),
kind: DeferralKind::OriginPolicyUnreachable,
constraint: host_key,
};
}
}
let mut binding: Option<Verdict> = None;
for key in self.keys_for(host) {
if let Some(wait) = self.saturated(&key, at) {
binding = Some(later(binding, wait));
}
}
match binding {
Some(verdict) => verdict,
None => Verdict::Admit,
}
}
fn saturated(&self, key: &ConstraintKey, at: Duration) -> Option<Verdict> {
let entry = self.state.get(key)?;
if entry.blocked_until > at {
Some(Verdict::Defer {
resume_at: entry.blocked_until,
kind: DeferralKind::RateLimited,
constraint: key.clone(),
})
} else if entry.in_flight >= entry.window {
Some(Verdict::Defer {
resume_at: at,
kind: DeferralKind::RateLimited,
constraint: key.clone(),
})
} else {
None
}
}
pub fn admit(&mut self, host: &str, at: Duration) -> Admission {
match self.decide(host, at).refusal() {
Some(refused) => refused,
None => self.reserve(host, at),
}
}
fn reserve(&mut self, host: &str, at: Duration) -> Admission {
let ordinal = self.next_ordinal;
let Some(next) = ordinal.checked_add(1) else {
return Admission::Reject {
kind: RejectKind::Exhausted,
};
};
let host_key = ConstraintKey::host(host);
self.promote_rules(&host_key, at);
let keys = self.keys_for(host);
for key in &keys {
let entry = self.entry(key);
entry.in_flight = entry.in_flight.saturating_add(1);
}
self.next_ordinal = next;
self.reservations.insert(ordinal);
Admission::Admit(InFlightPermit {
issuer: Arc::clone(&self.issuer),
ordinal,
host: host.to_owned(),
origin: self.constraints.get(host).cloned(),
egress: self.egress.clone(),
})
}
pub fn complete(
&mut self,
permit: InFlightPermit,
outcome: Outcome,
) -> Result<(), CompletionError> {
self.verify(&permit)?;
self.reservations.remove(&permit.ordinal);
for key in permit.keys() {
let entry = self.entry(&key);
entry.in_flight = entry.in_flight.saturating_sub(1);
}
let policy = self.policy;
for key in permit.server_keys() {
let entry = self.entry(&key);
match outcome {
Outcome::Success => {
entry.successes = entry.successes.saturating_add(1);
if entry.successes >= policy.increase_after() {
entry.successes = 0;
entry.window = policy.increase(entry.window);
}
}
Outcome::CacheHit => {}
Outcome::Failure => {
entry.successes = 0;
entry.window = policy.decrease(entry.window);
}
}
}
Ok(())
}
fn verify(&self, permit: &InFlightPermit) -> Result<(), CompletionError> {
if !Arc::ptr_eq(&permit.issuer, &self.issuer) {
let refusal = Err(CompletionError::ForeignPermit);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "verify: returning an error to the caller");
return refusal;
}
if !self.reservations.contains(&permit.ordinal) {
let refusal = Err(CompletionError::UnknownPermit);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "verify: returning an error to the caller");
return refusal;
}
Ok(())
}
pub fn abandon(&mut self, host: &str) {
self.abandons.insert(host.to_owned());
}
#[must_use]
pub fn next_admissible(&self, at: Duration) -> Option<String> {
self.resolutions
.keys()
.find(|host| matches!(self.decide(host, at), Verdict::Admit))
.cloned()
}
#[must_use]
pub fn next_prerequisite(&self, at: Duration) -> Option<(String, DeferralKind)> {
self.resolutions
.keys()
.find_map(|host| match self.decide(host, at) {
Verdict::Defer { kind, .. } if kind.needs_prerequisite() => {
Some((host.clone(), kind))
}
Verdict::Admit | Verdict::Defer { .. } | Verdict::Reject { .. } => None,
})
}
}
#[cfg(test)]
mod tests {
use super::*;
const HOUR: Duration = Duration::from_secs(3_600);
const THIRTY_ONE_DAYS: Duration = Duration::from_secs(2_678_400);
type TestResult = Result<(), Box<dyn std::error::Error>>;
fn frontier() -> Frontier {
Frontier::new(PolitenessPolicy::DEFAULT, "eth0")
}
fn ready(frontier: &mut Frontier, host: &str, origin: &str) {
frontier.observe_resolution(host, Resolved::Origin(origin.to_owned()));
frontier.observe_rules(&ConstraintKey::host(host), RulesState::Allowed);
}
fn deferral(admission: &Admission) -> Option<(Duration, DeferralKind, ConstraintKey)> {
match *admission {
Admission::Defer {
resume_at,
kind,
ref constraint,
} => Some((resume_at, kind, constraint.clone())),
Admission::Admit(_) | Admission::Reject { .. } => None,
}
}
fn admitted(
frontier: &mut Frontier,
host: &str,
at: Duration,
) -> Result<InFlightPermit, Box<dyn std::error::Error>> {
match frontier.admit(host, at) {
Admission::Admit(permit) => Ok(permit),
other => Err(format!("expected an admission for {host}, got {other:?}").into()),
}
}
fn completed(
frontier: &mut Frontier,
host: &str,
outcome: Outcome,
) -> Result<(), Box<dyn std::error::Error>> {
let permit = admitted(frontier, host, Duration::ZERO)?;
frontier.complete(permit, outcome)?;
Ok(())
}
#[test]
fn an_unresolved_host_defers_rather_than_rejects() {
let mut frontier = frontier();
assert_eq!(
deferral(&frontier.admit("a.example", Duration::ZERO)),
Some((
Duration::ZERO,
DeferralKind::OriginUnmapped,
ConstraintKey::host("a.example"),
)),
"a crawler discovers hosts as it goes, so an unknown host is a \
prerequisite to resolve and not a dead item"
);
}
#[test]
fn nxdomain_rejects_and_servfail_defers() {
let mut frontier = frontier();
frontier.observe_resolution("gone.example", Resolved::DoesNotExist);
frontier.observe_resolution("flaky.example", Resolved::Failed);
assert_eq!(
frontier.admit("gone.example", Duration::ZERO),
Admission::Reject {
kind: RejectKind::DoesNotExist
},
"RFC 8020 makes NXDOMAIN authoritative"
);
assert_eq!(
deferral(&frontier.admit("flaky.example", Duration::ZERO)),
Some((
PolitenessPolicy::DEFAULT.retry_after_failure(),
DeferralKind::ResolutionFailed,
ConstraintKey::host("flaky.example"),
)),
"SERVFAIL is not authoritative, so it must not be terminal"
);
}
#[test]
fn rules_states_follow_rfc_9309() -> TestResult {
let mut frontier = frontier();
ready(&mut frontier, "a.example", "origin-a");
let key = ConstraintKey::host("a.example");
frontier.observe_rules(&key, RulesState::NotFetched);
assert_eq!(
deferral(&frontier.admit("a.example", Duration::ZERO)).map(|parts| parts.1),
Some(DeferralKind::RulesNotFetched),
"unfetched rules are a prerequisite, not a wait"
);
frontier.observe_rules(&key, RulesState::NoRules);
let permit = admitted(&mut frontier, "a.example", HOUR)?;
frontier.complete(permit, Outcome::Success)?;
frontier.observe_rules(
&key,
RulesState::Unreachable {
since: Duration::ZERO,
},
);
assert_eq!(
deferral(&frontier.admit("a.example", HOUR)).map(|parts| parts.1),
Some(DeferralKind::OriginPolicyUnreachable),
"RFC 9309 2.3.1.4 requires assuming complete disallow"
);
assert!(
matches!(
frontier.admit("a.example", THIRTY_ONE_DAYS),
Admission::Admit(_)
),
"RFC 9309 2.3.1.4 lets a long-unreachable file be treated as unavailable"
);
assert_eq!(
frontier.rules(&key),
RulesState::NoRules,
"taking the admission is what records the expiry the verdict rested on"
);
Ok(())
}
#[test]
fn disallowed_is_terminal_and_unreachable_is_not() {
let mut frontier = frontier();
ready(&mut frontier, "a.example", "origin-a");
let key = ConstraintKey::host("a.example");
frontier.observe_rules(
&key,
RulesState::Disallowed {
rule: "Disallow: /private".to_owned(),
},
);
assert_eq!(
frontier.admit("a.example", Duration::ZERO),
Admission::Reject {
kind: RejectKind::Disallowed {
rule: "Disallow: /private".to_owned()
}
}
);
frontier.observe_rules(
&key,
RulesState::Unreachable {
since: Duration::ZERO,
},
);
assert_eq!(
deferral(&frontier.admit("a.example", HOUR)).map(|parts| parts.1),
Some(DeferralKind::OriginPolicyUnreachable),
"unreachability is provisional and must never be spelled Reject"
);
}
#[test]
fn hosts_sharing_an_origin_share_one_window() {
let mut frontier = frontier();
for index in 0..8_u32 {
ready(
&mut frontier,
&format!("shop{index}.example"),
"shopify-edge",
);
}
let mut admitted = 0_u32;
for index in 0..8_u32 {
if matches!(
frontier.admit(&format!("shop{index}.example"), Duration::ZERO),
Admission::Admit(_)
) {
admitted = admitted.saturating_add(1);
}
}
assert_eq!(
admitted,
PolitenessPolicy::DEFAULT.min_window(),
"eight 'independent' hosts behind one proxy must share one window, \
or the run is a Layer 7 flood with a polite-looking config"
);
assert_eq!(
frontier.in_flight(&ConstraintKey::origin("shopify-edge")),
PolitenessPolicy::DEFAULT.min_window()
);
}
#[test]
fn distinct_origins_are_limited_by_egress_alone() {
let mut frontier = frontier();
for index in 0..4_u32 {
let host = format!("h{index}.example");
ready(&mut frontier, &host, &format!("origin-{index}"));
}
let mut admitted = 0_u32;
for index in 0..4_u32 {
if matches!(
frontier.admit(&format!("h{index}.example"), Duration::ZERO),
Admission::Admit(_)
) {
admitted = admitted.saturating_add(1);
}
}
assert_eq!(
admitted,
PolitenessPolicy::DEFAULT.egress_window(),
"four hosts on four origins are bounded by the egress, not by each other"
);
assert_eq!(
deferral(&frontier.admit("h0.example", Duration::ZERO)).map(|parts| parts.2),
Some(ConstraintKey::egress("eth0")),
"and the binding constraint is the egress"
);
}
#[test]
fn the_binding_constraint_is_reported() {
let mut frontier = frontier();
ready(&mut frontier, "a.example", "origin-a");
assert!(
matches!(
frontier.admit("a.example", Duration::ZERO),
Admission::Admit(_)
),
"a resolved, permitted host on a free origin and egress admits"
);
assert_eq!(
deferral(&frontier.admit("a.example", Duration::ZERO)).map(|parts| parts.1),
Some(DeferralKind::RateLimited),
"a window of one admits exactly once"
);
}
#[test]
fn retry_after_binds_the_origin_not_just_the_host() -> TestResult {
let mut frontier = frontier();
ready(&mut frontier, "a.example", "origin-a");
ready(&mut frontier, "b.example", "origin-a");
let permit = admitted(&mut frontier, "a.example", Duration::ZERO)?;
frontier.observe_retry_after(&permit, Duration::ZERO, Duration::from_secs(30))?;
frontier.complete(permit, Outcome::Success)?;
assert_eq!(
deferral(&frontier.admit("b.example", Duration::from_secs(1))),
Some((
Duration::from_secs(30),
DeferralKind::RateLimited,
ConstraintKey::origin("origin-a"),
)),
"a pause asked of the client binds its sibling on the same origin"
);
assert!(
matches!(
frontier.admit("b.example", Duration::from_secs(31)),
Admission::Admit(_)
),
"and it lifts when the instant it named arrives"
);
Ok(())
}
#[test]
fn a_retry_after_does_not_bind_the_egress() -> TestResult {
let mut frontier = frontier();
ready(&mut frontier, "a.example", "origin-a");
ready(&mut frontier, "unrelated.example", "origin-z");
let permit = admitted(&mut frontier, "a.example", Duration::ZERO)?;
frontier.observe_retry_after(&permit, Duration::ZERO, Duration::from_secs(30))?;
assert!(
matches!(
frontier.admit("unrelated.example", Duration::from_secs(1)),
Admission::Admit(_)
),
"one server asking for a pause is not evidence about our network's \
capacity; binding the egress would let a single target halt every \
unrelated host in the run"
);
assert_eq!(
deferral(&frontier.admit("a.example", Duration::from_secs(1))).map(|parts| parts.2),
Some(ConstraintKey::origin("origin-a")),
"but it does bind the server-side keys the target actually shares"
);
Ok(())
}
#[test]
fn the_window_shrinks_on_failure_and_grows_on_success() -> TestResult {
let mut frontier = frontier();
ready(&mut frontier, "a.example", "origin-a");
let key = ConstraintKey::origin("origin-a");
for _ in 0..PolitenessPolicy::DEFAULT.increase_after() {
completed(&mut frontier, "a.example", Outcome::Success)?;
}
assert_eq!(
frontier.window(&key),
2,
"eight successes add one to a window of one"
);
completed(&mut frontier, "a.example", Outcome::Failure)?;
assert_eq!(frontier.window(&key), 1, "2 x 70% floors at min_window");
for _ in 0..5 {
completed(&mut frontier, "a.example", Outcome::Failure)?;
}
assert_eq!(
frontier.window(&key),
PolitenessPolicy::DEFAULT.min_window(),
"a window of zero would admit nothing forever"
);
Ok(())
}
#[test]
fn a_cache_hit_does_not_grow_the_window() -> TestResult {
let mut frontier = frontier();
ready(&mut frontier, "a.example", "origin-a");
let key = ConstraintKey::origin("origin-a");
for _ in 0..PolitenessPolicy::DEFAULT.increase_after() {
completed(&mut frontier, "a.example", Outcome::CacheHit)?;
}
assert_eq!(
frontier.window(&key),
PolitenessPolicy::DEFAULT.min_window(),
"an edge cache answering fast says nothing about origin capacity, and \
counting it is how a latency controller speeds up under load"
);
assert_eq!(
frontier.in_flight(&key),
0,
"a cache hit still returns the slots it held; only the growth signal \
is excluded"
);
Ok(())
}
#[test]
fn releasing_after_reresolution_preserves_other_requests() -> TestResult {
let mut frontier = frontier();
for (host, origin) in [
("a.example", "old"),
("b.example", "new"),
("c.example", "new"),
] {
frontier.observe_resolution(host, Resolved::Origin(origin.to_owned()));
frontier.observe_rules(&ConstraintKey::host(host), RulesState::Allowed);
}
let permit = admitted(&mut frontier, "a.example", Duration::ZERO)?;
assert!(
matches!(
frontier.admit("b.example", Duration::ZERO),
Admission::Admit(_)
),
"two hosts on two origins are bounded by the egress, not by each other"
);
frontier.observe_resolution("a.example", Resolved::Origin("new".to_owned()));
frontier.complete(permit, Outcome::Success)?;
assert_eq!(
frontier.in_flight(&ConstraintKey::origin("old")),
0,
"A's request was admitted against the old origin, so completing it \
must return that origin's slot — not one derived from the name's \
new topology"
);
assert_eq!(
frontier.in_flight(&ConstraintKey::origin("new")),
1,
"B is still physically running against the new origin and still owns \
its only slot"
);
assert!(
matches!(
frontier.admit("c.example", Duration::ZERO),
Admission::Defer { .. }
),
"the default origin window of one must never admit C while B owns it"
);
assert_eq!(
frontier.outstanding(),
1,
"one request is still in flight, so the ledger holds exactly one \
reservation"
);
Ok(())
}
#[test]
fn a_resolution_failure_during_flight_returns_the_original_slots() -> TestResult {
let mut frontier = frontier();
ready(&mut frontier, "a.example", "old");
let permit = admitted(&mut frontier, "a.example", Duration::ZERO)?;
assert_eq!(frontier.in_flight(&ConstraintKey::origin("old")), 1);
frontier.observe_resolution("a.example", Resolved::Failed);
frontier.complete(permit, Outcome::Failure)?;
assert_eq!(
frontier.in_flight(&ConstraintKey::origin("old")),
0,
"a resolution failure is an observation about the name, not evidence \
that the request stopped holding its origin's slot"
);
assert_eq!(
frontier.in_flight(&ConstraintKey::egress("eth0")),
0,
"and the egress slot comes back with it"
);
assert_eq!(frontier.outstanding(), 0, "the ledger is empty again");
Ok(())
}
#[test]
fn several_in_flight_across_topology_generations_conserve_every_slot() -> TestResult {
let mut frontier = frontier();
for (host, origin) in [
("a.example", "origin-a"),
("b.example", "origin-b"),
("c.example", "origin-c"),
] {
ready(&mut frontier, host, origin);
}
let first = admitted(&mut frontier, "a.example", Duration::ZERO)?;
let second = admitted(&mut frontier, "b.example", Duration::ZERO)?;
let third = admitted(&mut frontier, "c.example", Duration::ZERO)?;
assert_eq!(
frontier.in_flight(&ConstraintKey::egress("eth0")),
3,
"three requests are physically running"
);
frontier.observe_resolution("a.example", Resolved::Origin("gen2".to_owned()));
frontier.observe_resolution("a.example", Resolved::Origin("gen3".to_owned()));
assert_eq!(
frontier.in_flight(&ConstraintKey::origin("origin-a")),
1,
"the request holds the origin it was admitted against however many \
times its name has moved since"
);
frontier.complete(second, Outcome::Success)?;
frontier.complete(third, Outcome::Failure)?;
frontier.complete(first, Outcome::Success)?;
for origin in ["origin-a", "origin-b", "origin-c", "gen2", "gen3"] {
assert_eq!(
frontier.in_flight(&ConstraintKey::origin(origin)),
0,
"every origin's occupancy returns to zero once its requests have \
completed"
);
}
assert_eq!(frontier.in_flight(&ConstraintKey::egress("eth0")), 0);
assert_eq!(frontier.outstanding(), 0, "and no reservation is left over");
let moved = admitted(&mut frontier, "a.example", Duration::ZERO)?;
assert!(
moved.holds(&ConstraintKey::origin("gen3")),
"a new admission keys on the current topology"
);
assert!(
!moved.holds(&ConstraintKey::origin("origin-a")),
"and not on the one an earlier request for the same name used"
);
frontier.complete(moved, Outcome::Success)?;
Ok(())
}
#[test]
fn a_permit_holds_its_admission_time_keys() -> TestResult {
let mut frontier = frontier();
ready(&mut frontier, "a.example", "old");
let permit = admitted(&mut frontier, "a.example", Duration::ZERO)?;
assert!(
permit.holds(&ConstraintKey::host("a.example")),
"the permit holds the host it was admitted for"
);
assert!(
permit.holds(&ConstraintKey::origin("old")),
"and the origin the name resolved to at that instant"
);
assert!(
permit.holds(&ConstraintKey::egress("eth0")),
"and the egress every request in the run shares"
);
assert_eq!(permit.host(), "a.example");
frontier.observe_resolution("a.example", Resolved::Origin("new".to_owned()));
assert!(
!permit.holds(&ConstraintKey::origin("new")),
"a resolution observed while the request was in flight does not \
change what the reservation holds"
);
Ok(())
}
#[test]
fn response_attribution_follows_the_admission_not_the_current_name() -> TestResult {
let mut frontier = frontier();
ready(&mut frontier, "a.example", "old");
ready(&mut frontier, "b.example", "new");
let permit = admitted(&mut frontier, "a.example", Duration::ZERO)?;
frontier.observe_resolution("a.example", Resolved::Origin("new".to_owned()));
for _ in 0..PolitenessPolicy::DEFAULT.increase_after() {
completed(&mut frontier, "b.example", Outcome::Success)?;
}
let new_before = frontier.window(&ConstraintKey::origin("new"));
let old_before = frontier.window(&ConstraintKey::origin("old"));
frontier.complete(permit, Outcome::Failure)?;
assert_eq!(
frontier.window(&ConstraintKey::origin("old")),
PolitenessPolicy::DEFAULT.decrease(old_before),
"the failure shrank the origin the request was admitted against"
);
assert_eq!(
frontier.window(&ConstraintKey::origin("new")),
new_before,
"and left the origin the name moved to untouched"
);
Ok(())
}
#[test]
fn a_permit_from_another_frontier_is_refused() -> TestResult {
let mut issuer = frontier();
let mut stranger = frontier();
ready(&mut issuer, "a.example", "origin-a");
let permit = admitted(&mut issuer, "a.example", Duration::ZERO)?;
assert_eq!(
stranger.complete(permit, Outcome::Success),
Err(CompletionError::ForeignPermit),
"a permit is evidence only against the accounting that minted it, \
and a second frontier holds a different ledger"
);
assert_eq!(
issuer.outstanding(),
1,
"the refusal changed nothing: the issuer still holds the slot the \
request occupies"
);
assert_eq!(
issuer.in_flight(&ConstraintKey::origin("origin-a")),
1,
"and a foreign completion must not decrement it"
);
Ok(())
}
#[test]
fn a_permit_with_no_live_reservation_is_refused() -> TestResult {
let mut frontier = frontier();
ready(&mut frontier, "a.example", "origin-a");
let permit = admitted(&mut frontier, "a.example", Duration::ZERO)?;
let ordinal = permit.ordinal;
frontier.reservations.remove(&ordinal);
assert_eq!(
frontier.verify(&permit),
Err(CompletionError::UnknownPermit),
"a permit this ledger does not hold is a typed refusal, not a \
subtraction that quietly does nothing"
);
assert_eq!(
frontier.complete(permit, Outcome::Success),
Err(CompletionError::UnknownPermit),
"and completing it is refused before any window is touched"
);
Ok(())
}
#[test]
fn selection_skips_hosts_the_gate_will_not_admit() {
let mut frontier = frontier();
for host in ["a.example", "b.example"] {
frontier.observe_resolution(host, Resolved::Origin(host.to_owned()));
}
frontier.observe_rules(
&ConstraintKey::host("a.example"),
RulesState::Disallowed {
rule: "Disallow: /".to_owned(),
},
);
frontier.observe_rules(&ConstraintKey::host("b.example"), RulesState::Allowed);
assert!(
matches!(
frontier.admit("a.example", Duration::ZERO),
Admission::Reject { .. }
),
"the gate refuses the disallowed host"
);
assert_eq!(
frontier.next_admissible(Duration::ZERO).as_deref(),
Some("b.example"),
"so selection must not return it: rejection does not remove the \
host, and a selector with its own predicate returns it forever"
);
}
#[test]
fn selection_agrees_with_admission_on_every_rules_state() -> TestResult {
let mut frontier = frontier();
let hosts = ["a.example", "b.example", "c.example", "d.example"];
for host in hosts {
frontier.observe_resolution(host, Resolved::Origin(format!("origin-{host}")));
}
frontier.observe_rules(
&ConstraintKey::host("a.example"),
RulesState::Disallowed {
rule: "Disallow: /".to_owned(),
},
);
frontier.observe_rules(&ConstraintKey::host("b.example"), RulesState::NotFetched);
frontier.observe_rules(
&ConstraintKey::host("c.example"),
RulesState::Unreachable {
since: Duration::ZERO,
},
);
frontier.observe_rules(&ConstraintKey::host("d.example"), RulesState::Allowed);
let next = frontier.next_admissible(Duration::ZERO);
assert_eq!(
next.as_deref(),
Some("d.example"),
"disallowed, unfetched and unexpired-unreachable are all states the \
gate will not admit, so none of them is selectable"
);
let Some(selected) = next else {
return Err("the gate selected no host, so the round trip cannot be \
exercised"
.into());
};
let permit = match frontier.admit(&selected, Duration::ZERO) {
Admission::Admit(permit) => permit,
other => {
return Err(format!(
"selection returned {selected}, which the gate answered with {other:?}"
)
.into());
}
};
assert!(permit.holds(&ConstraintKey::host("d.example")));
Ok(())
}
#[test]
fn a_lone_unpermitted_host_is_not_selected() {
let mut frontier = frontier();
frontier.observe_resolution("only.example", Resolved::Origin("origin".to_owned()));
frontier.observe_rules(&ConstraintKey::host("only.example"), RulesState::NotFetched);
assert_eq!(
frontier.next_admissible(Duration::ZERO),
None,
"a host whose rules have not been fetched is a prerequisite, not an \
admission, and returning it as admissible deadlocks the scheduler \
that trusts the answer"
);
assert_eq!(
frontier.next_prerequisite(Duration::ZERO),
Some(("only.example".to_owned(), DeferralKind::RulesNotFetched)),
"the same predicate reports the work that would unblock it"
);
}
#[test]
fn selection_agrees_with_admission_across_cooldowns_and_saturation() -> TestResult {
let mut frontier = frontier();
ready(&mut frontier, "a.example", "origin-a");
ready(&mut frontier, "b.example", "origin-b");
ready(&mut frontier, "c.example", "origin-c");
let permit = admitted(&mut frontier, "a.example", Duration::ZERO)?;
frontier.observe_retry_after(&permit, Duration::ZERO, Duration::from_secs(30))?;
let held = admitted(&mut frontier, "b.example", Duration::ZERO)?;
assert_eq!(
frontier.next_admissible(Duration::from_secs(1)).as_deref(),
Some("c.example"),
"a host on cooldown and a host at its window are both held, so \
selection moves past them to one the gate will take"
);
assert!(
matches!(
frontier.admit("c.example", Duration::from_secs(1)),
Admission::Admit(_)
),
"and the host selection returned is one the gate admits"
);
frontier.complete(permit, Outcome::Success)?;
assert_eq!(
frontier.next_admissible(Duration::from_secs(1)),
None,
"with the egress saturated, nothing is admissible"
);
assert_eq!(
frontier.next_prerequisite(Duration::from_secs(1)),
None,
"and nothing is waiting on a prerequisite either: these are waits, \
not work the run can do"
);
frontier.complete(held, Outcome::Success)?;
assert_eq!(
frontier.next_admissible(Duration::from_secs(31)).as_deref(),
Some("a.example"),
"once the egress frees and the cooldown expires, the backlog is \
selectable again — no caller-maintained shadow filtering"
);
Ok(())
}
#[test]
fn a_grace_expiry_is_visible_to_selection_before_anything_acts_on_it() {
let mut frontier = frontier();
frontier.observe_resolution("a.example", Resolved::Origin("origin-a".to_owned()));
frontier.observe_rules(
&ConstraintKey::host("a.example"),
RulesState::Unreachable {
since: Duration::ZERO,
},
);
let inside_grace = PolitenessPolicy::DEFAULT
.unreachable_grace()
.saturating_sub(Duration::from_secs(1));
assert_eq!(
frontier.next_admissible(inside_grace),
None,
"inside the grace the complete disallow still holds"
);
assert_eq!(
frontier
.next_admissible(PolitenessPolicy::DEFAULT.unreachable_grace())
.as_deref(),
Some("a.example"),
"and at the boundary it lifts for selection, exactly where the gate \
would stop deferring"
);
assert_eq!(
frontier.rules(&ConstraintKey::host("a.example")),
RulesState::Unreachable {
since: Duration::ZERO
},
"selection is side-effect free: it reports the expiry without \
writing it down"
);
assert!(
matches!(
frontier.admit("a.example", PolitenessPolicy::DEFAULT.unreachable_grace()),
Admission::Admit(_)
),
"and the two agree once the grant is acted on"
);
}
#[test]
fn selection_makes_progress_through_a_disallowed_prefix() -> TestResult {
let mut frontier = frontier();
for host in ["a.example", "b.example", "c.example", "d.example"] {
frontier.observe_resolution(host, Resolved::Origin(format!("origin-{host}")));
frontier.observe_rules(
&ConstraintKey::host(host),
if host == "a.example" {
RulesState::Disallowed {
rule: "Disallow: /".to_owned(),
}
} else {
RulesState::Allowed
},
);
}
let mut dispatched = Vec::new();
let mut held = Vec::new();
for _ in 0..3 {
let Some(host) = frontier.next_admissible(Duration::ZERO) else {
break;
};
match frontier.admit(&host, Duration::ZERO) {
Admission::Admit(permit) => {
dispatched.push(host);
held.push(permit);
}
other => {
return Err(format!(
"selection returned {host}, which the gate answered with {other:?}"
)
.into());
}
}
}
assert_eq!(
dispatched,
vec!["b.example", "c.example", "d.example"],
"three cycles dispatch three distinct ready hosts behind a host that \
can never be dispatched"
);
assert_eq!(held.len(), 3, "each dispatch is still in flight");
Ok(())
}
#[test]
fn ten_thousand_requests_replay_identically() {
#[derive(Debug, PartialEq, Eq)]
enum Shape {
Admit,
Defer(Duration, DeferralKind, ConstraintKey),
Reject(RejectKind),
}
struct Replay {
verdicts: Vec<Shape>,
windows: Vec<u32>,
refusals: usize,
}
fn run(mut frontier: Frontier) -> Replay {
use std::collections::VecDeque;
let hosts: Vec<String> = (0..20_u32)
.map(|index| format!("h{index}.example"))
.collect();
let mut replay = Replay {
verdicts: Vec::new(),
windows: Vec::new(),
refusals: 0,
};
let mut parked: VecDeque<InFlightPermit> = VecDeque::new();
let mut at = Duration::ZERO;
let mut index = 0_usize;
for step in 0..10_000_usize {
let host = &hosts[index];
let outcome = if step.checked_rem(7) == Some(0) {
Outcome::Failure
} else if step.checked_rem(3) == Some(0) {
Outcome::CacheHit
} else {
Outcome::Success
};
if parked.len() >= 3 {
let oldest = parked.pop_front();
let refused =
oldest.is_some_and(|permit| frontier.complete(permit, outcome).is_err());
if refused {
replay.refusals = replay.refusals.saturating_add(1);
}
}
let shape = match frontier.admit(host, at) {
Admission::Admit(permit) => {
parked.push_back(permit);
Shape::Admit
}
Admission::Defer {
resume_at,
kind,
constraint,
} => Shape::Defer(resume_at, kind, constraint),
Admission::Reject { kind } => Shape::Reject(kind),
};
replay.verdicts.push(shape);
at = at.saturating_add(Duration::from_millis(5));
if step.checked_rem(500) == Some(0) {
replay
.windows
.push(frontier.window(&ConstraintKey::egress("eth0")));
}
index = if index.saturating_add(1) >= hosts.len() {
0
} else {
index.saturating_add(1)
};
}
replay
}
fn build() -> Frontier {
let mut frontier = frontier();
const ORIGINS: [&str; 3] = ["origin-0", "origin-1", "origin-2"];
for (index, origin) in (0..20_u32).zip(ORIGINS.iter().cycle()) {
let host = format!("h{index}.example");
ready(&mut frontier, &host, origin);
}
frontier
}
let first = run(build());
let second = run(build());
assert_eq!(first.verdicts.len(), 10_000);
assert_eq!(
first.refusals, 0,
"every admission in the run was released by the ledger that minted \
it, so the refusal counter is the run's own consistency check"
);
assert!(
first
.verdicts
.iter()
.any(|shape| matches!(*shape, Shape::Admit))
&& first
.verdicts
.iter()
.any(|shape| matches!(*shape, Shape::Defer(..))),
"the schedule exercised both dispatch and back-off, so the replay \
comparison is not two identical silences"
);
assert_eq!(
first.verdicts, second.verdicts,
"the same schedule on the same policy must produce identical verdicts"
);
assert_eq!(first.windows, second.windows);
}
#[test]
fn permits_from_identically_driven_frontiers_compare_equal() -> TestResult {
let mut first = frontier();
let mut second = frontier();
ready(&mut first, "a.example", "origin-a");
ready(&mut second, "a.example", "origin-a");
assert_eq!(
admitted(&mut first, "a.example", Duration::ZERO)?,
admitted(&mut second, "a.example", Duration::ZERO)?,
"identical schedules produce identical reservations"
);
Ok(())
}
#[test]
fn next_admissible_is_stable_and_skips_the_abandoned() {
let mut frontier = frontier();
ready(&mut frontier, "b.example", "origin-b");
ready(&mut frontier, "a.example", "origin-a");
ready(&mut frontier, "c.example", "origin-c");
assert_eq!(
frontier.next_admissible(Duration::ZERO).as_deref(),
Some("a.example"),
"selection is by key order, not by insertion order"
);
frontier.abandon("a.example");
assert_eq!(
frontier.next_admissible(Duration::ZERO).as_deref(),
Some("b.example")
);
assert!(frontier.is_abandoned("a.example"));
}
#[test]
fn releasing_returns_every_shared_slot() -> TestResult {
let mut frontier = frontier();
ready(&mut frontier, "a.example", "origin-a");
ready(&mut frontier, "b.example", "origin-a");
let permit = admitted(&mut frontier, "a.example", Duration::ZERO)?;
assert_eq!(
deferral(&frontier.admit("b.example", Duration::ZERO)).map(|parts| parts.1),
Some(DeferralKind::RateLimited)
);
frontier.complete(permit, Outcome::Success)?;
assert_eq!(frontier.in_flight(&ConstraintKey::origin("origin-a")), 0);
assert_eq!(frontier.in_flight(&ConstraintKey::egress("eth0")), 0);
assert_eq!(frontier.outstanding(), 0);
Ok(())
}
}