Skip to main content

mkit_server/timers/
budget.rs

1//! One timer allowance shared by every logical partition in a physical tick.
2
3use super::{TickBudget, TimerKind};
4use crate::rt::Clock;
5
6/// Physical-tick counters. Fetching a page reserves its rows before processing
7/// them, so an exhausted scan allowance does not discard an already read page.
8#[derive(Debug)]
9pub struct TickState {
10    budget: TickBudget,
11    start: i64,
12    /// Rows fetched, including unknown, deferred and lookahead rows.
13    pub examined: u32,
14    /// Successfully committed timer batches.
15    pub committed: u32,
16    /// Claimed handler invocations, including failures and races.
17    pub attempted: u32,
18    per_kind: [u32; 256],
19}
20
21impl TickState {
22    /// Begin one physical tick. Reuse this state for all its logical partitions.
23    #[must_use]
24    pub fn new(clock: &dyn Clock, budget: TickBudget) -> Self {
25        Self {
26            budget,
27            start: clock.now_ms(),
28            examined: 0,
29            committed: 0,
30            attempted: 0,
31            per_kind: [0; 256],
32        }
33    }
34
35    /// Rows that another page may fetch.
36    #[must_use]
37    pub fn remaining_scanned(&self) -> u32 {
38        self.budget.max_scanned.saturating_sub(self.examined)
39    }
40
41    /// Whether another page or partition must wait for a later physical tick.
42    #[must_use]
43    pub fn exhausted(&self, clock: &dyn Clock) -> bool {
44        self.remaining_scanned() == 0 || self.work_exhausted(clock)
45    }
46
47    /// Whether processing must stop, including inside an already reserved page.
48    /// The scan allowance is excluded because those rows were charged on fetch.
49    #[must_use]
50    pub fn work_exhausted(&self, clock: &dyn Clock) -> bool {
51        let elapsed = u64::try_from(clock.now_ms().saturating_sub(self.start)).unwrap_or(0);
52        self.committed >= self.budget.max_fired || elapsed >= self.budget.max_elapsed_ms
53    }
54
55    /// Reserve exactly the number of fetched rows. An oversized charge changes
56    /// no counters; the caller must limit its query to the remaining allowance.
57    #[must_use]
58    pub fn charge_scan(&mut self, rows: u32) -> bool {
59        if rows > self.remaining_scanned() {
60            return false;
61        }
62        self.examined += rows;
63        true
64    }
65
66    /// Claim one handler invocation before it runs. A handler can lower its
67    /// kind's allowance, but cannot enlarge the shared physical-tick ceiling.
68    #[must_use]
69    pub fn claim_attempt(&mut self, kind: TimerKind, handler_cap: Option<u32>) -> bool {
70        let cap = handler_cap
71            .unwrap_or(self.budget.max_per_kind)
72            .max(1)
73            .min(self.budget.max_per_kind);
74        let count = &mut self.per_kind[usize::from(kind.get())];
75        if *count >= cap {
76            return false;
77        }
78        *count += 1;
79        self.attempted = self.attempted.saturating_add(1);
80        true
81    }
82
83    /// Record one successful guarded commit.
84    pub fn committed(&mut self) {
85        self.committed = self.committed.saturating_add(1);
86    }
87
88    /// Invocations already claimed by this kind, for verification and reporting.
89    #[must_use]
90    pub fn attempted_for(&self, kind: TimerKind) -> u32 {
91        self.per_kind[usize::from(kind.get())]
92    }
93}
94
95#[cfg(test)]
96mod tests {
97    use super::*;
98    use crate::rt::ManualClock;
99
100    #[test]
101    fn frozen_clock_partitions_share_scans_and_failed_attempts() {
102        let clock = ManualClock::new(100);
103        let mut state = TickState::new(&clock, TickBudget::default());
104        let kind = TimerKind::new(12);
105        let mut attempts = 0;
106        for _partition in 0..16 {
107            assert!(state.charge_scan(32));
108            for _row in 0..32 {
109                // A failed handler still consumes its invocation. No commit
110                // is recorded, so failures cannot refresh the kind allowance.
111                attempts += u32::from(state.claim_attempt(kind, Some(128)));
112            }
113        }
114        assert_eq!(
115            (attempts, state.attempted, state.attempted_for(kind)),
116            (32, 32, 32)
117        );
118        assert_eq!((state.examined, state.committed), (512, 0));
119        assert!(state.exhausted(&clock));
120        assert!(!state.work_exhausted(&clock));
121        assert!(!state.charge_scan(1));
122        assert_eq!(state.examined, 512);
123        // Rows in the last reserved page may still invoke another kind.
124        assert!(state.claim_attempt(TimerKind::new(8), None));
125    }
126
127    #[test]
128    fn handler_cap_can_lower_but_cannot_raise_the_shared_kind_limit() {
129        let clock = ManualClock::new(0);
130        let mut state = TickState::new(&clock, TickBudget::new(128, 3, 512, 10_000));
131        let limited = TimerKind::new(1);
132        assert!(state.claim_attempt(limited, Some(0)));
133        assert!(!state.claim_attempt(limited, Some(0)));
134        let default = TimerKind::new(2);
135        let raised = TimerKind::new(3);
136        for _ in 0..3 {
137            assert!(state.claim_attempt(default, None));
138            assert!(state.claim_attempt(raised, Some(u32::MAX)));
139        }
140        assert!(!state.claim_attempt(default, None));
141        assert!(!state.claim_attempt(raised, Some(u32::MAX)));
142        assert_eq!(state.attempted, 7);
143    }
144
145    #[test]
146    fn final_reserved_page_can_finish_until_the_commit_ceiling() {
147        let clock = ManualClock::new(0);
148        let mut state = TickState::new(&clock, TickBudget::new(2, 32, 3, 100));
149        assert!(!state.charge_scan(4));
150        assert_eq!(state.examined, 0);
151        assert!(state.charge_scan(3));
152        assert!(state.exhausted(&clock));
153        assert!(!state.work_exhausted(&clock));
154        state.committed();
155        assert!(!state.work_exhausted(&clock));
156        state.committed();
157        assert!(state.work_exhausted(&clock));
158        assert_eq!(state.committed, 2);
159    }
160
161    #[test]
162    fn clock_deadline_is_shared_and_stops_at_the_exact_boundary() {
163        let clock = ManualClock::new(100);
164        let state = TickState::new(&clock, TickBudget::new(128, 32, 512, 10));
165        clock.set(99);
166        assert!(!state.work_exhausted(&clock));
167        clock.set(109);
168        assert!(!state.exhausted(&clock));
169        clock.set(110);
170        assert!(state.work_exhausted(&clock));
171        assert!(state.exhausted(&clock));
172        assert_eq!(state.remaining_scanned(), 512);
173    }
174}