Skip to main content

mkit_server/timers/
mod.rs

1//! Partition-local timers, shared by native and Durable Object drivers.
2
3mod budget;
4pub use budget::TickState;
5
6pub mod lease_sweep;
7pub mod outcome_delivery;
8pub mod publication_recheck;
9pub mod quota_rollup;
10pub mod registry;
11pub mod reservation_reconcile;
12#[cfg(feature = "test-faults")]
13pub mod test_kind;
14#[cfg(test)]
15mod tests;
16pub mod ticket_expiry;
17
18use crate::rt::Clock;
19use crate::store::{
20    Batch, BatchOutcome, Key, NamespaceStore, Partition, Precondition, StoreError, Value, Write,
21    keys,
22};
23use bytes::Bytes;
24pub use registry::{TimerHandler, TimerKind, TimerRegistry};
25
26/// Delay before retrying failed or unknown timers, avoiding a busy loop.
27pub const RETRY_BACKOFF_MS: u64 = 5_000;
28/// Largest persisted retry delay, including after an isolate restart.
29pub const MAX_RETRY_BACKOFF_MS: u64 = 600_000;
30const PAGE_SIZE: u32 = 64;
31
32/// The row handed to a kind-specific codec and handler.
33#[non_exhaustive]
34#[derive(Debug, Clone)]
35pub struct DueTimer {
36    /// Scheduled Unix epoch milliseconds.
37    pub due_at_ms: u64,
38    /// Stable handler identifier.
39    pub kind: TimerKind,
40    /// Opaque kind-specific identity.
41    pub reference: Bytes,
42    /// Opaque kind-specific payload.
43    pub value: Value,
44}
45/// Reads available to a handler; effects belong in the returned batch.
46#[non_exhaustive]
47#[derive(Debug)]
48pub struct TimerCtx<'a, S> {
49    /// The driver's concrete store.
50    pub store: &'a S,
51    /// Atomicity boundary for every returned effect.
52    pub partition: &'a Partition,
53    /// Business time for this tick.
54    pub now_ms: u64,
55}
56/// A handler's decision, committed by the core.
57#[non_exhaustive]
58#[derive(Debug)]
59pub enum Fired {
60    /// Commit `batch` and delete the timer, atomically.
61    Done(Batch),
62    /// Commit effects and move the timer atomically.
63    Reschedule {
64        /// New Unix epoch milliseconds; must change the timer key.
65        due_at_ms: u64,
66        /// New kind-specific payload.
67        value: Value,
68        /// Effects in this partition.
69        batch: Batch,
70    },
71    /// No effect now; try again later (counts as a failure for backoff).
72    Retry,
73}
74/// Work limits shared by every logical partition of one physical tick.
75#[non_exhaustive]
76#[derive(Debug, Clone, Copy)]
77pub struct TickBudget {
78    /// Maximum successful commits.
79    pub max_fired: u32,
80    /// Maximum handler invocations per kind, including failed attempts.
81    pub max_per_kind: u32,
82    /// Maximum examined rows, including unknown and deferred kinds.
83    pub max_scanned: u32,
84    /// Maximum elapsed injected-clock milliseconds.
85    pub max_elapsed_ms: u64,
86}
87impl TickBudget {
88    /// Construct limits, clamping each to at least one.
89    #[must_use]
90    pub const fn new(
91        max_fired: u32,
92        max_per_kind: u32,
93        max_scanned: u32,
94        max_elapsed_ms: u64,
95    ) -> Self {
96        Self {
97            max_fired: if max_fired == 0 { 1 } else { max_fired },
98            max_per_kind: if max_per_kind == 0 { 1 } else { max_per_kind },
99            max_scanned: if max_scanned == 0 { 1 } else { max_scanned },
100            max_elapsed_ms: if max_elapsed_ms == 0 {
101                1
102            } else {
103                max_elapsed_ms
104            },
105        }
106    }
107}
108impl Default for TickBudget {
109    fn default() -> Self {
110        Self::new(128, 32, 512, 10_000)
111    }
112}
113/// Counts and the next driver wake, never earlier than this tick's business time.
114#[non_exhaustive]
115#[derive(Debug, Clone, PartialEq, Eq, Default)]
116pub struct RunReport {
117    /// Successfully committed timer batches.
118    pub fired: u32,
119    /// Batches lost to a precondition race.
120    pub raced: u32,
121    /// Handler or commit failures.
122    pub failed: u32,
123    /// Rows owned by unregistered kinds.
124    pub unknown: u32,
125    /// Rows skipped by a kind's invocation cap.
126    pub deferred: u32,
127    /// Examined due rows.
128    pub scanned: u32,
129    /// A global work limit stopped the tick.
130    pub stopped_on_budget: bool,
131    /// Next scheduled wake for this partition.
132    pub next_wake_ms: Option<u64>,
133}
134
135/// Earliest timer Put in a batch, shared by both driver adapters.
136#[must_use]
137pub fn earliest_timer_put(batch: &Batch) -> Option<u64> {
138    batch
139        .writes
140        .iter()
141        .filter_map(|write| match write {
142            Write::Put(key, _) => match keys::parse(key) {
143                Some(keys::ParsedKey::Timer { due_at_ms, .. }) => Some(due_at_ms),
144                _ => None,
145            },
146            Write::Delete(_) => None,
147        })
148        .min()
149}
150fn min_due(a: Option<u64>, b: Option<u64>) -> Option<u64> {
151    match (a, b) {
152        (Some(a), Some(b)) => Some(a.min(b)),
153        (a, b) => a.or(b),
154    }
155}
156fn time_prefix(now: u64) -> Key {
157    Key::new([&b"w\0"[..], &now.to_be_bytes()].concat())
158}
159enum FireOutcome {
160    Committed(Option<u64>),
161    Raced,
162    Failed,
163}
164
165async fn fire_timer<S: NamespaceStore>(
166    handler: &dyn TimerHandler<S>,
167    ctx: &TimerCtx<'_, S>,
168    timer: &DueTimer,
169    key: Key,
170) -> FireOutcome {
171    let batch = match handler.fire(ctx, timer).await {
172        Ok(Fired::Done(batch)) => batch
173            .require(Precondition::Equals(key.clone(), timer.value.clone()))
174            .delete(key),
175        Ok(Fired::Reschedule {
176            due_at_ms,
177            value,
178            batch,
179        }) => {
180            let new_key = keys::timer(due_at_ms, timer.kind.get(), &timer.reference);
181            if due_at_ms == timer.due_at_ms || new_key == key {
182                return FireOutcome::Failed;
183            }
184            batch
185                .require(Precondition::Equals(key.clone(), timer.value.clone()))
186                .require(Precondition::Absent(new_key.clone()))
187                .delete(key)
188                .put(new_key, value)
189        }
190        Ok(Fired::Retry) => {
191            #[cfg(feature = "test-faults")]
192            tracing::warn!(kind = timer.kind.get(), "test timer requested retry");
193            return FireOutcome::Failed;
194        }
195        Err(error) => {
196            #[cfg(feature = "test-faults")]
197            tracing::warn!(kind = timer.kind.get(), %error, "test timer handler failed");
198            #[cfg(not(feature = "test-faults"))]
199            let _ = error;
200            return FireOutcome::Failed;
201        }
202    };
203    let put_due = earliest_timer_put(&batch);
204    match ctx.store.apply(ctx.partition, batch).await {
205        Ok(BatchOutcome::Committed) => FireOutcome::Committed(put_due),
206        Ok(BatchOutcome::PreconditionFailed { .. }) => FireOutcome::Raced,
207        Ok(BatchOutcome::DeadlinePassed { .. }) => {
208            #[cfg(feature = "test-faults")]
209            tracing::warn!(kind = timer.kind.get(), "test timer deadline passed");
210            FireOutcome::Failed
211        }
212        Err(error) => {
213            #[cfg(feature = "test-faults")]
214            tracing::warn!(kind = timer.kind.get(), %error, "test timer apply failed");
215            #[cfg(not(feature = "test-faults"))]
216            let _ = error;
217            FireOutcome::Failed
218        }
219    }
220}
221
222/// Fire one partition with its own allowance. Physical drivers should use
223/// [`run_due_with_state`] and share a single [`TickState`] across their heads.
224///
225/// # Errors
226/// Scan errors escape. Handler/apply failures remain durable timer work.
227pub async fn run_due<S: NamespaceStore>(
228    store: &S,
229    p: &Partition,
230    registry: &TimerRegistry<'_, S>,
231    clock: &dyn Clock,
232    now_ms: u64,
233    budget: &TickBudget,
234) -> Result<RunReport, StoreError> {
235    let mut state = TickState::new(clock, *budget);
236    run_due_with_state(store, p, registry, clock, now_ms, &mut state).await
237}
238
239/// Fire a partition without refreshing physical-alarm limits. Scans reserve
240/// their returned rows and one portable SQL lookahead row before processing.
241/// Retry moves preserve the handler's original due time and opaque payload.
242///
243/// # Errors
244/// Only scan errors escape; failed retry moves retain the guarded original row.
245pub async fn run_due_with_state<S: NamespaceStore>(
246    store: &S,
247    p: &Partition,
248    registry: &TimerRegistry<'_, S>,
249    clock: &dyn Clock,
250    now_ms: u64,
251    state: &mut TickState,
252) -> Result<RunReport, StoreError> {
253    let (start, class_end) = keys::class_range(keys::TAG_TIMER);
254    let end = now_ms
255        .checked_add(1)
256        .map_or_else(|| class_end.clone(), time_prefix);
257    let ctx = TimerCtx {
258        store,
259        partition: p,
260        now_ms,
261    };
262    let mut run = PartitionRun::new();
263    let mut cursor = None;
264    'pages: loop {
265        if state.exhausted(clock) || state.remaining_scanned() < 2 {
266            run.report.stopped_on_budget = true;
267            break;
268        }
269        let limit = PAGE_SIZE.min(state.remaining_scanned() - 1);
270        let page = store.scan(p, &start, &end, cursor.as_ref(), limit).await?;
271        // SQL scans use LIMIT limit+1. Charge conservatively on other backends.
272        let rows = u32::try_from(page.entries.len()).unwrap_or(limit);
273        let _ = state.charge_scan(rows + 1);
274        for (key, value) in page.entries {
275            if state.work_exhausted(clock) {
276                run.report.stopped_on_budget = true;
277                break 'pages;
278            }
279            run.report.scanned += 1;
280            process_row(&ctx, registry, key, value, state, &mut run).await;
281        }
282        match page.next {
283            Some(next) => cursor = Some(next),
284            None => break,
285        }
286    }
287    // If the future range cannot be examined, keep a wake even when the
288    // final due page happened to end exactly at the commit/clock allowance.
289    if now_ms < u64::MAX && (state.work_exhausted(clock) || state.remaining_scanned() < 2) {
290        run.report.stopped_on_budget = true;
291    }
292    let future_due =
293        if !state.work_exhausted(clock) && state.remaining_scanned() >= 2 && now_ms < u64::MAX {
294            let page = store.scan(p, &end, &class_end, None, 1).await?;
295            let _ = state.charge_scan(2);
296            page.entries
297                .first()
298                .and_then(|(key, _)| match keys::parse(key) {
299                    Some(keys::ParsedKey::Timer { due_at_ms, .. }) => Some(due_at_ms),
300                    _ => None,
301                })
302        } else {
303            None
304        };
305    let pending = run.report.stopped_on_budget || run.retained_due;
306    let next = if pending && run.progress {
307        Some(now_ms)
308    } else if pending {
309        Some(now_ms.saturating_add(RETRY_BACKOFF_MS))
310    } else {
311        future_due
312    };
313    run.report.next_wake_ms =
314        min_due(min_due(next, future_due), run.committed_due).map(|due| due.max(now_ms));
315    Ok(run.report)
316}
317
318struct PartitionRun {
319    report: RunReport,
320    warned: [bool; 256],
321    retained_due: bool,
322    committed_due: Option<u64>,
323    progress: bool,
324}
325impl PartitionRun {
326    fn new() -> Self {
327        Self {
328            report: RunReport::default(),
329            warned: [false; 256],
330            retained_due: false,
331            committed_due: None,
332            progress: false,
333        }
334    }
335}
336
337async fn process_row<S: NamespaceStore>(
338    ctx: &TimerCtx<'_, S>,
339    registry: &TimerRegistry<'_, S>,
340    key: Key,
341    value: Value,
342    state: &mut TickState,
343    run: &mut PartitionRun,
344) {
345    let Some(keys::ParsedKey::Timer {
346        kind, reference, ..
347    }) = keys::parse(&key)
348    else {
349        // Corrupt encodings cannot safely be moved; retain them for repair.
350        run.report.unknown += 1;
351        run.retained_due = true;
352        if !run.warned[0] {
353            tracing::warn!("malformed timer row");
354            run.warned[0] = true;
355        }
356        return;
357    };
358    let Some((original_due, attempt)) = keys::timer_retry_state(&key) else {
359        run.retained_due = true;
360        return;
361    };
362    let timer = DueTimer {
363        due_at_ms: original_due,
364        kind: TimerKind::new(kind),
365        reference,
366        value,
367    };
368    let Some(handler) = registry.get(timer.kind) else {
369        run.report.unknown += 1;
370        if !run.warned[usize::from(kind)] {
371            tracing::warn!(kind, "unknown timer kind");
372            run.warned[usize::from(kind)] = true;
373        }
374        match backoff(ctx, &timer, key, attempt).await {
375            FireOutcome::Committed(due) => {
376                state.committed();
377                run.progress = true;
378                run.committed_due = min_due(run.committed_due, due);
379            }
380            FireOutcome::Raced => {
381                run.report.raced += 1;
382                run.retained_due = true;
383            }
384            FireOutcome::Failed => {
385                run.report.failed += 1;
386                run.retained_due = true;
387            }
388        }
389        return;
390    };
391    if !state.claim_attempt(timer.kind, handler.max_per_tick()) {
392        run.report.deferred += 1;
393        run.retained_due = true;
394        return;
395    }
396    match fire_timer(handler, ctx, &timer, key.clone()).await {
397        FireOutcome::Committed(put_due) => {
398            state.committed();
399            run.report.fired += 1;
400            run.progress = true;
401            run.committed_due = min_due(run.committed_due, put_due);
402        }
403        FireOutcome::Raced => {
404            run.report.raced += 1;
405            // A failed ancillary guard can leave the exact timer unchanged.
406            // Move that retained work too; a replaced timer fails this guard.
407            match backoff(ctx, &timer, key, attempt).await {
408                FireOutcome::Committed(due) => {
409                    state.committed();
410                    run.progress = true;
411                    run.committed_due = min_due(run.committed_due, due);
412                }
413                FireOutcome::Raced | FireOutcome::Failed => run.retained_due = true,
414            }
415        }
416        FireOutcome::Failed => {
417            run.report.failed += 1;
418            // Failed handler effects are discarded; only its timer moves.
419            match backoff(ctx, &timer, key, attempt).await {
420                FireOutcome::Committed(due) => {
421                    state.committed();
422                    run.progress = true;
423                    run.committed_due = min_due(run.committed_due, due);
424                }
425                FireOutcome::Raced => {
426                    run.report.raced += 1;
427                    run.retained_due = true;
428                }
429                FireOutcome::Failed => run.retained_due = true,
430            }
431        }
432    }
433}
434
435async fn backoff<S: NamespaceStore>(
436    ctx: &TimerCtx<'_, S>,
437    timer: &DueTimer,
438    key: Key,
439    attempt: u8,
440) -> FireOutcome {
441    let next_attempt = attempt.saturating_add(1).min(keys::MAX_TIMER_RETRY_ATTEMPT);
442    let delay = RETRY_BACKOFF_MS
443        .saturating_mul(1_u64 << (next_attempt - 1))
444        .min(MAX_RETRY_BACKOFF_MS);
445    let due = ctx.now_ms.saturating_add(delay);
446    let next_key = keys::timer_retry(
447        due,
448        timer.kind.get(),
449        &timer.reference,
450        timer.due_at_ms,
451        next_attempt,
452    );
453    if next_key == key {
454        return FireOutcome::Failed;
455    }
456    let batch = Batch::new()
457        .require(Precondition::Equals(key.clone(), timer.value.clone()))
458        .require(Precondition::Absent(next_key.clone()))
459        .delete(key)
460        .put(next_key, timer.value.clone());
461    match ctx.store.apply(ctx.partition, batch).await {
462        Ok(BatchOutcome::Committed) => FireOutcome::Committed(Some(due)),
463        Ok(BatchOutcome::PreconditionFailed { .. }) => FireOutcome::Raced,
464        Ok(BatchOutcome::DeadlinePassed { .. }) | Err(_) => FireOutcome::Failed,
465    }
466}