use crate::cli::parser::QueueAction;
pub async fn run(action: &QueueAction, dsn: Option<&str>) -> Result<(), QueueError> {
let url = resolve_dsn(dsn)?;
let db = crate::database::Db::connect(
crate::database::DatabaseConfig::new(&url).map_err(QueueError::Config)?,
)
.await
.map_err(QueueError::Connect)?;
let pool = db.sqlx().clone();
crate::jobs::Jobs::new(pool.clone())
.migrate()
.await
.map_err(QueueError::Migrate)?;
match action {
QueueAction::Stats => {
print_stats(&pool).await?;
}
QueueAction::Drain => {
let requeued = requeue_all_dead(&pool).await?;
let swept = crate::jobs::admin::sweep_expired_leases(&pool, 1024)
.await
.map_err(QueueError::Admin)?;
println!("drained: {requeued} dead requeued, {swept} expired leases swept");
}
QueueAction::Work => {
eprintln!(
"note: `arc queue work` runs a no-handler worker; use the app's in-process worker for real dispatch"
);
let swept = crate::jobs::admin::sweep_expired_leases(&pool, 1024)
.await
.map_err(QueueError::Admin)?;
println!("swept {swept} expired leases");
}
}
db.close().await;
Ok(())
}
fn resolve_dsn(dsn: Option<&str>) -> Result<String, QueueError> {
if let Some(url) = dsn {
return Ok(url.to_string());
}
std::env::var("DATABASE_URL").map_err(|_| QueueError::NoDsn)
}
const COUNT_SQL: &str = r#"SELECT
COUNT(CASE WHEN status = 'pending' THEN 1 END) AS pending,
COUNT(CASE WHEN status = 'running' THEN 1 END) AS running,
COUNT(CASE WHEN status = 'dead' THEN 1 END) AS dead,
COUNT(CASE WHEN status = 'cancelled' THEN 1 END) AS cancelled
FROM arcature_jobs"#;
async fn print_stats(pool: &crate::database::Pool) -> Result<(), QueueError> {
let row: (i64, i64, i64, i64) = sqlx::query_as(COUNT_SQL)
.fetch_one(pool)
.await
.map_err(QueueError::Sqlx)?;
println!("pending: {}", row.0);
println!("running: {}", row.1);
println!("dead: {}", row.2);
println!("cancelled: {}", row.3);
Ok(())
}
const DEAD_BATCH: usize = 1024;
const DEAD_IDS_SQL: &str =
"SELECT id FROM arcature_jobs WHERE status = 'dead' ORDER BY id LIMIT 1024";
async fn requeue_all_dead(pool: &crate::database::Pool) -> Result<u64, QueueError> {
let mut total = 0u64;
loop {
let ids: Vec<uuid::Uuid> = sqlx::query_scalar(DEAD_IDS_SQL)
.fetch_all(pool)
.await
.map_err(QueueError::Sqlx)?;
if ids.is_empty() {
return Ok(total);
}
let mut requeued = 0u64;
for id in &ids {
requeued += crate::jobs::admin::requeue_dead(pool, *id)
.await
.map_err(QueueError::Admin)?;
}
if requeued == 0 {
return Ok(total);
}
total += requeued;
if ids.len() < DEAD_BATCH {
return Ok(total);
}
}
}
#[derive(Debug)]
pub enum QueueError {
NoDsn,
Config(crate::Error),
Connect(crate::Error),
Migrate(crate::jobs::MigrateError),
Admin(crate::jobs::WorkerError),
Sqlx(sqlx::Error),
}
impl std::fmt::Display for QueueError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::NoDsn => f.write_str(
"no --dsn given and DATABASE_URL is unset; \
pass --dsn <url> or set DATABASE_URL",
),
Self::Config(e) => write!(f, "invalid database config: {e}"),
Self::Connect(e) => write!(f, "failed to connect: {e}"),
Self::Migrate(e) => write!(f, "queue migration failed: {e}"),
Self::Admin(e) => write!(f, "queue operation failed: {e}"),
Self::Sqlx(e) => write!(f, "query failed: {e}"),
}
}
}
impl std::error::Error for QueueError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match self {
Self::Config(e) | Self::Connect(e) => Some(e),
Self::Migrate(e) => Some(e),
Self::Admin(e) => Some(e),
Self::Sqlx(e) => Some(e),
_ => None,
}
}
}