Skip to main content

apalis_sqlite/
sink.rs

1use std::{
2    pin::Pin,
3    task::{Context, Poll},
4};
5
6use futures::Sink;
7use sqlx::Executor;
8use ulid::Ulid;
9
10use crate::{SqliteStorage, SqliteTask, error::Error};
11
12/// Push a batch of tasks into the database
13pub async fn push_tasks<'a, E>(
14    executor: &'a mut E,
15    queue: &str,
16    tasks: &[SqliteTask],
17) -> Result<(), Error>
18where
19    for<'e> &'e mut E: Executor<'e, Database = sqlx::Sqlite>,
20{
21    for task in tasks {
22        let id = task
23            .task_id()
24            .as_ref()
25            .map(|id| id.to_string())
26            .unwrap_or(Ulid::generate().to_string());
27        let run_at = task.run_at().unwrap_or(0) as i64;
28        let max_attempts = task.max_attempts().unwrap_or(25) as i64;
29        let priority = task.priority().unwrap_or_default() as i64;
30        let args = &task.args;
31        let idempotency_key = task.idempotency_key();
32        let meta = serde_json::to_string(task.metadata()).unwrap_or_default();
33        sqlx::query_file!(
34            "queries/task/sink.sql",
35            args,
36            id,
37            queue,
38            max_attempts,
39            run_at,
40            priority,
41            meta,
42            idempotency_key
43        )
44        .execute(&mut *executor)
45        .await?;
46    }
47    Ok(())
48}
49
50impl<Args> Sink<SqliteTask> for SqliteStorage<Args> {
51    type Error = Error;
52
53    fn poll_ready(self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
54        Poll::Ready(Ok(()))
55    }
56
57    fn start_send(self: Pin<&mut Self>, item: SqliteTask) -> Result<(), Self::Error> {
58        self.project().persistence.start_send(item)
59    }
60
61    fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
62        Sink::poll_flush(self.project().persistence, cx)
63    }
64
65    fn poll_close(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
66        Sink::poll_close(self.project().persistence, cx)
67    }
68}