Skip to main content

apalis_postgres/queries/
fetch_next.rs

1use apalis_core::worker::context::WorkerContext;
2use sqlx::Executor;
3
4use crate::{PgTask, config::Config, error::Error};
5
6/// Fetch the next batch of tasks from the sqlite backend
7pub async fn fetch_next<E>(
8    conn: &mut E,
9    config: &Config,
10    worker: &WorkerContext,
11) -> Result<Vec<PgTask>, Error>
12where
13    for<'e> &'e mut E: Executor<'e, Database = sqlx::Postgres>,
14{
15    use crate::from_row::PgTaskRow;
16    let job_type = config.queue.as_ref();
17    let buffer_size = config.batch_size as i32;
18    let worker = worker.name();
19
20    sqlx::query_file_as!(
21        PgTaskRow,
22        "queries/task/fetch_next.sql",
23        worker,
24        job_type,
25        buffer_size
26    )
27    .fetch_all(conn)
28    .await?
29    .into_iter()
30    .map(|r| r.try_into())
31    .collect()
32}