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}