1use std::collections::{BTreeMap, BTreeSet};
5use std::sync::atomic::{AtomicBool, Ordering};
6use std::sync::{Arc, Mutex};
7use std::time::Duration;
8
9use jiff::Timestamp;
10use tokio::sync::watch;
11use tokio::task::{JoinHandle, JoinSet};
12use tokio::time::Instant;
13use tollgate_admission::{AdmissionCounters, AdmissionEngine, ArcSwapSnapshotMap, RequestContext};
14use tollgate_core::{
15 AccountId, DenyReason, LocalSharding, PermissionBits, Principal, ShardOccupancy,
16};
17use tollgate_store::{Clock, LeaseAllocator, SnapshotSource, UsageSink};
18
19use crate::registry::AccountBinding;
20pub use crate::registry::RuntimeFundingReport;
21use crate::{
22 AccountLeaseConfig, LeaseCounters, LeaseManager, LeaseManagerReport, LeaseStats, SlotRegistry,
23 SnapshotCounters, SnapshotManager, SnapshotManagerConfig, SnapshotManagerReport, SnapshotStats,
24 TrackedPrincipals, UsageRecorder, UsageWriter, UsageWriterConfig, WriterHealth,
25 WriterShutdownError, WriterStats,
26};
27
28#[derive(Debug, Clone)]
36pub struct InstanceRuntimeConfig {
37 pub snapshots: SnapshotManagerConfig,
42 pub leases: AccountLeaseConfig,
47 pub usage: UsageWriterConfig,
51 pub sharding: LocalSharding,
58 pub snapshot_history_capacity: std::num::NonZeroUsize,
68 pub idle_account_linger: Duration,
79 pub manager_restart_backoff: Duration,
90 pub shutdown_deadline: Duration,
100}
101
102#[derive(Debug, Clone, PartialEq, Eq)]
106pub struct InstanceRuntimeConfigError(pub String);
107impl std::fmt::Display for InstanceRuntimeConfigError {
108 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
109 f.write_str(&self.0)
110 }
111}
112impl std::error::Error for InstanceRuntimeConfigError {}
113
114impl InstanceRuntimeConfig {
115 pub fn validate(&self) -> Result<(), InstanceRuntimeConfigError> {
126 let error = |text: String| InstanceRuntimeConfigError(text);
127 self.snapshots
128 .validate()
129 .map_err(|e| error(e.to_string()))?;
130 if matches!(&self.snapshots.principals, TrackedPrincipals::Fixed(principals)
131 if principals.len() > self.snapshot_history_capacity.get())
132 {
133 return Err(error(
134 "fixed principals exceed snapshot_history_capacity".into(),
135 ));
136 }
137 self.leases.validate().map_err(|e| error(e.to_string()))?;
138 self.usage.validate().map_err(|e| error(e.to_string()))?;
139 let phases = self
140 .usage
141 .shutdown_drain_deadline
142 .checked_add(self.leases.shutdown_release_deadline)
143 .ok_or_else(|| error("shutdown phase budgets overflow".into()))?;
144 if self.shutdown_deadline < phases {
145 return Err(error(
146 "shutdown_deadline must cover usage drain and lease release".into(),
147 ));
148 }
149 if self.manager_restart_backoff.is_zero() {
150 return Err(error("manager_restart_backoff must be positive".into()));
151 }
152 for duration in [
153 self.shutdown_deadline,
154 self.idle_account_linger,
155 self.manager_restart_backoff,
156 self.snapshots.refresh_interval,
157 self.snapshots.retry_backoff,
158 self.snapshots.fetch_timeout,
159 self.snapshots.enumeration_timeout,
160 self.leases.poll_interval,
161 self.leases.store_call_timeout,
162 self.usage.flush_interval,
163 self.usage.retry_backoff,
164 self.usage.ingest_timeout,
165 ] {
166 if Instant::now().checked_add(duration).is_none() {
167 return Err(error(
168 "runtime duration exceeds the monotonic clock domain".into(),
169 ));
170 }
171 }
172 Ok(())
173 }
174}
175
176#[derive(Debug, Clone, Copy, PartialEq, Eq)]
182pub enum AccountPhase {
183 Running,
185 Lingering,
189 Retiring,
192 Backoff,
195 Dormant,
198 Faulted,
202}
203
204#[derive(Debug, Clone)]
208pub struct AccountReport {
209 pub account: AccountId,
211 pub phase: AccountPhase,
213 pub eligible: bool,
216 pub fundable: bool,
221 pub task_healthy: bool,
224 pub restarts: u64,
227 pub unrecovered_grants: u64,
230 pub uncertain_acquires: u64,
233 pub refill: Option<LeaseStats>,
235}
236
237#[derive(Debug, Clone)]
243pub struct RuntimeReadiness {
244 pub stopping: bool,
247 pub snapshots_ready: bool,
252 pub background_healthy: bool,
256 pub accounting_healthy: bool,
260 pub eligible_accounts: usize,
262 pub unfundable_accounts: usize,
267 pub unmanaged_accounts: usize,
271 pub unresolved_principals: u64,
274 ready: bool,
275}
276impl RuntimeReadiness {
277 #[must_use]
284 pub fn is_ready(&self) -> bool {
285 self.ready
286 }
287}
288
289#[derive(Debug, Clone)]
294pub struct RuntimeReport {
295 pub retained_accounts: usize,
299 pub managed_accounts: usize,
301 pub lingering_accounts: usize,
304 pub retiring_accounts: usize,
306 pub restarting_accounts: usize,
309 pub manager_restarts: u64,
313 pub unrecovered_grants: u64,
317 pub uncertain_acquires: u64,
319 pub counter_overflow: bool,
322 pub refill: Option<LeaseStats>,
325 pub snapshots: SnapshotStats,
327 pub accounting: WriterHealth,
329 pub sharding: ShardOccupancy,
339 pub contention: ContentionReport,
347}
348
349#[derive(Debug, Clone, PartialEq, Eq)]
356pub struct ContentionReport {
357 pub contended_exchanges: u64,
359 pub hottest: Vec<(AccountId, u64)>,
363}
364
365impl ContentionReport {
366 pub const HOTTEST: usize = 8;
368
369 fn collect(slots: Vec<(AccountId, Arc<tollgate_admission::LeaseSlot>)>) -> Self {
370 Self::rank(
371 slots
372 .into_iter()
373 .map(|(account, slot)| (account, slot.contended_exchanges())),
374 )
375 }
376
377 fn rank(counts: impl Iterator<Item = (AccountId, u64)>) -> Self {
378 let mut counted: Vec<(AccountId, u64)> = counts.filter(|&(_, count)| count > 0).collect();
379 let contended_exchanges = counted
380 .iter()
381 .fold(0u64, |total, &(_, count)| total.saturating_add(count));
382 counted.sort_unstable_by(|a, b| b.1.cmp(&a.1).then(a.0.cmp(&b.0)));
385 counted.truncate(Self::HOTTEST);
386 Self {
387 contended_exchanges,
388 hottest: counted,
389 }
390 }
391}
392
393#[derive(Debug)]
396pub enum RuntimeWriterError {
397 Task(WriterShutdownError),
400 Deadline(WriterHealth),
404}
405#[derive(Debug)]
408pub struct RuntimeShutdownReport {
409 pub usage: Result<WriterStats, RuntimeWriterError>,
412 pub snapshots: Option<SnapshotManagerReport>,
415 pub accounts: BTreeMap<AccountId, LeaseManagerReport>,
419 pub unfinished_accounts: Vec<AccountId>,
422 pub background_failed: bool,
427 pub deadline_expired: bool,
430}
431
432struct Observation {
433 phase: AccountPhase,
434 health: Option<watch::Receiver<bool>>,
435 counters: Option<Arc<LeaseCounters>>,
436 settled: Option<LeaseStats>,
437 inherited_grants: u64,
438 restarts: u64,
439 unrecovered: u64,
440 uncertain: u64,
441 overflow: bool,
442}
443fn add_counter(counter: &mut u64, value: u64, overflow: &mut bool) {
444 if let Some(sum) = counter.checked_add(value) {
445 *counter = sum;
446 } else {
447 *counter = u64::MAX;
448 *overflow = true;
449 }
450}
451
452impl Observation {
453 fn uncertain_acquires(&self) -> Option<u64> {
454 self.uncertain.checked_add(
455 self.counters
456 .as_ref()
457 .map_or(0, |c| c.snapshot().uncertain_acquires),
458 )
459 }
460
461 fn stats(&self) -> Option<LeaseStats> {
462 self.settled?.checked_add(
463 self.counters
464 .as_ref()
465 .map_or(LeaseStats::ZERO, |c| c.snapshot()),
466 )
467 }
468 fn healthy(&self) -> bool {
469 self.phase != AccountPhase::Faulted && self.health.as_ref().is_none_or(plane_healthy)
470 }
471}
472
473struct Shared {
474 engine: AdmissionEngine<Arc<ArcSwapSnapshotMap>>,
475 slots: Arc<SlotRegistry>,
476 recorder: UsageRecorder,
477 snapshots_ready: watch::Receiver<bool>,
478 snapshot_counters: Arc<SnapshotCounters>,
479 observations: Mutex<BTreeMap<AccountId, Observation>>,
480 stop: watch::Sender<Option<Instant>>,
481 stopping: AtomicBool,
482 failed: AtomicBool,
483 stopped: AtomicBool,
484 all: bool,
485}
486
487impl Shared {
488 fn request_shutdown(&self, budget: Duration) -> Instant {
489 self.stopping.store(true, Ordering::Release);
490 self.stop.send_if_modified(|current| {
493 if current.is_none() {
494 *current = Some(Instant::now() + budget);
495 true
496 } else {
497 false
498 }
499 });
500 self.stop
501 .borrow()
502 .expect("shutdown request published its deadline")
503 }
504}
505
506#[derive(Clone)]
509pub struct RuntimeHandle {
510 shared: Arc<Shared>,
511 budget: Duration,
512}
513impl RuntimeHandle {
514 #[must_use]
519 pub fn funding(&self, now: Timestamp) -> RuntimeFundingReport {
520 self.shared.slots.funding(now)
521 }
522 pub fn begin(
535 &self,
536 principal: Principal,
537 required: PermissionBits,
538 now: Timestamp,
539 ) -> Result<RequestContext, DenyReason> {
540 self.shared.engine.begin(principal, required, now)
541 }
542 #[must_use]
546 pub fn recorder(&self) -> &UsageRecorder {
547 &self.shared.recorder
548 }
549 #[must_use]
551 pub fn counters(&self) -> &AdmissionCounters {
552 self.shared.engine.counters()
553 }
554
555 pub fn request_shutdown(&self) -> Instant {
559 self.shared.request_shutdown(self.budget)
560 }
561
562 #[must_use]
568 pub fn readiness(&self, now: Timestamp) -> RuntimeReadiness {
569 let stopping = self.shared.stopping.load(Ordering::Acquire);
570 let (tracked, unresolved) = self.shared.slots.resolution_counts(now);
571 let snapshots_ready = plane_healthy(&self.shared.snapshots_ready)
572 && if self.shared.all {
573 tracked == 0 || unresolved < tracked
574 } else {
575 unresolved == 0
576 };
577 let bindings = self.shared.slots.bindings();
578 let eligible_accounts = bindings.iter().filter(|b| b.eligible(now)).count();
579 let unfundable_accounts = bindings
580 .iter()
581 .filter(|b| b.eligible(now) && !b.fundable(now))
582 .count();
583 let observations = self
584 .shared
585 .observations
586 .lock()
587 .expect("runtime observations poisoned");
588 let unmanaged_accounts = bindings
589 .iter()
590 .filter(|binding| {
591 binding.eligible(now)
592 && !observations.get(&binding.account).is_some_and(|o| {
593 matches!(o.phase, AccountPhase::Running | AccountPhase::Lingering)
594 && o.healthy()
595 })
596 })
597 .count();
598 let background_healthy = !self.shared.failed.load(Ordering::Acquire)
599 && !self.shared.stopped.load(Ordering::Acquire)
600 && self.shared.snapshots_ready.has_changed().is_ok()
601 && unmanaged_accounts == 0
602 && observations
603 .values()
604 .all(|o| o.phase != AccountPhase::Faulted);
605 let accounting = self.shared.recorder.health();
606 let accounting_healthy = !self.shared.recorder.is_closed()
607 && accounting.stats.lost == 0
608 && accounting.stats.rejected == 0
609 && accounting.queue_depth < accounting.queue_capacity;
610 let funded = if bindings.is_empty() {
611 true
612 } else if self.shared.all {
613 eligible_accounts > unfundable_accounts
614 } else {
615 eligible_accounts > 0 && unfundable_accounts == 0
616 };
617 let ready =
618 !stopping && snapshots_ready && background_healthy && accounting_healthy && funded;
619 RuntimeReadiness {
620 stopping,
621 snapshots_ready,
622 background_healthy,
623 accounting_healthy,
624 eligible_accounts,
625 unfundable_accounts,
626 unmanaged_accounts,
627 unresolved_principals: unresolved as u64,
628 ready,
629 }
630 }
631
632 #[must_use]
634 pub fn account_reports(&self, now: Timestamp) -> Vec<AccountReport> {
635 let bindings: BTreeMap<_, _> = self
636 .shared
637 .slots
638 .bindings()
639 .into_iter()
640 .map(|b| (b.account, b))
641 .collect();
642 self.shared
643 .observations
644 .lock()
645 .expect("runtime observations poisoned")
646 .iter()
647 .map(|(&account, o)| {
648 let binding = bindings.get(&account);
649 AccountReport {
650 account,
651 phase: o.phase,
652 eligible: binding.is_some_and(|b| b.eligible(now)),
653 fundable: binding.is_some_and(|b| b.fundable(now)),
654 task_healthy: o.healthy(),
655 restarts: o.restarts,
656 unrecovered_grants: o.unrecovered,
657 uncertain_acquires: o.uncertain_acquires().unwrap_or(u64::MAX),
658 refill: o.stats(),
659 }
660 })
661 .collect()
662 }
663
664 #[must_use]
668 pub fn report(&self) -> RuntimeReport {
669 let observations = self
670 .shared
671 .observations
672 .lock()
673 .expect("runtime observations poisoned");
674 let mut report = RuntimeReport {
675 retained_accounts: self.shared.slots.retained_slots(),
676 managed_accounts: 0,
677 lingering_accounts: 0,
678 retiring_accounts: 0,
679 restarting_accounts: 0,
680 manager_restarts: 0,
681 unrecovered_grants: 0,
682 uncertain_acquires: 0,
683 counter_overflow: false,
684 refill: Some(LeaseStats::ZERO),
685 snapshots: self.shared.snapshot_counters.snapshot(),
686 accounting: self.shared.recorder.health(),
687 sharding: self.shared.slots.sharding().occupancy(),
688 contention: ContentionReport::collect(self.shared.slots.slots()),
689 };
690 for o in observations.values() {
691 report.counter_overflow |= o.overflow;
692 let uncertain = o.uncertain_acquires();
693 report.counter_overflow |= uncertain.is_none();
694 match o.phase {
695 AccountPhase::Running => report.managed_accounts += 1,
696 AccountPhase::Lingering => {
697 report.managed_accounts += 1;
698 report.lingering_accounts += 1;
699 }
700 AccountPhase::Retiring => report.retiring_accounts += 1,
701 AccountPhase::Backoff => report.restarting_accounts += 1,
702 AccountPhase::Dormant | AccountPhase::Faulted => {}
703 }
704 report.refill = report.refill.and_then(|sum| sum.checked_add(o.stats()?));
705 for (sum, value) in [
706 (&mut report.manager_restarts, o.restarts),
707 (&mut report.unrecovered_grants, o.unrecovered),
708 (
709 &mut report.uncertain_acquires,
710 uncertain.unwrap_or(u64::MAX),
711 ),
712 ] {
713 if let Some(next) = sum.checked_add(value) {
714 *sum = next;
715 } else {
716 *sum = u64::MAX;
717 report.counter_overflow = true;
718 }
719 }
720 }
721 report.counter_overflow |= report.refill.is_none();
722 report
723 }
724}
725
726#[must_use = "retain the runtime and await shutdown to settle usage and leases"]
729pub struct InstanceRuntime {
730 handle: RuntimeHandle,
731 task: Option<JoinHandle<RuntimeShutdownReport>>,
732}
733impl InstanceRuntime {
734 #[must_use]
736 pub fn handle(&self) -> RuntimeHandle {
737 self.handle.clone()
738 }
739 pub fn spawn(
753 source: Arc<dyn SnapshotSource>,
754 allocator: Arc<dyn LeaseAllocator>,
755 sink: Arc<dyn UsageSink>,
756 clock: Arc<dyn Clock>,
757 config: InstanceRuntimeConfig,
758 ) -> Result<(Self, RuntimeHandle), InstanceRuntimeConfigError> {
759 config.validate()?;
760 let map = Arc::new(ArcSwapSnapshotMap::with_capacities(
761 config.sharding,
762 ArcSwapSnapshotMap::DEFAULT_MAX_NEGATIVE_ENTRIES,
763 config.snapshot_history_capacity,
764 ));
765 let (slots, changes) = SlotRegistry::observed(config.sharding);
766 let (recorder, writer) = UsageWriter::spawn(sink, Arc::clone(&clock), config.usage)
767 .map_err(|e| InstanceRuntimeConfigError(e.to_string()))?;
768 let snapshots = SnapshotManager::spawn(
769 source,
770 map.clone(),
771 Arc::clone(&slots),
772 Arc::clone(&clock),
773 config.snapshots.clone(),
774 )
775 .map_err(|e| InstanceRuntimeConfigError(e.to_string()))?;
776 let (stop, stopping) = watch::channel(None);
777 let shared = Arc::new(Shared {
778 engine: AdmissionEngine::new(map),
779 slots,
780 recorder,
781 snapshots_ready: snapshots.ready(),
782 snapshot_counters: snapshots.counters(),
783 observations: Mutex::new(BTreeMap::new()),
784 stop,
785 stopping: AtomicBool::new(false),
786 failed: AtomicBool::new(false),
787 stopped: AtomicBool::new(false),
788 all: matches!(config.snapshots.principals, TrackedPrincipals::All { .. }),
789 });
790 let handle = RuntimeHandle {
791 shared: Arc::clone(&shared),
792 budget: config.shutdown_deadline,
793 };
794 let task = tokio::spawn(supervise(
795 shared, allocator, clock, config, snapshots, writer, changes, stopping,
796 ));
797 Ok((
798 Self {
799 handle: handle.clone(),
800 task: Some(task),
801 },
802 handle,
803 ))
804 }
805 pub async fn shutdown(mut self) -> Result<RuntimeShutdownReport, tokio::task::JoinError> {
818 self.handle.request_shutdown();
819 let report = self
820 .task
821 .as_mut()
822 .expect("runtime owns its supervisor")
823 .await;
824 drop(self.task.take());
825 report
826 }
827}
828impl Drop for InstanceRuntime {
829 fn drop(&mut self) {
830 self.handle.shared.stopping.store(true, Ordering::Release);
831 if let Some(task) = self.task.take() {
832 tracing::warn!(
833 unaccounted = self.handle.shared.recorder.health().unaccounted,
834 "instance runtime dropped; aborting tasks, unfinished grants require TTL reclaim"
835 );
836 task.abort();
837 }
838 }
839}
840
841fn plane_healthy(health: &watch::Receiver<bool>) -> bool {
842 health.has_changed().is_ok() && *health.borrow()
843}
844
845struct Managed {
846 manager: Option<LeaseManager>,
847 binding: AccountBinding,
848 desired: bool,
849 timer: Option<Instant>,
850 restart_after_death: bool,
851 monitor: Option<tokio::task::Id>,
852}
853struct Exit {
854 account: AccountId,
855 report: LeaseManagerReport,
856}
857
858struct Supervisor {
859 shared: Arc<Shared>,
860 allocator: Arc<dyn LeaseAllocator>,
861 clock: Arc<dyn Clock>,
862 config: InstanceRuntimeConfig,
863 accounts: BTreeMap<AccountId, Managed>,
864 timers: BTreeSet<(Instant, AccountId)>,
865 monitors: JoinSet<AccountId>,
866 cleanup: JoinSet<Exit>,
867}
868impl Supervisor {
869 fn phase(&self, account: AccountId, phase: AccountPhase) {
870 self.shared
871 .observations
872 .lock()
873 .expect("runtime observations poisoned")
874 .get_mut(&account)
875 .expect("managed account has observations")
876 .phase = phase;
877 }
878 fn arm(&mut self, account: AccountId, after: Duration) {
879 let record = self
880 .accounts
881 .get_mut(&account)
882 .expect("timer belongs to account");
883 if let Some(old) = record.timer.take() {
884 self.timers.remove(&(old, account));
885 }
886 let at = Instant::now() + after;
887 record.timer = Some(at);
888 self.timers.insert((at, account));
889 }
890 fn disarm(&mut self, account: AccountId) {
891 if let Some(at) = self
892 .accounts
893 .get_mut(&account)
894 .expect("managed account")
895 .timer
896 .take()
897 {
898 self.timers.remove(&(at, account));
899 }
900 }
901 fn start(&mut self, account: AccountId) {
902 let record = self.accounts.get_mut(&account).expect("managed account");
903 let inherited_grants = u64::from(record.binding.slot.load_observed().is_some());
904 let manager = LeaseManager::spawn(
905 Arc::clone(&self.allocator),
906 Arc::clone(&record.binding.slot),
907 Arc::clone(&self.clock),
908 self.config.leases.for_account(account),
909 )
910 .expect("runtime validated the complete lease configuration before spawning");
911 let mut health = manager.health();
912 let mut observations = self
913 .shared
914 .observations
915 .lock()
916 .expect("runtime observations poisoned");
917 let o = observations.get_mut(&account).expect("managed observation");
918 if record.restart_after_death {
919 add_counter(&mut o.restarts, 1, &mut o.overflow);
920 }
921 o.inherited_grants = inherited_grants;
922 o.health = Some(health.clone());
923 o.counters = Some(manager.counters());
924 o.phase = AccountPhase::Running;
925 record.restart_after_death = false;
926 record.manager = Some(manager);
927 record.monitor = Some(
928 self.monitors
929 .spawn(async move {
930 while plane_healthy(&health) {
931 if health.changed().await.is_err() {
932 break;
933 }
934 }
935 account
936 })
937 .id(),
938 );
939 }
940 fn reconcile(&mut self, binding: AccountBinding, now: Timestamp) {
941 let account = binding.account;
942 let desired = binding.eligible(now);
943 let record = self.accounts.entry(account).or_insert_with(|| {
944 self.shared
945 .observations
946 .lock()
947 .expect("runtime observations poisoned")
948 .insert(
949 account,
950 Observation {
951 phase: AccountPhase::Dormant,
952 health: None,
953 counters: None,
954 settled: Some(LeaseStats::ZERO),
955 inherited_grants: 0,
956 restarts: 0,
957 unrecovered: 0,
958 uncertain: 0,
959 overflow: false,
960 },
961 );
962 Managed {
963 manager: None,
964 binding: binding.clone(),
965 desired: false,
966 timer: None,
967 restart_after_death: false,
968 monitor: None,
969 }
970 });
971 record.binding = binding;
972 record.desired = desired;
973 let phase = self
974 .shared
975 .observations
976 .lock()
977 .expect("runtime observations poisoned")[&account]
978 .phase;
979 match (desired, phase) {
980 (true, AccountPhase::Dormant) => {
981 self.disarm(account);
982 self.start(account);
983 }
984 (true, AccountPhase::Lingering) => {
985 self.disarm(account);
986 self.phase(account, AccountPhase::Running);
987 }
988 (false, AccountPhase::Running) => {
989 self.phase(account, AccountPhase::Lingering);
990 self.arm(account, self.config.idle_account_linger);
991 }
992 (false, AccountPhase::Backoff) => {
993 self.disarm(account);
994 self.phase(account, AccountPhase::Dormant);
995 }
996 _ => {}
997 }
998 }
999
1000 fn retire(&mut self, account: AccountId, deadline: Instant) {
1001 self.disarm(account);
1002 if let Some(manager) = self
1003 .accounts
1004 .get_mut(&account)
1005 .expect("managed account")
1006 .manager
1007 .take()
1008 {
1009 self.phase(account, AccountPhase::Retiring);
1010 manager.stop_at(deadline);
1011 self.cleanup.spawn(async move {
1012 Exit {
1013 account,
1014 report: manager.shutdown().await,
1015 }
1016 });
1017 }
1018 }
1019 fn completed(&mut self, exit: Exit, restart: bool) {
1020 let account = exit.account;
1021 let mut observations = self
1022 .shared
1023 .observations
1024 .lock()
1025 .expect("runtime observations poisoned");
1026 let o = observations.get_mut(&account).expect("managed observation");
1027 let faulted = o.counters.as_ref().is_some_and(|c| c.integrity_fault());
1031 if faulted {
1032 self.shared.failed.store(true, Ordering::Release);
1033 self.shared.request_shutdown(self.config.shutdown_deadline);
1034 }
1035 let pending = o.counters.as_ref().is_some_and(|c| c.acquire_pending());
1036 let last = o.counters.take().map_or(LeaseStats::ZERO, |c| c.snapshot());
1037 add_counter(&mut o.uncertain, u64::from(pending), &mut o.overflow);
1038 add_counter(&mut o.uncertain, last.uncertain_acquires, &mut o.overflow);
1039 if exit.report.task_died {
1040 let current = u64::from(
1044 self.accounts[&account]
1045 .binding
1046 .slot
1047 .load_observed()
1048 .is_some(),
1049 );
1050 let unresolved = (u128::from(last.acquired) + u128::from(o.inherited_grants))
1054 .saturating_sub(u128::from(last.released))
1055 .saturating_sub(u128::from(last.abandoned))
1056 .saturating_sub(u128::from(current));
1057 let unresolved = u64::try_from(unresolved).unwrap_or_else(|_| {
1058 o.overflow = true;
1059 u64::MAX
1060 });
1061 add_counter(&mut o.unrecovered, unresolved, &mut o.overflow);
1062 self.accounts
1063 .get_mut(&account)
1064 .expect("managed account")
1065 .restart_after_death = true;
1066 tracing::error!(%account, unrecovered_grants = unresolved, uncertain_acquires = o.uncertain,
1067 "lease manager died; unrecovered grants return only at TTL reclaim");
1068 }
1069 o.settled = o.settled.and_then(|old| old.checked_add(last));
1070 o.health = None;
1071 let retry = restart && !faulted && self.accounts[&account].desired;
1072 o.phase = if faulted {
1073 AccountPhase::Faulted
1074 } else if retry {
1075 AccountPhase::Backoff
1076 } else {
1077 AccountPhase::Dormant
1078 };
1079 drop(observations);
1080 if retry {
1081 self.arm(account, self.config.manager_restart_backoff);
1082 }
1083 }
1084}
1085
1086struct SupervisorLiveness(Arc<Shared>);
1088impl Drop for SupervisorLiveness {
1089 fn drop(&mut self) {
1090 self.0.stopped.store(true, Ordering::Release);
1091 }
1092}
1093
1094#[allow(
1099 clippy::too_many_arguments,
1100 reason = "supervisor entry point: every argument is a handle it must keep alive for the process"
1101)]
1102async fn supervise(
1103 shared: Arc<Shared>,
1104 allocator: Arc<dyn LeaseAllocator>,
1105 clock: Arc<dyn Clock>,
1106 config: InstanceRuntimeConfig,
1107 snapshots: SnapshotManager,
1108 writer: UsageWriter,
1109 mut changes: watch::Receiver<()>,
1110 mut stop: watch::Receiver<Option<Instant>>,
1111) -> RuntimeShutdownReport {
1112 let _liveness = SupervisorLiveness(Arc::clone(&shared));
1113 let mut ready = snapshots.ready();
1114 let mut supervisor = Supervisor {
1115 shared: Arc::clone(&shared),
1116 allocator,
1117 clock,
1118 config,
1119 accounts: BTreeMap::new(),
1120 timers: BTreeSet::new(),
1121 monitors: JoinSet::new(),
1122 cleanup: JoinSet::new(),
1123 };
1124 let deadline = loop {
1125 if let Some(deadline) = *stop.borrow() {
1126 break deadline;
1127 }
1128 let now = supervisor.clock.now();
1129 for binding in shared.slots.drain_changes(now) {
1130 supervisor.reconcile(binding, now);
1131 }
1132 let expiry = shared.slots.next_expiry().map(|at| {
1133 let delay =
1134 std::time::Duration::try_from(now.duration_until(at)).unwrap_or(Duration::ZERO);
1135 Instant::now() + delay.min(Duration::from_secs(3_600))
1136 });
1137 let timer = supervisor.timers.first().map(|(at, _)| *at);
1138 let wake = [expiry, timer]
1139 .into_iter()
1140 .flatten()
1141 .min()
1142 .unwrap_or_else(|| Instant::now() + Duration::from_secs(3_600));
1143 tokio::select! {
1144 changed = stop.changed() => { if changed.is_err() { break shared.request_shutdown(supervisor.config.shutdown_deadline); } }
1145 changed = changes.changed() => { if changed.is_err() { shared.failed.store(true, Ordering::Release); break shared.request_shutdown(supervisor.config.shutdown_deadline); } }
1146 () = async { while ready.changed().await.is_ok() {} } => { shared.failed.store(true, Ordering::Release); break shared.request_shutdown(supervisor.config.shutdown_deadline); }
1147 () = shared.recorder.closed() => { shared.failed.store(true, Ordering::Release); break shared.request_shutdown(supervisor.config.shutdown_deadline); }
1148 exit = supervisor.cleanup.join_next(), if !supervisor.cleanup.is_empty() => {
1149 match exit {
1150 Some(Ok(exit)) => supervisor.completed(exit, true),
1151 _ => { shared.failed.store(true, Ordering::Release); break shared.request_shutdown(supervisor.config.shutdown_deadline); }
1152 }
1153 }
1154 notice = supervisor.monitors.join_next_with_id(), if !supervisor.monitors.is_empty() => {
1155 if let Some(Ok((id, account))) = notice {
1156 if supervisor.accounts[&account].monitor != Some(id) { continue; }
1157 if let Some(manager) = &supervisor.accounts[&account].manager {
1158 if manager.counters().integrity_fault() {
1159 supervisor.phase(account, AccountPhase::Faulted);
1160 shared.failed.store(true, Ordering::Release);
1161 break shared.request_shutdown(supervisor.config.shutdown_deadline);
1162 }
1163 supervisor.retire(account, Instant::now() + supervisor.config.leases.shutdown_release_deadline);
1164 }
1165 } else { shared.failed.store(true, Ordering::Release); break shared.request_shutdown(supervisor.config.shutdown_deadline); }
1166 }
1167 () = tokio::time::sleep_until(wake) => {
1168 while let Some(&(at, account)) = supervisor.timers.first() {
1169 if at > Instant::now() { break; }
1170 supervisor.disarm(account);
1171 if supervisor.accounts[&account].desired {
1172 if supervisor.accounts[&account].manager.is_none() { supervisor.start(account); }
1173 } else { supervisor.retire(account, Instant::now() + supervisor.config.leases.shutdown_release_deadline); }
1174 }
1175 }
1176 }
1177 };
1178 shared.stopping.store(true, Ordering::Release);
1179 for record in supervisor.accounts.values() {
1180 if let Some(manager) = &record.manager {
1181 manager.pause_refills();
1182 }
1183 }
1184 supervisor.monitors.abort_all();
1185 writer.stop_at(deadline.min(Instant::now() + supervisor.config.usage.shutdown_drain_deadline));
1186 let snapshot_report = tokio::time::timeout_at(deadline, snapshots.shutdown())
1187 .await
1188 .ok();
1189 let usage = match tokio::time::timeout_at(deadline, writer.shutdown()).await {
1190 Ok(Ok(stats)) => Ok(stats),
1191 Ok(Err(error)) => Err(RuntimeWriterError::Task(error)),
1192 Err(_) => Err(RuntimeWriterError::Deadline(shared.recorder.health())),
1193 };
1194 let accounts: Vec<_> = supervisor.accounts.keys().copied().collect();
1195 for account in accounts {
1196 supervisor.retire(account, deadline);
1197 }
1198 let mut returned = BTreeMap::new();
1199 while !supervisor.cleanup.is_empty() {
1200 match tokio::time::timeout_at(deadline, supervisor.cleanup.join_next()).await {
1201 Ok(Some(Ok(exit))) => {
1202 returned.insert(exit.account, exit.report);
1203 supervisor.completed(exit, false);
1204 }
1205 Ok(Some(Err(_))) => {
1206 shared.failed.store(true, Ordering::Release);
1207 }
1208 Ok(None) => break,
1209 Err(_) => {
1210 supervisor.cleanup.abort_all();
1211 break;
1212 }
1213 }
1214 }
1215 let unfinished_accounts = shared
1216 .observations
1217 .lock()
1218 .expect("runtime observations poisoned")
1219 .iter()
1220 .filter(|(_, o)| o.phase == AccountPhase::Retiring)
1221 .map(|(&account, _)| account)
1222 .collect();
1223 RuntimeShutdownReport {
1224 usage,
1225 snapshots: snapshot_report,
1226 accounts: returned,
1227 unfinished_accounts,
1228 background_failed: shared.failed.load(Ordering::Acquire),
1229 deadline_expired: Instant::now() >= deadline,
1230 }
1231}
1232
1233#[cfg(test)]
1234mod health_tests {
1235 use super::*;
1236 fn diagnostic_handle() -> (RuntimeHandle, UsageWriter, watch::Sender<bool>) {
1237 let (recorder, writer) = UsageWriter::spawn(
1238 tollgate_store::MemoryStore::new(tollgate_store::GrantPolicy::default()).unwrap(),
1239 Arc::new(crate::ManualClock::new(
1240 Timestamp::from_second(100).unwrap(),
1241 )),
1242 UsageWriterConfig {
1243 queue_capacity: 1,
1244 max_batch: 1,
1245 flush_interval: Duration::from_millis(1),
1246 retry_backoff: Duration::from_millis(1),
1247 shutdown_drain_deadline: Duration::from_millis(10),
1248 ingest_timeout: Duration::from_millis(1),
1249 },
1250 )
1251 .unwrap();
1252 let (snapshots_alive, snapshots_ready) = watch::channel(true);
1253 let (stop, _) = watch::channel(None);
1254 let handle = RuntimeHandle {
1255 shared: Arc::new(Shared {
1256 engine: AdmissionEngine::new(Arc::new(ArcSwapSnapshotMap::default())),
1257 slots: Arc::new(SlotRegistry::default()),
1258 recorder,
1259 snapshots_ready,
1260 snapshot_counters: Arc::new(SnapshotCounters::default()),
1261 observations: Mutex::new(BTreeMap::new()),
1262 stop,
1263 stopping: AtomicBool::new(false),
1264 failed: AtomicBool::new(false),
1265 stopped: AtomicBool::new(false),
1266 all: true,
1267 }),
1268 budget: Duration::from_millis(20),
1269 };
1270 (handle, writer, snapshots_alive)
1271 }
1272
1273 #[tokio::test(start_paused = true)]
1274 async fn supervisor_exit_withdraws_readiness_before_children_receive_their_abort() {
1275 let (handle, writer, _snapshots_alive) = diagnostic_handle();
1276 let now = Timestamp::from_second(100).unwrap();
1277 let liveness = SupervisorLiveness(Arc::clone(&handle.shared));
1278 assert!(handle.readiness(now).is_ready());
1279 drop(liveness);
1280 assert!(!handle.readiness(now).background_healthy);
1281 assert!(!handle.readiness(now).is_ready());
1282 assert!(
1283 !handle.recorder().is_closed(),
1284 "child shutdown has not been polled yet"
1285 );
1286 writer.shutdown().await.unwrap();
1287 }
1288
1289 #[tokio::test(start_paused = true)]
1290 async fn a_component_counter_overflow_survives_aggregation_with_clean_lease_stats() {
1291 let (handle, writer, _snapshots_alive) = diagnostic_handle();
1292 let mut observation = Observation {
1293 phase: AccountPhase::Dormant,
1294 health: None,
1295 counters: None,
1296 settled: Some(LeaseStats::ZERO),
1297 inherited_grants: 0,
1298 restarts: u64::MAX,
1299 unrecovered: 0,
1300 uncertain: 0,
1301 overflow: false,
1302 };
1303 add_counter(&mut observation.restarts, 1, &mut observation.overflow);
1304 handle
1305 .shared
1306 .observations
1307 .lock()
1308 .unwrap()
1309 .insert(AccountId(1), observation);
1310 let report = handle.report();
1311 assert!(report.counter_overflow);
1312 assert_eq!(report.manager_restarts, u64::MAX);
1313 assert_eq!(report.refill, Some(LeaseStats::ZERO));
1314 writer.shutdown().await.unwrap();
1315 }
1316
1317 #[test]
1318 fn a_plane_that_died_while_healthy_is_not_healthy() {
1319 let (sender, receiver) = tokio::sync::watch::channel(true);
1320 assert!(plane_healthy(&receiver));
1321
1322 sender.send_replace(false);
1323 assert!(!plane_healthy(&receiver), "the plane said it is unhealthy");
1324
1325 let (sender, receiver) = tokio::sync::watch::channel(true);
1326 drop(sender);
1327 assert!(
1328 !plane_healthy(&receiver),
1329 "the last value still reads true; the closed channel is the evidence"
1330 );
1331 }
1332}
1333
1334#[cfg(test)]
1335mod contention_tests {
1336 use super::*;
1337
1338 #[test]
1342 fn contention_report_ranks_by_count_then_account() {
1343 let counts =
1344 (0..20u128).map(|account| (AccountId(account), u64::try_from(account % 4).unwrap()));
1345 let report = ContentionReport::rank(counts);
1346 assert_eq!(
1347 report.contended_exchanges,
1348 (0..20u64).map(|a| a % 4).sum::<u64>()
1349 );
1350 assert_eq!(report.hottest.len(), ContentionReport::HOTTEST);
1351 assert_eq!(
1352 report.hottest,
1353 [
1354 (AccountId(3), 3),
1355 (AccountId(7), 3),
1356 (AccountId(11), 3),
1357 (AccountId(15), 3),
1358 (AccountId(19), 3),
1359 (AccountId(2), 2),
1360 (AccountId(6), 2),
1361 (AccountId(10), 2),
1362 ]
1363 );
1364 assert_eq!(
1365 ContentionReport::rank(std::iter::empty()),
1366 ContentionReport {
1367 contended_exchanges: 0,
1368 hottest: Vec::new()
1369 },
1370 "an uncontended instance names no account"
1371 );
1372 assert_eq!(
1373 ContentionReport::rank([(AccountId(1), u64::MAX), (AccountId(2), 5)].into_iter())
1374 .contended_exchanges,
1375 u64::MAX,
1376 "the total saturates"
1377 );
1378 }
1379}