1use 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
29pub 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 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
132pub 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 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 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 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 pub fn freeze_clock_at(&self, ms: i64) {
346 *self.clock.lock().unwrap() = Clock::Frozen(ms);
347 }
348
349 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 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 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 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 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#[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 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
494pub trait EnqueuedJobs {
497 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
512pub 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
556pub 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 headgate_core::validate_enqueue(batch)?;
596 let mut inner = self.inner.lock().unwrap();
597 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 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 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 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 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); } 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); }
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(); j.env.scheduled_at_ms = now + delay;
973 }
974 Outcome::RateLimited => {
975 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(); 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 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 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 Caps(0)
1162 }
1163}
1164
1165impl MemStore {
1166 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
1369fn 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; };
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}