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