use chrono::{DateTime, Utc};
use sqlx::{PgConnection, PgPool, Postgres, Transaction};
use crate::error::{map_sqlx_error, AwaError};
use crate::job::{JobRow, JobState};
use crate::queue_storage::QueueStorage;
#[derive(Debug, Clone, Copy)]
pub enum Reschedule {
Snooze { delay_secs: f64 },
RetryAfter { delay_secs: f64 },
RetryBackoff,
}
impl Reschedule {
fn next_state(&self) -> JobState {
match self {
Reschedule::Snooze { .. } => JobState::Scheduled,
Reschedule::RetryAfter { .. } | Reschedule::RetryBackoff => JobState::Retryable,
}
}
fn decrements_attempt(&self) -> bool {
matches!(self, Reschedule::Snooze { .. })
}
fn delay_parameters(&self) -> (f64, bool) {
match self {
Reschedule::Snooze { delay_secs } | Reschedule::RetryAfter { delay_secs } => {
(*delay_secs, false)
}
Reschedule::RetryBackoff => (0.0, true),
}
}
}
#[derive(Debug, Clone)]
pub enum RescheduleOutcome {
Rescheduled {
job_id: i64,
run_at: DateTime<Utc>,
attempt: i16,
},
Migrated {
job_id: i64,
run_at: DateTime<Utc>,
attempt: i16,
},
CancelledDuplicate { job_id: i64 },
Stale,
}
#[derive(sqlx::FromRow)]
struct TakenJob {
kind: String,
queue: String,
args: serde_json::Value,
priority: i16,
attempt: i16,
run_lease: i64,
max_attempts: i16,
attempted_at: Option<DateTime<Utc>>,
created_at: DateTime<Utc>,
errors: Option<Vec<serde_json::Value>>,
metadata: serde_json::Value,
tags: Vec<String>,
unique_key: Option<Vec<u8>>,
unique_states: Option<String>,
}
pub async fn reschedule_canonical_attempt(
pool: &PgPool,
job_id: i64,
run_lease: i64,
reschedule: Reschedule,
error: Option<&serde_json::Value>,
progress: Option<&serde_json::Value>,
) -> Result<RescheduleOutcome, AwaError> {
let mut tx = pool.begin().await.map_err(map_sqlx_error)?;
let outcome =
reschedule_canonical_attempt_tx(&mut tx, job_id, run_lease, reschedule, error, progress)
.await?;
tx.commit().await.map_err(map_sqlx_error)?;
Ok(outcome)
}
pub async fn reschedule_canonical_attempt_tx(
tx: &mut Transaction<'_, Postgres>,
job_id: i64,
run_lease: i64,
reschedule: Reschedule,
error: Option<&serde_json::Value>,
progress: Option<&serde_json::Value>,
) -> Result<RescheduleOutcome, AwaError> {
let active_schema: Option<String> = sqlx::query_scalar(
"SELECT awa.active_queue_storage_schema() \
FROM awa.storage_transition_state WHERE singleton FOR SHARE",
)
.fetch_optional(tx.as_mut())
.await
.map_err(map_sqlx_error)?
.flatten();
match active_schema {
None => {
reschedule_canonical(tx.as_mut(), job_id, run_lease, reschedule, error, progress).await
}
Some(schema) => {
migrate_to_queue_storage(tx, &schema, job_id, run_lease, reschedule, error, progress)
.await
}
}
}
async fn reschedule_canonical(
conn: &mut PgConnection,
job_id: i64,
run_lease: i64,
reschedule: Reschedule,
error: Option<&serde_json::Value>,
progress: Option<&serde_json::Value>,
) -> Result<RescheduleOutcome, AwaError> {
let (delay_secs, use_backoff) = reschedule.delay_parameters();
let row: Option<(i64, DateTime<Utc>, i16)> = sqlx::query_as(
r#"
WITH deleted AS (
DELETE FROM awa.jobs_hot
WHERE id = $1 AND state = 'running' AND run_lease = $2
RETURNING *
), moved AS (
INSERT INTO awa.scheduled_jobs (
id, kind, queue, args, state, priority, attempt, max_attempts,
run_at, heartbeat_at, deadline_at, attempted_at, finalized_at,
created_at, errors, metadata, tags, unique_key, unique_states,
callback_id, callback_timeout_at, callback_filter,
callback_on_complete, callback_on_fail, callback_transform,
run_lease, progress
)
SELECT
id, kind, queue, args,
$3::awa.job_state,
priority,
CASE WHEN $4 THEN attempt - 1 ELSE attempt END,
max_attempts,
CASE WHEN $5 THEN now() + awa.backoff_duration(attempt, max_attempts)
ELSE now() + make_interval(secs => $6) END,
NULL, NULL, attempted_at,
CASE WHEN $3::awa.job_state = 'retryable' THEN now() ELSE finalized_at END,
created_at,
CASE WHEN $7::jsonb IS NULL THEN errors ELSE errors || $7::jsonb END,
metadata, tags, unique_key, unique_states,
callback_id, callback_timeout_at, callback_filter,
callback_on_complete, callback_on_fail, callback_transform,
run_lease, $8
FROM deleted
RETURNING id, run_at, attempt
)
SELECT id, run_at, attempt FROM moved
"#,
)
.bind(job_id)
.bind(run_lease)
.bind(reschedule.next_state())
.bind(reschedule.decrements_attempt())
.bind(use_backoff)
.bind(delay_secs)
.bind(error)
.bind(progress)
.fetch_optional(&mut *conn)
.await
.map_err(map_sqlx_error)?;
Ok(match row {
Some((job_id, run_at, attempt)) => RescheduleOutcome::Rescheduled {
job_id,
run_at,
attempt,
},
None => RescheduleOutcome::Stale,
})
}
async fn migrate_to_queue_storage(
tx: &mut Transaction<'_, Postgres>,
schema: &str,
job_id: i64,
run_lease: i64,
reschedule: Reschedule,
error: Option<&serde_json::Value>,
progress: Option<&serde_json::Value>,
) -> Result<RescheduleOutcome, AwaError> {
let taken: Option<TakenJob> = sqlx::query_as(
r#"
DELETE FROM awa.jobs_hot
WHERE id = $1 AND state = 'running' AND run_lease = $2
AND callback_id IS NULL
AND callback_timeout_at IS NULL
AND callback_filter IS NULL
AND callback_on_complete IS NULL
AND callback_on_fail IS NULL
AND callback_transform IS NULL
RETURNING kind, queue, args, priority, attempt, run_lease, max_attempts,
attempted_at, created_at, errors, metadata, tags,
unique_key, unique_states::text AS unique_states
"#,
)
.bind(job_id)
.bind(run_lease)
.fetch_optional(tx.as_mut())
.await
.map_err(map_sqlx_error)?;
let Some(job) = taken else {
return reschedule_canonical(tx.as_mut(), job_id, run_lease, reschedule, error, progress)
.await;
};
let (delay_secs, use_backoff) = reschedule.delay_parameters();
let next_state = reschedule.next_state();
let attempt = if reschedule.decrements_attempt() {
job.attempt.saturating_sub(1)
} else {
job.attempt
};
let (new_id, run_at, db_now): (i64, DateTime<Utc>, DateTime<Utc>) = sqlx::query_as(&format!(
r#"
SELECT nextval('{schema}.job_id_seq')::bigint,
CASE WHEN $1 THEN now() + awa.backoff_duration($2::smallint, $3::smallint)
ELSE now() + make_interval(secs => $4) END,
now()
"#
))
.bind(use_backoff)
.bind(job.attempt)
.bind(job.max_attempts)
.bind(delay_secs)
.fetch_one(tx.as_mut())
.await
.map_err(map_sqlx_error)?;
let mut errors: Vec<serde_json::Value> = job.errors.clone().unwrap_or_default();
if let Some(error) = error {
errors.push(error.clone());
}
let successor = JobRow {
id: new_id,
kind: job.kind,
queue: job.queue,
args: job.args,
state: next_state,
priority: job.priority,
attempt,
run_lease: job.run_lease,
max_attempts: job.max_attempts,
run_at,
heartbeat_at: None,
deadline_at: None,
attempted_at: job.attempted_at,
finalized_at: (next_state == JobState::Retryable).then_some(db_now),
created_at: job.created_at,
errors: (!errors.is_empty()).then_some(errors),
metadata: job.metadata,
tags: job.tags,
unique_key: job.unique_key,
unique_states: None,
callback_id: None,
callback_timeout_at: None,
callback_filter: None,
callback_on_complete: None,
callback_on_fail: None,
callback_transform: None,
progress: progress.cloned(),
};
let inserted = QueueStorage::from_existing_schema(schema)?
.insert_migrated_deferred_or_cancel_duplicate_tx(
tx,
successor,
job.unique_states,
"rescheduled as duplicate: unique claim held by a newer job",
)
.await?;
if inserted.state == JobState::Cancelled {
Ok(RescheduleOutcome::CancelledDuplicate {
job_id: inserted.id,
})
} else {
Ok(RescheduleOutcome::Migrated {
job_id: inserted.id,
run_at: inserted.run_at,
attempt: inserted.attempt,
})
}
}