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