Skip to main content

ocre/runtime/
jobs.rs

1use std::time::Duration;
2
3use serde::{Serialize, de::DeserializeOwned};
4use worker::{
5    BatchMessageBuilder, Env, MessageBatch, MessageBuilder, MessageExt, QueueContentType, QueueRetryOptionsBuilder,
6    ScheduledEvent, send::SendFuture,
7};
8
9use super::Ctx;
10use crate::{
11    Result,
12    jobs::{
13        CRON_LOG_PREFIX, DEFAULT_QUEUE, LOG_PREFIX, Payload, batches, binding, decode, decode_job, delay_seconds,
14        discards, encode, job_name, job_payload, missing_queue, retry_delay,
15    },
16    now,
17};
18
19/// Sends `job` to the `JOBS` queue, to run in the background through [`consume`](crate::jobs::consume).
20///
21/// Returns once Cloudflare stored the message; the `queue` event runs it
22/// within seconds (`maxBatchTimeout: 5` in cloudflare.config.ts). The job is
23/// serialized to JSON and wrapped as `{"at": <now>, "job": ...}`; the whole
24/// message must fit in 128 KB, so pass ids, not records. The returned future
25/// is `Send`, so axum handlers can await it.
26///
27/// Free-plan cost: each job is one message, 3 of the 10,000 daily Queues
28/// operations (a write now, a read and a delete when consumed), plus one
29/// read per retry and one write if it ends in the dead-letter queue: about
30/// 3,300 jobs a day.
31///
32/// # Errors
33///
34/// [`Error::Internal`](crate::Error::Internal) (500) when the job does not
35/// serialize to JSON, the message is over the 128 KB limit, the
36/// `JOBS: bindings.queue(...)` entry is missing from cloudflare.config.ts
37/// (run `ocre g job <Name>` once), or Queues refuses the message.
38///
39/// # Examples
40///
41/// ```no_run
42/// use axum::extract::State;
43/// use ocre::{Ctx, Result};
44/// use serde::Serialize;
45///
46/// #[derive(Serialize)]
47/// #[serde(rename_all = "snake_case")]
48/// enum Job {
49///     SendWelcome { user_id: i64 },
50/// }
51///
52/// async fn create(State(ctx): State<Ctx>) -> Result<&'static str> {
53///     ocre::jobs::enqueue(&ctx, &Job::SendWelcome { user_id: 1 }).await?;
54///     Ok("created")
55/// }
56/// ```
57pub fn enqueue<J: Serialize>(ctx: &Ctx, job: &J) -> impl Future<Output = Result<()>> + Send + use<J> {
58    queue(ctx, DEFAULT_QUEUE).enqueue(job)
59}
60
61/// Sends every job of `jobs` to the `JOBS` queue in as few calls as possible, like Rails' `perform_all_later`.
62///
63/// One `sendBatch` call carries up to 100 messages and 256 KB, so 250 jobs
64/// take 3 calls instead of 250: use it whenever a handler or a scheduled
65/// task enqueues more than a few jobs (a Worker invocation may only make a
66/// limited number of calls to bindings). The jobs may be different variants
67/// of the app's `Job` enum; they run in any order. Each call is atomic, the
68/// whole list is not: when a later call fails, the earlier jobs are queued.
69/// An empty list sends nothing.
70///
71/// Free-plan cost: the same as [`enqueue`](crate::jobs::enqueue) for each
72/// job (3 Queues operations each); only the number of calls shrinks.
73///
74/// # Errors
75///
76/// Every error of [`enqueue`](crate::jobs::enqueue); when one job cannot be
77/// serialized or is over 128 KB, nothing is sent.
78///
79/// # Examples
80///
81/// ```no_run
82/// use ocre::{Ctx, Result};
83/// use serde::Serialize;
84///
85/// #[derive(Serialize)]
86/// #[serde(rename_all = "snake_case")]
87/// enum Job {
88///     SendDigest { user_id: i64 },
89/// }
90///
91/// async fn digests(ctx: &Ctx, user_ids: &[i64]) -> Result<()> {
92///     let jobs: Vec<Job> = user_ids.iter().map(|&user_id| Job::SendDigest { user_id }).collect();
93///     ocre::jobs::enqueue_all(ctx, &jobs).await
94/// }
95/// ```
96pub fn enqueue_all<J: Serialize>(ctx: &Ctx, jobs: &[J]) -> impl Future<Output = Result<()>> + Send + use<J> {
97    queue(ctx, DEFAULT_QUEUE).enqueue_all(jobs)
98}
99
100/// A named queue, for jobs that must not wait behind others: `ocre::jobs::queue(&ctx, "urgent").enqueue(&job)`.
101///
102/// Cloudflare Queues has no priorities: Ocre gives urgent work its own
103/// queue instead, like Rails' `queue_as`/`set(queue:)` and Loco's named
104/// queues. Each queue has its own consumer settings in cloudflare.config.ts
105/// (`ocre g job <Name> --queue urgent` adds `<app>-jobs-urgent` with
106/// `maxBatchTimeout: 1`), so a backlog of slow jobs on `default` never
107/// delays it. Every queue is consumed by the same `queue` event and the
108/// same `perform`. The name `default` is the `JOBS` queue of
109/// [`enqueue`](crate::jobs::enqueue); `urgent` is bound as `JOBS_URGENT`.
110///
111/// Free-plan cost: queues are free to create; each message costs the same
112/// 3 operations whatever its queue.
113///
114/// # Examples
115///
116/// ```no_run
117/// use ocre::{Ctx, Result};
118/// use serde::Serialize;
119///
120/// #[derive(Serialize)]
121/// #[serde(rename_all = "snake_case")]
122/// enum Job {
123///     SendMagicLink { user_id: i64 },
124/// }
125///
126/// async fn sign_in(ctx: &Ctx) -> Result<()> {
127///     ocre::jobs::queue(ctx, "urgent").enqueue(&Job::SendMagicLink { user_id: 1 }).await
128/// }
129/// ```
130pub fn queue(ctx: &Ctx, name: &'static str) -> Queue {
131    Queue { env: ctx.env().clone(), name }
132}
133
134/// A job queue, from [`queue`](crate::jobs::queue): enqueue on it like on the default queue.
135///
136/// # Examples
137///
138/// ```no_run
139/// # async fn f(ctx: &ocre::Ctx) -> ocre::Result<()> {
140/// let urgent = ocre::jobs::queue(ctx, "urgent");
141/// urgent.enqueue(&"reindex").await?;
142/// # Ok(())
143/// # }
144/// ```
145#[derive(Debug, Clone)]
146pub struct Queue {
147    env: Env,
148    name: &'static str,
149}
150
151impl Queue {
152    /// Sends `job` to this queue, like [`enqueue`](crate::jobs::enqueue) does to `default`.
153    ///
154    /// # Errors
155    ///
156    /// Those of [`enqueue`](crate::jobs::enqueue); a missing binding names
157    /// the fix, `ocre g job <Name> --queue <name>`, and a queue name that is
158    /// not lowercase letters, digits and `-` is an
159    /// [`Error::Internal`](crate::Error::Internal).
160    ///
161    /// # Examples
162    ///
163    /// ```no_run
164    /// # async fn f(ctx: &ocre::Ctx) -> ocre::Result<()> {
165    /// ocre::jobs::queue(ctx, "urgent").enqueue(&serde_json::json!({"send_code": {"user_id": 1}})).await?;
166    /// # Ok(())
167    /// # }
168    /// ```
169    pub fn enqueue<J: Serialize>(&self, job: &J) -> impl Future<Output = Result<()>> + Send + use<J> {
170        send(self.env.clone(), self.name, job_payload(job), Duration::ZERO)
171    }
172
173    /// Sends `job` to this queue, to run after `delay` (24 hours at most), like [`enqueue_in`](crate::jobs::enqueue_in).
174    ///
175    /// # Errors
176    ///
177    /// Those of [`enqueue_in`](crate::jobs::enqueue_in) and [`Queue::enqueue`].
178    ///
179    /// # Examples
180    ///
181    /// ```no_run
182    /// # async fn f(ctx: &ocre::Ctx) -> ocre::Result<()> {
183    /// let later = std::time::Duration::from_secs(60);
184    /// ocre::jobs::queue(ctx, "urgent").enqueue_in(&"retry_payment", later).await?;
185    /// # Ok(())
186    /// # }
187    /// ```
188    pub fn enqueue_in<J: Serialize>(
189        &self,
190        job: &J,
191        delay: Duration,
192    ) -> impl Future<Output = Result<()>> + Send + use<J> {
193        send(self.env.clone(), self.name, job_payload(job), delay)
194    }
195
196    /// Sends every job of `jobs` to this queue in batches, like [`enqueue_all`](crate::jobs::enqueue_all).
197    ///
198    /// # Errors
199    ///
200    /// Those of [`enqueue_all`](crate::jobs::enqueue_all) and [`Queue::enqueue`].
201    ///
202    /// # Examples
203    ///
204    /// ```no_run
205    /// # async fn f(ctx: &ocre::Ctx) -> ocre::Result<()> {
206    /// ocre::jobs::queue(ctx, "urgent").enqueue_all(&["a", "b"]).await?;
207    /// # Ok(())
208    /// # }
209    /// ```
210    pub fn enqueue_all<J: Serialize>(&self, jobs: &[J]) -> impl Future<Output = Result<()>> + Send + use<J> {
211        let payloads: Vec<Result<Payload>> = jobs.iter().map(job_payload).collect();
212        send_all(self.env.clone(), self.name, payloads)
213    }
214}
215
216fn producer(env: &Env, name: &str) -> Result<worker::Queue> {
217    let binding = binding(name)?;
218    env.queue(&binding).map_err(|err| missing_queue(&binding, &err.to_string()))
219}
220
221/// Sends `job` to the `JOBS` queue like [`enqueue`](crate::jobs::enqueue), to run after `delay`.
222///
223/// The delay is whole seconds, 24 hours at most
224/// ([`MAX_DELAY`](crate::jobs::MAX_DELAY)); for later work, enqueue from a
225/// scheduled task or store the due time in D1. Costs the same Queues
226/// operations as [`enqueue`](crate::jobs::enqueue).
227///
228/// # Errors
229///
230/// [`Error::Internal`](crate::Error::Internal) (500) when `delay` is over
231/// 24 hours, plus every error of [`enqueue`](crate::jobs::enqueue).
232///
233/// # Examples
234///
235/// ```no_run
236/// use std::time::Duration;
237///
238/// use axum::extract::State;
239/// use ocre::{Ctx, Result};
240/// use serde::Serialize;
241///
242/// #[derive(Serialize)]
243/// #[serde(rename_all = "snake_case")]
244/// enum Job {
245///     SendReminder { user_id: i64 },
246/// }
247///
248/// async fn remind(State(ctx): State<Ctx>) -> Result<&'static str> {
249///     let reminder = Job::SendReminder { user_id: 1 };
250///     ocre::jobs::enqueue_in(&ctx, &reminder, Duration::from_secs(3600)).await?;
251///     Ok("reminder set")
252/// }
253/// ```
254pub fn enqueue_in<J: Serialize>(
255    ctx: &Ctx,
256    job: &J,
257    delay: Duration,
258) -> impl Future<Output = Result<()>> + Send + use<J> {
259    queue(ctx, DEFAULT_QUEUE).enqueue_in(job, delay)
260}
261
262/// Sends a message as JSON text to the queue `name`, due after `delay`.
263pub(crate) fn send(
264    env: Env,
265    name: &'static str,
266    payload: Result<Payload>,
267    delay: Duration,
268) -> impl Future<Output = Result<()>> + Send + use<> {
269    SendFuture::new(async move {
270        let delay = delay_seconds(delay)?;
271        let payload = payload?;
272        let recorded = cfg!(debug_assertions).then(|| payload.clone());
273        let body = encode(payload, now() + i64::from(delay))?;
274        let queue = producer(&env, name)?;
275        let message = MessageBuilder::new(body).content_type(QueueContentType::Text).delay_seconds(delay).build();
276        queue.send(message).await?;
277        if let Some(payload) = recorded {
278            crate::jobs::record_enqueued(name, &payload);
279        }
280        Ok(())
281    })
282}
283
284fn send_all(
285    env: Env,
286    name: &'static str,
287    payloads: Vec<Result<Payload>>,
288) -> impl Future<Output = Result<()>> + Send + use<> {
289    SendFuture::new(async move {
290        if payloads.is_empty() {
291            return Ok(());
292        }
293        let at = now();
294        let payloads = payloads.into_iter().collect::<Result<Vec<_>>>()?;
295        let bodies = payloads.iter().map(|payload| encode(payload.clone(), at)).collect::<Result<Vec<_>>>()?;
296        let queue = producer(&env, name)?;
297        for batch in batches(bodies) {
298            let messages =
299                batch.into_iter().map(|body| MessageBuilder::new(body).content_type(QueueContentType::Text).build());
300            queue.send_batch(BatchMessageBuilder::new().messages(messages).build()).await?;
301        }
302        for payload in &payloads {
303            crate::jobs::record_enqueued(name, payload);
304        }
305        Ok(())
306    })
307}
308
309/// Runs a batch of queue messages through the app's `perform`; the Worker's `queue` entry point.
310///
311/// Each job goes to `perform(ctx, job)`, each email from
312/// [`deliver_later`](crate::mail::deliver_later) to
313/// [`mail::send`](crate::mail::send). Messages run one after the other in one
314/// Worker invocation (10 per batch with the generated `max_batch_size`, all
315/// sharing its CPU limit: 10 ms on Free).
316///
317/// - `Ok` acknowledges the message and logs `[ocre jobs] <job> done`.
318/// - An error that another try cannot fix, like Rails' `discard_on`
319///   ([`Error::NotFound`](crate::Error::NotFound), `BadRequest`,
320///   `Unauthorized`, `Forbidden`, `Invalid`, `PayloadTooLarge`: a record
321///   deleted since the job was enqueued, bad input) logs
322///   `[ocre jobs] <job> discarded, not retried: <error>` and acknowledges it.
323/// - Any other `Err` (`Internal`, `TooManyRequests`) logs
324///   `[ocre jobs] <job> failed, retrying in <n> s: <error>` and
325///   retries the message after twice the time since it was due, between 30 s
326///   and 24 hours (30 s, 1 min, 3 min, 9 min, 27 min), until `max_retries = 5`
327///   sends it to the dead-letter queue `<app>-jobs-failed`.
328/// - A message that does not decode (not an Ocre message, or a job renamed
329///   or changed while messages were queued) is logged as
330///   `[ocre jobs] dropped message <id>: <reason>` and acknowledged, never retried.
331///
332/// Delivery is at-least-once, so a job can run twice: make it safe to repeat.
333///
334/// Free-plan cost: reading and acknowledging a message are 2 Queues
335/// operations (of 10,000 a day); each retry is one more read, a dead-lettered
336/// message one more write.
337///
338/// # Errors
339///
340/// Never returns `Err` itself: job errors are logged, then discarded or retried.
341///
342/// # Examples
343///
344/// `ocre g job` writes the entry point in `src/lib.rs` and `perform` in `src/jobs/mod.rs`:
345///
346/// ```no_run
347/// mod jobs {
348///     use ocre::{Ctx, Result};
349///     use serde::Deserialize;
350///
351///     #[derive(Deserialize)]
352///     #[serde(rename_all = "snake_case")]
353///     pub enum Job {
354///         SendWelcome { user_id: i64 },
355///     }
356///
357///     pub async fn perform(_ctx: Ctx, job: Job) -> Result<()> {
358///         match job {
359///             Job::SendWelcome { user_id } => {
360///                 # let _ = user_id;
361///                 Ok(())
362///             }
363///         }
364///     }
365/// }
366///
367/// #[worker::event(queue)]
368/// async fn queue(batch: worker::MessageBatch<String>, env: worker::Env, _ctx: worker::Context) -> worker::Result<()> {
369///     ocre::jobs::consume(batch, env, jobs::perform).await
370/// }
371/// # fn main() {}
372/// ```
373pub async fn consume<J, F, Fut>(batch: MessageBatch<String>, env: Env, perform: F) -> worker::Result<()>
374where
375    J: DeserializeOwned,
376    F: Fn(Ctx, J) -> Fut,
377    Fut: Future<Output = Result<()>>,
378{
379    let ctx = Ctx::new(env);
380    for message in batch.raw_iter() {
381        let envelope = match decode(message.body().as_string()) {
382            Ok(envelope) => envelope,
383            Err(reason) => {
384                worker::console_error!("{LOG_PREFIX} dropped message {}: {reason}", message.id());
385                message.ack();
386                continue;
387            }
388        };
389        let (name, result) = match envelope.payload {
390            Payload::Mail(email) => ("mail".to_owned(), super::mail::deliver(ctx.env(), *email).await),
391            Payload::Job(value) => {
392                let name = job_name(&value).to_owned();
393                match decode_job::<J>(value) {
394                    Ok(job) => (name, perform(ctx.fresh(), job).await),
395                    Err(reason) => {
396                        worker::console_error!("{LOG_PREFIX} dropped message {}: {reason}", message.id());
397                        message.ack();
398                        continue;
399                    }
400                }
401            }
402        };
403        let outcome = match &result {
404            Ok(()) => "done",
405            Err(err) if discards(err) => "discarded",
406            Err(_) => "retried",
407        };
408        crate::jobs::record_performed(&name, outcome);
409        match result {
410            Ok(()) => {
411                worker::console_log!("{LOG_PREFIX} {name} done");
412                message.ack();
413            }
414            Err(err) if discards(&err) => {
415                worker::console_error!("{LOG_PREFIX} {name} discarded, not retried: {err}");
416                report(&ctx, "ocre.job", &err, [("job", name.as_str()), ("retried", "false")]);
417                message.ack();
418            }
419            Err(err) => {
420                let delay = retry_delay(now(), envelope.at);
421                worker::console_error!("{LOG_PREFIX} {name} failed, retrying in {delay} s: {err}");
422                report(&ctx, "ocre.job", &err, [("job", name.as_str()), ("retried", "true")]);
423                message.retry_with_options(&QueueRetryOptionsBuilder::new().with_delay_seconds(delay).build());
424            }
425        }
426    }
427    super::errors::flush(&ctx).await;
428    Ok(())
429}
430
431/// Queues Ocre's own report of a failed job or cron run (already logged).
432pub(crate) fn report<const N: usize>(ctx: &Ctx, source: &str, err: &crate::Error, context: [(&str, &str); N]) {
433    let context = context.iter().map(|(key, value)| ((*key).to_owned(), serde_json::Value::from(*value))).collect();
434    ctx.errors().report_unhandled(&err.to_string(), source, context, false);
435}
436
437/// Runs the app's task for the Cron Trigger that fired; the Worker's `scheduled` entry point.
438///
439/// Calls `run(ctx, cron)`, where `cron` is the expression from
440/// a `triggers.scheduled({ schedule })` entry of cloudflare.config.ts (UTC), and logs
441/// `[ocre cron] <cron> done` or `[ocre cron] <cron> failed: <error>`.
442/// Cloudflare does not retry a failed run; the next one comes at the next
443/// scheduled time, so enqueue jobs from the task for work that must not be
444/// lost. The Free plan allows 5 Cron Triggers per account and 10 ms of CPU
445/// per run: run several tasks from one cron, and move heavy work to jobs.
446///
447/// # Examples
448///
449/// `ocre g schedule` writes the entry point in `src/lib.rs` and `run` in `src/schedules/mod.rs`:
450///
451/// ```no_run
452/// mod schedules {
453///     use ocre::{Ctx, Error, Result};
454///
455///     pub async fn run(_ctx: Ctx, cron: String) -> Result<()> {
456///         match cron.as_str() {
457///             "0 3 * * *" => Ok(()), // nightly_cleanup
458///             other => Err(Error::internal(format!("no task for cron {other}"))),
459///         }
460///     }
461/// }
462///
463/// #[worker::event(scheduled)]
464/// async fn scheduled(event: worker::ScheduledEvent, env: worker::Env, _ctx: worker::ScheduleContext) {
465///     ocre::jobs::cron(event, env, schedules::run).await
466/// }
467/// # fn main() {}
468/// ```
469pub async fn cron<F, Fut>(event: ScheduledEvent, env: Env, run: F)
470where
471    F: FnOnce(Ctx, String) -> Fut,
472    Fut: Future<Output = Result<()>>,
473{
474    let cron = event.cron();
475    let ctx = Ctx::new(env);
476    match run(ctx.clone(), cron.clone()).await {
477        Ok(()) => worker::console_log!("{CRON_LOG_PREFIX} {cron} done"),
478        Err(err) => {
479            worker::console_error!("{CRON_LOG_PREFIX} {cron} failed: {err}");
480            report(&ctx, "ocre.cron", &err, [("cron", cron.as_str())]);
481        }
482    }
483    super::errors::flush(&ctx).await;
484}
485
486/// Takes the lock `key` for `owner` until `ttl` seconds from now; `false`
487/// when another owner holds it and it has not expired. The owner holding it
488/// takes it again (extending it), so a job that continues in several queue
489/// messages keeps its lock from step to step. One D1 write.
490///
491/// The `job_locks` table comes from the migration of `ocre g job --lock`
492/// ([`LOCKS_TABLE_SQL`](crate::jobs::LOCKS_TABLE_SQL)).
493///
494/// # Errors
495///
496/// A D1 error (the table is missing: run the migration).
497///
498/// # Examples
499///
500/// ```no_run
501/// use ocre::{Ctx, Error, Result, jobs};
502///
503/// async fn import(ctx: &Ctx, account_id: i64, run: &str) -> Result<()> {
504///     let db = ctx.db()?;
505///     let key = format!("import:{account_id}");
506///     if !jobs::lock(&db, &key, run, 600).await? {
507///         return Err(Error::internal("another import of this account is running")); // retried later
508///     }
509///     // ... the work
510///     jobs::unlock(&db, &key, run).await
511/// }
512/// # let _ = import;
513/// ```
514pub async fn lock(db: &crate::Db, key: &str, owner: &str, ttl: i64) -> Result<bool> {
515    let now = crate::now();
516    let taken = db
517        .execute(
518            "INSERT INTO job_locks (key, owner, expires_at) VALUES (?1, ?2, ?3) \
519             ON CONFLICT (key) DO UPDATE SET owner = ?2, expires_at = ?3 \
520             WHERE job_locks.owner = ?2 OR job_locks.expires_at < ?4",
521            crate::params![key, owner, now + ttl, now],
522        )
523        .await?;
524    Ok(taken > 0)
525}
526
527/// Releases the lock `key` if `owner` holds it. One D1 write.
528///
529/// # Errors
530///
531/// A D1 error.
532pub async fn unlock(db: &crate::Db, key: &str, owner: &str) -> Result<()> {
533    db.execute("DELETE FROM job_locks WHERE key = ?1 AND owner = ?2", crate::params![key, owner]).await?;
534    Ok(())
535}