apalis_postgres/queries/
fetch_next.rs1use apalis_core::worker::context::WorkerContext;
2use sqlx::Executor;
3
4use crate::{PgTask, config::Config, error::Error};
5
6pub 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}