pub struct SqliteQueueBackend { /* private fields */ }Expand description
Implementations§
Source§impl SqliteQueueBackend
impl SqliteQueueBackend
Sourcepub async fn new(path: impl AsRef<Path>) -> Result<Self>
pub async fn new(path: impl AsRef<Path>) -> Result<Self>
Open a SQLite database at path (creates the file if missing).
See crate-level Mode 2 — Enqueue binary and Mode 2 — Worker binary.
§Examples
use std::sync::Arc;
use boson_backend_sqlite::SqliteQueueBackend;
use boson_core::JsonExecutionContextFactory;
use boson_runtime::Boson;
let path = std::env::var("BOSON_SQLITE_PATH").unwrap_or_else(|_| "/tmp/boson.db".into());
let backend = SqliteQueueBackend::new(&path).await?;
let _boson = Boson::builder()
.queue_backend(Arc::new(backend))
.execution_context_factory(JsonExecutionContextFactory)
.auto_registry()
.build()?;§Errors
Returns an error when the database cannot be opened or schema bootstrap fails.
Sourcepub async fn connect(url: &str) -> Result<Self>
pub async fn connect(url: &str) -> Result<Self>
Connect using a SQLite connection URL.
§Errors
Returns an error when the pool cannot connect or schema bootstrap fails.
Sourcepub async fn from_pool(pool: SqlitePool) -> Result<Self>
pub async fn from_pool(pool: SqlitePool) -> Result<Self>
Wrap an existing pool (schema bootstrap runs).
§Errors
Returns an error when schema bootstrap fails.
Sourcepub fn pool(&self) -> &SqlitePool
pub fn pool(&self) -> &SqlitePool
Underlying connection pool.
§Panics
Panics if the inner pool is not SQLite (internal invariant violation).
Trait Implementations§
Source§impl Debug for SqliteQueueBackend
impl Debug for SqliteQueueBackend
Source§impl QueueBackend for SqliteQueueBackend
impl QueueBackend for SqliteQueueBackend
Source§fn upsert_job<'life0, 'life1, 'async_trait>(
&'life0 self,
job: &'life1 Job,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn upsert_job<'life0, 'life1, 'async_trait>(
&'life0 self,
job: &'life1 Job,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Persist or update a job row.
Source§fn enqueue_with_policies<'life0, 'life1, 'async_trait>(
&'life0 self,
job: Job,
task_config: &'life1 TaskConfig,
) -> Pin<Box<dyn Future<Output = Result<(String, JobEnqueueDisposition)>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn enqueue_with_policies<'life0, 'life1, 'async_trait>(
&'life0 self,
job: Job,
task_config: &'life1 TaskConfig,
) -> Pin<Box<dyn Future<Output = Result<(String, JobEnqueueDisposition)>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Insert job with idempotency semantics (see trait-level contract).
Source§fn get_job<'life0, 'life1, 'async_trait>(
&'life0 self,
job_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<Job>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn get_job<'life0, 'life1, 'async_trait>(
&'life0 self,
job_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<Job>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Load one job by id.
Source§fn list_jobs<'life0, 'async_trait>(
&'life0 self,
status_filter: Option<JobStatus>,
offset: usize,
limit: usize,
) -> Pin<Box<dyn Future<Output = Result<Vec<Job>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn list_jobs<'life0, 'async_trait>(
&'life0 self,
status_filter: Option<JobStatus>,
offset: usize,
limit: usize,
) -> Pin<Box<dyn Future<Output = Result<Vec<Job>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
List jobs with optional status filter and pagination.
Source§fn cancel_job_if_active<'life0, 'life1, 'async_trait>(
&'life0 self,
job_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn cancel_job_if_active<'life0, 'life1, 'async_trait>(
&'life0 self,
job_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Cancel a job if it is still active (
queued or running).Source§fn try_claim_job<'life0, 'life1, 'async_trait>(
&'life0 self,
job_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<Job>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn try_claim_job<'life0, 'life1, 'async_trait>(
&'life0 self,
job_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<Job>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Atomically claim a queued job for execution.
Source§fn revert_job_to_queued<'life0, 'life1, 'async_trait>(
&'life0 self,
job_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn revert_job_to_queued<'life0, 'life1, 'async_trait>(
&'life0 self,
job_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Revert a job to
queued (retry path).Source§fn distinct_pools_queued<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<Vec<String>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn distinct_pools_queued<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<Vec<String>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Distinct pool names with queued jobs.
Source§fn list_queued_for_pool_sorted<'life0, 'life1, 'async_trait>(
&'life0 self,
pool: &'life1 str,
limit: usize,
) -> Pin<Box<dyn Future<Output = Result<Vec<Job>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn list_queued_for_pool_sorted<'life0, 'life1, 'async_trait>(
&'life0 self,
pool: &'life1 str,
limit: usize,
) -> Pin<Box<dyn Future<Output = Result<Vec<Job>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Queued jobs for one pool, sorted by priority then created time.
Source§fn count_jobs<'life0, 'async_trait>(
&'life0 self,
status_filter: Option<JobStatus>,
) -> Pin<Box<dyn Future<Output = Result<u64>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn count_jobs<'life0, 'async_trait>(
&'life0 self,
status_filter: Option<JobStatus>,
) -> Pin<Box<dyn Future<Output = Result<u64>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Count jobs, optionally filtered by status.
Source§fn count_jobs_for_task<'life0, 'life1, 'async_trait>(
&'life0 self,
task_name: &'life1 str,
status: Option<JobStatus>,
) -> Pin<Box<dyn Future<Output = Result<u64>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn count_jobs_for_task<'life0, 'life1, 'async_trait>(
&'life0 self,
task_name: &'life1 str,
status: Option<JobStatus>,
) -> Pin<Box<dyn Future<Output = Result<u64>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Count jobs for one task, optionally filtered by status.
Source§fn count_active_jobs_for_task<'life0, 'life1, 'async_trait>(
&'life0 self,
task_name: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<u32>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn count_active_jobs_for_task<'life0, 'life1, 'async_trait>(
&'life0 self,
task_name: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<u32>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Count active (
queued + running) jobs for rate-limit checks.Source§fn find_nonterminal_by_idempotency_key<'life0, 'life1, 'async_trait>(
&'life0 self,
key: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<String>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn find_nonterminal_by_idempotency_key<'life0, 'life1, 'async_trait>(
&'life0 self,
key: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<String>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Find non-terminal job by idempotency key. Read more
Source§fn upsert_run<'life0, 'life1, 'async_trait>(
&'life0 self,
run: &'life1 Run,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn upsert_run<'life0, 'life1, 'async_trait>(
&'life0 self,
run: &'life1 Run,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Persist or update a run row.
Source§fn get_run<'life0, 'life1, 'async_trait>(
&'life0 self,
run_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<Run>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn get_run<'life0, 'life1, 'async_trait>(
&'life0 self,
run_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<Run>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Load one run by id.
Source§fn list_runs<'life0, 'life1, 'async_trait>(
&'life0 self,
job_id_filter: Option<&'life1 str>,
offset: usize,
limit: usize,
) -> Pin<Box<dyn Future<Output = Result<Vec<Run>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn list_runs<'life0, 'life1, 'async_trait>(
&'life0 self,
job_id_filter: Option<&'life1 str>,
offset: usize,
limit: usize,
) -> Pin<Box<dyn Future<Output = Result<Vec<Run>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
List runs with optional job filter and pagination.
Source§fn finish_run<'life0, 'life1, 'async_trait>(
&'life0 self,
run_id: &'life1 str,
status: RunStatus,
duration_ms: Option<i64>,
error_message: Option<String>,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn finish_run<'life0, 'life1, 'async_trait>(
&'life0 self,
run_id: &'life1 str,
status: RunStatus,
duration_ms: Option<i64>,
error_message: Option<String>,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Mark a run terminal with outcome fields.
Source§fn count_runs<'life0, 'life1, 'async_trait>(
&'life0 self,
job_id_filter: Option<&'life1 str>,
) -> Pin<Box<dyn Future<Output = Result<u64>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn count_runs<'life0, 'life1, 'async_trait>(
&'life0 self,
job_id_filter: Option<&'life1 str>,
) -> Pin<Box<dyn Future<Output = Result<u64>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Count runs, optionally filtered by job id.
Source§fn count_runs_since<'life0, 'async_trait>(
&'life0 self,
since: DateTime<Utc>,
) -> Pin<Box<dyn Future<Output = Result<u64>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn count_runs_since<'life0, 'async_trait>(
&'life0 self,
since: DateTime<Utc>,
) -> Pin<Box<dyn Future<Output = Result<u64>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Count runs with
started_at >= since.Source§fn task_run_stats<'life0, 'life1, 'async_trait>(
&'life0 self,
task_name: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<TaskRunStats>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn task_run_stats<'life0, 'life1, 'async_trait>(
&'life0 self,
task_name: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<TaskRunStats>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Aggregate run totals for one task.
Source§fn get_task_config<'life0, 'life1, 'async_trait>(
&'life0 self,
task_name: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<TaskConfig>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn get_task_config<'life0, 'life1, 'async_trait>(
&'life0 self,
task_name: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<TaskConfig>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Load task config by name.
Source§fn upsert_task_config<'life0, 'life1, 'async_trait>(
&'life0 self,
config: &'life1 TaskConfig,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn upsert_task_config<'life0, 'life1, 'async_trait>(
&'life0 self,
config: &'life1 TaskConfig,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Persist task config.
Source§fn try_claim_run_lease<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
job_id: &'life1 str,
worker_id: &'life2 str,
ttl_secs: i64,
) -> Pin<Box<dyn Future<Output = Result<Option<String>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
fn try_claim_run_lease<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
job_id: &'life1 str,
worker_id: &'life2 str,
ttl_secs: i64,
) -> Pin<Box<dyn Future<Output = Result<Option<String>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
Attempt to claim a run lease for
job_id (stored as lease task id). Read moreSource§fn extend_lease<'life0, 'life1, 'async_trait>(
&'life0 self,
lease_id: &'life1 str,
ttl_secs: i64,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn extend_lease<'life0, 'life1, 'async_trait>(
&'life0 self,
lease_id: &'life1 str,
ttl_secs: i64,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Extend lease TTL for a held lease.
Source§fn release_lease<'life0, 'life1, 'async_trait>(
&'life0 self,
lease_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn release_lease<'life0, 'life1, 'async_trait>(
&'life0 self,
lease_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Release a held lease.
Source§fn expired_lease_job_pairs<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<Vec<(String, String)>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn expired_lease_job_pairs<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<Vec<(String, String)>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Expired leases as
(lease_record_id, job_id).Source§fn pop_claim_from_pool<'life0, 'life1, 'async_trait>(
&'life0 self,
_pool: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<Job>, BosonError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
Self: 'async_trait,
fn pop_claim_from_pool<'life0, 'life1, 'async_trait>(
&'life0 self,
_pool: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<Job>, BosonError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
Self: 'async_trait,
Atomically claim the highest-priority queued job from
pool when supported. Read moreAuto Trait Implementations§
impl !Freeze for SqliteQueueBackend
impl !RefUnwindSafe for SqliteQueueBackend
impl !UnwindSafe for SqliteQueueBackend
impl Send for SqliteQueueBackend
impl Sync for SqliteQueueBackend
impl Unpin for SqliteQueueBackend
impl UnsafeUnpin for SqliteQueueBackend
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Mutably borrows from an owned value. Read more
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
Source§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self>
fn into_either(self, into_left: bool) -> Either<Self, Self>
Converts
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
Converts
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more