Skip to main content

headgate_testkit/
lib.rs

1//! headgate-testkit — the in-process test double: a complete in-memory [`Store`] so
2//! handler code, retry behavior, steps, and runner wiring can be tested with no
3//! database at all (River's rivertest, asynq's asynqtest). The Go twin is
4//! `go/headgatetest`; both keep the same split:
5//!
6//! FAITHFUL: the transition table (every ack outcome, fence-gated identity,
7//! `LeaseRejected` for a superseded holder), attempts vs crash_attempts, quarantine at
8//! the crash limit, job uniqueness uniqueness in both modes, retention policy ephemeral retention-0 delete,
9//! retention and eviction contract retention eviction (quarantined exempt), per-partition round-robin admission,
10//! priority ordering, duty leases.
11//!
12//! SIMPLIFIED, capability-honestly (runtime capability boundary): `caps()` is 0 — no Transactional (so
13//! `once`/`step_once` error, as they must without a real transaction), no Inspect
14//! (scheduler/operations/quarantine duties idle), no Notifying (workers poll). Like
15//! the SQL backends it admits `state = available` only — pair `admit` with
16//! `promote_due` (the worker and `testing::drain` already do). An unconfigured rate
17//! class is UNLIMITED here; configure one with [`MemStore::set_rate_limit`]. Time is
18//! the store's own clock and tests can freeze or step it — no sleeps.
19
20use std::collections::HashMap;
21use std::sync::Mutex;
22use std::time::Duration;
23
24use headgate_core::{
25    AdmissionUnit, AdmitRequest, Caps, Checkpoint, Claim, Envelope, Inspect, LeaseRef, Outcome,
26    Reclaimed, Store, StoreError,
27};
28
29/// Shared live-backend proof for strict sticky routing. A 5,000-job high-priority
30/// backlog pinned to another worker must not fill the bounded candidate draw and hide
31/// the caller's pinned or ordinary work. Rate-limited requeue then proves the route is
32/// durable lifecycle state rather than lease metadata.
33pub async fn assert_sticky_routing(store: std::sync::Arc<dyn Store>, backend: &str) -> String {
34    static RUN: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1);
35    let run = RUN.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
36    let queue = format!("sticky-{backend}-{}-{run}", std::process::id());
37    let env = |id: String, sticky: &str, priority: i32| Envelope {
38        id,
39        kind: "test:sticky".into(),
40        payload: b"{}".to_vec(),
41        queue: queue.clone(),
42        partition_key: "tenant".into(),
43        fingerprint: format!("fp-sticky-{backend}"),
44        priority,
45        sticky_worker: sticky.into(),
46        scheduled_at_ms: 1,
47        retention_ms: 86_400_000,
48        ..Default::default()
49    };
50    let a_id = format!("{queue}-a");
51    let general_id = format!("{queue}-general");
52    let mut batch = Vec::with_capacity(5_002);
53    for i in 0..5_000 {
54        batch.push(env(format!("{queue}-b-{i:04}"), "worker-b", 10_000));
55    }
56    batch.push(env(a_id.clone(), "worker-a", 50));
57    batch.push(env(general_id.clone(), "", 1));
58    // Keep each producer batch portable: MySQL caps one prepared statement at 65,535
59    // placeholders, while the proof needs 5,000 rows in the partition.
60    for chunk in batch.chunks(500) {
61        store.enqueue(chunk).await.expect("sticky enqueue");
62    }
63
64    let req = |worker: &str, lease: &str, capacity| AdmitRequest {
65        worker: worker.into(),
66        lease_id: lease.into(),
67        queues: vec![queue.clone()],
68        capacity,
69        lease: Duration::from_secs(60),
70        quantum: 10_000,
71    };
72    let units = store
73        .admit(req("worker-a", "sticky-la", 2))
74        .await
75        .expect("worker-a admit");
76    let claims: Vec<_> = units.iter().flat_map(|u| &u.claims).collect();
77    let mut ids: Vec<_> = claims.iter().map(|c| c.envelope.id.as_str()).collect();
78    ids.sort_unstable();
79    assert_eq!(ids, vec![a_id.as_str(), general_id.as_str()]);
80    assert_eq!(
81        claims
82            .iter()
83            .find(|c| c.envelope.id == a_id)
84            .unwrap()
85            .envelope
86            .sticky_worker,
87        "worker-a"
88    );
89
90    let lease_for = |id: &str| {
91        claims
92            .iter()
93            .find(|c| c.envelope.id == id)
94            .unwrap()
95            .lease_ref()
96    };
97    store
98        .ack(&lease_for(&a_id), Outcome::RateLimited, None, None)
99        .await
100        .expect("route-preserving requeue");
101    store
102        .ack(&lease_for(&general_id), Outcome::Success, None, None)
103        .await
104        .expect("general completion");
105
106    assert!(
107        store
108            .admit(req("worker-c", "sticky-lc", 2))
109            .await
110            .expect("worker-c admit")
111            .is_empty(),
112        "another worker must not claim pinned work"
113    );
114    let a_again = store
115        .admit(req("worker-a", "sticky-la2", 1))
116        .await
117        .expect("worker-a re-admit");
118    assert_eq!(a_again[0].claims[0].envelope.id, a_id);
119    let b = store
120        .admit(req("worker-b", "sticky-lb", 1))
121        .await
122        .expect("worker-b admit");
123    assert!(
124        b[0].claims[0]
125            .envelope
126            .id
127            .starts_with(&format!("{queue}-b-"))
128    );
129    queue
130}
131
132/// Shared live-backend proof for bounded enqueue. It deliberately drives contention:
133/// 64 producers race for 25 slots, then the helper verifies idempotent replay, atomic
134/// batch rejection, capacity release on terminalization, lowering below current depth,
135/// and disabling the policy. Every assertion uses the public Store/Inspect ports.
136pub async fn assert_enqueue_backpressure(store: std::sync::Arc<dyn Inspect>, queue: &str) {
137    let queue = queue.to_string();
138    let envelope = |id: String| Envelope {
139        id,
140        kind: "test:backpressure".into(),
141        payload: b"{}".to_vec(),
142        queue: queue.clone(),
143        fingerprint: format!("fp-backpressure-{queue}"),
144        scheduled_at_ms: 1,
145        retention_ms: 86_400_000,
146        ..Default::default()
147    };
148
149    store
150        .set_enqueue_limit(&queue, Some(25))
151        .await
152        .expect("configure enqueue limit");
153    let mut tasks = Vec::new();
154    for i in 0..64 {
155        let store = store.clone();
156        let job = envelope(format!("{queue}-bp-{i}"));
157        tasks.push(tokio::spawn(async move {
158            let id = job.id.clone();
159            (id, store.enqueue(&[job]).await)
160        }));
161    }
162    let mut accepted = Vec::new();
163    let mut rejected = 0;
164    for task in tasks {
165        let (id, result) = task.await.expect("producer task");
166        match result {
167            Ok(()) => accepted.push(id),
168            Err(StoreError::Backpressure {
169                queue: rejected_queue,
170                limit,
171                current,
172                incoming,
173            }) => {
174                assert_eq!(rejected_queue, queue);
175                assert_eq!(limit, 25);
176                assert_eq!(incoming, 1);
177                assert!(current <= 25);
178                rejected += 1;
179            }
180            other => panic!("unexpected concurrent enqueue result: {other:?}"),
181        }
182    }
183    assert_eq!(accepted.len(), 25, "the store must never over-admit");
184    assert_eq!(rejected, 39);
185
186    let stats = store.queue_stats().await.expect("queue stats");
187    let stat = stats.iter().find(|s| s.queue == queue).expect("queue stat");
188    assert_eq!(stat.unfinished_jobs, 25);
189    assert_eq!(stat.max_unfinished_jobs, Some(25));
190
191    // Matching-id replay is success and consumes no additional capacity.
192    store
193        .enqueue(&[envelope(accepted[0].clone())])
194        .await
195        .expect("idempotent replay at limit");
196
197    let batch = [
198        envelope(format!("{queue}-batch-a")),
199        envelope(format!("{queue}-batch-b")),
200    ];
201    match store.enqueue(&batch).await {
202        Err(StoreError::Backpressure {
203            limit,
204            current,
205            incoming,
206            ..
207        }) => assert_eq!((limit, current, incoming), (25, 25, 2)),
208        other => panic!("full batch should be rejected atomically: {other:?}"),
209    }
210    for id in [&batch[0].id, &batch[1].id] {
211        assert!(
212            store
213                .get_job(id, false)
214                .await
215                .expect("rejected lookup")
216                .is_none(),
217            "a rejected batch must write no rows"
218        );
219    }
220
221    store
222        .operator_cancel(&accepted[0])
223        .await
224        .expect("terminalization releases one slot");
225    store
226        .enqueue(&[envelope(format!("{queue}-replacement"))])
227        .await
228        .expect("replacement after drain");
229
230    store
231        .set_enqueue_limit(&queue, Some(10))
232        .await
233        .expect("lower limit below current depth");
234    assert!(matches!(
235        store
236            .enqueue(&[envelope(format!("{queue}-still-full"))])
237            .await,
238        Err(StoreError::Backpressure {
239            limit: 10,
240            current: 25,
241            incoming: 1,
242            ..
243        })
244    ));
245    store
246        .set_enqueue_limit(&queue, None)
247        .await
248        .expect("disable enqueue limit");
249    store
250        .enqueue(&[envelope(format!("{queue}-unbounded"))])
251        .await
252        .expect("disabled policy accepts");
253    let stats = store.queue_stats().await.expect("final queue stats");
254    let stat = stats.iter().find(|s| s.queue == queue).expect("final stat");
255    assert_eq!(stat.unfinished_jobs, 26);
256    assert_eq!(stat.max_unfinished_jobs, None);
257}
258
259mod database;
260pub use database::{
261    MysqlTestDatabase, PostgresTestDatabase, RedisTestNamespace, TestDatabaseError,
262};
263
264#[derive(Default)]
265struct MemJob {
266    env: Envelope,
267    state: String,
268    fence: u64,
269    lease_id: String,
270    lease_expires: i64,
271    finalized_at: i64,
272    checkpoint: Checkpoint,
273    errs: Vec<String>,
274    /// Estimated cost actually charged for this attempt; zero is fail-open.
275    rate_charge: i64,
276    result: Option<headgate_core::JobResult>,
277    output: Option<headgate_core::JobOutput>,
278    progress: Option<headgate_core::JobProgress>,
279}
280
281struct RateBucket {
282    tokens: i64,
283    burst: i64,
284    limit: i64,
285    window: i64,
286    refilled: i64,
287}
288
289enum Clock {
290    System,
291    Frozen(i64),
292}
293
294#[derive(Default)]
295struct Inner {
296    jobs: HashMap<String, MemJob>,
297    unique: HashMap<Vec<u8>, String>,
298    throttle: HashMap<Vec<u8>, (String, i64)>,
299    quarantine: HashMap<String, bool>,
300    paused: HashMap<String, bool>,
301    rate: HashMap<String, RateBucket>,
302    duties: HashMap<String, (String, i64)>,
303    rr: HashMap<String, usize>,
304}
305
306pub struct MemStore {
307    inner: Mutex<Inner>,
308    clock: Mutex<Clock>,
309    /// crash quarantine quarantine threshold.
310    pub crash_limit: u32,
311    pub retry_base_ms: i64,
312    pub retry_cap_ms: i64,
313}
314
315impl Default for MemStore {
316    fn default() -> Self {
317        Self::new()
318    }
319}
320
321impl MemStore {
322    pub fn new() -> Self {
323        Self {
324            inner: Mutex::new(Inner::default()),
325            clock: Mutex::new(Clock::System),
326            crash_limit: 3,
327            retry_base_ms: 1_000,
328            retry_cap_ms: 3_600_000,
329        }
330    }
331
332    fn now(&self) -> i64 {
333        match *self.clock.lock().unwrap() {
334            Clock::System => std::time::SystemTime::now()
335                .duration_since(std::time::UNIX_EPOCH)
336                .unwrap()
337                .as_millis() as i64,
338            Clock::Frozen(ms) => ms,
339        }
340    }
341
342    // ---------- test-facing helpers ----------
343
344    /// Freeze the STORE clock (boundary validation: store-supplied time, even here) at `ms`.
345    pub fn freeze_clock_at(&self, ms: i64) {
346        *self.clock.lock().unwrap() = Clock::Frozen(ms);
347    }
348
349    /// Step a frozen clock forward — deterministic backoff/retention tests, no sleeps.
350    /// Freezes at system-now first if the clock was live.
351    pub fn advance_clock(&self, by_ms: i64) {
352        let now = self.now();
353        *self.clock.lock().unwrap() = Clock::Frozen(now + by_ms);
354    }
355
356    pub fn unfreeze_clock(&self) {
357        *self.clock.lock().unwrap() = Clock::System;
358    }
359
360    /// (envelope snapshot, state) — `None` if the job does not exist (deleted counts).
361    pub fn job_state(&self, id: &str) -> Option<(Envelope, String)> {
362        let inner = self.inner.lock().unwrap();
363        inner.jobs.get(id).map(|j| (j.env.clone(), j.state.clone()))
364    }
365
366    /// The per-attempt error history recorded for a job.
367    pub fn errors(&self, id: &str) -> Vec<String> {
368        let inner = self.inner.lock().unwrap();
369        inner
370            .jobs
371            .get(id)
372            .map(|j| j.errs.clone())
373            .unwrap_or_default()
374    }
375
376    /// state -> count for one queue (`None` = all queues).
377    pub fn counts(&self, queue: Option<&str>) -> HashMap<String, usize> {
378        let inner = self.inner.lock().unwrap();
379        let mut out = HashMap::new();
380        for j in inner.jobs.values() {
381            if queue.is_none_or(|q| q == j.env.queue) {
382                *out.entry(j.state.clone()).or_insert(0) += 1;
383            }
384        }
385        out
386    }
387
388    pub fn set_queue_paused(&self, queue: &str, paused: bool) {
389        self.inner
390            .lock()
391            .unwrap()
392            .paused
393            .insert(queue.into(), paused);
394    }
395
396    /// Configure a fleet token bucket. Unconfigured classes are unlimited here.
397    pub fn set_rate_limit(&self, name: &str, limit: i64, window_ms: i64, burst: i64) {
398        let now = self.now();
399        self.inner.lock().unwrap().rate.insert(
400            name.into(),
401            RateBucket {
402                tokens: burst,
403                burst,
404                limit,
405                window: window_ms,
406                refilled: now,
407            },
408        );
409    }
410}
411
412// ---------------------------------------------------------------------------
413// failure classification ASSERT-ENQUEUED . River ships `rivertest.RequireInserted`; this register
414// row claimed the same affordance and 's evidence linter found NO helper of any
415// name in either language, so every test that wanted "did the producer enqueue what I
416// think it did" hand-rolled a `job_state(id)` lookup — which needs the id, i.e. needs the
417// test to already know the answer to the question. This is the version that takes a
418// DESCRIPTION instead, and whose failure names what it found instead.
419// ---------------------------------------------------------------------------
420
421/// A description of an enqueue. `kind` is required; every other field is an optional
422/// matcher, and `None` means "do not care".
423#[derive(Clone, Debug, Default)]
424pub struct Enqueued {
425    pub kind: String,
426    pub queue: Option<String>,
427    pub payload: Option<Vec<u8>>,
428    pub scheduled_at_ms: Option<i64>,
429    pub partition_key: Option<String>,
430    /// Exactly this many matches. `None` means "at least one".
431    pub count: Option<usize>,
432}
433
434impl Enqueued {
435    pub fn of_kind(kind: &str) -> Self {
436        Self {
437            kind: kind.into(),
438            ..Default::default()
439        }
440    }
441    pub fn in_queue(mut self, q: &str) -> Self {
442        self.queue = Some(q.into());
443        self
444    }
445    pub fn with_payload(mut self, p: impl AsRef<[u8]>) -> Self {
446        self.payload = Some(p.as_ref().to_vec());
447        self
448    }
449    pub fn scheduled_at(mut self, ms: i64) -> Self {
450        self.scheduled_at_ms = Some(ms);
451        self
452    }
453    pub fn in_partition(mut self, k: &str) -> Self {
454        self.partition_key = Some(k.into());
455        self
456    }
457    pub fn times(mut self, n: usize) -> Self {
458        self.count = Some(n);
459        self
460    }
461
462    fn matches(&self, e: &Envelope) -> bool {
463        e.kind == self.kind
464            && self.queue.as_ref().is_none_or(|q| *q == e.queue)
465            && self.payload.as_ref().is_none_or(|p| *p == e.payload)
466            && self.scheduled_at_ms.is_none_or(|s| s == e.scheduled_at_ms)
467            && self
468                .partition_key
469                .as_ref()
470                .is_none_or(|k| *k == e.partition_key)
471    }
472
473    fn describe(&self) -> String {
474        let mut s = format!("kind `{}`", self.kind);
475        if let Some(q) = &self.queue {
476            s.push_str(&format!(", queue `{q}`"));
477        }
478        if let Some(p) = &self.payload {
479            s.push_str(&format!(", payload `{}`", String::from_utf8_lossy(p)));
480        }
481        if let Some(ms) = self.scheduled_at_ms {
482            s.push_str(&format!(", scheduled_at_ms {ms}"));
483        }
484        if let Some(k) = &self.partition_key {
485            s.push_str(&format!(", partition_key `{k}`"));
486        }
487        if let Some(n) = self.count {
488            s.push_str(&format!(", exactly {n} time(s)"));
489        }
490        s
491    }
492}
493
494/// Whatever a test double can list back. Implemented for [`MemStore`]; a live backend
495/// implements it over `Inspect::list_jobs` in the test that needs it.
496pub trait EnqueuedJobs {
497    /// Every job the store currently holds, id-ordered. A job DELETED (retention policy ephemeral
498    /// retention-0, retention and eviction contract eviction, `revoke`) is gone from here, which is the honest
499    /// answer: "was enqueued" is only observable while the row exists.
500    fn all_enqueued(&self) -> Vec<Envelope>;
501}
502
503impl EnqueuedJobs for MemStore {
504    fn all_enqueued(&self) -> Vec<Envelope> {
505        let inner = self.inner.lock().unwrap();
506        let mut out: Vec<Envelope> = inner.jobs.values().map(|j| j.env.clone()).collect();
507        out.sort_by(|a, b| a.id.cmp(&b.id));
508        out
509    }
510}
511
512/// Find the enqueued jobs matching `want`, or an error saying what was there instead.
513///
514/// The error is the deliverable. `assert!(store.job_state("x").is_some())` tells you a
515/// lookup failed; this tells you the store held two `mail:welcome` jobs on queue `default`
516/// when you expected one on `priority`, which is the difference between a failing test and
517/// a debugged one.
518pub fn find_enqueued<S: EnqueuedJobs>(store: &S, want: &Enqueued) -> Result<Vec<Envelope>, String> {
519    let all = store.all_enqueued();
520    let hits: Vec<Envelope> = all.iter().filter(|e| want.matches(e)).cloned().collect();
521    let ok = match want.count {
522        Some(n) => hits.len() == n,
523        None => !hits.is_empty(),
524    };
525    if ok {
526        return Ok(hits);
527    }
528    let mut msg = format!(
529        "assert_enqueued: no job matches {} — {} match(es) found among {} enqueued job(s)",
530        want.describe(),
531        hits.len(),
532        all.len()
533    );
534    if all.is_empty() {
535        msg.push_str("\n  the store is EMPTY: nothing was enqueued at all");
536    } else {
537        msg.push_str("\n  what IS enqueued:");
538        for e in all.iter().take(20) {
539            msg.push_str(&format!(
540                "\n    id=`{}` kind=`{}` queue=`{}` partition=`{}` scheduled_at_ms={} payload=`{}`",
541                e.id,
542                e.kind,
543                e.queue,
544                e.partition_key,
545                e.scheduled_at_ms,
546                String::from_utf8_lossy(&e.payload)
547            ));
548        }
549        if all.len() > 20 {
550            msg.push_str(&format!("\n    ... and {} more", all.len() - 20));
551        }
552    }
553    Err(msg)
554}
555
556/// [`find_enqueued`], panicking with that message. The assertion form.
557pub fn assert_enqueued<S: EnqueuedJobs>(store: &S, want: &Enqueued) -> Vec<Envelope> {
558    match find_enqueued(store, want) {
559        Ok(hits) => hits,
560        Err(msg) => panic!("{msg}"),
561    }
562}
563
564fn default_backoff(attempt: i64, base: i64, cap: i64) -> i64 {
565    let shift = attempt.saturating_sub(1).min(20) as u32;
566    (base << shift).min(cap)
567}
568
569fn release_unique(inner: &mut Inner, id: &str) {
570    let Some(j) = inner.jobs.get(id) else { return };
571    if let Some(k) = headgate_core::effective_unique_key(&j.env) {
572        if j.env.unique_window_ms == 0 && inner.unique.get(&k).map(String::as_str) == Some(id) {
573            inner.unique.remove(&k);
574        }
575    }
576}
577
578#[async_trait::async_trait]
579impl Store for MemStore {
580    fn as_result_store(&self) -> Option<&dyn headgate_core::ResultStore> {
581        Some(self)
582    }
583
584    fn as_output_store(&self) -> Option<&dyn headgate_core::OutputStore> {
585        Some(self)
586    }
587
588    fn as_progress_store(&self) -> Option<&dyn headgate_core::ProgressStore> {
589        Some(self)
590    }
591
592    async fn enqueue(&self, batch: &[Envelope]) -> Result<(), StoreError> {
593        let now = self.now();
594        // typed dispatch / boundary validation / idempotent enqueue identity one shared boundary check for every backend.
595        headgate_core::validate_enqueue(batch)?;
596        let mut inner = self.inner.lock().unwrap();
597        // idempotent enqueue identity the id pass, over the WHOLE batch before any other check so all four
598        // backends classify a mixed batch identically. Matching content is skipped —
599        // idempotent success, no re-write, and no unique-key check that would find the
600        // job conflicting with ITSELF. A terminal job's row still exists, so id reuse
601        // follows retention eviction.
602        let mut skip = vec![false; batch.len()];
603        for (i, e) in batch.iter().enumerate() {
604            if let Some(j) = inner.jobs.get(&e.id) {
605                if headgate_core::same_job_content(e, &j.env.kind, &j.env.fingerprint, &j.env.queue)
606                {
607                    skip[i] = true;
608                } else {
609                    return Err(StoreError::IdConflict {
610                        job_id: e.id.clone(),
611                    });
612                }
613            }
614        }
615        // Validate pass — all-or-nothing, like the batch enqueues in the real backends.
616        for (i, e) in batch.iter().enumerate() {
617            if skip[i] {
618                continue;
619            }
620            if !e.fingerprint.is_empty() && inner.quarantine.contains_key(&e.fingerprint) {
621                return Err(StoreError::Quarantined {
622                    fingerprint: e.fingerprint.clone(),
623                });
624            }
625            if let Some(k) = headgate_core::effective_unique_key(e) {
626                let holder = if e.unique_window_ms > 0 {
627                    if let Some((id, expiry)) = inner.throttle.get(&k) {
628                        if *expiry > now {
629                            Some(id.clone())
630                        } else {
631                            None
632                        }
633                    } else {
634                        None
635                    }
636                } else {
637                    inner.unique.get(&k).cloned()
638                };
639                if let Some(id) = holder {
640                    let mut replaced = false;
641                    if e.unique_replace != 0 || e.unique_debounce_ms > 0 {
642                        if let Some(job) = inner.jobs.get_mut(&id) {
643                            if matches!(job.state.as_str(), "scheduled" | "available" | "retryable")
644                            {
645                                let mask = e.unique_replace;
646                                if e.unique_debounce_ms > 0 {
647                                    job.env.schema_version = if e.schema_version == 0 {
648                                        1
649                                    } else {
650                                        e.schema_version
651                                    };
652                                    job.env.payload.clone_from(&e.payload);
653                                    job.env.fingerprint.clone_from(&e.fingerprint);
654                                    job.env.tags = headgate_core::canonical_tags(&e.tags);
655                                    job.env.scheduled_at_ms = now + e.unique_debounce_ms;
656                                    job.state = "scheduled".into();
657                                    replaced = true;
658                                }
659                                if mask & headgate_core::UNIQUE_REPLACE_PAYLOAD != 0 {
660                                    job.env.schema_version = if e.schema_version == 0 {
661                                        1
662                                    } else {
663                                        e.schema_version
664                                    };
665                                    job.env.payload.clone_from(&e.payload);
666                                    job.env.fingerprint.clone_from(&e.fingerprint);
667                                    replaced = true;
668                                }
669                                if mask & headgate_core::UNIQUE_REPLACE_SCHEDULED_AT != 0
670                                    && job.state == "scheduled"
671                                {
672                                    job.env.scheduled_at_ms = if e.scheduled_at_ms == 0 {
673                                        now
674                                    } else {
675                                        e.scheduled_at_ms
676                                    };
677                                    replaced = true;
678                                }
679                                if mask & headgate_core::UNIQUE_REPLACE_PRIORITY != 0 {
680                                    job.env.priority = e.priority;
681                                    replaced = true;
682                                }
683                                if mask & headgate_core::UNIQUE_REPLACE_MAX_ATTEMPTS != 0 {
684                                    job.env.max_attempts = if e.max_attempts == 0 {
685                                        25
686                                    } else {
687                                        e.max_attempts
688                                    };
689                                    replaced = true;
690                                }
691                            }
692                        }
693                    }
694                    return Err(StoreError::Duplicate {
695                        existing_id: id,
696                        replaced,
697                    });
698                }
699            }
700        }
701        for (i, e) in batch.iter().enumerate() {
702            if skip[i] {
703                continue;
704            }
705            let mut env = e.clone();
706            if env.queue.is_empty() {
707                env.queue = "default".into();
708            }
709            if env.max_attempts == 0 {
710                env.max_attempts = 25;
711            }
712            if env.schema_version == 0 {
713                env.schema_version = 1;
714            }
715            env.weight = headgate_core::effective_weight(env.weight);
716            env.tags = headgate_core::canonical_tags(&env.tags);
717            if env.unique_debounce_ms > 0 {
718                env.scheduled_at_ms = now + env.unique_debounce_ms;
719            } else if env.scheduled_at_ms == 0 {
720                env.scheduled_at_ms = now;
721            }
722            let state = if env.pending {
723                "pending"
724            } else if env.scheduled_at_ms > now {
725                "scheduled"
726            } else {
727                "available"
728            };
729            if let Some(k) = headgate_core::effective_unique_key(&env) {
730                if env.unique_window_ms > 0 {
731                    inner
732                        .throttle
733                        .insert(k, (env.id.clone(), now + env.unique_window_ms));
734                } else {
735                    inner.unique.insert(k, env.id.clone());
736                }
737            }
738            inner.jobs.insert(
739                env.id.clone(),
740                MemJob {
741                    state: state.into(),
742                    env,
743                    ..Default::default()
744                },
745            );
746        }
747        Ok(())
748    }
749
750    async fn admit(&self, req: AdmitRequest) -> Result<Vec<AdmissionUnit>, StoreError> {
751        if req.lease.is_zero() {
752            return Err(StoreError::Invalid("lease must be >= 1ms".into()));
753        }
754        let now = self.now();
755        let mut inner = self.inner.lock().unwrap();
756        let mut units = Vec::new();
757        let mut taken: HashMap<String, i64> = HashMap::new();
758        for queue in &req.queues {
759            if units.len() >= req.capacity as usize
760                || inner.paused.get(queue).copied().unwrap_or(false)
761            {
762                continue;
763            }
764            // tenant fairness draw per partition, never one flat window: group candidates, then a
765            // rotating round-robin across groups. Within a partition: priority DESC,
766            // then scheduled_at, then id.
767            let mut by_part: HashMap<String, Vec<String>> = HashMap::new();
768            for (id, j) in &inner.jobs {
769                if j.env.queue == *queue
770                    && j.state == "available"
771                    && j.env.scheduled_at_ms <= now
772                    && (j.env.sticky_worker.is_empty() || j.env.sticky_worker == req.worker)
773                {
774                    by_part
775                        .entry(j.env.partition_key.clone())
776                        .or_default()
777                        .push(id.clone());
778                }
779            }
780            let mut parts: Vec<String> = by_part.keys().cloned().collect();
781            parts.sort();
782            if parts.is_empty() {
783                continue;
784            }
785            for ids in by_part.values_mut() {
786                ids.sort_by(|a, b| {
787                    let (x, y) = (&inner.jobs[a].env, &inner.jobs[b].env);
788                    y.priority
789                        .cmp(&x.priority)
790                        .then(x.scheduled_at_ms.cmp(&y.scheduled_at_ms))
791                        .then(a.cmp(b))
792                });
793            }
794            let start = {
795                let r = inner.rr.entry(queue.clone()).or_insert(0);
796                let s = *r % parts.len();
797                *r += 1;
798                s
799            };
800            loop {
801                let mut progressed = false;
802                for i in 0..parts.len() {
803                    if units.len() >= req.capacity as usize {
804                        break;
805                    }
806                    let p = &parts[(start + i) % parts.len()];
807                    let Some(ids) = by_part.get_mut(p) else {
808                        continue;
809                    };
810                    let mut picked = None;
811                    while let Some(id) = ids.first().cloned() {
812                        ids.remove(0);
813                        if admissible(&mut inner, &id, &taken, now) {
814                            picked = Some(id);
815                            break;
816                        }
817                    }
818                    let Some(id) = picked else { continue };
819                    progressed = true;
820                    let expires = now + req.lease.as_millis() as i64;
821                    let (rate_class, cost) = {
822                        let e = &inner.jobs[&id].env;
823                        (
824                            e.rate_class.clone(),
825                            headgate_core::effective_weight(e.weight) as i64,
826                        )
827                    };
828                    let charged = !rate_class.is_empty() && inner.rate.contains_key(&rate_class);
829                    let j = inner.jobs.get_mut(&id).unwrap();
830                    j.fence += 1;
831                    j.state = "running".into();
832                    j.lease_id = req.lease_id.clone();
833                    j.lease_expires = expires;
834                    j.rate_charge = if charged { cost } else { 0 };
835                    if charged {
836                        *taken.entry(rate_class).or_insert(0) += cost;
837                    }
838                    units.push(AdmissionUnit {
839                        claims: vec![Claim {
840                            envelope: j.env.clone(),
841                            lease_id: req.lease_id.clone(),
842                            fence: j.fence,
843                            expires_at_ms: expires,
844                            checkpoint: j.checkpoint.clone(),
845                        }],
846                    });
847                }
848                if !progressed || units.len() >= req.capacity as usize {
849                    break;
850                }
851            }
852        }
853        // Spend the tokens actually consumed.
854        for (rc, n) in taken {
855            if let Some(b) = inner.rate.get_mut(&rc) {
856                b.tokens -= n;
857            }
858        }
859        Ok(units)
860    }
861
862    async fn ack_attempt_with_actual_weight(
863        &self,
864        lease: &LeaseRef,
865        outcome: Outcome,
866        err: Option<&str>,
867        delay_ms: Option<i64>,
868        logs: &[String],
869        actual_weight: Option<u32>,
870    ) -> Result<(), StoreError> {
871        let now = self.now();
872        let (base, cap, limit) = (self.retry_base_ms, self.retry_cap_ms, self.crash_limit);
873        let _ = limit;
874        let mut inner = self.inner.lock().unwrap();
875        identity(&inner, lease)?;
876        let id = lease.job_id.clone();
877        if let Some(actual) = actual_weight {
878            let (rc, charge) = {
879                let j = &inner.jobs[&id];
880                (j.env.rate_class.clone(), j.rate_charge)
881            };
882            if charge > 0 {
883                if let Some(b) = inner.rate.get_mut(&rc) {
884                    let gained = if b.limit > 0 && b.window > 0 {
885                        (now - b.refilled).max(0) * b.limit / b.window
886                    } else {
887                        0
888                    };
889                    let avail = b.burst.min(b.tokens + gained);
890                    b.tokens = b.burst.min(avail + charge - actual as i64);
891                    b.refilled = now;
892                }
893            }
894            inner.jobs.get_mut(&id).unwrap().rate_charge = 0;
895        }
896        // attempt-log contract per-attempt logs, rendered into the same history `errors()` returns.
897        let logline = if logs.is_empty() {
898            None
899        } else {
900            Some(format!("logs: {}", logs.join(" | ")))
901        };
902        match outcome {
903            Outcome::Success => {
904                release_unique(&mut inner, &id);
905                let j = inner.jobs.get_mut(&id).unwrap();
906                if j.env.retention_ms == 0 {
907                    inner.jobs.remove(&id); // retention policy ephemeral: delete, not keep
908                } else {
909                    drop_lease(j);
910                    j.state = "completed".into();
911                    j.finalized_at = now;
912                    if let Some(l) = &logline {
913                        j.errs.push(format!("success {l}"));
914                    }
915                }
916            }
917            Outcome::Retry => {
918                let j = inner.jobs.get_mut(&id).unwrap();
919                j.env.attempt += 1;
920                drop_lease(j);
921                j.errs.push(format!(
922                    "retry (attempt {}): {}",
923                    j.env.attempt,
924                    err.unwrap_or("")
925                ));
926                if let Some(l) = &logline {
927                    j.errs.push(l.clone());
928                }
929                if j.env.attempt < j.env.max_attempts {
930                    let backoff = match delay_ms {
931                        Some(d) if d > 0 => d,
932                        _ => default_backoff(j.env.attempt as i64, base, cap),
933                    };
934                    j.state = "retryable".into();
935                    j.env.scheduled_at_ms = now + backoff;
936                } else {
937                    j.state = "archived".into();
938                    j.finalized_at = now;
939                    release_unique(&mut inner, &id);
940                }
941            }
942            Outcome::Skip | Outcome::Undecodable => {
943                let state = if outcome == Outcome::Skip {
944                    "archived"
945                } else {
946                    "undecodable"
947                };
948                let j = inner.jobs.get_mut(&id).unwrap();
949                drop_lease(j);
950                j.state = state.into();
951                j.finalized_at = now;
952                if let Some(e) = err {
953                    j.errs.push(format!("{state}: {e}"));
954                }
955                if let Some(l) = &logline {
956                    j.errs.push(l.clone());
957                }
958                release_unique(&mut inner, &id);
959            }
960            Outcome::Revoke => {
961                release_unique(&mut inner, &id);
962                inner.jobs.remove(&id); // transition table: revoke -> deleted
963            }
964            Outcome::Snooze => {
965                let delay = delay_ms.unwrap_or(0);
966                if delay <= 0 {
967                    return Err(StoreError::Invalid("snooze requires delay_ms > 0".into()));
968                }
969                let j = inner.jobs.get_mut(&id).unwrap();
970                drop_lease(j);
971                j.state = "scheduled".into(); // surveyed policy behavior no attempt consumed
972                j.env.scheduled_at_ms = now + delay;
973            }
974            Outcome::RateLimited => {
975                // surveyed policy behavior NOT a failure: back to available, neither counter moves.
976                let j = inner.jobs.get_mut(&id).unwrap();
977                drop_lease(j);
978                j.state = "available".into();
979                if j.env.scheduled_at_ms > now {
980                    j.env.scheduled_at_ms = now;
981                }
982            }
983            Outcome::LeaseLost => {
984                return Err(StoreError::Invalid(
985                    "lease_lost is applied by the reclaimer, not acked".into(),
986                ));
987            }
988        }
989        Ok(())
990    }
991
992    async fn renew(&self, leases: &[LeaseRef], lease: Duration) -> Result<Vec<String>, StoreError> {
993        if lease.is_zero() {
994            return Err(StoreError::Invalid("lease must be >= 1ms".into()));
995        }
996        let now = self.now();
997        let mut inner = self.inner.lock().unwrap();
998        let mut lost = Vec::new();
999        for l in leases {
1000            match inner.jobs.get_mut(&l.job_id) {
1001                Some(j)
1002                    if j.state == "running" && j.lease_id == l.lease_id && j.fence == l.fence =>
1003                {
1004                    j.lease_expires = now + lease.as_millis() as i64;
1005                }
1006                _ => lost.push(l.job_id.clone()),
1007            }
1008        }
1009        Ok(lost)
1010    }
1011
1012    async fn checkpoint(&self, lease: &LeaseRef, cp: &Checkpoint) -> Result<(), StoreError> {
1013        let mut inner = self.inner.lock().unwrap();
1014        identity(&inner, lease)?;
1015        inner.jobs.get_mut(&lease.job_id).unwrap().checkpoint = cp.clone();
1016        Ok(())
1017    }
1018
1019    async fn reclaim_expired(&self, limit: i64) -> Result<Vec<Reclaimed>, StoreError> {
1020        let now = self.now();
1021        let (crash_limit, base, cap) = (self.crash_limit, self.retry_base_ms, self.retry_cap_ms);
1022        let mut inner = self.inner.lock().unwrap();
1023        let mut ids: Vec<String> = inner.jobs.keys().cloned().collect();
1024        ids.sort(); // deterministic sweep order (map iteration is not)
1025        let mut out = Vec::new();
1026        for id in ids {
1027            if out.len() as i64 >= limit {
1028                break;
1029            }
1030            {
1031                let j = inner.jobs.get(&id).unwrap();
1032                if j.state != "running" || j.lease_expires > now {
1033                    continue;
1034                }
1035            }
1036            let quarantined;
1037            let (fp, ca);
1038            {
1039                let j = inner.jobs.get_mut(&id).unwrap();
1040                j.env.crash_attempt += 1;
1041                drop_lease(j);
1042                j.errs
1043                    .push(format!("lease_lost (crash {})", j.env.crash_attempt));
1044                // crash quarantine step attribution: the checkpoint was durable BEFORE the
1045                // in-progress step's side effects; the crash lands on that step.
1046                if let Some(s) = j.checkpoint.in_progress_step.clone() {
1047                    match j
1048                        .checkpoint
1049                        .crashes_by_step
1050                        .iter_mut()
1051                        .find(|(k, _)| *k == s)
1052                    {
1053                        Some((_, n)) => *n += 1,
1054                        None => j.checkpoint.crashes_by_step.push((s, 1)),
1055                    }
1056                }
1057                quarantined = j.env.crash_attempt >= crash_limit;
1058                fp = j.env.fingerprint.clone();
1059                ca = j.env.crash_attempt;
1060                if quarantined {
1061                    j.state = "quarantined".into();
1062                    j.finalized_at = now;
1063                } else {
1064                    j.state = "retryable".into();
1065                    j.env.scheduled_at_ms = now + default_backoff(ca as i64, base, cap);
1066                }
1067            }
1068            if quarantined {
1069                release_unique(&mut inner, &id);
1070                if !fp.is_empty() {
1071                    inner.quarantine.insert(fp.clone(), true);
1072                }
1073            }
1074            out.push(Reclaimed {
1075                job_id: id,
1076                fingerprint: fp,
1077                crash_attempt: ca,
1078                quarantined,
1079            });
1080        }
1081        Ok(out)
1082    }
1083
1084    async fn promote_due(&self, limit: i64) -> Result<u64, StoreError> {
1085        let now = self.now();
1086        let mut inner = self.inner.lock().unwrap();
1087        let mut n = 0u64;
1088        for j in inner.jobs.values_mut() {
1089            if n as i64 >= limit {
1090                break;
1091            }
1092            if (j.state == "scheduled" || j.state == "retryable") && j.env.scheduled_at_ms <= now {
1093                j.state = "available".into();
1094                n += 1;
1095            }
1096        }
1097        Ok(n)
1098    }
1099
1100    async fn evict_retained(&self, limit: i64) -> Result<u64, StoreError> {
1101        let now = self.now();
1102        let mut inner = self.inner.lock().unwrap();
1103        let lapsed: Vec<String> = inner
1104            .jobs
1105            .iter()
1106            .filter(|(_, j)| {
1107                matches!(
1108                    j.state.as_str(),
1109                    "completed" | "archived" | "cancelled" | "undecodable"
1110                ) && j.env.retention_ms > 0
1111                    && j.finalized_at + j.env.retention_ms <= now
1112            })
1113            .take(limit.max(0) as usize)
1114            .map(|(id, _)| id.clone())
1115            .collect();
1116        // quarantined exempt by design (retention and eviction contract): it parks visibly until an operator acts.
1117        for id in &lapsed {
1118            inner.jobs.remove(id);
1119        }
1120        Ok(lapsed.len() as u64)
1121    }
1122
1123    async fn claim_duty(
1124        &self,
1125        name: &str,
1126        holder: &str,
1127        lease: Duration,
1128    ) -> Result<bool, StoreError> {
1129        if lease.is_zero() {
1130            return Err(StoreError::Invalid("duty lease must be >= 1ms".into()));
1131        }
1132        let now = self.now();
1133        let mut inner = self.inner.lock().unwrap();
1134        if let Some((h, expires)) = inner.duties.get(name) {
1135            if *expires > now && h != holder {
1136                return Ok(false);
1137            }
1138        }
1139        inner
1140            .duties
1141            .insert(name.into(), (holder.into(), now + lease.as_millis() as i64));
1142        Ok(true)
1143    }
1144
1145    async fn release_duty(&self, name: &str, holder: &str) -> Result<(), StoreError> {
1146        let mut inner = self.inner.lock().unwrap();
1147        if inner
1148            .duties
1149            .get(name)
1150            .map(|(h, _)| h == holder)
1151            .unwrap_or(false)
1152        {
1153            inner.duties.remove(name);
1154        }
1155        Ok(())
1156    }
1157
1158    fn caps(&self) -> Caps {
1159        // runtime capability boundary capability honesty: no Transactional (once/step_once error), no Inspect
1160        // (those duties idle), no Notifying (workers poll). See the crate docs.
1161        Caps(0)
1162    }
1163}
1164
1165impl MemStore {
1166    /// Test-only, call-scoped uniqueness bypass. IDs remain strict and no mutable flag
1167    /// can leak into another parallel test.
1168    pub async fn enqueue_without_uniqueness(&self, batch: &[Envelope]) -> Result<(), StoreError> {
1169        let mut cloned = batch.to_vec();
1170        for e in &mut cloned {
1171            e.unique_key = None;
1172            e.unique_window_ms = 0;
1173            e.unique_replace = 0;
1174            e.unique_debounce_ms = 0;
1175        }
1176        self.enqueue(&cloned).await
1177    }
1178}
1179
1180#[async_trait::async_trait]
1181impl headgate_core::ResultStore for MemStore {
1182    async fn ack_success_with_result(
1183        &self,
1184        lease: &LeaseRef,
1185        logs: &[String],
1186        actual_weight: Option<u32>,
1187        result: &headgate_core::JobResult,
1188    ) -> Result<(), StoreError> {
1189        if result.schema_version == 0 {
1190            return Err(StoreError::Invalid(
1191                "result schema_version must be greater than zero".into(),
1192            ));
1193        }
1194        if result.schema_version > headgate_core::MAX_OPAQUE_SCHEMA_VERSION {
1195            return Err(StoreError::Invalid(
1196                "result schema_version exceeds the portable signed-integer limit".into(),
1197            ));
1198        }
1199        if result.bytes.len() > 32 * 1024 * 1024 {
1200            return Err(StoreError::Invalid(
1201                "result bytes exceed the 32 MiB limit".into(),
1202            ));
1203        }
1204        let now = self.now();
1205        let mut inner = self.inner.lock().unwrap();
1206        identity(&inner, lease)?;
1207        let id = lease.job_id.clone();
1208        if let Some(actual) = actual_weight {
1209            let (rc, charge) = {
1210                let job = &inner.jobs[&id];
1211                (job.env.rate_class.clone(), job.rate_charge)
1212            };
1213            if charge > 0 {
1214                if let Some(bucket) = inner.rate.get_mut(&rc) {
1215                    let gained = if bucket.limit > 0 && bucket.window > 0 {
1216                        (now - bucket.refilled).max(0) * bucket.limit / bucket.window
1217                    } else {
1218                        0
1219                    };
1220                    let available = bucket.burst.min(bucket.tokens + gained);
1221                    bucket.tokens = bucket.burst.min(available + charge - actual as i64);
1222                    bucket.refilled = now;
1223                }
1224            }
1225            inner.jobs.get_mut(&id).unwrap().rate_charge = 0;
1226        }
1227        release_unique(&mut inner, &id);
1228        let job = inner.jobs.get_mut(&id).unwrap();
1229        if job.env.retention_ms == 0 {
1230            inner.jobs.remove(&id);
1231        } else {
1232            drop_lease(job);
1233            job.state = "completed".into();
1234            job.finalized_at = now;
1235            job.result = Some(result.clone());
1236            if !logs.is_empty() {
1237                job.errs.push(format!("success logs: {}", logs.join(" | ")));
1238            }
1239        }
1240        Ok(())
1241    }
1242}
1243
1244#[async_trait::async_trait]
1245impl headgate_core::OutputStore for MemStore {
1246    async fn write_job_output(
1247        &self,
1248        lease: &LeaseRef,
1249        output: &headgate_core::JobResult,
1250    ) -> Result<headgate_core::JobOutput, StoreError> {
1251        if output.schema_version == 0 {
1252            return Err(StoreError::Invalid(
1253                "output schema_version must be greater than zero".into(),
1254            ));
1255        }
1256        if output.schema_version > headgate_core::MAX_OPAQUE_SCHEMA_VERSION {
1257            return Err(StoreError::Invalid(
1258                "output schema_version exceeds the portable signed-integer limit".into(),
1259            ));
1260        }
1261        if output.bytes.len() > 32 * 1024 * 1024 {
1262            return Err(StoreError::Invalid(
1263                "output bytes exceed the 32 MiB limit".into(),
1264            ));
1265        }
1266        let now = self.now();
1267        let mut inner = self.inner.lock().unwrap();
1268        identity(&inner, lease)?;
1269        let persisted = headgate_core::JobOutput {
1270            schema_version: output.schema_version,
1271            bytes: output.bytes.clone(),
1272            fence: lease.fence,
1273            updated_at_ms: now,
1274        };
1275        inner.jobs.get_mut(&lease.job_id).unwrap().output = Some(persisted.clone());
1276        Ok(persisted)
1277    }
1278}
1279
1280#[async_trait::async_trait]
1281impl headgate_core::ProgressStore for MemStore {
1282    async fn write_job_progress(
1283        &self,
1284        lease: &LeaseRef,
1285        update: &headgate_core::ProgressUpdate,
1286    ) -> Result<headgate_core::JobProgress, StoreError> {
1287        headgate_core::validate_progress(update)?;
1288        let now = self.now();
1289        let mut inner = self.inner.lock().unwrap();
1290        identity(&inner, lease)?;
1291        let persisted = headgate_core::JobProgress {
1292            current: update.current,
1293            total: update.total,
1294            message: update.message.clone(),
1295            fence: lease.fence,
1296            updated_at_ms: now,
1297        };
1298        inner.jobs.get_mut(&lease.job_id).unwrap().progress = Some(persisted.clone());
1299        Ok(persisted)
1300    }
1301}
1302
1303#[async_trait::async_trait]
1304impl headgate_core::ResultInspect for MemStore {
1305    async fn get_job_result(
1306        &self,
1307        id: &str,
1308    ) -> Result<Option<headgate_core::JobResult>, StoreError> {
1309        Ok(self
1310            .inner
1311            .lock()
1312            .unwrap()
1313            .jobs
1314            .get(id)
1315            .and_then(|job| job.result.clone()))
1316    }
1317}
1318
1319#[async_trait::async_trait]
1320impl headgate_core::OutputInspect for MemStore {
1321    async fn get_job_output(
1322        &self,
1323        id: &str,
1324    ) -> Result<Option<headgate_core::JobOutput>, StoreError> {
1325        Ok(self
1326            .inner
1327            .lock()
1328            .unwrap()
1329            .jobs
1330            .get(id)
1331            .and_then(|job| job.output.clone()))
1332    }
1333}
1334
1335#[async_trait::async_trait]
1336impl headgate_core::ProgressInspect for MemStore {
1337    async fn get_job_progress(
1338        &self,
1339        id: &str,
1340    ) -> Result<Option<headgate_core::JobProgress>, StoreError> {
1341        Ok(self
1342            .inner
1343            .lock()
1344            .unwrap()
1345            .jobs
1346            .get(id)
1347            .and_then(|job| job.progress.clone()))
1348    }
1349}
1350
1351fn drop_lease(j: &mut MemJob) {
1352    j.lease_id.clear();
1353    j.lease_expires = 0;
1354}
1355
1356fn identity(inner: &Inner, lease: &LeaseRef) -> Result<(), StoreError> {
1357    match inner.jobs.get(&lease.job_id) {
1358        Some(j)
1359            if j.state == "running" && j.lease_id == lease.lease_id && j.fence == lease.fence =>
1360        {
1361            Ok(())
1362        }
1363        _ => Err(StoreError::LeaseRejected {
1364            job_id: lease.job_id.clone(),
1365        }),
1366    }
1367}
1368
1369/// The gate's clause order, minus what this store honestly does not model: quarantine,
1370/// then the fleet rate limit (lazy refill, same math as the real buckets).
1371fn admissible(inner: &mut Inner, id: &str, taken: &HashMap<String, i64>, now: i64) -> bool {
1372    let (fp, rc, cost) = {
1373        let j = &inner.jobs[id];
1374        (
1375            j.env.fingerprint.clone(),
1376            j.env.rate_class.clone(),
1377            headgate_core::effective_weight(j.env.weight) as i64,
1378        )
1379    };
1380    if !fp.is_empty() && inner.quarantine.contains_key(&fp) {
1381        return false;
1382    }
1383    if rc.is_empty() {
1384        return true;
1385    }
1386    let Some(b) = inner.rate.get_mut(&rc) else {
1387        return true; // unconfigured class is unlimited HERE (see crate docs)
1388    };
1389    if b.limit > 0 && b.window > 0 {
1390        let gained = (now - b.refilled) * b.limit / b.window;
1391        if gained > 0 {
1392            b.tokens = b.burst.min(b.tokens + gained);
1393            b.refilled = now;
1394        }
1395    }
1396    taken.get(&rc).copied().unwrap_or(0) + cost <= b.tokens
1397}