1use serde_json::Value;
20use sqlx::Row;
21use sqlx::sqlite::SqliteRow;
22use tracing::{debug, info, warn};
23
24use crate::sqlite::db::Database;
25use crate::sqlite::nonce::now_secs;
26
27#[derive(Debug, Clone)]
35pub struct Job {
36 pub id: String,
37 pub kind: String,
38 pub dedup_key: String,
39 pub payload: Value,
40 pub status: String,
41 pub run_at: i64,
42 pub attempts: i64,
43 pub max_attempts: i64,
44 pub deadline: Option<i64>,
45 pub lease_until: Option<i64>,
46 pub lease_owner: Option<String>,
47 pub last_error: Option<String>,
48 pub created_at: i64,
49 pub updated_at: i64,
50}
51
52const COLUMNS: &str = "id, kind, dedup_key, payload, status, run_at, attempts, max_attempts, \
55 deadline, lease_until, lease_owner, last_error, created_at, updated_at";
56
57#[derive(Debug)]
66pub struct NewJob<'a> {
67 pub id: &'a str,
68 pub kind: &'a str,
69 pub dedup_key: &'a str,
70 pub payload: &'a Value,
71 pub run_at: i64,
72 pub deadline: Option<i64>,
73 pub max_attempts: i64,
76}
77
78impl Job {
79 fn from_row(row: SqliteRow) -> Result<Self, sqlx::Error> {
80 let payload_json: String = row.try_get("payload")?;
81 let payload: Value = serde_json::from_str(&payload_json)
82 .map_err(|error| sqlx::Error::Decode(Box::new(error)))?;
83
84 Ok(Self {
85 id: row.try_get("id")?,
86 kind: row.try_get("kind")?,
87 dedup_key: row.try_get("dedup_key")?,
88 payload,
89 status: row.try_get("status")?,
90 run_at: row.try_get("run_at")?,
91 attempts: row.try_get("attempts")?,
92 max_attempts: row.try_get("max_attempts")?,
93 deadline: row.try_get("deadline")?,
94 lease_until: row.try_get("lease_until")?,
95 lease_owner: row.try_get("lease_owner")?,
96 last_error: row.try_get("last_error")?,
97 created_at: row.try_get("created_at")?,
98 updated_at: row.try_get("updated_at")?,
99 })
100 }
101
102 pub async fn enqueue(row: NewJob<'_>, database: &Database) -> Result<bool, sqlx::Error> {
111 let now = now_secs();
112 let queued = sqlx::query(
113 "INSERT OR IGNORE INTO jobs \
114 (id, kind, dedup_key, payload, status, run_at, attempts, max_attempts, \
115 deadline, created_at, updated_at) \
116 VALUES (?, ?, ?, ?, 'ready', ?, 0, ?, ?, ?, ?);",
117 )
118 .bind(row.id)
119 .bind(row.kind)
120 .bind(row.dedup_key)
121 .bind(row.payload.to_string())
122 .bind(row.run_at)
123 .bind(row.max_attempts)
124 .bind(row.deadline)
125 .bind(now)
126 .bind(now)
127 .execute(&database.pool)
128 .await?
129 .rows_affected()
130 == 1;
131
132 if queued {
133 debug!(
134 event = "db_job_enqueued",
135 outcome = "success",
136 job_id = %row.id,
137 job_kind = %row.kind,
138 dedup_key = %row.dedup_key,
139 run_at = row.run_at,
140 );
141 }
142 Ok(queued)
143 }
144
145 pub async fn claim_next(
161 runner_id: &str,
162 kinds: &[&str],
163 lease_until: i64,
164 now: i64,
165 database: &Database,
166 ) -> Result<Option<Self>, sqlx::Error> {
167 if kinds.is_empty() {
168 return Ok(None);
169 }
170
171 let placeholders = std::iter::repeat_n("?", kinds.len())
177 .collect::<Vec<_>>()
178 .join(", ");
179 let sql = format!(
180 "UPDATE jobs \
181 SET status = 'running', attempts = attempts + 1, lease_owner = ?, \
182 lease_until = ?, updated_at = ? \
183 WHERE id = (SELECT id FROM jobs \
184 WHERE status = 'ready' AND run_at <= ? AND kind IN ({placeholders}) \
185 ORDER BY run_at ASC, created_at ASC LIMIT 1) \
186 AND status = 'ready' \
187 RETURNING {COLUMNS};"
188 );
189
190 let mut query = sqlx::query(sqlx::AssertSqlSafe(sql))
191 .bind(runner_id)
192 .bind(lease_until)
193 .bind(now)
194 .bind(now);
195 for kind in kinds {
196 query = query.bind(*kind);
197 }
198 let row = query.fetch_optional(&database.pool).await?;
199 let job = row.map(Self::from_row).transpose()?;
200
201 if let Some(job) = &job {
202 debug!(
203 event = "db_job_claimed",
204 outcome = "success",
205 job_id = %job.id,
206 job_kind = %job.kind,
207 attempts = job.attempts,
208 lease_until = lease_until,
209 );
210 }
211 Ok(job)
212 }
213
214 pub async fn complete(
216 id: &str,
217 runner_id: &str,
218 database: &Database,
219 ) -> Result<bool, sqlx::Error> {
220 Self::settle(id, runner_id, "done", None, None, false, database).await
221 }
222
223 pub async fn retry(
225 id: &str,
226 runner_id: &str,
227 run_at: i64,
228 error: &str,
229 database: &Database,
230 ) -> Result<bool, sqlx::Error> {
231 Self::settle(
232 id,
233 runner_id,
234 "ready",
235 Some(run_at),
236 Some(error),
237 false,
238 database,
239 )
240 .await
241 }
242
243 pub async fn reschedule(
251 id: &str,
252 runner_id: &str,
253 run_at: i64,
254 database: &Database,
255 ) -> Result<bool, sqlx::Error> {
256 Self::settle(id, runner_id, "ready", Some(run_at), None, true, database).await
257 }
258
259 pub async fn abandon(
261 id: &str,
262 runner_id: &str,
263 error: &str,
264 database: &Database,
265 ) -> Result<bool, sqlx::Error> {
266 Self::settle(id, runner_id, "failed", None, Some(error), false, database).await
267 }
268
269 async fn settle(
277 id: &str,
278 runner_id: &str,
279 status: &str,
280 run_at: Option<i64>,
281 error: Option<&str>,
282 reset_attempts: bool,
283 database: &Database,
284 ) -> Result<bool, sqlx::Error> {
285 let now = now_secs();
286 let sql = format!(
287 "UPDATE jobs \
288 SET status = ?, last_error = ?, lease_owner = NULL, lease_until = NULL, \
289 run_at = COALESCE(?, run_at), updated_at = ?{} \
290 WHERE id = ? AND status = 'running' AND lease_owner = ?;",
291 if reset_attempts { ", attempts = 0" } else { "" }
292 );
293 let settled = sqlx::query(sqlx::AssertSqlSafe(sql))
294 .bind(status)
295 .bind(error)
296 .bind(run_at)
297 .bind(now)
298 .bind(id)
299 .bind(runner_id)
300 .execute(&database.pool)
301 .await?
302 .rows_affected()
303 == 1;
304 Ok(settled)
305 }
306
307 pub async fn reclaim_expired(now: i64, database: &Database) -> Result<u64, sqlx::Error> {
313 let reclaimed = sqlx::query(
314 "UPDATE jobs SET status = 'ready', lease_owner = NULL, lease_until = NULL, \
315 updated_at = ? WHERE status = 'running' AND lease_until <= ?;",
316 )
317 .bind(now)
318 .bind(now)
319 .execute(&database.pool)
320 .await?
321 .rows_affected();
322
323 if reclaimed > 0 {
324 warn!(
325 event = "db_job_leases_reclaimed",
326 outcome = "advisory",
327 rows_reclaimed = reclaimed,
328 );
329 }
330 Ok(reclaimed)
331 }
332
333 pub async fn release_owned(runner_id: &str, database: &Database) -> Result<u64, sqlx::Error> {
338 let released = sqlx::query(
339 "UPDATE jobs SET status = 'ready', lease_owner = NULL, lease_until = NULL, \
340 updated_at = ? WHERE status = 'running' AND lease_owner = ?;",
341 )
342 .bind(now_secs())
343 .bind(runner_id)
344 .execute(&database.pool)
345 .await?
346 .rows_affected();
347
348 if released > 0 {
349 debug!(
350 event = "db_job_leases_released",
351 outcome = "success",
352 rows_released = released,
353 );
354 }
355 Ok(released)
356 }
357
358 pub async fn find_by_id(id: &str, database: &Database) -> Result<Option<Self>, sqlx::Error> {
360 let mut query = sqlx::QueryBuilder::new(format!("SELECT {COLUMNS} FROM jobs WHERE id = "));
364 query.push_bind(id);
365 let row = query.build().fetch_optional(&database.pool).await?;
366 row.map(Self::from_row).transpose()
367 }
368
369 pub async fn find_live(
371 kind: &str,
372 dedup_key: &str,
373 database: &Database,
374 ) -> Result<Option<Self>, sqlx::Error> {
375 let mut query = sqlx::QueryBuilder::new(format!(
376 "SELECT {COLUMNS} FROM jobs \
377 WHERE status IN ('ready', 'running') AND kind = "
378 ));
379 query.push_bind(kind);
380 query.push(" AND dedup_key = ");
381 query.push_bind(dedup_key);
382 let row = query.build().fetch_optional(&database.pool).await?;
383 row.map(Self::from_row).transpose()
384 }
385
386 pub async fn count_live(kind: &str, database: &Database) -> Result<i64, sqlx::Error> {
393 let (count,): (i64,) = sqlx::query_as(
394 "SELECT COUNT(*) FROM jobs WHERE status IN ('ready', 'running') AND kind = ?;",
395 )
396 .bind(kind)
397 .fetch_one(&database.pool)
398 .await?;
399 Ok(count)
400 }
401
402 pub async fn cleanup(cutoff: i64, database: &Database) -> Result<u64, sqlx::Error> {
407 let deleted = sqlx::query(
408 "DELETE FROM jobs \
409 WHERE status IN ('done', 'failed', 'cancelled') AND updated_at < ?;",
410 )
411 .bind(cutoff)
412 .execute(&database.pool)
413 .await?
414 .rows_affected();
415 info!(
416 event = "db_job_cleanup_completed",
417 outcome = "success",
418 rows_removed = deleted,
419 cutoff = cutoff,
420 );
421 Ok(deleted)
422 }
423}
424
425#[cfg(test)]
426mod tests {
427 use super::*;
428 use serde_json::json;
429 use std::sync::Arc;
430
431 async fn db() -> Arc<Database> {
432 Arc::new(Database::connect_in_memory().await.unwrap())
433 }
434
435 async fn enqueue(id: &str, key: &str, run_at: i64, database: &Database) -> bool {
437 Job::enqueue(
438 NewJob {
439 id,
440 kind: "test",
441 dedup_key: key,
442 payload: &json!({"n": 1}),
443 run_at,
444 deadline: None,
445 max_attempts: 3,
446 },
447 database,
448 )
449 .await
450 .unwrap()
451 }
452
453 async fn claim(runner: &str, database: &Database) -> Option<Job> {
454 Job::claim_next(runner, &["test"], now_secs() + 60, now_secs(), database)
455 .await
456 .unwrap()
457 }
458
459 #[tokio::test]
460 async fn a_queued_job_round_trips_through_every_column() {
461 let database = db().await;
462 let deadline = now_secs() + 900;
463 assert!(
464 Job::enqueue(
465 NewJob {
466 id: "job-1",
467 kind: "relay",
468 dedup_key: "ord-1",
469 payload: &json!({"order_id": "ord-1"}),
470 run_at: 1_234,
471 deadline: Some(deadline),
472 max_attempts: 7,
473 },
474 &database,
475 )
476 .await
477 .unwrap()
478 );
479
480 let job = Job::find_by_id("job-1", &database).await.unwrap().unwrap();
481 assert_eq!(job.kind, "relay");
482 assert_eq!(job.dedup_key, "ord-1");
483 assert_eq!(job.payload, json!({"order_id": "ord-1"}));
484 assert_eq!(job.status, "ready");
485 assert_eq!(job.run_at, 1_234);
486 assert_eq!(job.attempts, 0);
487 assert_eq!(job.max_attempts, 7);
488 assert_eq!(job.deadline, Some(deadline));
489 assert!(job.lease_until.is_none());
490 assert!(job.lease_owner.is_none());
491 assert!(job.last_error.is_none());
492 }
493
494 #[tokio::test]
497 async fn a_live_job_holds_its_identity_and_a_settled_one_releases_it() {
498 let database = db().await;
499 assert!(enqueue("job-1", "ord-1", now_secs(), &database).await);
500 assert!(
501 !enqueue("job-2", "ord-1", now_secs(), &database).await,
502 "a second live job must not take an identity already held"
503 );
504
505 let job = claim("runner-a", &database).await.unwrap();
507 assert!(
508 !enqueue("job-3", "ord-1", now_secs(), &database).await,
509 "'running' holds the identity exactly as 'ready' does"
510 );
511
512 Job::complete(&job.id, "runner-a", &database).await.unwrap();
513 assert!(
514 enqueue("job-4", "ord-1", now_secs(), &database).await,
515 "a settled job releases its identity"
516 );
517 }
518
519 #[tokio::test]
520 async fn claiming_takes_the_oldest_eligible_row_and_only_once() {
521 let database = db().await;
522 assert!(enqueue("job-new", "b", now_secs() - 10, &database).await);
523 assert!(enqueue("job-old", "a", now_secs() - 100, &database).await);
524
525 let first = claim("runner-a", &database).await.unwrap();
526 assert_eq!(first.id, "job-old", "oldest `run_at` first");
527 assert_eq!(first.attempts, 1, "the counter moves at claim");
528 assert_eq!(first.lease_owner.as_deref(), Some("runner-a"));
529
530 let second = claim("runner-b", &database).await.unwrap();
531 assert_eq!(second.id, "job-new");
532 assert!(
533 claim("runner-c", &database).await.is_none(),
534 "a claimed row is not claimable again"
535 );
536 }
537
538 #[tokio::test]
539 async fn claiming_skips_a_job_whose_run_at_has_not_arrived() {
540 let database = db().await;
541 assert!(enqueue("job-1", "a", now_secs() + 3_600, &database).await);
542 assert!(claim("runner-a", &database).await.is_none());
543 }
544
545 #[tokio::test]
546 async fn claiming_skips_a_kind_the_runner_does_not_hold() {
547 let database = db().await;
548 assert!(enqueue("job-1", "a", now_secs(), &database).await);
549
550 let claimed = Job::claim_next(
551 "runner-a",
552 &["something-else"],
553 now_secs() + 60,
554 now_secs(),
555 &database,
556 )
557 .await
558 .unwrap();
559 assert!(
560 claimed.is_none(),
561 "an unregistered kind is left alone, not mis-run"
562 );
563 }
564
565 #[tokio::test]
566 async fn claiming_with_no_registered_kinds_asks_the_database_nothing() {
567 let database = db().await;
568 assert!(enqueue("job-1", "a", now_secs(), &database).await);
569 assert!(
570 Job::claim_next("runner-a", &[], now_secs() + 60, now_secs(), &database)
571 .await
572 .unwrap()
573 .is_none()
574 );
575 }
576
577 #[tokio::test]
580 async fn a_settlement_from_a_runner_that_lost_the_lease_writes_nothing() {
581 let database = db().await;
582 assert!(enqueue("job-1", "a", now_secs(), &database).await);
583 let job = claim("runner-a", &database).await.unwrap();
584
585 for settled in [
586 Job::complete(&job.id, "runner-b", &database).await.unwrap(),
587 Job::retry(&job.id, "runner-b", now_secs(), "x", &database)
588 .await
589 .unwrap(),
590 Job::reschedule(&job.id, "runner-b", now_secs(), &database)
591 .await
592 .unwrap(),
593 Job::abandon(&job.id, "runner-b", "x", &database)
594 .await
595 .unwrap(),
596 ] {
597 assert!(!settled, "the lease guard must refuse a foreign settlement");
598 }
599
600 let after = Job::find_by_id("job-1", &database).await.unwrap().unwrap();
601 assert_eq!(after.status, "running");
602 assert_eq!(after.lease_owner.as_deref(), Some("runner-a"));
603 }
604
605 #[tokio::test]
606 async fn retry_returns_the_job_at_its_new_time_and_keeps_the_attempt_count() {
607 let database = db().await;
608 assert!(enqueue("job-1", "a", now_secs(), &database).await);
609 let job = claim("runner-a", &database).await.unwrap();
610
611 assert!(
612 Job::retry(&job.id, "runner-a", 9_999, "upstream timed out", &database)
613 .await
614 .unwrap()
615 );
616
617 let after = Job::find_by_id("job-1", &database).await.unwrap().unwrap();
618 assert_eq!(after.status, "ready");
619 assert_eq!(after.run_at, 9_999);
620 assert_eq!(after.attempts, 1, "a retry does not forgive the attempt");
621 assert_eq!(after.last_error.as_deref(), Some("upstream timed out"));
622 assert!(after.lease_owner.is_none());
623 }
624
625 #[tokio::test]
626 async fn reschedule_resets_the_attempt_count_and_clears_the_last_error() {
627 let database = db().await;
628 assert!(enqueue("job-1", "a", now_secs(), &database).await);
629 let job = claim("runner-a", &database).await.unwrap();
630 Job::retry(&job.id, "runner-a", now_secs(), "a failure", &database)
631 .await
632 .unwrap();
633 let job = claim("runner-a", &database).await.unwrap();
634 assert_eq!(job.attempts, 2);
635
636 assert!(
637 Job::reschedule(&job.id, "runner-a", 5_000, &database)
638 .await
639 .unwrap()
640 );
641
642 let after = Job::find_by_id("job-1", &database).await.unwrap().unwrap();
643 assert_eq!(after.status, "ready");
644 assert_eq!(after.run_at, 5_000);
645 assert_eq!(after.attempts, 0, "a fresh occurrence starts fresh");
646 assert!(after.last_error.is_none());
647 }
648
649 #[tokio::test]
650 async fn abandon_is_terminal_and_records_why() {
651 let database = db().await;
652 assert!(enqueue("job-1", "a", now_secs(), &database).await);
653 let job = claim("runner-a", &database).await.unwrap();
654
655 assert!(
656 Job::abandon(&job.id, "runner-a", "the order no longer exists", &database)
657 .await
658 .unwrap()
659 );
660
661 let after = Job::find_by_id("job-1", &database).await.unwrap().unwrap();
662 assert_eq!(after.status, "failed");
663 assert_eq!(
664 after.last_error.as_deref(),
665 Some("the order no longer exists")
666 );
667 }
668
669 #[tokio::test]
670 async fn reclaim_takes_only_leases_that_have_expired() {
671 let database = db().await;
672 assert!(enqueue("job-live", "a", now_secs(), &database).await);
673 assert!(enqueue("job-dead", "b", now_secs(), &database).await);
674
675 Job::claim_next(
677 "runner-a",
678 &["test"],
679 now_secs() + 600,
680 now_secs(),
681 &database,
682 )
683 .await
684 .unwrap();
685 Job::claim_next("runner-a", &["test"], now_secs() - 1, now_secs(), &database)
686 .await
687 .unwrap();
688
689 assert_eq!(
690 Job::reclaim_expired(now_secs(), &database).await.unwrap(),
691 1
692 );
693
694 let reclaimed = Job::find_by_id("job-dead", &database)
695 .await
696 .unwrap()
697 .unwrap();
698 assert_eq!(reclaimed.status, "ready");
699 assert!(reclaimed.lease_owner.is_none());
700 assert_eq!(
701 reclaimed.attempts, 1,
702 "the attempt was spent, and the row must keep saying so"
703 );
704
705 let held = Job::find_by_id("job-live", &database)
706 .await
707 .unwrap()
708 .unwrap();
709 assert_eq!(held.status, "running");
710 }
711
712 #[tokio::test]
713 async fn release_owned_is_scoped_to_one_runner() {
714 let database = db().await;
715 assert!(enqueue("job-1", "a", now_secs(), &database).await);
716 assert!(enqueue("job-2", "b", now_secs(), &database).await);
717 claim("runner-a", &database).await.unwrap();
718 claim("runner-b", &database).await.unwrap();
719
720 assert_eq!(Job::release_owned("runner-a", &database).await.unwrap(), 1);
721 assert_eq!(
722 Job::find_by_id("job-1", &database)
723 .await
724 .unwrap()
725 .unwrap()
726 .status,
727 "ready"
728 );
729 assert_eq!(
730 Job::find_by_id("job-2", &database)
731 .await
732 .unwrap()
733 .unwrap()
734 .status,
735 "running"
736 );
737 }
738
739 #[tokio::test]
740 async fn find_live_sees_a_queued_job_and_not_a_settled_one() {
741 let database = db().await;
742 assert!(enqueue("job-1", "ord-1", now_secs(), &database).await);
743 assert!(
744 Job::find_live("test", "ord-1", &database)
745 .await
746 .unwrap()
747 .is_some()
748 );
749
750 let job = claim("runner-a", &database).await.unwrap();
751 Job::complete(&job.id, "runner-a", &database).await.unwrap();
752 assert!(
753 Job::find_live("test", "ord-1", &database)
754 .await
755 .unwrap()
756 .is_none()
757 );
758 }
759
760 #[tokio::test]
761 async fn cleanup_removes_settled_rows_strictly_older_than_the_cutoff() {
762 let database = db().await;
763 assert!(enqueue("job-done", "a", now_secs(), &database).await);
764 assert!(enqueue("job-live", "b", now_secs() + 3_600, &database).await);
765 let job = claim("runner-a", &database).await.unwrap();
766 Job::complete(&job.id, "runner-a", &database).await.unwrap();
767
768 assert_eq!(Job::cleanup(now_secs(), &database).await.unwrap(), 0);
770 assert_eq!(Job::cleanup(now_secs() + 1, &database).await.unwrap(), 1);
771 assert!(
772 Job::find_by_id("job-done", &database)
773 .await
774 .unwrap()
775 .is_none()
776 );
777 assert!(
778 Job::find_by_id("job-live", &database)
779 .await
780 .unwrap()
781 .is_some(),
782 "a job still queued is not old, however long ago it was written"
783 );
784 }
785
786 #[tokio::test]
787 async fn an_unknown_status_is_refused_by_the_check_constraint() {
788 let database = db().await;
789 let error = sqlx::query(
790 "INSERT INTO jobs (id, kind, dedup_key, payload, status, run_at, attempts, \
791 max_attempts, created_at, updated_at) \
792 VALUES ('x', 'test', 'k', '{}', 'halfway', 0, 0, 1, 0, 0);",
793 )
794 .execute(&database.pool)
795 .await
796 .unwrap_err();
797 assert!(error.to_string().contains("CHECK constraint failed"));
798 }
799}