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, SchedulerPolicy, SchedulerSnapshot, SchedulerState};
13use super::{DispatchOutcome, TriggerEvent};
14
15#[cfg(test)]
16mod exclusion_tests;
17mod scheduling;
18mod state;
19
20pub use scheduling::{
21    WorkerQueuePriority, WorkerQueueSchedulingDecision, WorkerQueueSchedulingReceipt,
22    DEFERRABLE_PROMOTION_AGE_MS,
23};
24pub use state::{WorkerQueueJobState, WorkerQueueState, WorkerQueueSummary};
25
26pub const WORKER_QUEUE_CATALOG_TOPIC: &str = "worker.queues";
27const WORKER_QUEUE_CLAIMS_SUFFIX: &str = ".claims";
28const WORKER_QUEUE_RESPONSES_SUFFIX: &str = ".responses";
29
30#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
31pub struct WorkerQueueJob {
32    pub queue: String,
33    pub trigger_id: String,
34    pub binding_key: String,
35    pub binding_version: u32,
36    pub event: TriggerEvent,
37    #[serde(default)]
38    pub replay_of_event_id: Option<String>,
39    #[serde(default)]
40    pub priority: WorkerQueuePriority,
41}
42
43#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
44pub struct WorkerQueueEnqueueReceipt {
45    pub queue: String,
46    pub job_event_id: u64,
47    pub response_topic: String,
48}
49
50#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
51pub struct WorkerQueueClaimHandle {
52    pub queue: String,
53    pub job_event_id: u64,
54    pub claim_id: String,
55    pub consumer_id: String,
56    pub expires_at_ms: i64,
57}
58
59#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
60pub struct ClaimedWorkerJob {
61    pub handle: WorkerQueueClaimHandle,
62    pub job: WorkerQueueJob,
63    pub scheduling: WorkerQueueSchedulingReceipt,
64}
65
66#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
67pub struct WorkerQueueResponseRecord {
68    pub queue: String,
69    pub job_event_id: u64,
70    pub consumer_id: String,
71    pub handled_at_ms: i64,
72    pub outcome: Option<DispatchOutcome>,
73    pub error: Option<String>,
74}
75
76#[derive(Clone)]
77pub struct WorkerQueue {
78    event_log: Arc<AnyEventLog>,
79    /// Active scheduler policy. Reads on every claim so it can be hot-swapped
80    /// at runtime without rebuilding the queue.
81    policy: Arc<RwLock<SchedulerPolicy>>,
82    /// Per-queue ephemeral scheduler state. Keyed by queue name; entries are
83    /// created lazily on first claim. Self-correcting — safe to lose on
84    /// process restart.
85    scheduler_states: Arc<Mutex<BTreeMap<String, SchedulerState>>>,
86}
87
88#[derive(Clone, Debug, Serialize)]
89pub struct WorkerQueueInspectSnapshot {
90    pub summary: WorkerQueueSummary,
91    pub scheduler: SchedulerSnapshot,
92}
93
94impl WorkerQueue {
95    /// Construct a `WorkerQueue` using the policy derived from the
96    /// `HARN_SCHEDULER_*` environment variables (see
97    /// [`SchedulerPolicy::from_env`]). Defaults to FIFO so single-tenant
98    /// deployments behave exactly as before unless they opt in.
99    pub fn new(event_log: Arc<AnyEventLog>) -> Self {
100        Self::with_policy(event_log, SchedulerPolicy::from_env())
101    }
102
103    pub fn with_policy(event_log: Arc<AnyEventLog>, policy: SchedulerPolicy) -> Self {
104        Self {
105            event_log,
106            policy: Arc::new(RwLock::new(policy)),
107            scheduler_states: Arc::new(Mutex::new(BTreeMap::new())),
108        }
109    }
110
111    /// Replace the active scheduler policy. Existing per-queue state is
112    /// preserved (deficits self-correct against the new weights).
113    pub fn set_policy(&self, policy: SchedulerPolicy) {
114        *self.policy.write().expect("scheduler policy poisoned") = policy;
115    }
116
117    pub fn policy(&self) -> SchedulerPolicy {
118        self.policy
119            .read()
120            .expect("scheduler policy poisoned")
121            .clone()
122    }
123
124    pub async fn enqueue(
125        &self,
126        job: &WorkerQueueJob,
127    ) -> Result<WorkerQueueEnqueueReceipt, LogError> {
128        let queue = job.queue.trim();
129        if queue.is_empty() {
130            return Err(LogError::Config(
131                "worker queue name cannot be empty".to_string(),
132            ));
133        }
134        let queue_name = queue.to_string();
135        let catalog_topic = Topic::new(WORKER_QUEUE_CATALOG_TOPIC)
136            .expect("static worker queue catalog topic should always be valid");
137        self.event_log
138            .append(
139                &catalog_topic,
140                LogEvent::new(
141                    "queue_seen",
142                    serde_json::to_value(WorkerQueueCatalogRecord {
143                        queue: queue_name.clone(),
144                    })
145                    .map_err(|error| LogError::Serde(error.to_string()))?,
146                ),
147            )
148            .await?;
149
150        let job_topic = job_topic(&queue_name)?;
151        let mut headers = BTreeMap::new();
152        headers.insert("queue".to_string(), queue_name.clone());
153        headers.insert("trigger_id".to_string(), job.trigger_id.clone());
154        headers.insert("binding_key".to_string(), job.binding_key.clone());
155        headers.insert("event_id".to_string(), job.event.id.0.clone());
156        headers.insert("priority".to_string(), job.priority.as_str().to_string());
157        let job_event_id = self
158            .event_log
159            .append(
160                &job_topic,
161                LogEvent::new(
162                    "trigger_dispatch",
163                    serde_json::to_value(job)
164                        .map_err(|error| LogError::Serde(error.to_string()))?,
165                )
166                .with_headers(headers),
167            )
168            .await?;
169        if let Some(metrics) = crate::active_metrics_registry() {
170            if let Ok(state) = self.queue_state(&queue_name).await {
171                let summary = state.summary(now_ms());
172                metrics.set_worker_queue_depth(
173                    &queue_name,
174                    (summary.ready + summary.in_flight) as u64,
175                );
176            }
177        }
178        Ok(WorkerQueueEnqueueReceipt {
179            queue: queue_name.clone(),
180            job_event_id,
181            response_topic: response_topic_name(&queue_name),
182        })
183    }
184
185    pub async fn known_queues(&self) -> Result<Vec<String>, LogError> {
186        let topic = Topic::new(WORKER_QUEUE_CATALOG_TOPIC)
187            .expect("static worker queue catalog topic should always be valid");
188        let events = self.event_log.read_range(&topic, None, usize::MAX).await?;
189        let mut queues = BTreeSet::new();
190        for (_, event) in events {
191            if event.kind != "queue_seen" {
192                continue;
193            }
194            let record: WorkerQueueCatalogRecord = serde_json::from_value(event.payload)
195                .map_err(|error| LogError::Serde(error.to_string()))?;
196            if !record.queue.trim().is_empty() {
197                queues.insert(record.queue);
198            }
199        }
200        Ok(queues.into_iter().collect())
201    }
202
203    pub async fn queue_state(&self, queue: &str) -> Result<WorkerQueueState, LogError> {
204        let queue_name = queue.trim();
205        if queue_name.is_empty() {
206            return Err(LogError::Config(
207                "worker queue name cannot be empty".to_string(),
208            ));
209        }
210        let now_ms = now_ms();
211        let job_events = self
212            .event_log
213            .read_range(&job_topic(queue_name)?, None, usize::MAX)
214            .await?;
215        let claim_events = self
216            .event_log
217            .read_range(&claims_topic(queue_name)?, None, usize::MAX)
218            .await?;
219        let response_events = self
220            .event_log
221            .read_range(&responses_topic(queue_name)?, None, usize::MAX)
222            .await?;
223
224        let mut jobs = BTreeMap::<u64, WorkerQueueJobStateInternal>::new();
225        for (job_event_id, event) in job_events {
226            if event.kind != "trigger_dispatch" {
227                continue;
228            }
229            let job: WorkerQueueJob = serde_json::from_value(event.payload)
230                .map_err(|error| LogError::Serde(error.to_string()))?;
231            jobs.insert(
232                job_event_id,
233                WorkerQueueJobStateInternal {
234                    job_event_id,
235                    enqueued_at_ms: event.occurred_at_ms,
236                    job,
237                    active_claim: None,
238                    acked: false,
239                    purged: false,
240                    seen_claim_ids: BTreeSet::new(),
241                },
242            );
243        }
244
245        for (_, event) in claim_events {
246            match event.kind.as_str() {
247                "job_claimed" => {
248                    let claim: WorkerQueueClaimRecord = serde_json::from_value(event.payload)
249                        .map_err(|error| LogError::Serde(error.to_string()))?;
250                    let Some(job) = jobs.get_mut(&claim.job_event_id) else {
251                        continue;
252                    };
253                    if job.acked || job.purged {
254                        continue;
255                    }
256                    job.seen_claim_ids.insert(claim.claim_id.clone());
257                    let can_take = job
258                        .active_claim
259                        .as_ref()
260                        .is_none_or(|active| active.expires_at_ms <= claim.claimed_at_ms);
261                    if can_take {
262                        job.active_claim = Some(WorkerQueueClaimHandle {
263                            queue: queue_name.to_string(),
264                            job_event_id: claim.job_event_id,
265                            claim_id: claim.claim_id,
266                            consumer_id: claim.consumer_id,
267                            expires_at_ms: claim.expires_at_ms,
268                        });
269                    }
270                }
271                "claim_renewed" => {
272                    let renewal: WorkerQueueClaimRenewalRecord =
273                        serde_json::from_value(event.payload)
274                            .map_err(|error| LogError::Serde(error.to_string()))?;
275                    let Some(job) = jobs.get_mut(&renewal.job_event_id) else {
276                        continue;
277                    };
278                    if let Some(active) = job.active_claim.as_mut() {
279                        if active.claim_id == renewal.claim_id {
280                            active.expires_at_ms = renewal.expires_at_ms;
281                        }
282                    }
283                }
284                "job_released" => {
285                    let release: WorkerQueueReleaseRecord =
286                        serde_json::from_value(event.payload)
287                            .map_err(|error| LogError::Serde(error.to_string()))?;
288                    let Some(job) = jobs.get_mut(&release.job_event_id) else {
289                        continue;
290                    };
291                    if job
292                        .active_claim
293                        .as_ref()
294                        .is_some_and(|active| active.claim_id == release.claim_id)
295                    {
296                        job.active_claim = None;
297                    }
298                }
299                "job_acked" => {
300                    let ack: WorkerQueueAckRecord = serde_json::from_value(event.payload)
301                        .map_err(|error| LogError::Serde(error.to_string()))?;
302                    let Some(job) = jobs.get_mut(&ack.job_event_id) else {
303                        continue;
304                    };
305                    if ack.claim_id.is_empty() || job.seen_claim_ids.contains(&ack.claim_id) {
306                        job.acked = true;
307                        job.active_claim = None;
308                    }
309                }
310                "job_purged" => {
311                    let purge: WorkerQueuePurgeRecord = serde_json::from_value(event.payload)
312                        .map_err(|error| LogError::Serde(error.to_string()))?;
313                    let Some(job) = jobs.get_mut(&purge.job_event_id) else {
314                        continue;
315                    };
316                    if !job.acked {
317                        job.purged = true;
318                        job.active_claim = None;
319                    }
320                }
321                _ => {}
322            }
323        }
324
325        let responses = response_events
326            .into_iter()
327            .filter(|(_, event)| event.kind == "job_response")
328            .map(|(_, event)| {
329                serde_json::from_value::<WorkerQueueResponseRecord>(event.payload)
330                    .map_err(|error| LogError::Serde(error.to_string()))
331            })
332            .collect::<Result<Vec<_>, _>>()?;
333
334        let mut queue_state = WorkerQueueState {
335            queue: queue_name.to_string(),
336            responses,
337            jobs: jobs
338                .into_values()
339                .map(|mut job| {
340                    if job
341                        .active_claim
342                        .as_ref()
343                        .is_some_and(|active| active.expires_at_ms <= now_ms)
344                    {
345                        job.active_claim = None;
346                    }
347                    WorkerQueueJobState {
348                        job_event_id: job.job_event_id,
349                        enqueued_at_ms: job.enqueued_at_ms,
350                        job: job.job,
351                        active_claim: job.active_claim,
352                        acked: job.acked,
353                        purged: job.purged,
354                    }
355                })
356                .collect(),
357        };
358        queue_state
359            .jobs
360            .sort_by_key(|job| (job.enqueued_at_ms, job.job_event_id));
361        Ok(queue_state)
362    }
363
364    pub async fn queue_summaries(&self) -> Result<Vec<WorkerQueueSummary>, LogError> {
365        let now_ms = now_ms();
366        let mut summaries = Vec::new();
367        for queue in self.known_queues().await? {
368            let state = self.queue_state(&queue).await?;
369            summaries.push(state.summary(now_ms));
370        }
371        summaries.sort_by(|left, right| left.queue.cmp(&right.queue));
372        Ok(summaries)
373    }
374
375    pub async fn claim_next(
376        &self,
377        queue: &str,
378        consumer_id: &str,
379        ttl: StdDuration,
380    ) -> Result<Option<ClaimedWorkerJob>, LogError> {
381        self.claim_next_excluding(queue, consumer_id, ttl, &BTreeSet::new())
382            .await
383    }
384
385    /// Claim the next eligible job except those already attempted by this
386    /// consumer operation. Exclusions affect selection only; a later caller
387    /// can reclaim an unacknowledged job through [`Self::claim_next`].
388    pub async fn claim_next_excluding(
389        &self,
390        queue: &str,
391        consumer_id: &str,
392        ttl: StdDuration,
393        excluded_job_event_ids: &BTreeSet<u64>,
394    ) -> Result<Option<ClaimedWorkerJob>, LogError> {
395        let queue_name = queue.trim();
396        if queue_name.is_empty() {
397            return Err(LogError::Config(
398                "worker queue name cannot be empty".to_string(),
399            ));
400        }
401        if consumer_id.trim().is_empty() {
402            return Err(LogError::InvalidConsumer(
403                "worker queue consumer id cannot be empty".to_string(),
404            ));
405        }
406        let policy = self.policy();
407        for _ in 0..8 {
408            let now_ms = now_ms();
409            let state = self.queue_state(queue_name).await?;
410            let (job, selection) = {
411                let mut states = self
412                    .scheduler_states
413                    .lock()
414                    .expect("scheduler state poisoned");
415                let scheduler_state = states.entry(queue_name.to_string()).or_default();
416                let Some((job, selection)) = state.next_ready_job_with_scheduler(
417                    scheduler_state,
418                    &policy,
419                    now_ms,
420                    excluded_job_event_ids,
421                ) else {
422                    return Ok(None);
423                };
424                (job.clone(), selection)
425            };
426            let scheduling = WorkerQueueSchedulingReceipt {
427                selected_at_ms: now_ms,
428                enqueued_at_ms: job.enqueued_at_ms,
429                waited_ms: now_ms.saturating_sub(job.enqueued_at_ms).max(0) as u64,
430                priority: job.job.priority,
431                decision: selection.decision,
432                promotion_deadline_at_ms: selection.promotion_deadline_at_ms,
433                fairness_key: selection.fairness_key.clone(),
434            };
435            let claim = WorkerQueueClaimRecord {
436                job_event_id: job.job_event_id,
437                claim_id: Uuid::new_v4().to_string(),
438                consumer_id: consumer_id.to_string(),
439                claimed_at_ms: now_ms,
440                expires_at_ms: expiry_ms(now_ms, ttl),
441                scheduling: Some(scheduling.clone()),
442            };
443            self.event_log
444                .append(
445                    &claims_topic(queue_name)?,
446                    LogEvent::new(
447                        "job_claimed",
448                        serde_json::to_value(&claim)
449                            .map_err(|error| LogError::Serde(error.to_string()))?,
450                    ),
451                )
452                .await?;
453            let refreshed = self.queue_state(queue_name).await?;
454            if refreshed
455                .active_claim_for(job.job_event_id)
456                .is_some_and(|active| active.claim_id == claim.claim_id)
457            {
458                {
459                    let mut states = self
460                        .scheduler_states
461                        .lock()
462                        .expect("scheduler state poisoned");
463                    let scheduler_state = states.entry(queue_name.to_string()).or_default();
464                    scheduler_state.note_claim_committed(&selection.fairness_key);
465                }
466                if let Some(metrics) = crate::active_metrics_registry() {
467                    let summary = refreshed.summary(now_ms);
468                    metrics.record_worker_queue_claim_age(
469                        queue_name,
470                        now_ms.saturating_sub(job.enqueued_at_ms) as f64 / 1000.0,
471                    );
472                    metrics.set_worker_queue_depth(
473                        queue_name,
474                        (summary.ready + summary.in_flight) as u64,
475                    );
476                    metrics.record_scheduler_selection(
477                        queue_name,
478                        policy.fairness_key.as_str(),
479                        &selection.fairness_key,
480                    );
481                    if let Ok(snap) = self.inspect_queue(queue_name).await {
482                        for stat in &snap.scheduler.keys {
483                            metrics.set_scheduler_deficit(
484                                queue_name,
485                                policy.fairness_key.as_str(),
486                                &stat.fairness_key,
487                                stat.deficit,
488                            );
489                            metrics.set_scheduler_oldest_eligible_age(
490                                queue_name,
491                                policy.fairness_key.as_str(),
492                                &stat.fairness_key,
493                                stat.oldest_ready_age_ms,
494                            );
495                        }
496                    }
497                }
498                return Ok(Some(ClaimedWorkerJob {
499                    handle: WorkerQueueClaimHandle {
500                        queue: queue_name.to_string(),
501                        job_event_id: claim.job_event_id,
502                        claim_id: claim.claim_id,
503                        consumer_id: claim.consumer_id,
504                        expires_at_ms: claim.expires_at_ms,
505                    },
506                    job: job.job,
507                    scheduling,
508                }));
509            }
510        }
511        Ok(None)
512    }
513
514    /// Build a fairness-aware inspect snapshot for `queue` that includes
515    /// scheduler state alongside the standard summary.
516    pub async fn inspect_queue(&self, queue: &str) -> Result<WorkerQueueInspectSnapshot, LogError> {
517        let queue_name = queue.trim();
518        if queue_name.is_empty() {
519            return Err(LogError::Config(
520                "worker queue name cannot be empty".to_string(),
521            ));
522        }
523        let now_ms = now_ms();
524        let state = self.queue_state(queue_name).await?;
525        let summary = state.summary(now_ms);
526        let policy = self.policy();
527        let ready = scheduler::ready_stats_by_key(&state.jobs, &policy, now_ms);
528        // Make sure in-flight stays authoritative against the rebuilt log.
529        let in_flight = scheduler::in_flight_by_key(&state.jobs, &policy);
530        let scheduler_snapshot = {
531            let mut states = self
532                .scheduler_states
533                .lock()
534                .expect("scheduler state poisoned");
535            let scheduler_state = states.entry(queue_name.to_string()).or_default();
536            scheduler_state.replace_in_flight(in_flight);
537            scheduler_state.snapshot(&policy, &ready)
538        };
539        Ok(WorkerQueueInspectSnapshot {
540            summary,
541            scheduler: scheduler_snapshot,
542        })
543    }
544
545    /// Inspect snapshots for every known queue.
546    pub async fn inspect_all_queues(&self) -> Result<Vec<WorkerQueueInspectSnapshot>, LogError> {
547        let mut snapshots = Vec::new();
548        for queue in self.known_queues().await? {
549            snapshots.push(self.inspect_queue(&queue).await?);
550        }
551        snapshots.sort_by(|left, right| left.summary.queue.cmp(&right.summary.queue));
552        Ok(snapshots)
553    }
554
555    pub async fn renew_claim(
556        &self,
557        handle: &WorkerQueueClaimHandle,
558        ttl: StdDuration,
559    ) -> Result<bool, LogError> {
560        let now_ms = now_ms();
561        let renewal = WorkerQueueClaimRenewalRecord {
562            job_event_id: handle.job_event_id,
563            claim_id: handle.claim_id.clone(),
564            consumer_id: handle.consumer_id.clone(),
565            renewed_at_ms: now_ms,
566            expires_at_ms: expiry_ms(now_ms, ttl),
567        };
568        self.event_log
569            .append(
570                &claims_topic(&handle.queue)?,
571                LogEvent::new(
572                    "claim_renewed",
573                    serde_json::to_value(&renewal)
574                        .map_err(|error| LogError::Serde(error.to_string()))?,
575                ),
576            )
577            .await?;
578        let refreshed = self.queue_state(&handle.queue).await?;
579        Ok(refreshed
580            .active_claim_for(handle.job_event_id)
581            .is_some_and(|active| active.claim_id == handle.claim_id))
582    }
583
584    pub async fn release_claim(
585        &self,
586        handle: &WorkerQueueClaimHandle,
587        reason: &str,
588    ) -> Result<(), LogError> {
589        let release = WorkerQueueReleaseRecord {
590            job_event_id: handle.job_event_id,
591            claim_id: handle.claim_id.clone(),
592            consumer_id: handle.consumer_id.clone(),
593            released_at_ms: now_ms(),
594            reason: if reason.trim().is_empty() {
595                None
596            } else {
597                Some(reason.to_string())
598            },
599        };
600        self.event_log
601            .append(
602                &claims_topic(&handle.queue)?,
603                LogEvent::new(
604                    "job_released",
605                    serde_json::to_value(&release)
606                        .map_err(|error| LogError::Serde(error.to_string()))?,
607                ),
608            )
609            .await?;
610        Ok(())
611    }
612
613    pub async fn append_response(
614        &self,
615        queue: &str,
616        response: &WorkerQueueResponseRecord,
617    ) -> Result<u64, LogError> {
618        self.event_log
619            .append(
620                &responses_topic(queue)?,
621                LogEvent::new(
622                    "job_response",
623                    serde_json::to_value(response)
624                        .map_err(|error| LogError::Serde(error.to_string()))?,
625                ),
626            )
627            .await
628    }
629
630    pub async fn ack_claim(&self, handle: &WorkerQueueClaimHandle) -> Result<u64, LogError> {
631        self.event_log
632            .append(
633                &claims_topic(&handle.queue)?,
634                LogEvent::new(
635                    "job_acked",
636                    serde_json::to_value(WorkerQueueAckRecord {
637                        job_event_id: handle.job_event_id,
638                        claim_id: handle.claim_id.clone(),
639                        consumer_id: handle.consumer_id.clone(),
640                        acked_at_ms: now_ms(),
641                    })
642                    .map_err(|error| LogError::Serde(error.to_string()))?,
643                ),
644            )
645            .await
646    }
647
648    pub async fn ack_job(
649        &self,
650        queue: &str,
651        job_event_id: u64,
652        consumer_id: &str,
653    ) -> Result<bool, LogError> {
654        let queue_name = queue.trim();
655        if queue_name.is_empty() {
656            return Err(LogError::Config(
657                "worker queue name cannot be empty".to_string(),
658            ));
659        }
660        let state = self.queue_state(queue_name).await?;
661        let Some(job) = state
662            .jobs
663            .iter()
664            .find(|job| job.job_event_id == job_event_id)
665        else {
666            return Ok(false);
667        };
668        if job.acked || job.purged {
669            return Ok(false);
670        }
671        self.event_log
672            .append(
673                &claims_topic(queue_name)?,
674                LogEvent::new(
675                    "job_acked",
676                    serde_json::to_value(WorkerQueueAckRecord {
677                        job_event_id,
678                        claim_id: String::new(),
679                        consumer_id: consumer_id.to_string(),
680                        acked_at_ms: now_ms(),
681                    })
682                    .map_err(|error| LogError::Serde(error.to_string()))?,
683                ),
684            )
685            .await?;
686        Ok(true)
687    }
688
689    pub async fn purge_unclaimed(
690        &self,
691        queue: &str,
692        purged_by: &str,
693        reason: Option<&str>,
694    ) -> Result<usize, LogError> {
695        let state = self.queue_state(queue).await?;
696        let ready_jobs: Vec<_> = state
697            .jobs
698            .into_iter()
699            .filter(|job| job.is_ready())
700            .map(|job| job.job_event_id)
701            .collect();
702        for job_event_id in &ready_jobs {
703            self.event_log
704                .append(
705                    &claims_topic(queue)?,
706                    LogEvent::new(
707                        "job_purged",
708                        serde_json::to_value(WorkerQueuePurgeRecord {
709                            job_event_id: *job_event_id,
710                            purged_by: purged_by.to_string(),
711                            purged_at_ms: now_ms(),
712                            reason: reason
713                                .filter(|value| !value.trim().is_empty())
714                                .map(|value| value.to_string()),
715                        })
716                        .map_err(|error| LogError::Serde(error.to_string()))?,
717                    ),
718                )
719                .await?;
720        }
721        Ok(ready_jobs.len())
722    }
723}
724
725#[derive(Clone, Debug)]
726struct WorkerQueueJobStateInternal {
727    job_event_id: u64,
728    enqueued_at_ms: i64,
729    job: WorkerQueueJob,
730    active_claim: Option<WorkerQueueClaimHandle>,
731    acked: bool,
732    purged: bool,
733    seen_claim_ids: BTreeSet<String>,
734}
735
736#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
737struct WorkerQueueCatalogRecord {
738    queue: String,
739}
740
741#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
742struct WorkerQueueClaimRecord {
743    job_event_id: u64,
744    claim_id: String,
745    consumer_id: String,
746    claimed_at_ms: i64,
747    expires_at_ms: i64,
748    #[serde(default, skip_serializing_if = "Option::is_none")]
749    scheduling: Option<WorkerQueueSchedulingReceipt>,
750}
751
752#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
753struct WorkerQueueClaimRenewalRecord {
754    job_event_id: u64,
755    claim_id: String,
756    consumer_id: String,
757    renewed_at_ms: i64,
758    expires_at_ms: i64,
759}
760
761#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
762struct WorkerQueueReleaseRecord {
763    job_event_id: u64,
764    claim_id: String,
765    consumer_id: String,
766    released_at_ms: i64,
767    #[serde(default)]
768    reason: Option<String>,
769}
770
771#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
772struct WorkerQueueAckRecord {
773    job_event_id: u64,
774    claim_id: String,
775    consumer_id: String,
776    acked_at_ms: i64,
777}
778
779#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
780struct WorkerQueuePurgeRecord {
781    job_event_id: u64,
782    purged_by: String,
783    purged_at_ms: i64,
784    #[serde(default)]
785    reason: Option<String>,
786}
787
788pub fn job_topic_name(queue: &str) -> String {
789    format!("worker.{}", sanitize_topic_component(queue))
790}
791
792pub fn claims_topic_name(queue: &str) -> String {
793    format!("{}{}", job_topic_name(queue), WORKER_QUEUE_CLAIMS_SUFFIX)
794}
795
796pub fn response_topic_name(queue: &str) -> String {
797    format!("{}{}", job_topic_name(queue), WORKER_QUEUE_RESPONSES_SUFFIX)
798}
799
800fn job_topic(queue: &str) -> Result<Topic, LogError> {
801    Topic::new(job_topic_name(queue))
802}
803
804fn claims_topic(queue: &str) -> Result<Topic, LogError> {
805    Topic::new(claims_topic_name(queue))
806}
807
808fn responses_topic(queue: &str) -> Result<Topic, LogError> {
809    Topic::new(response_topic_name(queue))
810}
811
812fn now_ms() -> i64 {
813    std::time::SystemTime::now()
814        .duration_since(std::time::UNIX_EPOCH)
815        .map(|duration| duration.as_millis() as i64)
816        .unwrap_or(0)
817}
818
819fn expiry_ms(now_ms: i64, ttl: StdDuration) -> i64 {
820    now_ms.saturating_add(ttl.as_millis().min(i64::MAX as u128) as i64)
821}
822
823#[cfg(test)]
824mod tests;