Skip to main content

apalis_sqlite/queries/
fetch_next.rs

1use 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
8/// Fetch the next batch of tasks from the sqlite backend
9pub 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}