use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use atomic_waker::AtomicWaker;
use jiff::SignedDuration;
use tokio::sync::watch;
use tracing::Instrument as _;
use tollgate_admission::LeaseSlot;
use tollgate_core::{AccountId, CostUnits, LeaseGrant, LocalLease, RefillSignal, RefillVerdict};
use tollgate_store::{AllocateError, LeaseAllocator};
use tollgate_store::Clock;
#[derive(Debug, Clone, Copy)]
pub struct LeaseManagerConfig {
pub account: AccountId,
pub target_grant: CostUnits,
pub low_water: CostUnits,
pub lease_ttl: SignedDuration,
pub expiry_safety_margin: SignedDuration,
pub poll_interval: std::time::Duration,
pub store_call_timeout: std::time::Duration,
pub shutdown_release_deadline: std::time::Duration,
}
#[derive(Debug, Clone, Copy)]
pub struct AccountLeaseConfig {
pub target_grant: CostUnits,
pub low_water: CostUnits,
pub lease_ttl: SignedDuration,
pub expiry_safety_margin: SignedDuration,
pub poll_interval: std::time::Duration,
pub store_call_timeout: std::time::Duration,
pub shutdown_release_deadline: std::time::Duration,
}
impl AccountLeaseConfig {
#[must_use]
pub fn for_account(self, account: AccountId) -> LeaseManagerConfig {
LeaseManagerConfig {
account,
target_grant: self.target_grant,
low_water: self.low_water,
lease_ttl: self.lease_ttl,
expiry_safety_margin: self.expiry_safety_margin,
poll_interval: self.poll_interval,
store_call_timeout: self.store_call_timeout,
shutdown_release_deadline: self.shutdown_release_deadline,
}
}
pub fn validate(&self) -> Result<(), LeaseManagerConfigError> {
self.for_account(AccountId(0)).validate()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct LeaseManagerReport {
pub released: u64,
pub abandoned: u64,
pub task_died: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct LeaseManagerConfigError(pub &'static str);
impl std::fmt::Display for LeaseManagerConfigError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.0)
}
}
impl std::error::Error for LeaseManagerConfigError {}
impl LeaseManagerConfig {
pub fn validate(&self) -> Result<(), LeaseManagerConfigError> {
if self.target_grant.is_zero() {
return Err(LeaseManagerConfigError("target_grant must be positive"));
}
if self.low_water >= self.target_grant {
return Err(LeaseManagerConfigError(
"low_water must be below target_grant",
));
}
if self.lease_ttl <= SignedDuration::ZERO {
return Err(LeaseManagerConfigError("lease_ttl must be positive"));
}
if self.expiry_safety_margin < SignedDuration::ZERO {
return Err(LeaseManagerConfigError(
"expiry_safety_margin must not be negative",
));
}
if self.expiry_safety_margin >= self.lease_ttl {
return Err(LeaseManagerConfigError(
"expiry_safety_margin must be shorter than lease_ttl",
));
}
if self.poll_interval.is_zero() {
return Err(LeaseManagerConfigError("poll_interval must be positive"));
}
if self.store_call_timeout.is_zero() {
return Err(LeaseManagerConfigError(
"store_call_timeout must be positive",
));
}
if self.shutdown_release_deadline.is_zero() {
return Err(LeaseManagerConfigError(
"shutdown_release_deadline must be positive",
));
}
Ok(())
}
}
#[derive(Debug, Default)]
struct RefillRequests {
waker: AtomicWaker,
pending: AtomicBool,
}
impl RefillSignal for RefillRequests {
fn request_refill(&self) {
self.pending.store(true, Ordering::Release);
self.waker.wake();
}
}
impl RefillRequests {
async fn requested(&self) {
std::future::poll_fn(|cx| {
if self.pending.swap(false, Ordering::Acquire) {
return std::task::Poll::Ready(());
}
self.waker.register(cx.waker());
if self.pending.swap(false, Ordering::Acquire) {
std::task::Poll::Ready(())
} else {
std::task::Poll::Pending
}
})
.await;
}
}
#[derive(Debug)]
pub struct LeaseCounters {
acquired: AtomicU64,
acquired_units: AtomicU64,
acquire_timeouts: AtomicU64,
uncertain_acquires: AtomicU64,
acquire_pending: AtomicBool,
integrity_fault: AtomicBool,
acquire_refused: [AtomicU64; AllocateError::COUNT],
released: AtomicU64,
abandoned: AtomicU64,
consolidated: AtomicU64,
consolidations_deferred: AtomicU64,
}
impl LeaseCounters {
#[must_use]
pub const fn new() -> Self {
LeaseCounters {
acquired: AtomicU64::new(0),
acquired_units: AtomicU64::new(0),
acquire_timeouts: AtomicU64::new(0),
uncertain_acquires: AtomicU64::new(0),
acquire_pending: AtomicBool::new(false),
integrity_fault: AtomicBool::new(false),
acquire_refused: [const { AtomicU64::new(0) }; AllocateError::COUNT],
released: AtomicU64::new(0),
abandoned: AtomicU64::new(0),
consolidated: AtomicU64::new(0),
consolidations_deferred: AtomicU64::new(0),
}
}
pub(crate) fn acquire_pending(&self) -> bool {
self.acquire_pending.load(Ordering::Acquire)
}
pub(crate) fn integrity_fault(&self) -> bool {
self.integrity_fault.load(Ordering::Acquire)
}
fn record_integrity_fault(&self, health: &watch::Sender<bool>) {
self.integrity_fault.store(true, Ordering::Release);
crate::signal(health, false, "lease-manager health");
}
fn record_acquired(&self, units: CostUnits) {
self.acquired.fetch_add(1, Ordering::Relaxed);
self.acquired_units
.fetch_add(units.get(), Ordering::Relaxed);
}
fn record_acquire_timeout(&self) {
self.acquire_timeouts.fetch_add(1, Ordering::Relaxed);
self.uncertain_acquires.fetch_add(1, Ordering::Relaxed);
}
fn record_consolidated(&self, units: CostUnits) {
self.record_acquired(units);
self.record_released();
self.consolidated.fetch_add(1, Ordering::Relaxed);
}
fn record_consolidation_deferred(&self) {
self.consolidations_deferred.fetch_add(1, Ordering::Relaxed);
}
fn record_acquire_refused(&self, error: &AllocateError) {
self.acquire_refused[error.index()].fetch_add(1, Ordering::Relaxed);
if matches!(error, AllocateError::Storage(_)) {
self.uncertain_acquires.fetch_add(1, Ordering::Relaxed);
}
}
fn record_released(&self) {
self.released.fetch_add(1, Ordering::Relaxed);
}
fn record_abandoned(&self) {
self.abandoned.fetch_add(1, Ordering::Relaxed);
}
#[must_use]
pub fn snapshot(&self) -> LeaseStats {
LeaseStats {
acquired: self.acquired.load(Ordering::Relaxed),
acquired_units: self.acquired_units.load(Ordering::Relaxed),
acquire_timeouts: self.acquire_timeouts.load(Ordering::Relaxed),
uncertain_acquires: self.uncertain_acquires.load(Ordering::Relaxed),
acquire_refused: std::array::from_fn(|slot| {
self.acquire_refused[slot].load(Ordering::Relaxed)
}),
released: self.released.load(Ordering::Relaxed),
abandoned: self.abandoned.load(Ordering::Relaxed),
consolidated: self.consolidated.load(Ordering::Relaxed),
consolidations_deferred: self.consolidations_deferred.load(Ordering::Relaxed),
}
}
}
impl Default for LeaseCounters {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct LeaseStats {
pub acquired: u64,
pub acquired_units: u64,
pub acquire_timeouts: u64,
pub uncertain_acquires: u64,
pub acquire_refused: [u64; AllocateError::COUNT],
pub released: u64,
pub abandoned: u64,
pub consolidated: u64,
pub consolidations_deferred: u64,
}
impl LeaseStats {
pub const ZERO: Self = Self {
acquired: 0,
acquired_units: 0,
acquire_timeouts: 0,
uncertain_acquires: 0,
acquire_refused: [0; AllocateError::COUNT],
released: 0,
abandoned: 0,
consolidated: 0,
consolidations_deferred: 0,
};
pub fn checked_add(self, other: Self) -> Option<Self> {
let mut acquire_refused = [0; AllocateError::COUNT];
for (i, value) in acquire_refused.iter_mut().enumerate() {
*value = self.acquire_refused[i].checked_add(other.acquire_refused[i])?;
}
acquire_refused
.iter()
.try_fold(0_u64, |sum, value| sum.checked_add(*value))?;
Some(Self {
acquired: self.acquired.checked_add(other.acquired)?,
acquired_units: self.acquired_units.checked_add(other.acquired_units)?,
acquire_timeouts: self.acquire_timeouts.checked_add(other.acquire_timeouts)?,
uncertain_acquires: self
.uncertain_acquires
.checked_add(other.uncertain_acquires)?,
acquire_refused,
released: self.released.checked_add(other.released)?,
abandoned: self.abandoned.checked_add(other.abandoned)?,
consolidated: self.consolidated.checked_add(other.consolidated)?,
consolidations_deferred: self
.consolidations_deferred
.checked_add(other.consolidations_deferred)?,
})
}
pub fn refusals_by_name(&self) -> impl Iterator<Item = (&'static str, u64)> + '_ {
AllocateError::NAMES
.iter()
.copied()
.zip(self.acquire_refused.iter().copied())
}
#[must_use]
pub fn refused(&self) -> u64 {
self.acquire_refused.iter().sum()
}
}
pub struct LeaseManager {
shutdown: watch::Sender<bool>,
health: watch::Receiver<bool>,
handle: Option<tokio::task::JoinHandle<LeaseManagerReport>>,
counters: Arc<LeaseCounters>,
deadline: Arc<crate::ShutdownDeadline>,
paused: Arc<AtomicBool>,
}
impl LeaseManager {
pub fn spawn(
allocator: Arc<dyn LeaseAllocator>,
slot: Arc<LeaseSlot>,
clock: Arc<dyn Clock>,
config: LeaseManagerConfig,
) -> Result<Self, LeaseManagerConfigError> {
config.validate()?;
let (shutdown, shutdown_rx) = watch::channel(false);
let (health_tx, health) = crate::task_health::TaskHealth::channel(true);
let account = config.account;
let counters = Arc::new(LeaseCounters::new());
let task_counters = Arc::clone(&counters);
let deadline = Arc::new(crate::ShutdownDeadline::default());
let task_deadline = Arc::clone(&deadline);
let paused = Arc::new(AtomicBool::new(false));
let task_paused = Arc::clone(&paused);
let handle = tokio::spawn(
async move {
run(
allocator,
slot,
clock,
config,
shutdown_rx,
health_tx.sender(),
&task_counters,
&task_deadline,
&task_paused,
)
.await
}
.instrument(tracing::info_span!("lease_manager", %account)),
);
Ok(LeaseManager {
shutdown,
health,
handle: Some(handle),
counters,
deadline,
paused,
})
}
#[must_use]
pub fn counters(&self) -> Arc<LeaseCounters> {
Arc::clone(&self.counters)
}
#[must_use]
pub fn health(&self) -> watch::Receiver<bool> {
self.health.clone()
}
pub(crate) fn pause_refills(&self) {
self.paused.store(true, Ordering::Release);
}
pub(crate) fn stop_at(&self, deadline: tokio::time::Instant) {
self.deadline.constrain(deadline);
crate::signal(&self.shutdown, true, "lease-manager shutdown");
}
pub async fn shutdown(mut self) -> LeaseManagerReport {
crate::signal(&self.shutdown, true, "lease-manager shutdown");
let died = LeaseManagerReport {
released: 0,
abandoned: 0,
task_died: true,
};
match self.handle.as_mut() {
Some(handle) => handle.await.unwrap_or(died),
None => died,
}
}
}
impl Drop for LeaseManager {
fn drop(&mut self) {
if let Some(handle) = self.handle.take() {
handle.abort();
}
}
}
#[allow(
clippy::too_many_arguments,
reason = "task entry point: every argument is a handle the loop owns for its lifetime, assembled once by spawn"
)]
async fn run(
allocator: Arc<dyn LeaseAllocator>,
slot: Arc<LeaseSlot>,
clock: Arc<dyn Clock>,
config: LeaseManagerConfig,
mut shutdown: watch::Receiver<bool>,
health: &watch::Sender<bool>,
counters: &LeaseCounters,
shutdown_deadline: &crate::ShutdownDeadline,
paused: &AtomicBool,
) -> LeaseManagerReport {
let mut tick = tokio::time::interval(config.poll_interval);
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
let refill = Arc::new(RefillRequests::default());
let mut parked: Vec<Arc<LocalLease>> = Vec::new();
loop {
tokio::select! {
_ = tick.tick() => {}
() = refill.requested() => {}
changed = shutdown.changed() => {
if changed.is_err() {
break;
}
}
}
if *shutdown.borrow() {
break;
}
if paused.load(Ordering::Acquire) {
continue;
}
let now = clock.now();
let pass_deadline = tokio::time::Instant::now() + config.store_call_timeout;
if release_quiesced(
&allocator,
&mut parked,
&clock,
&config,
&slot,
health,
counters,
&mut shutdown,
pass_deadline,
)
.await
== ReleasePass::ShutdownObserved
{
break;
}
let rotation = match slot.load_observed() {
None => Rotation::Acquire,
Some(lease) if now >= lease.usable_until() => {
drop(lease);
if let Some(old) = slot.take() {
parked.push(old);
}
if release_quiesced(
&allocator,
&mut parked,
&clock,
&config,
&slot,
health,
counters,
&mut shutdown,
pass_deadline,
)
.await
== ReleasePass::ShutdownObserved
{
break;
}
Rotation::Acquire
}
Some(lease) => match lease.refill_due_or_rearm() {
RefillVerdict::Idle => Rotation::Idle,
RefillVerdict::Draining => Rotation::Acquire,
RefillVerdict::Refused => Rotation::Consolidate,
},
};
if rotation == Rotation::Idle {
continue;
}
if paused.load(Ordering::Acquire) || *shutdown.borrow() {
continue;
}
if rotation == Rotation::Consolidate {
match consolidate_live_lease(
&allocator,
&mut parked,
&clock,
&config,
&slot,
&refill,
health,
counters,
&mut shutdown,
)
.await
{
Consolidation::ShutdownObserved => break,
Consolidation::Installed | Consolidation::KeptServing => continue,
Consolidation::AcquireInstead => {}
}
}
let funding_attempt = slot.funding_attempt();
counters.acquire_pending.store(true, Ordering::Release);
let acquire = tokio::time::timeout(
config.store_call_timeout,
allocator.acquire(config.account, config.target_grant, config.lease_ttl, now),
);
let acquired = tokio::select! {
outcome = acquire => outcome,
_ = shutdown.changed() => {
tracing::debug!("shutdown observed during an acquire; abandoning the refill");
break;
}
};
counters.acquire_pending.store(false, Ordering::Release);
match acquired {
Ok(Ok(allocation)) => {
counters.record_acquired(allocation.grant.units);
let fresh = install_lease(allocation.grant, &config, &slot, &refill);
park(
funding_attempt.granted(fresh, allocation.funding),
&mut parked,
);
}
outcome => {
let serving = slot.load_observed().is_some();
let reason: &dyn std::fmt::Display = match &outcome {
Ok(Err(error)) => {
if let Some(evidence) = refusal_evidence(error) {
funding_attempt.shortfall(evidence);
}
counters.record_acquire_refused(error);
error
}
_ => {
counters.record_acquire_timeout();
&"allocator timed out"
}
};
if serving {
tracing::debug!(%reason, "lease acquire refused; still serving");
} else {
tracing::warn!(
%reason,
"lease acquire refused with an empty slot; requests are denied"
);
}
}
}
}
if let Some(lease) = slot.take() {
parked.push(lease);
}
let deadline = shutdown_deadline.within(config.shutdown_release_deadline);
let mut report = LeaseManagerReport {
released: 0,
abandoned: 0,
task_died: false,
};
for lease in parked {
let grant = lease.grant();
if !wait_for_quiescence(&lease, deadline).await {
report.abandoned += 1;
counters.record_abandoned();
tracing::warn!(
lease = %grant.lease_id,
units = lease.remaining().get(),
"lease still held by an in-flight request at the shutdown \
deadline; abandoned rather than released, because releasing \
units that may still be spent cannot be undone"
);
continue;
}
let call_deadline = deadline.min(tokio::time::Instant::now() + config.store_call_timeout);
match tokio::time::timeout_at(
call_deadline,
allocator.release(
grant.lease_id,
grant.fencing_token,
lease.remaining(),
clock.now(),
),
)
.await
{
Ok(Ok(())) => {
report.released += 1;
counters.record_released();
}
Ok(Err(error)) => match release_failure(&error) {
ReleaseFailure::Settled | ReleaseFailure::Fenced => {
report.released += 1;
counters.record_released();
tracing::warn!(lease = %grant.lease_id, %error,
"lease no longer belongs to this manager; considered settled");
}
ReleaseFailure::Retryable => {
report.abandoned += 1;
counters.record_abandoned();
tracing::warn!(lease = %grant.lease_id, %error,
"release unconfirmed at shutdown; grant requires TTL reclaim");
}
ReleaseFailure::Integrity => {
report.abandoned += 1;
counters.record_abandoned();
counters.record_integrity_fault(health);
tracing::error!(lease = %grant.lease_id, %error,
"release violated the allocator contract; grant abandoned and integrity fault retained");
}
},
Err(_) => {
report.abandoned += 1;
counters.record_abandoned();
tracing::warn!(
lease = %grant.lease_id,
units = lease.remaining().get(),
"shutdown budget expired with a lease unreturned; \
its units settle at TTL reclaim"
);
}
}
}
report
}
#[derive(Clone, Copy)]
enum ReleaseFailure {
Settled,
Fenced,
Retryable,
Integrity,
}
fn release_failure(error: &AllocateError) -> ReleaseFailure {
match error {
AllocateError::UnknownLease | AllocateError::LeaseNotActive => ReleaseFailure::Settled,
AllocateError::Fenced => ReleaseFailure::Fenced,
AllocateError::Storage(_) => ReleaseFailure::Retryable,
AllocateError::InvalidRelease
| AllocateError::UnknownAccount
| AllocateError::AccountInactive
| AllocateError::InsufficientBalance
| AllocateError::BalanceExhausted(_)
| AllocateError::BalanceInsufficient(_)
| AllocateError::BalanceOverflow
| AllocateError::InvalidTtl => ReleaseFailure::Integrity,
}
}
fn install_lease(
grant: LeaseGrant,
config: &LeaseManagerConfig,
slot: &Arc<LeaseSlot>,
refill: &Arc<RefillRequests>,
) -> Arc<LocalLease> {
let low_water = CostUnits(
config
.low_water
.get()
.min(grant.units.get().saturating_sub(1)),
);
Arc::new(
LocalLease::with_sharding(
grant,
low_water,
config.expiry_safety_margin,
slot.sharding(),
)
.with_refill(Arc::clone(refill) as Arc<dyn RefillSignal>),
)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Rotation {
Idle,
Acquire,
Consolidate,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ConsolidationFailure {
RolledBack,
Settled,
Ambiguous,
Integrity,
}
fn consolidation_failure(error: &AllocateError) -> ConsolidationFailure {
match error {
AllocateError::UnknownLease | AllocateError::LeaseNotActive | AllocateError::Fenced => {
ConsolidationFailure::Settled
}
AllocateError::InsufficientBalance
| AllocateError::BalanceExhausted(_)
| AllocateError::BalanceInsufficient(_)
| AllocateError::UnknownAccount
| AllocateError::AccountInactive
| AllocateError::InvalidTtl => ConsolidationFailure::RolledBack,
AllocateError::InvalidRelease | AllocateError::BalanceOverflow => {
ConsolidationFailure::Integrity
}
AllocateError::Storage(_) => ConsolidationFailure::Ambiguous,
}
}
#[allow(
clippy::too_many_arguments,
reason = "one consolidation step over the loop's own borrowed state"
)]
async fn consolidate_live_lease(
allocator: &Arc<dyn LeaseAllocator>,
parked: &mut Vec<Arc<LocalLease>>,
clock: &Arc<dyn Clock>,
config: &LeaseManagerConfig,
slot: &Arc<LeaseSlot>,
refill: &Arc<RefillRequests>,
health: &watch::Sender<bool>,
counters: &LeaseCounters,
shutdown: &mut watch::Receiver<bool>,
) -> Consolidation {
let Some(live) = slot.take() else {
return Consolidation::AcquireInstead;
};
if Arc::strong_count(&live) > 1 || !live.is_only_local_view() {
counters.record_consolidation_deferred();
publish_and_park(slot, live, parked);
return Consolidation::KeptServing;
}
let unspent = live.remaining();
let needed = live.largest_refused_quote();
let grant = *live.grant();
let funding_attempt = slot.funding_attempt();
counters.acquire_pending.store(true, Ordering::Release);
let call = tokio::time::timeout(
config.store_call_timeout,
allocator.consolidate(
grant.lease_id,
grant.fencing_token,
unspent,
config.target_grant,
needed,
config.lease_ttl,
clock.now(),
),
);
let outcome = tokio::select! {
outcome = call => outcome,
_ = shutdown.changed() => {
tracing::debug!(
lease = %grant.lease_id,
"shutdown observed during a consolidation; parking the grant"
);
parked.push(live);
return Consolidation::ShutdownObserved;
}
};
counters.acquire_pending.store(false, Ordering::Release);
let error: AllocateError = match outcome {
Ok(Ok(allocation)) => {
let fresh = allocation.grant;
counters.record_consolidated(fresh.units);
tracing::debug!(
superseded = %grant.lease_id,
lease = %fresh.lease_id,
folded = unspent.get(),
units = fresh.units.get(),
"consolidated a refused lease into a larger grant"
);
drop(live);
let fresh = install_lease(fresh, config, slot, refill);
park(funding_attempt.granted(fresh, allocation.funding), parked);
return Consolidation::Installed;
}
Ok(Err(error)) => error,
Err(_) => {
counters.record_acquire_timeout();
tracing::warn!(
lease = %grant.lease_id,
"consolidation timed out; parking the grant because the store may have settled it"
);
parked.push(live);
return Consolidation::KeptServing;
}
};
if let Some(evidence) = refusal_evidence(&error) {
funding_attempt.shortfall(evidence);
}
counters.record_acquire_refused(&error);
match consolidation_failure(&error) {
ConsolidationFailure::RolledBack => {
tracing::debug!(
lease = %grant.lease_id, %error,
"consolidation refused; the grant is untouched and keeps serving"
);
publish_and_park(slot, live, parked);
Consolidation::KeptServing
}
ConsolidationFailure::Integrity => {
counters.record_integrity_fault(health);
tracing::error!(
lease = %grant.lease_id, units = unspent.get(), %error,
"consolidation violated the allocator contract; readiness withdrawn"
);
publish_and_park(slot, live, parked);
Consolidation::KeptServing
}
ConsolidationFailure::Settled => {
counters.record_released();
tracing::warn!(
lease = %grant.lease_id, %error,
"the store no longer holds this grant open; acquiring a fresh one"
);
drop(live);
Consolidation::AcquireInstead
}
ConsolidationFailure::Ambiguous => {
tracing::warn!(
lease = %grant.lease_id, %error,
"consolidation failed without saying whether it committed; parking the grant"
);
parked.push(live);
Consolidation::KeptServing
}
}
}
fn publish_and_park(slot: &LeaseSlot, lease: Arc<LocalLease>, parked: &mut Vec<Arc<LocalLease>>) {
park(slot.replace(lease), parked);
}
fn park(previous: Option<Arc<LocalLease>>, parked: &mut Vec<Arc<LocalLease>>) {
if let Some(previous) = previous {
parked.push(previous);
}
}
fn refusal_evidence(error: &AllocateError) -> Option<tollgate_core::BalanceShortfall> {
match error {
AllocateError::BalanceExhausted(evidence) => Some((*evidence).into()),
AllocateError::BalanceInsufficient(evidence) => Some(*evidence),
_ => None,
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Consolidation {
Installed,
KeptServing,
AcquireInstead,
ShutdownObserved,
}
async fn wait_for_quiescence(lease: &Arc<LocalLease>, deadline: tokio::time::Instant) -> bool {
const POLL: std::time::Duration = std::time::Duration::from_millis(5);
loop {
if Arc::strong_count(lease) == 1 && lease.is_only_local_view() {
return true;
}
if tokio::time::Instant::now() >= deadline {
return false;
}
tokio::time::sleep(POLL.min(deadline - tokio::time::Instant::now())).await;
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ReleasePass {
Complete,
BudgetExpired,
ShutdownObserved,
}
#[allow(
clippy::too_many_arguments,
reason = "one release step over the loop's own borrowed state"
)]
async fn release_quiesced(
allocator: &Arc<dyn LeaseAllocator>,
parked: &mut Vec<Arc<LocalLease>>,
clock: &Arc<dyn Clock>,
config: &LeaseManagerConfig,
slot: &Arc<LeaseSlot>,
health: &watch::Sender<bool>,
counters: &LeaseCounters,
shutdown: &mut watch::Receiver<bool>,
deadline: tokio::time::Instant,
) -> ReleasePass {
let mut queue = std::mem::take(parked).into_iter();
let mut retry = Vec::new();
while let Some(lease) = queue.next() {
if tokio::time::Instant::now() >= deadline {
*parked = std::iter::once(lease).chain(queue).chain(retry).collect();
return ReleasePass::BudgetExpired;
}
if Arc::strong_count(&lease) > 1 || !lease.is_only_local_view() {
retry.push(lease);
continue;
}
let grant = lease.grant();
let call_deadline = deadline.min(tokio::time::Instant::now() + config.store_call_timeout);
let call = tokio::time::timeout_at(
call_deadline,
allocator.release(
grant.lease_id,
grant.fencing_token,
lease.remaining(),
clock.now(),
),
);
let lease_id = grant.lease_id;
let call_outcome = tokio::select! {
outcome = call => outcome,
_ = shutdown.changed() => {
tracing::debug!(
lease = %lease_id,
"shutdown observed during a release; entering the shutdown release phase"
);
*parked = std::iter::once(lease).chain(queue).chain(retry).collect();
return ReleasePass::ShutdownObserved;
}
};
match call_outcome {
Err(_) => {
tracing::debug!(lease = %lease_id, "release timed out; retrying next tick");
retry.push(lease);
}
Ok(Ok(())) => counters.record_released(),
Ok(Err(error)) => match release_failure(&error) {
ReleaseFailure::Retryable => {
tracing::warn!(lease = %lease_id, %error, "release failed; retrying next tick");
retry.push(lease);
}
ReleaseFailure::Settled => {
counters.record_released();
tracing::debug!(lease = %lease_id, %error, "lease was already settled");
}
ReleaseFailure::Fenced => {
counters.record_released();
tracing::warn!(lease = %lease_id,
"store rejected this lease's capability; clearing the slot so this instance stops serving");
if let Some(current) = slot.take() {
retry.push(current);
}
}
ReleaseFailure::Integrity => {
counters.record_abandoned();
counters.record_integrity_fault(health);
tracing::error!(lease = %lease_id, units = lease.remaining().get(), %error,
"release violated the allocator contract; grant abandoned and readiness withdrawn");
}
},
}
}
*parked = retry;
ReleasePass::Complete
}
#[cfg(test)]
mod tests {
#[test]
fn aggregate_refill_counters_reject_overflow_in_fields_and_totals() {
let mut left = super::LeaseStats::ZERO;
left.acquired_units = u64::MAX;
let mut right = super::LeaseStats::ZERO;
right.acquired_units = 1;
assert!(left.checked_add(right).is_none());
left = super::LeaseStats::ZERO;
right = super::LeaseStats::ZERO;
left.acquire_refused[0] = u64::MAX;
right.acquire_refused[1] = 1;
assert!(left.checked_add(right).is_none());
right.acquire_refused[1] = 0;
assert_eq!(left.checked_add(right).unwrap().refused(), u64::MAX);
left = super::LeaseStats::ZERO;
right = super::LeaseStats::ZERO;
left.uncertain_acquires = u64::MAX;
right.uncertain_acquires = 1;
assert!(left.checked_add(right).is_none());
}
use std::collections::HashMap;
use std::sync::Mutex;
use async_trait::async_trait;
use jiff::Timestamp;
use tollgate_core::{FencingToken, LeaseGrant, LeaseId};
use tollgate_store::{AllocateError, Allocation, ReclaimBatch, StoreError, SystemClock};
use super::*;
#[test]
fn only_ambiguous_acquire_outcomes_increase_uncertainty() {
let counters = LeaseCounters::new();
for error in [
AllocateError::UnknownAccount,
AllocateError::AccountInactive,
AllocateError::InsufficientBalance,
AllocateError::UnknownLease,
AllocateError::Fenced,
AllocateError::LeaseNotActive,
AllocateError::InvalidRelease,
AllocateError::InvalidTtl,
AllocateError::BalanceOverflow,
] {
counters.record_acquire_refused(&error);
}
assert_eq!(counters.snapshot().uncertain_acquires, 0);
counters.record_acquire_refused(&AllocateError::Storage(StoreError("lost reply".into())));
assert_eq!(counters.snapshot().uncertain_acquires, 1);
counters.record_acquire_timeout();
assert_eq!(counters.snapshot().uncertain_acquires, 2);
}
#[tokio::test]
async fn a_request_raised_before_the_wait_is_not_lost() {
let refill = RefillRequests::default();
refill.request_refill();
tokio::time::timeout(std::time::Duration::from_secs(5), refill.requested())
.await
.expect("a request raised before the wait must complete it");
}
#[tokio::test]
async fn a_request_raised_during_the_wait_wakes_it() {
let refill = Arc::new(RefillRequests::default());
let signal = Arc::clone(&refill);
tokio::spawn(async move {
tokio::task::yield_now().await;
signal.request_refill();
});
tokio::time::timeout(std::time::Duration::from_secs(5), refill.requested())
.await
.expect("a request raised while waiting must wake the waiter");
}
#[tokio::test(start_paused = true)]
async fn a_request_is_consumed_by_the_wait_it_completes() {
let refill = RefillRequests::default();
refill.request_refill();
refill.requested().await;
assert!(
tokio::time::timeout(std::time::Duration::from_secs(5), refill.requested())
.await
.is_err(),
"the request was already answered; waiting again must block"
);
}
#[tokio::test(start_paused = true)]
async fn repeated_requests_collapse_into_one_wake() {
let refill = RefillRequests::default();
for _ in 0..10 {
refill.request_refill();
}
refill.requested().await;
assert!(
tokio::time::timeout(std::time::Duration::from_secs(5), refill.requested())
.await
.is_err(),
"ten requests are still one outstanding refill, not ten"
);
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum Refusal {
Storage,
Fenced,
InvalidRelease,
LeaseNotActive,
UnknownLease,
Hang,
}
struct ScriptedAllocator {
released: Mutex<Vec<LeaseId>>,
refusals: HashMap<LeaseId, Refusal>,
}
impl ScriptedAllocator {
fn new(refusals: impl IntoIterator<Item = (LeaseId, Refusal)>) -> Arc<Self> {
Arc::new(Self {
released: Mutex::new(Vec::new()),
refusals: refusals.into_iter().collect(),
})
}
fn released(&self) -> Vec<LeaseId> {
self.released.lock().unwrap().clone()
}
}
#[async_trait]
impl LeaseAllocator for ScriptedAllocator {
async fn acquire(
&self,
_account: AccountId,
_requested: CostUnits,
_ttl: SignedDuration,
_now: Timestamp,
) -> Result<Allocation, AllocateError> {
unreachable!("release_quiesced never acquires")
}
async fn release(
&self,
lease_id: LeaseId,
_fencing_token: FencingToken,
_unspent: CostUnits,
_now: Timestamp,
) -> Result<(), AllocateError> {
match self.refusals.get(&lease_id) {
Some(Refusal::Storage) => {
Err(AllocateError::Storage(StoreError("scripted outage".into())))
}
Some(Refusal::Fenced) => Err(AllocateError::Fenced),
Some(Refusal::InvalidRelease) => Err(AllocateError::InvalidRelease),
Some(Refusal::LeaseNotActive) => Err(AllocateError::LeaseNotActive),
Some(Refusal::UnknownLease) => Err(AllocateError::UnknownLease),
Some(Refusal::Hang) => std::future::pending().await,
None => {
self.released.lock().unwrap().push(lease_id);
Ok(())
}
}
}
async fn consolidate(
&self,
_lease_id: LeaseId,
_fencing_token: FencingToken,
_unspent: CostUnits,
_requested: CostUnits,
_needed: CostUnits,
_ttl: SignedDuration,
_now: Timestamp,
) -> Result<Allocation, AllocateError> {
unreachable!("release_quiesced never consolidates")
}
async fn reclaim_expired_batch(
&self,
_now: Timestamp,
_limit: std::num::NonZeroUsize,
) -> Result<ReclaimBatch, StoreError> {
unreachable!("release_quiesced never reclaims")
}
}
struct ConsolidatingAllocator {
answer: Mutex<Option<Result<Allocation, AllocateError>>>,
calls: AtomicU64,
needed: AtomicU64,
publish_during_call: Option<Arc<LeaseSlot>>,
}
impl ConsolidatingAllocator {
fn new(answer: Result<Allocation, AllocateError>) -> Arc<Self> {
Arc::new(Self {
answer: Mutex::new(Some(answer)),
calls: AtomicU64::new(0),
needed: AtomicU64::new(0),
publish_during_call: None,
})
}
}
#[async_trait]
impl LeaseAllocator for ConsolidatingAllocator {
async fn acquire(
&self,
_account: AccountId,
_requested: CostUnits,
_ttl: SignedDuration,
_now: Timestamp,
) -> Result<Allocation, AllocateError> {
unreachable!("these tests drive consolidation directly")
}
async fn release(
&self,
_lease_id: LeaseId,
_fencing_token: FencingToken,
_unspent: CostUnits,
_now: Timestamp,
) -> Result<(), AllocateError> {
Ok(())
}
async fn consolidate(
&self,
_lease_id: LeaseId,
_fencing_token: FencingToken,
_unspent: CostUnits,
_requested: CostUnits,
needed: CostUnits,
_ttl: SignedDuration,
_now: Timestamp,
) -> Result<Allocation, AllocateError> {
self.calls.fetch_add(1, Ordering::SeqCst);
self.needed.store(needed.get(), Ordering::SeqCst);
if let Some(slot) = &self.publish_during_call {
assert!(slot.replace(parked_lease(99)).is_none());
}
self.answer
.lock()
.unwrap()
.take()
.expect("one scripted consolidation per test")
}
async fn reclaim_expired_batch(
&self,
_now: Timestamp,
_limit: std::num::NonZeroUsize,
) -> Result<ReclaimBatch, StoreError> {
unreachable!("these tests never reclaim")
}
}
fn grant(id: u128, units: u64) -> LeaseGrant {
LeaseGrant {
lease_id: LeaseId(id),
account_id: AccountId(1),
fencing_token: FencingToken(1),
units: CostUnits(units),
expires_at: Timestamp::from_second(3_600).unwrap(),
}
}
fn allocation(id: u128, units: u64) -> Allocation {
Allocation {
grant: grant(id, units),
funding: None,
}
}
async fn consolidate_once(
allocator: Arc<dyn LeaseAllocator>,
harness: &Harness,
lease: Arc<LocalLease>,
) -> (Consolidation, Option<u128>, Vec<u128>) {
drop(harness.slot.replace(lease));
let mut parked = Vec::new();
let (_tx, mut shutdown) = watch::channel(false);
let outcome = consolidate_live_lease(
&allocator,
&mut parked,
&clock(),
&harness.config,
&harness.slot,
&Arc::new(RefillRequests::default()),
&harness.health,
&harness.counters,
&mut shutdown,
)
.await;
let served = harness.slot.load().map(|l| l.grant().lease_id.0);
let parked_ids = parked.iter().map(|l| l.grant().lease_id.0).collect();
(outcome, served, parked_ids)
}
#[tokio::test(start_paused = true)]
async fn a_successful_consolidation_installs_the_grant_and_parks_nothing() {
let harness = Harness::new();
let allocator = ConsolidatingAllocator::new(Ok(allocation(2, 500)));
let (outcome, served, parked) = consolidate_once(
Arc::clone(&allocator) as Arc<dyn LeaseAllocator>,
&harness,
parked_lease(1),
)
.await;
assert_eq!(outcome, Consolidation::Installed);
assert_eq!(served, Some(2), "the larger grant is serving");
assert!(
parked.is_empty(),
"the superseded lease was already settled"
);
let stats = harness.stats();
assert_eq!(stats.consolidated, 1);
assert_eq!(stats.acquired, 1);
assert_eq!(stats.released, 1, "consolidation settled its predecessor");
assert!(!harness.counters.acquire_pending());
assert_eq!(stats.acquired_units, 500);
}
#[tokio::test(start_paused = true)]
async fn consolidation_carries_the_refused_quote() {
for refusals in [&[][..], &[150, 252, 101][..]] {
let harness = Harness::new();
let allocator = ConsolidatingAllocator::new(Ok(allocation(2, 500)));
let lease = parked_lease(1);
for "e in refusals {
assert!(
lease
.try_debit(CostUnits(quote), Timestamp::UNIX_EPOCH)
.is_err()
);
}
let (outcome, _, _) = consolidate_once(
Arc::clone(&allocator) as Arc<dyn LeaseAllocator>,
&harness,
lease,
)
.await;
assert_eq!(outcome, Consolidation::Installed);
assert_eq!(
allocator.needed.load(Ordering::SeqCst),
refusals.iter().copied().max().unwrap_or(0)
);
}
}
#[tokio::test(start_paused = true)]
async fn a_consolidated_tail_publishes_the_accounts_remaining_funding() {
let harness = Harness::new();
let mut tail = grant(2, 1);
tail.fencing_token = FencingToken(2);
let allocator = ConsolidatingAllocator::new(Ok(Allocation {
grant: tail,
funding: Some(tollgate_core::BalanceShortfall {
remaining: CostUnits(1),
period_end: None,
}),
}));
let (outcome, served, _) = consolidate_once(allocator, &harness, parked_lease(1)).await;
assert_eq!(outcome, Consolidation::Installed);
assert_eq!(served, Some(2));
assert_eq!(
harness.slot.funding_evidence(Timestamp::UNIX_EPOCH),
Some(CostUnits(1))
);
}
#[tokio::test(start_paused = true)]
async fn consolidation_retains_a_grant_published_during_the_store_call() {
for (answer, expected_outcome, expected_served) in [
(Ok(allocation(2, 500)), Consolidation::Installed, 2),
(
Err(AllocateError::InsufficientBalance),
Consolidation::KeptServing,
1,
),
(
Err(AllocateError::InvalidRelease),
Consolidation::KeptServing,
1,
),
] {
let harness = Harness::new();
let mut allocator = ConsolidatingAllocator::new(answer);
Arc::get_mut(&mut allocator).unwrap().publish_during_call =
Some(Arc::clone(&harness.slot));
let (outcome, served, parked) =
consolidate_once(allocator, &harness, parked_lease(1)).await;
assert_eq!(outcome, expected_outcome);
assert_eq!(served, Some(expected_served));
assert_eq!(
parked,
vec![99],
"a concurrent publisher's grant must remain available for quiesced release"
);
}
}
#[tokio::test(start_paused = true)]
async fn a_rolled_back_consolidation_returns_the_lease_to_the_slot() {
for error in [
AllocateError::InsufficientBalance,
AllocateError::BalanceExhausted(tollgate_core::BalanceExhaustion { period_end: None }),
AllocateError::BalanceInsufficient(tollgate_core::BalanceShortfall {
remaining: CostUnits(3),
period_end: None,
}),
AllocateError::AccountInactive,
AllocateError::UnknownAccount,
] {
let harness = Harness::new();
let allocator = ConsolidatingAllocator::new(Err(error.clone()));
let (outcome, served, parked) = consolidate_once(
Arc::clone(&allocator) as Arc<dyn LeaseAllocator>,
&harness,
parked_lease(1),
)
.await;
assert_eq!(outcome, Consolidation::KeptServing, "{error}");
assert_eq!(
served,
Some(1),
"still serving the untouched lease: {error}"
);
assert!(parked.is_empty(), "{error}");
assert_eq!(
harness.slot.funding_evidence(Timestamp::UNIX_EPOCH),
match error {
AllocateError::BalanceExhausted(_) => Some(CostUnits::ZERO),
AllocateError::BalanceInsufficient(evidence) => Some(evidence.remaining),
_ => None,
},
"{error}"
);
assert!(harness.is_healthy(), "an ordinary refusal is not a fault");
}
}
#[tokio::test(start_paused = true)]
async fn an_ambiguous_consolidation_parks_the_grant_rather_than_reinstating_it() {
let harness = Harness::new();
let allocator = ConsolidatingAllocator::new(Err(AllocateError::Storage(StoreError(
"connection reset".into(),
))));
let (outcome, served, parked) = consolidate_once(
Arc::clone(&allocator) as Arc<dyn LeaseAllocator>,
&harness,
parked_lease(1),
)
.await;
assert_eq!(outcome, Consolidation::KeptServing);
assert_eq!(
served, None,
"the slot fails closed rather than double-spending"
);
assert_eq!(parked, vec![1], "and the release pass settles it");
}
#[tokio::test(start_paused = true)]
async fn a_settled_lease_falls_through_to_an_ordinary_acquire() {
for error in [
AllocateError::UnknownLease,
AllocateError::LeaseNotActive,
AllocateError::Fenced,
] {
let harness = Harness::new();
let allocator = ConsolidatingAllocator::new(Err(error.clone()));
let (outcome, served, parked) = consolidate_once(
Arc::clone(&allocator) as Arc<dyn LeaseAllocator>,
&harness,
parked_lease(1),
)
.await;
assert_eq!(outcome, Consolidation::AcquireInstead, "{error}");
assert_eq!(served, None, "{error}");
assert!(parked.is_empty(), "the store already settled it: {error}");
}
}
#[tokio::test(start_paused = true)]
async fn an_over_claimed_fold_withdraws_readiness_and_keeps_serving() {
let harness = Harness::new();
let allocator = ConsolidatingAllocator::new(Err(AllocateError::InvalidRelease));
let (outcome, served, parked) = consolidate_once(
Arc::clone(&allocator) as Arc<dyn LeaseAllocator>,
&harness,
parked_lease(1),
)
.await;
assert_eq!(outcome, Consolidation::KeptServing);
assert_eq!(served, Some(1));
assert!(parked.is_empty());
assert!(!harness.is_healthy(), "an accounting fault is never silent");
}
#[tokio::test(start_paused = true)]
async fn a_consolidation_defers_while_a_reservation_is_in_flight() {
let harness = Harness::new();
let allocator = ConsolidatingAllocator::new(Ok(allocation(2, 500)));
let lease = parked_lease(1);
let _in_flight = Arc::clone(&lease);
let (outcome, served, parked) = consolidate_once(
Arc::clone(&allocator) as Arc<dyn LeaseAllocator>,
&harness,
lease,
)
.await;
assert_eq!(outcome, Consolidation::KeptServing);
assert_eq!(served, Some(1), "the slot keeps serving it");
assert!(parked.is_empty());
assert_eq!(
allocator.calls.load(Ordering::SeqCst),
0,
"no store call is made against an inexact aggregate"
);
assert_eq!(harness.stats().consolidations_deferred, 1);
}
struct Harness {
config: LeaseManagerConfig,
slot: Arc<LeaseSlot>,
health: watch::Sender<bool>,
healthy: watch::Receiver<bool>,
counters: LeaseCounters,
}
impl Harness {
fn new() -> Self {
let (health, healthy) = watch::channel(true);
Self {
counters: LeaseCounters::new(),
config: LeaseManagerConfig {
account: AccountId(1),
target_grant: CostUnits(1_000),
low_water: CostUnits(250),
lease_ttl: SignedDuration::from_secs(60),
expiry_safety_margin: SignedDuration::ZERO,
poll_interval: std::time::Duration::from_millis(5),
store_call_timeout: std::time::Duration::from_millis(50),
shutdown_release_deadline: std::time::Duration::from_secs(10),
},
slot: LeaseSlot::for_account(AccountId(1)),
health,
healthy,
}
}
async fn release_quiesced(
&self,
allocator: &Arc<dyn LeaseAllocator>,
parked: &mut Vec<Arc<LocalLease>>,
) -> ReleasePass {
self.release_pass(allocator, parked, std::time::Duration::from_secs(3_600))
.await
}
async fn release_pass(
&self,
allocator: &Arc<dyn LeaseAllocator>,
parked: &mut Vec<Arc<LocalLease>>,
budget: std::time::Duration,
) -> ReleasePass {
let (_tx, mut shutdown) = watch::channel(false);
release_quiesced(
allocator,
parked,
&clock(),
&self.config,
&self.slot,
&self.health,
&self.counters,
&mut shutdown,
tokio::time::Instant::now() + budget,
)
.await
}
fn stats(&self) -> LeaseStats {
self.counters.snapshot()
}
fn is_healthy(&self) -> bool {
*self.healthy.borrow()
}
}
fn parked_lease(id: u128) -> Arc<LocalLease> {
Arc::new(LocalLease::new(
LeaseGrant {
lease_id: LeaseId(id),
account_id: AccountId(1),
fencing_token: FencingToken(1),
units: CostUnits(100),
expires_at: Timestamp::from_second(3_600).unwrap(),
},
CostUnits(10),
))
}
fn clock() -> Arc<dyn Clock> {
Arc::new(SystemClock)
}
#[tokio::test(start_paused = true)]
async fn shutdown_distinguishes_settled_leases_from_unconfirmed_or_invalid_releases() {
for (refusal, released, faulted) in [
(Refusal::Storage, 0, false),
(Refusal::Fenced, 1, false),
(Refusal::InvalidRelease, 0, true),
(Refusal::LeaseNotActive, 1, false),
(Refusal::UnknownLease, 1, false),
(Refusal::Hang, 0, false),
] {
let harness = Harness::new();
drop(harness.slot.replace(parked_lease(1)));
let manager = LeaseManager::spawn(
ScriptedAllocator::new([(LeaseId(1), refusal)]),
harness.slot,
Arc::new(crate::ManualClock::new(Timestamp::from_second(0).unwrap())),
harness.config,
)
.unwrap();
let counters = manager.counters();
let health = manager.health();
let report = manager.shutdown().await;
assert_eq!(report.released, released);
assert_eq!(report.abandoned, 1 - released);
assert_eq!(counters.snapshot().released, released);
assert_eq!(counters.snapshot().abandoned, 1 - released);
assert_eq!(counters.integrity_fault(), faulted);
assert!(
!*health.borrow(),
"every stopped task is unhealthy, including clean stops"
);
}
}
#[tokio::test(start_paused = true)]
async fn runtime_deadline_shortens_the_managers_actual_release_pass() {
let harness = Harness::new();
drop(harness.slot.replace(parked_lease(1)));
let manager = LeaseManager::spawn(
ScriptedAllocator::new([(LeaseId(1), Refusal::Hang)]),
harness.slot,
Arc::new(crate::ManualClock::new(Timestamp::from_second(0).unwrap())),
harness.config,
)
.unwrap();
let began = tokio::time::Instant::now();
manager.stop_at(began + std::time::Duration::from_millis(5));
let report = tokio::time::timeout(std::time::Duration::from_millis(6), manager.shutdown())
.await
.unwrap();
assert_eq!(report.abandoned, 1);
assert!(!report.task_died);
assert_eq!(began.elapsed(), std::time::Duration::from_millis(5));
}
#[tokio::test]
async fn storage_failure_does_not_skip_the_next_parked_lease() {
let scripted = ScriptedAllocator::new([(LeaseId(1), Refusal::Storage)]);
let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
let harness = Harness::new();
let mut parked = vec![parked_lease(1), parked_lease(2)];
harness.release_quiesced(&allocator, &mut parked).await;
assert_eq!(scripted.released(), [LeaseId(2)]);
assert_eq!(parked.len(), 1, "only the failed lease stays parked");
assert_eq!(parked[0].grant().lease_id, LeaseId(1));
}
#[tokio::test]
async fn unquiesced_lease_is_never_released() {
let scripted = ScriptedAllocator::new([]);
let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
let harness = Harness::new();
let lease = parked_lease(7);
let in_flight = Arc::clone(&lease);
let mut parked = vec![lease];
harness.release_quiesced(&allocator, &mut parked).await;
assert!(scripted.released().is_empty());
assert_eq!(parked.len(), 1, "held lease stays parked");
assert_eq!(
harness.stats().released,
0,
"a lease still on our books has not been released"
);
drop(in_flight);
harness.release_quiesced(&allocator, &mut parked).await;
assert_eq!(scripted.released(), [LeaseId(7)]);
assert!(parked.is_empty());
assert_eq!(harness.stats().released, 1);
}
#[tokio::test]
async fn settled_refusals_drop_silently() {
for refusal in [Refusal::LeaseNotActive, Refusal::UnknownLease] {
let scripted = ScriptedAllocator::new([(LeaseId(3), refusal)]);
let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
let harness = Harness::new();
drop(harness.slot.replace(parked_lease(99)));
let mut parked = vec![parked_lease(3)];
harness.release_quiesced(&allocator, &mut parked).await;
assert!(scripted.released().is_empty());
assert!(parked.is_empty(), "settled lease is not retried");
assert!(harness.is_healthy(), "settlement is not a health event");
assert!(harness.slot.load().is_some(), "the slot is untouched");
assert_eq!(
harness.stats().released,
1,
"already settled still means the store no longer holds it"
);
}
}
#[tokio::test]
async fn fenced_release_clears_the_slot() {
let scripted = ScriptedAllocator::new([(LeaseId(3), Refusal::Fenced)]);
let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
let harness = Harness::new();
drop(harness.slot.replace(parked_lease(99)));
let mut parked = vec![parked_lease(3)];
harness.release_quiesced(&allocator, &mut parked).await;
assert!(
harness.slot.load().is_none(),
"an instance with divergent lease identity must stop serving"
);
let ids: Vec<_> = parked.iter().map(|l| l.grant().lease_id).collect();
assert_eq!(
ids,
[LeaseId(99)],
"the fenced lease is not retried, and the slot's live lease is not lost"
);
assert_eq!(
parked[0].remaining(),
CostUnits(100),
"its unspent units are still accounted for"
);
}
#[tokio::test]
async fn invalid_release_fails_readiness() {
let scripted = ScriptedAllocator::new([(LeaseId(3), Refusal::InvalidRelease)]);
let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
let harness = Harness::new();
let mut parked = vec![parked_lease(3)];
assert!(harness.is_healthy());
harness.release_quiesced(&allocator, &mut parked).await;
assert!(
!harness.is_healthy(),
"accounting divergence must drop readiness"
);
assert!(harness.counters.integrity_fault());
assert_eq!(harness.counters.snapshot().abandoned, 1);
assert_eq!(harness.counters.snapshot().released, 0);
}
#[tokio::test(start_paused = true)]
async fn hung_release_times_out_and_reparks() {
let scripted = ScriptedAllocator::new([(LeaseId(5), Refusal::Hang)]);
let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
let harness = Harness::new();
let mut parked = vec![parked_lease(5), parked_lease(6)];
tokio::time::timeout(
std::time::Duration::from_secs(30),
harness.release_quiesced(&allocator, &mut parked),
)
.await
.expect("a hung release must not stall the pass");
assert_eq!(scripted.released(), [LeaseId(6)], "the pass continues");
assert_eq!(parked.len(), 1, "the timed-out lease is retried, not lost");
assert_eq!(parked[0].grant().lease_id, LeaseId(5));
assert!(
harness.is_healthy(),
"a timeout is not accounting divergence"
);
}
#[tokio::test(start_paused = true)]
async fn a_release_pass_costs_one_budget_whatever_the_parked_count() {
let scripted = ScriptedAllocator::new((1..=6).map(|id| (LeaseId(id), Refusal::Hang)));
let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
let harness = Harness::new();
let mut parked: Vec<_> = (1..=6).map(parked_lease).collect();
let budget = std::time::Duration::from_millis(50);
let began = tokio::time::Instant::now();
let outcome = harness.release_pass(&allocator, &mut parked, budget).await;
let elapsed = began.elapsed();
assert_eq!(outcome, ReleasePass::BudgetExpired);
assert!(
elapsed < budget * 2,
"a pass over six hung leases took {elapsed:?}, which is per-call not per-pass"
);
assert_eq!(parked.len(), 6, "every lease is still parked, none dropped");
}
#[tokio::test(start_paused = true)]
async fn a_lease_that_eats_the_budget_yields_its_place() {
let scripted = ScriptedAllocator::new([(LeaseId(1), Refusal::Hang)]);
let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
let harness = Harness::new();
let mut parked = vec![parked_lease(1), parked_lease(2), parked_lease(3)];
let budget = std::time::Duration::from_millis(50);
let outcome = harness.release_pass(&allocator, &mut parked, budget).await;
assert_eq!(outcome, ReleasePass::BudgetExpired);
let order: Vec<_> = parked.iter().map(|l| l.grant().lease_id).collect();
assert_eq!(
order,
[LeaseId(2), LeaseId(3), LeaseId(1)],
"the hung lease must not hold the front of the queue every pass"
);
assert_eq!(scripted.released(), [], "nothing settled under a hung head");
}
#[tokio::test(start_paused = true)]
async fn shutdown_during_a_release_reparks_every_lease() {
let scripted = ScriptedAllocator::new([(LeaseId(1), Refusal::Hang)]);
let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
let harness = Harness::new();
let mut parked = vec![parked_lease(1), parked_lease(2), parked_lease(3)];
let (tx, mut shutdown) = watch::channel(false);
let clock = clock();
let pass = release_quiesced(
&allocator,
&mut parked,
&clock,
&harness.config,
&harness.slot,
&harness.health,
&harness.counters,
&mut shutdown,
tokio::time::Instant::now() + std::time::Duration::from_secs(3_600),
);
let signal = async {
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
tx.send(true).unwrap();
};
let (outcome, ()) = tokio::join!(pass, signal);
assert_eq!(outcome, ReleasePass::ShutdownObserved);
assert_eq!(
parked.len(),
3,
"the in-flight lease and the unexamined ones all survive the pass"
);
assert_eq!(parked[0].grant().lease_id, LeaseId(1));
}
}