Skip to main content

mkit_server/
quota.rs

1//! Write quotas: the value types, the default limits and the fixed-window
2//! evaluation.
3//!
4//! [`evaluate_quota`] is the canonical copy of vcs-worker's former
5//! `write_quota.rs` (removed when vcs-worker moved onto this crate in
6//! WP-M0-17). `apps/repo-worker` keeps its own copy (planner decision Q11).
7
8use mkit_core::hash::to_hex;
9
10use crate::error::{AbortCause, ServerError};
11use crate::repo::NamespaceKey;
12use crate::store::{Key, Precondition, Value, Write, codec, keys};
13use crate::timers::registry::kinds;
14
15/// Limits for one fixed quota window.
16#[derive(Debug, Clone, Copy, PartialEq, Eq)]
17pub struct QuotaLimits {
18    /// Window length in milliseconds.
19    pub window_ms: i64,
20    /// Most write operations allowed per window.
21    pub max_ops: u32,
22    /// Most bytes allowed per window.
23    pub max_bytes: u64,
24}
25
26/// Usage recorded in the current window.
27#[derive(Debug, Clone, Copy, PartialEq, Eq)]
28pub struct QuotaState {
29    /// Window start, Unix epoch milliseconds.
30    pub window_start: i64,
31    /// Operations counted so far.
32    pub ops: u32,
33    /// Bytes counted so far.
34    pub bytes: u64,
35}
36
37/// The key a quota is counted under. In M0 that is one signer within one
38/// namespace (planner default Q14). Namespace totals use separate `qs`/`qt`
39/// rows and leave this per-signer scope unchanged.
40#[derive(Debug, Clone, PartialEq, Eq, Hash)]
41pub struct QuotaScope(String);
42
43impl QuotaScope {
44    /// The scope for `signer` in `ns`: `"<namespace>\n<signer hex>"`.
45    #[must_use]
46    pub fn for_signer(ns: &NamespaceKey, signer: &[u8; 32]) -> Self {
47        Self(format!("{}\n{}", ns.as_str(), to_hex(signer)))
48    }
49
50    /// The scope key as a string.
51    #[must_use]
52    pub fn as_str(&self) -> &str {
53        &self.0
54    }
55}
56
57/// One charge an Admission decision asks the write's batch to apply: one
58/// operation and `bytes` against `scope`'s current window under `limits`.
59/// The planner evaluates it with [`evaluate_quota`] on the value it read
60/// and guards that read, so the charge is exact within one partition.
61#[derive(Debug, Clone, PartialEq, Eq)]
62pub struct QuotaCharge {
63    /// The counter charged.
64    pub scope: QuotaScope,
65    /// Bytes charged: 0 for ref writes, the declared size for `UploadPack`.
66    pub bytes: u64,
67    /// The window and caps.
68    pub limits: QuotaLimits,
69}
70
71/// Today's per-signer write quota (planner decision Q14): 300 writes and
72/// 128 MiB of `UploadPack` bytes per one-hour window.
73pub const DEFAULT_WRITE_QUOTA: QuotaLimits = QuotaLimits {
74    window_ms: 3_600_000,
75    max_ops: 300,
76    max_bytes: 128 * 1024 * 1024,
77};
78
79/// How often an active ref shard reconciles its fixed-window usage.
80pub const QUOTA_ROLLUP_MS: u64 = 60_000;
81
82// plan_namespace asserts that ticketed advances have no namespace charge.
83// An admitted write may add a qs guard/put, timer, and initial view seed.
84const _: () = assert!(crate::store::outbox::ADVANCE_SHARED_OPS + 4 <= crate::store::MAX_BATCH_OPS);
85
86/// A fixed-window namespace counter. `u64` also holds sums across shards.
87#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
88pub struct NamespaceUsage {
89    /// Admitted writes.
90    pub ops: u64,
91    /// Admitted upload bytes.
92    pub bytes: u64,
93}
94
95impl NamespaceUsage {
96    /// Add a nondecreasing cumulative counter's delta.
97    #[must_use]
98    pub fn delta_from(self, older: Self) -> Option<Self> {
99        Some(Self {
100            ops: self.ops.checked_sub(older.ops)?,
101            bytes: self.bytes.checked_sub(older.bytes)?,
102        })
103    }
104
105    /// Add two counters, failing closed on overflow.
106    #[must_use]
107    pub fn checked_add(self, other: Self) -> Option<Self> {
108        Some(Self {
109            ops: self.ops.checked_add(other.ops)?,
110            bytes: self.bytes.checked_add(other.bytes)?,
111        })
112    }
113}
114
115/// A ref shard's unguarded copy of the namespace aggregate.
116#[derive(Debug, Clone, Copy, PartialEq, Eq)]
117pub struct NamespaceView {
118    /// Coordinator total at `observed_at_ms`.
119    pub total: NamespaceUsage,
120    /// This shard's contribution already included in `total`.
121    pub pushed: NamespaceUsage,
122    /// Time of the coordinator read, for stale-view metrics.
123    pub observed_at_ms: u64,
124}
125
126/// One default-admission namespace charge, prepared for the write planner.
127#[derive(Debug, Clone, Copy, PartialEq, Eq)]
128pub struct NamespaceCharge {
129    /// The fixed window and caps.
130    pub limits: QuotaLimits,
131    /// Window selected by the request's one read-ahead.
132    pub window: u64,
133    /// Bytes from the admission charge.
134    pub bytes: u64,
135    /// Ref shards roll up; coordinator and Single partitions are exact.
136    pub rollup: bool,
137}
138
139/// Freshness of the unguarded local view.
140#[derive(Debug, Clone, Copy, PartialEq, Eq)]
141pub enum ViewStatus {
142    /// An aggregate was read within two rollup periods.
143    Fresh,
144    /// No view has been installed for this window.
145    Missing,
146    /// The last coordinator read is older than two periods.
147    Stale,
148}
149
150/// The namespace quota's pure admission result.
151#[derive(Debug, Clone, Copy, PartialEq, Eq)]
152pub enum NamespaceDecision {
153    /// Charge this exact local counter in the write batch.
154    Allowed {
155        /// New local cumulative usage.
156        usage: NamespaceUsage,
157        /// Whether the view was usable.
158        view: ViewStatus,
159    },
160    /// Neither the write nor any quota state is allocated.
161    Exhausted,
162}
163
164/// The fixed window number, with floor division at the epoch.
165#[must_use]
166pub fn namespace_window(now_ms: i64, window_ms: i64) -> u64 {
167    let window = window_ms.max(1);
168    u64::try_from(now_ms.div_euclid(window)).unwrap_or(0)
169}
170
171/// Test local exact usage first, then the most recent coordinator estimate.
172/// Missing or stale views admit under the local cap and are metered by the
173/// caller. A view contains this shard's last pushed count, so its estimate
174/// never double-counts that contribution.
175#[must_use]
176pub fn evaluate_namespace(
177    current: NamespaceUsage,
178    view: Option<NamespaceView>,
179    now_ms: i64,
180    charge: NamespaceCharge,
181) -> NamespaceDecision {
182    let Some(usage) = current.checked_add(NamespaceUsage {
183        ops: 1,
184        bytes: charge.bytes,
185    }) else {
186        return NamespaceDecision::Exhausted;
187    };
188    let cap = NamespaceUsage {
189        ops: u64::from(charge.limits.max_ops),
190        bytes: charge.limits.max_bytes,
191    };
192    if usage.ops > cap.ops || usage.bytes > cap.bytes {
193        return NamespaceDecision::Exhausted;
194    }
195    if !charge.rollup {
196        return NamespaceDecision::Allowed {
197            usage,
198            view: ViewStatus::Fresh,
199        };
200    }
201    let status = match view {
202        None => ViewStatus::Missing,
203        Some(v)
204            if u64::try_from(now_ms)
205                .unwrap_or(0)
206                .saturating_sub(v.observed_at_ms)
207                > 2 * QUOTA_ROLLUP_MS =>
208        {
209            ViewStatus::Stale
210        }
211        Some(_) => ViewStatus::Fresh,
212    };
213    if let (ViewStatus::Fresh, Some(view)) = (status, view) {
214        let Some(other) = view.total.delta_from(view.pushed) else {
215            return NamespaceDecision::Exhausted;
216        };
217        let Some(estimate) = other.checked_add(usage) else {
218            return NamespaceDecision::Exhausted;
219        };
220        if estimate.ops > cap.ops || estimate.bytes > cap.bytes {
221            return NamespaceDecision::Exhausted;
222        }
223    }
224    NamespaceDecision::Allowed {
225        usage,
226        view: status,
227    }
228}
229
230/// Decode a snapshot observation and apply the pure check. Corrupt quota
231/// values fail closed like a failed quota read.
232pub(crate) fn check_namespace(
233    current: Option<&Value>,
234    view: Option<&Value>,
235    now_ms: i64,
236    charge: NamespaceCharge,
237) -> Result<NamespaceDecision, ServerError> {
238    let current = current
239        .map(codec::decode_namespace_usage)
240        .transpose()
241        .map_err(|e| ServerError::internal("namespace quota read failed", e.to_string()))?
242        .unwrap_or_default();
243    let view = if charge.rollup {
244        view.map(codec::decode_namespace_view)
245            .transpose()
246            .map_err(|e| ServerError::internal("namespace quota read failed", e.to_string()))?
247    } else {
248        None
249    };
250    Ok(evaluate_namespace(current, view, now_ms, charge))
251}
252
253/// The key of this charge's exact local counter.
254#[must_use]
255pub(crate) fn counter_key(charge: NamespaceCharge, window: u64) -> Key {
256    if charge.rollup {
257        keys::quota_shard(window)
258    } else {
259        keys::quota_total(window)
260    }
261}
262
263/// Append the namespace charge to the same batch as the write. The qv read
264/// has no guard: timer refreshes cannot force a hot write to re-plan.
265pub(crate) fn plan_namespace_charge(
266    charge: NamespaceCharge,
267    current: Option<&Value>,
268    view: Option<&Value>,
269    now_ms: i64,
270    server_now_ms: u64,
271    pre: &mut Vec<Precondition>,
272    puts: &mut Vec<Write>,
273) -> Result<(), ServerError> {
274    let window = charge.window;
275    let key = counter_key(charge, window);
276    let decision = check_namespace(current, view, now_ms, charge)?;
277    let NamespaceDecision::Allowed { usage, .. } = decision else {
278        return Err(ServerError::resource_exhausted(
279            "namespace write op/byte quota exceeded for this window; try again later",
280        ));
281    };
282    pre.push(match current {
283        Some(value) => Precondition::Equals(key.clone(), value.clone()),
284        None => Precondition::Absent(key.clone()),
285    });
286    puts.push(Write::Put(key, codec::encode_namespace_usage(usage)));
287    if current.is_none() {
288        let window_ms = u64::try_from(charge.limits.window_ms.max(1)).unwrap_or(1);
289        let due = if charge.rollup {
290            server_now_ms.saturating_add(QUOTA_ROLLUP_MS)
291        } else {
292            window
293                .saturating_add(1)
294                .saturating_mul(window_ms)
295                .saturating_add(mkit_core::write_auth::MAX_CLOCK_LEAD_MS.unsigned_abs())
296        };
297        puts.push(Write::Put(
298            keys::timer(due, kinds::QUOTA_ROLLUP.get(), &window.to_be_bytes()),
299            codec::encode_u64(window_ms),
300        ));
301    }
302    Ok(())
303}
304
305/// Plan the post-admission charge. A changed window or a quota race after a
306/// lease grant asks the client to retry instead of denying after allocation.
307pub(crate) fn plan_namespace_after_admission(
308    charge: NamespaceCharge,
309    current: Option<&Value>,
310    view: Option<&Value>,
311    now_ms: i64,
312    server_now_ms: u64,
313    lease_committed: bool,
314    pre: &mut Vec<Precondition>,
315    puts: &mut Vec<Write>,
316) -> Result<(), ServerError> {
317    if namespace_window(now_ms, charge.limits.window_ms) != charge.window {
318        return Err(
319            ServerError::aborted_retryable("namespace quota window advanced; retry")
320                .with_abort_cause(AbortCause::QuotaWindow),
321        );
322    }
323    plan_namespace_charge(charge, current, view, now_ms, server_now_ms, pre, puts).map_err(
324        |error| {
325            if lease_committed && error.code() == crate::Code::ResourceExhausted {
326                ServerError::aborted_retryable("namespace quota changed; retry")
327            } else {
328                error
329            }
330        },
331    )
332}
333
334/// The outcome of charging one write against a quota.
335#[derive(Debug, Clone, Copy, PartialEq, Eq)]
336pub enum QuotaDecision {
337    /// Under budget: persist this state in place of the old one and let the
338    /// write proceed.
339    Allowed(QuotaState),
340    /// Over budget: reject the write with `resource_exhausted` and leave the
341    /// stored state untouched.
342    Exhausted {
343        /// Client-safe reason.
344        reason: &'static str,
345    },
346}
347
348/// Charge one write of `incoming_bytes` at server time `now` (Unix epoch
349/// milliseconds) against `current`, the stored state (`None` for a first
350/// write or a pruned row).
351///
352/// A window that has fully elapsed (`now - window_start >= window_ms`)
353/// resets the counters before this write is applied, so a signer is never
354/// charged for activity outside the current window. That permits a burst of
355/// up to twice the budget across a window boundary, never more. Ops are
356/// checked before bytes, so a flood of zero-byte ref writes hits the op cap
357/// on its own. Ref writes charge 0 bytes; `UploadPack` charges its declared
358/// `total_bytes`.
359#[must_use]
360pub fn evaluate_quota(
361    current: Option<QuotaState>,
362    now: i64,
363    incoming_bytes: u64,
364    limits: &QuotaLimits,
365) -> QuotaDecision {
366    let base = match current {
367        Some(s) if now.saturating_sub(s.window_start) < limits.window_ms => s,
368        _ => QuotaState {
369            window_start: now,
370            ops: 0,
371            bytes: 0,
372        },
373    };
374    let ops = base.ops.saturating_add(1);
375    let bytes = base.bytes.saturating_add(incoming_bytes);
376    if ops > limits.max_ops {
377        return QuotaDecision::Exhausted {
378            reason: "write op quota exceeded for this window; try again later",
379        };
380    }
381    if bytes > limits.max_bytes {
382        return QuotaDecision::Exhausted {
383            reason: "write byte quota exceeded for this window; try again later",
384        };
385    }
386    QuotaDecision::Allowed(QuotaState {
387        window_start: base.window_start,
388        ops,
389        bytes,
390    })
391}
392
393#[cfg(test)]
394mod tests {
395    use super::*;
396
397    // The quota tests below are ported verbatim from vcs-worker's former
398    // write_quota.rs, reading its constants from DEFAULT_WRITE_QUOTA.
399    const WRITE_QUOTA_WINDOW_MS: i64 = DEFAULT_WRITE_QUOTA.window_ms;
400    const WRITE_QUOTA_MAX_OPS: u32 = DEFAULT_WRITE_QUOTA.max_ops;
401    const WRITE_QUOTA_MAX_BYTES: u64 = DEFAULT_WRITE_QUOTA.max_bytes;
402
403    fn evaluate(current: Option<QuotaState>, now: i64, incoming_bytes: u64) -> QuotaDecision {
404        evaluate_quota(current, now, incoming_bytes, &DEFAULT_WRITE_QUOTA)
405    }
406
407    #[test]
408    fn default_quota_is_todays_vcs_worker_limits() {
409        assert_eq!(WRITE_QUOTA_WINDOW_MS, 60 * 60 * 1_000);
410        assert_eq!(WRITE_QUOTA_MAX_OPS, 300);
411        assert_eq!(WRITE_QUOTA_MAX_BYTES, 128 * 1024 * 1024);
412    }
413
414    #[test]
415    fn exhausted_reasons_are_verbatim() {
416        let full_ops = QuotaState {
417            window_start: 0,
418            ops: WRITE_QUOTA_MAX_OPS,
419            bytes: 0,
420        };
421        assert_eq!(
422            evaluate(Some(full_ops), 1, 0),
423            QuotaDecision::Exhausted {
424                reason: "write op quota exceeded for this window; try again later"
425            }
426        );
427        assert_eq!(
428            evaluate(None, 1, WRITE_QUOTA_MAX_BYTES + 1),
429            QuotaDecision::Exhausted {
430                reason: "write byte quota exceeded for this window; try again later"
431            }
432        );
433    }
434
435    #[test]
436    fn first_write_from_a_fresh_key_is_allowed() {
437        let d = evaluate(None, 10_000, 1_000);
438        assert_eq!(
439            d,
440            QuotaDecision::Allowed(QuotaState {
441                window_start: 10_000,
442                ops: 1,
443                bytes: 1_000
444            })
445        );
446    }
447
448    #[test]
449    fn under_the_op_cap_stays_allowed() {
450        let state = QuotaState {
451            window_start: 0,
452            ops: WRITE_QUOTA_MAX_OPS - 1,
453            bytes: 0,
454        };
455        let d = evaluate(Some(state), 100, 0);
456        assert_eq!(
457            d,
458            QuotaDecision::Allowed(QuotaState {
459                window_start: 0,
460                ops: WRITE_QUOTA_MAX_OPS,
461                bytes: 0
462            })
463        );
464    }
465
466    #[test]
467    fn at_the_op_cap_is_rejected() {
468        // Already AT the cap: one more op would push it over.
469        let state = QuotaState {
470            window_start: 0,
471            ops: WRITE_QUOTA_MAX_OPS,
472            bytes: 0,
473        };
474        let d = evaluate(Some(state), 100, 0);
475        assert!(matches!(d, QuotaDecision::Exhausted { .. }));
476    }
477
478    #[test]
479    fn exactly_at_the_byte_cap_is_allowed() {
480        // Modeling a single UploadPack landing exactly at the byte cap.
481        let state = QuotaState {
482            window_start: 0,
483            ops: 0,
484            bytes: 0,
485        };
486        let d = evaluate(Some(state), 100, WRITE_QUOTA_MAX_BYTES);
487        assert_eq!(
488            d,
489            QuotaDecision::Allowed(QuotaState {
490                window_start: 0,
491                ops: 1,
492                bytes: WRITE_QUOTA_MAX_BYTES
493            })
494        );
495    }
496
497    #[test]
498    fn one_byte_over_the_cap_is_rejected() {
499        let state = QuotaState {
500            window_start: 0,
501            ops: 0,
502            bytes: WRITE_QUOTA_MAX_BYTES,
503        };
504        let d = evaluate(Some(state), 100, 1);
505        assert!(matches!(d, QuotaDecision::Exhausted { .. }));
506    }
507
508    #[test]
509    fn window_resets_after_it_elapses() {
510        // Exhausted at the tail of a window...
511        let state = QuotaState {
512            window_start: 0,
513            ops: WRITE_QUOTA_MAX_OPS,
514            bytes: 0,
515        };
516        let still_current = evaluate(Some(state), WRITE_QUOTA_WINDOW_MS - 1, 0);
517        assert!(matches!(still_current, QuotaDecision::Exhausted { .. }));
518        // ...but once the window has fully elapsed, a fresh window starts and
519        // the SAME author is allowed again — quota state resets over time.
520        let reset = evaluate(Some(state), WRITE_QUOTA_WINDOW_MS, 0);
521        assert_eq!(
522            reset,
523            QuotaDecision::Allowed(QuotaState {
524                window_start: WRITE_QUOTA_WINDOW_MS,
525                ops: 1,
526                bytes: 0
527            })
528        );
529    }
530
531    #[test]
532    fn ops_and_bytes_are_independent_caps() {
533        // Many zero-byte UpdateRef/AdvanceRefs calls can hit the op cap well
534        // under the byte cap (which only UploadPack ever touches).
535        let mut state = None;
536        let mut now = 0i64;
537        for _ in 0..WRITE_QUOTA_MAX_OPS {
538            match evaluate(state, now, 0) {
539                QuotaDecision::Allowed(s) => state = Some(s),
540                QuotaDecision::Exhausted { .. } => panic!("should still be under the op cap"),
541            }
542            now += 1;
543        }
544        assert!(matches!(
545            evaluate(state, now, 0),
546            QuotaDecision::Exhausted { .. }
547        ));
548    }
549
550    #[test]
551    fn a_single_max_size_pack_is_allowed_but_a_second_is_not() {
552        // service::MAX_PACK_BYTES is 64 MiB; WRITE_QUOTA_MAX_BYTES is 128 MiB,
553        // so exactly two max-size packs fit in one window and a third does not.
554        const MAX_PACK_BYTES: u64 = 64 * 1024 * 1024;
555        let first = evaluate(None, 0, MAX_PACK_BYTES);
556        let state = match first {
557            QuotaDecision::Allowed(s) => s,
558            QuotaDecision::Exhausted { .. } => panic!("first max-size pack should be allowed"),
559        };
560        let second = evaluate(Some(state), 1, MAX_PACK_BYTES);
561        assert!(matches!(second, QuotaDecision::Allowed(_)));
562        let state = match second {
563            QuotaDecision::Allowed(s) => s,
564            QuotaDecision::Exhausted { .. } => unreachable!(),
565        };
566        let third = evaluate(Some(state), 2, 1);
567        assert!(matches!(third, QuotaDecision::Exhausted { .. }));
568    }
569
570    #[test]
571    fn different_authors_are_independent() {
572        // Not modeled in this module (the DO keys the table by author), but
573        // documented here: a fresh `current = None` for a distinct key always
574        // starts a clean window regardless of any other key's state.
575        let exhausted = QuotaState {
576            window_start: 0,
577            ops: WRITE_QUOTA_MAX_OPS,
578            bytes: WRITE_QUOTA_MAX_BYTES,
579        };
580        let _ = exhausted; // another author's state; irrelevant to a fresh `None`
581        let d = evaluate(None, 0, 1);
582        assert_eq!(
583            d,
584            QuotaDecision::Allowed(QuotaState {
585                window_start: 0,
586                ops: 1,
587                bytes: 1
588            })
589        );
590    }
591
592    #[test]
593    fn signer_scope_is_namespace_newline_hex() {
594        let ns = NamespaceKey::deployment_default();
595        let scope = QuotaScope::for_signer(&ns, &[0xab; 32]);
596        assert_eq!(scope.as_str(), format!("root\n{}", "ab".repeat(32)));
597        assert_ne!(scope, QuotaScope::for_signer(&ns, &[0xac; 32]));
598    }
599
600    #[test]
601    fn namespace_local_exhaustion_is_exact_even_without_a_view() {
602        let charge = NamespaceCharge {
603            limits: QuotaLimits {
604                window_ms: 1_000,
605                max_ops: 2,
606                max_bytes: 5,
607            },
608            window: 0,
609            bytes: 2,
610            rollup: true,
611        };
612        let first = evaluate_namespace(NamespaceUsage::default(), None, 0, charge);
613        assert_eq!(
614            first,
615            NamespaceDecision::Allowed {
616                usage: NamespaceUsage { ops: 1, bytes: 2 },
617                view: ViewStatus::Missing,
618            }
619        );
620        let second = evaluate_namespace(NamespaceUsage { ops: 1, bytes: 2 }, None, 1, charge);
621        assert!(matches!(second, NamespaceDecision::Allowed { .. }));
622        assert_eq!(
623            evaluate_namespace(NamespaceUsage { ops: 2, bytes: 4 }, None, 2, charge),
624            NamespaceDecision::Exhausted
625        );
626        let too_many_bytes = NamespaceCharge { bytes: 6, ..charge };
627        assert_eq!(
628            evaluate_namespace(NamespaceUsage::default(), None, 0, too_many_bytes),
629            NamespaceDecision::Exhausted
630        );
631    }
632
633    #[test]
634    fn namespace_fixed_window_rolls_over_without_reusing_the_old_counter() {
635        let limits = QuotaLimits {
636            window_ms: 600_000,
637            max_ops: 1,
638            max_bytes: 0,
639        };
640        let old = namespace_window(599_999, limits.window_ms);
641        let new = namespace_window(600_000, limits.window_ms);
642        assert_eq!((old, new), (0, 1));
643        assert_ne!(keys::quota_shard(old), keys::quota_shard(new));
644        let charge = NamespaceCharge {
645            limits,
646            window: new,
647            bytes: 0,
648            rollup: true,
649        };
650        assert!(matches!(
651            evaluate_namespace(NamespaceUsage::default(), None, 600_000, charge),
652            NamespaceDecision::Allowed { .. }
653        ));
654    }
655
656    #[test]
657    fn namespace_view_subtracts_own_pushed_count_and_expires() {
658        let charge = NamespaceCharge {
659            limits: QuotaLimits {
660                window_ms: 1_000_000,
661                max_ops: 5,
662                max_bytes: 10,
663            },
664            window: 0,
665            bytes: 0,
666            rollup: true,
667        };
668        let local = NamespaceUsage { ops: 2, bytes: 0 };
669        let view = NamespaceView {
670            total: NamespaceUsage { ops: 4, bytes: 0 },
671            pushed: local,
672            observed_at_ms: 0,
673        };
674        assert!(matches!(
675            evaluate_namespace(local, Some(view), 1, charge),
676            NamespaceDecision::Allowed {
677                usage: NamespaceUsage { ops: 3, .. },
678                view: ViewStatus::Fresh
679            }
680        ));
681        assert_eq!(
682            evaluate_namespace(NamespaceUsage { ops: 3, bytes: 0 }, Some(view), 1, charge),
683            NamespaceDecision::Exhausted
684        );
685        assert!(matches!(
686            evaluate_namespace(
687                NamespaceUsage { ops: 3, bytes: 0 },
688                Some(view),
689                (2 * QUOTA_ROLLUP_MS + 1).cast_signed(),
690                charge
691            ),
692            NamespaceDecision::Allowed {
693                view: ViewStatus::Stale,
694                ..
695            }
696        ));
697        let exact = NamespaceCharge {
698            rollup: false,
699            ..charge
700        };
701        assert_eq!(
702            evaluate_namespace(NamespaceUsage { ops: 5, bytes: 0 }, None, 0, exact),
703            NamespaceDecision::Exhausted
704        );
705    }
706
707    #[test]
708    fn simulated_overshoot_is_bounded_by_other_shards_last_three_periods() {
709        const SHARDS: usize = 5;
710        const CAP: u32 = 80;
711        let charge = NamespaceCharge {
712            limits: QuotaLimits {
713                window_ms: 1_000_000,
714                max_ops: CAP,
715                max_bytes: 0,
716            },
717            window: 0,
718            bytes: 0,
719            rollup: true,
720        };
721        let mut local = [NamespaceUsage::default(); SHARDS];
722        let mut views = [None; SHARDS];
723        let mut writes: Vec<(u64, usize)> = Vec::new();
724        for second in 0..180_u64 {
725            let now = second * 1_000;
726            if second > 0 && now.is_multiple_of(QUOTA_ROLLUP_MS) {
727                let total = NamespaceUsage {
728                    ops: writes.len() as u64,
729                    bytes: 0,
730                };
731                for shard in 0..SHARDS {
732                    views[shard] = Some(NamespaceView {
733                        total,
734                        pushed: local[shard],
735                        observed_at_ms: now,
736                    });
737                }
738            }
739            for shard in 0..SHARDS {
740                if let NamespaceDecision::Allowed { usage, .. } =
741                    evaluate_namespace(local[shard], views[shard], now.cast_signed(), charge)
742                {
743                    local[shard] = usage;
744                    writes.push((now, shard));
745                    let overshoot = writes.len().saturating_sub(CAP as usize);
746                    let other_recent = writes
747                        .iter()
748                        .filter(|(at, source)| {
749                            *source != shard && now.saturating_sub(*at) <= 3 * QUOTA_ROLLUP_MS
750                        })
751                        .count();
752                    assert!(
753                        overshoot <= other_recent,
754                        "t={now} shard={shard}: {overshoot} > {other_recent}"
755                    );
756                }
757            }
758        }
759        assert!(
760            writes.len() > CAP as usize,
761            "simulation must exercise overshoot"
762        );
763    }
764
765    #[test]
766    fn rollup_timer_uses_server_time_even_with_business_clock_skew() {
767        let business_now_ms = 180_000;
768        let server_now_ms = 1_000;
769        let charge = NamespaceCharge {
770            limits: QuotaLimits {
771                window_ms: 60_000,
772                max_ops: 2,
773                max_bytes: 0,
774            },
775            window: namespace_window(business_now_ms, 60_000),
776            bytes: 0,
777            rollup: true,
778        };
779        let (mut pre, mut writes) = (Vec::new(), Vec::new());
780        plan_namespace_charge(
781            charge,
782            None,
783            None,
784            business_now_ms,
785            server_now_ms,
786            &mut pre,
787            &mut writes,
788        )
789        .unwrap();
790        assert!(writes.iter().any(|write| matches!(write,
791            Write::Put(key, _) if matches!(keys::parse(key), Some(keys::ParsedKey::Timer { due_at_ms: 61_000, kind: 5, .. }))
792        )));
793    }
794
795    #[cfg(feature = "memory")]
796    #[test]
797    fn single_and_coordinator_exact_paths_commit_the_total() {
798        use crate::store::{Batch, BatchOutcome, NamespaceStore, Partition};
799        use crate::{MemoryKv, NamespaceKey};
800        use futures_executor::block_on;
801
802        let store = MemoryKv::default();
803        let charge = NamespaceCharge {
804            limits: QuotaLimits {
805                window_ms: 60_000,
806                max_ops: 2,
807                max_bytes: 4,
808            },
809            window: 0,
810            bytes: 2,
811            rollup: false,
812        };
813        for partition in [
814            Partition::Namespace(NamespaceKey::deployment_default()),
815            Partition::Coordinator(NamespaceKey::deployment_default()),
816        ] {
817            let key = keys::quota_total(0);
818            for _ in 0..2 {
819                let current = block_on(store.get(&partition, &key)).unwrap();
820                let (mut pre, mut writes) = (Vec::new(), Vec::new());
821                plan_namespace_charge(charge, current.as_ref(), None, 0, 0, &mut pre, &mut writes)
822                    .unwrap();
823                assert_eq!(
824                    block_on(store.apply(
825                        &partition,
826                        Batch {
827                            preconditions: pre,
828                            writes
829                        }
830                    ))
831                    .unwrap(),
832                    BatchOutcome::Committed
833                );
834            }
835            let current = block_on(store.get(&partition, &key)).unwrap();
836            let (mut pre, mut writes) = (Vec::new(), Vec::new());
837            assert_eq!(
838                plan_namespace_charge(charge, current.as_ref(), None, 0, 0, &mut pre, &mut writes)
839                    .unwrap_err()
840                    .code(),
841                crate::Code::ResourceExhausted
842            );
843            assert!(pre.is_empty() && writes.is_empty());
844            assert_eq!(
845                codec::decode_namespace_usage(current.as_ref().unwrap()).unwrap(),
846                NamespaceUsage { ops: 2, bytes: 4 }
847            );
848            assert!(
849                block_on(store.get(&partition, &keys::quota_shard(0)))
850                    .unwrap()
851                    .is_none()
852            );
853        }
854    }
855}