Skip to main content

mkit_server/timers/
registry.rs

1//! Timer kind allocations and startup handler registration.
2
3use super::{DueTimer, Fired, TimerCtx};
4use crate::rt::{BoxFuture, MaybeSend, MaybeSync};
5use crate::store::{NamespaceStore, StoreError};
6use std::collections::BTreeMap;
7
8/// A timer's stable codec and handler identifier.
9#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
10pub struct TimerKind(u8);
11impl TimerKind {
12    /// Wrap a kind number. Zero cannot be registered.
13    #[must_use]
14    pub const fn new(kind: u8) -> Self {
15        Self(kind)
16    }
17    /// The encoded number.
18    #[must_use]
19    pub fn get(self) -> u8 {
20        self.0
21    }
22}
23
24/// Kind allocations. A new kind takes the next free number; shipped kinds are never reused.
25/// R-190 allocates kind 15 to checkpointed takedown and preservation work.
26///
27/// | Numbers | Allocation |
28/// |---|---|
29/// | 0 | Invalid |
30/// | 1 | LEASE_SWEEP (WP-1.25) |
31/// | 2 | TICKET_EXPIRY (handler added in WP-1.14) |
32/// | 3 | RELAY (WP-1.23a) |
33/// | 4 | BACKUP (Worker only, WP-1.29b) |
34/// | 5 | QUOTA_ROLLUP (WP-1.26a) |
35/// | 6 | Reserved |
36/// | 7 | VERIFY (scheduled indexed verification, WP-4.8) |
37/// | 8 | OUTCOME_DELIVERY (WP-3.3) |
38/// | 9 | RESERVATION_RECONCILE (WP-3.3) |
39/// | 10 | PUBLISHED_VIEW (Worker only, WP-1.21) |
40/// | 11 | CACHE_PURGE (WP-5.10) |
41/// | 12 | PUBLICATION_RECHECK (WP-5.4, R-182) |
42/// | 13 | CONTENT_TAKEDOWN_REQUEST (WP-4.10b, R-186) |
43/// | 14 | INSPECTION (reserved for WP-5.5a, R-198) |
44/// | 15 | TAKEDOWN_WORK (WP-5.6a, R-190) |
45/// | 16..=0xEF | Production, unallocated |
46/// | 0xF0..=0xFE | Reserved for tests |
47/// | 0xFF | TEST (`test-faults` only) |
48pub mod kinds {
49    /// Expired coordinator epoch-lease table rows.
50    pub const LEASE_SWEEP: super::TimerKind = super::TimerKind::new(1);
51    /// Ticket expiry.
52    pub const TICKET_EXPIRY: super::TimerKind = super::TimerKind::new(2);
53    /// Source-side outbox delivery.
54    pub const RELAY: super::TimerKind = super::TimerKind::new(3);
55    /// Per-partition Worker snapshot export.
56    pub const BACKUP: super::TimerKind = super::TimerKind::new(4);
57    /// Ref-shard namespace quota reconciliation.
58    pub const QUOTA_ROLLUP: super::TimerKind = super::TimerKind::new(5);
59    /// Checkpointed slices of a ticketed pack's scheduled verification.
60    pub const VERIFY: super::TimerKind = super::TimerKind::new(7);
61    /// Deliver durable terminal outcomes.
62    pub const OUTCOME_DELIVERY: super::TimerKind = super::TimerKind::new(8);
63    /// Settle abandoned pending reservations.
64    pub const RESERVATION_RECONCILE: super::TimerKind = super::TimerKind::new(9);
65    /// Published ref-index snapshots (explicit Worker opt-in only).
66    pub const PUBLISHED_VIEW: super::TimerKind = super::TimerKind::new(10);
67    /// Materialize a durable late-holder takedown handoff; not takedown completion.
68    pub const CONTENT_TAKEDOWN_REQUEST: super::TimerKind = super::TimerKind::new(13);
69    /// Durable local and shared cache purge.
70    pub const CACHE_PURGE: super::TimerKind = super::TimerKind::new(11);
71    /// Retained inspection and published-membership dependency clearance.
72    pub const PUBLICATION_RECHECK: super::TimerKind = super::TimerKind::new(12);
73    /// Preservation acquisition, holder discovery and audited retention purge.
74    pub const TAKEDOWN_WORK: super::TimerKind = super::TimerKind::new(15);
75    /// Ref deletion used only by test drivers and directives.
76    #[cfg(feature = "test-faults")]
77    pub const TEST: super::TimerKind = super::TimerKind::new(0xFF);
78}
79
80/// Handles a due timer in its own partition.
81///
82/// `fire` may run more than once for the same timer. Every effect outside
83/// the returned batch MUST be idempotent. Effects inside the batch are
84/// applied at most once per timer row (guarded by the row's original value).
85/// Allow room for the core's Equals/Delete and a reschedule Absent/Put
86/// within the store's batch limits. Cross-partition effects require an outbox.
87pub trait TimerHandler<S: NamespaceStore>: MaybeSend + MaybeSync {
88    /// The kind this handler decodes.
89    fn kind(&self) -> TimerKind;
90    /// Lowers the shared `TickBudget::max_per_kind` allowance for this kind.
91    fn max_per_tick(&self) -> Option<u32> {
92        None
93    }
94    /// Prepare effects; the core atomically guards and removes the timer.
95    fn fire<'a>(
96        &'a self,
97        ctx: &'a TimerCtx<'a, S>,
98        timer: &'a DueTimer,
99    ) -> BoxFuture<'a, Result<Fired, StoreError>>;
100}
101
102/// Handlers registered once at driver startup, indexed by stable kind number.
103pub struct TimerRegistry<'a, S> {
104    handlers: BTreeMap<TimerKind, Box<dyn TimerHandler<S> + 'a>>,
105}
106impl<'a, S: NamespaceStore> TimerRegistry<'a, S> {
107    /// An empty registry; unknown kinds remain stored for newer binaries.
108    #[must_use]
109    pub fn new() -> Self {
110        Self {
111            handlers: BTreeMap::new(),
112        }
113    }
114    /// Register a handler.
115    ///
116    /// # Panics
117    /// If its kind is already registered or is zero (a startup programming error).
118    #[must_use]
119    pub fn register(mut self, handler: impl TimerHandler<S> + 'a) -> Self {
120        let kind = handler.kind();
121        assert_ne!(kind.get(), 0, "timer kind zero is invalid");
122        assert!(
123            !self.handlers.contains_key(&kind),
124            "timer kind already registered"
125        );
126        self.handlers.insert(kind, Box::new(handler));
127        self
128    }
129    pub(super) fn get(&self, kind: TimerKind) -> Option<&dyn TimerHandler<S>> {
130        self.handlers.get(&kind).map(Box::as_ref)
131    }
132}
133impl<S: NamespaceStore> Default for TimerRegistry<'_, S> {
134    fn default() -> Self {
135        Self::new()
136    }
137}
138
139impl<S> core::fmt::Debug for TimerRegistry<'_, S> {
140    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
141        f.debug_struct("TimerRegistry")
142            .field("kinds", &self.handlers.keys().collect::<Vec<_>>())
143            .finish()
144    }
145}