Skip to main content

agentos_runtime/
fairness.rs

1//! Deterministic hierarchical deficit round robin for VM-owned capabilities.
2
3use std::collections::{BTreeMap, VecDeque};
4use std::fmt;
5use std::sync::{Arc, Mutex};
6
7use crate::metrics::{FairnessLevel, RuntimeMetrics};
8
9/// Independent count and byte dimensions for one scheduling turn.
10#[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/// Bounded policy for the process -> VM -> capability scheduler.
53#[derive(Clone, Copy, Debug, Eq, PartialEq)]
54pub struct FairnessConfig {
55    /// Credit added whenever a ready VM reaches the front of the process ring.
56    pub vm_quantum: FairBudget,
57    /// Credit added whenever a ready capability reaches the front of its VM ring.
58    pub capability_quantum: FairBudget,
59    /// Maximum credit a VM may retain between selections, including while idle.
60    pub max_vm_deficit: FairBudget,
61    /// Maximum credit a capability may retain between selections, including while idle.
62    pub max_capability_deficit: FairBudget,
63    /// Maximum number of VM states retained by the scheduler.
64    pub max_vms: usize,
65    /// Maximum number of capability states retained for one VM.
66    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/// One deterministic grant. Exactly one grant may be outstanding at a time.
142#[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    /// Hard per-turn ceiling in both dimensions.
148    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/// Process-level VM rotation containing a second capability rotation per VM.
303///
304/// The scheduler is intentionally serial: callers select one bounded turn,
305/// perform no more than its allowance, then report actual use. This makes the
306/// selection order reproducible in tests and keeps requeueing atomic.
307#[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/// Process-owned async admission facade over the deterministic HDRR state.
318///
319/// There is exactly one broker per [`crate::SidecarRuntime`]. A ready handle
320/// first joins the coalesced VM/capability rings, then receives one bounded
321/// turn. At most one grant is outstanding process-wide, which makes completion
322/// and requeueing atomic and prevents independent handle tasks from bypassing
323/// tenant rotation by racing on Tokio workers.
324#[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    /// VM generations are monotonic process-wide. Retain compact ranges so a
340    /// stale task can never recreate scheduler membership after VM teardown.
341    retired_vm_generations: RetiredIdRanges,
342    /// Capability IDs are monotonic within one VM generation. Merged retired
343    /// ranges prevent a late transport task from recreating scheduler state
344    /// without retaining one tombstone per closed capability.
345    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/// One process-fair work turn. Dropping an unfinished turn reconciles the
397/// scheduler with zero work and without requeueing, so task cancellation can
398/// never strand the global scheduler in an in-flight state.
399#[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    /// Wait for this ready capability's next process-fair turn. `requested`
428    /// carries the VM's configured per-turn ceilings; the returned allowance
429    /// is the lower of those ceilings and the process scheduler's grant.
430    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            // Arm before observing state so a completion between the state
455            // probe and await cannot be lost.
456            // `Notify::notified()` does not register its waiter until the
457            // future is first polled. Enable it before observing scheduler
458            // state, otherwise another completion can notify between the
459            // state probe and `.await`, permanently stranding this acquire.
460            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                // Multiple protocol tasks for one full-duplex description can
494                // wait on the same capability concurrently (for example, one
495                // reader and one writer). A grant is consumed by only one of
496                // them, so every still-waiting acquire must reassert the
497                // coalesced ready edge after it wakes. `mark_ready` is
498                // idempotent and rearms an in-flight capability.
499                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    /// Permanently retire one monotonic capability identity for this VM
520    /// generation. Queued or granted work is revoked immediately. If native
521    /// work already owns the bounded turn, its completion removes membership;
522    /// later acquire attempts fail instead of recreating stale scheduler state.
523    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    /// Insert one ready membership edge. Repeated marks never add queue entries.
752    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    /// Remove queued readiness without discarding retained, capped deficit.
804    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    /// Choose the next VM and capability in activation order.
826    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    /// Charge actual work and atomically requeue a source that remains ready.
882    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    /// Revoke every queued/rearmed turn for a VM without disturbing a turn
984    /// already issued to native work. VM retirement uses this before attempting
985    /// removal so no queued capability can run while issued work settles.
986    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}