1use 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#[derive(Debug)]
20pub struct QuotaRollup<T, M = NoopMetrics> {
21 pub coordinator: T,
23 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
48async 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 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
167async 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)] 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 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 ®istry,
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 ®istry,
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 ®istry,
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 ®istry,
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 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 ®istry,
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 ®istry,
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 ®istry,
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 ®istry,
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}