pub struct Job {Show 14 fields
pub id: String,
pub kind: String,
pub dedup_key: String,
pub payload: Value,
pub status: String,
pub run_at: i64,
pub attempts: i64,
pub max_attempts: i64,
pub deadline: Option<i64>,
pub lease_until: Option<i64>,
pub lease_owner: Option<String>,
pub last_error: Option<String>,
pub created_at: i64,
pub updated_at: i64,
}Expand description
One stored job row.
status and kind come back as the strings they were stored as, for the
reason crate::sqlite::audit::AuditEntry keeps event a String: an
older binary meeting a row a newer one wrote should render it, not refuse to
load. The runner never claims a kind its registry does not hold, so an
unrecognised one is simply left alone.
Fields§
§id: String§kind: String§dedup_key: String§payload: Value§status: String§run_at: i64§attempts: i64§max_attempts: i64§deadline: Option<i64>§lease_until: Option<i64>§lease_owner: Option<String>§last_error: Option<String>§created_at: i64§updated_at: i64Implementations§
Source§impl Job
impl Job
Sourcepub async fn enqueue(
row: NewJob<'_>,
database: &Database,
) -> Result<bool, Error>
pub async fn enqueue( row: NewJob<'_>, database: &Database, ) -> Result<bool, Error>
Queues one job, unless a live one already holds (kind, dedup_key).
Ok(false) is “already queued”, not an error: it is what two racing
callers must both be able to survive, and the partial unique index is
what decides which of them won. A job that already reached done or
failed releases the identity, so the same key can be queued again —
which is what makes a retried order and a periodic sweep both expressible
without a second table.
Sourcepub async fn claim_next(
runner_id: &str,
kinds: &[&str],
lease_until: i64,
now: i64,
database: &Database,
) -> Result<Option<Self>, Error>
pub async fn claim_next( runner_id: &str, kinds: &[&str], lease_until: i64, now: i64, database: &Database, ) -> Result<Option<Self>, Error>
Claims the oldest eligible job of one of kinds, or None.
One statement, and that is the point: the subselect picks a candidate and
the outer AND status = 'ready' is what makes the pick binding, so two
runners choosing the same row end with exactly one write. None collapses
“the queue was empty” and “somebody else won” — which is correct, because
the caller does the same thing either way.
attempts increments here rather than at completion: a job that reliably
kills the process must still exhaust its budget, and nothing reports back
from a process that died.
deadline is deliberately not filtered here. Skipping a dead row
would leave it ready and re-read on every tick for ever; the runner
claims it and retires it on the spot.
Sourcepub async fn complete(
id: &str,
runner_id: &str,
database: &Database,
) -> Result<bool, Error>
pub async fn complete( id: &str, runner_id: &str, database: &Database, ) -> Result<bool, Error>
Marks a claimed job finished. Ok(false) means the lease was lost.
Sourcepub async fn retry(
id: &str,
runner_id: &str,
run_at: i64,
error: &str,
database: &Database,
) -> Result<bool, Error>
pub async fn retry( id: &str, runner_id: &str, run_at: i64, error: &str, database: &Database, ) -> Result<bool, Error>
Returns a claimed job to the queue, to run again at run_at.
Sourcepub async fn reschedule(
id: &str,
runner_id: &str,
run_at: i64,
database: &Database,
) -> Result<bool, Error>
pub async fn reschedule( id: &str, runner_id: &str, run_at: i64, database: &Database, ) -> Result<bool, Error>
Returns a claimed job to the queue as a fresh occurrence: the attempt counter goes back to zero and the last error is cleared.
That reset is what separates a periodic job from a retried one. A sweep that runs every day for a year must not accumulate 365 attempts and retire itself, and a successful occurrence must not leave the previous failure’s text sitting on the row as though it were current.
Sourcepub async fn abandon(
id: &str,
runner_id: &str,
error: &str,
database: &Database,
) -> Result<bool, Error>
pub async fn abandon( id: &str, runner_id: &str, error: &str, database: &Database, ) -> Result<bool, Error>
Retires a claimed job permanently, recording why.
Sourcepub async fn reclaim_expired(
now: i64,
database: &Database,
) -> Result<u64, Error>
pub async fn reclaim_expired( now: i64, database: &Database, ) -> Result<u64, Error>
Returns to the queue every job whose runner died holding its lease.
attempts is deliberately left alone: the attempt really was spent, and
a counter rewritten to look better would let a job that crashes the
process loop for ever.
Sourcepub async fn release_owned(
runner_id: &str,
database: &Database,
) -> Result<u64, Error>
pub async fn release_owned( runner_id: &str, database: &Database, ) -> Result<u64, Error>
Releases every lease this runner holds, without settling the jobs.
What a graceful shutdown runs, so a restart re-claims its own work immediately instead of waiting out a full lease.
Sourcepub async fn find_by_id(
id: &str,
database: &Database,
) -> Result<Option<Self>, Error>
pub async fn find_by_id( id: &str, database: &Database, ) -> Result<Option<Self>, Error>
One row by id.
Sourcepub async fn find_live(
kind: &str,
dedup_key: &str,
database: &Database,
) -> Result<Option<Self>, Error>
pub async fn find_live( kind: &str, dedup_key: &str, database: &Database, ) -> Result<Option<Self>, Error>
The live job holding (kind, dedup_key), if there is one.
Sourcepub async fn count_live(kind: &str, database: &Database) -> Result<i64, Error>
pub async fn count_live(kind: &str, database: &Database) -> Result<i64, Error>
How many live jobs of kind are queued or running.
The kind-wide counterpart to Self::find_live, for the callers that
want “is there work of this sort outstanding?” without naming a
dedup_key — a notify_deliver key is a per-occurrence uuid, so there
is no single key to ask about.
Trait Implementations§
Auto Trait Implementations§
impl Freeze for Job
impl RefUnwindSafe for Job
impl Send for Job
impl Sync for Job
impl Unpin for Job
impl UnsafeUnpin for Job
impl UnwindSafe for Job
Blanket Implementations§
Source§impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
Source§impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
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
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
impl<A, B, T> HttpServerConnExec<A, B> for Twhere
B: Body,
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> ⓘ
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> ⓘ
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