use super::dialect::{JobDb, JobPool};
use super::enqueue::{EnqueuedJob, JobRequest, insert_job};
use super::error::{EnqueueError, MigrateError};
use super::migrate;
#[derive(Clone)]
pub struct Jobs {
pool: JobPool,
}
impl Jobs {
pub fn new(pool: JobPool) -> Self {
Self { pool }
}
pub fn pool(&self) -> &JobPool {
&self.pool
}
pub async fn migrate(&self) -> Result<(), MigrateError> {
migrate::apply(&self.pool).await
}
pub async fn migrate_tx(
&self,
tx: &mut sqlx::Transaction<'_, JobDb>,
) -> Result<(), MigrateError> {
migrate::apply_tx(tx).await
}
pub async fn enqueue<J>(&self, request: &JobRequest<J>) -> Result<EnqueuedJob, EnqueueError>
where
J: serde::Serialize + serde::de::DeserializeOwned,
{
insert_job(&self.pool, request).await
}
pub async fn enqueue_with<'c, E, J>(
&self,
executor: E,
request: &JobRequest<J>,
) -> Result<EnqueuedJob, EnqueueError>
where
E: sqlx::Executor<'c, Database = JobDb>,
J: serde::Serialize + serde::de::DeserializeOwned,
{
insert_job(executor, request).await
}
pub async fn enqueue_tx<J>(
&self,
tx: &mut sqlx::Transaction<'_, JobDb>,
request: &JobRequest<J>,
) -> Result<EnqueuedJob, EnqueueError>
where
J: serde::Serialize + serde::de::DeserializeOwned,
{
insert_job(&mut **tx, request).await
}
}