Skip to main content

ocre/
jobs.rs

1//! Background jobs (Cloudflare Queues) and scheduled tasks (Cron Triggers).
2//!
3//! Rails' Active Job, on [Cloudflare Queues](https://developers.cloudflare.com/queues/)
4//! (in the Workers Free plan since February 2026). The app's Worker is both
5//! the producer and the consumer of its queues: `<app>-jobs`, bound as
6//! [`QUEUE_BINDING`], plus one per named queue (`<app>-jobs-urgent`, bound
7//! as `JOBS_URGENT`, see [`queue`]). A job is a serde value, usually the
8//! app's `Job` enum (`src/jobs/mod.rs`, written by `ocre g job`), and
9//! dispatch is a plain `match` in the app's `perform` function, not a registry.
10//!
11//! A handler enqueues with [`enqueue`], [`enqueue_in`] or [`enqueue_all`]
12//! (generated jobs wrap it as `job.perform_later(&ctx)`) and returns as soon
13//! as Cloudflare stored the message; the Worker's `queue` event runs it
14//! moments later through [`consume`]. Scheduled tasks run from the
15//! `scheduled` event through [`cron`]:
16//!
17//! ```no_run
18//! mod jobs {
19//!     use ocre::{Ctx, Result};
20//!     use serde::{Deserialize, Serialize};
21//!
22//!     #[derive(Serialize, Deserialize)]
23//!     #[serde(rename_all = "snake_case")]
24//!     pub enum Job {
25//!         SendWelcome { user_id: i64 },
26//!     }
27//!
28//!     pub async fn perform(_ctx: Ctx, job: Job) -> Result<()> {
29//!         match job {
30//!             Job::SendWelcome { user_id } => {
31//!                 # let _ = user_id;
32//!                 Ok(())
33//!             }
34//!         }
35//!     }
36//! }
37//!
38//! mod schedules {
39//!     pub async fn run(_ctx: ocre::Ctx, cron: String) -> ocre::Result<()> {
40//!         # let _ = cron;
41//!         Ok(())
42//!     }
43//! }
44//!
45//! // src/lib.rs
46//! #[worker::event(queue)]
47//! async fn queue(batch: worker::MessageBatch<String>, env: worker::Env, _ctx: worker::Context) -> worker::Result<()> {
48//!     ocre::jobs::consume(batch, env, jobs::perform).await
49//! }
50//!
51//! #[worker::event(scheduled)]
52//! async fn scheduled(event: worker::ScheduledEvent, env: worker::Env, _ctx: worker::ScheduleContext) {
53//!     ocre::jobs::cron(event, env, schedules::run).await
54//! }
55//!
56//! // A handler
57//! async fn create(axum::extract::State(ctx): axum::extract::State<ocre::Ctx>) -> ocre::Result<&'static str> {
58//!     ocre::jobs::enqueue(&ctx, &jobs::Job::SendWelcome { user_id: 1 }).await?;
59//!     Ok("created")
60//! }
61//! # fn main() {}
62//! ```
63//!
64//! A message is JSON text, `{"at": <due unix time>, "job": {"send_welcome": {"user_id": 1}}}`.
65//! [`consume`] acknowledges a job that returns `Ok`; drops, with a log line,
66//! one that returns an error another try cannot fix (a 4xx
67//! [`Error`] such as `NotFound`: Rails' `discard_on`); retries any other
68//! `Err` with a growing delay (30 s, 1 min, 3 min, 9 min, 27 min: twice the
69//! time since it was due); and drops a message it cannot decode (an unknown
70//! or changed job), so it is never retried forever. After `maxRetries: 5`
71//! (cloudflare.config.ts) Cloudflare moves a failing message to the dead-letter
72//! queue `<app>-jobs-failed`, kept 24 hours. Delivery is at-least-once: write
73//! jobs to be safe to repeat. Every line Ocre logs starts with
74//! [`LOG_PREFIX`] or [`CRON_LOG_PREFIX`].
75//!
76//! # Free-plan budget (September 2026)
77//!
78//! - **Queues operations**: 10,000 a day. A message costs 3 (write, read,
79//!   delete), each retry 1 more read, a dead-lettered message 1 more write.
80//!   Ocre sends one message per job: about 3,300 jobs a day.
81//! - **Retention**: 24 hours on Free; the last retry comes after about 40 minutes.
82//! - **Message size**: 128 KB; [`enqueue`] refuses larger jobs (pass ids).
83//!   [`enqueue_all`] sends 100 messages (256 KB) per call.
84//! - **Delay**: 24 hours at most, on send and on retry;
85//!   [`enqueue_in`] refuses longer delays
86//!   ([`MAX_DELAY`]).
87//! - **Batches**: up to 100 messages and 60 s wait; the generated
88//!   `max_batch_size = 10`, `max_batch_timeout = 5` make one consumer run
89//!   (one Worker invocation) per 10 jobs.
90//! - **CPU**: 10 ms per invocation, consumer batches and cron runs included:
91//!   jobs should be I/O (D1, mail, `fetch`); lower `max_batch_size` for CPU-heavy jobs.
92//! - **Cron Triggers**: 5 per account; run several tasks from one cron.
93
94use std::time::Duration;
95
96use serde::{Deserialize, Serialize, de::DeserializeOwned};
97use serde_json::Value;
98
99pub use crate::runtime::jobs::{Queue, consume, cron, enqueue, enqueue_all, enqueue_in, lock, queue, unlock};
100
101/// The `job_locks` table of [`lock`] and [`unlock`], created by the
102/// migration of `ocre g job --lock`.
103pub const LOCKS_TABLE_SQL: &str = "CREATE TABLE job_locks (
104  key TEXT PRIMARY KEY,
105  owner TEXT NOT NULL,
106  expires_at INTEGER NOT NULL
107);
108";
109
110use crate::{Error, Result, mail::Email};
111
112/// Name of the queue producer binding every Ocre app sends jobs to.
113///
114/// `ocre g job` adds `JOBS: bindings.queue({ name: "<app>-jobs" })` (and the
115/// consumer trigger) to cloudflare.config.ts; without it, enqueueing is an
116/// [`Error::Internal`] naming that entry.
117///
118/// # Examples
119///
120/// ```
121/// assert_eq!(ocre::jobs::QUEUE_BINDING, "JOBS");
122/// ```
123pub const QUEUE_BINDING: &str = "JOBS";
124
125/// Name of the queue [`enqueue`] and [`enqueue_in`] use: `default`, bound as [`QUEUE_BINDING`].
126///
127/// Other queues (`ocre g job <Name> --queue urgent`) are bound as
128/// `JOBS_<NAME>` and used through [`queue`].
129///
130/// # Examples
131///
132/// ```
133/// assert_eq!(ocre::jobs::DEFAULT_QUEUE, "default");
134/// ```
135pub const DEFAULT_QUEUE: &str = "default";
136
137/// Prefix of every line Ocre logs about jobs, e.g. `[ocre jobs] send_welcome done`.
138///
139/// [`consume`] logs `<prefix> <job> done`,
140/// `<prefix> <job> failed, retrying in <n> s: <error>`,
141/// `<prefix> <job> discarded, not retried: <error>` and
142/// `<prefix> dropped message <id>: <reason>`.
143///
144/// # Examples
145///
146/// ```
147/// assert!("[ocre jobs] send_welcome done".starts_with(ocre::jobs::LOG_PREFIX));
148/// ```
149pub const LOG_PREFIX: &str = "[ocre jobs]";
150
151/// Prefix of every line Ocre logs about Cron Triggers, e.g. `[ocre cron] 0 3 * * * done`.
152///
153/// [`cron`] logs `<prefix> <cron> done` or `<prefix> <cron> failed: <error>`.
154///
155/// # Examples
156///
157/// ```
158/// assert!("[ocre cron] 0 3 * * * done".starts_with(ocre::jobs::CRON_LOG_PREFIX));
159/// ```
160pub const CRON_LOG_PREFIX: &str = "[ocre cron]";
161
162/// Longest delay Cloudflare Queues accepts: 24 hours, for [`enqueue_in`] and retries.
163///
164/// A longer delay makes [`enqueue_in`] fail; retry
165/// delays are capped to it.
166///
167/// # Examples
168///
169/// ```
170/// assert_eq!(ocre::jobs::MAX_DELAY.as_secs(), 24 * 60 * 60);
171/// ```
172pub const MAX_DELAY: Duration = Duration::from_secs(24 * 60 * 60);
173
174/// Largest message Ocre sends: Queues' 128 KB limit, minus room for the
175/// ~100 bytes of metadata Cloudflare adds.
176pub(crate) const MAX_MESSAGE_BYTES: usize = 127_000;
177
178/// First retry delay, in seconds.
179pub(crate) const RETRY_BASE_SECONDS: u32 = 30;
180
181/// The JSON text of a queue message: `{"at": 1727000000, "job": {...}}` or
182/// `{"at": ..., "mail": {...}}` for [`deliver_later`](crate::mail::deliver_later).
183#[derive(Debug, PartialEq, Serialize, Deserialize)]
184pub(crate) struct Envelope {
185    /// Unix time the message was due; retries back off from it.
186    pub at: i64,
187    #[serde(flatten)]
188    pub payload: Payload,
189}
190
191#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
192#[serde(rename_all = "snake_case")]
193pub(crate) enum Payload {
194    /// An app job, as serialized by the app (`{"send_welcome": {"user_id": 1}}`).
195    Job(Value),
196    /// An email to send with [`mail::send`](crate::mail::send).
197    Mail(Box<Email>),
198}
199
200/// The job as JSON.
201pub(crate) fn job_payload<J: Serialize>(job: &J) -> Result<Payload> {
202    serde_json::to_value(job)
203        .map(Payload::Job)
204        .map_err(|err| Error::internal(format!("cannot enqueue the job: it does not serialize to JSON ({err})")))
205}
206
207/// Seconds for `delaySeconds`; at most [`MAX_DELAY`].
208pub(crate) fn delay_seconds(delay: Duration) -> Result<u32> {
209    if delay > MAX_DELAY {
210        return Err(Error::internal(format!(
211            "cannot delay a job by {} s: Cloudflare Queues delays messages by 24 hours (86400 s) at most. Fix: \
212             enqueue it later from a scheduled task (`ocre g schedule`), or store the due time in D1",
213            delay.as_secs()
214        )));
215    }
216    Ok(u32::try_from(delay.as_secs()).expect("24 hours fit in u32"))
217}
218
219/// The message text, due at `at`. Too large a message is an error that names the fix.
220pub(crate) fn encode(payload: Payload, at: i64) -> Result<String> {
221    let text = serde_json::to_string(&Envelope { at, payload }).expect("JSON values and emails serialize");
222    if text.len() > MAX_MESSAGE_BYTES {
223        return Err(Error::internal(format!(
224            "cannot enqueue a {} byte message: Cloudflare Queues messages hold 128 KB at most. Fix: store large \
225             data in D1 or R2 and put its id in the job",
226            text.len()
227        )));
228    }
229    Ok(text)
230}
231
232/// Reads a message body. `Err` explains why it is dropped.
233pub(crate) fn decode(body: Option<String>) -> std::result::Result<Envelope, String> {
234    let text = body.ok_or("the body is not text; Ocre sends jobs as JSON text with `ocre::jobs::enqueue`")?;
235    serde_json::from_str(&text).map_err(|err| format!("not an Ocre job message ({err}): {}", preview(&text)))
236}
237
238/// A job's name in logs: the variant of an externally tagged enum
239/// (`send_welcome` for `{"send_welcome": {...}}` or `"send_welcome"`), else `job`.
240pub(crate) fn job_name(job: &Value) -> &str {
241    match job {
242        Value::String(name) => name,
243        Value::Object(map) if map.len() == 1 => map.keys().next().expect("one key"),
244        _ => "job",
245    }
246}
247
248/// The app's job, from the message JSON. `Err` explains why it is dropped.
249pub(crate) fn decode_job<J: DeserializeOwned>(job: Value) -> std::result::Result<J, String> {
250    let text = job.to_string();
251    serde_json::from_value(job).map_err(|err| {
252        format!(
253            "{} does not match the app's jobs ({err}); a job renamed or changed while messages were queued? \
254             Fix: keep accepting the old form in src/jobs/mod.rs",
255            preview(&text)
256        )
257    })
258}
259
260/// Delay before retrying a job due at `at` that failed at `now`: twice the
261/// time already waited, between 30 s and 24 hours.
262pub(crate) fn retry_delay(now: i64, at: i64) -> u32 {
263    let waited = u32::try_from(now.saturating_sub(at).max(0)).unwrap_or(u32::MAX);
264    let max = u32::try_from(MAX_DELAY.as_secs()).expect("24 hours fit in u32");
265    waited.saturating_mul(2).clamp(RETRY_BASE_SECONDS, max)
266}
267
268/// Error for a missing queue producer binding.
269pub(crate) fn missing_queue(binding: &str, detail: &str) -> Error {
270    let fix = if binding == QUEUE_BINDING {
271        "run `ocre g job <Name>` once; it adds `JOBS: bindings.queue(...)` and its `triggers.queue(...)` consumer"
272            .to_owned()
273    } else {
274        let name = binding.trim_start_matches("JOBS_").to_ascii_lowercase().replace('_', "-");
275        format!(
276            "run `ocre g job <Name> --queue {name}`; it adds `{binding}: bindings.queue(...)` and its `triggers.queue(...)` consumer"
277        )
278    };
279    Error::internal(format!("the queue binding `{binding}` is missing ({detail}). Fix: {fix} to cloudflare.config.ts"))
280}
281
282/// The producer binding of a named queue: `default` is `JOBS`, `urgent` is
283/// `JOBS_URGENT`, `low-priority` is `JOBS_LOW_PRIORITY`. A name that is not
284/// lowercase letters, digits and `-` is an error.
285pub(crate) fn binding(queue: &str) -> Result<String> {
286    let valid = !queue.is_empty()
287        && queue.chars().all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || c == '-')
288        && !queue.starts_with('-')
289        && !queue.ends_with('-');
290    if !valid {
291        return Err(Error::internal(format!(
292            "invalid queue name {queue:?}. Fix: use lowercase letters, digits and `-`, e.g. \"urgent\", as in \
293             `ocre g job <Name> --queue urgent`"
294        )));
295    }
296    Ok(match queue {
297        DEFAULT_QUEUE => QUEUE_BINDING.to_owned(),
298        name => format!("{QUEUE_BINDING}_{}", name.to_ascii_uppercase().replace('-', "_")),
299    })
300}
301
302/// Most messages in one `sendBatch` call.
303pub(crate) const MAX_BATCH_MESSAGES: usize = 100;
304/// Most bytes in one `sendBatch` call: Queues' 256 KB, minus room for metadata.
305pub(crate) const MAX_BATCH_BYTES: usize = 250_000;
306
307/// Splits message bodies into `sendBatch` calls: 100 messages and 256 KB at most each, in order.
308pub(crate) fn batches(bodies: Vec<String>) -> Vec<Vec<String>> {
309    let mut out: Vec<Vec<String>> = Vec::new();
310    let mut size = 0;
311    for body in bodies {
312        let full = out
313            .last()
314            .is_none_or(|batch| batch.len() == MAX_BATCH_MESSAGES || size + body.len() + 100 > MAX_BATCH_BYTES);
315        if full {
316            out.push(Vec::new());
317            size = 0;
318        }
319        size += body.len() + 100;
320        out.last_mut().expect("pushed above").push(body);
321    }
322    out
323}
324
325/// Whether a failed job is dropped instead of retried: errors that another
326/// try cannot fix (a missing record, bad input, a refused permission), like
327/// Rails' `discard_on`. `Internal` and `TooManyRequests` are retried.
328pub(crate) fn discards(err: &Error) -> bool {
329    matches!(
330        err,
331        Error::NotFound
332            | Error::BadRequest(_)
333            | Error::Unauthorized
334            | Error::Forbidden
335            | Error::Invalid(_)
336            | Error::PayloadTooLarge(_)
337            | Error::Conflict(_)
338    )
339}
340
341/// The first 200 characters, for log lines.
342fn preview(text: &str) -> String {
343    match text.char_indices().nth(200) {
344        Some((end, _)) => format!("{}...", &text[..end]),
345        None => text.to_owned(),
346    }
347}
348
349/// How many calls an invocation may still make: D1 queries, `fetch`es, KV
350/// and R2 operations, as the job counts them.
351///
352/// On the free plan a Worker invocation (a request, a queue batch, a cron
353/// run) may run 50 D1 queries and make 50 subrequests; past that the call
354/// fails. A job over many rows takes a budget, spends it per call, and
355/// stops early with a cursor to continue from ([`run_steps`]). Ocre does not
356/// count for you: [`take`](Self::take) what each step is about to use.
357///
358/// # Examples
359///
360/// ```
361/// use ocre::jobs::Budget;
362///
363/// let budget = Budget::new(Budget::FREE_D1_QUERIES - 5); // 5 for the rest of the job
364/// assert!(budget.take(40));
365/// assert!(!budget.take(10), "not enough left: nothing taken");
366/// assert_eq!(budget.left(), 5);
367/// ```
368#[derive(Debug)]
369pub struct Budget {
370    left: std::sync::atomic::AtomicU32,
371}
372
373impl Budget {
374    /// D1 queries per invocation on the free plan: 50.
375    pub const FREE_D1_QUERIES: u32 = 50;
376    /// Subrequests (`fetch`) per invocation on the free plan: 50.
377    pub const FREE_SUBREQUESTS: u32 = 50;
378
379    /// A budget of `calls`.
380    pub fn new(calls: u32) -> Self {
381        Self { left: std::sync::atomic::AtomicU32::new(calls) }
382    }
383
384    /// Takes `calls` from the budget when that many are left; `false` (and nothing taken) otherwise.
385    pub fn take(&self, calls: u32) -> bool {
386        use std::sync::atomic::Ordering;
387        self.left.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |left| left.checked_sub(calls)).is_ok()
388    }
389
390    /// The calls not taken yet.
391    pub fn left(&self) -> u32 {
392        self.left.load(std::sync::atomic::Ordering::SeqCst)
393    }
394}
395
396/// What a step of [`run_steps`] did: more to do from the cursor, or done.
397#[derive(Debug, Clone, PartialEq, Eq)]
398pub enum Step<C> {
399    /// Continue from this cursor (the next page, the last id handled).
400    Next(C),
401    /// Nothing left.
402    Done,
403}
404
405/// Runs `step` from `cursor` while `budget` covers its `cost` (the calls one
406/// step makes), for work too large for one invocation: a page of rows per
407/// step, the next page's cursor between them (Active Job's continuations).
408///
409/// Returns `Some(cursor)` when the budget ran out first: enqueue the job
410/// again with it, and the next run continues there. `None` when a step
411/// answered [`Step::Done`]. A failing step stops the run with its error; a
412/// retried job starts again from the cursor it was enqueued with, so make
413/// steps safe to repeat.
414///
415/// # Errors
416///
417/// The first error a step returns.
418///
419/// # Examples
420///
421/// ```no_run
422/// use ocre::jobs::{Budget, Step, run_steps};
423/// use ocre::{Ctx, Query, Result};
424/// use serde::{Deserialize, Serialize};
425///
426/// #[derive(Deserialize)]
427/// struct Row {
428///     id: i64,
429/// }
430///
431/// #[derive(Serialize, Deserialize)]
432/// pub struct Reindex {
433///     after_id: i64,
434/// }
435///
436/// impl Reindex {
437///     pub async fn perform(self, ctx: &Ctx) -> Result<()> {
438///         let budget = Budget::new(Budget::FREE_D1_QUERIES - 1); // 1 left to enqueue the rest
439///         // Each step reads a page of 100 rows and writes them back: 2 queries.
440///         let rest = run_steps(&budget, 2, self.after_id, |after_id| async move {
441///             let db = ctx.db()?;
442///             let page: Vec<Row> = Query::table("tracks").gt("id", after_id).order_asc("id").limit(100).all(&db).await?;
443///             let Some(last) = page.last() else { return Ok(Step::Done) };
444///             // ... one bulk::update of the page ...
445///             Ok(Step::Next(last.id))
446///         })
447///         .await?;
448///         if let Some(after_id) = rest {
449///             ocre::jobs::enqueue(ctx, &Reindex { after_id }).await?;
450///         }
451///         Ok(())
452///     }
453/// }
454/// ```
455pub async fn run_steps<C, F, Fut>(budget: &Budget, cost: u32, mut cursor: C, mut step: F) -> Result<Option<C>>
456where
457    F: FnMut(C) -> Fut,
458    Fut: Future<Output = Result<Step<C>>>,
459{
460    loop {
461        if !budget.take(cost) {
462            return Ok(Some(cursor));
463        }
464        match step(cursor).await? {
465            Step::Next(next) => cursor = next,
466            Step::Done => return Ok(None),
467        }
468    }
469}
470
471/// Development endpoint listing recent jobs, served by `ocre dev` only, for
472/// tests (Rails' `assert_enqueued_with` and `assert_performed_jobs`).
473///
474/// `GET /ocre/dev/jobs.json` answers the last 50 jobs this Worker instance
475/// enqueued and the last 50 it ran, oldest first:
476/// `{"enqueued": [{"id": 1, "queue": "JOBS", "job": {"send_welcome": {"user_id": 7}}}],
477/// "performed": [{"id": 2, "job": "send_welcome", "outcome": "done"}]}`. The
478/// outcome is `done`, `discarded` or `retried`, as [`consume`] logs it.
479/// `ocre::testing::Client::jobs` reads it.
480///
481/// Debug builds only (`ocre dev`); release builds (`ocre deploy`) get an
482/// empty router, so it is a 404 in production. The first `ocre g job`
483/// merges it into `routes()`. It uses no billed resource: the lists live
484/// in the Worker's memory.
485///
486/// # Examples
487///
488/// ```
489/// use axum::Router;
490/// use ocre::Ctx;
491///
492/// fn routes() -> Router<Ctx> {
493///     Router::new().merge(ocre::jobs::dev_routes())
494/// }
495/// # let _ = routes;
496/// ```
497pub fn dev_routes<S: Clone + Send + Sync + 'static>() -> axum::Router<S> {
498    #[cfg(not(debug_assertions))]
499    {
500        axum::Router::new()
501    }
502    #[cfg(debug_assertions)]
503    axum::Router::new().route("/ocre/dev/jobs.json", axum::routing::get(|| async { axum::Json(dev::snapshot()) }))
504}
505
506/// Remembers a job sent to the queue bound as `binding` (debug builds only; emails are listed by the mail pages).
507pub(crate) fn record_enqueued(binding: &str, payload: &Payload) {
508    #[cfg(debug_assertions)]
509    if let Payload::Job(job) = payload {
510        dev::enqueued(binding, job);
511    }
512    #[cfg(not(debug_assertions))]
513    let _ = (binding, payload);
514}
515
516/// Remembers how a job run ended: `done`, `discarded` or `retried` (debug builds only).
517pub(crate) fn record_performed(job: &str, outcome: &str) {
518    #[cfg(debug_assertions)]
519    dev::performed(job, outcome);
520    #[cfg(not(debug_assertions))]
521    let _ = (job, outcome);
522}
523
524#[cfg(debug_assertions)]
525mod dev {
526    use std::sync::{Mutex, PoisonError};
527
528    use serde::Serialize;
529    use serde_json::Value;
530
531    /// How many jobs each list keeps.
532    const KEEP: usize = 50;
533
534    #[derive(Debug, Clone, Serialize)]
535    pub(super) struct Enqueued {
536        id: u64,
537        queue: String,
538        job: Value,
539    }
540
541    #[derive(Debug, Clone, Serialize)]
542    pub(super) struct Performed {
543        id: u64,
544        job: String,
545        outcome: String,
546    }
547
548    #[derive(Debug, Clone, Serialize)]
549    pub(super) struct Snapshot {
550        enqueued: Vec<Enqueued>,
551        performed: Vec<Performed>,
552    }
553
554    static JOBS: Mutex<(u64, Snapshot)> = Mutex::new((1, Snapshot { enqueued: Vec::new(), performed: Vec::new() }));
555
556    fn keep<T>(list: &mut Vec<T>, item: T) {
557        list.push(item);
558        if list.len() > KEEP {
559            list.remove(0);
560        }
561    }
562
563    pub(super) fn enqueued(queue: &str, job: &Value) {
564        let mut jobs = JOBS.lock().unwrap_or_else(PoisonError::into_inner);
565        let id = jobs.0;
566        jobs.0 += 1;
567        keep(&mut jobs.1.enqueued, Enqueued { id, queue: queue.to_owned(), job: job.clone() });
568    }
569
570    pub(super) fn performed(job: &str, outcome: &str) {
571        let mut jobs = JOBS.lock().unwrap_or_else(PoisonError::into_inner);
572        let id = jobs.0;
573        jobs.0 += 1;
574        keep(&mut jobs.1.performed, Performed { id, job: job.to_owned(), outcome: outcome.to_owned() });
575    }
576
577    pub(super) fn snapshot() -> Snapshot {
578        JOBS.lock().unwrap_or_else(PoisonError::into_inner).1.clone()
579    }
580}
581
582#[cfg(test)]
583#[path = "../tests/jobs.rs"]
584mod tests;