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