diff --git a/migrations/postgres/20260513000000_add_backfill_stage.sql b/migrations/postgres/20260513000000_add_backfill_stage.sql new file mode 100644 index 0000000..dc47b0f --- /dev/null +++ b/migrations/postgres/20260513000000_add_backfill_stage.sql @@ -0,0 +1,4 @@ +ALTER TABLE backfill_jobs ADD COLUMN stage TEXT NOT NULL DEFAULT 'pending'; + +UPDATE backfill_jobs SET stage = status WHERE status IN ('completed', 'failed'); +UPDATE backfill_jobs SET stage = 'failed', status = 'failed', error = 'interrupted by restart' WHERE status = 'running'; diff --git a/migrations/postgres/20260513000001_create_backfill_repos.sql b/migrations/postgres/20260513000001_create_backfill_repos.sql new file mode 100644 index 0000000..9c0d25b --- /dev/null +++ b/migrations/postgres/20260513000001_create_backfill_repos.sql @@ -0,0 +1,7 @@ +CREATE TABLE IF NOT EXISTS backfill_repos ( + job_id UUID NOT NULL REFERENCES backfill_jobs(id) ON DELETE CASCADE, + did TEXT NOT NULL, + pds_endpoint TEXT, + status TEXT NOT NULL DEFAULT 'pending', + PRIMARY KEY (job_id, did) +); diff --git a/migrations/sqlite/20260513000000_add_backfill_stage.sql b/migrations/sqlite/20260513000000_add_backfill_stage.sql new file mode 100644 index 0000000..dc47b0f --- /dev/null +++ b/migrations/sqlite/20260513000000_add_backfill_stage.sql @@ -0,0 +1,4 @@ +ALTER TABLE backfill_jobs ADD COLUMN stage TEXT NOT NULL DEFAULT 'pending'; + +UPDATE backfill_jobs SET stage = status WHERE status IN ('completed', 'failed'); +UPDATE backfill_jobs SET stage = 'failed', status = 'failed', error = 'interrupted by restart' WHERE status = 'running'; diff --git a/migrations/sqlite/20260513000001_create_backfill_repos.sql b/migrations/sqlite/20260513000001_create_backfill_repos.sql new file mode 100644 index 0000000..0e613eb --- /dev/null +++ b/migrations/sqlite/20260513000001_create_backfill_repos.sql @@ -0,0 +1,7 @@ +CREATE TABLE IF NOT EXISTS backfill_repos ( + job_id TEXT NOT NULL REFERENCES backfill_jobs(id) ON DELETE CASCADE, + did TEXT NOT NULL, + pds_endpoint TEXT, + status TEXT NOT NULL DEFAULT 'pending', + PRIMARY KEY (job_id, did) +); diff --git a/src/admin/backfill.rs b/src/admin/backfill.rs index 6b8ce1e..b7c62a5 100644 --- a/src/admin/backfill.rs +++ b/src/admin/backfill.rs @@ -22,7 +22,7 @@ use super::permissions::Permission; use super::types::{BackfillJob, CreateBackfillBody}; // --------------------------------------------------------------------------- -// Relay discovery (reused from old backfill module) +// Response types // --------------------------------------------------------------------------- #[derive(Deserialize)] @@ -36,10 +36,6 @@ struct RepoEntry { did: String, } -// --------------------------------------------------------------------------- -// PDS record types -// --------------------------------------------------------------------------- - #[derive(Deserialize)] struct ListRecordsResponse { records: Vec, @@ -53,15 +49,133 @@ struct RecordEntry { value: serde_json::Value, } -/// Discover all DIDs that have records in `collection` via the relay's -/// `com.atproto.sync.listReposByCollection` endpoint. Paginates until done. -async fn list_repos_by_collection( - http: &reqwest::Client, - relay_url: &str, +// --------------------------------------------------------------------------- +// Helpers +// --------------------------------------------------------------------------- + +async fn set_stage(state: &AppState, job_id: &str, stage: &str) { + let sql = adapt_sql( + "UPDATE backfill_jobs SET stage = ? WHERE id = ?", + state.db_backend, + ); + let _ = sqlx::query(&sql) + .bind(stage) + .bind(job_id) + .execute(&state.db) + .await; +} + +async fn update_job_counter(state: &AppState, job_id: &str, column: &str, value: i32) { + let sql = adapt_sql( + &format!("UPDATE backfill_jobs SET {column} = ? WHERE id = ?"), + state.db_backend, + ); + let _ = sqlx::query(&sql) + .bind(value) + .bind(job_id) + .execute(&state.db) + .await; +} + +async fn count_repos(state: &AppState, job_id: &str) -> i32 { + let sql = adapt_sql( + "SELECT COUNT(*) FROM backfill_repos WHERE job_id = ?", + state.db_backend, + ); + sqlx::query_as::<_, (i32,)>(&sql) + .bind(job_id) + .fetch_one(&state.db) + .await + .map(|(c,)| c) + .unwrap_or(0) +} + +async fn cleanup_repos(state: &AppState, job_id: &str) { + let sql = adapt_sql( + "DELETE FROM backfill_repos WHERE job_id = ?", + state.db_backend, + ); + let _ = sqlx::query(&sql).bind(job_id).execute(&state.db).await; +} + +async fn fail_job(state: &AppState, job_id: &str, error: &str) { + let now = now_rfc3339(); + let sql = adapt_sql( + "UPDATE backfill_jobs SET status = 'failed', stage = 'failed', completed_at = ?, error = ? WHERE id = ?", + state.db_backend, + ); + let _ = sqlx::query(&sql) + .bind(&now) + .bind(error) + .bind(job_id) + .execute(&state.db) + .await; + cleanup_repos(state, job_id).await; +} + +async fn complete_job( + state: &AppState, + job_id: &str, + processed_repos: i32, + total_records: i32, + error: Option<&str>, +) { + let now = now_rfc3339(); + let sql = adapt_sql( + "UPDATE backfill_jobs SET status = 'completed', stage = 'completed', completed_at = ?, processed_repos = ?, total_records = ?, error = ? WHERE id = ?", + state.db_backend, + ); + let _ = sqlx::query(&sql) + .bind(&now) + .bind(processed_repos) + .bind(total_records) + .bind(error) + .bind(job_id) + .execute(&state.db) + .await; + cleanup_repos(state, job_id).await; +} + +// --------------------------------------------------------------------------- +// Phase 1: Discover repos via relay +// --------------------------------------------------------------------------- + +async fn run_discovery_phase( + state: &AppState, + job_id: &str, + collections: &[String], + specific_did: Option<&str>, +) { + set_stage(state, job_id, "discovering_repos").await; + + if let Some(did) = specific_did { + let sql = adapt_sql( + "INSERT INTO backfill_repos (job_id, did) VALUES (?, ?) ON CONFLICT DO NOTHING", + state.db_backend, + ); + let _ = sqlx::query(&sql) + .bind(job_id) + .bind(did) + .execute(&state.db) + .await; + } else { + for collection in collections { + if let Err(e) = discover_repos_from_relay(state, job_id, collection).await { + tracing::warn!(collection, error = %e, "failed to discover repos, skipping"); + } + } + } + + let total = count_repos(state, job_id).await; + update_job_counter(state, job_id, "total_repos", total).await; +} + +async fn discover_repos_from_relay( + state: &AppState, + job_id: &str, collection: &str, -) -> Result, String> { - let base = relay_url.trim_end_matches('/'); - let mut dids = Vec::new(); +) -> Result<(), String> { + let base = state.config.relay_url.trim_end_matches('/'); let mut cursor: Option = None; loop { @@ -72,7 +186,8 @@ async fn list_repos_by_collection( url.push_str(&format!("&cursor={c}")); } - let resp = http + let resp = state + .http .get(&url) .send() .await @@ -88,23 +203,210 @@ async fn list_repos_by_collection( .map_err(|e| format!("invalid relay response: {e}"))?; let page_count = body.repos.len(); - for repo in body.repos { - dids.push(repo.did); + + for repo in &body.repos { + let sql = adapt_sql( + "INSERT INTO backfill_repos (job_id, did) VALUES (?, ?) ON CONFLICT DO NOTHING", + state.db_backend, + ); + let _ = sqlx::query(&sql) + .bind(job_id) + .bind(&repo.did) + .execute(&state.db) + .await; } + let total = count_repos(state, job_id).await; + update_job_counter(state, job_id, "total_repos", total).await; + match body.cursor { Some(c) if page_count > 0 => cursor = Some(c), _ => break, } } - Ok(dids) + Ok(()) +} + +// --------------------------------------------------------------------------- +// Phase 2: Resolve PDS endpoints +// --------------------------------------------------------------------------- + +async fn run_resolution_phase(state: &AppState, job_id: &str) { + set_stage(state, job_id, "resolving_pds").await; + + let sql = adapt_sql( + "SELECT did FROM backfill_repos WHERE job_id = ? AND pds_endpoint IS NULL", + state.db_backend, + ); + let unresolved: Vec<(String,)> = sqlx::query_as(&sql) + .bind(job_id) + .fetch_all(&state.db) + .await + .unwrap_or_default(); + + let sql = adapt_sql( + "SELECT COUNT(*) FROM backfill_repos WHERE job_id = ? AND pds_endpoint IS NOT NULL", + state.db_backend, + ); + let already_resolved: i32 = sqlx::query_as::<_, (i32,)>(&sql) + .bind(job_id) + .fetch_one(&state.db) + .await + .map(|(c,)| c) + .unwrap_or(0); + + let mut resolved_count = already_resolved; + + for (did,) in &unresolved { + match profile::resolve_pds_endpoint(&state.http, &state.config.plc_url, did).await { + Ok(pds) => { + let sql = adapt_sql( + "UPDATE backfill_repos SET pds_endpoint = ? WHERE job_id = ? AND did = ?", + state.db_backend, + ); + let _ = sqlx::query(&sql) + .bind(&pds) + .bind(job_id) + .bind(did) + .execute(&state.db) + .await; + } + Err(e) => { + tracing::warn!(did, error = %e, "failed to resolve PDS endpoint, skipping DID"); + } + } + resolved_count += 1; + if resolved_count % 100 == 0 { + update_job_counter(state, job_id, "processed_repos", resolved_count).await; + } + } + + update_job_counter(state, job_id, "processed_repos", resolved_count).await; } // --------------------------------------------------------------------------- -// PDS record fetching +// Phase 3: Fetch records from PDS instances // --------------------------------------------------------------------------- +async fn run_fetching_phase(state: &AppState, job_id: &str, collections: &[String]) { + set_stage(state, job_id, "fetching_records").await; + + // Load pending repos grouped by PDS + let sql = adapt_sql( + "SELECT did, pds_endpoint FROM backfill_repos WHERE job_id = ? AND status = 'pending' AND pds_endpoint IS NOT NULL", + state.db_backend, + ); + let rows: Vec<(String, String)> = sqlx::query_as(&sql) + .bind(job_id) + .fetch_all(&state.db) + .await + .unwrap_or_default(); + + let mut pds_to_dids: HashMap> = HashMap::new(); + for (did, pds) in rows { + pds_to_dids.entry(pds).or_default().push(did); + } + + // Count already-completed repos for accurate progress + let sql = adapt_sql( + "SELECT COUNT(*) FROM backfill_repos WHERE job_id = ? AND status = 'completed'", + state.db_backend, + ); + let already_completed: i32 = sqlx::query_as::<_, (i32,)>(&sql) + .bind(job_id) + .fetch_one(&state.db) + .await + .map(|(c,)| c) + .unwrap_or(0); + + let processed_repos = Arc::new(AtomicI32::new(already_completed)); + let total_records = Arc::new(AtomicI32::new(0)); + let state = Arc::new(state.clone()); + let collections = Arc::new(collections.to_vec()); + let job_id_arc = Arc::new(job_id.to_string()); + + let pds_entries: Vec<(String, Vec)> = pds_to_dids.into_iter().collect(); + + stream::iter(pds_entries) + .for_each_concurrent(10, |(pds_endpoint, dids)| { + let state = Arc::clone(&state); + let collections = Arc::clone(&collections); + let processed_repos = Arc::clone(&processed_repos); + let total_records = Arc::clone(&total_records); + let job_id = Arc::clone(&job_id_arc); + + async move { + stream::iter(dids) + .for_each_concurrent(3, |did| { + let state = Arc::clone(&state); + let collections = Arc::clone(&collections); + let processed_repos = Arc::clone(&processed_repos); + let total_records = Arc::clone(&total_records); + let pds_endpoint = pds_endpoint.clone(); + let job_id = Arc::clone(&job_id); + + async move { + for collection in collections.iter() { + match fetch_records_from_pds( + &state, + &pds_endpoint, + &did, + collection, + ) + .await + { + Ok(count) => { + total_records + .fetch_add(count as i32, Ordering::Relaxed); + } + Err(e) => { + tracing::warn!( + did, + collection, + pds = %pds_endpoint, + error = %e, + "failed to fetch records from PDS" + ); + } + } + } + + // Mark DID as completed + let sql = adapt_sql( + "UPDATE backfill_repos SET status = 'completed' WHERE job_id = ? AND did = ?", + state.db_backend, + ); + let _ = sqlx::query(&sql) + .bind(job_id.as_str()) + .bind(&did) + .execute(&state.db) + .await; + + let repos = processed_repos.fetch_add(1, Ordering::Relaxed) + 1; + + if repos % 100 == 0 { + let records = total_records.load(Ordering::Relaxed); + let backend = state.db_backend; + let sql = adapt_sql( + "UPDATE backfill_jobs SET processed_repos = ?, total_records = ? WHERE id = ?", + backend, + ); + let _ = sqlx::query(&sql) + .bind(repos) + .bind(records) + .bind(job_id.as_str()) + .execute(&state.db) + .await; + } + } + }) + .await; + } + }) + .await; +} + /// Fetch all records for a given DID and collection from a PDS via /// `com.atproto.repo.listRecords`, paginating and handling rate limits. async fn fetch_records_from_pds( @@ -147,7 +449,7 @@ async fn fetch_records_from_pds( "rate limited by PDS, sleeping" ); tokio::time::sleep(tokio::time::Duration::from_secs(retry_after)).await; - continue; // retry same page + continue; } if !resp.status().is_success() { @@ -187,77 +489,31 @@ async fn fetch_records_from_pds( } // --------------------------------------------------------------------------- -// Admin handlers +// Background backfill worker // --------------------------------------------------------------------------- -/// POST /admin/backfill — create a backfill job and spawn background work. -pub(super) async fn create_backfill( - State(state): State, - admin: UserAuth, - Json(body): Json, -) -> Result<(StatusCode, Json), AppError> { - admin.require(Permission::BackfillCreate).await?; +async fn run_backfill_job(state: AppState, job_id: String) { let backend = state.db_backend; - let now = now_rfc3339(); - let job_id = Uuid::new_v4().to_string(); + // Load job metadata let sql = adapt_sql( - "INSERT INTO backfill_jobs (id, collection, did, status, started_at, created_at) VALUES (?, ?, ?, 'running', ?, ?) RETURNING id", + "SELECT collection, did, stage FROM backfill_jobs WHERE id = ?", backend, ); - let row: (String,) = sqlx::query_as(&sql) + let job: Option<(Option, Option, String)> = sqlx::query_as(&sql) .bind(&job_id) - .bind(&body.collection) - .bind(&body.did) - .bind(&now) - .bind(&now) - .fetch_one(&state.db) + .fetch_optional(&state.db) .await - .map_err(|e| AppError::Internal(format!("failed to create backfill job: {e}")))?; - - let job_id = row.0.clone(); - - log_event( - &state.db, - EventLog { - event_type: "backfill.started".to_string(), - severity: Severity::Info, - actor_did: Some(admin.did.clone()), - subject: body.collection.clone(), - detail: serde_json::json!({ - "job_id": job_id.clone(), - }), - }, - backend, - ) - .await; + .ok() + .flatten(); - // Clone what we need and spawn the background job - let spawn_state = state.clone(); - let spawn_job_id = job_id.clone(); - let spawn_body = body.clone(); - tokio::spawn(async move { - run_backfill_job(spawn_state, spawn_job_id, spawn_body).await; - }); - - Ok(( - StatusCode::CREATED, - Json(serde_json::json!({ - "id": job_id, - "status": "running", - })), - )) -} - -// --------------------------------------------------------------------------- -// Background backfill worker -// --------------------------------------------------------------------------- - -async fn run_backfill_job(state: AppState, job_id: String, body: CreateBackfillBody) { - let backend = state.db_backend; + let Some((collection, did, stage)) = job else { + tracing::error!(job_id, "backfill job not found"); + return; + }; // Determine target collections - let collections: Vec = if let Some(ref col) = body.collection { + let collections: Vec = if let Some(ref col) = collection { let lexicon_exists: bool = state .lexicons .get(col) @@ -297,157 +553,64 @@ async fn run_backfill_job(state: AppState, job_id: String, body: CreateBackfillB return; } - // Discover DIDs - let mut all_dids = Vec::new(); - - for collection in &collections { - let dids = if let Some(ref did) = body.did { - vec![did.clone()] - } else { - match list_repos_by_collection(&state.http, &state.config.relay_url, collection).await { - Ok(dids) => dids, - Err(e) => { - tracing::warn!(collection, error = %e, "failed to discover repos, skipping"); - continue; - } - } - }; - - all_dids.extend(dids); + // Run phases, skipping those already completed + if matches!(stage.as_str(), "pending" | "discovering_repos") { + run_discovery_phase(&state, &job_id, &collections, did.as_deref()).await; + + let total = count_repos(&state, &job_id).await; + if total == 0 { + complete_job(&state, &job_id, 0, 0, None).await; + log_event( + &state.db, + EventLog { + event_type: "backfill.completed".to_string(), + severity: Severity::Info, + actor_did: None, + subject: collection.clone(), + detail: serde_json::json!({ + "job_id": job_id, + "total_repos": 0, + "total_records": 0, + }), + }, + backend, + ) + .await; + return; + } } - all_dids.sort(); - all_dids.dedup(); + if matches!( + stage.as_str(), + "pending" | "discovering_repos" | "resolving_pds" + ) { + run_resolution_phase(&state, &job_id).await; + } - let total_repos = all_dids.len() as i32; + run_fetching_phase(&state, &job_id, &collections).await; - // Update total_repos in DB + // Read final counters from backfill_repos before cleanup let sql = adapt_sql( - "UPDATE backfill_jobs SET total_repos = ? WHERE id = ?", - backend, + "SELECT COUNT(*) FROM backfill_repos WHERE job_id = ? AND status = 'completed'", + state.db_backend, ); - let _ = sqlx::query(&sql) - .bind(total_repos) + let final_processed: i32 = sqlx::query_as::<_, (i32,)>(&sql) .bind(&job_id) - .execute(&state.db) - .await; - - if all_dids.is_empty() { - complete_job(&state, &job_id, 0, 0, None).await; - - log_event( - &state.db, - EventLog { - event_type: "backfill.completed".to_string(), - severity: Severity::Info, - actor_did: None, - subject: body.collection.clone(), - detail: serde_json::json!({ - "job_id": job_id, - "total_repos": 0, - "total_records": 0, - }), - }, - backend, - ) - .await; - return; - } - - // Resolve DIDs to PDS endpoints and group by PDS - let mut pds_to_dids: HashMap> = HashMap::new(); - - for did in &all_dids { - match profile::resolve_pds_endpoint(&state.http, &state.config.plc_url, did).await { - Ok(pds) => { - pds_to_dids.entry(pds).or_default().push(did.clone()); - } - Err(e) => { - tracing::warn!(did, error = %e, "failed to resolve PDS endpoint, skipping DID"); - } - } - } - - let processed_repos = Arc::new(AtomicI32::new(0)); - let total_records = Arc::new(AtomicI32::new(0)); - - let state = Arc::new(state); - let collections = Arc::new(collections); - let job_id_arc = Arc::new(job_id.clone()); - - // Process PDSes with nested concurrency - let pds_entries: Vec<(String, Vec)> = pds_to_dids.into_iter().collect(); - - stream::iter(pds_entries) - .for_each_concurrent(10, |(pds_endpoint, dids)| { - let state = Arc::clone(&state); - let collections = Arc::clone(&collections); - let processed_repos = Arc::clone(&processed_repos); - let total_records = Arc::clone(&total_records); - let job_id = Arc::clone(&job_id_arc); - - async move { - stream::iter(dids) - .for_each_concurrent(3, |did| { - let state = Arc::clone(&state); - let collections = Arc::clone(&collections); - let processed_repos = Arc::clone(&processed_repos); - let total_records = Arc::clone(&total_records); - let pds_endpoint = pds_endpoint.clone(); - let job_id = Arc::clone(&job_id); - - async move { - for collection in collections.iter() { - match fetch_records_from_pds( - &state, - &pds_endpoint, - &did, - collection, - ) - .await - { - Ok(count) => { - total_records - .fetch_add(count as i32, Ordering::Relaxed); - } - Err(e) => { - tracing::warn!( - did, - collection, - pds = %pds_endpoint, - error = %e, - "failed to fetch records from PDS" - ); - } - } - } - - let repos = processed_repos.fetch_add(1, Ordering::Relaxed) + 1; - - // Update DB progress every 100 repos - if repos % 100 == 0 { - let records = total_records.load(Ordering::Relaxed); - let backend = state.db_backend; - let sql = adapt_sql( - "UPDATE backfill_jobs SET processed_repos = ?, total_records = ? WHERE id = ?", - backend, - ); - let _ = sqlx::query(&sql) - .bind(repos) - .bind(records) - .bind(job_id.as_str()) - .execute(&state.db) - .await; - } - } - }) - .await; - } - }) - .await; + .fetch_one(&state.db) + .await + .map(|(c,)| c) + .unwrap_or(0); - let final_processed = processed_repos.load(Ordering::Relaxed); - let final_records = total_records.load(Ordering::Relaxed); + let sql = adapt_sql( + "SELECT total_records FROM backfill_jobs WHERE id = ?", + state.db_backend, + ); + let final_records: i32 = sqlx::query_as::<_, (i32,)>(&sql) + .bind(&job_id) + .fetch_one(&state.db) + .await + .map(|(c,)| c) + .unwrap_or(0); complete_job(&state, &job_id, final_processed, final_records, None).await; @@ -457,7 +620,7 @@ async fn run_backfill_job(state: AppState, job_id: String, body: CreateBackfillB event_type: "backfill.completed".to_string(), severity: Severity::Info, actor_did: None, - subject: body.collection.clone(), + subject: collection, detail: serde_json::json!({ "job_id": job_id, "total_repos": final_processed, @@ -470,46 +633,64 @@ async fn run_backfill_job(state: AppState, job_id: String, body: CreateBackfillB } // --------------------------------------------------------------------------- -// Helper functions +// Admin handlers // --------------------------------------------------------------------------- -async fn fail_job(state: &AppState, job_id: &str, error: &str) { - let now = now_rfc3339(); +/// POST /admin/backfill — create a backfill job and spawn background work. +pub(super) async fn create_backfill( + State(state): State, + admin: UserAuth, + Json(body): Json, +) -> Result<(StatusCode, Json), AppError> { + admin.require(Permission::BackfillCreate).await?; let backend = state.db_backend; + + let now = now_rfc3339(); + let job_id = Uuid::new_v4().to_string(); let sql = adapt_sql( - "UPDATE backfill_jobs SET status = 'failed', completed_at = ?, error = ? WHERE id = ?", + "INSERT INTO backfill_jobs (id, collection, did, status, stage, started_at, created_at) VALUES (?, ?, ?, 'running', 'pending', ?, ?) RETURNING id", backend, ); - let _ = sqlx::query(&sql) + let row: (String,) = sqlx::query_as(&sql) + .bind(&job_id) + .bind(&body.collection) + .bind(&body.did) .bind(&now) - .bind(error) - .bind(job_id) - .execute(&state.db) - .await; -} + .bind(&now) + .fetch_one(&state.db) + .await + .map_err(|e| AppError::Internal(format!("failed to create backfill job: {e}")))?; -async fn complete_job( - state: &AppState, - job_id: &str, - processed_repos: i32, - total_records: i32, - error: Option<&str>, -) { - let now = now_rfc3339(); - let backend = state.db_backend; + let job_id = row.0.clone(); - let sql = adapt_sql( - "UPDATE backfill_jobs SET status = 'completed', completed_at = ?, processed_repos = ?, total_records = ?, error = ? WHERE id = ?", + log_event( + &state.db, + EventLog { + event_type: "backfill.started".to_string(), + severity: Severity::Info, + actor_did: Some(admin.did.clone()), + subject: body.collection.clone(), + detail: serde_json::json!({ + "job_id": job_id.clone(), + }), + }, backend, - ); - let _ = sqlx::query(&sql) - .bind(&now) - .bind(processed_repos) - .bind(total_records) - .bind(error) - .bind(job_id) - .execute(&state.db) - .await; + ) + .await; + + let spawn_state = state.clone(); + let spawn_job_id = job_id.clone(); + tokio::spawn(async move { + run_backfill_job(spawn_state, spawn_job_id).await; + }); + + Ok(( + StatusCode::CREATED, + Json(serde_json::json!({ + "id": job_id, + "status": "running", + })), + )) } /// GET /admin/backfill/status — list all backfill jobs. @@ -521,7 +702,7 @@ pub(super) async fn backfill_status( let backend = state.db_backend; let sql = adapt_sql( - "SELECT id, collection, did, status, total_repos, processed_repos, total_records, error, started_at, completed_at, created_at FROM backfill_jobs ORDER BY created_at DESC", + "SELECT id, collection, did, status, stage, total_repos, processed_repos, total_records, error, started_at, completed_at, created_at FROM backfill_jobs ORDER BY created_at DESC", backend, ); #[allow(clippy::type_complexity)] @@ -530,6 +711,7 @@ pub(super) async fn backfill_status( Option, Option, String, + String, Option, Option, Option, @@ -550,6 +732,7 @@ pub(super) async fn backfill_status( collection, did, status, + stage, total_repos, processed_repos, total_records, @@ -563,6 +746,7 @@ pub(super) async fn backfill_status( collection, did, status, + stage, total_repos, processed_repos, total_records, @@ -577,3 +761,27 @@ pub(super) async fn backfill_status( Ok(Json(jobs)) } + +// --------------------------------------------------------------------------- +// Startup resumption +// --------------------------------------------------------------------------- + +/// Resume any backfill jobs that were running when the server last stopped. +pub async fn resume_backfill_jobs(state: &AppState) { + let sql = adapt_sql( + "SELECT id FROM backfill_jobs WHERE status = 'running'", + state.db_backend, + ); + let rows: Vec<(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; + }); + } +} diff --git a/src/admin/mod.rs b/src/admin/mod.rs index 6a65c66..fb308f4 100644 --- a/src/admin/mod.rs +++ b/src/admin/mod.rs @@ -1,7 +1,7 @@ mod api_clients; mod api_keys; pub(crate) mod auth; -mod backfill; +pub mod backfill; mod dead_letters; mod domains; mod events; diff --git a/src/admin/types.rs b/src/admin/types.rs index c75bcaf..4871b2e 100644 --- a/src/admin/types.rs +++ b/src/admin/types.rs @@ -69,18 +69,19 @@ pub(super) struct CreateBackfillBody { } #[derive(Serialize)] -pub(super) struct BackfillJob { - pub(super) id: String, - pub(super) collection: Option, - pub(super) did: Option, - pub(super) status: String, - pub(super) total_repos: Option, - pub(super) processed_repos: Option, - pub(super) total_records: Option, - pub(super) error: Option, - pub(super) started_at: Option, - pub(super) completed_at: Option, - pub(super) created_at: String, +pub(crate) struct BackfillJob { + pub(crate) id: String, + pub(crate) collection: Option, + pub(crate) did: Option, + pub(crate) status: String, + pub(crate) stage: String, + pub(crate) total_repos: Option, + pub(crate) processed_repos: Option, + pub(crate) total_records: Option, + pub(crate) error: Option, + pub(crate) started_at: Option, + pub(crate) completed_at: Option, + pub(crate) created_at: String, } // --------------------------------------------------------------------------- diff --git a/src/main.rs b/src/main.rs index 33b998a..16f92fe 100644 --- a/src/main.rs +++ b/src/main.rs @@ -644,6 +644,8 @@ async fn main() { state.db_backend, )); + happyview::admin::backfill::resume_backfill_jobs(&state).await; + let app = server::router(state); let addr = config.listen_addr(); diff --git a/web/src/app/dashboard/backfill/page.tsx b/web/src/app/dashboard/backfill/page.tsx index 74585b5..89e92bc 100644 --- a/web/src/app/dashboard/backfill/page.tsx +++ b/web/src/app/dashboard/backfill/page.tsx @@ -77,8 +77,7 @@ export default function BackfillPage() { ID Collection DID - Progress - Records + Status Started @@ -86,7 +85,7 @@ export default function BackfillPage() { {jobs.length === 0 && ( No backfill jobs yet. @@ -105,12 +104,7 @@ export default function BackfillPage() { {job.did ?? "All"} - {job.processed_repos != null && job.total_repos != null - ? `${job.processed_repos} / ${job.total_repos}` - : "--"} - - - {job.total_records?.toLocaleString() ?? "--"} + {job.started_at @@ -127,11 +121,51 @@ export default function BackfillPage() { ); } -function CreateDialog({ - onSuccess, -}: { - onSuccess: () => void; -}) { +function StageDisplay({ job }: { job: BackfillJob }) { + const repos = job.total_repos?.toLocaleString() ?? "0"; + const processed = job.processed_repos?.toLocaleString() ?? "0"; + const records = job.total_records?.toLocaleString() ?? "0"; + + switch (job.stage) { + case "pending": + return Pending; + case "discovering_repos": + return ( + + Discovering repos… +
{repos} found +
+ ); + case "resolving_pds": + return ( + + Resolving PDS… ({processed} / {repos}) + + ); + case "fetching_records": + return ( + + Fetching records… ({processed} / {repos} repos, {records} records) + + ); + case "completed": + return ( + + Completed — {repos} repos, {records} records + + ); + case "failed": + return ( + + Failed{job.error ? ` — ${job.error}` : ""} + + ); + default: + return {job.stage}; + } +} + +function CreateDialog({ onSuccess }: { onSuccess: () => void }) { const [collection, setCollection] = useState(null); const [did, setDid] = useState(""); const [error, setError] = useState(null); @@ -177,7 +211,11 @@ function CreateDialog({ { const target = e.target as HTMLElement; - if (target.closest("[data-slot='combobox-item'], [data-slot='combobox-content']")) { + if ( + target.closest( + "[data-slot='combobox-item'], [data-slot='combobox-content']", + ) + ) { e.preventDefault(); } }} diff --git a/web/src/types/backfill.ts b/web/src/types/backfill.ts index f65ea7c..5fea8a7 100644 --- a/web/src/types/backfill.ts +++ b/web/src/types/backfill.ts @@ -3,6 +3,7 @@ export interface BackfillJob { collection: string | null did: string | null status: string + stage: string total_repos: number | null processed_repos: number | null total_records: number | null