Something went wrong. Try again.
A lexicon-driven AppView for ATProto.
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330use serde_json::Value;use uuid::Uuid;
use crate::AppState;use crate::db::{DatabaseBackend, adapt_sql, now_rfc3339};use crate::error::AppError;
use super::Job;
const JOB_COLUMNS: &str = "id, job_type, status, input, progress, result, error, created_by, started_at, completed_at, created_at, CAST(inherit_auth AS INTEGER), api_client_id, dpop_key_id";
type JobRow = ( String, String, String, String, String, Option<String>, Option<String>, String, Option<String>, Option<String>, String, i64, Option<String>, Option<String>,);
fn row_to_job( ( id, job_type, status, input, progress, result, error, created_by, started_at, completed_at, created_at, inherit_auth, api_client_id, dpop_key_id, ): JobRow,) -> Job { Job { id, job_type, status, input: serde_json::from_str(&input).unwrap_or(Value::Null), progress: serde_json::from_str(&progress).unwrap_or(Value::Null), result: result.and_then(|r| serde_json::from_str(&r).ok()), error, created_by, started_at, completed_at, created_at, inherit_auth: inherit_auth != 0, api_client_id, dpop_key_id, }}
pub async fn create_job( state: &AppState, job_type: &str, input: &Value, created_by: &str, inherit_auth: bool, api_client_id: Option<&str>, dpop_key_id: Option<&str>,) -> Result<String, AppError> { let id = Uuid::new_v4().to_string(); let now = now_rfc3339(); let input_str = serde_json::to_string(input) .map_err(|e| AppError::Internal(format!("failed to serialize job input: {e}")))?;
let sql = adapt_sql( "INSERT INTO happyview_jobs (id, job_type, status, input, created_by, created_at, inherit_auth, api_client_id, dpop_key_id) VALUES (?, ?, 'pending', ?, ?, ?, ?, ?, ?)", state.db_backend, ); crate::db::query(&sql) .bind(&id) .bind(job_type) .bind(&input_str) .bind(created_by) .bind(&now) .bind(inherit_auth) .bind(api_client_id) .bind(dpop_key_id) .execute(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to create job: {e}")))?;
Ok(id)}
pub async fn get_job(state: &AppState, id: &str) -> Result<Option<Job>, AppError> { let sql = adapt_sql( &format!("SELECT {JOB_COLUMNS} FROM happyview_jobs WHERE id = ?"), state.db_backend, ); let row: Option<JobRow> = crate::db::query_as(&sql) .bind(id) .fetch_optional(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to fetch job: {e}")))?;
Ok(row.map(row_to_job))}
pub async fn list_jobs( state: &AppState, status_filter: Option<&str>, limit: i64, cursor: Option<&str>,) -> Result<(Vec<Job>, Option<String>), AppError> { let sql = if status_filter.is_some() { let base = if cursor.is_some() { format!( "SELECT {JOB_COLUMNS} FROM happyview_jobs WHERE status = ? AND created_at < ? ORDER BY created_at DESC LIMIT ?" ) } else { format!( "SELECT {JOB_COLUMNS} FROM happyview_jobs WHERE status = ? ORDER BY created_at DESC LIMIT ?" ) }; adapt_sql(&base, state.db_backend) } else { let base = if cursor.is_some() { format!( "SELECT {JOB_COLUMNS} FROM happyview_jobs WHERE created_at < ? ORDER BY created_at DESC LIMIT ?" ) } else { format!("SELECT {JOB_COLUMNS} FROM happyview_jobs ORDER BY created_at DESC LIMIT ?") }; adapt_sql(&base, state.db_backend) };
let mut query = crate::db::query_as::<JobRow>(&sql);
if let Some(status) = status_filter { query = query.bind(status); } if let Some(cursor) = cursor { query = query.bind(cursor); } query = query.bind(limit + 1);
let rows = query .fetch_all(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to list jobs: {e}")))?;
let has_more = rows.len() as i64 > limit; let jobs: Vec<Job> = rows .into_iter() .take(limit as usize) .map(row_to_job) .collect();
let next_cursor = if has_more { jobs.last().map(|j| j.created_at.clone()) } else { None };
Ok((jobs, next_cursor))}
pub async fn set_status(state: &AppState, id: &str, status: &str) -> Result<(), AppError> { let now = now_rfc3339(); let sql = match status { "running" => adapt_sql( "UPDATE happyview_jobs SET status = ?, started_at = ? WHERE id = ?", state.db_backend, ), "completed" | "failed" | "cancelled" => adapt_sql( "UPDATE happyview_jobs SET status = ?, completed_at = ? WHERE id = ?", state.db_backend, ), _ => adapt_sql( "UPDATE happyview_jobs SET status = ? WHERE id = ? AND 1=1", state.db_backend, ), };
match status { "running" | "completed" | "failed" | "cancelled" => { crate::db::query(&sql) .bind(status) .bind(&now) .bind(id) .execute(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to update job status: {e}")))?; } _ => { let sql = adapt_sql( "UPDATE happyview_jobs SET status = ? WHERE id = ?", state.db_backend, ); crate::db::query(&sql) .bind(status) .bind(id) .execute(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to update job status: {e}")))?; } }
Ok(())}
pub async fn update_progress(state: &AppState, id: &str, progress: &Value) -> Result<(), AppError> { let progress_str = serde_json::to_string(progress) .map_err(|e| AppError::Internal(format!("failed to serialize progress: {e}")))?; let sql = adapt_sql( "UPDATE happyview_jobs SET progress = ? WHERE id = ?", state.db_backend, ); crate::db::query(&sql) .bind(&progress_str) .bind(id) .execute(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to update job progress: {e}")))?; Ok(())}
pub async fn set_result(state: &AppState, id: &str, result: &Value) -> Result<(), AppError> { let result_str = serde_json::to_string(result) .map_err(|e| AppError::Internal(format!("failed to serialize result: {e}")))?; let now = now_rfc3339(); let sql = adapt_sql( "UPDATE happyview_jobs SET status = 'completed', result = ?, completed_at = ? WHERE id = ?", state.db_backend, ); crate::db::query(&sql) .bind(&result_str) .bind(&now) .bind(id) .execute(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to set job result: {e}")))?; Ok(())}
pub async fn set_error(state: &AppState, id: &str, error: &str) -> Result<(), AppError> { let now = now_rfc3339(); let sql = adapt_sql( "UPDATE happyview_jobs SET status = 'failed', error = ?, completed_at = ? WHERE id = ?", state.db_backend, ); crate::db::query(&sql) .bind(error) .bind(&now) .bind(id) .execute(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to set job error: {e}")))?; Ok(())}
/// Check if a job should stop (status changed to cancelling or pausing)./// Same cooperative cancellation pattern as the backfill system.pub async fn should_stop(state: &AppState, id: &str) -> Option<&'static str> { let sql = adapt_sql( "SELECT status FROM happyview_jobs WHERE id = ?", state.db_backend, ); let status = crate::db::query_as::<(String,)>(&sql) .bind(id) .fetch_optional(&state.db) .await .ok() .flatten() .map(|(s,)| s); match status.as_deref() { Some("cancelling") => Some("cancelling"), Some("pausing") => Some("pausing"), _ => None, }}
/// Find jobs that were interrupted by a server restart.pub async fn find_interrupted_jobs(state: &AppState) -> Vec<Job> { let sql = adapt_sql( &format!( "SELECT {JOB_COLUMNS} FROM happyview_jobs WHERE status IN ('running', 'cancelling', 'pausing')" ), state.db_backend, ); let rows: Vec<JobRow> = crate::db::query_as(&sql) .fetch_all(&state.db) .await .unwrap_or_default();
rows.into_iter().map(row_to_job).collect()}
/// Pick the next pending job and atomically set it to running.pub async fn claim_next_job(state: &AppState) -> Result<Option<Job>, AppError> { let now = now_rfc3339();
let sql = match state.db_backend { DatabaseBackend::Postgres => adapt_sql( &format!( "UPDATE happyview_jobs SET status = 'running', started_at = ? WHERE id = (SELECT id FROM happyview_jobs WHERE status = 'pending' ORDER BY created_at ASC LIMIT 1 FOR UPDATE SKIP LOCKED) RETURNING {JOB_COLUMNS}" ), state.db_backend, ), DatabaseBackend::Sqlite => adapt_sql( &format!( "UPDATE happyview_jobs SET status = 'running', started_at = ? WHERE id = (SELECT id FROM happyview_jobs WHERE status = 'pending' ORDER BY created_at ASC LIMIT 1) AND status = 'pending' RETURNING {JOB_COLUMNS}" ), state.db_backend, ), };
let row: Option<JobRow> = crate::db::query_as(&sql) .bind(&now) .fetch_optional(&state.db) .await .map_err(|e| AppError::Internal(format!("failed to claim job: {e}")))?;
Ok(row.map(row_to_job))}