Skip to main content

renox_core/queue/
worker.rs

1use std::sync::Arc;
2use std::sync::atomic::{AtomicI64, Ordering};
3use std::time::Duration;
4
5use tokio::sync::watch;
6use tokio::task::JoinSet;
7
8use super::{Handlers, JobContext, Middleware, MiddlewareKind, unix_now};
9use crate::AppState;
10use crate::db::Dialect;
11
12/// A reserved job older than this is assumed abandoned (e.g. the process
13/// crashed) and becomes available again. Jobs with a longer `TIMEOUT` get
14/// their reservation extended (`extend_reservation`).
15const RESERVATION: i64 = 15 * 60;
16const POLL: Duration = Duration::from_secs(1);
17/// The error a job gets when its last attempt never finished.
18const CRASHED: &str = "the worker stopped during the last attempt (crash, kill or out of memory)";
19
20struct Reserved {
21    id: i64,
22    queue: String,
23    job: String,
24    payload: String,
25    attempts: u32,
26    max_attempts: u32,
27    chain: Option<String>,
28    batch_id: Option<i64>,
29    /// The batch this `then`/`catch`/`finally` job follows.
30    callback_of: Option<i64>,
31}
32
33/// Runs queued jobs.
34#[derive(Clone)]
35pub struct Worker {
36    state: AppState,
37    handlers: Handlers,
38    queues: Vec<String>,
39    /// Unix seconds of the last `sweep_exhausted`.
40    last_sweep: Arc<AtomicI64>,
41}
42
43/// Why an attempt failed, and whether another attempt could succeed.
44struct Failure {
45    error: String,
46    permanent: bool,
47}
48
49impl Failure {
50    fn retry(error: String) -> Self {
51        Self {
52            error,
53            permanent: false,
54        }
55    }
56
57    fn permanent(error: String) -> Self {
58        Self {
59            error,
60            permanent: true,
61        }
62    }
63}
64
65/// What became of a reserved job.
66enum Outcome {
67    Done,
68    /// Its batch was cancelled: not run, counted as done.
69    Skipped,
70    /// Middleware held it back: queued again after the wait, the attempt
71    /// not counted.
72    Released(Duration),
73    Failed(Failure),
74}
75
76fn panic_message(err: tokio::task::JoinError) -> String {
77    match err.try_into_panic() {
78        Ok(panic) => format!("the job panicked: {}", crate::error::panic_message(&*panic)),
79        Err(err) => format!("the job was cancelled: {err}"),
80    }
81}
82
83async fn fail(tx: &mut crate::db::Transaction, job: &Reserved, error: &str) -> crate::Result {
84    crate::db::sql(
85        "INSERT INTO failed_jobs (queue, job, payload, max_attempts, error, failed_at, chain, \
86         batch_id, callback_of) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
87    )
88    .bind(&job.queue)
89    .bind(&job.job)
90    .bind(&job.payload)
91    .bind(i64::from(job.max_attempts))
92    .bind(error)
93    .bind(unix_now())
94    .bind(job.chain.clone())
95    .bind(job.batch_id)
96    .bind(job.callback_of)
97    .execute(&mut *tx)
98    .await?;
99    Ok(())
100}
101
102/// Runs a bookkeeping write, retrying for about 10 seconds while the
103/// database is briefly unavailable (restarting, or SQLite busy).
104async fn retry_write<F, Fut>(mut write: F) -> crate::Result
105where
106    F: FnMut() -> Fut,
107    Fut: std::future::Future<Output = crate::Result>,
108{
109    let mut delays = [100u64, 400, 1000, 2500, 6000].into_iter();
110    loop {
111        match write().await {
112            Ok(()) => return Ok(()),
113            Err(err) => match delays.next() {
114                Some(ms) => {
115                    tracing::warn!(error = ?err, "queue bookkeeping failed, retrying");
116                    tokio::time::sleep(Duration::from_millis(ms)).await;
117                }
118                None => return Err(err),
119            },
120        }
121    }
122}
123
124impl Worker {
125    pub(crate) fn new(state: AppState, handlers: Handlers, queues: Vec<String>) -> Self {
126        Self {
127            state,
128            handlers,
129            queues,
130            last_sweep: Arc::default(),
131        }
132    }
133
134    async fn reserve(&self) -> crate::Result<Option<Reserved>> {
135        let now = unix_now();
136        let (queues, order) = if self.queues.is_empty() {
137            (String::new(), String::new())
138        } else {
139            let marks = vec!["?"; self.queues.len()].join(", ");
140            // Listed first, drained first.
141            let whens: String = (0..self.queues.len())
142                .map(|i| format!(" WHEN ? THEN {i}"))
143                .collect();
144            (
145                format!(" AND queue IN ({marks})"),
146                format!("CASE queue{whens} END, "),
147            )
148        };
149        // SQLite runs one write at a time, so the UPDATE is enough. On
150        // PostgreSQL, workers on several servers must not pick the same row.
151        let lock = match self.state.db.dialect() {
152            Dialect::Sqlite => "",
153            Dialect::Postgres => " FOR UPDATE SKIP LOCKED",
154        };
155        let sql = format!(
156            "UPDATE jobs SET reserved_at = ?, attempts = attempts + 1 WHERE id = (\
157                SELECT id FROM jobs WHERE available_at <= ? \
158                AND (reserved_at IS NULL OR reserved_at <= ?) \
159                AND attempts < max_attempts{queues} \
160                ORDER BY {order}available_at, id LIMIT 1{lock}) \
161             RETURNING id, queue, job, payload, attempts, max_attempts, chain, batch_id, callback_of"
162        );
163        let mut query = crate::db::sql(sql)
164            .bind(now)
165            .bind(now)
166            .bind(now - RESERVATION);
167        for queue in self.queues.iter().chain(&self.queues) {
168            query = query.bind(queue);
169        }
170        let Some(row) = query.fetch_optional(&self.state.db).await? else {
171            return Ok(None);
172        };
173        Ok(Some(Reserved {
174            id: row.try_get("id")?,
175            queue: row.try_get("queue")?,
176            job: row.try_get("job")?,
177            payload: row.try_get("payload")?,
178            attempts: row.try_get::<i64>("attempts")? as u32,
179            max_attempts: row.try_get::<i64>("max_attempts")? as u32,
180            chain: row.try_get("chain")?,
181            batch_id: row.try_get("batch_id")?,
182            callback_of: row.try_get("callback_of")?,
183        }))
184    }
185
186    async fn batch_cancelled(&self, batch_id: Option<i64>) -> crate::Result<bool> {
187        let Some(id) = batch_id else {
188            return Ok(false);
189        };
190        let cancelled: Option<Option<i64>> =
191            crate::db::sql("SELECT cancelled_at FROM job_batches WHERE id = ?")
192                .bind(id)
193                .scalar_optional(&self.state.db)
194                .await?;
195        Ok(cancelled.flatten().is_some())
196    }
197
198    /// Checks the job's middleware; `Some(wait)` puts it back for `wait`.
199    /// A lock it takes is returned to be held while the job runs.
200    async fn check_middleware(
201        &self,
202        middleware: Vec<Middleware>,
203        timeout: Duration,
204    ) -> crate::Result<(Option<Duration>, Vec<crate::cache::LockGuard>)> {
205        let mut guards = Vec::new();
206        for Middleware(kind) in middleware {
207            match kind {
208                MiddlewareKind::WithoutOverlapping { key, release_after } => {
209                    let lock = self
210                        .state
211                        .cache
212                        .lock(&format!("overlap:{key}"), timeout + Duration::from_secs(60));
213                    match lock.try_acquire().await? {
214                        Some(guard) => guards.push(guard),
215                        None => return Ok((Some(release_after), guards)),
216                    }
217                }
218                MiddlewareKind::RateLimited { key, max, per } => {
219                    let per = per.as_secs() as i64;
220                    let now = unix_now();
221                    let window = now - now.rem_euclid(per);
222                    let counter = format!("renox:rate:{key}:{window}");
223                    let cache = &self.state.cache;
224                    let ttl = Duration::from_secs(per as u64 + 1);
225                    cache.add(&counter, &0, Some(ttl)).await?;
226                    if cache.increment(&counter, 1).await? > i64::from(max) {
227                        let wait = (window + per - now).max(1) as u64;
228                        return Ok((Some(Duration::from_secs(wait)), guards));
229                    }
230                }
231            }
232        }
233        Ok((None, guards))
234    }
235
236    /// Runs one available job; returns whether there was one.
237    pub async fn run_next(&self) -> crate::Result<bool> {
238        self.sweep_exhausted().await?;
239        let Some(job) = self.reserve().await? else {
240            return Ok(false);
241        };
242        let handler = self.handlers.get(job.job.as_str()).cloned();
243        if let Some(handler) = &handler {
244            self.extend_reservation(&job, handler.timeout).await?;
245        }
246        let plain = self.state.queue.open(&job.payload);
247        let unique = match (&handler, &plain) {
248            (Some(handler), Ok(plain)) => (handler.unique_key)(plain),
249            _ => None,
250        };
251        let mut guards = Vec::new();
252        let outcome = match (&handler, plain.as_ref()) {
253            (None, _) => Outcome::Failed(Failure::permanent(format!(
254                "no handler registered for job `{}`",
255                job.job
256            ))),
257            (Some(_), Err(err)) => Outcome::Failed(Failure::permanent(format!("{err:?}"))),
258            (Some(_), Ok(_)) if self.batch_cancelled(job.batch_id).await? => Outcome::Skipped,
259            (Some(handler), Ok(plain)) => {
260                let (release, held) = self
261                    .check_middleware((handler.middleware)(plain), handler.timeout)
262                    .await?;
263                guards = held;
264                match release {
265                    Some(wait) => Outcome::Released(wait),
266                    None => self.attempt(handler, &job, plain.clone()).await,
267                }
268            }
269        };
270
271        // The job ran; don't let a brief database hiccup strand it until its
272        // reservation expires.
273        let backoff = handler
274            .as_ref()
275            .map(|h| (h.backoff)(job.attempts))
276            .unwrap_or_default();
277        retry_write(|| self.record(&job, &outcome, backoff, unique.as_deref())).await?;
278        for guard in guards {
279            let _ = guard.release().await;
280        }
281        match outcome {
282            Outcome::Done => tracing::info!(job = %job.job, id = job.id, "job done"),
283            Outcome::Skipped => {
284                tracing::info!(job = %job.job, id = job.id, "job skipped: its batch was cancelled");
285            }
286            Outcome::Released(wait) => {
287                tracing::debug!(job = %job.job, id = job.id, ?wait, "job held back by its middleware");
288            }
289            Outcome::Failed(failure) if !failure.permanent && job.attempts < job.max_attempts => {
290                tracing::warn!(job = %job.job, id = job.id, attempt = job.attempts, error = %failure.error, "job failed, will retry");
291            }
292            Outcome::Failed(failure) => {
293                tracing::error!(job = %job.job, id = job.id, error = %failure.error, "job failed for good");
294                self.failed_for_good(&job, plain.ok(), failure.error).await;
295            }
296        }
297        Ok(true)
298    }
299
300    /// Reports a job that failed for good (`App::report`) and runs its
301    /// `failed` hook (when its handler is known and its payload opens).
302    async fn failed_for_good(&self, job: &Reserved, plain: Option<String>, error: String) {
303        let report = crate::report::ErrorReport::new(
304            &self.state,
305            crate::report::ReportKind::Job,
306            error.lines().next().unwrap_or_default().to_owned(),
307            error.clone(),
308            Some(format!("{} #{}", job.job, job.id)),
309        );
310        crate::report::send(&self.state, report);
311        if let (Some(handler), Some(plain)) = (self.handlers.get(job.job.as_str()), plain) {
312            let ctx = JobContext {
313                state: self.state.clone(),
314                attempt: job.attempts,
315                id: job.id,
316                batch_id: job.batch_id.or(job.callback_of),
317            };
318            let hook = (handler.failed)(plain, ctx, error);
319            let hook = crate::context::scope_app(self.state.clone(), hook);
320            if tokio::spawn(crate::clock::carry(hook)).await.is_err() {
321                tracing::error!(job = %job.job, id = job.id, "the job's failed hook panicked");
322            }
323        }
324    }
325
326    /// One attempt, in its own task so a panic is a failed attempt, not a
327    /// dead worker.
328    async fn attempt(
329        &self,
330        handler: &super::JobHandler,
331        job: &Reserved,
332        payload: String,
333    ) -> Outcome {
334        let ctx = JobContext {
335            state: self.state.clone(),
336            attempt: job.attempts,
337            id: job.id,
338            batch_id: job.batch_id.or(job.callback_of),
339        };
340        let run = (handler.run)(payload, ctx);
341        let run = crate::context::scope_app(self.state.clone(), run);
342        let mut task = tokio::spawn(crate::clock::carry(run));
343        match tokio::time::timeout(handler.timeout, &mut task).await {
344            Ok(Ok(Ok(()))) => Outcome::Done,
345            Ok(Ok(Err(err))) => Outcome::Failed(Failure {
346                permanent: err.is_permanent(),
347                error: format!("{err:?}"),
348            }),
349            Ok(Err(join)) => Outcome::Failed(Failure::retry(panic_message(join))),
350            Err(_) => {
351                task.abort();
352                Outcome::Failed(Failure::retry(format!(
353                    "timed out after {:?}",
354                    handler.timeout
355                )))
356            }
357        }
358    }
359
360    /// Writes a job's outcome: deleted (with its chain, batch and unique
361    /// claim moved on), back on the queue, or failed.
362    async fn record(
363        &self,
364        job: &Reserved,
365        outcome: &Outcome,
366        backoff: Duration,
367        unique: Option<&str>,
368    ) -> crate::Result {
369        let db = &self.state.db;
370        let failure = match outcome {
371            Outcome::Released(wait) => {
372                crate::db::sql(
373                    "UPDATE jobs SET reserved_at = NULL, available_at = ?, \
374                     attempts = attempts - 1 WHERE id = ?",
375                )
376                .bind(unix_now() + wait.as_secs() as i64)
377                .bind(job.id)
378                .execute(db)
379                .await?;
380                return Ok(());
381            }
382            Outcome::Failed(failure) if !failure.permanent && job.attempts < job.max_attempts => {
383                crate::db::sql("UPDATE jobs SET reserved_at = NULL, available_at = ? WHERE id = ?")
384                    .bind(unix_now() + backoff.as_secs() as i64)
385                    .bind(job.id)
386                    .execute(db)
387                    .await?;
388                return Ok(());
389            }
390            Outcome::Failed(failure) => Some(failure),
391            Outcome::Done | Outcome::Skipped => None,
392        };
393        let mut tx = db.begin().await?;
394        let deleted = crate::db::sql("DELETE FROM jobs WHERE id = ?")
395            .bind(job.id)
396            .execute(&mut tx)
397            .await?;
398        if deleted > 0 {
399            match failure {
400                Some(failure) => fail(&mut tx, job, &failure.error).await?,
401                None => {
402                    if let (Outcome::Done, Some(chain)) = (outcome, &job.chain) {
403                        super::continue_chain(&mut tx, chain).await?;
404                    }
405                }
406            }
407            if let Some(batch) = job.batch_id {
408                super::batch_job_done(&mut tx, batch, failure.is_some()).await?;
409            }
410            if let Some(key) = unique {
411                super::release_unique(&mut tx, key).await?;
412            }
413            if !matches!(outcome, Outcome::Skipped) {
414                super::dashboard::count_finished(&mut tx, failure.is_some()).await?;
415            }
416        }
417        tx.commit().await?;
418        if deleted > 0 {
419            self.state.queue.wake_workers();
420        }
421        Ok(())
422    }
423
424    /// A job whose attempts may outlast the reservation keeps its row for
425    /// the attempt's timeout plus a minute, so no other worker starts it again.
426    async fn extend_reservation(&self, job: &Reserved, timeout: Duration) -> crate::Result {
427        let needed = timeout.as_secs() as i64 + 60;
428        if needed > RESERVATION {
429            crate::db::sql("UPDATE jobs SET reserved_at = ? WHERE id = ?")
430                .bind(unix_now() + (needed - RESERVATION))
431                .bind(job.id)
432                .execute(&self.state.db)
433                .await?;
434        }
435        Ok(())
436    }
437
438    /// Moves jobs whose last attempt never finished (the process crashed,
439    /// was killed or ran out of memory) to `failed_jobs`, at most once a minute.
440    async fn sweep_exhausted(&self) -> crate::Result {
441        let now = unix_now();
442        let last = self.last_sweep.load(Ordering::Relaxed);
443        if now - last < 60
444            || self
445                .last_sweep
446                .compare_exchange(last, now, Ordering::Relaxed, Ordering::Relaxed)
447                .is_err()
448        {
449            return Ok(());
450        }
451        let mut tx = self.state.db.begin().await?;
452        let rows = crate::db::sql(
453            "DELETE FROM jobs WHERE attempts >= max_attempts AND reserved_at IS NOT NULL \
454             AND reserved_at <= ? \
455             RETURNING id, queue, job, payload, attempts, max_attempts, chain, batch_id, callback_of",
456        )
457        .bind(now - RESERVATION)
458        .fetch_all(&mut tx)
459        .await?;
460        let mut swept = Vec::new();
461        for row in &rows {
462            let job = Reserved {
463                id: row.try_get("id")?,
464                queue: row.try_get("queue")?,
465                job: row.try_get("job")?,
466                payload: row.try_get("payload")?,
467                attempts: row.try_get::<i64>("attempts")? as u32,
468                max_attempts: row.try_get::<i64>("max_attempts")? as u32,
469                chain: row.try_get("chain")?,
470                batch_id: row.try_get("batch_id")?,
471                callback_of: row.try_get("callback_of")?,
472            };
473            fail(&mut tx, &job, CRASHED).await?;
474            if let Some(batch) = job.batch_id {
475                super::batch_job_done(&mut tx, batch, true).await?;
476            }
477            let unique = self.handlers.get(job.job.as_str()).and_then(|handler| {
478                let plain = self.state.queue.open(&job.payload).ok()?;
479                (handler.unique_key)(&plain)
480            });
481            if let Some(key) = unique {
482                super::release_unique(&mut tx, &key).await?;
483            }
484            tracing::error!(job = %job.job, "job's last attempt never finished; moved to failed_jobs");
485            swept.push(job);
486        }
487        tx.commit().await?;
488        // Like a last attempt that returned an error: reported, and the
489        // job's `failed` hook runs.
490        for job in swept {
491            let plain = self.state.queue.open(&job.payload).ok();
492            self.failed_for_good(&job, plain, CRASHED.to_owned()).await;
493        }
494        Ok(())
495    }
496
497    /// Runs every job that is available now; returns how many ran.
498    pub async fn drain(&self) -> crate::Result<usize> {
499        let mut ran = 0;
500        while self.run_next().await? {
501            ran += 1;
502        }
503        Ok(ran)
504    }
505
506    /// Runs `concurrency` loops until `shutdown` flips to true, then lets
507    /// running jobs finish.
508    pub async fn run(self, concurrency: usize, mut shutdown: watch::Receiver<bool>) {
509        let mut loops = JoinSet::new();
510        for _ in 0..concurrency.max(1) {
511            let worker = self.clone();
512            let mut stop = shutdown.clone();
513            loops.spawn(async move {
514                let wake = worker.state.queue.wake();
515                loop {
516                    if *stop.borrow() {
517                        break;
518                    }
519                    // Listen before looking for work: a job dispatched while
520                    // the query runs still wakes this loop instead of waiting
521                    // for the next poll.
522                    let notified = wake.notified();
523                    tokio::pin!(notified);
524                    notified.as_mut().enable();
525                    match worker.run_next().await {
526                        Ok(true) => continue,
527                        Ok(false) => {}
528                        Err(err) => tracing::error!(error = ?err, "queue worker error"),
529                    }
530                    tokio::select! {
531                        _ = notified => {}
532                        _ = tokio::time::sleep(POLL) => {}
533                        _ = stop.changed() => {}
534                    }
535                }
536            });
537        }
538        let _ = shutdown.changed().await;
539        while loops.join_next().await.is_some() {}
540    }
541}
542
543#[cfg(test)]
544mod tests {
545    use std::sync::Arc;
546    use std::sync::atomic::{AtomicUsize, Ordering};
547
548    use super::*;
549
550    #[tokio::test]
551    async fn an_aborted_attempt_says_it_was_cancelled() {
552        let task = tokio::spawn(std::future::pending::<()>());
553        task.abort();
554        let err = task.await.unwrap_err();
555        assert!(panic_message(err).starts_with("the job was cancelled"));
556    }
557
558    /// Bookkeeping retries while the database is briefly away: a write
559    /// that works on its second try succeeds; one that never works gives
560    /// up after the last delay with its error.
561    #[tokio::test(start_paused = true)]
562    async fn bookkeeping_writes_are_retried_then_given_up() {
563        let tries = Arc::new(AtomicUsize::new(0));
564        let counted = tries.clone();
565        retry_write(|| {
566            let n = counted.fetch_add(1, Ordering::SeqCst);
567            async move {
568                if n == 0 {
569                    Err(crate::Error::Internal(anyhow::anyhow!(
570                        "database is locked"
571                    )))
572                } else {
573                    Ok(())
574                }
575            }
576        })
577        .await
578        .unwrap();
579        assert_eq!(tries.load(Ordering::SeqCst), 2);
580
581        let tries = Arc::new(AtomicUsize::new(0));
582        let counted = tries.clone();
583        let started = tokio::time::Instant::now();
584        let err = retry_write(|| {
585            counted.fetch_add(1, Ordering::SeqCst);
586            async { Err(crate::Error::Internal(anyhow::anyhow!("database is gone"))) }
587        })
588        .await
589        .unwrap_err();
590        assert!(format!("{err:?}").contains("database is gone"));
591        assert_eq!(
592            tries.load(Ordering::SeqCst),
593            6,
594            "the first try and five retries"
595        );
596        assert_eq!(started.elapsed(), Duration::from_millis(10_000));
597    }
598}