1use crate::error::{QueueError, QueueResult};
4use crate::job::{Job, JobData, JobId, JobPriority, JobState, JobStatus};
5use armature_log::{debug, info, warn};
6use chrono::{DateTime, Utc};
7use redis::{AsyncCommands, Client, aio::ConnectionManager};
8use std::sync::LazyLock;
9use std::time::Duration;
10
11const SCAN_COUNT: usize = 500;
13
14const PRIORITY_LUA: &str = r#"
22local function priority_of(job_json)
23 local pname = 'normal'
24 local pscore = -1
25 local ok, job = pcall(cjson.decode, job_json)
26 if ok and type(job) == 'table' and job.priority then
27 local p = job.priority
28 if p == 'Low' then pname = 'low'; pscore = 0
29 elseif p == 'Normal' then pname = 'normal'; pscore = -1
30 elseif p == 'High' then pname = 'high'; pscore = -2
31 elseif p == 'Critical' then pname = 'critical'; pscore = -3
32 end
33 end
34 return pname, pscore
35end
36"#;
37
38const MOVE_DELAYED_BODY: &str = r#"
56local delayed_key = KEYS[1]
57local prefix = ARGV[1]
58local now = ARGV[2]
59
60local job_ids = redis.call('ZRANGEBYSCORE', delayed_key, '-inf', now)
61local promoted = 0
62local dropped = 0
63
64for _, job_id in ipairs(job_ids) do
65 local job_json = redis.call('GET', prefix .. ':job:' .. job_id)
66 if job_json then
67 -- Claim the job atomically; skip if another pass already took it.
68 if redis.call('ZREM', delayed_key, job_id) == 1 then
69 local pname, pscore = priority_of(job_json)
70 redis.call('ZADD', prefix .. ':pending:' .. pname, pscore, job_id)
71 promoted = promoted + 1
72 end
73 elseif redis.call('ZREM', delayed_key, job_id) == 1 then
74 -- Body expired before the job came due: it can never run, so drop the
75 -- id instead of rescanning it on every future pass.
76 dropped = dropped + 1
77 end
78end
79
80return {promoted, dropped}
81"#;
82
83const RECLAIM_STALE_BODY: &str = r#"
94local processing_key = KEYS[1]
95local prefix = ARGV[1]
96local cutoff = ARGV[2]
97
98local job_ids = redis.call('ZRANGEBYSCORE', processing_key, '-inf', cutoff)
99local reclaimed = 0
100local dropped = 0
101
102for _, job_id in ipairs(job_ids) do
103 local job_json = redis.call('GET', prefix .. ':job:' .. job_id)
104 if job_json then
105 if redis.call('ZREM', processing_key, job_id) == 1 then
106 local pname, pscore = priority_of(job_json)
107 redis.call('ZADD', prefix .. ':pending:' .. pname, pscore, job_id)
108 reclaimed = reclaimed + 1
109 end
110 elseif redis.call('ZREM', processing_key, job_id) == 1 then
111 dropped = dropped + 1
112 end
113end
114
115return {reclaimed, dropped}
116"#;
117
118static MOVE_DELAYED_SCRIPT: LazyLock<String> =
120 LazyLock::new(|| format!("{PRIORITY_LUA}{MOVE_DELAYED_BODY}"));
121
122static RECLAIM_STALE_SCRIPT: LazyLock<String> =
124 LazyLock::new(|| format!("{PRIORITY_LUA}{RECLAIM_STALE_BODY}"));
125
126const DEQUEUE_POP_SCRIPT: &str = r#"
144local prefix = ARGV[1]
145local now = ARGV[2]
146for i = 1, #KEYS do
147 while true do
148 local popped = redis.call('ZPOPMIN', KEYS[i], 1)
149 if not popped or not popped[1] then
150 break
151 end
152 local job_id = popped[1]
153 local job_json = redis.call('GET', prefix .. ':job:' .. job_id)
154 if job_json then
155 redis.call('ZADD', prefix .. ':processing', now, job_id)
156 return {job_id, job_json}
157 end
158 -- Job body expired between enqueue and dequeue: the id is discarded
159 -- (already popped) and we keep draining this same queue.
160 end
161end
162return nil
163"#;
164
165#[derive(Debug, Clone)]
167pub struct QueueConfig {
168 pub redis_url: String,
170
171 pub queue_name: String,
173
174 pub key_prefix: String,
176
177 pub max_size: usize,
179
180 pub retention_time: Duration,
182}
183
184impl QueueConfig {
185 pub fn new(redis_url: impl Into<String>, queue_name: impl Into<String>) -> Self {
187 let queue_name = queue_name.into();
188 Self {
189 redis_url: redis_url.into(),
190 key_prefix: format!("armature:queue:{}", queue_name),
191 queue_name,
192 max_size: 0,
193 retention_time: Duration::from_secs(86400), }
195 }
196
197 pub fn with_key_prefix(mut self, prefix: impl Into<String>) -> Self {
199 self.key_prefix = prefix.into();
200 self
201 }
202
203 pub fn with_max_size(mut self, max_size: usize) -> Self {
205 self.max_size = max_size;
206 self
207 }
208
209 pub fn with_retention_time(mut self, retention_time: Duration) -> Self {
211 self.retention_time = retention_time;
212 self
213 }
214
215 fn key(&self, suffix: &str) -> String {
217 format!("{}:{}", self.key_prefix, suffix)
218 }
219
220 fn body_ttl_secs(&self, wait: Duration) -> u64 {
230 self.retention_time.as_secs().saturating_add(wait.as_secs())
231 }
232}
233
234fn wait_until(scheduled_at: Option<DateTime<Utc>>) -> Duration {
236 scheduled_at
237 .map(|at| Duration::from_secs((at - Utc::now()).num_seconds().max(0) as u64))
238 .unwrap_or(Duration::ZERO)
239}
240
241#[derive(Clone)]
243pub struct Queue {
244 connection: ConnectionManager,
245 config: QueueConfig,
246}
247
248impl Queue {
249 pub async fn new(
251 redis_url: impl Into<String>,
252 queue_name: impl Into<String>,
253 ) -> QueueResult<Self> {
254 let config = QueueConfig::new(redis_url, queue_name);
255 Self::with_config(config).await
256 }
257
258 pub async fn with_config(config: QueueConfig) -> QueueResult<Self> {
260 info!("Initializing job queue: {}", config.queue_name);
261 debug!(
262 "Queue config - prefix: {}, max_size: {}",
263 config.key_prefix, config.max_size
264 );
265
266 let client = Client::open(config.redis_url.as_str())
267 .map_err(|e| QueueError::Config(e.to_string()))?;
268
269 let connection = ConnectionManager::new(client).await?;
270
271 info!("Job queue '{}' ready", config.queue_name);
272 Ok(Self { connection, config })
273 }
274
275 pub async fn enqueue(&self, job_type: impl Into<String>, data: JobData) -> QueueResult<JobId> {
277 let job_type = job_type.into();
278 debug!(
279 "Enqueueing job: {} on queue '{}'",
280 job_type, self.config.queue_name
281 );
282 let job = Job::new(&self.config.queue_name, &job_type, data);
283 self.enqueue_job(job).await
284 }
285
286 pub async fn enqueue_in(
293 &self,
294 delay: chrono::Duration,
295 job_type: impl Into<String>,
296 data: JobData,
297 ) -> QueueResult<JobId> {
298 let job = Job::new(&self.config.queue_name, job_type, data).schedule_after(delay);
299 self.enqueue_job(job).await
300 }
301
302 pub async fn enqueue_at(
309 &self,
310 when: DateTime<Utc>,
311 job_type: impl Into<String>,
312 data: JobData,
313 ) -> QueueResult<JobId> {
314 let job = Job::new(&self.config.queue_name, job_type, data).schedule_at(when);
315 self.enqueue_job(job).await
316 }
317
318 pub async fn enqueue_job(&self, job: Job) -> QueueResult<JobId> {
320 if self.config.max_size > 0 {
325 let size = self.backlog_size().await?;
326 if size >= self.config.max_size {
327 return Err(QueueError::QueueFull);
328 }
329 }
330
331 let job_id = job.id;
332 let mut conn = self.connection.clone();
333
334 let job_json =
336 serde_json::to_string(&job).map_err(|e| QueueError::Serialization(e.to_string()))?;
337
338 let job_key = self.config.key(&format!("job:{}", job_id));
340 let _: () = conn
341 .set_ex(
342 &job_key,
343 job_json,
344 self.config.body_ttl_secs(wait_until(job.scheduled_at)),
345 )
346 .await?;
347
348 if job.is_ready() {
350 let queue_key = self.priority_queue_key(job.priority);
351 let score = -(job.priority as i64); let _: () = conn.zadd(&queue_key, job_id.to_string(), score).await?;
353 } else {
354 let delayed_key = self.config.key("delayed");
356 let score = job.scheduled_at.unwrap().timestamp();
357 let _: () = conn.zadd(&delayed_key, job_id.to_string(), score).await?;
358 }
359
360 Ok(job_id)
361 }
362
363 pub async fn dequeue(&self) -> QueueResult<Option<Job>> {
365 self.move_delayed_jobs().await?;
366
367 let mut conn = self.connection.clone();
368
369 let script = redis::Script::new(DEQUEUE_POP_SCRIPT);
377 let popped: Option<(String, String)> = script
378 .key(self.priority_queue_key(JobPriority::Critical))
379 .key(self.priority_queue_key(JobPriority::High))
380 .key(self.priority_queue_key(JobPriority::Normal))
381 .key(self.priority_queue_key(JobPriority::Low))
382 .arg(&self.config.key_prefix)
383 .arg(Utc::now().timestamp())
384 .invoke_async(&mut conn)
385 .await?;
386
387 let Some((job_id_str, job_json)) = popped else {
388 return Ok(None);
389 };
390 let job_id = job_id_str
393 .parse::<JobId>()
394 .map_err(|e| QueueError::Deserialization(e.to_string()))?;
395 let mut job: Job = serde_json::from_str(&job_json)
396 .map_err(|e| QueueError::Deserialization(e.to_string()))?;
397
398 job.start_processing();
399 let job_key = self.config.key(&format!("job:{}", job_id));
400 let updated_json =
401 serde_json::to_string(&job).map_err(|e| QueueError::Serialization(e.to_string()))?;
402
403 let _: () = conn
408 .set_ex(&job_key, updated_json, self.config.retention_time.as_secs())
409 .await?;
410
411 Ok(Some(job))
412 }
413
414 pub async fn complete(&self, job_id: JobId) -> QueueResult<()> {
416 if let Some(mut job) = self.get_job(job_id).await? {
421 job.complete();
422 let job_key = self.config.key(&format!("job:{}", job_id));
423 let processing_key = self.config.key("processing");
424 let job_json = serde_json::to_string(&job)
425 .map_err(|e| QueueError::Serialization(e.to_string()))?;
426
427 let mut conn = self.connection.clone();
428 let _: () = redis::pipe()
429 .set_ex(&job_key, job_json, self.config.retention_time.as_secs())
430 .ignore()
431 .zrem(&processing_key, job_id.to_string())
432 .ignore()
433 .query_async(&mut conn)
434 .await?;
435 } else {
436 self.remove_from_processing(job_id).await?;
437 }
438 Ok(())
439 }
440
441 pub async fn fail(&self, job_id: JobId, error: String) -> QueueResult<()> {
443 if let Some(mut job) = self.get_job(job_id).await? {
447 job.fail(error);
448
449 let job_key = self.config.key(&format!("job:{}", job_id));
450 let processing_key = self.config.key("processing");
451 let job_json = serde_json::to_string(&job)
452 .map_err(|e| QueueError::Serialization(e.to_string()))?;
453 let mut conn = self.connection.clone();
454
455 if job.status.state == JobState::Failed && job.can_retry() {
456 let retry_at = Utc::now() + job.backoff_delay();
458 job.scheduled_at = Some(retry_at);
459 let job_json = serde_json::to_string(&job)
460 .map_err(|e| QueueError::Serialization(e.to_string()))?;
461
462 let delayed_key = self.config.key("delayed");
466 let _: () = redis::pipe()
467 .set_ex(
468 &job_key,
469 job_json,
470 self.config.body_ttl_secs(wait_until(Some(retry_at))),
471 )
472 .ignore()
473 .zadd(&delayed_key, job_id.to_string(), retry_at.timestamp())
474 .ignore()
475 .zrem(&processing_key, job_id.to_string())
476 .ignore()
477 .query_async(&mut conn)
478 .await?;
479 } else {
480 let dead_key = self.config.key("dead");
482 let _: () = redis::pipe()
483 .set_ex(&job_key, job_json, self.config.retention_time.as_secs())
484 .ignore()
485 .zadd(&dead_key, job_id.to_string(), Utc::now().timestamp())
486 .ignore()
487 .zrem(&processing_key, job_id.to_string())
488 .ignore()
489 .query_async(&mut conn)
490 .await?;
491 }
492 } else {
493 self.remove_from_processing(job_id).await?;
494 }
495 Ok(())
496 }
497
498 pub async fn requeue(&self, job: &Job) -> QueueResult<()> {
510 let mut job = job.clone();
511
512 job.status = JobStatus::pending();
514 job.started_at = None;
515 job.attempts = job.attempts.saturating_sub(1);
516 self.save_job(&job).await?;
517
518 let mut conn = self.connection.clone();
519 let queue_key = self.priority_queue_key(job.priority);
520 let score = -(job.priority as i64);
521 let _: () = conn.zadd(&queue_key, job.id.to_string(), score).await?;
522
523 self.remove_from_processing(job.id).await?;
524 Ok(())
525 }
526
527 pub async fn get_job(&self, job_id: JobId) -> QueueResult<Option<Job>> {
529 let mut conn = self.connection.clone();
530 let job_key = self.config.key(&format!("job:{}", job_id));
531
532 let job_json: Option<String> = conn.get(&job_key).await?;
533
534 if let Some(json) = job_json {
535 let job: Job = serde_json::from_str(&json)
536 .map_err(|e| QueueError::Deserialization(e.to_string()))?;
537 Ok(Some(job))
538 } else {
539 Ok(None)
540 }
541 }
542
543 async fn save_job(&self, job: &Job) -> QueueResult<()> {
545 let mut conn = self.connection.clone();
546 let job_key = self.config.key(&format!("job:{}", job.id));
547 let job_json =
548 serde_json::to_string(job).map_err(|e| QueueError::Serialization(e.to_string()))?;
549
550 let _: () = conn
551 .set_ex(&job_key, job_json, self.config.retention_time.as_secs())
552 .await?;
553 Ok(())
554 }
555
556 pub async fn size(&self) -> QueueResult<usize> {
558 let mut conn = self.connection.clone();
559
560 let mut pipe = redis::pipe();
563 for priority in [
564 JobPriority::Critical,
565 JobPriority::High,
566 JobPriority::Normal,
567 JobPriority::Low,
568 ] {
569 pipe.zcard(self.priority_queue_key(priority));
570 }
571
572 let counts: Vec<usize> = pipe.query_async(&mut conn).await?;
573 Ok(counts.iter().sum())
574 }
575
576 pub async fn backlog_size(&self) -> QueueResult<usize> {
585 let mut conn = self.connection.clone();
586
587 let mut pipe = redis::pipe();
588 for priority in [
589 JobPriority::Critical,
590 JobPriority::High,
591 JobPriority::Normal,
592 JobPriority::Low,
593 ] {
594 pipe.zcard(self.priority_queue_key(priority));
595 }
596 pipe.zcard(self.config.key("delayed"));
597 pipe.zcard(self.config.key("processing"));
598
599 let counts: Vec<usize> = pipe.query_async(&mut conn).await?;
600 Ok(counts.iter().sum())
601 }
602
603 pub async fn processing_len(&self) -> QueueResult<usize> {
605 let mut conn = self.connection.clone();
606 let processing_key = self.config.key("processing");
607 let count: usize = conn.zcard(&processing_key).await?;
608 Ok(count)
609 }
610
611 async fn move_delayed_jobs(&self) -> QueueResult<()> {
619 let mut conn = self.connection.clone();
620 let delayed_key = self.config.key("delayed");
621 let now = Utc::now().timestamp();
622
623 let earliest: Vec<(String, i64)> = conn.zrange_withscores(&delayed_key, 0, 0).await?;
629 match earliest.first() {
630 Some((_, score)) if *score <= now => {}
631 _ => return Ok(()),
632 }
633
634 let script = redis::Script::new(&MOVE_DELAYED_SCRIPT);
635 let (_promoted, dropped): (i64, i64) = script
636 .key(&delayed_key)
637 .arg(&self.config.key_prefix)
638 .arg(now)
639 .invoke_async(&mut conn)
640 .await?;
641
642 if dropped > 0 {
643 warn!(
644 "Dropped {} delayed job(s) from queue '{}' whose body expired before their scheduled time",
645 dropped, self.config.queue_name
646 );
647 }
648
649 Ok(())
650 }
651
652 pub async fn reclaim_stale(&self, visibility_timeout: Duration) -> QueueResult<usize> {
674 let mut conn = self.connection.clone();
675 let processing_key = self.config.key("processing");
676 let cutoff = Utc::now().timestamp() - visibility_timeout.as_secs() as i64;
677
678 let script = redis::Script::new(&RECLAIM_STALE_SCRIPT);
679 let (reclaimed, dropped): (i64, i64) = script
680 .key(&processing_key)
681 .arg(&self.config.key_prefix)
682 .arg(cutoff)
683 .invoke_async(&mut conn)
684 .await?;
685
686 if reclaimed > 0 {
687 warn!(
688 "Reclaimed {} stale in-flight job(s) on queue '{}' (claim older than {:?})",
689 reclaimed, self.config.queue_name, visibility_timeout
690 );
691 }
692 if dropped > 0 {
693 warn!(
694 "Dropped {} stale in-flight job(s) on queue '{}' whose body had already expired",
695 dropped, self.config.queue_name
696 );
697 }
698
699 Ok(reclaimed as usize)
700 }
701
702 async fn remove_from_processing(&self, job_id: JobId) -> QueueResult<()> {
704 let mut conn = self.connection.clone();
705 let processing_key = self.config.key("processing");
706 let _: () = conn.zrem(&processing_key, job_id.to_string()).await?;
707 Ok(())
708 }
709
710 fn priority_queue_key(&self, priority: JobPriority) -> String {
712 self.config
713 .key(&format!("pending:{:?}", priority).to_lowercase())
714 }
715
716 pub async fn clear(&self) -> QueueResult<()> {
718 let mut conn = self.connection.clone();
719 let pattern = format!("{}:*", self.config.key_prefix);
720
721 let mut cursor: u64 = 0;
724 loop {
725 let (next, keys): (u64, Vec<String>) = redis::cmd("SCAN")
726 .arg(cursor)
727 .arg("MATCH")
728 .arg(&pattern)
729 .arg("COUNT")
730 .arg(SCAN_COUNT)
731 .query_async(&mut conn)
732 .await?;
733
734 if !keys.is_empty() {
735 let _: () = redis::cmd("UNLINK")
736 .arg(&keys)
737 .query_async(&mut conn)
738 .await?;
739 }
740
741 cursor = next;
742 if cursor == 0 {
743 break;
744 }
745 }
746
747 Ok(())
748 }
749}
750
751#[cfg(test)]
752mod tests {
753 use super::*;
754
755 #[test]
756 fn test_queue_config() {
757 let config = QueueConfig::new("redis://localhost:6379", "test");
758 assert_eq!(config.queue_name, "test");
759 assert!(config.key_prefix.contains("test"));
760 }
761
762 #[test]
767 fn test_lua_priority_mapping_matches_rust() {
768 for priority in [
769 JobPriority::Low,
770 JobPriority::Normal,
771 JobPriority::High,
772 JobPriority::Critical,
773 ] {
774 let rust_score = -(priority as i64);
776 let variant = format!("{priority:?}");
778 let name = variant.to_lowercase();
780
781 let expected = format!("'{variant}' then pname = '{name}'; pscore = {rust_score}");
782 assert!(
783 PRIORITY_LUA.contains(&expected),
784 "script missing mapping for {priority:?}: expected `{expected}`"
785 );
786 }
787 }
788
789 #[test]
792 fn test_refiling_scripts_share_priority_helper() {
793 for script in [&*MOVE_DELAYED_SCRIPT, &*RECLAIM_STALE_SCRIPT] {
794 assert!(script.contains("local function priority_of"));
795 assert!(script.contains("priority_of(job_json)"));
796 }
797 }
798
799 #[test]
803 fn test_move_delayed_drops_bodyless_ids() {
804 assert!(MOVE_DELAYED_BODY.contains("elseif redis.call('ZREM', delayed_key, job_id) == 1"));
805 }
806
807 #[test]
809 fn test_reclaim_script_refiles_to_pending() {
810 assert!(RECLAIM_STALE_BODY.contains("ZRANGEBYSCORE"));
811 assert!(RECLAIM_STALE_BODY.contains("ZADD', prefix .. ':pending:' .. pname"));
812 assert!(RECLAIM_STALE_BODY.contains("ZREM', processing_key"));
813 }
814
815 #[test]
819 fn test_dequeue_script_claims_atomically() {
820 assert!(DEQUEUE_POP_SCRIPT.contains("ZADD', prefix .. ':processing', now, job_id"));
821 }
822
823 #[test]
824 fn test_priority_queue_key_layout() {
825 let config = QueueConfig::new("redis://localhost:6379", "jobs");
827 for (priority, name) in [
828 (JobPriority::Low, "low"),
829 (JobPriority::Normal, "normal"),
830 (JobPriority::High, "high"),
831 (JobPriority::Critical, "critical"),
832 ] {
833 let key = config.key(&format!("pending:{priority:?}").to_lowercase());
834 assert_eq!(key, format!("{}:pending:{}", config.key_prefix, name));
835 }
836 }
837
838 #[test]
839 fn test_dequeue_script_scans_all_keys() {
840 assert!(DEQUEUE_POP_SCRIPT.contains("for i = 1, #KEYS do"));
842 assert!(DEQUEUE_POP_SCRIPT.contains("ZPOPMIN"));
843 }
844
845 #[tokio::test]
847 #[ignore = "requires a running Redis instance"]
848 async fn test_move_delayed_promotes_due_jobs() {
849 use crate::job::Job;
850
851 let queue = Queue::new("redis://localhost:6379", "test_move_delayed")
852 .await
853 .unwrap();
854 queue.clear().await.unwrap();
855
856 let past = Utc::now() - chrono::Duration::seconds(30);
858 let job = Job::new("test_move_delayed", "task", serde_json::json!({}))
859 .with_priority(JobPriority::High)
860 .schedule_at(past);
861 queue.enqueue_job(job).await.unwrap();
862
863 assert_eq!(queue.size().await.unwrap(), 0);
865
866 let dequeued = queue.dequeue().await.unwrap();
868 assert!(dequeued.is_some());
869 assert_eq!(dequeued.unwrap().priority, JobPriority::High);
870
871 queue.clear().await.unwrap();
872 }
873
874 #[test]
878 fn test_body_ttl_covers_long_horizon_schedule() {
879 let config = QueueConfig::new("redis://localhost:6379", "test"); let retention = config.retention_time.as_secs();
881
882 assert_eq!(config.body_ttl_secs(Duration::ZERO), retention);
884
885 let wait = Duration::from_secs(40 * 86_400);
887 let ttl = config.body_ttl_secs(wait);
888 assert!(
889 ttl > wait.as_secs(),
890 "body would expire {} s before its scheduled time",
891 wait.as_secs() - ttl
892 );
893 assert_eq!(ttl, wait.as_secs() + retention);
894 }
895
896 #[test]
897 fn test_wait_until_clamps_past_times() {
898 let past = Utc::now() - chrono::Duration::days(3);
901 assert_eq!(wait_until(Some(past)), Duration::ZERO);
902 assert_eq!(wait_until(None), Duration::ZERO);
903 assert!(wait_until(Some(Utc::now() + chrono::Duration::hours(2))) > Duration::from_secs(0));
904 }
905
906 #[tokio::test]
908 #[ignore = "requires a running Redis instance"]
909 async fn test_long_horizon_scheduled_job_body_survives() {
910 use crate::job::Job;
911
912 let queue = Queue::new("redis://localhost:6379", "test_long_horizon")
913 .await
914 .unwrap();
915 queue.clear().await.unwrap();
916
917 let far_future = Utc::now() + chrono::Duration::days(40);
919 let job =
920 Job::new("test_long_horizon", "task", serde_json::json!({})).schedule_at(far_future);
921 let job_id = queue.enqueue_job(job).await.unwrap();
922
923 assert!(queue.get_job(job_id).await.unwrap().is_some());
926
927 let mut conn = queue.connection.clone();
928 let ttl: i64 = conn
929 .ttl(queue.config.key(&format!("job:{job_id}")))
930 .await
931 .unwrap();
932 assert!(
933 ttl > 40 * 86_400,
934 "body TTL {ttl}s expires before the job's scheduled time"
935 );
936
937 queue.clear().await.unwrap();
938 }
939
940 #[tokio::test]
942 #[ignore = "requires a running Redis instance"]
943 async fn test_reclaim_stale_returns_orphaned_job() {
944 use crate::job::Job;
945
946 let queue = Queue::new("redis://localhost:6379", "test_reclaim")
947 .await
948 .unwrap();
949 queue.clear().await.unwrap();
950
951 let job = Job::new("test_reclaim", "task", serde_json::json!({}))
952 .with_priority(JobPriority::High);
953 let job_id = queue.enqueue_job(job).await.unwrap();
954
955 let dequeued = queue.dequeue().await.unwrap().unwrap();
958 assert_eq!(dequeued.id, job_id);
959 assert_eq!(queue.size().await.unwrap(), 0);
960 assert_eq!(queue.processing_len().await.unwrap(), 1);
961
962 assert_eq!(
964 queue
965 .reclaim_stale(Duration::from_secs(3600))
966 .await
967 .unwrap(),
968 0
969 );
970 assert_eq!(queue.processing_len().await.unwrap(), 1);
971
972 assert_eq!(queue.reclaim_stale(Duration::ZERO).await.unwrap(), 1);
975 assert_eq!(queue.processing_len().await.unwrap(), 0);
976 assert_eq!(queue.size().await.unwrap(), 1);
977
978 let again = queue.dequeue().await.unwrap().unwrap();
979 assert_eq!(again.id, job_id);
980 assert_eq!(again.priority, JobPriority::High);
981
982 queue.clear().await.unwrap();
983 }
984
985 #[test]
986 fn test_priority_queue_key() {
987 let config = QueueConfig::new("redis://localhost:6379", "test");
988 assert!(config.key("pending:high").contains("high"));
989 }
990
991 #[test]
992 fn test_queue_config_with_custom_prefix() {
993 let config = QueueConfig::new("redis://localhost:6379", "myqueue").with_key_prefix("app");
994 assert!(config.key_prefix.contains("app"));
995 }
996
997 #[test]
998 fn test_queue_config_default_retention() {
999 let config = QueueConfig::new("redis://localhost:6379", "test");
1000 assert_eq!(config.retention_time, Duration::from_secs(86400)); }
1002
1003 #[test]
1004 fn test_queue_config_custom_retention() {
1005 let retention = Duration::from_secs(3600);
1006 let config =
1007 QueueConfig::new("redis://localhost:6379", "test").with_retention_time(retention);
1008 assert_eq!(config.retention_time, retention);
1009 }
1010
1011 #[test]
1012 fn test_queue_config_default_max_size() {
1013 let config = QueueConfig::new("redis://localhost:6379", "test");
1014 assert_eq!(config.max_size, 0); }
1016
1017 #[test]
1018 fn test_queue_config_custom_max_size() {
1019 let config = QueueConfig::new("redis://localhost:6379", "test").with_max_size(1000);
1020 assert_eq!(config.max_size, 1000);
1021 }
1022
1023 #[test]
1024 fn test_queue_key_generation() {
1025 let config = QueueConfig::new("redis://localhost:6379", "jobs");
1026
1027 let pending_key = config.key("pending:normal");
1028 let processing_key = config.key("processing");
1029 let completed_key = config.key("completed");
1030
1031 assert!(pending_key.contains("jobs"));
1032 assert!(processing_key.contains("jobs"));
1033 assert!(completed_key.contains("jobs"));
1034 }
1035
1036 #[test]
1037 fn test_queue_config_clone() {
1038 let config1 = QueueConfig::new("redis://localhost:6379", "test");
1039 let config2 = config1.clone();
1040
1041 assert_eq!(config1.queue_name, config2.queue_name);
1042 assert_eq!(config1.redis_url, config2.redis_url);
1043 }
1044
1045 #[test]
1046 fn test_queue_config_different_queues() {
1047 let config1 = QueueConfig::new("redis://localhost:6379", "queue1");
1048 let config2 = QueueConfig::new("redis://localhost:6379", "queue2");
1049
1050 assert_ne!(config1.key_prefix, config2.key_prefix);
1051 }
1052
1053 #[test]
1054 fn test_queue_config_key_consistency() {
1055 let config = QueueConfig::new("redis://localhost:6379", "test");
1056
1057 let key1 = config.key("pending");
1058 let key2 = config.key("pending");
1059
1060 assert_eq!(key1, key2);
1061 }
1062
1063 #[test]
1064 fn test_queue_config_builder_pattern() {
1065 let config = QueueConfig::new("redis://localhost:6379", "test")
1066 .with_key_prefix("app")
1067 .with_retention_time(Duration::from_secs(7200))
1068 .with_max_size(500);
1069
1070 assert!(config.key_prefix.contains("app"));
1071 assert_eq!(config.retention_time, Duration::from_secs(7200));
1072 assert_eq!(config.max_size, 500);
1073 }
1074
1075 #[test]
1076 fn test_queue_config_redis_url() {
1077 let url = "redis://user:pass@host:6380/2";
1078 let config = QueueConfig::new(url, "test");
1079 assert_eq!(config.redis_url, url);
1080 }
1081
1082 #[test]
1083 fn test_queue_config_key_with_empty_suffix() {
1084 let config = QueueConfig::new("redis://localhost:6379", "test");
1085 let key = config.key("");
1086 assert!(key.contains("test"));
1087 }
1088
1089 #[test]
1090 fn test_queue_config_key_with_special_characters() {
1091 let config = QueueConfig::new("redis://localhost:6379", "test");
1092 let key = config.key("pending:high:priority");
1093 assert!(key.contains("pending:high:priority"));
1094 }
1095
1096 #[test]
1097 fn test_queue_config_multiple_prefixes() {
1098 let config1 =
1099 QueueConfig::new("redis://localhost:6379", "app1").with_key_prefix("production");
1100 let config2 =
1101 QueueConfig::new("redis://localhost:6379", "app2").with_key_prefix("development");
1102
1103 let key1 = config1.key("jobs");
1104 let key2 = config2.key("jobs");
1105
1106 assert_ne!(key1, key2);
1107 }
1108
1109 #[test]
1110 fn test_queue_config_unlimited_max_size() {
1111 let config = QueueConfig::new("redis://localhost:6379", "test").with_max_size(0);
1112 assert_eq!(config.max_size, 0);
1113 }
1114
1115 #[test]
1116 fn test_queue_config_large_retention() {
1117 let week = Duration::from_secs(7 * 24 * 3600);
1118 let config = QueueConfig::new("redis://localhost:6379", "test").with_retention_time(week);
1119 assert_eq!(config.retention_time, week);
1120 }
1121}