1mod 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
26pub const RETRY_BACKOFF_MS: u64 = 5_000;
28pub const MAX_RETRY_BACKOFF_MS: u64 = 600_000;
30const PAGE_SIZE: u32 = 64;
31
32#[non_exhaustive]
34#[derive(Debug, Clone)]
35pub struct DueTimer {
36 pub due_at_ms: u64,
38 pub kind: TimerKind,
40 pub reference: Bytes,
42 pub value: Value,
44}
45#[non_exhaustive]
47#[derive(Debug)]
48pub struct TimerCtx<'a, S> {
49 pub store: &'a S,
51 pub partition: &'a Partition,
53 pub now_ms: u64,
55}
56#[non_exhaustive]
58#[derive(Debug)]
59pub enum Fired {
60 Done(Batch),
62 Reschedule {
64 due_at_ms: u64,
66 value: Value,
68 batch: Batch,
70 },
71 Retry,
73}
74#[non_exhaustive]
76#[derive(Debug, Clone, Copy)]
77pub struct TickBudget {
78 pub max_fired: u32,
80 pub max_per_kind: u32,
82 pub max_scanned: u32,
84 pub max_elapsed_ms: u64,
86}
87impl TickBudget {
88 #[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#[non_exhaustive]
115#[derive(Debug, Clone, PartialEq, Eq, Default)]
116pub struct RunReport {
117 pub fired: u32,
119 pub raced: u32,
121 pub failed: u32,
123 pub unknown: u32,
125 pub deferred: u32,
127 pub scanned: u32,
129 pub stopped_on_budget: bool,
131 pub next_wake_ms: Option<u64>,
133}
134
135#[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
222pub 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
239pub 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 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 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 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 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 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}