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
12pub 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}