aide_de_camp_sqlite/
job_handle.rs

1use crate::types::JobRow;
2use aide_de_camp::core::job_handle::JobHandle;
3use aide_de_camp::core::queue::QueueError;
4use aide_de_camp::core::{Bytes, Xid};
5use anyhow::Context;
6use async_trait::async_trait;
7use sqlx::SqlitePool;
8
9pub struct SqliteJobHandle {
10    pool: SqlitePool,
11    row: JobRow,
12}
13
14#[async_trait]
15impl JobHandle for SqliteJobHandle {
16    fn id(&self) -> Xid {
17        self.row.jid
18    }
19
20    fn job_type(&self) -> &str {
21        &self.row.job_type
22    }
23
24    fn payload(&self) -> Bytes {
25        self.row.payload.clone()
26    }
27
28    fn retries(&self) -> u32 {
29        self.row.retries
30    }
31
32    async fn complete(mut self) -> Result<(), QueueError> {
33        let jid = self.row.jid.to_string();
34        sqlx::query!("DELETE FROM adc_queue where jid = ?1", jid)
35            .execute(&self.pool)
36            .await
37            .context("Failed to mark job as completed")?;
38        Ok(())
39    }
40
41    async fn fail(mut self) -> Result<(), QueueError> {
42        let jid = self.row.jid.to_string();
43        sqlx::query!("UPDATE adc_queue SET started_at=null WHERE jid = ?1", jid)
44            .execute(&self.pool)
45            .await
46            .context("Failed to mark job as failed")?;
47        Ok(())
48    }
49
50    async fn dead_queue(mut self) -> Result<(), QueueError> {
51        let jid = self.row.jid.to_string();
52        let retries = self.row.retries;
53        let job_type = self.row.job_type.clone();
54        let payload = self.row.payload.as_ref();
55        let scheduled_at = self.row.scheduled_at;
56        let enqueued_at = self.row.enqueued_at;
57
58        let mut tx = self
59            .pool
60            .begin()
61            .await
62            .context("Failed to start transaction")?;
63        sqlx::query!("DELETE FROM adc_queue WHERE jid = ?1", jid)
64            .execute(&mut tx)
65            .await
66            .context("Failed to delete job from the queue")?;
67
68        sqlx::query!("INSERT INTO adc_dead_queue (jid, job_type, payload, retries, scheduled_at, enqueued_at) VALUES (?1, ?2, ?3, ?4,?5,?6)",
69            jid,
70            job_type,
71            payload,
72            retries,
73            scheduled_at,
74            enqueued_at)
75            .execute(&mut tx).await
76            .context("Failed to move job to dead queue")?;
77        Ok(())
78    }
79}
80
81impl SqliteJobHandle {
82    pub(crate) fn new(row: JobRow, pool: SqlitePool) -> Self {
83        Self { row, pool }
84    }
85}