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}