Skip to main content

harn_vm/triggers/
worker_queue.rs

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    /// Select the next ready job by consulting `scheduler` under `policy`.
134    ///
135    /// Under `Fifo` this is equivalent to picking the job with the lowest
136    /// `(priority_rank, enqueued_at_ms, job_event_id)` — the historical
137    /// behaviour. Under `DeficitRoundRobin`, candidates are grouped by the
138    /// configured fairness key and the scheduler rotates so a hot
139    /// tenant/binding cannot monopolise the queue.
140    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        // Refresh authoritative in-flight count from the rebuilt queue state.
157        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    /// Active scheduler policy. Reads on every claim so it can be hot-swapped
179    /// at runtime without rebuilding the queue.
180    policy: Arc<RwLock<SchedulerPolicy>>,
181    /// Per-queue ephemeral scheduler state. Keyed by queue name; entries are
182    /// created lazily on first claim. Self-correcting — safe to lose on
183    /// process restart.
184    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    /// Construct a `WorkerQueue` using the policy derived from the
195    /// `HARN_SCHEDULER_*` environment variables (see
196    /// [`SchedulerPolicy::from_env`]). Defaults to FIFO so single-tenant
197    /// deployments behave exactly as before unless they opt in.
198    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    /// Replace the active scheduler policy. Existing per-queue state is
211    /// preserved (deficits self-correct against the new weights).
212    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    /// Build a fairness-aware inspect snapshot for `queue` that includes
597    /// scheduler state alongside the standard summary.
598    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        // Make sure in-flight stays authoritative against the rebuilt log.
611        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    /// Inspect snapshots for every known queue.
628    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        // Tenant A enqueues 8 jobs before tenant B enqueues a single job.
1219        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        // Claim+ack 4 jobs back-to-back. Under FIFO, tenant B would never be
1243        // touched. Under fair-share, B must be served within the first two
1244        // claims.
1245        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        // Fill queue with tenant-a jobs first, then a single tenant-b job.
1278        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        // FIFO must drain all of tenant-a before touching tenant-b.
1302        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        // Without releasing the first claim, the second pick must skip the
1409        // capped tenant-a and serve tenant-b instead.
1410        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        // Unknown key falls back to default_weight.
1459        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}