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
12const RESERVATION: i64 = 15 * 60;
16const POLL: Duration = Duration::from_secs(1);
17const 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 callback_of: Option<i64>,
31}
32
33#[derive(Clone)]
35pub struct Worker {
36 state: AppState,
37 handlers: Handlers,
38 queues: Vec<String>,
39 last_sweep: Arc<AtomicI64>,
41}
42
43struct 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
65enum Outcome {
67 Done,
68 Skipped,
70 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
102async 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 #[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}