1use std::sync::Arc;
51use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
52
53use atomic_waker::AtomicWaker;
54use jiff::SignedDuration;
55use tokio::sync::watch;
56use tracing::Instrument as _;
57
58use tollgate_admission::LeaseSlot;
59use tollgate_core::{AccountId, CostUnits, LeaseGrant, LocalLease, RefillSignal, RefillVerdict};
60use tollgate_store::{AllocateError, LeaseAllocator};
61
62use tollgate_store::Clock;
63
64#[derive(Debug, Clone, Copy)]
71pub struct LeaseManagerConfig {
72 pub account: AccountId,
75 pub target_grant: CostUnits,
86 pub low_water: CostUnits,
98 pub lease_ttl: SignedDuration,
107 pub expiry_safety_margin: SignedDuration,
118 pub poll_interval: std::time::Duration,
127 pub store_call_timeout: std::time::Duration,
138 pub shutdown_release_deadline: std::time::Duration,
147}
148
149#[derive(Debug, Clone, Copy)]
152pub struct AccountLeaseConfig {
153 pub target_grant: CostUnits,
155 pub low_water: CostUnits,
158 pub lease_ttl: SignedDuration,
160 pub expiry_safety_margin: SignedDuration,
163 pub poll_interval: std::time::Duration,
165 pub store_call_timeout: std::time::Duration,
168 pub shutdown_release_deadline: std::time::Duration,
173}
174
175impl AccountLeaseConfig {
176 #[must_use]
178 pub fn for_account(self, account: AccountId) -> LeaseManagerConfig {
179 LeaseManagerConfig {
180 account,
181 target_grant: self.target_grant,
182 low_water: self.low_water,
183 lease_ttl: self.lease_ttl,
184 expiry_safety_margin: self.expiry_safety_margin,
185 poll_interval: self.poll_interval,
186 store_call_timeout: self.store_call_timeout,
187 shutdown_release_deadline: self.shutdown_release_deadline,
188 }
189 }
190
191 pub fn validate(&self) -> Result<(), LeaseManagerConfigError> {
197 self.for_account(AccountId(0)).validate()
198 }
199}
200
201#[derive(Debug, Clone, Copy, PartialEq, Eq)]
203pub struct LeaseManagerReport {
204 pub released: u64,
206 pub abandoned: u64,
209 pub task_died: bool,
213}
214
215#[derive(Debug, Clone, Copy, PartialEq, Eq)]
217pub struct LeaseManagerConfigError(pub &'static str);
218
219impl std::fmt::Display for LeaseManagerConfigError {
220 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
221 f.write_str(self.0)
222 }
223}
224
225impl std::error::Error for LeaseManagerConfigError {}
226
227impl LeaseManagerConfig {
228 pub fn validate(&self) -> Result<(), LeaseManagerConfigError> {
239 if self.target_grant.is_zero() {
240 return Err(LeaseManagerConfigError("target_grant must be positive"));
241 }
242 if self.low_water >= self.target_grant {
243 return Err(LeaseManagerConfigError(
244 "low_water must be below target_grant",
245 ));
246 }
247 if self.lease_ttl <= SignedDuration::ZERO {
248 return Err(LeaseManagerConfigError("lease_ttl must be positive"));
249 }
250 if self.expiry_safety_margin < SignedDuration::ZERO {
251 return Err(LeaseManagerConfigError(
252 "expiry_safety_margin must not be negative",
253 ));
254 }
255 if self.expiry_safety_margin >= self.lease_ttl {
256 return Err(LeaseManagerConfigError(
257 "expiry_safety_margin must be shorter than lease_ttl",
258 ));
259 }
260 if self.poll_interval.is_zero() {
261 return Err(LeaseManagerConfigError("poll_interval must be positive"));
262 }
263 if self.store_call_timeout.is_zero() {
264 return Err(LeaseManagerConfigError(
265 "store_call_timeout must be positive",
266 ));
267 }
268 if self.shutdown_release_deadline.is_zero() {
269 return Err(LeaseManagerConfigError(
270 "shutdown_release_deadline must be positive",
271 ));
272 }
273 Ok(())
274 }
275}
276
277#[derive(Debug, Default)]
291struct RefillRequests {
292 waker: AtomicWaker,
293 pending: AtomicBool,
294}
295
296impl RefillSignal for RefillRequests {
297 fn request_refill(&self) {
298 self.pending.store(true, Ordering::Release);
301 self.waker.wake();
302 }
303}
304
305impl RefillRequests {
306 async fn requested(&self) {
312 std::future::poll_fn(|cx| {
313 if self.pending.swap(false, Ordering::Acquire) {
314 return std::task::Poll::Ready(());
315 }
316 self.waker.register(cx.waker());
317 if self.pending.swap(false, Ordering::Acquire) {
322 std::task::Poll::Ready(())
323 } else {
324 std::task::Poll::Pending
325 }
326 })
327 .await;
328 }
329}
330
331#[derive(Debug)]
342pub struct LeaseCounters {
343 acquired: AtomicU64,
344 acquired_units: AtomicU64,
345 acquire_timeouts: AtomicU64,
346 uncertain_acquires: AtomicU64,
347 acquire_pending: AtomicBool,
348 integrity_fault: AtomicBool,
349 acquire_refused: [AtomicU64; AllocateError::COUNT],
350 released: AtomicU64,
351 abandoned: AtomicU64,
352 consolidated: AtomicU64,
353 consolidations_deferred: AtomicU64,
354}
355
356impl LeaseCounters {
357 #[must_use]
359 pub const fn new() -> Self {
360 LeaseCounters {
361 acquired: AtomicU64::new(0),
362 acquired_units: AtomicU64::new(0),
363 acquire_timeouts: AtomicU64::new(0),
364 uncertain_acquires: AtomicU64::new(0),
365 acquire_pending: AtomicBool::new(false),
366 integrity_fault: AtomicBool::new(false),
367 acquire_refused: [const { AtomicU64::new(0) }; AllocateError::COUNT],
368 released: AtomicU64::new(0),
369 abandoned: AtomicU64::new(0),
370 consolidated: AtomicU64::new(0),
371 consolidations_deferred: AtomicU64::new(0),
372 }
373 }
374
375 pub(crate) fn acquire_pending(&self) -> bool {
376 self.acquire_pending.load(Ordering::Acquire)
377 }
378
379 pub(crate) fn integrity_fault(&self) -> bool {
380 self.integrity_fault.load(Ordering::Acquire)
381 }
382
383 fn record_integrity_fault(&self, health: &watch::Sender<bool>) {
384 self.integrity_fault.store(true, Ordering::Release);
385 crate::signal(health, false, "lease-manager health");
386 }
387
388 fn record_acquired(&self, units: CostUnits) {
389 self.acquired.fetch_add(1, Ordering::Relaxed);
390 self.acquired_units
391 .fetch_add(units.get(), Ordering::Relaxed);
392 }
393
394 fn record_acquire_timeout(&self) {
398 self.acquire_timeouts.fetch_add(1, Ordering::Relaxed);
399 self.uncertain_acquires.fetch_add(1, Ordering::Relaxed);
400 }
401
402 fn record_consolidated(&self, units: CostUnits) {
406 self.record_acquired(units);
407 self.record_released();
408 self.consolidated.fetch_add(1, Ordering::Relaxed);
409 }
410
411 fn record_consolidation_deferred(&self) {
412 self.consolidations_deferred.fetch_add(1, Ordering::Relaxed);
413 }
414
415 fn record_acquire_refused(&self, error: &AllocateError) {
416 self.acquire_refused[error.index()].fetch_add(1, Ordering::Relaxed);
417 if matches!(error, AllocateError::Storage(_)) {
418 self.uncertain_acquires.fetch_add(1, Ordering::Relaxed);
419 }
420 }
421
422 fn record_released(&self) {
425 self.released.fetch_add(1, Ordering::Relaxed);
426 }
427
428 fn record_abandoned(&self) {
429 self.abandoned.fetch_add(1, Ordering::Relaxed);
430 }
431
432 #[must_use]
435 pub fn snapshot(&self) -> LeaseStats {
436 LeaseStats {
437 acquired: self.acquired.load(Ordering::Relaxed),
438 acquired_units: self.acquired_units.load(Ordering::Relaxed),
439 acquire_timeouts: self.acquire_timeouts.load(Ordering::Relaxed),
440 uncertain_acquires: self.uncertain_acquires.load(Ordering::Relaxed),
441 acquire_refused: std::array::from_fn(|slot| {
442 self.acquire_refused[slot].load(Ordering::Relaxed)
443 }),
444 released: self.released.load(Ordering::Relaxed),
445 abandoned: self.abandoned.load(Ordering::Relaxed),
446 consolidated: self.consolidated.load(Ordering::Relaxed),
447 consolidations_deferred: self.consolidations_deferred.load(Ordering::Relaxed),
448 }
449 }
450}
451
452impl Default for LeaseCounters {
453 fn default() -> Self {
454 Self::new()
455 }
456}
457
458#[derive(Debug, Clone, Copy, PartialEq, Eq)]
460pub struct LeaseStats {
461 pub acquired: u64,
463 pub acquired_units: u64,
466 pub acquire_timeouts: u64,
468 pub uncertain_acquires: u64,
473 pub acquire_refused: [u64; AllocateError::COUNT],
475 pub released: u64,
477 pub abandoned: u64,
481 pub consolidated: u64,
487 pub consolidations_deferred: u64,
492}
493
494impl LeaseStats {
495 pub const ZERO: Self = Self {
497 acquired: 0,
498 acquired_units: 0,
499 acquire_timeouts: 0,
500 uncertain_acquires: 0,
501 acquire_refused: [0; AllocateError::COUNT],
502 released: 0,
503 abandoned: 0,
504 consolidated: 0,
505 consolidations_deferred: 0,
506 };
507
508 pub fn checked_add(self, other: Self) -> Option<Self> {
510 let mut acquire_refused = [0; AllocateError::COUNT];
511 for (i, value) in acquire_refused.iter_mut().enumerate() {
512 *value = self.acquire_refused[i].checked_add(other.acquire_refused[i])?;
513 }
514 acquire_refused
516 .iter()
517 .try_fold(0_u64, |sum, value| sum.checked_add(*value))?;
518 Some(Self {
519 acquired: self.acquired.checked_add(other.acquired)?,
520 acquired_units: self.acquired_units.checked_add(other.acquired_units)?,
521 acquire_timeouts: self.acquire_timeouts.checked_add(other.acquire_timeouts)?,
522 uncertain_acquires: self
523 .uncertain_acquires
524 .checked_add(other.uncertain_acquires)?,
525 acquire_refused,
526 released: self.released.checked_add(other.released)?,
527 abandoned: self.abandoned.checked_add(other.abandoned)?,
528 consolidated: self.consolidated.checked_add(other.consolidated)?,
529 consolidations_deferred: self
530 .consolidations_deferred
531 .checked_add(other.consolidations_deferred)?,
532 })
533 }
534 pub fn refusals_by_name(&self) -> impl Iterator<Item = (&'static str, u64)> + '_ {
536 AllocateError::NAMES
537 .iter()
538 .copied()
539 .zip(self.acquire_refused.iter().copied())
540 }
541
542 #[must_use]
544 pub fn refused(&self) -> u64 {
545 self.acquire_refused.iter().sum()
546 }
547}
548
549pub struct LeaseManager {
551 shutdown: watch::Sender<bool>,
552 health: watch::Receiver<bool>,
553 handle: Option<tokio::task::JoinHandle<LeaseManagerReport>>,
554 counters: Arc<LeaseCounters>,
555 deadline: Arc<crate::ShutdownDeadline>,
556 paused: Arc<AtomicBool>,
557}
558
559impl LeaseManager {
560 pub fn spawn(
575 allocator: Arc<dyn LeaseAllocator>,
576 slot: Arc<LeaseSlot>,
577 clock: Arc<dyn Clock>,
578 config: LeaseManagerConfig,
579 ) -> Result<Self, LeaseManagerConfigError> {
580 config.validate()?;
581 let (shutdown, shutdown_rx) = watch::channel(false);
582 let (health_tx, health) = crate::task_health::TaskHealth::channel(true);
583 let account = config.account;
584 let counters = Arc::new(LeaseCounters::new());
585 let task_counters = Arc::clone(&counters);
586 let deadline = Arc::new(crate::ShutdownDeadline::default());
587 let task_deadline = Arc::clone(&deadline);
588 let paused = Arc::new(AtomicBool::new(false));
589 let task_paused = Arc::clone(&paused);
590 let handle = tokio::spawn(
591 async move {
592 run(
593 allocator,
594 slot,
595 clock,
596 config,
597 shutdown_rx,
598 health_tx.sender(),
599 &task_counters,
600 &task_deadline,
601 &task_paused,
602 )
603 .await
604 }
605 .instrument(tracing::info_span!("lease_manager", %account)),
606 );
607 Ok(LeaseManager {
608 shutdown,
609 health,
610 handle: Some(handle),
611 counters,
612 deadline,
613 paused,
614 })
615 }
616
617 #[must_use]
624 pub fn counters(&self) -> Arc<LeaseCounters> {
625 Arc::clone(&self.counters)
626 }
627
628 #[must_use]
633 pub fn health(&self) -> watch::Receiver<bool> {
634 self.health.clone()
635 }
636
637 pub(crate) fn pause_refills(&self) {
640 self.paused.store(true, Ordering::Release);
641 }
642
643 pub(crate) fn stop_at(&self, deadline: tokio::time::Instant) {
644 self.deadline.constrain(deadline);
645 crate::signal(&self.shutdown, true, "lease-manager shutdown");
646 }
647
648 pub async fn shutdown(mut self) -> LeaseManagerReport {
652 crate::signal(&self.shutdown, true, "lease-manager shutdown");
653 let died = LeaseManagerReport {
654 released: 0,
655 abandoned: 0,
656 task_died: true,
657 };
658 match self.handle.as_mut() {
660 Some(handle) => handle.await.unwrap_or(died),
663 None => died,
664 }
665 }
666}
667
668impl Drop for LeaseManager {
669 fn drop(&mut self) {
670 if let Some(handle) = self.handle.take() {
674 handle.abort();
675 }
676 }
677}
678
679#[allow(
682 clippy::too_many_arguments,
683 reason = "task entry point: every argument is a handle the loop owns for its lifetime, assembled once by spawn"
684)]
685async fn run(
686 allocator: Arc<dyn LeaseAllocator>,
687 slot: Arc<LeaseSlot>,
688 clock: Arc<dyn Clock>,
689 config: LeaseManagerConfig,
690 mut shutdown: watch::Receiver<bool>,
691 health: &watch::Sender<bool>,
692 counters: &LeaseCounters,
693 shutdown_deadline: &crate::ShutdownDeadline,
694 paused: &AtomicBool,
695) -> LeaseManagerReport {
696 let mut tick = tokio::time::interval(config.poll_interval);
697 tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
698 let refill = Arc::new(RefillRequests::default());
703 let mut parked: Vec<Arc<LocalLease>> = Vec::new();
706 loop {
707 tokio::select! {
708 _ = tick.tick() => {}
709 () = refill.requested() => {}
714 changed = shutdown.changed() => {
715 if changed.is_err() {
718 break;
719 }
720 }
721 }
722 if *shutdown.borrow() {
723 break;
724 }
725
726 if paused.load(Ordering::Acquire) {
727 continue;
728 }
729
730 let now = clock.now();
731 let pass_deadline = tokio::time::Instant::now() + config.store_call_timeout;
735 if release_quiesced(
736 &allocator,
737 &mut parked,
738 &clock,
739 &config,
740 &slot,
741 health,
742 counters,
743 &mut shutdown,
744 pass_deadline,
745 )
746 .await
747 == ReleasePass::ShutdownObserved
748 {
749 break;
750 }
751 let rotation = match slot.load_observed() {
752 None => Rotation::Acquire,
753 Some(lease) if now >= lease.usable_until() => {
754 drop(lease);
758 if let Some(old) = slot.take() {
763 parked.push(old);
764 }
765 if release_quiesced(
768 &allocator,
769 &mut parked,
770 &clock,
771 &config,
772 &slot,
773 health,
774 counters,
775 &mut shutdown,
776 pass_deadline,
777 )
778 .await
779 == ReleasePass::ShutdownObserved
780 {
781 break;
782 }
783 Rotation::Acquire
784 }
785 Some(lease) => match lease.refill_due_or_rearm() {
786 RefillVerdict::Idle => Rotation::Idle,
787 RefillVerdict::Draining => Rotation::Acquire,
788 RefillVerdict::Refused => Rotation::Consolidate,
794 },
795 };
796 if rotation == Rotation::Idle {
797 continue;
798 }
799
800 if paused.load(Ordering::Acquire) || *shutdown.borrow() {
807 continue;
808 }
809
810 if rotation == Rotation::Consolidate {
811 match consolidate_live_lease(
812 &allocator,
813 &mut parked,
814 &clock,
815 &config,
816 &slot,
817 &refill,
818 health,
819 counters,
820 &mut shutdown,
821 )
822 .await
823 {
824 Consolidation::ShutdownObserved => break,
825 Consolidation::Installed | Consolidation::KeptServing => continue,
826 Consolidation::AcquireInstead => {}
830 }
831 }
832 let funding_attempt = slot.funding_attempt();
833 counters.acquire_pending.store(true, Ordering::Release);
834 let acquire = tokio::time::timeout(
835 config.store_call_timeout,
836 allocator.acquire(config.account, config.target_grant, config.lease_ttl, now),
837 );
838 let acquired = tokio::select! {
839 outcome = acquire => outcome,
840 _ = shutdown.changed() => {
841 tracing::debug!("shutdown observed during an acquire; abandoning the refill");
842 break;
843 }
844 };
845 counters.acquire_pending.store(false, Ordering::Release);
846 match acquired {
847 Ok(Ok(allocation)) => {
848 counters.record_acquired(allocation.grant.units);
849 let fresh = install_lease(allocation.grant, &config, &slot, &refill);
853 park(
854 funding_attempt.granted(fresh, allocation.funding),
855 &mut parked,
856 );
857 }
858 outcome => {
867 let serving = slot.load_observed().is_some();
868 let reason: &dyn std::fmt::Display = match &outcome {
872 Ok(Err(error)) => {
873 if let Some(evidence) = refusal_evidence(error) {
874 funding_attempt.shortfall(evidence);
875 }
876 counters.record_acquire_refused(error);
877 error
878 }
879 _ => {
880 counters.record_acquire_timeout();
881 &"allocator timed out"
882 }
883 };
884 if serving {
885 tracing::debug!(%reason, "lease acquire refused; still serving");
886 } else {
887 tracing::warn!(
888 %reason,
889 "lease acquire refused with an empty slot; requests are denied"
890 );
891 }
892 }
893 }
894 }
895
896 if let Some(lease) = slot.take() {
915 parked.push(lease);
916 }
917 let deadline = shutdown_deadline.within(config.shutdown_release_deadline);
918 let mut report = LeaseManagerReport {
919 released: 0,
920 abandoned: 0,
921 task_died: false,
922 };
923 for lease in parked {
924 let grant = lease.grant();
925 if !wait_for_quiescence(&lease, deadline).await {
931 report.abandoned += 1;
932 counters.record_abandoned();
933 tracing::warn!(
934 lease = %grant.lease_id,
935 units = lease.remaining().get(),
936 "lease still held by an in-flight request at the shutdown \
937 deadline; abandoned rather than released, because releasing \
938 units that may still be spent cannot be undone"
939 );
940 continue;
941 }
942 let call_deadline = deadline.min(tokio::time::Instant::now() + config.store_call_timeout);
946 match tokio::time::timeout_at(
947 call_deadline,
948 allocator.release(
949 grant.lease_id,
950 grant.fencing_token,
951 lease.remaining(),
952 clock.now(),
953 ),
954 )
955 .await
956 {
957 Ok(Ok(())) => {
958 report.released += 1;
959 counters.record_released();
960 }
961 Ok(Err(error)) => match release_failure(&error) {
962 ReleaseFailure::Settled | ReleaseFailure::Fenced => {
963 report.released += 1;
964 counters.record_released();
965 tracing::warn!(lease = %grant.lease_id, %error,
966 "lease no longer belongs to this manager; considered settled");
967 }
968 ReleaseFailure::Retryable => {
969 report.abandoned += 1;
970 counters.record_abandoned();
971 tracing::warn!(lease = %grant.lease_id, %error,
972 "release unconfirmed at shutdown; grant requires TTL reclaim");
973 }
974 ReleaseFailure::Integrity => {
975 report.abandoned += 1;
976 counters.record_abandoned();
977 counters.record_integrity_fault(health);
978 tracing::error!(lease = %grant.lease_id, %error,
979 "release violated the allocator contract; grant abandoned and integrity fault retained");
980 }
981 },
982 Err(_) => {
983 report.abandoned += 1;
984 counters.record_abandoned();
985 tracing::warn!(
986 lease = %grant.lease_id,
987 units = lease.remaining().get(),
988 "shutdown budget expired with a lease unreturned; \
989 its units settle at TTL reclaim"
990 );
991 }
992 }
993 }
994 report
995}
996
997#[derive(Clone, Copy)]
1000enum ReleaseFailure {
1001 Settled,
1002 Fenced,
1003 Retryable,
1004 Integrity,
1005}
1006fn release_failure(error: &AllocateError) -> ReleaseFailure {
1007 match error {
1008 AllocateError::UnknownLease | AllocateError::LeaseNotActive => ReleaseFailure::Settled,
1009 AllocateError::Fenced => ReleaseFailure::Fenced,
1010 AllocateError::Storage(_) => ReleaseFailure::Retryable,
1011 AllocateError::InvalidRelease
1012 | AllocateError::UnknownAccount
1013 | AllocateError::AccountInactive
1014 | AllocateError::InsufficientBalance
1015 | AllocateError::BalanceExhausted(_)
1016 | AllocateError::BalanceInsufficient(_)
1017 | AllocateError::BalanceOverflow
1018 | AllocateError::InvalidTtl => ReleaseFailure::Integrity,
1019 }
1020}
1021
1022fn install_lease(
1037 grant: LeaseGrant,
1038 config: &LeaseManagerConfig,
1039 slot: &Arc<LeaseSlot>,
1040 refill: &Arc<RefillRequests>,
1041) -> Arc<LocalLease> {
1042 let low_water = CostUnits(
1043 config
1044 .low_water
1045 .get()
1046 .min(grant.units.get().saturating_sub(1)),
1047 );
1048 Arc::new(
1049 LocalLease::with_sharding(
1050 grant,
1051 low_water,
1052 config.expiry_safety_margin,
1053 slot.sharding(),
1054 )
1055 .with_refill(Arc::clone(refill) as Arc<dyn RefillSignal>),
1056 )
1057}
1058
1059#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1061enum Rotation {
1062 Idle,
1064 Acquire,
1068 Consolidate,
1071}
1072
1073#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1081enum ConsolidationFailure {
1082 RolledBack,
1085 Settled,
1087 Ambiguous,
1091 Integrity,
1096}
1097
1098fn consolidation_failure(error: &AllocateError) -> ConsolidationFailure {
1099 match error {
1100 AllocateError::UnknownLease | AllocateError::LeaseNotActive | AllocateError::Fenced => {
1101 ConsolidationFailure::Settled
1102 }
1103 AllocateError::InsufficientBalance
1108 | AllocateError::BalanceExhausted(_)
1109 | AllocateError::BalanceInsufficient(_)
1110 | AllocateError::UnknownAccount
1111 | AllocateError::AccountInactive
1112 | AllocateError::InvalidTtl => ConsolidationFailure::RolledBack,
1113 AllocateError::InvalidRelease | AllocateError::BalanceOverflow => {
1114 ConsolidationFailure::Integrity
1115 }
1116 AllocateError::Storage(_) => ConsolidationFailure::Ambiguous,
1117 }
1118}
1119
1120#[allow(
1130 clippy::too_many_arguments,
1131 reason = "one consolidation step over the loop's own borrowed state"
1132)]
1133async fn consolidate_live_lease(
1134 allocator: &Arc<dyn LeaseAllocator>,
1135 parked: &mut Vec<Arc<LocalLease>>,
1136 clock: &Arc<dyn Clock>,
1137 config: &LeaseManagerConfig,
1138 slot: &Arc<LeaseSlot>,
1139 refill: &Arc<RefillRequests>,
1140 health: &watch::Sender<bool>,
1141 counters: &LeaseCounters,
1142 shutdown: &mut watch::Receiver<bool>,
1143) -> Consolidation {
1144 let Some(live) = slot.take() else {
1145 return Consolidation::AcquireInstead;
1148 };
1149 if Arc::strong_count(&live) > 1 || !live.is_only_local_view() {
1150 counters.record_consolidation_deferred();
1151 publish_and_park(slot, live, parked);
1152 return Consolidation::KeptServing;
1153 }
1154 let unspent = live.remaining();
1157 let needed = live.largest_refused_quote();
1159 let grant = *live.grant();
1160 let funding_attempt = slot.funding_attempt();
1161 counters.acquire_pending.store(true, Ordering::Release);
1162 let call = tokio::time::timeout(
1163 config.store_call_timeout,
1164 allocator.consolidate(
1165 grant.lease_id,
1166 grant.fencing_token,
1167 unspent,
1168 config.target_grant,
1169 needed,
1170 config.lease_ttl,
1171 clock.now(),
1172 ),
1173 );
1174 let outcome = tokio::select! {
1175 outcome = call => outcome,
1176 _ = shutdown.changed() => {
1177 tracing::debug!(
1178 lease = %grant.lease_id,
1179 "shutdown observed during a consolidation; parking the grant"
1180 );
1181 parked.push(live);
1186 return Consolidation::ShutdownObserved;
1187 }
1188 };
1189 counters.acquire_pending.store(false, Ordering::Release);
1190
1191 let error: AllocateError = match outcome {
1192 Ok(Ok(allocation)) => {
1193 let fresh = allocation.grant;
1194 counters.record_consolidated(fresh.units);
1195 tracing::debug!(
1196 superseded = %grant.lease_id,
1197 lease = %fresh.lease_id,
1198 folded = unspent.get(),
1199 units = fresh.units.get(),
1200 "consolidated a refused lease into a larger grant"
1201 );
1202 drop(live);
1206 let fresh = install_lease(fresh, config, slot, refill);
1207 park(funding_attempt.granted(fresh, allocation.funding), parked);
1208 return Consolidation::Installed;
1209 }
1210 Ok(Err(error)) => error,
1211 Err(_) => {
1212 counters.record_acquire_timeout();
1213 tracing::warn!(
1214 lease = %grant.lease_id,
1215 "consolidation timed out; parking the grant because the store may have settled it"
1216 );
1217 parked.push(live);
1218 return Consolidation::KeptServing;
1219 }
1220 };
1221 if let Some(evidence) = refusal_evidence(&error) {
1222 funding_attempt.shortfall(evidence);
1223 }
1224 counters.record_acquire_refused(&error);
1225 match consolidation_failure(&error) {
1226 ConsolidationFailure::RolledBack => {
1227 tracing::debug!(
1228 lease = %grant.lease_id, %error,
1229 "consolidation refused; the grant is untouched and keeps serving"
1230 );
1231 publish_and_park(slot, live, parked);
1232 Consolidation::KeptServing
1233 }
1234 ConsolidationFailure::Integrity => {
1235 counters.record_integrity_fault(health);
1236 tracing::error!(
1237 lease = %grant.lease_id, units = unspent.get(), %error,
1238 "consolidation violated the allocator contract; readiness withdrawn"
1239 );
1240 publish_and_park(slot, live, parked);
1241 Consolidation::KeptServing
1242 }
1243 ConsolidationFailure::Settled => {
1244 counters.record_released();
1245 tracing::warn!(
1246 lease = %grant.lease_id, %error,
1247 "the store no longer holds this grant open; acquiring a fresh one"
1248 );
1249 drop(live);
1250 Consolidation::AcquireInstead
1251 }
1252 ConsolidationFailure::Ambiguous => {
1253 tracing::warn!(
1254 lease = %grant.lease_id, %error,
1255 "consolidation failed without saying whether it committed; parking the grant"
1256 );
1257 parked.push(live);
1258 Consolidation::KeptServing
1259 }
1260 }
1261}
1262
1263fn publish_and_park(slot: &LeaseSlot, lease: Arc<LocalLease>, parked: &mut Vec<Arc<LocalLease>>) {
1267 park(slot.replace(lease), parked);
1268}
1269
1270fn park(previous: Option<Arc<LocalLease>>, parked: &mut Vec<Arc<LocalLease>>) {
1271 if let Some(previous) = previous {
1272 parked.push(previous);
1273 }
1274}
1275
1276fn refusal_evidence(error: &AllocateError) -> Option<tollgate_core::BalanceShortfall> {
1279 match error {
1280 AllocateError::BalanceExhausted(evidence) => Some((*evidence).into()),
1281 AllocateError::BalanceInsufficient(evidence) => Some(*evidence),
1282 _ => None,
1283 }
1284}
1285
1286#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1288enum Consolidation {
1289 Installed,
1291 KeptServing,
1294 AcquireInstead,
1296 ShutdownObserved,
1297}
1298
1299async fn wait_for_quiescence(lease: &Arc<LocalLease>, deadline: tokio::time::Instant) -> bool {
1311 const POLL: std::time::Duration = std::time::Duration::from_millis(5);
1312 loop {
1313 if Arc::strong_count(lease) == 1 && lease.is_only_local_view() {
1314 return true;
1315 }
1316 if tokio::time::Instant::now() >= deadline {
1317 return false;
1318 }
1319 tokio::time::sleep(POLL.min(deadline - tokio::time::Instant::now())).await;
1320 }
1321}
1322
1323#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1328enum ReleasePass {
1329 Complete,
1331 BudgetExpired,
1334 ShutdownObserved,
1337}
1338
1339#[allow(
1366 clippy::too_many_arguments,
1367 reason = "one release step over the loop's own borrowed state"
1368)]
1369async fn release_quiesced(
1370 allocator: &Arc<dyn LeaseAllocator>,
1371 parked: &mut Vec<Arc<LocalLease>>,
1372 clock: &Arc<dyn Clock>,
1373 config: &LeaseManagerConfig,
1374 slot: &Arc<LeaseSlot>,
1375 health: &watch::Sender<bool>,
1376 counters: &LeaseCounters,
1377 shutdown: &mut watch::Receiver<bool>,
1378 deadline: tokio::time::Instant,
1379) -> ReleasePass {
1380 let mut queue = std::mem::take(parked).into_iter();
1387 let mut retry = Vec::new();
1388 while let Some(lease) = queue.next() {
1389 if tokio::time::Instant::now() >= deadline {
1390 *parked = std::iter::once(lease).chain(queue).chain(retry).collect();
1392 return ReleasePass::BudgetExpired;
1393 }
1394 if Arc::strong_count(&lease) > 1 || !lease.is_only_local_view() {
1399 retry.push(lease);
1400 continue;
1401 }
1402 let grant = lease.grant();
1403 let call_deadline = deadline.min(tokio::time::Instant::now() + config.store_call_timeout);
1407 let call = tokio::time::timeout_at(
1408 call_deadline,
1409 allocator.release(
1410 grant.lease_id,
1411 grant.fencing_token,
1412 lease.remaining(),
1413 clock.now(),
1414 ),
1415 );
1416 let lease_id = grant.lease_id;
1417 let call_outcome = tokio::select! {
1424 outcome = call => outcome,
1425 _ = shutdown.changed() => {
1426 tracing::debug!(
1427 lease = %lease_id,
1428 "shutdown observed during a release; entering the shutdown release phase"
1429 );
1430 *parked = std::iter::once(lease).chain(queue).chain(retry).collect();
1431 return ReleasePass::ShutdownObserved;
1432 }
1433 };
1434 match call_outcome {
1435 Err(_) => {
1436 tracing::debug!(lease = %lease_id, "release timed out; retrying next tick");
1437 retry.push(lease);
1438 }
1439 Ok(Ok(())) => counters.record_released(),
1440 Ok(Err(error)) => match release_failure(&error) {
1441 ReleaseFailure::Retryable => {
1442 tracing::warn!(lease = %lease_id, %error, "release failed; retrying next tick");
1443 retry.push(lease);
1444 }
1445 ReleaseFailure::Settled => {
1446 counters.record_released();
1447 tracing::debug!(lease = %lease_id, %error, "lease was already settled");
1448 }
1449 ReleaseFailure::Fenced => {
1450 counters.record_released();
1451 tracing::warn!(lease = %lease_id,
1452 "store rejected this lease's capability; clearing the slot so this instance stops serving");
1453 if let Some(current) = slot.take() {
1457 retry.push(current);
1458 }
1459 }
1460 ReleaseFailure::Integrity => {
1461 counters.record_abandoned();
1462 counters.record_integrity_fault(health);
1463 tracing::error!(lease = %lease_id, units = lease.remaining().get(), %error,
1464 "release violated the allocator contract; grant abandoned and readiness withdrawn");
1465 }
1466 },
1467 }
1468 }
1469 *parked = retry;
1472 ReleasePass::Complete
1473}
1474
1475#[cfg(test)]
1476mod tests {
1477 #[test]
1478 fn aggregate_refill_counters_reject_overflow_in_fields_and_totals() {
1479 let mut left = super::LeaseStats::ZERO;
1480 left.acquired_units = u64::MAX;
1481 let mut right = super::LeaseStats::ZERO;
1482 right.acquired_units = 1;
1483 assert!(left.checked_add(right).is_none());
1484 left = super::LeaseStats::ZERO;
1485 right = super::LeaseStats::ZERO;
1486 left.acquire_refused[0] = u64::MAX;
1487 right.acquire_refused[1] = 1;
1488 assert!(left.checked_add(right).is_none());
1489 right.acquire_refused[1] = 0;
1490 assert_eq!(left.checked_add(right).unwrap().refused(), u64::MAX);
1491 left = super::LeaseStats::ZERO;
1492 right = super::LeaseStats::ZERO;
1493 left.uncertain_acquires = u64::MAX;
1494 right.uncertain_acquires = 1;
1495 assert!(left.checked_add(right).is_none());
1496 }
1497 use std::collections::HashMap;
1498 use std::sync::Mutex;
1499
1500 use async_trait::async_trait;
1501 use jiff::Timestamp;
1502 use tollgate_core::{FencingToken, LeaseGrant, LeaseId};
1503 use tollgate_store::{AllocateError, Allocation, ReclaimBatch, StoreError, SystemClock};
1504
1505 use super::*;
1506
1507 #[test]
1508 fn only_ambiguous_acquire_outcomes_increase_uncertainty() {
1509 let counters = LeaseCounters::new();
1510 for error in [
1511 AllocateError::UnknownAccount,
1512 AllocateError::AccountInactive,
1513 AllocateError::InsufficientBalance,
1514 AllocateError::UnknownLease,
1515 AllocateError::Fenced,
1516 AllocateError::LeaseNotActive,
1517 AllocateError::InvalidRelease,
1518 AllocateError::InvalidTtl,
1519 AllocateError::BalanceOverflow,
1520 ] {
1521 counters.record_acquire_refused(&error);
1522 }
1523 assert_eq!(counters.snapshot().uncertain_acquires, 0);
1524 counters.record_acquire_refused(&AllocateError::Storage(StoreError("lost reply".into())));
1525 assert_eq!(counters.snapshot().uncertain_acquires, 1);
1526 counters.record_acquire_timeout();
1527 assert_eq!(counters.snapshot().uncertain_acquires, 2);
1528 }
1529
1530 #[tokio::test]
1537 async fn a_request_raised_before_the_wait_is_not_lost() {
1538 let refill = RefillRequests::default();
1539 refill.request_refill();
1540
1541 tokio::time::timeout(std::time::Duration::from_secs(5), refill.requested())
1542 .await
1543 .expect("a request raised before the wait must complete it");
1544 }
1545
1546 #[tokio::test]
1549 async fn a_request_raised_during_the_wait_wakes_it() {
1550 let refill = Arc::new(RefillRequests::default());
1551 let signal = Arc::clone(&refill);
1552 tokio::spawn(async move {
1553 tokio::task::yield_now().await;
1554 signal.request_refill();
1555 });
1556
1557 tokio::time::timeout(std::time::Duration::from_secs(5), refill.requested())
1558 .await
1559 .expect("a request raised while waiting must wake the waiter");
1560 }
1561
1562 #[tokio::test(start_paused = true)]
1566 async fn a_request_is_consumed_by_the_wait_it_completes() {
1567 let refill = RefillRequests::default();
1568 refill.request_refill();
1569 refill.requested().await;
1570
1571 assert!(
1572 tokio::time::timeout(std::time::Duration::from_secs(5), refill.requested())
1573 .await
1574 .is_err(),
1575 "the request was already answered; waiting again must block"
1576 );
1577 }
1578
1579 #[tokio::test(start_paused = true)]
1582 async fn repeated_requests_collapse_into_one_wake() {
1583 let refill = RefillRequests::default();
1584 for _ in 0..10 {
1585 refill.request_refill();
1586 }
1587 refill.requested().await;
1588
1589 assert!(
1590 tokio::time::timeout(std::time::Duration::from_secs(5), refill.requested())
1591 .await
1592 .is_err(),
1593 "ten requests are still one outstanding refill, not ten"
1594 );
1595 }
1596
1597 #[derive(Clone, Copy, PartialEq, Eq)]
1599 enum Refusal {
1600 Storage,
1601 Fenced,
1602 InvalidRelease,
1603 LeaseNotActive,
1604 UnknownLease,
1605 Hang,
1607 }
1608
1609 struct ScriptedAllocator {
1612 released: Mutex<Vec<LeaseId>>,
1613 refusals: HashMap<LeaseId, Refusal>,
1614 }
1615
1616 impl ScriptedAllocator {
1617 fn new(refusals: impl IntoIterator<Item = (LeaseId, Refusal)>) -> Arc<Self> {
1618 Arc::new(Self {
1619 released: Mutex::new(Vec::new()),
1620 refusals: refusals.into_iter().collect(),
1621 })
1622 }
1623
1624 fn released(&self) -> Vec<LeaseId> {
1625 self.released.lock().unwrap().clone()
1626 }
1627 }
1628
1629 #[async_trait]
1630 impl LeaseAllocator for ScriptedAllocator {
1631 async fn acquire(
1632 &self,
1633 _account: AccountId,
1634 _requested: CostUnits,
1635 _ttl: SignedDuration,
1636 _now: Timestamp,
1637 ) -> Result<Allocation, AllocateError> {
1638 unreachable!("release_quiesced never acquires")
1639 }
1640
1641 async fn release(
1642 &self,
1643 lease_id: LeaseId,
1644 _fencing_token: FencingToken,
1645 _unspent: CostUnits,
1646 _now: Timestamp,
1647 ) -> Result<(), AllocateError> {
1648 match self.refusals.get(&lease_id) {
1649 Some(Refusal::Storage) => {
1650 Err(AllocateError::Storage(StoreError("scripted outage".into())))
1651 }
1652 Some(Refusal::Fenced) => Err(AllocateError::Fenced),
1653 Some(Refusal::InvalidRelease) => Err(AllocateError::InvalidRelease),
1654 Some(Refusal::LeaseNotActive) => Err(AllocateError::LeaseNotActive),
1655 Some(Refusal::UnknownLease) => Err(AllocateError::UnknownLease),
1656 Some(Refusal::Hang) => std::future::pending().await,
1657 None => {
1658 self.released.lock().unwrap().push(lease_id);
1659 Ok(())
1660 }
1661 }
1662 }
1663
1664 async fn consolidate(
1665 &self,
1666 _lease_id: LeaseId,
1667 _fencing_token: FencingToken,
1668 _unspent: CostUnits,
1669 _requested: CostUnits,
1670 _needed: CostUnits,
1671 _ttl: SignedDuration,
1672 _now: Timestamp,
1673 ) -> Result<Allocation, AllocateError> {
1674 unreachable!("release_quiesced never consolidates")
1675 }
1676
1677 async fn reclaim_expired_batch(
1678 &self,
1679 _now: Timestamp,
1680 _limit: std::num::NonZeroUsize,
1681 ) -> Result<ReclaimBatch, StoreError> {
1682 unreachable!("release_quiesced never reclaims")
1683 }
1684 }
1685
1686 struct ConsolidatingAllocator {
1690 answer: Mutex<Option<Result<Allocation, AllocateError>>>,
1691 calls: AtomicU64,
1692 needed: AtomicU64,
1693 publish_during_call: Option<Arc<LeaseSlot>>,
1694 }
1695
1696 impl ConsolidatingAllocator {
1697 fn new(answer: Result<Allocation, AllocateError>) -> Arc<Self> {
1698 Arc::new(Self {
1699 answer: Mutex::new(Some(answer)),
1700 calls: AtomicU64::new(0),
1701 needed: AtomicU64::new(0),
1702 publish_during_call: None,
1703 })
1704 }
1705 }
1706
1707 #[async_trait]
1708 impl LeaseAllocator for ConsolidatingAllocator {
1709 async fn acquire(
1710 &self,
1711 _account: AccountId,
1712 _requested: CostUnits,
1713 _ttl: SignedDuration,
1714 _now: Timestamp,
1715 ) -> Result<Allocation, AllocateError> {
1716 unreachable!("these tests drive consolidation directly")
1717 }
1718
1719 async fn release(
1720 &self,
1721 _lease_id: LeaseId,
1722 _fencing_token: FencingToken,
1723 _unspent: CostUnits,
1724 _now: Timestamp,
1725 ) -> Result<(), AllocateError> {
1726 Ok(())
1727 }
1728
1729 async fn consolidate(
1730 &self,
1731 _lease_id: LeaseId,
1732 _fencing_token: FencingToken,
1733 _unspent: CostUnits,
1734 _requested: CostUnits,
1735 needed: CostUnits,
1736 _ttl: SignedDuration,
1737 _now: Timestamp,
1738 ) -> Result<Allocation, AllocateError> {
1739 self.calls.fetch_add(1, Ordering::SeqCst);
1740 self.needed.store(needed.get(), Ordering::SeqCst);
1741 if let Some(slot) = &self.publish_during_call {
1742 assert!(slot.replace(parked_lease(99)).is_none());
1743 }
1744 self.answer
1745 .lock()
1746 .unwrap()
1747 .take()
1748 .expect("one scripted consolidation per test")
1749 }
1750
1751 async fn reclaim_expired_batch(
1752 &self,
1753 _now: Timestamp,
1754 _limit: std::num::NonZeroUsize,
1755 ) -> Result<ReclaimBatch, StoreError> {
1756 unreachable!("these tests never reclaim")
1757 }
1758 }
1759
1760 fn grant(id: u128, units: u64) -> LeaseGrant {
1761 LeaseGrant {
1762 lease_id: LeaseId(id),
1763 account_id: AccountId(1),
1764 fencing_token: FencingToken(1),
1765 units: CostUnits(units),
1766 expires_at: Timestamp::from_second(3_600).unwrap(),
1767 }
1768 }
1769
1770 fn allocation(id: u128, units: u64) -> Allocation {
1771 Allocation {
1772 grant: grant(id, units),
1773 funding: None,
1774 }
1775 }
1776
1777 async fn consolidate_once(
1780 allocator: Arc<dyn LeaseAllocator>,
1781 harness: &Harness,
1782 lease: Arc<LocalLease>,
1783 ) -> (Consolidation, Option<u128>, Vec<u128>) {
1784 drop(harness.slot.replace(lease));
1785 let mut parked = Vec::new();
1786 let (_tx, mut shutdown) = watch::channel(false);
1787 let outcome = consolidate_live_lease(
1788 &allocator,
1789 &mut parked,
1790 &clock(),
1791 &harness.config,
1792 &harness.slot,
1793 &Arc::new(RefillRequests::default()),
1794 &harness.health,
1795 &harness.counters,
1796 &mut shutdown,
1797 )
1798 .await;
1799 let served = harness.slot.load().map(|l| l.grant().lease_id.0);
1800 let parked_ids = parked.iter().map(|l| l.grant().lease_id.0).collect();
1801 (outcome, served, parked_ids)
1802 }
1803
1804 #[tokio::test(start_paused = true)]
1808 async fn a_successful_consolidation_installs_the_grant_and_parks_nothing() {
1809 let harness = Harness::new();
1810 let allocator = ConsolidatingAllocator::new(Ok(allocation(2, 500)));
1811 let (outcome, served, parked) = consolidate_once(
1812 Arc::clone(&allocator) as Arc<dyn LeaseAllocator>,
1813 &harness,
1814 parked_lease(1),
1815 )
1816 .await;
1817
1818 assert_eq!(outcome, Consolidation::Installed);
1819 assert_eq!(served, Some(2), "the larger grant is serving");
1820 assert!(
1821 parked.is_empty(),
1822 "the superseded lease was already settled"
1823 );
1824 let stats = harness.stats();
1825 assert_eq!(stats.consolidated, 1);
1826 assert_eq!(stats.acquired, 1);
1827 assert_eq!(stats.released, 1, "consolidation settled its predecessor");
1828 assert!(!harness.counters.acquire_pending());
1829 assert_eq!(stats.acquired_units, 500);
1830 }
1831
1832 #[tokio::test(start_paused = true)]
1835 async fn consolidation_carries_the_refused_quote() {
1836 for refusals in [&[][..], &[150, 252, 101][..]] {
1837 let harness = Harness::new();
1838 let allocator = ConsolidatingAllocator::new(Ok(allocation(2, 500)));
1839 let lease = parked_lease(1);
1840 for "e in refusals {
1841 assert!(
1842 lease
1843 .try_debit(CostUnits(quote), Timestamp::UNIX_EPOCH)
1844 .is_err()
1845 );
1846 }
1847 let (outcome, _, _) = consolidate_once(
1848 Arc::clone(&allocator) as Arc<dyn LeaseAllocator>,
1849 &harness,
1850 lease,
1851 )
1852 .await;
1853 assert_eq!(outcome, Consolidation::Installed);
1854 assert_eq!(
1855 allocator.needed.load(Ordering::SeqCst),
1856 refusals.iter().copied().max().unwrap_or(0)
1857 );
1858 }
1859 }
1860
1861 #[tokio::test(start_paused = true)]
1865 async fn a_consolidated_tail_publishes_the_accounts_remaining_funding() {
1866 let harness = Harness::new();
1867 let mut tail = grant(2, 1);
1868 tail.fencing_token = FencingToken(2);
1869 let allocator = ConsolidatingAllocator::new(Ok(Allocation {
1870 grant: tail,
1871 funding: Some(tollgate_core::BalanceShortfall {
1872 remaining: CostUnits(1),
1873 period_end: None,
1874 }),
1875 }));
1876 let (outcome, served, _) = consolidate_once(allocator, &harness, parked_lease(1)).await;
1877 assert_eq!(outcome, Consolidation::Installed);
1878 assert_eq!(served, Some(2));
1879 assert_eq!(
1880 harness.slot.funding_evidence(Timestamp::UNIX_EPOCH),
1881 Some(CostUnits(1))
1882 );
1883 }
1884
1885 #[tokio::test(start_paused = true)]
1886 async fn consolidation_retains_a_grant_published_during_the_store_call() {
1887 for (answer, expected_outcome, expected_served) in [
1888 (Ok(allocation(2, 500)), Consolidation::Installed, 2),
1889 (
1890 Err(AllocateError::InsufficientBalance),
1891 Consolidation::KeptServing,
1892 1,
1893 ),
1894 (
1895 Err(AllocateError::InvalidRelease),
1896 Consolidation::KeptServing,
1897 1,
1898 ),
1899 ] {
1900 let harness = Harness::new();
1901 let mut allocator = ConsolidatingAllocator::new(answer);
1902 Arc::get_mut(&mut allocator).unwrap().publish_during_call =
1903 Some(Arc::clone(&harness.slot));
1904 let (outcome, served, parked) =
1905 consolidate_once(allocator, &harness, parked_lease(1)).await;
1906 assert_eq!(outcome, expected_outcome);
1907 assert_eq!(served, Some(expected_served));
1908 assert_eq!(
1909 parked,
1910 vec![99],
1911 "a concurrent publisher's grant must remain available for quiesced release"
1912 );
1913 }
1914 }
1915
1916 #[tokio::test(start_paused = true)]
1920 async fn a_rolled_back_consolidation_returns_the_lease_to_the_slot() {
1921 for error in [
1922 AllocateError::InsufficientBalance,
1923 AllocateError::BalanceExhausted(tollgate_core::BalanceExhaustion { period_end: None }),
1924 AllocateError::BalanceInsufficient(tollgate_core::BalanceShortfall {
1925 remaining: CostUnits(3),
1926 period_end: None,
1927 }),
1928 AllocateError::AccountInactive,
1929 AllocateError::UnknownAccount,
1930 ] {
1931 let harness = Harness::new();
1932 let allocator = ConsolidatingAllocator::new(Err(error.clone()));
1933 let (outcome, served, parked) = consolidate_once(
1934 Arc::clone(&allocator) as Arc<dyn LeaseAllocator>,
1935 &harness,
1936 parked_lease(1),
1937 )
1938 .await;
1939
1940 assert_eq!(outcome, Consolidation::KeptServing, "{error}");
1941 assert_eq!(
1942 served,
1943 Some(1),
1944 "still serving the untouched lease: {error}"
1945 );
1946 assert!(parked.is_empty(), "{error}");
1947 assert_eq!(
1948 harness.slot.funding_evidence(Timestamp::UNIX_EPOCH),
1949 match error {
1950 AllocateError::BalanceExhausted(_) => Some(CostUnits::ZERO),
1951 AllocateError::BalanceInsufficient(evidence) => Some(evidence.remaining),
1952 _ => None,
1953 },
1954 "{error}"
1955 );
1956 assert!(harness.is_healthy(), "an ordinary refusal is not a fault");
1957 }
1958 }
1959
1960 #[tokio::test(start_paused = true)]
1964 async fn an_ambiguous_consolidation_parks_the_grant_rather_than_reinstating_it() {
1965 let harness = Harness::new();
1966 let allocator = ConsolidatingAllocator::new(Err(AllocateError::Storage(StoreError(
1967 "connection reset".into(),
1968 ))));
1969 let (outcome, served, parked) = consolidate_once(
1970 Arc::clone(&allocator) as Arc<dyn LeaseAllocator>,
1971 &harness,
1972 parked_lease(1),
1973 )
1974 .await;
1975
1976 assert_eq!(outcome, Consolidation::KeptServing);
1977 assert_eq!(
1978 served, None,
1979 "the slot fails closed rather than double-spending"
1980 );
1981 assert_eq!(parked, vec![1], "and the release pass settles it");
1982 }
1983
1984 #[tokio::test(start_paused = true)]
1988 async fn a_settled_lease_falls_through_to_an_ordinary_acquire() {
1989 for error in [
1990 AllocateError::UnknownLease,
1991 AllocateError::LeaseNotActive,
1992 AllocateError::Fenced,
1993 ] {
1994 let harness = Harness::new();
1995 let allocator = ConsolidatingAllocator::new(Err(error.clone()));
1996 let (outcome, served, parked) = consolidate_once(
1997 Arc::clone(&allocator) as Arc<dyn LeaseAllocator>,
1998 &harness,
1999 parked_lease(1),
2000 )
2001 .await;
2002
2003 assert_eq!(outcome, Consolidation::AcquireInstead, "{error}");
2004 assert_eq!(served, None, "{error}");
2005 assert!(parked.is_empty(), "the store already settled it: {error}");
2006 }
2007 }
2008
2009 #[tokio::test(start_paused = true)]
2013 async fn an_over_claimed_fold_withdraws_readiness_and_keeps_serving() {
2014 let harness = Harness::new();
2015 let allocator = ConsolidatingAllocator::new(Err(AllocateError::InvalidRelease));
2016 let (outcome, served, parked) = consolidate_once(
2017 Arc::clone(&allocator) as Arc<dyn LeaseAllocator>,
2018 &harness,
2019 parked_lease(1),
2020 )
2021 .await;
2022
2023 assert_eq!(outcome, Consolidation::KeptServing);
2024 assert_eq!(served, Some(1));
2025 assert!(parked.is_empty());
2026 assert!(!harness.is_healthy(), "an accounting fault is never silent");
2027 }
2028
2029 #[tokio::test(start_paused = true)]
2033 async fn a_consolidation_defers_while_a_reservation_is_in_flight() {
2034 let harness = Harness::new();
2035 let allocator = ConsolidatingAllocator::new(Ok(allocation(2, 500)));
2036 let lease = parked_lease(1);
2037 let _in_flight = Arc::clone(&lease);
2039
2040 let (outcome, served, parked) = consolidate_once(
2041 Arc::clone(&allocator) as Arc<dyn LeaseAllocator>,
2042 &harness,
2043 lease,
2044 )
2045 .await;
2046
2047 assert_eq!(outcome, Consolidation::KeptServing);
2048 assert_eq!(served, Some(1), "the slot keeps serving it");
2049 assert!(parked.is_empty());
2050 assert_eq!(
2051 allocator.calls.load(Ordering::SeqCst),
2052 0,
2053 "no store call is made against an inexact aggregate"
2054 );
2055 assert_eq!(harness.stats().consolidations_deferred, 1);
2056 }
2057
2058 struct Harness {
2061 config: LeaseManagerConfig,
2062 slot: Arc<LeaseSlot>,
2063 health: watch::Sender<bool>,
2064 healthy: watch::Receiver<bool>,
2065 counters: LeaseCounters,
2066 }
2067
2068 impl Harness {
2069 fn new() -> Self {
2070 let (health, healthy) = watch::channel(true);
2071 Self {
2072 counters: LeaseCounters::new(),
2073 config: LeaseManagerConfig {
2074 account: AccountId(1),
2075 target_grant: CostUnits(1_000),
2076 low_water: CostUnits(250),
2077 lease_ttl: SignedDuration::from_secs(60),
2078 expiry_safety_margin: SignedDuration::ZERO,
2079 poll_interval: std::time::Duration::from_millis(5),
2080 store_call_timeout: std::time::Duration::from_millis(50),
2081 shutdown_release_deadline: std::time::Duration::from_secs(10),
2082 },
2083 slot: LeaseSlot::for_account(AccountId(1)),
2084 health,
2085 healthy,
2086 }
2087 }
2088
2089 async fn release_quiesced(
2092 &self,
2093 allocator: &Arc<dyn LeaseAllocator>,
2094 parked: &mut Vec<Arc<LocalLease>>,
2095 ) -> ReleasePass {
2096 self.release_pass(allocator, parked, std::time::Duration::from_secs(3_600))
2097 .await
2098 }
2099
2100 async fn release_pass(
2101 &self,
2102 allocator: &Arc<dyn LeaseAllocator>,
2103 parked: &mut Vec<Arc<LocalLease>>,
2104 budget: std::time::Duration,
2105 ) -> ReleasePass {
2106 let (_tx, mut shutdown) = watch::channel(false);
2107 release_quiesced(
2108 allocator,
2109 parked,
2110 &clock(),
2111 &self.config,
2112 &self.slot,
2113 &self.health,
2114 &self.counters,
2115 &mut shutdown,
2116 tokio::time::Instant::now() + budget,
2117 )
2118 .await
2119 }
2120
2121 fn stats(&self) -> LeaseStats {
2122 self.counters.snapshot()
2123 }
2124
2125 fn is_healthy(&self) -> bool {
2126 *self.healthy.borrow()
2127 }
2128 }
2129
2130 fn parked_lease(id: u128) -> Arc<LocalLease> {
2131 Arc::new(LocalLease::new(
2132 LeaseGrant {
2133 lease_id: LeaseId(id),
2134 account_id: AccountId(1),
2135 fencing_token: FencingToken(1),
2136 units: CostUnits(100),
2137 expires_at: Timestamp::from_second(3_600).unwrap(),
2138 },
2139 CostUnits(10),
2140 ))
2141 }
2142
2143 fn clock() -> Arc<dyn Clock> {
2144 Arc::new(SystemClock)
2145 }
2146
2147 #[tokio::test(start_paused = true)]
2148 async fn shutdown_distinguishes_settled_leases_from_unconfirmed_or_invalid_releases() {
2149 for (refusal, released, faulted) in [
2150 (Refusal::Storage, 0, false),
2151 (Refusal::Fenced, 1, false),
2152 (Refusal::InvalidRelease, 0, true),
2153 (Refusal::LeaseNotActive, 1, false),
2154 (Refusal::UnknownLease, 1, false),
2155 (Refusal::Hang, 0, false),
2156 ] {
2157 let harness = Harness::new();
2158 drop(harness.slot.replace(parked_lease(1)));
2159 let manager = LeaseManager::spawn(
2160 ScriptedAllocator::new([(LeaseId(1), refusal)]),
2161 harness.slot,
2162 Arc::new(crate::ManualClock::new(Timestamp::from_second(0).unwrap())),
2163 harness.config,
2164 )
2165 .unwrap();
2166 let counters = manager.counters();
2167 let health = manager.health();
2168 let report = manager.shutdown().await;
2169 assert_eq!(report.released, released);
2170 assert_eq!(report.abandoned, 1 - released);
2171 assert_eq!(counters.snapshot().released, released);
2172 assert_eq!(counters.snapshot().abandoned, 1 - released);
2173 assert_eq!(counters.integrity_fault(), faulted);
2174 assert!(
2175 !*health.borrow(),
2176 "every stopped task is unhealthy, including clean stops"
2177 );
2178 }
2179 }
2180
2181 #[tokio::test(start_paused = true)]
2182 async fn runtime_deadline_shortens_the_managers_actual_release_pass() {
2183 let harness = Harness::new();
2184 drop(harness.slot.replace(parked_lease(1)));
2185 let manager = LeaseManager::spawn(
2186 ScriptedAllocator::new([(LeaseId(1), Refusal::Hang)]),
2187 harness.slot,
2188 Arc::new(crate::ManualClock::new(Timestamp::from_second(0).unwrap())),
2189 harness.config,
2190 )
2191 .unwrap();
2192 let began = tokio::time::Instant::now();
2193 manager.stop_at(began + std::time::Duration::from_millis(5));
2194 let report = tokio::time::timeout(std::time::Duration::from_millis(6), manager.shutdown())
2195 .await
2196 .unwrap();
2197 assert_eq!(report.abandoned, 1);
2198 assert!(!report.task_died);
2199 assert_eq!(began.elapsed(), std::time::Duration::from_millis(5));
2200 }
2201
2202 #[tokio::test]
2203 async fn storage_failure_does_not_skip_the_next_parked_lease() {
2204 let scripted = ScriptedAllocator::new([(LeaseId(1), Refusal::Storage)]);
2205 let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
2206 let harness = Harness::new();
2207 let mut parked = vec![parked_lease(1), parked_lease(2)];
2208
2209 harness.release_quiesced(&allocator, &mut parked).await;
2210
2211 assert_eq!(scripted.released(), [LeaseId(2)]);
2212 assert_eq!(parked.len(), 1, "only the failed lease stays parked");
2213 assert_eq!(parked[0].grant().lease_id, LeaseId(1));
2214 }
2215
2216 #[tokio::test]
2217 async fn unquiesced_lease_is_never_released() {
2218 let scripted = ScriptedAllocator::new([]);
2219 let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
2220 let harness = Harness::new();
2221 let lease = parked_lease(7);
2222 let in_flight = Arc::clone(&lease);
2223 let mut parked = vec![lease];
2224
2225 harness.release_quiesced(&allocator, &mut parked).await;
2226 assert!(scripted.released().is_empty());
2227 assert_eq!(parked.len(), 1, "held lease stays parked");
2228 assert_eq!(
2229 harness.stats().released,
2230 0,
2231 "a lease still on our books has not been released"
2232 );
2233
2234 drop(in_flight);
2235 harness.release_quiesced(&allocator, &mut parked).await;
2236 assert_eq!(scripted.released(), [LeaseId(7)]);
2237 assert!(parked.is_empty());
2238 assert_eq!(harness.stats().released, 1);
2239 }
2240
2241 #[tokio::test]
2244 async fn settled_refusals_drop_silently() {
2245 for refusal in [Refusal::LeaseNotActive, Refusal::UnknownLease] {
2246 let scripted = ScriptedAllocator::new([(LeaseId(3), refusal)]);
2247 let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
2248 let harness = Harness::new();
2249 drop(harness.slot.replace(parked_lease(99)));
2250 let mut parked = vec![parked_lease(3)];
2251
2252 harness.release_quiesced(&allocator, &mut parked).await;
2253
2254 assert!(scripted.released().is_empty());
2255 assert!(parked.is_empty(), "settled lease is not retried");
2256 assert!(harness.is_healthy(), "settlement is not a health event");
2257 assert!(harness.slot.load().is_some(), "the slot is untouched");
2258 assert_eq!(
2259 harness.stats().released,
2260 1,
2261 "already settled still means the store no longer holds it"
2262 );
2263 }
2264 }
2265
2266 #[tokio::test]
2270 async fn fenced_release_clears_the_slot() {
2271 let scripted = ScriptedAllocator::new([(LeaseId(3), Refusal::Fenced)]);
2272 let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
2273 let harness = Harness::new();
2274 drop(harness.slot.replace(parked_lease(99)));
2275 let mut parked = vec![parked_lease(3)];
2276
2277 harness.release_quiesced(&allocator, &mut parked).await;
2278
2279 assert!(
2280 harness.slot.load().is_none(),
2281 "an instance with divergent lease identity must stop serving"
2282 );
2283 let ids: Vec<_> = parked.iter().map(|l| l.grant().lease_id).collect();
2293 assert_eq!(
2294 ids,
2295 [LeaseId(99)],
2296 "the fenced lease is not retried, and the slot's live lease is not lost"
2297 );
2298 assert_eq!(
2299 parked[0].remaining(),
2300 CostUnits(100),
2301 "its unspent units are still accounted for"
2302 );
2303 }
2304
2305 #[tokio::test]
2308 async fn invalid_release_fails_readiness() {
2309 let scripted = ScriptedAllocator::new([(LeaseId(3), Refusal::InvalidRelease)]);
2310 let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
2311 let harness = Harness::new();
2312 let mut parked = vec![parked_lease(3)];
2313
2314 assert!(harness.is_healthy());
2315 harness.release_quiesced(&allocator, &mut parked).await;
2316
2317 assert!(
2318 !harness.is_healthy(),
2319 "accounting divergence must drop readiness"
2320 );
2321 assert!(harness.counters.integrity_fault());
2322 assert_eq!(harness.counters.snapshot().abandoned, 1);
2323 assert_eq!(harness.counters.snapshot().released, 0);
2324 }
2325
2326 #[tokio::test(start_paused = true)]
2329 async fn hung_release_times_out_and_reparks() {
2330 let scripted = ScriptedAllocator::new([(LeaseId(5), Refusal::Hang)]);
2331 let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
2332 let harness = Harness::new();
2333 let mut parked = vec![parked_lease(5), parked_lease(6)];
2334
2335 tokio::time::timeout(
2336 std::time::Duration::from_secs(30),
2337 harness.release_quiesced(&allocator, &mut parked),
2338 )
2339 .await
2340 .expect("a hung release must not stall the pass");
2341
2342 assert_eq!(scripted.released(), [LeaseId(6)], "the pass continues");
2343 assert_eq!(parked.len(), 1, "the timed-out lease is retried, not lost");
2344 assert_eq!(parked[0].grant().lease_id, LeaseId(5));
2345 assert!(
2346 harness.is_healthy(),
2347 "a timeout is not accounting divergence"
2348 );
2349 }
2350
2351 #[tokio::test(start_paused = true)]
2355 async fn a_release_pass_costs_one_budget_whatever_the_parked_count() {
2356 let scripted = ScriptedAllocator::new((1..=6).map(|id| (LeaseId(id), Refusal::Hang)));
2357 let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
2358 let harness = Harness::new();
2359 let mut parked: Vec<_> = (1..=6).map(parked_lease).collect();
2360 let budget = std::time::Duration::from_millis(50);
2361
2362 let began = tokio::time::Instant::now();
2363 let outcome = harness.release_pass(&allocator, &mut parked, budget).await;
2364 let elapsed = began.elapsed();
2365
2366 assert_eq!(outcome, ReleasePass::BudgetExpired);
2367 assert!(
2368 elapsed < budget * 2,
2369 "a pass over six hung leases took {elapsed:?}, which is per-call not per-pass"
2370 );
2371 assert_eq!(parked.len(), 6, "every lease is still parked, none dropped");
2372 }
2373
2374 #[tokio::test(start_paused = true)]
2377 async fn a_lease_that_eats_the_budget_yields_its_place() {
2378 let scripted = ScriptedAllocator::new([(LeaseId(1), Refusal::Hang)]);
2379 let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
2380 let harness = Harness::new();
2381 let mut parked = vec![parked_lease(1), parked_lease(2), parked_lease(3)];
2382 let budget = std::time::Duration::from_millis(50);
2383
2384 let outcome = harness.release_pass(&allocator, &mut parked, budget).await;
2385
2386 assert_eq!(outcome, ReleasePass::BudgetExpired);
2387 let order: Vec<_> = parked.iter().map(|l| l.grant().lease_id).collect();
2388 assert_eq!(
2389 order,
2390 [LeaseId(2), LeaseId(3), LeaseId(1)],
2391 "the hung lease must not hold the front of the queue every pass"
2392 );
2393 assert_eq!(scripted.released(), [], "nothing settled under a hung head");
2394 }
2395
2396 #[tokio::test(start_paused = true)]
2400 async fn shutdown_during_a_release_reparks_every_lease() {
2401 let scripted = ScriptedAllocator::new([(LeaseId(1), Refusal::Hang)]);
2402 let allocator: Arc<dyn LeaseAllocator> = Arc::clone(&scripted) as _;
2403 let harness = Harness::new();
2404 let mut parked = vec![parked_lease(1), parked_lease(2), parked_lease(3)];
2405 let (tx, mut shutdown) = watch::channel(false);
2406 let clock = clock();
2407
2408 let pass = release_quiesced(
2409 &allocator,
2410 &mut parked,
2411 &clock,
2412 &harness.config,
2413 &harness.slot,
2414 &harness.health,
2415 &harness.counters,
2416 &mut shutdown,
2417 tokio::time::Instant::now() + std::time::Duration::from_secs(3_600),
2419 );
2420 let signal = async {
2421 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
2422 tx.send(true).unwrap();
2423 };
2424 let (outcome, ()) = tokio::join!(pass, signal);
2425
2426 assert_eq!(outcome, ReleasePass::ShutdownObserved);
2427 assert_eq!(
2428 parked.len(),
2429 3,
2430 "the in-flight lease and the unexamined ones all survive the pass"
2431 );
2432 assert_eq!(parked[0].grant().lease_id, LeaseId(1));
2433 }
2434}