1mod dashboard;
63mod worker;
64
65use std::collections::HashMap;
66use std::future::Future;
67use std::pin::Pin;
68use std::sync::Arc;
69use std::time::Duration;
70
71use serde::de::DeserializeOwned;
72use serde::{Deserialize, Serialize};
73use tokio::sync::Notify;
74
75pub use dashboard::{Dashboard, GATE as DASHBOARD_GATE, QueueCounts, QueueStats};
76pub use worker::Worker;
77
78use crate::db::{Db, Migration, Transaction};
79use crate::{AppState, Error, Result};
80
81pub(crate) const MIGRATIONS: &[Migration] = &[
82 crate::db::framework_migration!("queue", "00010101000100_create_jobs_table"),
83 crate::db::framework_migration!("queue", "00010101000110_add_chains_and_batches_to_jobs"),
84 crate::db::framework_migration!("queue", "00010101000120_add_callback_of_to_jobs"),
85];
86
87const SEALED: &str = "enc:";
89const UNIQUE_PREFIX: &str = "renox:unique:";
91
92pub trait Job: Serialize + DeserializeOwned + Send + Sync + 'static {
95 const NAME: &'static str;
97 const QUEUE: &'static str = "default";
99 const MAX_ATTEMPTS: u32 = 3;
101 const TIMEOUT: Duration = Duration::from_secs(60);
103 const ENCRYPTED: bool = false;
106 const UNIQUE_FOR: Option<Duration> = None;
110
111 fn backoff(attempt: u32) -> Duration {
113 Duration::from_secs(10 * u64::from(attempt))
114 }
115
116 fn unique_id(&self) -> String {
119 String::new()
120 }
121
122 fn middleware(&self) -> Vec<Middleware> {
125 Vec::new()
126 }
127
128 fn handle(self, ctx: JobContext) -> impl Future<Output = Result> + Send;
130
131 fn failed(self, ctx: JobContext, error: String) -> impl Future<Output = ()> + Send {
135 let _ = (ctx, error);
136 async {}
137 }
138}
139
140#[derive(Debug, Clone)]
167pub struct Middleware(MiddlewareKind);
168
169#[derive(Debug, Clone)]
170enum MiddlewareKind {
171 WithoutOverlapping {
172 key: String,
173 release_after: Duration,
174 },
175 RateLimited {
176 key: String,
177 max: u32,
178 per: Duration,
179 },
180}
181
182impl Middleware {
183 pub fn without_overlapping(key: impl Into<String>) -> Self {
186 Self(MiddlewareKind::WithoutOverlapping {
187 key: key.into(),
188 release_after: Duration::from_secs(5),
189 })
190 }
191
192 pub fn rate_limited(key: impl Into<String>, max: u32, per: Duration) -> Self {
195 Self(MiddlewareKind::RateLimited {
196 key: key.into(),
197 max,
198 per: per.max(Duration::from_secs(1)),
199 })
200 }
201
202 pub fn release_after(mut self, wait: Duration) -> Self {
204 if let MiddlewareKind::WithoutOverlapping { release_after, .. } = &mut self.0 {
205 *release_after = wait;
206 }
207 self
208 }
209}
210
211#[non_exhaustive]
213pub struct JobContext {
214 pub state: AppState,
216 pub attempt: u32,
218 pub id: i64,
220 pub batch_id: Option<i64>,
223}
224
225type RunFn =
226 Arc<dyn Fn(String, JobContext) -> Pin<Box<dyn Future<Output = Result> + Send>> + Send + Sync>;
227type FailedFn = Arc<
228 dyn Fn(String, JobContext, String) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + Sync,
229>;
230
231#[derive(Clone)]
233pub(crate) struct JobHandler {
234 pub(crate) run: RunFn,
235 pub(crate) failed: FailedFn,
236 pub(crate) backoff: fn(u32) -> Duration,
237 pub(crate) timeout: Duration,
238 pub(crate) unique_key: fn(&str) -> Option<String>,
240 pub(crate) middleware: fn(&str) -> Vec<Middleware>,
241}
242
243pub(crate) fn handler<J: Job>() -> JobHandler {
245 JobHandler {
246 run: Arc::new(|payload, ctx| {
247 Box::pin(async move {
248 let job: J = serde_json::from_str(&payload).map_err(crate::Error::permanent)?;
250 job.handle(ctx).await
251 })
252 }),
253 failed: Arc::new(|payload, ctx, error| {
254 Box::pin(async move {
255 if let Ok(job) = serde_json::from_str::<J>(&payload) {
256 job.failed(ctx, error).await;
257 }
258 })
259 }),
260 backoff: J::backoff,
261 timeout: J::TIMEOUT,
262 unique_key: |payload| {
263 J::UNIQUE_FOR?;
264 let job: J = serde_json::from_str(payload).ok()?;
265 Some(unique_key::<J>(&job))
266 },
267 middleware: |payload| {
268 serde_json::from_str::<J>(payload)
269 .map(|job| job.middleware())
270 .unwrap_or_default()
271 },
272 }
273}
274
275fn unique_key<J: Job>(job: &J) -> String {
276 format!("{UNIQUE_PREFIX}{}:{}", J::NAME, job.unique_id())
277}
278
279pub(crate) type Handlers = Arc<HashMap<&'static str, JobHandler>>;
280
281#[derive(Debug, Clone, Serialize, Deserialize)]
283pub(crate) struct Encoded {
284 queue: String,
285 job: String,
286 payload: String,
287 max_attempts: u32,
288}
289
290#[derive(Clone)]
292pub struct Queue {
293 db: Db,
294 key: cookie::Key,
295 wake: Arc<Notify>,
296}
297
298pub(crate) fn unix_now() -> i64 {
299 crate::clock::unix_secs()
300}
301
302impl Queue {
303 pub(crate) fn new(db: Db, key: cookie::Key) -> Self {
304 Self {
305 db,
306 key,
307 wake: Arc::new(Notify::new()),
308 }
309 }
310
311 pub(crate) fn wake(&self) -> &Notify {
312 &self.wake
313 }
314
315 pub(crate) fn wake_workers(&self) {
316 self.wake.notify_waiters();
317 }
318
319 fn encode<J: Job>(&self, job: &J, queue: &str) -> Result<Encoded> {
320 let json = serde_json::to_string(job)?;
321 let payload = if J::ENCRYPTED {
322 format!("{SEALED}{}", crate::crypto::seal(&self.key, &json))
323 } else {
324 json
325 };
326 Ok(Encoded {
327 queue: queue.to_owned(),
328 job: J::NAME.to_owned(),
329 payload,
330 max_attempts: J::MAX_ATTEMPTS.max(1),
331 })
332 }
333
334 pub(crate) fn open(&self, payload: &str) -> Result<String> {
336 match payload.strip_prefix(SEALED) {
337 Some(sealed) => Ok(crate::crypto::open(&self.key, sealed).map_err(Error::permanent)?),
338 None => Ok(payload.to_owned()),
339 }
340 }
341
342 pub async fn dispatch<J: Job>(&self, job: J) -> Result<i64> {
344 self.dispatch_after(job, Duration::ZERO).await
345 }
346
347 pub async fn dispatch_after<J: Job>(&self, job: J, delay: Duration) -> Result<i64> {
349 self.push(job, J::QUEUE, delay).await
350 }
351
352 pub async fn dispatch_on<J: Job>(&self, queue: &str, job: J) -> Result<i64> {
354 self.push(job, queue, Duration::ZERO).await
355 }
356
357 async fn push<J: Job>(&self, job: J, queue: &str, delay: Duration) -> Result<i64> {
358 let encoded = self.encode(&job, queue)?;
359 let id = match J::UNIQUE_FOR {
360 None => insert(&self.db, &encoded, delay, None, None, None).await?,
361 Some(ttl) => {
362 let mut tx = self.db.begin().await?;
363 let id =
364 insert_unique(&mut tx, &unique_key::<J>(&job), ttl, &encoded, delay).await?;
365 tx.commit().await?;
366 id
367 }
368 };
369 self.wake.notify_waiters();
370 Ok(id)
371 }
372
373 pub async fn dispatch_in<J: Job>(&self, tx: &mut Transaction, job: J) -> Result<i64> {
393 let encoded = self.encode(&job, J::QUEUE)?;
394 match J::UNIQUE_FOR {
395 None => insert(tx, &encoded, Duration::ZERO, None, None, None).await,
396 Some(ttl) => {
397 insert_unique(tx, &unique_key::<J>(&job), ttl, &encoded, Duration::ZERO).await
398 }
399 }
400 }
401
402 pub fn chain(&self) -> Chain {
406 Chain {
407 queue: self.clone(),
408 jobs: Vec::new(),
409 error: None,
410 }
411 }
412
413 pub fn batch(&self, name: &str) -> Batch {
415 Batch {
416 queue: self.clone(),
417 name: name.to_owned(),
418 jobs: Vec::new(),
419 then: None,
420 catch: None,
421 finally: None,
422 allow_failures: false,
423 error: None,
424 }
425 }
426
427 pub async fn batch_status(&self, id: i64) -> Result<Option<BatchStatus>> {
429 let row = crate::db::sql(
430 "SELECT id, name, total, pending, failed, cancelled_at, finished_at, created_at \
431 FROM job_batches WHERE id = ?",
432 )
433 .bind(id)
434 .fetch_optional(&self.db)
435 .await?;
436 let Some(row) = row else {
437 return Ok(None);
438 };
439 Ok(Some(BatchStatus {
440 id: row.try_get("id")?,
441 name: row.try_get("name")?,
442 total: row.try_get("total")?,
443 pending: row.try_get("pending")?,
444 failed: row.try_get("failed")?,
445 cancelled: row.try_get::<Option<i64>>("cancelled_at")?.is_some(),
446 finished: row.try_get::<Option<i64>>("finished_at")?.is_some(),
447 created_at: crate::db::from_unix(row.try_get("created_at")?),
448 }))
449 }
450
451 pub async fn cancel_batch(&self, id: i64) -> Result<bool> {
453 Ok(crate::db::sql(
454 "UPDATE job_batches SET cancelled_at = ? WHERE id = ? AND cancelled_at IS NULL",
455 )
456 .bind(unix_now())
457 .bind(id)
458 .execute(&self.db)
459 .await?
460 == 1)
461 }
462
463 pub async fn pending(&self) -> Result<i64> {
465 Ok(crate::db::sql("SELECT COUNT(*) FROM jobs")
466 .scalar(&self.db)
467 .await?)
468 }
469
470 pub async fn failed(&self) -> Result<Vec<FailedJob>> {
472 let rows = crate::db::sql(
473 "SELECT id, queue, job, payload, error, failed_at FROM failed_jobs ORDER BY id",
474 )
475 .fetch_all(&self.db)
476 .await?;
477 rows.iter()
478 .map(|r| {
479 Ok(FailedJob {
480 id: r.try_get("id")?,
481 queue: r.try_get("queue")?,
482 job: r.try_get("job")?,
483 payload: r.try_get("payload")?,
484 error: r.try_get("error")?,
485 failed_at: crate::db::from_unix(r.try_get("failed_at")?),
486 })
487 })
488 .collect()
489 }
490
491 pub async fn retry(&self, id: i64) -> Result<bool> {
494 Ok(self.move_failed(Some(id)).await? == 1)
495 }
496
497 pub async fn retry_all(&self) -> Result<u64> {
499 self.move_failed(None).await
500 }
501
502 async fn move_failed(&self, id: Option<i64>) -> Result<u64> {
503 let mut tx = self.db.begin().await?;
504 let filter = if id.is_some() { " WHERE id = ?" } else { "" };
505 let mut batches = crate::db::sql(format!(
506 "UPDATE job_batches SET failed = failed - 1, pending = pending + 1, finished_at = NULL \
507 WHERE id IN (SELECT batch_id FROM failed_jobs{filter})"
508 ));
509 if let Some(id) = id {
510 batches = batches.bind(id);
511 }
512 batches.execute(&mut tx).await?;
513 let insert = format!(
514 "INSERT INTO jobs (queue, job, payload, max_attempts, available_at, created_at, chain, \
515 batch_id, callback_of) SELECT queue, job, payload, max_attempts, ?, ?, chain, batch_id, \
516 callback_of FROM failed_jobs{filter}"
517 );
518 let mut query = crate::db::sql(insert).bind(unix_now()).bind(unix_now());
519 if let Some(id) = id {
520 query = query.bind(id);
521 }
522 let moved = query.execute(&mut tx).await?;
523 let mut delete = crate::db::sql(format!("DELETE FROM failed_jobs{filter}"));
524 if let Some(id) = id {
525 delete = delete.bind(id);
526 }
527 delete.execute(&mut tx).await?;
528 tx.commit().await?;
529 self.wake.notify_waiters();
530 Ok(moved)
531 }
532
533 pub async fn forget_failed(&self, id: i64) -> Result<bool> {
535 Ok(crate::db::sql("DELETE FROM failed_jobs WHERE id = ?")
536 .bind(id)
537 .execute(&self.db)
538 .await?
539 == 1)
540 }
541
542 pub async fn prune_failed(&self, age: Duration) -> Result<u64> {
544 Ok(
545 crate::db::sql("DELETE FROM failed_jobs WHERE failed_at < ?")
546 .bind(unix_now() - age.as_secs() as i64)
547 .execute(&self.db)
548 .await?,
549 )
550 }
551
552 pub async fn prune_batches(&self, age: Duration) -> Result<u64> {
554 let before = unix_now() - age.as_secs() as i64;
555 Ok(crate::db::sql(
556 "DELETE FROM job_batches WHERE (finished_at IS NOT NULL AND finished_at < ?) \
557 OR (cancelled_at IS NOT NULL AND cancelled_at < ? AND pending = 0)",
558 )
559 .bind(before)
560 .bind(before)
561 .execute(&self.db)
562 .await?)
563 }
564
565 pub async fn flush_failed(&self) -> Result<u64> {
567 Ok(crate::db::sql("DELETE FROM failed_jobs")
568 .execute(&self.db)
569 .await?)
570 }
571}
572
573async fn insert<'c>(
575 db: impl crate::db::Executor<'c>,
576 job: &Encoded,
577 delay: Duration,
578 chain: Option<String>,
579 batch_id: Option<i64>,
580 callback_of: Option<i64>,
581) -> Result<i64> {
582 let now = unix_now();
583 let id: i64 = crate::db::sql(
584 "INSERT INTO jobs (queue, job, payload, max_attempts, available_at, created_at, chain, \
585 batch_id, callback_of) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) RETURNING id",
586 )
587 .bind(&job.queue)
588 .bind(&job.job)
589 .bind(&job.payload)
590 .bind(i64::from(job.max_attempts))
591 .bind(now + delay.as_secs() as i64)
592 .bind(now)
593 .bind(chain)
594 .bind(batch_id)
595 .bind(callback_of)
596 .scalar(db)
597 .await?;
598 Ok(id)
599}
600
601async fn insert_unique(
605 tx: &mut Transaction,
606 key: &str,
607 ttl: Duration,
608 job: &Encoded,
609 delay: Duration,
610) -> Result<i64> {
611 let now = unix_now();
612 let claimed = crate::db::sql(
613 "INSERT INTO cache (key, value, expires_at) VALUES (?, '0', ?) \
614 ON CONFLICT (key) DO UPDATE SET value = '0', expires_at = excluded.expires_at \
615 WHERE cache.expires_at IS NOT NULL AND cache.expires_at <= ?",
616 )
617 .bind(key)
618 .bind(now + ttl.as_secs().max(1) as i64)
619 .bind(now)
620 .execute(&mut *tx)
621 .await?;
622 if claimed == 0 {
623 let holder: String = crate::db::sql("SELECT value FROM cache WHERE key = ?")
624 .bind(key)
625 .scalar(&mut *tx)
626 .await?;
627 return Ok(holder.parse().unwrap_or_default());
628 }
629 let id = insert(&mut *tx, job, delay, None, None, None).await?;
630 crate::db::sql("UPDATE cache SET value = ? WHERE key = ?")
631 .bind(id.to_string())
632 .bind(key)
633 .execute(&mut *tx)
634 .await?;
635 Ok(id)
636}
637
638#[must_use = "a chain does nothing until dispatched"]
640pub struct Chain {
641 queue: Queue,
642 jobs: Vec<Encoded>,
643 error: Option<Error>,
644}
645
646impl Chain {
647 pub fn then<J: Job>(mut self, job: J) -> Self {
649 match self.queue.encode(&job, J::QUEUE) {
650 Ok(encoded) => self.jobs.push(encoded),
651 Err(err) => {
652 self.error.get_or_insert(err);
653 }
654 }
655 self
656 }
657
658 pub async fn dispatch(self) -> Result<i64> {
660 if let Some(err) = self.error {
661 return Err(err);
662 }
663 let mut jobs = self.jobs.into_iter();
664 let Some(first) = jobs.next() else {
665 return Ok(0);
666 };
667 let rest: Vec<Encoded> = jobs.collect();
668 let chain = (!rest.is_empty())
669 .then(|| serde_json::to_string(&rest))
670 .transpose()?;
671 let id = insert(&self.queue.db, &first, Duration::ZERO, chain, None, None).await?;
672 self.queue.wake_workers();
673 Ok(id)
674 }
675}
676
677#[must_use = "a batch does nothing until dispatched"]
681pub struct Batch {
682 queue: Queue,
683 name: String,
684 jobs: Vec<Encoded>,
685 then: Option<Encoded>,
686 catch: Option<Encoded>,
687 finally: Option<Encoded>,
688 allow_failures: bool,
689 error: Option<Error>,
690}
691
692impl Batch {
693 fn encode<J: Job>(&mut self, job: J) -> Option<Encoded> {
694 match self.queue.encode(&job, J::QUEUE) {
695 Ok(encoded) => Some(encoded),
696 Err(err) => {
697 self.error.get_or_insert(err);
698 None
699 }
700 }
701 }
702
703 pub fn push<J: Job>(mut self, job: J) -> Self {
705 if let Some(job) = self.encode(job) {
706 self.jobs.push(job);
707 }
708 self
709 }
710
711 pub fn then<J: Job>(mut self, job: J) -> Self {
713 self.then = self.encode(job);
714 self
715 }
716
717 pub fn catch<J: Job>(mut self, job: J) -> Self {
719 self.catch = self.encode(job);
720 self
721 }
722
723 pub fn finally<J: Job>(mut self, job: J) -> Self {
725 self.finally = self.encode(job);
726 self
727 }
728
729 pub fn allow_failures(mut self) -> Self {
731 self.allow_failures = true;
732 self
733 }
734
735 pub async fn dispatch(self) -> Result<i64> {
737 if let Some(err) = self.error {
738 return Err(err);
739 }
740 let json = |job: &Option<Encoded>| job.as_ref().map(serde_json::to_string).transpose();
741 let (then, catch, finally) = (json(&self.then)?, json(&self.catch)?, json(&self.finally)?);
742 let total = self.jobs.len() as i64;
743 let mut tx = self.queue.db.begin().await?;
744 let id: i64 = crate::db::sql(
745 "INSERT INTO job_batches (name, total, pending, allow_failures, then_job, catch_job, \
746 finally_job, created_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?) RETURNING id",
747 )
748 .bind(&self.name)
749 .bind(total)
750 .bind(total)
751 .bind(self.allow_failures)
752 .bind(then)
753 .bind(catch)
754 .bind(finally)
755 .bind(unix_now())
756 .scalar(&mut tx)
757 .await?;
758 for job in &self.jobs {
759 insert(&mut tx, job, Duration::ZERO, None, Some(id), None).await?;
760 }
761 if self.jobs.is_empty() {
762 finish_batch(&mut tx, id, 0).await?;
763 }
764 tx.commit().await?;
765 self.queue.wake_workers();
766 Ok(id)
767 }
768}
769
770pub(crate) async fn batch_job_done(tx: &mut Transaction, batch_id: i64, failed: bool) -> Result {
773 let row: Option<(i64, i64, Option<String>)> = crate::db::sql(
774 "UPDATE job_batches SET pending = pending - 1, failed = failed + ?, \
775 cancelled_at = CASE WHEN ? AND NOT allow_failures THEN COALESCE(cancelled_at, ?) \
776 ELSE cancelled_at END \
777 WHERE id = ? RETURNING pending, failed, catch_job",
778 )
779 .bind(i64::from(failed))
780 .bind(failed)
781 .bind(unix_now())
782 .bind(batch_id)
783 .fetch_as(&mut *tx)
784 .await?
785 .into_iter()
786 .next();
787 let Some((pending, failures, catch)) = row else {
788 return Ok(()); };
790 if failed && failures == 1 {
791 queue_follow_up(tx, catch, batch_id).await?;
792 }
793 if pending <= 0 {
794 finish_batch(tx, batch_id, failures).await?;
795 }
796 Ok(())
797}
798
799async fn finish_batch(tx: &mut Transaction, batch_id: i64, failures: i64) -> Result {
800 let (then, finally, cancelled): (Option<String>, Option<String>, Option<i64>) = crate::db::sql(
801 "UPDATE job_batches SET finished_at = ? WHERE id = ? \
802 RETURNING then_job, finally_job, cancelled_at",
803 )
804 .bind(unix_now())
805 .bind(batch_id)
806 .fetch_as(&mut *tx)
807 .await?
808 .into_iter()
809 .next()
810 .unwrap_or_default();
811 if failures == 0 && cancelled.is_none() {
812 queue_follow_up(tx, then, batch_id).await?;
813 }
814 queue_follow_up(tx, finally, batch_id).await
815}
816
817async fn queue_follow_up(tx: &mut Transaction, job: Option<String>, batch_id: i64) -> Result {
820 if let Some(job) = job {
821 let job: Encoded = serde_json::from_str(&job)?;
822 insert(&mut *tx, &job, Duration::ZERO, None, None, Some(batch_id)).await?;
823 }
824 Ok(())
825}
826
827pub(crate) async fn continue_chain(tx: &mut Transaction, chain: &str) -> Result {
829 let rest: Vec<Encoded> = serde_json::from_str(chain)?;
830 let mut rest = rest.into_iter();
831 if let Some(next) = rest.next() {
832 let rest: Vec<Encoded> = rest.collect();
833 let chain = (!rest.is_empty())
834 .then(|| serde_json::to_string(&rest))
835 .transpose()?;
836 insert(&mut *tx, &next, Duration::ZERO, chain, None, None).await?;
837 }
838 Ok(())
839}
840
841pub(crate) async fn release_unique(tx: &mut Transaction, key: &str) -> Result {
843 crate::db::sql("DELETE FROM cache WHERE key = ?")
844 .bind(key)
845 .execute(&mut *tx)
846 .await?;
847 Ok(())
848}
849
850#[derive(Debug, Clone, Serialize)]
852#[non_exhaustive]
853pub struct BatchStatus {
854 pub id: i64,
856 pub name: String,
858 pub total: i64,
860 pub pending: i64,
862 pub failed: i64,
864 pub cancelled: bool,
866 pub finished: bool,
868 pub created_at: crate::db::DateTime,
870}
871
872impl BatchStatus {
873 pub fn progress(&self) -> u8 {
875 if self.total == 0 {
876 return 100;
877 }
878 ((self.total - self.pending) * 100 / self.total).clamp(0, 100) as u8
879 }
880}
881
882#[derive(Debug, Clone, Serialize)]
884#[non_exhaustive]
885pub struct FailedJob {
886 pub id: i64,
888 pub queue: String,
890 pub job: String,
892 pub payload: String,
894 pub error: String,
896 pub failed_at: crate::db::DateTime,
898}
899
900impl AppState {
901 pub async fn dispatch<J: Job>(&self, job: J) -> Result<i64> {
903 self.queue.dispatch(job).await
904 }
905
906 pub async fn dispatch_sync<J: Job>(&self, job: J) -> Result {
909 let ctx = JobContext {
910 state: self.clone(),
911 attempt: 1,
912 id: 0,
913 batch_id: None,
914 };
915 crate::context::scope_app(self.clone(), job.handle(ctx)).await
916 }
917}