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 policy: Arc<RwLock<SchedulerPolicy>>,
83 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 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 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 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 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 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 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 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 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;