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;