From 5a3a4b92b966cadf33efe21b3c8d62c1ba2b6c47 Mon Sep 17 00:00:00 2001 From: Trezy Date: Thu, 14 May 2026 08:25:25 -0500 Subject: [PATCH] feat: allow backfills to be cancelled Signed-off-by: Trezy --- src/admin/backfill.rs | 147 ++++++++++++++++++++++-- src/admin/mod.rs | 1 + web/src/app/dashboard/backfill/page.tsx | 68 ++++++++++- web/src/lib/api.ts | 7 ++ 4 files changed, 210 insertions(+), 13 deletions(-) diff --git a/src/admin/backfill.rs b/src/admin/backfill.rs index 026da56..633096a 100644 --- a/src/admin/backfill.rs +++ b/src/admin/backfill.rs @@ -1,9 +1,9 @@ use std::collections::HashMap; use std::sync::Arc; -use std::sync::atomic::{AtomicI32, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicI32, Ordering}; use axum::Json; -use axum::extract::State; +use axum::extract::{Path, State}; use axum::http::StatusCode; use futures_util::stream::{self, StreamExt}; use serde::Deserialize; @@ -137,6 +137,42 @@ async fn fail_job(state: &AppState, job_id: &str, error: &str) { cleanup_repos(state, job_id).await; } +async fn is_cancelled(state: &AppState, job_id: &str) -> bool { + let sql = adapt_sql( + "SELECT status FROM backfill_jobs WHERE id = ?", + state.db_backend, + ); + sqlx::query_as::<_, (String,)>(&sql) + .bind(job_id) + .fetch_optional(&state.db) + .await + .ok() + .flatten() + .is_some_and(|(status,)| status == "cancelling") +} + +async fn request_cancel(state: &AppState, job_id: &str) { + let sql = adapt_sql( + "UPDATE backfill_jobs SET status = 'cancelling' WHERE id = ? AND status = 'running'", + state.db_backend, + ); + let _ = sqlx::query(&sql).bind(job_id).execute(&state.db).await; +} + +async fn finalise_cancel(state: &AppState, job_id: &str) { + let now = now_rfc3339(); + let sql = adapt_sql( + "UPDATE backfill_jobs SET status = 'cancelled', stage = 'cancelled', completed_at = ?, error = 'cancelled by user' WHERE id = ?", + state.db_backend, + ); + let _ = sqlx::query(&sql) + .bind(&now) + .bind(job_id) + .execute(&state.db) + .await; + cleanup_repos(state, job_id).await; +} + async fn complete_job( state: &AppState, job_id: &str, @@ -184,6 +220,9 @@ async fn run_discovery_phase( .await; } else { for collection in collections { + if is_cancelled(state, job_id).await { + return; + } if let Err(e) = discover_repos_from_relay(state, job_id, collection).await { tracing::warn!(collection, error = %e, "failed to discover repos, skipping"); } @@ -254,6 +293,10 @@ async fn discover_repos_from_relay( let total = count_repos(state, job_id).await; update_job_counter(state, job_id, "total_repos", total).await; + if is_cancelled(state, job_id).await { + return Ok(()); + } + match body.cursor { Some(c) if page_count > 0 => cursor = Some(c), _ => break, @@ -314,6 +357,9 @@ async fn run_resolution_phase(state: &AppState, job_id: &str) { resolved_count += 1; if resolved_count % 100 == 0 { update_job_counter(state, job_id, "processed_repos", resolved_count).await; + if is_cancelled(state, job_id).await { + return; + } } } @@ -357,6 +403,7 @@ async fn run_fetching_phase(state: &AppState, job_id: &str, collections: &[Strin let processed_repos = Arc::new(AtomicI32::new(already_completed)); let total_records = Arc::new(AtomicI32::new(0)); + let cancelled = Arc::new(AtomicBool::new(false)); let state = Arc::new(state.clone()); let collections = Arc::new(collections.to_vec()); let job_id_arc = Arc::new(job_id.to_string()); @@ -369,6 +416,7 @@ async fn run_fetching_phase(state: &AppState, job_id: &str, collections: &[Strin let collections = Arc::clone(&collections); let processed_repos = Arc::clone(&processed_repos); let total_records = Arc::clone(&total_records); + let cancelled = Arc::clone(&cancelled); let job_id = Arc::clone(&job_id_arc); async move { @@ -378,10 +426,15 @@ async fn run_fetching_phase(state: &AppState, job_id: &str, collections: &[Strin let collections = Arc::clone(&collections); let processed_repos = Arc::clone(&processed_repos); let total_records = Arc::clone(&total_records); + let cancelled = Arc::clone(&cancelled); let pds_endpoint = pds_endpoint.clone(); let job_id = Arc::clone(&job_id); async move { + if cancelled.load(Ordering::Relaxed) { + return; + } + for collection in collections.iter() { match fetch_records_from_pds( &state, @@ -433,6 +486,10 @@ async fn run_fetching_phase(state: &AppState, job_id: &str, collections: &[Strin .bind(job_id.as_str()) .execute(&state.db) .await; + + if is_cancelled(&state, job_id.as_str()).await { + cancelled.store(true, Ordering::Relaxed); + } } } }) @@ -581,6 +638,12 @@ async fn run_backfill_job(state: AppState, job_id: String) { if matches!(stage.as_str(), "pending" | "discovering_repos") { run_discovery_phase(&state, &job_id, &collections, did.as_deref()).await; + if is_cancelled(&state, &job_id).await { + tracing::info!(job_id, "backfill job cancelled"); + finalise_cancel(&state, &job_id).await; + return; + } + let total = count_repos(&state, &job_id).await; if total == 0 { complete_job(&state, &job_id, 0, 0, None).await; @@ -609,10 +672,21 @@ async fn run_backfill_job(state: AppState, job_id: String) { "pending" | "discovering_repos" | "resolving_pds" ) { run_resolution_phase(&state, &job_id).await; + + if is_cancelled(&state, &job_id).await { + tracing::info!(job_id, "backfill job cancelled"); + finalise_cancel(&state, &job_id).await; + return; + } } run_fetching_phase(&state, &job_id, &collections).await; + if is_cancelled(&state, &job_id).await { + tracing::info!(job_id, "backfill job cancelled"); + return; + } + // Read final counters from backfill_repos before cleanup let sql = adapt_sql( "SELECT COUNT(*) FROM backfill_repos WHERE job_id = ? AND status = 'completed'", @@ -717,6 +791,50 @@ pub(super) async fn create_backfill( )) } +/// POST /admin/backfill/{id}/cancel — cancel a running backfill job. +pub(super) async fn cancel_backfill( + State(state): State, + admin: UserAuth, + Path(job_id): Path, +) -> Result, AppError> { + admin.require(Permission::BackfillCreate).await?; + + let sql = adapt_sql( + "SELECT status FROM backfill_jobs WHERE id = ?", + state.db_backend, + ); + let row: Option<(String,)> = sqlx::query_as(&sql) + .bind(&job_id) + .fetch_optional(&state.db) + .await + .map_err(|e| AppError::Internal(format!("failed to query backfill job: {e}")))?; + + match row { + None => Err(AppError::NotFound("backfill job not found".into())), + Some((status,)) if status != "running" => Err(AppError::BadRequest(format!( + "job is not running (status: {status})" + ))), + Some(_) => { + request_cancel(&state, &job_id).await; + log_event( + &state.db, + EventLog { + event_type: "backfill.cancelling".to_string(), + severity: Severity::Info, + actor_did: Some(admin.did.clone()), + subject: None, + detail: serde_json::json!({ "job_id": job_id }), + }, + state.db_backend, + ) + .await; + Ok(Json( + serde_json::json!({ "id": job_id, "status": "cancelling" }), + )) + } + } +} + /// GET /admin/backfill/status — list all backfill jobs. pub(super) async fn backfill_status( State(state): State, @@ -791,21 +909,30 @@ pub(super) async fn backfill_status( // --------------------------------------------------------------------------- /// Resume any backfill jobs that were running when the server last stopped. +/// Jobs stuck in `cancelling` are finalised immediately. pub async fn resume_backfill_jobs(state: &AppState) { let sql = adapt_sql( - "SELECT id FROM backfill_jobs WHERE status = 'running'", + "SELECT id, status FROM backfill_jobs WHERE status IN ('running', 'cancelling')", state.db_backend, ); - let rows: Vec<(String,)> = sqlx::query_as(&sql) + let rows: Vec<(String, String)> = sqlx::query_as(&sql) .fetch_all(&state.db) .await .unwrap_or_default(); - for (job_id,) in rows { - tracing::info!(job_id, "resuming interrupted backfill job"); - let spawn_state = state.clone(); - tokio::spawn(async move { - run_backfill_job(spawn_state, job_id).await; - }); + for (job_id, status) in rows { + if status == "cancelling" { + tracing::info!( + job_id, + "finalising cancelled backfill job from previous run" + ); + finalise_cancel(state, &job_id).await; + } else { + tracing::info!(job_id, "resuming interrupted backfill job"); + let spawn_state = state.clone(); + tokio::spawn(async move { + run_backfill_job(spawn_state, job_id).await; + }); + } } } diff --git a/src/admin/mod.rs b/src/admin/mod.rs index fb308f4..6ad2260 100644 --- a/src/admin/mod.rs +++ b/src/admin/mod.rs @@ -37,6 +37,7 @@ pub fn admin_routes(_state: AppState) -> Router { .route("/stats", get(stats::stats)) .route("/backfill", post(backfill::create_backfill)) .route("/backfill/status", get(backfill::backfill_status)) + .route("/backfill/{id}/cancel", post(backfill::cancel_backfill)) .route("/events", get(events::list_events)) .route("/users", post(users::create_user).get(users::list_users)) .route("/users/transfer-super", post(users::transfer_super)) diff --git a/web/src/app/dashboard/backfill/page.tsx b/web/src/app/dashboard/backfill/page.tsx index 4f7d91d..df1d918 100644 --- a/web/src/app/dashboard/backfill/page.tsx +++ b/web/src/app/dashboard/backfill/page.tsx @@ -3,7 +3,12 @@ import { useCallback, useEffect, useState } from "react"; import { useCurrentUser } from "@/hooks/use-current-user"; -import { createBackfillJob, getBackfillJobs, getLexicons } from "@/lib/api"; +import { + cancelBackfillJob, + createBackfillJob, + getBackfillJobs, + getLexicons, +} from "@/lib/api"; import type { BackfillJob } from "@/types/backfill"; import { CheckCircle2, Circle, Loader2 } from "lucide-react"; import { SiteHeader } from "@/components/site-header"; @@ -32,6 +37,7 @@ import { Label } from "@/components/ui/label"; import { Sheet, SheetContent, + SheetFooter, SheetHeader, SheetTitle, } from "@/components/ui/sheet"; @@ -51,6 +57,8 @@ const STAGES = [ "fetching_records", "completed", "failed", + "cancelled", + "cancelling", ] as const; const STAGE_LABELS: Record = { @@ -60,6 +68,8 @@ const STAGE_LABELS: Record = { fetching_records: "Fetching records", completed: "Completed", failed: "Failed", + cancelled: "Cancelled", + cancelling: "Cancelling", }; function stageBadge(stage: string) { @@ -72,6 +82,18 @@ function stageBadge(stage: string) { ); case "failed": return failed; + case "cancelled": + return ( + + cancelled + + ); + case "cancelling": + return ( + + cancelling + + ); case "pending": return pending; default: @@ -180,7 +202,16 @@ export default function BackfillPage() { }} > - {selectedJob && } + {selectedJob && ( + { + await cancelBackfillJob(selectedJob.id); + load(); + }} + /> + )} @@ -188,8 +219,27 @@ export default function BackfillPage() { ); } -function JobDetail({ job }: { job: BackfillJob }) { +function JobDetail({ + job, + canCancel, + onCancel, +}: { + job: BackfillJob; + canCancel: boolean; + onCancel: () => Promise; +}) { + const [cancelling, setCancelling] = useState(false); const current = stageIndex(job.stage); + const isActive = job.status === "running" || job.status === "cancelling"; + + async function handleCancel() { + setCancelling(true); + try { + await onCancel(); + } finally { + setCancelling(false); + } + } return ( <> @@ -284,6 +334,18 @@ function JobDetail({ job }: { job: BackfillJob }) { + {canCancel && isActive && ( + + + + )} ); } diff --git a/web/src/lib/api.ts b/web/src/lib/api.ts index f05a348..d923355 100644 --- a/web/src/lib/api.ts +++ b/web/src/lib/api.ts @@ -170,6 +170,13 @@ export function createBackfillJob(body: { collection?: string; did?: string }) { }); } +export function cancelBackfillJob(id: string) { + return apiFetch<{ id: string; status: string }>( + `/admin/backfill/${id}/cancel`, + { method: "POST" }, + ); +} + // Users export function getUsers() { return apiFetch("/admin/users"); -- 2.51.2