Skip to main content

mkit_server/timers/
quota_rollup.rs

1//! Reconcile each ref shard's cumulative fixed-window quota into its
2//! coordinator. The coordinator apply precedes the local view write; a crash
3//! between them is safe because `qc` stores the last applied cumulative count.
4
5use super::{DueTimer, Fired, TimerCtx, TimerHandler, TimerKind, registry::kinds};
6use crate::quota::{NamespaceUsage, NamespaceView, QUOTA_ROLLUP_MS};
7use crate::rt::BoxFuture;
8use crate::store::{
9    Batch, BatchOutcome, NamespaceStore, Partition, Precondition, StoreError, Value, codec, keys,
10};
11use crate::telemetry::{
12    METRIC_NAMESPACE_QUOTA_REBASE, METRIC_NAMESPACE_QUOTA_ROLLUP_ERROR, Metrics, NoopMetrics,
13};
14
15const MAX_AGGREGATE_REPLANS: usize = 8;
16const CLOCK_GRACE_MS: u64 = mkit_core::write_auth::MAX_CLOCK_LEAD_MS.unsigned_abs();
17
18/// A source-side kind-5 handler with a client for the namespace coordinator.
19#[derive(Debug)]
20pub struct QuotaRollup<T, M = NoopMetrics> {
21    /// Store client capable of reaching the coordinator partition.
22    pub coordinator: T,
23    /// Metrics sink shared with the serving adapter.
24    pub metrics: M,
25}
26
27fn rollup_error<M: Metrics>(metrics: &M, error: &StoreError) {
28    let reason = match error {
29        StoreError::Corrupt(_) => "corrupt",
30        StoreError::Invalid(_) => "contention",
31        _ => "storage",
32    };
33    tracing::error!(%error, reason, "namespace quota rollup failed");
34    metrics.incr(
35        METRIC_NAMESPACE_QUOTA_ROLLUP_ERROR,
36        &[("reason", reason)],
37        1,
38    );
39}
40
41fn guard(key: crate::store::Key, value: Option<&Value>) -> Precondition {
42    match value {
43        Some(value) => Precondition::Equals(key, value.clone()),
44        None => Precondition::Absent(key),
45    }
46}
47
48/// Drain an older-window range in guarded batches of at most eight keys.
49/// The timer remains present until the caller commits its final batch, so a
50/// crash or a raced page can resume without retaining stale windows.
51async fn prune_older<S: NamespaceStore>(
52    store: &S,
53    partition: &Partition,
54    tag: &str,
55    window: u64,
56) -> Result<(), StoreError> {
57    if window == 0 {
58        return Ok(());
59    }
60    let (start, end) = keys::quota_namespace_before(tag, window);
61    let mut races = 0;
62    loop {
63        let page = store.scan(partition, &start, &end, None, 8).await?;
64        if page.entries.is_empty() {
65            return Ok(());
66        }
67        let mut batch = Batch::new();
68        for (key, value) in page.entries {
69            batch = batch
70                .require(Precondition::Equals(key.clone(), value))
71                .delete(key);
72        }
73        debug_assert!(
74            batch.preconditions.len() + batch.writes.len() <= crate::store::MAX_BATCH_OPS
75        );
76        match store.apply(partition, batch).await? {
77            BatchOutcome::Committed => races = 0,
78            BatchOutcome::PreconditionFailed { .. } if races < 8 => races += 1,
79            BatchOutcome::PreconditionFailed { .. } => {
80                return Err(StoreError::Invalid("namespace prune contention".into()));
81            }
82            BatchOutcome::DeadlinePassed { .. } => {
83                return Err(StoreError::Invalid(
84                    "namespace prune had no deadline".into(),
85                ));
86            }
87        }
88    }
89}
90
91async fn aggregate<T: NamespaceStore, M: Metrics>(
92    target: &T,
93    metrics: &M,
94    coordinator: &Partition,
95    source: &Partition,
96    window: u64,
97    local: NamespaceUsage,
98) -> Result<NamespaceUsage, StoreError> {
99    let source_key = keys::quota_contribution(window, source)?;
100    let total_key = keys::quota_total(window);
101    for _ in 0..MAX_AGGREGATE_REPLANS {
102        let values = target
103            .get_many(coordinator, &[source_key.clone(), total_key.clone()])
104            .await?;
105        let old_value = values.first().and_then(Option::as_ref);
106        let total_value = values.get(1).and_then(Option::as_ref);
107        let old = old_value
108            .map(codec::decode_namespace_usage)
109            .transpose()?
110            .unwrap_or_default();
111        let total = total_value
112            .map(codec::decode_namespace_usage)
113            .transpose()?
114            .unwrap_or_default();
115        if local == old {
116            return Ok(total);
117        }
118        let rebased = local.delta_from(old).is_none();
119        let removed = NamespaceUsage {
120            ops: old.ops.saturating_sub(local.ops),
121            bytes: old.bytes.saturating_sub(local.bytes),
122        };
123        let added = NamespaceUsage {
124            ops: local.ops.saturating_sub(old.ops),
125            bytes: local.bytes.saturating_sub(old.bytes),
126        };
127        let next = NamespaceUsage {
128            ops: total.ops.saturating_sub(removed.ops),
129            bytes: total.bytes.saturating_sub(removed.bytes),
130        }
131        .checked_add(added)
132        .ok_or_else(|| StoreError::Corrupt("namespace total overflow".into()))?;
133        // An inconsistent coordinator (a total below this source's own old
134        // contribution) must never leave the total below the local usage, or
135        // this shard's view would deny every write until the window ends.
136        let next = NamespaceUsage {
137            ops: next.ops.max(local.ops),
138            bytes: next.bytes.max(local.bytes),
139        };
140        let batch = Batch::new()
141            .require(guard(source_key.clone(), old_value))
142            .require(guard(total_key.clone(), total_value))
143            .put(source_key.clone(), codec::encode_namespace_usage(local))
144            .put(total_key.clone(), codec::encode_namespace_usage(next));
145        match target.apply(coordinator, batch).await? {
146            BatchOutcome::Committed => {
147                if rebased {
148                    tracing::warn!(
149                        window,
150                        "namespace quota contribution decreased; re-baselined"
151                    );
152                    metrics.incr(METRIC_NAMESPACE_QUOTA_REBASE, &[], 1);
153                }
154                return Ok(next);
155            }
156            BatchOutcome::PreconditionFailed { .. } => {}
157            BatchOutcome::DeadlinePassed { .. } => {
158                return Err(StoreError::Invalid(
159                    "namespace aggregate had no deadline".into(),
160                ));
161            }
162        }
163    }
164    Err(StoreError::Invalid("namespace aggregate contention".into()))
165}
166
167/// Delete this source's old contribution. The last source also deletes the
168/// obsolete aggregate. A crash before local cleanup safely re-adds and
169/// removes the cumulative contribution on the next fire.
170async fn prune_coordinator<T: NamespaceStore>(
171    target: &T,
172    coordinator: &Partition,
173    source: &Partition,
174    window: u64,
175) -> Result<(), StoreError> {
176    let source_key = keys::quota_contribution(window, source)?;
177    let total_key = keys::quota_total(window);
178    let source_value = target.get(coordinator, &source_key).await?;
179    let (start, end) = keys::quota_namespace_window(keys::TAG_QUOTA_CONTRIBUTION, window);
180    let page = target.scan(coordinator, &start, &end, None, 2).await?;
181    let only_source = page.next.is_none() && page.entries.iter().all(|(key, _)| *key == source_key);
182    let total_value = if only_source {
183        target.get(coordinator, &total_key).await?
184    } else {
185        None
186    };
187    let mut batch = Batch::new()
188        .require(guard(source_key.clone(), source_value.as_ref()))
189        .delete(source_key);
190    if only_source {
191        batch = batch
192            .require(guard(total_key.clone(), total_value.as_ref()))
193            .delete(total_key);
194    }
195    prune_older(target, coordinator, keys::TAG_QUOTA_CONTRIBUTION, window).await?;
196    prune_older(target, coordinator, keys::TAG_QUOTA_TOTAL, window).await?;
197    match target.apply(coordinator, batch).await? {
198        BatchOutcome::Committed => Ok(()),
199        BatchOutcome::PreconditionFailed { .. } => {
200            Err(StoreError::Invalid("namespace prune contention".into()))
201        }
202        BatchOutcome::DeadlinePassed { .. } => Err(StoreError::Invalid(
203            "namespace prune had no deadline".into(),
204        )),
205    }
206}
207
208async fn fire_exact<S: NamespaceStore>(
209    ctx: &TimerCtx<'_, S>,
210    timer: &DueTimer,
211    window: u64,
212    end: u64,
213    expired: bool,
214) -> Result<Fired, StoreError> {
215    let key = keys::quota_total(window);
216    let value = ctx.store.get(ctx.partition, &key).await?;
217    if expired {
218        let batch = Batch::new()
219            .require(guard(key.clone(), value.as_ref()))
220            .delete(key);
221        prune_older(ctx.store, ctx.partition, keys::TAG_QUOTA_TOTAL, window).await?;
222        prune_older(
223            ctx.store,
224            ctx.partition,
225            keys::TAG_QUOTA_CONTRIBUTION,
226            window,
227        )
228        .await?;
229        Ok(Fired::Done(batch))
230    } else {
231        Ok(Fired::Reschedule {
232            due_at_ms: end
233                .saturating_add(CLOCK_GRACE_MS)
234                .max(timer.due_at_ms.saturating_add(1)),
235            value: timer.value.clone(),
236            batch: Batch::new(),
237        })
238    }
239}
240
241impl<S: NamespaceStore, T: NamespaceStore, M: Metrics> TimerHandler<S> for QuotaRollup<T, M> {
242    fn kind(&self) -> TimerKind {
243        kinds::QUOTA_ROLLUP
244    }
245
246    #[allow(clippy::too_many_lines)] // One fire must finish coordinator cleanup before its local timer batch.
247    fn fire<'a>(
248        &'a self,
249        ctx: &'a TimerCtx<'a, S>,
250        timer: &'a DueTimer,
251    ) -> BoxFuture<'a, Result<Fired, StoreError>> {
252        Box::pin(async move {
253            let Ok(reference) = <[u8; 8]>::try_from(timer.reference.as_ref()) else {
254                tracing::warn!("discarding malformed quota timer reference");
255                return Ok(Fired::Done(Batch::new()));
256            };
257            let window = u64::from_be_bytes(reference);
258            let window_ms = codec::decode_u64(&timer.value)?;
259            if window_ms == 0 {
260                return Err(StoreError::Corrupt("zero quota window".into()));
261            }
262            let end = window.saturating_add(1).saturating_mul(window_ms);
263            let expired = ctx.now_ms >= end.saturating_add(CLOCK_GRACE_MS);
264            match ctx.partition {
265                Partition::Ref { ns, .. } => {
266                    let key = keys::quota_shard(window);
267                    let coordinator = Partition::Coordinator(ns.clone());
268                    let backoff = || Fired::Reschedule {
269                        due_at_ms: ctx
270                            .now_ms
271                            .saturating_add(QUOTA_ROLLUP_MS)
272                            .max(timer.due_at_ms.saturating_add(1)),
273                        value: timer.value.clone(),
274                        batch: Batch::new(),
275                    };
276                    let Some(value) = ctx.store.get(ctx.partition, &key).await? else {
277                        if expired {
278                            if let Err(error) = prune_coordinator(
279                                &self.coordinator,
280                                &coordinator,
281                                ctx.partition,
282                                window,
283                            )
284                            .await
285                            {
286                                rollup_error(&self.metrics, &error);
287                                return Ok(backoff());
288                            }
289                            prune_older(ctx.store, ctx.partition, keys::TAG_QUOTA_SHARD, window)
290                                .await?;
291                            prune_older(ctx.store, ctx.partition, keys::TAG_QUOTA_VIEW, window)
292                                .await?;
293                        }
294                        return Ok(Fired::Done(Batch::new().delete(keys::quota_view(window))));
295                    };
296                    let local = codec::decode_namespace_usage(&value)?;
297                    let aggregated = aggregate(
298                        &self.coordinator,
299                        &self.metrics,
300                        &coordinator,
301                        ctx.partition,
302                        window,
303                        local,
304                    )
305                    .await;
306                    if expired {
307                        if let Err(error) = aggregated {
308                            rollup_error(&self.metrics, &error);
309                        }
310                        if let Err(error) = prune_coordinator(
311                            &self.coordinator,
312                            &coordinator,
313                            ctx.partition,
314                            window,
315                        )
316                        .await
317                        {
318                            rollup_error(&self.metrics, &error);
319                            return Ok(backoff());
320                        }
321                        prune_older(ctx.store, ctx.partition, keys::TAG_QUOTA_SHARD, window)
322                            .await?;
323                        prune_older(ctx.store, ctx.partition, keys::TAG_QUOTA_VIEW, window).await?;
324                        let batch = Batch::new()
325                            .require(Precondition::Equals(key.clone(), value))
326                            .delete(key)
327                            .delete(keys::quota_view(window));
328                        Ok(Fired::Done(batch))
329                    } else {
330                        let total = match aggregated {
331                            Ok(total) => total,
332                            Err(error) => {
333                                rollup_error(&self.metrics, &error);
334                                return Ok(backoff());
335                            }
336                        };
337                        let view = NamespaceView {
338                            total,
339                            pushed: local,
340                            observed_at_ms: ctx.now_ms,
341                        };
342                        Ok(Fired::Reschedule {
343                            due_at_ms: ctx
344                                .now_ms
345                                .saturating_add(QUOTA_ROLLUP_MS)
346                                .max(timer.due_at_ms.saturating_add(1)),
347                            value: timer.value.clone(),
348                            batch: Batch::new()
349                                .put(keys::quota_view(window), codec::encode_namespace_view(view)),
350                        })
351                    }
352                }
353                Partition::Coordinator(_) | Partition::Namespace(_) => {
354                    fire_exact(ctx, timer, window, end, expired).await
355                }
356                _ => {
357                    tracing::warn!(partition = ?ctx.partition, "discarding quota timer on unexpected partition");
358                    Ok(Fired::Done(Batch::new()))
359                }
360            }
361        })
362    }
363}
364
365#[cfg(all(test, feature = "memory"))]
366mod tests {
367    use super::*;
368    use crate::MemoryKv;
369    use crate::repo::{NamespaceKey, RepoName};
370    use crate::rt::ManualClock;
371    use crate::timers::{TickBudget, TimerRegistry, run_due};
372    use bytes::Bytes;
373    use std::sync::atomic::{AtomicBool, Ordering};
374    use std::sync::{Arc, Mutex};
375
376    const WINDOW_MS: u64 = 600_000;
377
378    #[derive(Debug, Default, Clone)]
379    struct CountMetrics(Arc<Mutex<Vec<(&'static str, String)>>>);
380
381    impl Metrics for CountMetrics {
382        fn incr(&self, name: &'static str, labels: &[(&'static str, &str)], _by: u64) {
383            self.0.lock().expect("metrics lock").push((
384                name,
385                labels
386                    .iter()
387                    .find(|(key, _)| *key == "reason")
388                    .map_or(String::new(), |(_, value)| (*value).to_owned()),
389            ));
390        }
391
392        fn observe_ms(&self, _: &'static str, _: &[(&'static str, &str)], _: f64) {}
393    }
394
395    #[derive(Debug, Clone)]
396    struct SharedStore {
397        inner: Arc<MemoryKv>,
398        batches: Arc<Mutex<Vec<(Partition, usize)>>>,
399        contend: Arc<AtomicBool>,
400    }
401
402    impl SharedStore {
403        fn new(clock: Arc<ManualClock>) -> Self {
404            Self {
405                inner: Arc::new(MemoryKv::with_clock(clock)),
406                batches: Arc::new(Mutex::new(Vec::new())),
407                contend: Arc::new(AtomicBool::new(false)),
408            }
409        }
410    }
411
412    impl NamespaceStore for SharedStore {
413        fn capabilities(&self) -> crate::store::StoreCapabilities {
414            self.inner.capabilities()
415        }
416
417        async fn get(
418            &self,
419            p: &Partition,
420            key: &crate::store::Key,
421        ) -> Result<Option<Value>, StoreError> {
422            self.inner.get(p, key).await
423        }
424
425        async fn get_many(
426            &self,
427            p: &Partition,
428            keys: &[crate::store::Key],
429        ) -> Result<Vec<Option<Value>>, StoreError> {
430            self.inner.get_many(p, keys).await
431        }
432
433        async fn scan(
434            &self,
435            p: &Partition,
436            start: &crate::store::Key,
437            end: &crate::store::Key,
438            after: Option<&crate::store::Cursor>,
439            limit: u32,
440        ) -> Result<crate::store::ScanPage, StoreError> {
441            self.inner.scan(p, start, end, after, limit).await
442        }
443
444        async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
445            self.batches
446                .lock()
447                .expect("batch log lock")
448                .push((p.clone(), batch.preconditions.len() + batch.writes.len()));
449            if self.contend.load(Ordering::SeqCst) {
450                return Ok(BatchOutcome::PreconditionFailed {
451                    index: 0,
452                    observed: None,
453                });
454            }
455            self.inner.apply(p, batch).await
456        }
457
458        async fn stats(&self, p: &Partition) -> Result<crate::store::PartitionStats, StoreError> {
459            self.inner.stats(p).await
460        }
461
462        async fn probe(&self) -> Result<(), StoreError> {
463            self.inner.probe().await
464        }
465    }
466
467    fn source(i: usize) -> Partition {
468        Partition::Ref {
469            ns: NamespaceKey::deployment_default(),
470            repo: RepoName::new("room").expect("valid test repository"),
471            shard_ref: format!("refs/heads/b{i}"),
472        }
473    }
474
475    fn coordinator() -> Partition {
476        Partition::Coordinator(NamespaceKey::deployment_default())
477    }
478
479    fn timer(due: u64, window: u64) -> (crate::store::Key, DueTimer) {
480        let reference = Bytes::copy_from_slice(&window.to_be_bytes());
481        let value = codec::encode_u64(WINDOW_MS);
482        (
483            keys::timer(due, kinds::QUOTA_ROLLUP.get(), &reference),
484            DueTimer {
485                due_at_ms: due,
486                kind: kinds::QUOTA_ROLLUP,
487                reference,
488                value,
489            },
490        )
491    }
492
493    async fn usage<S: NamespaceStore>(
494        store: &S,
495        p: &Partition,
496        key: crate::store::Key,
497    ) -> NamespaceUsage {
498        store
499            .get(p, &key)
500            .await
501            .expect("read usage")
502            .as_ref()
503            .map(codec::decode_namespace_usage)
504            .transpose()
505            .expect("decode usage")
506            .unwrap_or_default()
507    }
508
509    #[tokio::test]
510    async fn converges_across_shards_and_refire_after_coordinator_crash_is_idempotent() {
511        let clock = Arc::new(ManualClock::new(60_000));
512        let local = MemoryKv::with_clock(clock.clone());
513        let handler = QuotaRollup {
514            coordinator: MemoryKv::with_clock(clock.clone()),
515            metrics: NoopMetrics,
516        };
517        let (_, fired) = timer(60_000, 0);
518        for i in 0..4 {
519            let (key, _) = timer(60_000, 0);
520            local
521                .apply(
522                    &source(i),
523                    Batch::new()
524                        .put(
525                            keys::quota_shard(0),
526                            codec::encode_namespace_usage(NamespaceUsage {
527                                ops: (i + 1) as u64,
528                                bytes: i as u64,
529                            }),
530                        )
531                        .put(key, codec::encode_u64(WINDOW_MS)),
532                )
533                .await
534                .unwrap();
535        }
536        // Simulate a crash after target apply but before the returned local
537        // view/timer batch. The same fire sees qc and adds zero delta.
538        let first_source = source(0);
539        let ctx = TimerCtx {
540            store: &local,
541            partition: &first_source,
542            now_ms: 60_000,
543        };
544        let first = handler.fire(&ctx, &fired).await.unwrap();
545        let Fired::Reschedule { batch, .. } = first else {
546            panic!("must reschedule")
547        };
548        assert!(
549            batch.preconditions.is_empty(),
550            "view refresh does not guard qs"
551        );
552        assert_eq!(
553            usage(&handler.coordinator, &coordinator(), keys::quota_total(0))
554                .await
555                .ops,
556            1
557        );
558        let again = handler.fire(&ctx, &fired).await.unwrap();
559        assert!(matches!(again, Fired::Reschedule { .. }));
560        assert_eq!(
561            usage(&handler.coordinator, &coordinator(), keys::quota_total(0))
562                .await
563                .ops,
564            1
565        );
566
567        for i in 0..4 {
568            let partition = source(i);
569            let ctx = TimerCtx {
570                store: &local,
571                partition: &partition,
572                now_ms: 60_000,
573            };
574            assert!(matches!(
575                handler.fire(&ctx, &fired).await.unwrap(),
576                Fired::Reschedule { .. }
577            ));
578        }
579        assert_eq!(
580            usage(&handler.coordinator, &coordinator(), keys::quota_total(0)).await,
581            NamespaceUsage { ops: 10, bytes: 6 }
582        );
583        for i in 0..4 {
584            let partition = source(i);
585            let ctx = TimerCtx {
586                store: &local,
587                partition: &partition,
588                now_ms: 60_000,
589            };
590            let Fired::Reschedule { batch, .. } = handler.fire(&ctx, &fired).await.unwrap() else {
591                panic!("must reschedule")
592            };
593            let view_value = batch
594                .writes
595                .iter()
596                .find_map(|write| match write {
597                    crate::store::Write::Put(key, value) if *key == keys::quota_view(0) => {
598                        Some(value)
599                    }
600                    _ => None,
601                })
602                .unwrap();
603            let view = codec::decode_namespace_view(view_value).unwrap();
604            assert_eq!(view.total.ops, 10);
605            assert_eq!(view.pushed.ops, (i + 1) as u64);
606        }
607    }
608
609    #[tokio::test]
610    async fn ended_window_prunes_both_partitions_and_leaves_no_timer() {
611        let clock = Arc::new(ManualClock::new(60_000));
612        let local = MemoryKv::with_clock(clock.clone());
613        let handler = QuotaRollup {
614            coordinator: MemoryKv::with_clock(clock.clone()),
615            metrics: NoopMetrics,
616        };
617        let shard = source(0);
618        let (timer_key, _) = timer(60_000, 0);
619        local
620            .apply(
621                &shard,
622                Batch::new()
623                    .put(
624                        keys::quota_shard(0),
625                        codec::encode_namespace_usage(NamespaceUsage { ops: 3, bytes: 12 }),
626                    )
627                    .put(timer_key, codec::encode_u64(WINDOW_MS)),
628            )
629            .await
630            .unwrap();
631        let registry = TimerRegistry::new().register(handler);
632        let first = run_due(
633            &local,
634            &shard,
635            &registry,
636            clock.as_ref(),
637            60_000,
638            &TickBudget::default(),
639        )
640        .await
641        .unwrap();
642        assert_eq!(first.fired, 1);
643        assert!(
644            local
645                .get(&shard, &keys::quota_view(0))
646                .await
647                .unwrap()
648                .is_some()
649        );
650        let before = local.stats(&shard).await.unwrap().keys.unwrap();
651        clock.set(700_000);
652        let last = run_due(
653            &local,
654            &shard,
655            &registry,
656            clock.as_ref(),
657            700_000,
658            &TickBudget::default(),
659        )
660        .await
661        .unwrap();
662        assert_eq!(last.fired, 1);
663        assert_eq!(last.next_wake_ms, None);
664        assert_eq!(local.stats(&shard).await.unwrap().keys, Some(0));
665        assert!(before >= 3);
666    }
667
668    #[tokio::test]
669    async fn final_fire_removes_old_coordinator_rows() {
670        let clock = Arc::new(ManualClock::new(700_000));
671        let local = MemoryKv::with_clock(clock.clone());
672        let handler = QuotaRollup {
673            coordinator: MemoryKv::with_clock(clock),
674            metrics: NoopMetrics,
675        };
676        let shard = source(0);
677        let (_, fired) = timer(60_000, 0);
678        local
679            .apply(
680                &shard,
681                Batch::new().put(
682                    keys::quota_shard(0),
683                    codec::encode_namespace_usage(NamespaceUsage { ops: 2, bytes: 0 }),
684                ),
685            )
686            .await
687            .unwrap();
688        let ctx = TimerCtx {
689            store: &local,
690            partition: &shard,
691            now_ms: 700_000,
692        };
693        let Fired::Done(_) = handler.fire(&ctx, &fired).await.unwrap() else {
694            panic!("ended window must stop")
695        };
696        assert_eq!(
697            handler
698                .coordinator
699                .stats(&coordinator())
700                .await
701                .unwrap()
702                .keys,
703            Some(0)
704        );
705    }
706
707    #[tokio::test]
708    async fn coordinator_direct_charge_remains_in_the_aggregate() {
709        let clock = Arc::new(ManualClock::new(60_000));
710        let local = MemoryKv::with_clock(clock.clone());
711        let handler = QuotaRollup {
712            coordinator: MemoryKv::with_clock(clock),
713            metrics: NoopMetrics,
714        };
715        let shard = source(0);
716        local
717            .apply(
718                &shard,
719                Batch::new().put(
720                    keys::quota_shard(0),
721                    codec::encode_namespace_usage(NamespaceUsage { ops: 2, bytes: 0 }),
722                ),
723            )
724            .await
725            .unwrap();
726        handler
727            .coordinator
728            .apply(
729                &coordinator(),
730                Batch::new().put(
731                    keys::quota_total(0),
732                    codec::encode_namespace_usage(NamespaceUsage { ops: 1, bytes: 9 }),
733                ),
734            )
735            .await
736            .unwrap();
737        let (_, fired) = timer(60_000, 0);
738        let ctx = TimerCtx {
739            store: &local,
740            partition: &shard,
741            now_ms: 60_000,
742        };
743        assert!(matches!(
744            handler.fire(&ctx, &fired).await.unwrap(),
745            Fired::Reschedule { .. }
746        ));
747        assert_eq!(
748            usage(&handler.coordinator, &coordinator(), keys::quota_total(0)).await,
749            NamespaceUsage { ops: 3, bytes: 9 }
750        );
751    }
752
753    #[tokio::test]
754    async fn decreased_contribution_rebaselines_then_prunes_after_idle() {
755        let window = 2;
756        let first_due = window * WINDOW_MS + QUOTA_ROLLUP_MS;
757        let clock = Arc::new(ManualClock::new(first_due.cast_signed()));
758        let local = SharedStore::new(clock.clone());
759        let coordinator_store = SharedStore::new(clock.clone());
760        let metrics = CountMetrics::default();
761        let shard = source(0);
762        let old = NamespaceUsage { ops: 5, bytes: 10 };
763        let restored = NamespaceUsage { ops: 2, bytes: 12 };
764        let (timer_key, _) = timer(first_due, window);
765        local
766            .apply(
767                &shard,
768                Batch::new()
769                    .put(
770                        keys::quota_shard(window),
771                        codec::encode_namespace_usage(restored),
772                    )
773                    .put(timer_key, codec::encode_u64(WINDOW_MS)),
774            )
775            .await
776            .unwrap();
777        coordinator_store
778            .apply(
779                &coordinator(),
780                Batch::new()
781                    .put(
782                        keys::quota_contribution(window, &shard).unwrap(),
783                        codec::encode_namespace_usage(old),
784                    )
785                    .put(
786                        keys::quota_total(window),
787                        codec::encode_namespace_usage(old),
788                    ),
789            )
790            .await
791            .unwrap();
792        let registry = TimerRegistry::new().register(QuotaRollup {
793            coordinator: coordinator_store.clone(),
794            metrics: metrics.clone(),
795        });
796        let first = run_due(
797            &local,
798            &shard,
799            &registry,
800            clock.as_ref(),
801            first_due,
802            &TickBudget::default(),
803        )
804        .await
805        .unwrap();
806        assert_eq!(first.fired, 1);
807        assert_eq!(
808            usage(
809                &coordinator_store,
810                &coordinator(),
811                keys::quota_total(window)
812            )
813            .await,
814            restored
815        );
816        assert_eq!(
817            usage(
818                &coordinator_store,
819                &coordinator(),
820                keys::quota_contribution(window, &shard).unwrap()
821            )
822            .await,
823            restored
824        );
825        assert_eq!(
826            metrics.0.lock().unwrap().as_slice(),
827            &[(METRIC_NAMESPACE_QUOTA_REBASE, String::new())]
828        );
829        let expired = (window + 1) * WINDOW_MS + CLOCK_GRACE_MS + 1;
830        clock.set(expired.cast_signed());
831        let last = run_due(
832            &local,
833            &shard,
834            &registry,
835            clock.as_ref(),
836            expired,
837            &TickBudget::default(),
838        )
839        .await
840        .unwrap();
841        assert_eq!(last.fired, 1);
842        assert_eq!(last.next_wake_ms, None);
843        assert_eq!(local.stats(&shard).await.unwrap().keys, Some(0));
844        assert_eq!(
845            coordinator_store.stats(&coordinator()).await.unwrap().keys,
846            Some(0)
847        );
848    }
849
850    #[tokio::test]
851    async fn expired_window_drains_more_than_eight_older_rows_per_class() {
852        let window = 12;
853        let due = (window + 1) * WINDOW_MS + CLOCK_GRACE_MS + 1;
854        let clock = Arc::new(ManualClock::new(due.cast_signed()));
855        let local = SharedStore::new(clock.clone());
856        let coordinator_store = SharedStore::new(clock.clone());
857        let shard = source(0);
858        let one = codec::encode_namespace_usage(NamespaceUsage { ops: 1, bytes: 0 });
859        let view = codec::encode_namespace_view(NamespaceView {
860            total: NamespaceUsage { ops: 1, bytes: 0 },
861            pushed: NamespaceUsage::default(),
862            observed_at_ms: due,
863        });
864        let mut local_rows = Batch::new().put(keys::quota_shard(window), one.clone());
865        let mut coordinator_rows = Batch::new();
866        for old in 0..window {
867            local_rows = local_rows
868                .put(keys::quota_shard(old), one.clone())
869                .put(keys::quota_view(old), view.clone());
870            coordinator_rows = coordinator_rows
871                .put(keys::quota_contribution(old, &shard).unwrap(), one.clone())
872                .put(keys::quota_total(old), one.clone());
873        }
874        // Rows of the next window must survive the prune.
875        local_rows = local_rows
876            .put(keys::quota_shard(window + 1), one.clone())
877            .put(keys::quota_view(window + 1), view.clone());
878        coordinator_rows = coordinator_rows
879            .put(
880                keys::quota_contribution(window + 1, &shard).unwrap(),
881                one.clone(),
882            )
883            .put(keys::quota_total(window + 1), one.clone());
884        let (timer_key, _) = timer(due, window);
885        local_rows = local_rows.put(timer_key, codec::encode_u64(WINDOW_MS));
886        assert_eq!(
887            local.apply(&shard, local_rows).await.unwrap(),
888            BatchOutcome::Committed
889        );
890        assert_eq!(
891            coordinator_store
892                .apply(&coordinator(), coordinator_rows)
893                .await
894                .unwrap(),
895            BatchOutcome::Committed
896        );
897        let registry = TimerRegistry::new().register(QuotaRollup {
898            coordinator: coordinator_store.clone(),
899            metrics: NoopMetrics,
900        });
901        let report = run_due(
902            &local,
903            &shard,
904            &registry,
905            clock.as_ref(),
906            due,
907            &TickBudget::default(),
908        )
909        .await
910        .unwrap();
911        let coord = coordinator();
912        for (store, partition) in [(&local, &shard), (&coordinator_store, &coord)] {
913            assert!(
914                store
915                    .batches
916                    .lock()
917                    .unwrap()
918                    .iter()
919                    .all(|(p, ops)| { p == partition && *ops <= crate::store::MAX_BATCH_OPS })
920            );
921        }
922        assert_eq!(report.fired, 1);
923        assert_eq!(report.failed, 0);
924        assert_eq!(report.next_wake_ms, None);
925        assert_eq!(local.stats(&shard).await.unwrap().keys, Some(2));
926        assert_eq!(
927            coordinator_store.stats(&coordinator()).await.unwrap().keys,
928            Some(2)
929        );
930        let kept = local
931            .get_many(
932                &shard,
933                &[keys::quota_shard(window + 1), keys::quota_view(window + 1)],
934            )
935            .await
936            .unwrap();
937        assert!(kept.iter().all(Option::is_some));
938        let kept = coordinator_store
939            .get_many(
940                &coordinator(),
941                &[
942                    keys::quota_contribution(window + 1, &shard).unwrap(),
943                    keys::quota_total(window + 1),
944                ],
945            )
946            .await
947            .unwrap();
948        assert!(kept.iter().all(Option::is_some));
949    }
950
951    #[tokio::test]
952    async fn corrupt_coordinator_row_backs_off_for_a_full_period() {
953        let clock = Arc::new(ManualClock::new(QUOTA_ROLLUP_MS.cast_signed()));
954        let local = SharedStore::new(clock.clone());
955        let coordinator_store = SharedStore::new(clock.clone());
956        let metrics = CountMetrics::default();
957        let shard = source(0);
958        let (timer_key, _) = timer(QUOTA_ROLLUP_MS, 0);
959        local
960            .apply(
961                &shard,
962                Batch::new()
963                    .put(
964                        keys::quota_shard(0),
965                        codec::encode_namespace_usage(NamespaceUsage { ops: 1, bytes: 0 }),
966                    )
967                    .put(timer_key, codec::encode_u64(WINDOW_MS)),
968            )
969            .await
970            .unwrap();
971        coordinator_store
972            .apply(
973                &coordinator(),
974                Batch::new().put(keys::quota_total(0), Value::new(vec![1])),
975            )
976            .await
977            .unwrap();
978        let registry = TimerRegistry::new().register(QuotaRollup {
979            coordinator: coordinator_store.clone(),
980            metrics: metrics.clone(),
981        });
982        let report = run_due(
983            &local,
984            &shard,
985            &registry,
986            clock.as_ref(),
987            QUOTA_ROLLUP_MS,
988            &TickBudget::default(),
989        )
990        .await
991        .unwrap();
992        assert_eq!(report.fired, 1);
993        assert_eq!(report.next_wake_ms, Some(2 * QUOTA_ROLLUP_MS));
994        assert_eq!(
995            metrics.0.lock().unwrap().as_slice(),
996            &[(METRIC_NAMESPACE_QUOTA_ROLLUP_ERROR, "corrupt".to_owned())]
997        );
998        let expired = WINDOW_MS + CLOCK_GRACE_MS + 1;
999        clock.set(expired.cast_signed());
1000        let final_fire = run_due(
1001            &local,
1002            &shard,
1003            &registry,
1004            clock.as_ref(),
1005            expired,
1006            &TickBudget::default(),
1007        )
1008        .await
1009        .unwrap();
1010        assert_eq!(final_fire.fired, 1);
1011        assert_eq!(final_fire.next_wake_ms, None);
1012        assert_eq!(local.stats(&shard).await.unwrap().keys, Some(0));
1013        assert_eq!(
1014            coordinator_store.stats(&coordinator()).await.unwrap().keys,
1015            Some(0)
1016        );
1017    }
1018
1019    #[tokio::test]
1020    async fn exhausted_aggregate_replans_log_and_back_off() {
1021        let clock = Arc::new(ManualClock::new(QUOTA_ROLLUP_MS.cast_signed()));
1022        let local = MemoryKv::with_clock(clock.clone());
1023        let coordinator_store = SharedStore::new(clock);
1024        coordinator_store.contend.store(true, Ordering::SeqCst);
1025        let metrics = CountMetrics::default();
1026        let shard = source(0);
1027        local
1028            .apply(
1029                &shard,
1030                Batch::new().put(
1031                    keys::quota_shard(0),
1032                    codec::encode_namespace_usage(NamespaceUsage { ops: 1, bytes: 0 }),
1033                ),
1034            )
1035            .await
1036            .unwrap();
1037        let handler = QuotaRollup {
1038            coordinator: coordinator_store.clone(),
1039            metrics: metrics.clone(),
1040        };
1041        let (_, fired) = timer(QUOTA_ROLLUP_MS, 0);
1042        let ctx = TimerCtx {
1043            store: &local,
1044            partition: &shard,
1045            now_ms: QUOTA_ROLLUP_MS,
1046        };
1047        let Fired::Reschedule { due_at_ms, .. } = handler.fire(&ctx, &fired).await.unwrap() else {
1048            panic!("aggregate contention must back off")
1049        };
1050        assert_eq!(due_at_ms, 2 * QUOTA_ROLLUP_MS);
1051        assert_eq!(
1052            coordinator_store.batches.lock().unwrap().len(),
1053            MAX_AGGREGATE_REPLANS
1054        );
1055        assert_eq!(
1056            metrics.0.lock().unwrap().as_slice(),
1057            &[(METRIC_NAMESPACE_QUOTA_ROLLUP_ERROR, "contention".to_owned())]
1058        );
1059    }
1060
1061    #[tokio::test]
1062    async fn exact_partition_timer_stops_after_window_end() {
1063        let clock = Arc::new(ManualClock::new(700_000));
1064        let store = MemoryKv::with_clock(clock.clone());
1065        let registry = TimerRegistry::new().register(QuotaRollup {
1066            coordinator: MemoryKv::with_clock(clock.clone()),
1067            metrics: NoopMetrics,
1068        });
1069        for partition in [
1070            Partition::Namespace(NamespaceKey::deployment_default()),
1071            coordinator(),
1072        ] {
1073            let (key, _) = timer(630_000, 0);
1074            store
1075                .apply(
1076                    &partition,
1077                    Batch::new()
1078                        .put(
1079                            keys::quota_total(0),
1080                            codec::encode_namespace_usage(NamespaceUsage { ops: 1, bytes: 0 }),
1081                        )
1082                        .put(key, codec::encode_u64(WINDOW_MS)),
1083                )
1084                .await
1085                .unwrap();
1086            let report = run_due(
1087                &store,
1088                &partition,
1089                &registry,
1090                clock.as_ref(),
1091                700_000,
1092                &TickBudget::default(),
1093            )
1094            .await
1095            .unwrap();
1096            assert_eq!(report.fired, 1);
1097            assert_eq!(report.next_wake_ms, None);
1098            assert_eq!(store.stats(&partition).await.unwrap().keys, Some(0));
1099        }
1100    }
1101}