1use std::collections::{BTreeMap, VecDeque};
4use std::fmt;
5use std::sync::{Arc, Mutex};
6
7use crate::metrics::{FairnessLevel, RuntimeMetrics};
8
9#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
11pub struct FairBudget {
12 pub operations: usize,
13 pub bytes: usize,
14}
15
16impl FairBudget {
17 pub const fn new(operations: usize, bytes: usize) -> Self {
18 Self { operations, bytes }
19 }
20
21 fn add_capped(&mut self, quantum: Self, maximum: Self) {
22 self.operations = self
23 .operations
24 .saturating_add(quantum.operations)
25 .min(maximum.operations);
26 self.bytes = self.bytes.saturating_add(quantum.bytes).min(maximum.bytes);
27 }
28
29 fn turn_allowance(self, other: Self, turn_cap: Self) -> Self {
30 Self {
31 operations: self
32 .operations
33 .min(other.operations)
34 .min(turn_cap.operations),
35 bytes: self.bytes.min(other.bytes).min(turn_cap.bytes),
36 }
37 }
38
39 fn consume(&mut self, used: Self) {
40 self.operations -= used.operations;
41 self.bytes -= used.bytes;
42 }
43
44 fn min(self, other: Self) -> Self {
45 Self {
46 operations: self.operations.min(other.operations),
47 bytes: self.bytes.min(other.bytes),
48 }
49 }
50}
51
52#[derive(Clone, Copy, Debug, Eq, PartialEq)]
54pub struct FairnessConfig {
55 pub vm_quantum: FairBudget,
57 pub capability_quantum: FairBudget,
59 pub max_vm_deficit: FairBudget,
61 pub max_capability_deficit: FairBudget,
63 pub max_vms: usize,
65 pub max_capabilities_per_vm: usize,
67}
68
69impl FairnessConfig {
70 fn validate(self) -> Result<Self, FairnessError> {
71 for (field, value) in [
72 ("fairness.vmQuantum.operations", self.vm_quantum.operations),
73 ("fairness.vmQuantum.bytes", self.vm_quantum.bytes),
74 (
75 "fairness.capabilityQuantum.operations",
76 self.capability_quantum.operations,
77 ),
78 (
79 "fairness.capabilityQuantum.bytes",
80 self.capability_quantum.bytes,
81 ),
82 (
83 "fairness.maxVmDeficit.operations",
84 self.max_vm_deficit.operations,
85 ),
86 ("fairness.maxVmDeficit.bytes", self.max_vm_deficit.bytes),
87 (
88 "fairness.maxCapabilityDeficit.operations",
89 self.max_capability_deficit.operations,
90 ),
91 (
92 "fairness.maxCapabilityDeficit.bytes",
93 self.max_capability_deficit.bytes,
94 ),
95 ("fairness.maxVms", self.max_vms),
96 (
97 "fairness.maxCapabilitiesPerVm",
98 self.max_capabilities_per_vm,
99 ),
100 ] {
101 if value == 0 {
102 return Err(FairnessError::InvalidConfig { field });
103 }
104 }
105 for (field, maximum, quantum) in [
106 (
107 "fairness.maxVmDeficit.operations",
108 self.max_vm_deficit.operations,
109 self.vm_quantum.operations,
110 ),
111 (
112 "fairness.maxVmDeficit.bytes",
113 self.max_vm_deficit.bytes,
114 self.vm_quantum.bytes,
115 ),
116 (
117 "fairness.maxCapabilityDeficit.operations",
118 self.max_capability_deficit.operations,
119 self.capability_quantum.operations,
120 ),
121 (
122 "fairness.maxCapabilityDeficit.bytes",
123 self.max_capability_deficit.bytes,
124 self.capability_quantum.bytes,
125 ),
126 ] {
127 if maximum < quantum {
128 return Err(FairnessError::DeficitBelowQuantum { field });
129 }
130 }
131 Ok(self)
132 }
133}
134
135#[derive(Clone, Copy, Debug, Eq, PartialEq)]
136pub enum MembershipUpdate {
137 Enqueued,
138 Coalesced,
139}
140
141#[derive(Clone, Debug, Eq, PartialEq)]
143pub struct FairSelection<VmId, CapabilityId> {
144 pub sequence: u64,
145 pub vm_id: VmId,
146 pub capability_id: CapabilityId,
147 pub allowance: FairBudget,
149}
150
151#[derive(Clone, Copy, Debug, Eq, PartialEq)]
152pub struct FairnessSnapshot {
153 pub vm_deficit: FairBudget,
154 pub capability_deficit: FairBudget,
155 pub vm_queued: bool,
156 pub capability_queued: bool,
157 pub capability_in_flight: bool,
158 pub capability_rearmed: bool,
159}
160
161#[derive(Clone, Debug, Eq, PartialEq)]
162pub enum FairnessError {
163 InvalidConfig {
164 field: &'static str,
165 },
166 DeficitBelowQuantum {
167 field: &'static str,
168 },
169 VmLimit {
170 limit: usize,
171 },
172 CapabilityLimit {
173 limit: usize,
174 },
175 SelectionInFlight {
176 sequence: u64,
177 },
178 StaleSelection {
179 supplied: u64,
180 outstanding: Option<u64>,
181 },
182 OperationBudgetExceeded {
183 used: usize,
184 allowance: usize,
185 },
186 ByteBudgetExceeded {
187 used: usize,
188 allowance: usize,
189 },
190 CapabilityInFlight,
191 CapabilityRetired {
192 vm_generation: u64,
193 capability_id: u64,
194 },
195 SequenceExhausted,
196 Invariant,
197}
198
199impl fmt::Display for FairnessError {
200 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
201 match self {
202 Self::InvalidConfig { field } => write!(
203 formatter,
204 "ERR_AGENTOS_FAIRNESS_CONFIG: {field} must be greater than zero"
205 ),
206 Self::DeficitBelowQuantum { field } => write!(
207 formatter,
208 "ERR_AGENTOS_FAIRNESS_CONFIG: {field} must be at least its quantum"
209 ),
210 Self::VmLimit { limit } => write!(
211 formatter,
212 "ERR_AGENTOS_FAIRNESS_VM_LIMIT: scheduler VM state exceeded {limit}; raise runtime.fairness.maxVms"
213 ),
214 Self::CapabilityLimit { limit } => write!(
215 formatter,
216 "ERR_AGENTOS_FAIRNESS_CAPABILITY_LIMIT: scheduler capability state exceeded {limit}; raise runtime.fairness.maxCapabilitiesPerVm"
217 ),
218 Self::SelectionInFlight { sequence } => write!(
219 formatter,
220 "ERR_AGENTOS_FAIRNESS_SELECTION_IN_FLIGHT: selection {sequence} must complete before selecting again"
221 ),
222 Self::StaleSelection {
223 supplied,
224 outstanding,
225 } => write!(
226 formatter,
227 "ERR_AGENTOS_FAIRNESS_STALE_SELECTION: supplied selection {supplied}, outstanding {outstanding:?}"
228 ),
229 Self::OperationBudgetExceeded { used, allowance } => write!(
230 formatter,
231 "ERR_AGENTOS_FAIRNESS_OPERATION_BUDGET: used {used} operations, allowance {allowance}"
232 ),
233 Self::ByteBudgetExceeded { used, allowance } => write!(
234 formatter,
235 "ERR_AGENTOS_FAIRNESS_BYTE_BUDGET: used {used} bytes, allowance {allowance}"
236 ),
237 Self::CapabilityInFlight => formatter.write_str(
238 "ERR_AGENTOS_FAIRNESS_CAPABILITY_IN_FLIGHT: cannot remove an in-flight capability",
239 ),
240 Self::CapabilityRetired {
241 vm_generation,
242 capability_id,
243 } => write!(
244 formatter,
245 "ERR_AGENTOS_FAIRNESS_CAPABILITY_RETIRED: capability {capability_id} in VM generation {vm_generation} is closed",
246 ),
247 Self::SequenceExhausted => formatter.write_str(
248 "ERR_AGENTOS_FAIRNESS_SEQUENCE_EXHAUSTED: selection sequence exhausted",
249 ),
250 Self::Invariant => formatter.write_str(
251 "ERR_AGENTOS_FAIRNESS_INVARIANT: scheduler membership state is inconsistent",
252 ),
253 }
254 }
255}
256
257impl std::error::Error for FairnessError {}
258
259#[derive(Debug)]
260struct CapabilityState {
261 deficit: FairBudget,
262 queued: bool,
263 in_flight: bool,
264 rearmed: bool,
265}
266
267impl CapabilityState {
268 fn new() -> Self {
269 Self {
270 deficit: FairBudget::default(),
271 queued: false,
272 in_flight: false,
273 rearmed: false,
274 }
275 }
276}
277
278#[derive(Debug)]
279struct VmState<CapabilityId> {
280 deficit: FairBudget,
281 queued: bool,
282 ready_capabilities: VecDeque<CapabilityId>,
283 capabilities: BTreeMap<CapabilityId, CapabilityState>,
284}
285
286impl<CapabilityId> VmState<CapabilityId> {
287 fn new() -> Self {
288 Self {
289 deficit: FairBudget::default(),
290 queued: false,
291 ready_capabilities: VecDeque::new(),
292 capabilities: BTreeMap::new(),
293 }
294 }
295}
296
297#[derive(Debug)]
298struct InFlight<VmId, CapabilityId> {
299 selection: FairSelection<VmId, CapabilityId>,
300}
301
302#[derive(Debug)]
308pub struct HierarchicalDeficitRoundRobin<VmId, CapabilityId> {
309 config: FairnessConfig,
310 ready_vms: VecDeque<VmId>,
311 vms: BTreeMap<VmId, VmState<CapabilityId>>,
312 in_flight: Option<InFlight<VmId, CapabilityId>>,
313 next_sequence: u64,
314 metrics: Option<RuntimeMetrics>,
315}
316
317#[derive(Clone, Debug)]
325pub struct FairWorkBroker {
326 inner: Arc<FairWorkInner>,
327}
328
329#[derive(Debug)]
330struct FairWorkInner {
331 state: Mutex<FairWorkState>,
332 changed: tokio::sync::Notify,
333}
334
335#[derive(Debug)]
336struct FairWorkState {
337 scheduler: HierarchicalDeficitRoundRobin<u64, u64>,
338 granted: BTreeMap<(u64, u64), FairSelection<u64, u64>>,
339 retired_vm_generations: RetiredIdRanges,
342 retired_capabilities: BTreeMap<u64, RetiredIdRanges>,
346}
347
348#[derive(Debug, Default)]
349struct RetiredIdRanges {
350 ranges: BTreeMap<u64, u64>,
351}
352
353impl RetiredIdRanges {
354 fn contains(&self, id: u64) -> bool {
355 self.ranges
356 .range(..=id)
357 .next_back()
358 .is_some_and(|(_, end)| id <= *end)
359 }
360
361 fn insert(&mut self, id: u64) -> bool {
362 if self.contains(id) {
363 return false;
364 }
365
366 let mut start = id;
367 let mut end = id;
368 if let Some((previous_start, previous_end)) = self
369 .ranges
370 .range(..id)
371 .next_back()
372 .map(|(start, end)| (*start, *end))
373 {
374 if previous_end.checked_add(1) == Some(id) {
375 start = previous_start;
376 self.ranges.remove(&previous_start);
377 }
378 }
379 while let Some((next_start, next_end)) = self
380 .ranges
381 .range(start..)
382 .next()
383 .map(|(start, end)| (*start, *end))
384 {
385 if next_start > end.saturating_add(1) {
386 break;
387 }
388 end = end.max(next_end);
389 self.ranges.remove(&next_start);
390 }
391 self.ranges.insert(start, end);
392 true
393 }
394}
395
396#[derive(Debug)]
400pub struct FairWorkTurn {
401 broker: FairWorkBroker,
402 selection: Option<FairSelection<u64, u64>>,
403 allowance: FairBudget,
404}
405
406struct FairAcquireGuard {
407 broker: FairWorkBroker,
408 key: (u64, u64),
409 armed: bool,
410}
411
412impl FairWorkBroker {
413 pub fn new(config: FairnessConfig, metrics: RuntimeMetrics) -> Result<Self, FairnessError> {
414 Ok(Self {
415 inner: Arc::new(FairWorkInner {
416 state: Mutex::new(FairWorkState {
417 scheduler: HierarchicalDeficitRoundRobin::new_with_metrics(config, metrics)?,
418 granted: BTreeMap::new(),
419 retired_vm_generations: RetiredIdRanges::default(),
420 retired_capabilities: BTreeMap::new(),
421 }),
422 changed: tokio::sync::Notify::new(),
423 }),
424 })
425 }
426
427 pub async fn acquire(
431 &self,
432 vm_generation: u64,
433 capability_id: u64,
434 requested: FairBudget,
435 ) -> Result<FairWorkTurn, FairnessError> {
436 if requested.operations == 0 {
437 return Err(FairnessError::InvalidConfig {
438 field: "limits.reactor.perHandleOperationQuantum",
439 });
440 }
441 if requested.bytes == 0 {
442 return Err(FairnessError::InvalidConfig {
443 field: "limits.reactor.byteQuantum",
444 });
445 }
446
447 let key = (vm_generation, capability_id);
448 let mut cancellation = FairAcquireGuard {
449 broker: self.clone(),
450 key,
451 armed: true,
452 };
453 loop {
454 let mut changed = Box::pin(self.inner.changed.notified());
461 changed.as_mut().enable();
462 let published = {
463 let mut state = self
464 .inner
465 .state
466 .lock()
467 .map_err(|_| FairnessError::Invariant)?;
468 if state.retired_vm_generations.contains(vm_generation) {
469 return Err(FairnessError::CapabilityRetired {
470 vm_generation,
471 capability_id,
472 });
473 }
474 if state
475 .retired_capabilities
476 .get(&vm_generation)
477 .is_some_and(|retired| retired.contains(capability_id))
478 {
479 return Err(FairnessError::CapabilityRetired {
480 vm_generation,
481 capability_id,
482 });
483 }
484 if let Some(selection) = state.granted.remove(&key) {
485 let allowance = selection.allowance.min(requested);
486 cancellation.armed = false;
487 return Ok(FairWorkTurn {
488 broker: self.clone(),
489 selection: Some(selection),
490 allowance,
491 });
492 }
493 state.scheduler.mark_ready(vm_generation, capability_id)?;
500 let published = Self::publish_next_locked(&mut state)?;
501 if let Some(selection) = state.granted.remove(&key) {
502 let allowance = selection.allowance.min(requested);
503 cancellation.armed = false;
504 return Ok(FairWorkTurn {
505 broker: self.clone(),
506 selection: Some(selection),
507 allowance,
508 });
509 }
510 published
511 };
512 if published {
513 self.inner.changed.notify_waiters();
514 }
515 changed.await;
516 }
517 }
518
519 pub fn retire_capability(
524 &self,
525 vm_generation: u64,
526 capability_id: u64,
527 ) -> Result<bool, FairnessError> {
528 let mut state = self
529 .inner
530 .state
531 .lock()
532 .map_err(|_| FairnessError::Invariant)?;
533 if state.retired_vm_generations.contains(vm_generation) {
534 return Ok(false);
535 }
536 let key = (vm_generation, capability_id);
537 let newly_retired = state
538 .retired_capabilities
539 .entry(vm_generation)
540 .or_default()
541 .insert(capability_id);
542 if let Some(selection) = state.granted.remove(&key) {
543 state
544 .scheduler
545 .complete(&selection, FairBudget::default(), false)?;
546 }
547 state.scheduler.clear_ready(&vm_generation, &capability_id);
548 let removed = match state
549 .scheduler
550 .remove_capability(&vm_generation, &capability_id)
551 {
552 Ok(removed) => removed,
553 Err(FairnessError::CapabilityInFlight) => false,
554 Err(error) => return Err(error),
555 };
556 drop(state);
557 self.inner.changed.notify_waiters();
558 Ok(newly_retired || removed)
559 }
560
561 pub fn retire_vm(&self, vm_generation: u64) -> Result<bool, FairnessError> {
562 let mut state = self
563 .inner
564 .state
565 .lock()
566 .map_err(|_| FairnessError::Invariant)?;
567 let newly_retired = state.retired_vm_generations.insert(vm_generation);
568 let granted_key = state
569 .granted
570 .keys()
571 .find(|(generation, _)| *generation == vm_generation)
572 .copied();
573 if let Some(key) = granted_key {
574 let selection = state.granted.remove(&key).ok_or(FairnessError::Invariant)?;
575 state
576 .scheduler
577 .complete(&selection, FairBudget::default(), false)?;
578 }
579 state.scheduler.clear_vm_ready(&vm_generation);
580 let removed = match state.scheduler.remove_vm(&vm_generation) {
581 Ok(removed) => removed,
582 Err(FairnessError::CapabilityInFlight) => false,
583 Err(error) => return Err(error),
584 };
585 state.retired_capabilities.remove(&vm_generation);
586 drop(state);
587 self.inner.changed.notify_waiters();
588 Ok(newly_retired || removed)
589 }
590
591 fn publish_next_locked(state: &mut FairWorkState) -> Result<bool, FairnessError> {
592 if !state.granted.is_empty() {
593 return Ok(false);
594 }
595 let selection = match state.scheduler.select_next() {
596 Ok(Some(selection)) => selection,
597 Ok(None) | Err(FairnessError::SelectionInFlight { .. }) => return Ok(false),
598 Err(error) => return Err(error),
599 };
600 let key = (selection.vm_id, selection.capability_id);
601 if state.granted.insert(key, selection).is_some() {
602 return Err(FairnessError::Invariant);
603 }
604 Ok(true)
605 }
606
607 fn finish(
608 &self,
609 selection: FairSelection<u64, u64>,
610 used: FairBudget,
611 still_ready: bool,
612 ) -> Result<(), FairnessError> {
613 let mut state = self
614 .inner
615 .state
616 .lock()
617 .map_err(|_| FairnessError::Invariant)?;
618 let vm_retired = state.retired_vm_generations.contains(selection.vm_id);
619 let capability_retired = state
620 .retired_capabilities
621 .get(&selection.vm_id)
622 .is_some_and(|retired| retired.contains(selection.capability_id));
623 state.scheduler.complete(
624 &selection,
625 used,
626 still_ready && !vm_retired && !capability_retired,
627 )?;
628 if vm_retired {
629 state.scheduler.remove_vm(&selection.vm_id)?;
630 } else if capability_retired {
631 state
632 .scheduler
633 .remove_capability(&selection.vm_id, &selection.capability_id)?;
634 }
635 let _ = Self::publish_next_locked(&mut state)?;
636 drop(state);
637 self.inner.changed.notify_waiters();
638 Ok(())
639 }
640
641 fn cancel_waiter(&self, key: (u64, u64)) -> Result<(), FairnessError> {
642 let mut state = self
643 .inner
644 .state
645 .lock()
646 .map_err(|_| FairnessError::Invariant)?;
647 if let Some(selection) = state.granted.remove(&key) {
648 state
649 .scheduler
650 .complete(&selection, FairBudget::default(), false)?;
651 } else {
652 state.scheduler.clear_ready(&key.0, &key.1);
653 }
654 if state.retired_vm_generations.contains(key.0) {
655 match state.scheduler.remove_vm(&key.0) {
656 Ok(_) | Err(FairnessError::CapabilityInFlight) => {}
657 Err(error) => return Err(error),
658 }
659 } else if state
660 .retired_capabilities
661 .get(&key.0)
662 .is_some_and(|retired| retired.contains(key.1))
663 {
664 state.scheduler.remove_capability(&key.0, &key.1)?;
665 }
666 let _ = Self::publish_next_locked(&mut state)?;
667 drop(state);
668 self.inner.changed.notify_waiters();
669 Ok(())
670 }
671}
672
673impl FairWorkTurn {
674 pub fn allowance(&self) -> FairBudget {
675 self.allowance
676 }
677
678 pub fn complete(mut self, used: FairBudget, still_ready: bool) -> Result<(), FairnessError> {
679 if used.operations > self.allowance.operations {
680 return Err(FairnessError::OperationBudgetExceeded {
681 used: used.operations,
682 allowance: self.allowance.operations,
683 });
684 }
685 if used.bytes > self.allowance.bytes {
686 return Err(FairnessError::ByteBudgetExceeded {
687 used: used.bytes,
688 allowance: self.allowance.bytes,
689 });
690 }
691 let selection = self.selection.take().ok_or(FairnessError::Invariant)?;
692 self.broker.finish(selection, used, still_ready)
693 }
694}
695
696impl Drop for FairWorkTurn {
697 fn drop(&mut self) {
698 let Some(selection) = self.selection.take() else {
699 return;
700 };
701 if let Err(error) = self.broker.finish(selection, FairBudget::default(), false) {
702 eprintln!("ERR_AGENTOS_FAIRNESS_TURN_DROP: {error}");
703 }
704 }
705}
706
707impl Drop for FairAcquireGuard {
708 fn drop(&mut self) {
709 if self.armed {
710 if let Err(error) = self.broker.cancel_waiter(self.key) {
711 eprintln!("ERR_AGENTOS_FAIRNESS_ACQUIRE_CANCEL: {error}");
712 }
713 }
714 }
715}
716
717impl<VmId, CapabilityId> HierarchicalDeficitRoundRobin<VmId, CapabilityId>
718where
719 VmId: Clone + Ord,
720 CapabilityId: Clone + Ord,
721{
722 pub fn new(config: FairnessConfig) -> Result<Self, FairnessError> {
723 Self::new_inner(config, None)
724 }
725
726 pub fn new_with_metrics(
727 config: FairnessConfig,
728 metrics: RuntimeMetrics,
729 ) -> Result<Self, FairnessError> {
730 Self::new_inner(config, Some(metrics))
731 }
732
733 fn new_inner(
734 config: FairnessConfig,
735 metrics: Option<RuntimeMetrics>,
736 ) -> Result<Self, FairnessError> {
737 Ok(Self {
738 config: config.validate()?,
739 ready_vms: VecDeque::new(),
740 vms: BTreeMap::new(),
741 in_flight: None,
742 next_sequence: 1,
743 metrics,
744 })
745 }
746
747 pub fn config(&self) -> FairnessConfig {
748 self.config
749 }
750
751 pub fn mark_ready(
753 &mut self,
754 vm_id: VmId,
755 capability_id: CapabilityId,
756 ) -> Result<MembershipUpdate, FairnessError> {
757 if !self.vms.contains_key(&vm_id) {
758 if self.vms.len() >= self.config.max_vms {
759 return Err(FairnessError::VmLimit {
760 limit: self.config.max_vms,
761 });
762 }
763 self.vms.insert(vm_id.clone(), VmState::new());
764 }
765
766 let vm = self.vms.get_mut(&vm_id).ok_or(FairnessError::Invariant)?;
767 if !vm.capabilities.contains_key(&capability_id) {
768 if vm.capabilities.len() >= self.config.max_capabilities_per_vm {
769 return Err(FairnessError::CapabilityLimit {
770 limit: self.config.max_capabilities_per_vm,
771 });
772 }
773 vm.capabilities
774 .insert(capability_id.clone(), CapabilityState::new());
775 }
776
777 let capability = vm
778 .capabilities
779 .get_mut(&capability_id)
780 .ok_or(FairnessError::Invariant)?;
781 if capability.in_flight {
782 capability.rearmed = true;
783 if let Some(metrics) = &self.metrics {
784 metrics.record_fairness_yield(FairnessLevel::Capability);
785 }
786 return Ok(MembershipUpdate::Coalesced);
787 }
788 if capability.queued {
789 if let Some(metrics) = &self.metrics {
790 metrics.record_fairness_yield(FairnessLevel::Capability);
791 }
792 return Ok(MembershipUpdate::Coalesced);
793 }
794 capability.queued = true;
795 vm.ready_capabilities.push_back(capability_id);
796 if !vm.queued {
797 vm.queued = true;
798 self.ready_vms.push_back(vm_id);
799 }
800 Ok(MembershipUpdate::Enqueued)
801 }
802
803 pub fn clear_ready(&mut self, vm_id: &VmId, capability_id: &CapabilityId) -> bool {
805 let Some(vm) = self.vms.get_mut(vm_id) else {
806 return false;
807 };
808 let Some(capability) = vm.capabilities.get_mut(capability_id) else {
809 return false;
810 };
811 capability.rearmed = false;
812 if !capability.queued {
813 return capability.in_flight;
814 }
815 capability.queued = false;
816 vm.ready_capabilities
817 .retain(|queued| queued != capability_id);
818 if vm.ready_capabilities.is_empty() && vm.queued {
819 vm.queued = false;
820 self.ready_vms.retain(|queued| queued != vm_id);
821 }
822 true
823 }
824
825 pub fn select_next(
827 &mut self,
828 ) -> Result<Option<FairSelection<VmId, CapabilityId>>, FairnessError> {
829 if let Some(in_flight) = &self.in_flight {
830 return Err(FairnessError::SelectionInFlight {
831 sequence: in_flight.selection.sequence,
832 });
833 }
834
835 while let Some(vm_id) = self.ready_vms.pop_front() {
836 let Some(vm) = self.vms.get_mut(&vm_id) else {
837 continue;
838 };
839 vm.queued = false;
840 while let Some(capability_id) = vm.ready_capabilities.pop_front() {
841 let Some(capability) = vm.capabilities.get_mut(&capability_id) else {
842 continue;
843 };
844 if !capability.queued || capability.in_flight {
845 continue;
846 }
847 let sequence = self.next_sequence;
848 let Some(next_sequence) = self.next_sequence.checked_add(1) else {
849 vm.ready_capabilities.push_front(capability_id);
850 vm.queued = true;
851 self.ready_vms.push_front(vm_id);
852 return Err(FairnessError::SequenceExhausted);
853 };
854 capability.queued = false;
855 capability.in_flight = true;
856 vm.deficit
857 .add_capped(self.config.vm_quantum, self.config.max_vm_deficit);
858 capability.deficit.add_capped(
859 self.config.capability_quantum,
860 self.config.max_capability_deficit,
861 );
862 let allowance = vm
863 .deficit
864 .turn_allowance(capability.deficit, self.config.capability_quantum);
865 self.next_sequence = next_sequence;
866 let selection = FairSelection {
867 sequence,
868 vm_id,
869 capability_id,
870 allowance,
871 };
872 self.in_flight = Some(InFlight {
873 selection: selection.clone(),
874 });
875 return Ok(Some(selection));
876 }
877 }
878 Ok(None)
879 }
880
881 pub fn complete(
883 &mut self,
884 selection: &FairSelection<VmId, CapabilityId>,
885 used: FairBudget,
886 still_ready: bool,
887 ) -> Result<(), FairnessError> {
888 let Some(outstanding) = self
889 .in_flight
890 .as_ref()
891 .map(|in_flight| &in_flight.selection)
892 else {
893 return Err(FairnessError::StaleSelection {
894 supplied: selection.sequence,
895 outstanding: None,
896 });
897 };
898 if outstanding.sequence != selection.sequence
899 || outstanding.vm_id != selection.vm_id
900 || outstanding.capability_id != selection.capability_id
901 {
902 return Err(FairnessError::StaleSelection {
903 supplied: selection.sequence,
904 outstanding: Some(outstanding.sequence),
905 });
906 }
907 let allowance = outstanding.allowance;
908 if used.operations > allowance.operations {
909 return Err(FairnessError::OperationBudgetExceeded {
910 used: used.operations,
911 allowance: allowance.operations,
912 });
913 }
914 if used.bytes > allowance.bytes {
915 return Err(FairnessError::ByteBudgetExceeded {
916 used: used.bytes,
917 allowance: allowance.bytes,
918 });
919 }
920
921 self.in_flight = None;
922 let vm = self
923 .vms
924 .get_mut(&selection.vm_id)
925 .ok_or(FairnessError::Invariant)?;
926 let capability = vm
927 .capabilities
928 .get_mut(&selection.capability_id)
929 .ok_or(FairnessError::Invariant)?;
930 if !capability.in_flight {
931 return Err(FairnessError::Invariant);
932 }
933 vm.deficit.consume(used);
934 capability.deficit.consume(used);
935 capability.in_flight = false;
936 let requeue = still_ready || capability.rearmed;
937 capability.rearmed = false;
938 if requeue && !capability.queued {
939 capability.queued = true;
940 vm.ready_capabilities
941 .push_back(selection.capability_id.clone());
942 }
943 if !vm.ready_capabilities.is_empty() && !vm.queued {
944 vm.queued = true;
945 self.ready_vms.push_back(selection.vm_id.clone());
946 }
947 if requeue
948 || used.operations == selection.allowance.operations
949 || used.bytes == selection.allowance.bytes
950 {
951 if let Some(metrics) = &self.metrics {
952 metrics.record_fairness_yield(FairnessLevel::Capability);
953 }
954 }
955 Ok(())
956 }
957
958 pub fn remove_capability(
959 &mut self,
960 vm_id: &VmId,
961 capability_id: &CapabilityId,
962 ) -> Result<bool, FairnessError> {
963 let Some(vm) = self.vms.get_mut(vm_id) else {
964 return Ok(false);
965 };
966 if vm
967 .capabilities
968 .get(capability_id)
969 .is_some_and(|capability| capability.in_flight)
970 {
971 return Err(FairnessError::CapabilityInFlight);
972 }
973 let removed = vm.capabilities.remove(capability_id).is_some();
974 vm.ready_capabilities
975 .retain(|queued| queued != capability_id);
976 if vm.ready_capabilities.is_empty() && vm.queued {
977 vm.queued = false;
978 self.ready_vms.retain(|queued| queued != vm_id);
979 }
980 Ok(removed)
981 }
982
983 pub fn clear_vm_ready(&mut self, vm_id: &VmId) -> bool {
987 self.ready_vms.retain(|queued| queued != vm_id);
988 let Some(vm) = self.vms.get_mut(vm_id) else {
989 return false;
990 };
991 vm.queued = false;
992 vm.ready_capabilities.clear();
993 for capability in vm.capabilities.values_mut() {
994 capability.queued = false;
995 capability.rearmed = false;
996 }
997 true
998 }
999
1000 pub fn remove_vm(&mut self, vm_id: &VmId) -> Result<bool, FairnessError> {
1001 if self
1002 .in_flight
1003 .as_ref()
1004 .is_some_and(|in_flight| &in_flight.selection.vm_id == vm_id)
1005 {
1006 return Err(FairnessError::CapabilityInFlight);
1007 }
1008 self.ready_vms.retain(|queued| queued != vm_id);
1009 Ok(self.vms.remove(vm_id).is_some())
1010 }
1011
1012 pub fn snapshot(&self, vm_id: &VmId, capability_id: &CapabilityId) -> Option<FairnessSnapshot> {
1013 let vm = self.vms.get(vm_id)?;
1014 let capability = vm.capabilities.get(capability_id)?;
1015 Some(FairnessSnapshot {
1016 vm_deficit: vm.deficit,
1017 capability_deficit: capability.deficit,
1018 vm_queued: vm.queued,
1019 capability_queued: capability.queued,
1020 capability_in_flight: capability.in_flight,
1021 capability_rearmed: capability.rearmed,
1022 })
1023 }
1024
1025 pub fn ready_vm_count(&self) -> usize {
1026 self.ready_vms.len()
1027 }
1028
1029 pub fn ready_capability_count(&self, vm_id: &VmId) -> usize {
1030 self.vms
1031 .get(vm_id)
1032 .map_or(0, |vm| vm.ready_capabilities.len())
1033 }
1034}
1035
1036#[cfg(test)]
1037mod tests {
1038 use super::*;
1039
1040 fn config() -> FairnessConfig {
1041 FairnessConfig {
1042 vm_quantum: FairBudget::new(4, 4_096),
1043 capability_quantum: FairBudget::new(2, 1_024),
1044 max_vm_deficit: FairBudget::new(16, 16_384),
1045 max_capability_deficit: FairBudget::new(8, 4_096),
1046 max_vms: 8,
1047 max_capabilities_per_vm: 8,
1048 }
1049 }
1050
1051 fn broker() -> FairWorkBroker {
1052 FairWorkBroker::new(config(), RuntimeMetrics::new()).expect("fair work broker")
1053 }
1054
1055 #[tokio::test]
1056 async fn broker_applies_vm_ceiling_and_drop_releases_global_turn() {
1057 let broker = broker();
1058 let turn = broker
1059 .acquire(1, 10, FairBudget::new(1, 128))
1060 .await
1061 .expect("first turn");
1062 assert_eq!(turn.allowance(), FairBudget::new(1, 128));
1063 drop(turn);
1064
1065 let next = tokio::time::timeout(
1066 std::time::Duration::from_secs(1),
1067 broker.acquire(2, 20, FairBudget::new(2, 512)),
1068 )
1069 .await
1070 .expect("dropped turn must not strand scheduler")
1071 .expect("next turn");
1072 next.complete(FairBudget::new(1, 64), false)
1073 .expect("complete next turn");
1074 }
1075
1076 #[tokio::test]
1077 async fn cancelled_waiter_removes_membership_and_cannot_strand_scheduler() {
1078 let broker = broker();
1079 let held = broker
1080 .acquire(1, 10, FairBudget::new(1, 128))
1081 .await
1082 .expect("held turn");
1083 let waiting = tokio::spawn({
1084 let broker = broker.clone();
1085 async move { broker.acquire(2, 20, FairBudget::new(1, 128)).await }
1086 });
1087 tokio::task::yield_now().await;
1088 assert!(!waiting.is_finished());
1089 waiting.abort();
1090 assert!(waiting.await.expect_err("waiter cancelled").is_cancelled());
1091 held.complete(FairBudget::new(1, 64), false)
1092 .expect("complete held turn");
1093
1094 let next = tokio::time::timeout(
1095 std::time::Duration::from_secs(1),
1096 broker.acquire(3, 30, FairBudget::new(1, 128)),
1097 )
1098 .await
1099 .expect("cancelled waiter must not retain the grant")
1100 .expect("next turn");
1101 next.complete(FairBudget::new(1, 64), false)
1102 .expect("complete next turn");
1103 }
1104
1105 #[tokio::test]
1106 async fn concurrent_waiters_for_one_full_duplex_capability_each_receive_a_turn() {
1107 let broker = broker();
1108 let held = broker
1109 .acquire(1, 10, FairBudget::new(1, 128))
1110 .await
1111 .expect("hold capability turn");
1112 let waiter = |broker: FairWorkBroker| {
1113 tokio::spawn(async move {
1114 let turn = broker
1115 .acquire(1, 10, FairBudget::new(1, 128))
1116 .await
1117 .expect("full-duplex waiter turn");
1118 turn.complete(FairBudget::new(1, 1), false)
1119 .expect("complete full-duplex waiter turn");
1120 })
1121 };
1122 let reader = waiter(broker.clone());
1123 let writer = waiter(broker.clone());
1124 for _ in 0..4 {
1125 tokio::task::yield_now().await;
1126 }
1127 held.complete(FairBudget::new(1, 1), false)
1128 .expect("release held capability turn");
1129
1130 tokio::time::timeout(std::time::Duration::from_secs(1), async {
1131 reader.await.expect("reader waiter joins");
1132 writer.await.expect("writer waiter joins");
1133 })
1134 .await
1135 .expect("same-capability waiters must not lose their shared ready edge");
1136 }
1137
1138 #[tokio::test]
1139 async fn retiring_an_active_turn_defers_removal_and_rejects_stale_reacquire() {
1140 let broker = broker();
1141 let turn = broker
1142 .acquire(7, 11, FairBudget::new(1, 128))
1143 .await
1144 .expect("active turn");
1145
1146 assert!(broker
1147 .retire_capability(7, 11)
1148 .expect("begin active capability retirement"));
1149 assert!(broker
1150 .inner
1151 .state
1152 .lock()
1153 .expect("fairness state")
1154 .scheduler
1155 .snapshot(&7, &11)
1156 .is_some_and(|snapshot| snapshot.capability_in_flight));
1157
1158 turn.complete(FairBudget::new(1, 64), true)
1159 .expect("settle retired turn");
1160 assert!(broker
1161 .inner
1162 .state
1163 .lock()
1164 .expect("fairness state")
1165 .scheduler
1166 .snapshot(&7, &11)
1167 .is_none());
1168 assert_eq!(
1169 broker
1170 .acquire(7, 11, FairBudget::new(1, 128))
1171 .await
1172 .expect_err("retired capability must not reacquire"),
1173 FairnessError::CapabilityRetired {
1174 vm_generation: 7,
1175 capability_id: 11,
1176 }
1177 );
1178 }
1179
1180 #[tokio::test]
1181 async fn retiring_vm_with_active_turn_defers_removal_and_rejects_generation() {
1182 let broker = broker();
1183 let turn = broker
1184 .acquire(7, 11, FairBudget::new(1, 128))
1185 .await
1186 .expect("active VM turn");
1187
1188 assert!(broker.retire_vm(7).expect("begin active VM retirement"));
1189 {
1190 let state = broker.inner.state.lock().expect("fairness state");
1191 assert!(state.retired_vm_generations.contains(7));
1192 assert!(state
1193 .scheduler
1194 .snapshot(&7, &11)
1195 .is_some_and(|snapshot| snapshot.capability_in_flight));
1196 assert_eq!(state.scheduler.ready_capability_count(&7), 0);
1197 }
1198 assert_eq!(
1199 broker
1200 .acquire(7, 12, FairBudget::new(1, 128))
1201 .await
1202 .expect_err("retired VM generation must reject new capabilities"),
1203 FairnessError::CapabilityRetired {
1204 vm_generation: 7,
1205 capability_id: 12,
1206 }
1207 );
1208
1209 turn.complete(FairBudget::new(1, 64), true)
1210 .expect("settle retired VM turn");
1211 {
1212 let state = broker.inner.state.lock().expect("fairness state");
1213 assert!(state.scheduler.snapshot(&7, &11).is_none());
1214 assert_eq!(state.scheduler.ready_vm_count(), 0);
1215 }
1216 assert_eq!(
1217 broker
1218 .acquire(7, 11, FairBudget::new(1, 128))
1219 .await
1220 .expect_err("settled retired VM must not resurrect"),
1221 FairnessError::CapabilityRetired {
1222 vm_generation: 7,
1223 capability_id: 11,
1224 }
1225 );
1226 }
1227
1228 #[tokio::test]
1229 async fn vm_generation_churn_reclaims_membership_and_compacts_tombstones() {
1230 let bounded = FairnessConfig {
1231 max_vms: 1,
1232 max_capabilities_per_vm: 1,
1233 ..config()
1234 };
1235 let broker = FairWorkBroker::new(bounded, RuntimeMetrics::new()).expect("bounded broker");
1236
1237 for vm_generation in 1..=1_024 {
1238 let turn = broker
1239 .acquire(vm_generation, 1, FairBudget::new(1, 128))
1240 .await
1241 .expect("VM churn turn");
1242 assert!(broker.retire_vm(vm_generation).expect("retire churn VM"));
1243 turn.complete(FairBudget::new(1, 64), true)
1244 .expect("settle churn VM turn");
1245 }
1246
1247 {
1248 let state = broker.inner.state.lock().expect("fairness state");
1249 assert_eq!(state.scheduler.ready_vm_count(), 0);
1250 assert_eq!(
1251 state.retired_vm_generations.ranges,
1252 BTreeMap::from([(1, 1_024)])
1253 );
1254 }
1255 assert_eq!(
1256 broker
1257 .acquire(512, 2, FairBudget::new(1, 128))
1258 .await
1259 .expect_err("retired VM generation must remain tombstoned"),
1260 FairnessError::CapabilityRetired {
1261 vm_generation: 512,
1262 capability_id: 2,
1263 }
1264 );
1265 let next = broker
1266 .acquire(1_025, 1, FairBudget::new(1, 128))
1267 .await
1268 .expect("VM churn must release bounded scheduler membership");
1269 next.complete(FairBudget::new(1, 64), false)
1270 .expect("complete post-churn VM turn");
1271 }
1272
1273 #[tokio::test]
1274 async fn capability_churn_reclaims_scheduler_membership_and_compacts_tombstones() {
1275 let bounded = FairnessConfig {
1276 max_vms: 1,
1277 max_capabilities_per_vm: 1,
1278 ..config()
1279 };
1280 let broker = FairWorkBroker::new(bounded, RuntimeMetrics::new()).expect("bounded broker");
1281
1282 for capability_id in 1..=1_024 {
1283 let turn = broker
1284 .acquire(1, capability_id, FairBudget::new(1, 128))
1285 .await
1286 .expect("churn turn");
1287 turn.complete(FairBudget::new(1, 64), false)
1288 .expect("complete churn turn");
1289 assert!(broker
1290 .retire_capability(1, capability_id)
1291 .expect("retire churn capability"));
1292 }
1293
1294 {
1295 let state = broker.inner.state.lock().expect("fairness state");
1296 assert_eq!(state.scheduler.ready_capability_count(&1), 0);
1297 assert!(state.scheduler.snapshot(&1, &1_024).is_none());
1298 assert_eq!(
1299 state
1300 .retired_capabilities
1301 .get(&1)
1302 .expect("retired VM range")
1303 .ranges,
1304 BTreeMap::from([(1, 1_024)])
1305 );
1306 }
1307
1308 let next = broker
1309 .acquire(1, 1_025, FairBudget::new(1, 128))
1310 .await
1311 .expect("churn must release bounded scheduler membership");
1312 next.complete(FairBudget::new(1, 64), false)
1313 .expect("complete post-churn turn");
1314 }
1315
1316 #[test]
1317 fn hot_vm_and_hot_capability_rotate_deterministically() {
1318 let mut scheduler = HierarchicalDeficitRoundRobin::new(config()).expect("scheduler");
1319 assert_eq!(
1320 scheduler.mark_ready("vm-a", 1),
1321 Ok(MembershipUpdate::Enqueued)
1322 );
1323 assert_eq!(
1324 scheduler.mark_ready("vm-a", 2),
1325 Ok(MembershipUpdate::Enqueued)
1326 );
1327 assert_eq!(
1328 scheduler.mark_ready("vm-b", 7),
1329 Ok(MembershipUpdate::Enqueued)
1330 );
1331
1332 let mut order = Vec::new();
1333 for _ in 0..12 {
1334 let selection = scheduler
1335 .select_next()
1336 .expect("select")
1337 .expect("ready selection");
1338 assert_eq!(selection.sequence, order.len() as u64 + 1);
1339 order.push((selection.vm_id, selection.capability_id));
1340 scheduler
1341 .complete(&selection, selection.allowance, true)
1342 .expect("complete");
1343 }
1344
1345 assert_eq!(
1346 order,
1347 vec![
1348 ("vm-a", 1),
1349 ("vm-b", 7),
1350 ("vm-a", 2),
1351 ("vm-b", 7),
1352 ("vm-a", 1),
1353 ("vm-b", 7),
1354 ("vm-a", 2),
1355 ("vm-b", 7),
1356 ("vm-a", 1),
1357 ("vm-b", 7),
1358 ("vm-a", 2),
1359 ("vm-b", 7),
1360 ]
1361 );
1362 }
1363
1364 #[test]
1365 fn duplicate_membership_is_coalesced() {
1366 let mut scheduler = HierarchicalDeficitRoundRobin::new(config()).expect("scheduler");
1367 assert_eq!(scheduler.mark_ready(1, 10), Ok(MembershipUpdate::Enqueued));
1368 for _ in 0..10_000 {
1369 assert_eq!(scheduler.mark_ready(1, 10), Ok(MembershipUpdate::Coalesced));
1370 }
1371 assert_eq!(scheduler.ready_vm_count(), 1);
1372 assert_eq!(scheduler.ready_capability_count(&1), 1);
1373 let selection = scheduler.select_next().expect("select").expect("selection");
1374 for _ in 0..10_000 {
1375 assert_eq!(scheduler.mark_ready(1, 10), Ok(MembershipUpdate::Coalesced));
1376 }
1377 scheduler
1378 .complete(&selection, FairBudget::new(1, 1), false)
1379 .expect("complete");
1380 assert_eq!(scheduler.ready_vm_count(), 1);
1381 assert_eq!(scheduler.ready_capability_count(&1), 1);
1382 let rearmed = scheduler
1383 .select_next()
1384 .expect("select rearmed")
1385 .expect("rearmed selection");
1386 scheduler
1387 .complete(&rearmed, FairBudget::new(1, 1), false)
1388 .expect("complete rearmed");
1389 assert!(scheduler.select_next().expect("empty selection").is_none());
1390 }
1391
1392 #[test]
1393 fn membership_state_is_bounded() {
1394 let bounded = FairnessConfig {
1395 max_vms: 1,
1396 max_capabilities_per_vm: 1,
1397 ..config()
1398 };
1399 let mut scheduler = HierarchicalDeficitRoundRobin::new(bounded).expect("scheduler");
1400 scheduler.mark_ready("vm-a", 1).expect("first member");
1401 assert_eq!(
1402 scheduler.mark_ready("vm-a", 2),
1403 Err(FairnessError::CapabilityLimit { limit: 1 })
1404 );
1405 assert_eq!(
1406 scheduler.mark_ready("vm-b", 1),
1407 Err(FairnessError::VmLimit { limit: 1 })
1408 );
1409 }
1410
1411 #[test]
1412 fn idle_credit_is_capped_and_does_not_expand_a_turn() {
1413 let mut scheduler = HierarchicalDeficitRoundRobin::new(config()).expect("scheduler");
1414 scheduler.mark_ready("idle", 1).expect("mark idle");
1415 for _ in 0..100 {
1416 let selection = scheduler
1417 .select_next()
1418 .expect("select idle")
1419 .expect("idle selection");
1420 scheduler
1421 .complete(&selection, FairBudget::default(), true)
1422 .expect("bank bounded credit");
1423 }
1424 let selection = scheduler
1425 .select_next()
1426 .expect("select idle")
1427 .expect("idle selection");
1428 scheduler
1429 .complete(&selection, FairBudget::default(), false)
1430 .expect("make idle");
1431
1432 let snapshot = scheduler.snapshot(&"idle", &1).expect("idle snapshot");
1433 assert_eq!(snapshot.vm_deficit, config().max_vm_deficit);
1434 assert_eq!(snapshot.capability_deficit, config().max_capability_deficit);
1435
1436 scheduler.mark_ready("hot", 2).expect("mark hot");
1437 scheduler.mark_ready("idle", 1).expect("wake idle");
1438 let hot = scheduler
1439 .select_next()
1440 .expect("select hot")
1441 .expect("hot selection");
1442 assert_eq!(hot.vm_id, "hot");
1443 scheduler
1444 .complete(&hot, hot.allowance, true)
1445 .expect("complete hot");
1446 let idle = scheduler
1447 .select_next()
1448 .expect("select idle")
1449 .expect("idle selection");
1450 assert_eq!(idle.vm_id, "idle");
1451 assert_eq!(idle.allowance, config().capability_quantum);
1452 let capped = scheduler.snapshot(&"idle", &1).expect("capped snapshot");
1453 assert_eq!(capped.vm_deficit, config().max_vm_deficit);
1454 assert_eq!(capped.capability_deficit, config().max_capability_deficit);
1455 }
1456
1457 #[test]
1458 fn completion_enforces_both_budget_dimensions() {
1459 let mut scheduler = HierarchicalDeficitRoundRobin::new(config()).expect("scheduler");
1460 scheduler.mark_ready(1, 1).expect("mark ready");
1461 let selection = scheduler.select_next().expect("select").expect("selection");
1462
1463 assert!(matches!(
1464 scheduler.complete(
1465 &selection,
1466 FairBudget::new(selection.allowance.operations + 1, 0),
1467 true,
1468 ),
1469 Err(FairnessError::OperationBudgetExceeded { .. })
1470 ));
1471 assert!(matches!(
1472 scheduler.complete(
1473 &selection,
1474 FairBudget::new(0, selection.allowance.bytes + 1),
1475 true,
1476 ),
1477 Err(FairnessError::ByteBudgetExceeded { .. })
1478 ));
1479 let mut forged = selection.clone();
1480 forged.allowance = FairBudget::new(usize::MAX, usize::MAX);
1481 assert!(matches!(
1482 scheduler.complete(
1483 &forged,
1484 FairBudget::new(selection.allowance.operations + 1, 0),
1485 true,
1486 ),
1487 Err(FairnessError::OperationBudgetExceeded { .. })
1488 ));
1489 scheduler
1490 .complete(&selection, selection.allowance, false)
1491 .expect("valid completion remains possible");
1492 }
1493}