1use std::collections::{BTreeMap, BTreeSet};
2use std::sync::{Arc, Mutex, RwLock};
3use std::time::Duration as StdDuration;
4
5use serde::{Deserialize, Serialize};
6use uuid::Uuid;
7
8use crate::event_log::{
9 sanitize_topic_component, AnyEventLog, EventLog, LogError, LogEvent, Topic,
10};
11
12use super::scheduler::{self, SchedulableJob, SchedulerPolicy, SchedulerSnapshot, SchedulerState};
13use super::{DispatchOutcome, TriggerEvent};
14
15mod scheduling;
16
17pub use scheduling::{
18 WorkerQueuePriority, WorkerQueueSchedulingDecision, WorkerQueueSchedulingReceipt,
19 DEFERRABLE_PROMOTION_AGE_MS,
20};
21
22pub const WORKER_QUEUE_CATALOG_TOPIC: &str = "worker.queues";
23const WORKER_QUEUE_CLAIMS_SUFFIX: &str = ".claims";
24const WORKER_QUEUE_RESPONSES_SUFFIX: &str = ".responses";
25
26#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
27pub struct WorkerQueueJob {
28 pub queue: String,
29 pub trigger_id: String,
30 pub binding_key: String,
31 pub binding_version: u32,
32 pub event: TriggerEvent,
33 #[serde(default)]
34 pub replay_of_event_id: Option<String>,
35 #[serde(default)]
36 pub priority: WorkerQueuePriority,
37}
38
39#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
40pub struct WorkerQueueEnqueueReceipt {
41 pub queue: String,
42 pub job_event_id: u64,
43 pub response_topic: String,
44}
45
46#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
47pub struct WorkerQueueClaimHandle {
48 pub queue: String,
49 pub job_event_id: u64,
50 pub claim_id: String,
51 pub consumer_id: String,
52 pub expires_at_ms: i64,
53}
54
55#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
56pub struct ClaimedWorkerJob {
57 pub handle: WorkerQueueClaimHandle,
58 pub job: WorkerQueueJob,
59 pub scheduling: WorkerQueueSchedulingReceipt,
60}
61
62#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
63pub struct WorkerQueueResponseRecord {
64 pub queue: String,
65 pub job_event_id: u64,
66 pub consumer_id: String,
67 pub handled_at_ms: i64,
68 pub outcome: Option<DispatchOutcome>,
69 pub error: Option<String>,
70}
71
72#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
73pub struct WorkerQueueSummary {
74 pub queue: String,
75 pub ready: usize,
76 pub in_flight: usize,
77 pub acked: usize,
78 pub purged: usize,
79 pub responses: usize,
80 pub oldest_unclaimed_age_ms: Option<u64>,
81}
82
83#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
84pub struct WorkerQueueJobState {
85 pub job_event_id: u64,
86 pub enqueued_at_ms: i64,
87 pub job: WorkerQueueJob,
88 pub active_claim: Option<WorkerQueueClaimHandle>,
89 pub acked: bool,
90 pub purged: bool,
91}
92
93impl WorkerQueueJobState {
94 pub fn is_ready(&self) -> bool {
95 !self.acked && !self.purged && self.active_claim.is_none()
96 }
97}
98
99#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
100pub struct WorkerQueueState {
101 pub queue: String,
102 pub responses: Vec<WorkerQueueResponseRecord>,
103 pub jobs: Vec<WorkerQueueJobState>,
104}
105
106impl WorkerQueueState {
107 pub fn summary(&self, now_ms: i64) -> WorkerQueueSummary {
108 let ready = self.jobs.iter().filter(|job| job.is_ready()).count();
109 let in_flight = self
110 .jobs
111 .iter()
112 .filter(|job| !job.acked && !job.purged && job.active_claim.is_some())
113 .count();
114 let acked = self.jobs.iter().filter(|job| job.acked).count();
115 let purged = self.jobs.iter().filter(|job| job.purged).count();
116 let oldest_unclaimed_age_ms = self
117 .jobs
118 .iter()
119 .filter(|job| job.is_ready())
120 .map(|job| now_ms.saturating_sub(job.enqueued_at_ms).max(0) as u64)
121 .max();
122 WorkerQueueSummary {
123 queue: self.queue.clone(),
124 ready,
125 in_flight,
126 acked,
127 purged,
128 responses: self.responses.len(),
129 oldest_unclaimed_age_ms,
130 }
131 }
132
133 fn next_ready_job_with_scheduler(
141 &self,
142 scheduler_state: &mut SchedulerState,
143 policy: &SchedulerPolicy,
144 now_ms: i64,
145 ) -> Option<(&WorkerQueueJobState, scheduler::SchedulerSelection)> {
146 let candidates: Vec<&WorkerQueueJobState> =
147 self.jobs.iter().filter(|job| job.is_ready()).collect();
148 if candidates.is_empty() {
149 return None;
150 }
151 let views: Vec<SchedulableJob<'_>> = candidates
152 .iter()
153 .map(|state| SchedulableJob::from_state(state))
154 .collect();
155
156 let in_flight = scheduler::in_flight_by_key(&self.jobs, policy);
158 scheduler_state.replace_in_flight(in_flight);
159
160 let pick = scheduler_state.select(&views, policy, now_ms)?;
161 candidates
162 .into_iter()
163 .find(|job| job.job_event_id == pick.job_event_id)
164 .map(|job| (job, pick))
165 }
166
167 fn active_claim_for(&self, job_event_id: u64) -> Option<&WorkerQueueClaimHandle> {
168 self.jobs
169 .iter()
170 .find(|job| job.job_event_id == job_event_id)
171 .and_then(|job| job.active_claim.as_ref())
172 }
173}
174
175#[derive(Clone)]
176pub struct WorkerQueue {
177 event_log: Arc<AnyEventLog>,
178 policy: Arc<RwLock<SchedulerPolicy>>,
181 scheduler_states: Arc<Mutex<BTreeMap<String, SchedulerState>>>,
185}
186
187#[derive(Clone, Debug, Serialize)]
188pub struct WorkerQueueInspectSnapshot {
189 pub summary: WorkerQueueSummary,
190 pub scheduler: SchedulerSnapshot,
191}
192
193impl WorkerQueue {
194 pub fn new(event_log: Arc<AnyEventLog>) -> Self {
199 Self::with_policy(event_log, SchedulerPolicy::from_env())
200 }
201
202 pub fn with_policy(event_log: Arc<AnyEventLog>, policy: SchedulerPolicy) -> Self {
203 Self {
204 event_log,
205 policy: Arc::new(RwLock::new(policy)),
206 scheduler_states: Arc::new(Mutex::new(BTreeMap::new())),
207 }
208 }
209
210 pub fn set_policy(&self, policy: SchedulerPolicy) {
213 *self.policy.write().expect("scheduler policy poisoned") = policy;
214 }
215
216 pub fn policy(&self) -> SchedulerPolicy {
217 self.policy
218 .read()
219 .expect("scheduler policy poisoned")
220 .clone()
221 }
222
223 pub async fn enqueue(
224 &self,
225 job: &WorkerQueueJob,
226 ) -> Result<WorkerQueueEnqueueReceipt, LogError> {
227 let queue = job.queue.trim();
228 if queue.is_empty() {
229 return Err(LogError::Config(
230 "worker queue name cannot be empty".to_string(),
231 ));
232 }
233 let queue_name = queue.to_string();
234 let catalog_topic = Topic::new(WORKER_QUEUE_CATALOG_TOPIC)
235 .expect("static worker queue catalog topic should always be valid");
236 self.event_log
237 .append(
238 &catalog_topic,
239 LogEvent::new(
240 "queue_seen",
241 serde_json::to_value(WorkerQueueCatalogRecord {
242 queue: queue_name.clone(),
243 })
244 .map_err(|error| LogError::Serde(error.to_string()))?,
245 ),
246 )
247 .await?;
248
249 let job_topic = job_topic(&queue_name)?;
250 let mut headers = BTreeMap::new();
251 headers.insert("queue".to_string(), queue_name.clone());
252 headers.insert("trigger_id".to_string(), job.trigger_id.clone());
253 headers.insert("binding_key".to_string(), job.binding_key.clone());
254 headers.insert("event_id".to_string(), job.event.id.0.clone());
255 headers.insert("priority".to_string(), job.priority.as_str().to_string());
256 let job_event_id = self
257 .event_log
258 .append(
259 &job_topic,
260 LogEvent::new(
261 "trigger_dispatch",
262 serde_json::to_value(job)
263 .map_err(|error| LogError::Serde(error.to_string()))?,
264 )
265 .with_headers(headers),
266 )
267 .await?;
268 if let Some(metrics) = crate::active_metrics_registry() {
269 if let Ok(state) = self.queue_state(&queue_name).await {
270 let summary = state.summary(now_ms());
271 metrics.set_worker_queue_depth(
272 &queue_name,
273 (summary.ready + summary.in_flight) as u64,
274 );
275 }
276 }
277 Ok(WorkerQueueEnqueueReceipt {
278 queue: queue_name.clone(),
279 job_event_id,
280 response_topic: response_topic_name(&queue_name),
281 })
282 }
283
284 pub async fn known_queues(&self) -> Result<Vec<String>, LogError> {
285 let topic = Topic::new(WORKER_QUEUE_CATALOG_TOPIC)
286 .expect("static worker queue catalog topic should always be valid");
287 let events = self.event_log.read_range(&topic, None, usize::MAX).await?;
288 let mut queues = BTreeSet::new();
289 for (_, event) in events {
290 if event.kind != "queue_seen" {
291 continue;
292 }
293 let record: WorkerQueueCatalogRecord = serde_json::from_value(event.payload)
294 .map_err(|error| LogError::Serde(error.to_string()))?;
295 if !record.queue.trim().is_empty() {
296 queues.insert(record.queue);
297 }
298 }
299 Ok(queues.into_iter().collect())
300 }
301
302 pub async fn queue_state(&self, queue: &str) -> Result<WorkerQueueState, LogError> {
303 let queue_name = queue.trim();
304 if queue_name.is_empty() {
305 return Err(LogError::Config(
306 "worker queue name cannot be empty".to_string(),
307 ));
308 }
309 let now_ms = now_ms();
310 let job_events = self
311 .event_log
312 .read_range(&job_topic(queue_name)?, None, usize::MAX)
313 .await?;
314 let claim_events = self
315 .event_log
316 .read_range(&claims_topic(queue_name)?, None, usize::MAX)
317 .await?;
318 let response_events = self
319 .event_log
320 .read_range(&responses_topic(queue_name)?, None, usize::MAX)
321 .await?;
322
323 let mut jobs = BTreeMap::<u64, WorkerQueueJobStateInternal>::new();
324 for (job_event_id, event) in job_events {
325 if event.kind != "trigger_dispatch" {
326 continue;
327 }
328 let job: WorkerQueueJob = serde_json::from_value(event.payload)
329 .map_err(|error| LogError::Serde(error.to_string()))?;
330 jobs.insert(
331 job_event_id,
332 WorkerQueueJobStateInternal {
333 job_event_id,
334 enqueued_at_ms: event.occurred_at_ms,
335 job,
336 active_claim: None,
337 acked: false,
338 purged: false,
339 seen_claim_ids: BTreeSet::new(),
340 },
341 );
342 }
343
344 for (_, event) in claim_events {
345 match event.kind.as_str() {
346 "job_claimed" => {
347 let claim: WorkerQueueClaimRecord = serde_json::from_value(event.payload)
348 .map_err(|error| LogError::Serde(error.to_string()))?;
349 let Some(job) = jobs.get_mut(&claim.job_event_id) else {
350 continue;
351 };
352 if job.acked || job.purged {
353 continue;
354 }
355 job.seen_claim_ids.insert(claim.claim_id.clone());
356 let can_take = job
357 .active_claim
358 .as_ref()
359 .is_none_or(|active| active.expires_at_ms <= claim.claimed_at_ms);
360 if can_take {
361 job.active_claim = Some(WorkerQueueClaimHandle {
362 queue: queue_name.to_string(),
363 job_event_id: claim.job_event_id,
364 claim_id: claim.claim_id,
365 consumer_id: claim.consumer_id,
366 expires_at_ms: claim.expires_at_ms,
367 });
368 }
369 }
370 "claim_renewed" => {
371 let renewal: WorkerQueueClaimRenewalRecord =
372 serde_json::from_value(event.payload)
373 .map_err(|error| LogError::Serde(error.to_string()))?;
374 let Some(job) = jobs.get_mut(&renewal.job_event_id) else {
375 continue;
376 };
377 if let Some(active) = job.active_claim.as_mut() {
378 if active.claim_id == renewal.claim_id {
379 active.expires_at_ms = renewal.expires_at_ms;
380 }
381 }
382 }
383 "job_released" => {
384 let release: WorkerQueueReleaseRecord =
385 serde_json::from_value(event.payload)
386 .map_err(|error| LogError::Serde(error.to_string()))?;
387 let Some(job) = jobs.get_mut(&release.job_event_id) else {
388 continue;
389 };
390 if job
391 .active_claim
392 .as_ref()
393 .is_some_and(|active| active.claim_id == release.claim_id)
394 {
395 job.active_claim = None;
396 }
397 }
398 "job_acked" => {
399 let ack: WorkerQueueAckRecord = serde_json::from_value(event.payload)
400 .map_err(|error| LogError::Serde(error.to_string()))?;
401 let Some(job) = jobs.get_mut(&ack.job_event_id) else {
402 continue;
403 };
404 if ack.claim_id.is_empty() || job.seen_claim_ids.contains(&ack.claim_id) {
405 job.acked = true;
406 job.active_claim = None;
407 }
408 }
409 "job_purged" => {
410 let purge: WorkerQueuePurgeRecord = serde_json::from_value(event.payload)
411 .map_err(|error| LogError::Serde(error.to_string()))?;
412 let Some(job) = jobs.get_mut(&purge.job_event_id) else {
413 continue;
414 };
415 if !job.acked {
416 job.purged = true;
417 job.active_claim = None;
418 }
419 }
420 _ => {}
421 }
422 }
423
424 let responses = response_events
425 .into_iter()
426 .filter(|(_, event)| event.kind == "job_response")
427 .map(|(_, event)| {
428 serde_json::from_value::<WorkerQueueResponseRecord>(event.payload)
429 .map_err(|error| LogError::Serde(error.to_string()))
430 })
431 .collect::<Result<Vec<_>, _>>()?;
432
433 let mut queue_state = WorkerQueueState {
434 queue: queue_name.to_string(),
435 responses,
436 jobs: jobs
437 .into_values()
438 .map(|mut job| {
439 if job
440 .active_claim
441 .as_ref()
442 .is_some_and(|active| active.expires_at_ms <= now_ms)
443 {
444 job.active_claim = None;
445 }
446 WorkerQueueJobState {
447 job_event_id: job.job_event_id,
448 enqueued_at_ms: job.enqueued_at_ms,
449 job: job.job,
450 active_claim: job.active_claim,
451 acked: job.acked,
452 purged: job.purged,
453 }
454 })
455 .collect(),
456 };
457 queue_state
458 .jobs
459 .sort_by_key(|job| (job.enqueued_at_ms, job.job_event_id));
460 Ok(queue_state)
461 }
462
463 pub async fn queue_summaries(&self) -> Result<Vec<WorkerQueueSummary>, LogError> {
464 let now_ms = now_ms();
465 let mut summaries = Vec::new();
466 for queue in self.known_queues().await? {
467 let state = self.queue_state(&queue).await?;
468 summaries.push(state.summary(now_ms));
469 }
470 summaries.sort_by(|left, right| left.queue.cmp(&right.queue));
471 Ok(summaries)
472 }
473
474 pub async fn claim_next(
475 &self,
476 queue: &str,
477 consumer_id: &str,
478 ttl: StdDuration,
479 ) -> Result<Option<ClaimedWorkerJob>, LogError> {
480 let queue_name = queue.trim();
481 if queue_name.is_empty() {
482 return Err(LogError::Config(
483 "worker queue name cannot be empty".to_string(),
484 ));
485 }
486 if consumer_id.trim().is_empty() {
487 return Err(LogError::InvalidConsumer(
488 "worker queue consumer id cannot be empty".to_string(),
489 ));
490 }
491 let policy = self.policy();
492 for _ in 0..8 {
493 let now_ms = now_ms();
494 let state = self.queue_state(queue_name).await?;
495 let (job, selection) = {
496 let mut states = self
497 .scheduler_states
498 .lock()
499 .expect("scheduler state poisoned");
500 let scheduler_state = states.entry(queue_name.to_string()).or_default();
501 let Some((job, selection)) =
502 state.next_ready_job_with_scheduler(scheduler_state, &policy, now_ms)
503 else {
504 return Ok(None);
505 };
506 (job.clone(), selection)
507 };
508 let scheduling = WorkerQueueSchedulingReceipt {
509 selected_at_ms: now_ms,
510 enqueued_at_ms: job.enqueued_at_ms,
511 waited_ms: now_ms.saturating_sub(job.enqueued_at_ms).max(0) as u64,
512 priority: job.job.priority,
513 decision: selection.decision,
514 promotion_deadline_at_ms: selection.promotion_deadline_at_ms,
515 fairness_key: selection.fairness_key.clone(),
516 };
517 let claim = WorkerQueueClaimRecord {
518 job_event_id: job.job_event_id,
519 claim_id: Uuid::new_v4().to_string(),
520 consumer_id: consumer_id.to_string(),
521 claimed_at_ms: now_ms,
522 expires_at_ms: expiry_ms(now_ms, ttl),
523 scheduling: Some(scheduling.clone()),
524 };
525 self.event_log
526 .append(
527 &claims_topic(queue_name)?,
528 LogEvent::new(
529 "job_claimed",
530 serde_json::to_value(&claim)
531 .map_err(|error| LogError::Serde(error.to_string()))?,
532 ),
533 )
534 .await?;
535 let refreshed = self.queue_state(queue_name).await?;
536 if refreshed
537 .active_claim_for(job.job_event_id)
538 .is_some_and(|active| active.claim_id == claim.claim_id)
539 {
540 {
541 let mut states = self
542 .scheduler_states
543 .lock()
544 .expect("scheduler state poisoned");
545 let scheduler_state = states.entry(queue_name.to_string()).or_default();
546 scheduler_state.note_claim_committed(&selection.fairness_key);
547 }
548 if let Some(metrics) = crate::active_metrics_registry() {
549 let summary = refreshed.summary(now_ms);
550 metrics.record_worker_queue_claim_age(
551 queue_name,
552 now_ms.saturating_sub(job.enqueued_at_ms) as f64 / 1000.0,
553 );
554 metrics.set_worker_queue_depth(
555 queue_name,
556 (summary.ready + summary.in_flight) as u64,
557 );
558 metrics.record_scheduler_selection(
559 queue_name,
560 policy.fairness_key.as_str(),
561 &selection.fairness_key,
562 );
563 if let Ok(snap) = self.inspect_queue(queue_name).await {
564 for stat in &snap.scheduler.keys {
565 metrics.set_scheduler_deficit(
566 queue_name,
567 policy.fairness_key.as_str(),
568 &stat.fairness_key,
569 stat.deficit,
570 );
571 metrics.set_scheduler_oldest_eligible_age(
572 queue_name,
573 policy.fairness_key.as_str(),
574 &stat.fairness_key,
575 stat.oldest_ready_age_ms,
576 );
577 }
578 }
579 }
580 return Ok(Some(ClaimedWorkerJob {
581 handle: WorkerQueueClaimHandle {
582 queue: queue_name.to_string(),
583 job_event_id: claim.job_event_id,
584 claim_id: claim.claim_id,
585 consumer_id: claim.consumer_id,
586 expires_at_ms: claim.expires_at_ms,
587 },
588 job: job.job,
589 scheduling,
590 }));
591 }
592 }
593 Ok(None)
594 }
595
596 pub async fn inspect_queue(&self, queue: &str) -> Result<WorkerQueueInspectSnapshot, LogError> {
599 let queue_name = queue.trim();
600 if queue_name.is_empty() {
601 return Err(LogError::Config(
602 "worker queue name cannot be empty".to_string(),
603 ));
604 }
605 let now_ms = now_ms();
606 let state = self.queue_state(queue_name).await?;
607 let summary = state.summary(now_ms);
608 let policy = self.policy();
609 let ready = scheduler::ready_stats_by_key(&state.jobs, &policy, now_ms);
610 let in_flight = scheduler::in_flight_by_key(&state.jobs, &policy);
612 let scheduler_snapshot = {
613 let mut states = self
614 .scheduler_states
615 .lock()
616 .expect("scheduler state poisoned");
617 let scheduler_state = states.entry(queue_name.to_string()).or_default();
618 scheduler_state.replace_in_flight(in_flight);
619 scheduler_state.snapshot(&policy, &ready)
620 };
621 Ok(WorkerQueueInspectSnapshot {
622 summary,
623 scheduler: scheduler_snapshot,
624 })
625 }
626
627 pub async fn inspect_all_queues(&self) -> Result<Vec<WorkerQueueInspectSnapshot>, LogError> {
629 let mut snapshots = Vec::new();
630 for queue in self.known_queues().await? {
631 snapshots.push(self.inspect_queue(&queue).await?);
632 }
633 snapshots.sort_by(|left, right| left.summary.queue.cmp(&right.summary.queue));
634 Ok(snapshots)
635 }
636
637 pub async fn renew_claim(
638 &self,
639 handle: &WorkerQueueClaimHandle,
640 ttl: StdDuration,
641 ) -> Result<bool, LogError> {
642 let now_ms = now_ms();
643 let renewal = WorkerQueueClaimRenewalRecord {
644 job_event_id: handle.job_event_id,
645 claim_id: handle.claim_id.clone(),
646 consumer_id: handle.consumer_id.clone(),
647 renewed_at_ms: now_ms,
648 expires_at_ms: expiry_ms(now_ms, ttl),
649 };
650 self.event_log
651 .append(
652 &claims_topic(&handle.queue)?,
653 LogEvent::new(
654 "claim_renewed",
655 serde_json::to_value(&renewal)
656 .map_err(|error| LogError::Serde(error.to_string()))?,
657 ),
658 )
659 .await?;
660 let refreshed = self.queue_state(&handle.queue).await?;
661 Ok(refreshed
662 .active_claim_for(handle.job_event_id)
663 .is_some_and(|active| active.claim_id == handle.claim_id))
664 }
665
666 pub async fn release_claim(
667 &self,
668 handle: &WorkerQueueClaimHandle,
669 reason: &str,
670 ) -> Result<(), LogError> {
671 let release = WorkerQueueReleaseRecord {
672 job_event_id: handle.job_event_id,
673 claim_id: handle.claim_id.clone(),
674 consumer_id: handle.consumer_id.clone(),
675 released_at_ms: now_ms(),
676 reason: if reason.trim().is_empty() {
677 None
678 } else {
679 Some(reason.to_string())
680 },
681 };
682 self.event_log
683 .append(
684 &claims_topic(&handle.queue)?,
685 LogEvent::new(
686 "job_released",
687 serde_json::to_value(&release)
688 .map_err(|error| LogError::Serde(error.to_string()))?,
689 ),
690 )
691 .await?;
692 Ok(())
693 }
694
695 pub async fn append_response(
696 &self,
697 queue: &str,
698 response: &WorkerQueueResponseRecord,
699 ) -> Result<u64, LogError> {
700 self.event_log
701 .append(
702 &responses_topic(queue)?,
703 LogEvent::new(
704 "job_response",
705 serde_json::to_value(response)
706 .map_err(|error| LogError::Serde(error.to_string()))?,
707 ),
708 )
709 .await
710 }
711
712 pub async fn ack_claim(&self, handle: &WorkerQueueClaimHandle) -> Result<u64, LogError> {
713 self.event_log
714 .append(
715 &claims_topic(&handle.queue)?,
716 LogEvent::new(
717 "job_acked",
718 serde_json::to_value(WorkerQueueAckRecord {
719 job_event_id: handle.job_event_id,
720 claim_id: handle.claim_id.clone(),
721 consumer_id: handle.consumer_id.clone(),
722 acked_at_ms: now_ms(),
723 })
724 .map_err(|error| LogError::Serde(error.to_string()))?,
725 ),
726 )
727 .await
728 }
729
730 pub async fn ack_job(
731 &self,
732 queue: &str,
733 job_event_id: u64,
734 consumer_id: &str,
735 ) -> Result<bool, LogError> {
736 let queue_name = queue.trim();
737 if queue_name.is_empty() {
738 return Err(LogError::Config(
739 "worker queue name cannot be empty".to_string(),
740 ));
741 }
742 let state = self.queue_state(queue_name).await?;
743 let Some(job) = state
744 .jobs
745 .iter()
746 .find(|job| job.job_event_id == job_event_id)
747 else {
748 return Ok(false);
749 };
750 if job.acked || job.purged {
751 return Ok(false);
752 }
753 self.event_log
754 .append(
755 &claims_topic(queue_name)?,
756 LogEvent::new(
757 "job_acked",
758 serde_json::to_value(WorkerQueueAckRecord {
759 job_event_id,
760 claim_id: String::new(),
761 consumer_id: consumer_id.to_string(),
762 acked_at_ms: now_ms(),
763 })
764 .map_err(|error| LogError::Serde(error.to_string()))?,
765 ),
766 )
767 .await?;
768 Ok(true)
769 }
770
771 pub async fn purge_unclaimed(
772 &self,
773 queue: &str,
774 purged_by: &str,
775 reason: Option<&str>,
776 ) -> Result<usize, LogError> {
777 let state = self.queue_state(queue).await?;
778 let ready_jobs: Vec<_> = state
779 .jobs
780 .into_iter()
781 .filter(|job| job.is_ready())
782 .map(|job| job.job_event_id)
783 .collect();
784 for job_event_id in &ready_jobs {
785 self.event_log
786 .append(
787 &claims_topic(queue)?,
788 LogEvent::new(
789 "job_purged",
790 serde_json::to_value(WorkerQueuePurgeRecord {
791 job_event_id: *job_event_id,
792 purged_by: purged_by.to_string(),
793 purged_at_ms: now_ms(),
794 reason: reason
795 .filter(|value| !value.trim().is_empty())
796 .map(|value| value.to_string()),
797 })
798 .map_err(|error| LogError::Serde(error.to_string()))?,
799 ),
800 )
801 .await?;
802 }
803 Ok(ready_jobs.len())
804 }
805}
806
807#[derive(Clone, Debug)]
808struct WorkerQueueJobStateInternal {
809 job_event_id: u64,
810 enqueued_at_ms: i64,
811 job: WorkerQueueJob,
812 active_claim: Option<WorkerQueueClaimHandle>,
813 acked: bool,
814 purged: bool,
815 seen_claim_ids: BTreeSet<String>,
816}
817
818#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
819struct WorkerQueueCatalogRecord {
820 queue: String,
821}
822
823#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
824struct WorkerQueueClaimRecord {
825 job_event_id: u64,
826 claim_id: String,
827 consumer_id: String,
828 claimed_at_ms: i64,
829 expires_at_ms: i64,
830 #[serde(default, skip_serializing_if = "Option::is_none")]
831 scheduling: Option<WorkerQueueSchedulingReceipt>,
832}
833
834#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
835struct WorkerQueueClaimRenewalRecord {
836 job_event_id: u64,
837 claim_id: String,
838 consumer_id: String,
839 renewed_at_ms: i64,
840 expires_at_ms: i64,
841}
842
843#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
844struct WorkerQueueReleaseRecord {
845 job_event_id: u64,
846 claim_id: String,
847 consumer_id: String,
848 released_at_ms: i64,
849 #[serde(default)]
850 reason: Option<String>,
851}
852
853#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
854struct WorkerQueueAckRecord {
855 job_event_id: u64,
856 claim_id: String,
857 consumer_id: String,
858 acked_at_ms: i64,
859}
860
861#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
862struct WorkerQueuePurgeRecord {
863 job_event_id: u64,
864 purged_by: String,
865 purged_at_ms: i64,
866 #[serde(default)]
867 reason: Option<String>,
868}
869
870pub fn job_topic_name(queue: &str) -> String {
871 format!("worker.{}", sanitize_topic_component(queue))
872}
873
874pub fn claims_topic_name(queue: &str) -> String {
875 format!("{}{}", job_topic_name(queue), WORKER_QUEUE_CLAIMS_SUFFIX)
876}
877
878pub fn response_topic_name(queue: &str) -> String {
879 format!("{}{}", job_topic_name(queue), WORKER_QUEUE_RESPONSES_SUFFIX)
880}
881
882fn job_topic(queue: &str) -> Result<Topic, LogError> {
883 Topic::new(job_topic_name(queue))
884}
885
886fn claims_topic(queue: &str) -> Result<Topic, LogError> {
887 Topic::new(claims_topic_name(queue))
888}
889
890fn responses_topic(queue: &str) -> Result<Topic, LogError> {
891 Topic::new(response_topic_name(queue))
892}
893
894fn now_ms() -> i64 {
895 std::time::SystemTime::now()
896 .duration_since(std::time::UNIX_EPOCH)
897 .map(|duration| duration.as_millis() as i64)
898 .unwrap_or(0)
899}
900
901fn expiry_ms(now_ms: i64, ttl: StdDuration) -> i64 {
902 now_ms.saturating_add(ttl.as_millis().min(i64::MAX as u128) as i64)
903}
904
905#[cfg(test)]
906mod tests {
907 use super::*;
908
909 use crate::event_log::{AnyEventLog, MemoryEventLog};
910 use crate::triggers::{
911 event::{GenericWebhookPayload, KnownProviderPayload},
912 scheduler::{self, SchedulerStrategy},
913 ProviderId, ProviderPayload, SignatureStatus, TraceId, TriggerEvent,
914 };
915
916 fn test_event(id: &str) -> TriggerEvent {
917 TriggerEvent {
918 id: crate::triggers::TriggerEventId(id.to_string()),
919 provider: ProviderId::from("github"),
920 kind: "issues.opened".to_string(),
921 trace_id: TraceId("trace-test".to_string()),
922 dedupe_key: id.to_string(),
923 tenant_id: None,
924 headers: BTreeMap::new(),
925 batch: None,
926 raw_body: None,
927 provider_payload: ProviderPayload::Known(KnownProviderPayload::Webhook(
928 GenericWebhookPayload {
929 source: Some("worker-queue-test".to_string()),
930 content_type: Some("application/json".to_string()),
931 raw: serde_json::json!({"id": id}),
932 },
933 )),
934 signature_status: SignatureStatus::Verified,
935 received_at: time::OffsetDateTime::now_utc(),
936 occurred_at: None,
937 dedupe_claimed: false,
938 }
939 }
940
941 fn test_job(
942 queue: &str,
943 trigger_id: &str,
944 event_id: &str,
945 priority: WorkerQueuePriority,
946 ) -> WorkerQueueJob {
947 WorkerQueueJob {
948 queue: queue.to_string(),
949 trigger_id: trigger_id.to_string(),
950 binding_key: format!("{trigger_id}@v1"),
951 binding_version: 1,
952 event: test_event(event_id),
953 replay_of_event_id: None,
954 priority,
955 }
956 }
957
958 #[tokio::test(flavor = "current_thread")]
959 async fn enqueue_and_summarize_queue() {
960 let log = Arc::new(AnyEventLog::Memory(MemoryEventLog::new(32)));
961 let queue = WorkerQueue::new(log);
962 queue
963 .enqueue(&test_job(
964 "triage",
965 "incoming-review-task",
966 "evt-1",
967 WorkerQueuePriority::Normal,
968 ))
969 .await
970 .unwrap();
971 let summaries = queue.queue_summaries().await.unwrap();
972 assert_eq!(summaries.len(), 1);
973 assert_eq!(summaries[0].queue, "triage");
974 assert_eq!(summaries[0].ready, 1);
975 assert_eq!(summaries[0].in_flight, 0);
976 }
977
978 #[tokio::test(flavor = "current_thread")]
979 async fn claim_and_ack_remove_job_from_ready_pool() {
980 let log = Arc::new(AnyEventLog::Memory(MemoryEventLog::new(32)));
981 let queue = WorkerQueue::new(log);
982 queue
983 .enqueue(&test_job(
984 "triage",
985 "incoming-review-task",
986 "evt-1",
987 WorkerQueuePriority::Normal,
988 ))
989 .await
990 .unwrap();
991 let claimed = queue
992 .claim_next("triage", "consumer-a", StdDuration::from_mins(1))
993 .await
994 .unwrap()
995 .unwrap();
996 assert_eq!(
997 claimed.scheduling.decision,
998 WorkerQueueSchedulingDecision::Priority
999 );
1000 assert_eq!(claimed.scheduling.priority, WorkerQueuePriority::Normal);
1001 assert_eq!(claimed.scheduling.fairness_key, "_no_tenant");
1002 assert!(
1003 claimed.scheduling.promotion_deadline_at_ms > Some(claimed.scheduling.enqueued_at_ms)
1004 );
1005 let before_ack = queue.queue_state("triage").await.unwrap();
1006 assert_eq!(before_ack.summary(now_ms()).ready, 0);
1007 assert_eq!(before_ack.summary(now_ms()).in_flight, 1);
1008 queue
1009 .append_response(
1010 "triage",
1011 &WorkerQueueResponseRecord {
1012 queue: "triage".to_string(),
1013 job_event_id: claimed.handle.job_event_id,
1014 consumer_id: "consumer-a".to_string(),
1015 handled_at_ms: now_ms(),
1016 outcome: Some(DispatchOutcome {
1017 trigger_id: "incoming-review-task".to_string(),
1018 binding_key: "incoming-review-task@v1".to_string(),
1019 event_id: "evt-1".to_string(),
1020 attempt_count: 1,
1021 status: super::super::DispatchStatus::Succeeded,
1022 handler_kind: "local".to_string(),
1023 target_uri: "handlers::on_review".to_string(),
1024 replay_of_event_id: None,
1025 result: Some(serde_json::json!({"ok": true})),
1026 error: None,
1027 }),
1028 error: None,
1029 },
1030 )
1031 .await
1032 .unwrap();
1033 queue.ack_claim(&claimed.handle).await.unwrap();
1034 let after_ack = queue.queue_state("triage").await.unwrap();
1035 let summary = after_ack.summary(now_ms());
1036 assert_eq!(summary.ready, 0);
1037 assert_eq!(summary.in_flight, 0);
1038 assert_eq!(summary.acked, 1);
1039 assert_eq!(summary.responses, 1);
1040 }
1041
1042 #[tokio::test(flavor = "current_thread")]
1043 async fn ack_job_acknowledges_without_active_claim() {
1044 let log = Arc::new(AnyEventLog::Memory(MemoryEventLog::new(32)));
1045 let queue = WorkerQueue::new(log);
1046 let receipt = queue
1047 .enqueue(&test_job(
1048 "triage",
1049 "incoming-review-task",
1050 "evt-1",
1051 WorkerQueuePriority::Normal,
1052 ))
1053 .await
1054 .unwrap();
1055
1056 assert!(queue
1057 .ack_job("triage", receipt.job_event_id, "pipeline_lifecycle")
1058 .await
1059 .unwrap());
1060 let state = queue.queue_state("triage").await.unwrap();
1061 let summary = state.summary(now_ms());
1062 assert_eq!(summary.ready, 0);
1063 assert_eq!(summary.acked, 1);
1064 assert!(
1065 !queue
1066 .ack_job("triage", receipt.job_event_id, "pipeline_lifecycle")
1067 .await
1068 .unwrap(),
1069 "already acknowledged jobs should not produce a second settlement"
1070 );
1071 }
1072
1073 #[tokio::test(flavor = "current_thread")]
1074 async fn expired_claim_allows_reclaim() {
1075 let log = Arc::new(AnyEventLog::Memory(MemoryEventLog::new(32)));
1076 let queue = WorkerQueue::new(log.clone());
1077 let receipt = queue
1078 .enqueue(&test_job(
1079 "triage",
1080 "incoming-review-task",
1081 "evt-1",
1082 WorkerQueuePriority::Normal,
1083 ))
1084 .await
1085 .unwrap();
1086 let expired_claim = WorkerQueueClaimRecord {
1087 job_event_id: receipt.job_event_id,
1088 claim_id: "expired-claim".to_string(),
1089 consumer_id: "consumer-a".to_string(),
1090 claimed_at_ms: now_ms().saturating_sub(2),
1091 expires_at_ms: now_ms().saturating_sub(1),
1092 scheduling: None,
1093 };
1094 log.append(
1095 &claims_topic("triage").unwrap(),
1096 LogEvent::new("job_claimed", serde_json::to_value(&expired_claim).unwrap()),
1097 )
1098 .await
1099 .unwrap();
1100 let second = queue
1101 .claim_next("triage", "consumer-b", StdDuration::from_mins(1))
1102 .await
1103 .unwrap()
1104 .unwrap();
1105 assert_eq!(second.job.event.id.0, "evt-1");
1106 assert_ne!(second.handle.claim_id, expired_claim.claim_id);
1107 assert_eq!(second.handle.consumer_id, "consumer-b");
1108 }
1109
1110 #[tokio::test(flavor = "current_thread")]
1111 async fn high_priority_and_aged_normal_are_selected_first() {
1112 let log = Arc::new(AnyEventLog::Memory(MemoryEventLog::new(32)));
1113 let queue = WorkerQueue::new(log.clone());
1114
1115 let catalog_topic = Topic::new(WORKER_QUEUE_CATALOG_TOPIC).unwrap();
1116 log.append(
1117 &catalog_topic,
1118 LogEvent::new("queue_seen", serde_json::json!({"queue":"triage"})),
1119 )
1120 .await
1121 .unwrap();
1122
1123 let topic = job_topic("triage").unwrap();
1124 let mut old_normal = LogEvent::new(
1125 "trigger_dispatch",
1126 serde_json::to_value(test_job(
1127 "triage",
1128 "incoming-review-task",
1129 "evt-old-normal",
1130 WorkerQueuePriority::Normal,
1131 ))
1132 .unwrap(),
1133 );
1134 old_normal.occurred_at_ms = now_ms() - scheduling::NORMAL_PROMOTION_AGE_MS - 1_000;
1135 log.append(&topic, old_normal).await.unwrap();
1136
1137 let high = LogEvent::new(
1138 "trigger_dispatch",
1139 serde_json::to_value(test_job(
1140 "triage",
1141 "incoming-review-task",
1142 "evt-high",
1143 WorkerQueuePriority::High,
1144 ))
1145 .unwrap(),
1146 );
1147 log.append(&topic, high).await.unwrap();
1148
1149 let claimed = queue
1150 .claim_next("triage", "consumer-a", StdDuration::from_mins(1))
1151 .await
1152 .unwrap()
1153 .unwrap();
1154 assert_eq!(claimed.job.event.id.0, "evt-old-normal");
1155 }
1156
1157 fn tenant_event(id: &str, tenant: &str) -> TriggerEvent {
1158 let mut event = test_event(id);
1159 event.tenant_id = Some(crate::triggers::TenantId::new(tenant));
1160 event
1161 }
1162
1163 fn tenant_job(
1164 queue: &str,
1165 trigger_id: &str,
1166 event_id: &str,
1167 tenant: &str,
1168 priority: WorkerQueuePriority,
1169 ) -> WorkerQueueJob {
1170 WorkerQueueJob {
1171 queue: queue.to_string(),
1172 trigger_id: trigger_id.to_string(),
1173 binding_key: format!("{trigger_id}@v1"),
1174 binding_version: 1,
1175 event: tenant_event(event_id, tenant),
1176 replay_of_event_id: None,
1177 priority,
1178 }
1179 }
1180
1181 async fn ack_and_respond(queue: &WorkerQueue, queue_name: &str, claim: &ClaimedWorkerJob) {
1182 queue
1183 .append_response(
1184 queue_name,
1185 &WorkerQueueResponseRecord {
1186 queue: queue_name.to_string(),
1187 job_event_id: claim.handle.job_event_id,
1188 consumer_id: claim.handle.consumer_id.clone(),
1189 handled_at_ms: now_ms(),
1190 outcome: Some(DispatchOutcome {
1191 trigger_id: claim.job.trigger_id.clone(),
1192 binding_key: claim.job.binding_key.clone(),
1193 event_id: claim.job.event.id.0.clone(),
1194 attempt_count: 1,
1195 status: super::super::DispatchStatus::Succeeded,
1196 handler_kind: "local".to_string(),
1197 target_uri: "test::handler".to_string(),
1198 replay_of_event_id: None,
1199 result: None,
1200 error: None,
1201 }),
1202 error: None,
1203 },
1204 )
1205 .await
1206 .unwrap();
1207 queue.ack_claim(&claim.handle).await.unwrap();
1208 }
1209
1210 #[tokio::test(flavor = "current_thread")]
1211 async fn drr_policy_rotates_across_tenants_through_claim_next() {
1212 let log = Arc::new(AnyEventLog::Memory(MemoryEventLog::new(256)));
1213 let queue = WorkerQueue::with_policy(
1214 log,
1215 SchedulerPolicy::deficit_round_robin(scheduler::FairnessKey::Tenant),
1216 );
1217
1218 for idx in 0..8 {
1220 queue
1221 .enqueue(&tenant_job(
1222 "triage",
1223 "trigger",
1224 &format!("a-{idx}"),
1225 "tenant-a",
1226 WorkerQueuePriority::Normal,
1227 ))
1228 .await
1229 .unwrap();
1230 }
1231 queue
1232 .enqueue(&tenant_job(
1233 "triage",
1234 "trigger",
1235 "b-1",
1236 "tenant-b",
1237 WorkerQueuePriority::Normal,
1238 ))
1239 .await
1240 .unwrap();
1241
1242 let mut tenants_seen = Vec::new();
1246 for n in 0..4 {
1247 let consumer = format!("c-{n}");
1248 let claim = queue
1249 .claim_next("triage", &consumer, StdDuration::from_mins(1))
1250 .await
1251 .unwrap()
1252 .expect("queue should still have ready jobs");
1253 tenants_seen.push(
1254 claim
1255 .job
1256 .event
1257 .tenant_id
1258 .as_ref()
1259 .map(|t| t.0.clone())
1260 .unwrap_or_default(),
1261 );
1262 ack_and_respond(&queue, "triage", &claim).await;
1263 }
1264
1265 let saw_b = tenants_seen.iter().any(|t| t == "tenant-b");
1266 assert!(
1267 saw_b,
1268 "tenant-b should have been served within the first 4 claims, got {tenants_seen:?}",
1269 );
1270 }
1271
1272 #[tokio::test(flavor = "current_thread")]
1273 async fn fifo_policy_preserves_legacy_behavior() {
1274 let log = Arc::new(AnyEventLog::Memory(MemoryEventLog::new(64)));
1275 let queue = WorkerQueue::with_policy(log, SchedulerPolicy::fifo());
1276
1277 for idx in 0..4 {
1279 queue
1280 .enqueue(&tenant_job(
1281 "triage",
1282 "trigger",
1283 &format!("a-{idx}"),
1284 "tenant-a",
1285 WorkerQueuePriority::Normal,
1286 ))
1287 .await
1288 .unwrap();
1289 }
1290 queue
1291 .enqueue(&tenant_job(
1292 "triage",
1293 "trigger",
1294 "b-1",
1295 "tenant-b",
1296 WorkerQueuePriority::Normal,
1297 ))
1298 .await
1299 .unwrap();
1300
1301 let first = queue
1303 .claim_next("triage", "c-0", StdDuration::from_mins(1))
1304 .await
1305 .unwrap()
1306 .unwrap();
1307 assert_eq!(first.job.event.tenant_id.unwrap().0, "tenant-a");
1308 }
1309
1310 #[tokio::test(flavor = "current_thread")]
1311 async fn inspect_queue_reports_per_tenant_fairness_state() {
1312 let log = Arc::new(AnyEventLog::Memory(MemoryEventLog::new(64)));
1313 let queue = WorkerQueue::with_policy(
1314 log,
1315 SchedulerPolicy::deficit_round_robin(scheduler::FairnessKey::Tenant)
1316 .with_weight("tenant-a", 2)
1317 .with_weight("tenant-b", 1),
1318 );
1319
1320 for idx in 0..3 {
1321 queue
1322 .enqueue(&tenant_job(
1323 "triage",
1324 "trigger",
1325 &format!("a-{idx}"),
1326 "tenant-a",
1327 WorkerQueuePriority::Normal,
1328 ))
1329 .await
1330 .unwrap();
1331 }
1332 queue
1333 .enqueue(&tenant_job(
1334 "triage",
1335 "trigger",
1336 "b-1",
1337 "tenant-b",
1338 WorkerQueuePriority::Normal,
1339 ))
1340 .await
1341 .unwrap();
1342
1343 for n in 0..2 {
1344 let consumer = format!("c-{n}");
1345 let claim = queue
1346 .claim_next("triage", &consumer, StdDuration::from_mins(1))
1347 .await
1348 .unwrap()
1349 .unwrap();
1350 ack_and_respond(&queue, "triage", &claim).await;
1351 }
1352
1353 let snap = queue.inspect_queue("triage").await.unwrap();
1354 assert_eq!(snap.scheduler.strategy, "drr");
1355 assert_eq!(snap.scheduler.fairness_key, "tenant");
1356 assert!(snap
1357 .scheduler
1358 .keys
1359 .iter()
1360 .any(|k| k.fairness_key == "tenant-a"));
1361 let weights: BTreeMap<String, u32> = snap
1362 .scheduler
1363 .keys
1364 .iter()
1365 .map(|k| (k.fairness_key.clone(), k.weight))
1366 .collect();
1367 assert_eq!(weights.get("tenant-a").copied(), Some(2));
1368 assert_eq!(weights.get("tenant-b").copied(), Some(1));
1369 }
1370
1371 #[tokio::test(flavor = "current_thread")]
1372 async fn drr_with_max_concurrent_per_key_throttles_hot_tenant() {
1373 let log = Arc::new(AnyEventLog::Memory(MemoryEventLog::new(128)));
1374 let queue = WorkerQueue::with_policy(
1375 log,
1376 SchedulerPolicy::deficit_round_robin(scheduler::FairnessKey::Tenant)
1377 .with_max_concurrent_per_key(1),
1378 );
1379
1380 for idx in 0..4 {
1381 queue
1382 .enqueue(&tenant_job(
1383 "triage",
1384 "trigger",
1385 &format!("a-{idx}"),
1386 "tenant-a",
1387 WorkerQueuePriority::Normal,
1388 ))
1389 .await
1390 .unwrap();
1391 }
1392 queue
1393 .enqueue(&tenant_job(
1394 "triage",
1395 "trigger",
1396 "b-1",
1397 "tenant-b",
1398 WorkerQueuePriority::Normal,
1399 ))
1400 .await
1401 .unwrap();
1402
1403 let first = queue
1404 .claim_next("triage", "consumer-a", StdDuration::from_mins(1))
1405 .await
1406 .unwrap()
1407 .unwrap();
1408 let second = queue
1411 .claim_next("triage", "consumer-b", StdDuration::from_mins(1))
1412 .await
1413 .unwrap()
1414 .unwrap();
1415 let pair = [
1416 first.job.event.tenant_id.clone().unwrap().0,
1417 second.job.event.tenant_id.unwrap().0,
1418 ];
1419 assert!(
1420 pair.contains(&"tenant-a".to_string()) && pair.contains(&"tenant-b".to_string()),
1421 "max_concurrent_per_key=1 must release tenant-b within two claims, got {pair:?}",
1422 );
1423 }
1424
1425 #[test]
1426 fn from_env_parses_drr_policy_from_lookup() {
1427 let lookup = |name: &str| -> Option<String> {
1428 match name {
1429 "HARN_SCHEDULER_STRATEGY" => Some("drr".to_string()),
1430 "HARN_SCHEDULER_FAIRNESS_KEY" => Some("tenant-and-binding".to_string()),
1431 "HARN_SCHEDULER_QUANTUM" => Some("3".to_string()),
1432 "HARN_SCHEDULER_STARVATION_AGE_MS" => Some("750".to_string()),
1433 "HARN_SCHEDULER_MAX_CONCURRENT_PER_KEY" => Some("4".to_string()),
1434 "HARN_SCHEDULER_DEFAULT_WEIGHT" => Some("2".to_string()),
1435 "HARN_SCHEDULER_WEIGHTS" => Some("tenant-a:5,tenant-b:1, : ,bad:abc".to_string()),
1436 _ => None,
1437 }
1438 };
1439 let policy = SchedulerPolicy::from_env_lookup(lookup);
1440 match policy.strategy {
1441 SchedulerStrategy::DeficitRoundRobin {
1442 quantum,
1443 starvation_age_ms,
1444 } => {
1445 assert_eq!(quantum, 3);
1446 assert_eq!(starvation_age_ms, Some(750));
1447 }
1448 other => panic!("expected DRR strategy, got {other:?}"),
1449 }
1450 assert_eq!(
1451 policy.fairness_key,
1452 scheduler::FairnessKey::TenantAndBinding
1453 );
1454 assert_eq!(policy.max_concurrent_per_key, 4);
1455 assert_eq!(policy.default_weight, 2);
1456 assert_eq!(policy.weight_for("tenant-a"), 5);
1457 assert_eq!(policy.weight_for("tenant-b"), 1);
1458 assert_eq!(policy.weight_for("tenant-c"), 2);
1460 }
1461
1462 #[test]
1463 fn from_env_defaults_to_fifo_when_missing() {
1464 let policy = SchedulerPolicy::from_env_lookup(|_| None);
1465 assert!(matches!(policy.strategy, SchedulerStrategy::Fifo));
1466 }
1467}