apalis_sqlite/queries/
fetch_next.rs1use std::time::{SystemTime, UNIX_EPOCH};
2
3use apalis_core::{task::Task, worker::context::WorkerContext};
4use sqlx::Executor;
5
6use crate::{Error, config::Config, from_row::SqliteTaskRow};
7
8pub async fn fetch_next<'a, E>(
10 executor: E,
11 config: &Config,
12 worker: &WorkerContext,
13) -> Result<Vec<Task<Vec<u8>>>, Error>
14where
15 E: Executor<'a, Database = sqlx::Sqlite>,
16{
17 let job_type = config.queue.as_ref();
18 let buffer_size = config.batch_size as i32;
19 let worker = worker.name();
20 let now: i64 = SystemTime::now()
21 .duration_since(UNIX_EPOCH)
22 .unwrap()
23 .as_secs() as i64;
24 sqlx::query_file_as!(
25 SqliteTaskRow,
26 "queries/backend/fetch_next.sql",
27 worker,
28 job_type,
29 buffer_size,
30 now
31 )
32 .fetch_all(executor)
33 .await?
34 .into_iter()
35 .map(|r| r.try_into())
36 .collect()
37}