From 274ab01d24a71f5fddb8da1c65c02623e4603104 Mon Sep 17 00:00:00 2001 From: Trezy Date: Thu, 14 May 2026 16:09:32 +0000 Subject: [PATCH] fix: bulk insert in discovery, seed records on resume, dedupe retry helper, add cancel tests Signed-off-by: Trezy --- src/http_retry.rs | 23 +++++++++++++++++++++++ src/lib.rs | 1 + src/profile.rs | 22 +--------------------- tests/e2e_admin.rs | 131 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ src/admin/backfill.rs | 78 ++++++++++++++++++++++++++++++++++++++++++------------------------------------ 5 file(s) changed, 198 insertion(s)(+), 57 deletion(s)(-) diff --git a/src/http_retry.rs b/src/http_retry.rs new file mode 100644 --- /dev/null +++ b/src/http_retry.rs @@ -0,0 +1,23 @@ +/// Parse rate-limit sleep duration from response headers. +/// Checks `RateLimit-Reset` (Unix timestamp, used by XRPC servers) first, +/// then `retry-after` (seconds), defaulting to 5s. +pub fn parse_retry_after(headers: &reqwest::header::HeaderMap) -> u64 { + if let Some(reset) = headers + .get("ratelimit-reset") + .and_then(|v| v.to_str().ok()) + .and_then(|v| v.parse::().ok()) + { + let now = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_secs() as i64; + let wait = (reset - now).max(1) as u64; + return wait.min(120); + } + + headers + .get("retry-after") + .and_then(|v| v.to_str().ok()) + .and_then(|v| v.parse::().ok()) + .unwrap_or(5) +} diff --git a/src/lib.rs b/src/lib.rs --- a/src/lib.rs +++ b/src/lib.rs @@ -12,6 +12,7 @@ pub mod external_auth; pub mod feature_flags; pub mod feature_middleware; +pub mod http_retry; pub mod jetstream; pub mod labeler; pub mod lexicon; diff --git a/src/profile.rs b/src/profile.rs --- a/src/profile.rs +++ b/src/profile.rs @@ -1,27 +1,7 @@ use serde::{Deserialize, Serialize}; use crate::error::AppError; - -fn parse_retry_after(headers: &reqwest::header::HeaderMap) -> u64 { - if let Some(reset) = headers - .get("ratelimit-reset") - .and_then(|v| v.to_str().ok()) - .and_then(|v| v.parse::().ok()) - { - let now = std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .unwrap_or_default() - .as_secs() as i64; - let wait = (reset - now).max(1) as u64; - return wait.min(120); - } - - headers - .get("retry-after") - .and_then(|v| v.to_str().ok()) - .and_then(|v| v.parse::().ok()) - .unwrap_or(5) -} +use crate::http_retry::parse_retry_after; #[derive(Serialize)] pub struct Profile { diff --git a/tests/e2e_admin.rs b/tests/e2e_admin.rs --- a/tests/e2e_admin.rs +++ b/tests/e2e_admin.rs @@ -561,6 +561,137 @@ assert_eq!(json.as_array().unwrap().len(), 1); } +#[tokio::test] +#[serial] +async fn backfill_cancel_running_job() { + common::require_db!(); + let app = TestApp::new().await; + let backend = app.state.db_backend; + + // Insert a running job directly so we don't need a real relay. + let job_id = uuid::Uuid::new_v4().to_string(); + let now = now_rfc3339(); + let sql = adapt_sql( + "INSERT INTO backfill_jobs (id, status, stage, started_at, created_at) VALUES (?, 'running', 'discovering_repos', ?, ?)", + backend, + ); + sqlx::query(&sql) + .bind(&job_id) + .bind(&now) + .bind(&now) + .execute(&app.state.db) + .await + .unwrap(); + + let resp = app + .router + .clone() + .oneshot(admin_post( + &format!("/admin/backfill/{job_id}/cancel"), + app.admin_cookie(), + &json!({}), + )) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::OK); + let json = json_body(resp).await; + assert_eq!(json["id"], job_id); + assert_eq!(json["status"], "cancelling"); +} + +#[tokio::test] +#[serial] +async fn backfill_cancel_already_cancelling_is_idempotent() { + common::require_db!(); + let app = TestApp::new().await; + let backend = app.state.db_backend; + + let job_id = uuid::Uuid::new_v4().to_string(); + let now = now_rfc3339(); + let sql = adapt_sql( + "INSERT INTO backfill_jobs (id, status, stage, started_at, created_at) VALUES (?, 'cancelling', 'fetching_records', ?, ?)", + backend, + ); + sqlx::query(&sql) + .bind(&job_id) + .bind(&now) + .bind(&now) + .execute(&app.state.db) + .await + .unwrap(); + + let resp = app + .router + .clone() + .oneshot(admin_post( + &format!("/admin/backfill/{job_id}/cancel"), + app.admin_cookie(), + &json!({}), + )) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::OK); + let json = json_body(resp).await; + assert_eq!(json["status"], "cancelling"); +} + +#[tokio::test] +#[serial] +async fn backfill_cancel_completed_returns_400() { + common::require_db!(); + let app = TestApp::new().await; + let backend = app.state.db_backend; + + let job_id = uuid::Uuid::new_v4().to_string(); + let now = now_rfc3339(); + let sql = adapt_sql( + "INSERT INTO backfill_jobs (id, status, stage, completed_at, created_at) VALUES (?, 'completed', 'completed', ?, ?)", + backend, + ); + sqlx::query(&sql) + .bind(&job_id) + .bind(&now) + .bind(&now) + .execute(&app.state.db) + .await + .unwrap(); + + let resp = app + .router + .clone() + .oneshot(admin_post( + &format!("/admin/backfill/{job_id}/cancel"), + app.admin_cookie(), + &json!({}), + )) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::BAD_REQUEST); +} + +#[tokio::test] +#[serial] +async fn backfill_cancel_not_found_returns_404() { + common::require_db!(); + let app = TestApp::new().await; + + let resp = app + .router + .clone() + .oneshot(admin_post( + "/admin/backfill/nonexistent-id/cancel", + app.admin_cookie(), + &json!({}), + )) + .await + .unwrap(); + + assert_eq!(resp.status(), StatusCode::NOT_FOUND); +} + // --------------------------------------------------------------------------- // Admin management // --------------------------------------------------------------------------- diff --git a/src/admin/backfill.rs b/src/admin/backfill.rs --- a/src/admin/backfill.rs +++ b/src/admin/backfill.rs @@ -14,32 +14,9 @@ use crate::db::{adapt_sql, now_rfc3339}; use crate::error::AppError; use crate::event_log::{EventLog, Severity, log_event}; +use crate::http_retry::parse_retry_after; use crate::profile; use crate::record_handler::{self, RecordEvent}; - -/// Parse rate-limit sleep duration from response headers. -/// Checks `RateLimit-Reset` (Unix timestamp, used by XRPC servers) first, -/// then `retry-after` (seconds), defaulting to 5s. -fn parse_retry_after(headers: &reqwest::header::HeaderMap) -> u64 { - if let Some(reset) = headers - .get("ratelimit-reset") - .and_then(|v| v.to_str().ok()) - .and_then(|v| v.parse::().ok()) - { - let now = std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .unwrap_or_default() - .as_secs() as i64; - let wait = (reset - now).max(1) as u64; - return wait.min(120); - } - - headers - .get("retry-after") - .and_then(|v| v.to_str().ok()) - .and_then(|v| v.parse::().ok()) - .unwrap_or(5) -} use super::auth::UserAuth; use super::permissions::Permission; @@ -249,6 +226,7 @@ ) -> Result<(), String> { let base = state.config.relay_url.trim_end_matches('/'); let mut cursor: Option = None; + let mut running_total: i32 = count_repos(state, job_id).await; loop { let mut url = format!( @@ -287,20 +265,34 @@ let page_count = body.repos.len(); - for repo in &body.repos { - let sql = adapt_sql( - "INSERT INTO backfill_repos (job_id, did) VALUES (?, ?) ON CONFLICT DO NOTHING", - state.db_backend, + if !body.repos.is_empty() { + let base_sql = "INSERT INTO backfill_repos (job_id, did) VALUES "; + let placeholders: Vec = body + .repos + .iter() + .enumerate() + .map(|(i, _)| { + if state.db_backend == crate::db::DatabaseBackend::Postgres { + format!("(${}, ${})", i * 2 + 1, i * 2 + 2) + } else { + "(?, ?)".to_string() + } + }) + .collect(); + let sql = format!( + "{base_sql}{} ON CONFLICT DO NOTHING", + placeholders.join(", ") ); - let _ = sqlx::query(&sql) - .bind(job_id) - .bind(&repo.did) - .execute(&state.db) - .await; + + let mut query = sqlx::query(&sql); + for repo in &body.repos { + query = query.bind(job_id).bind(&repo.did); + } + let _ = query.execute(&state.db).await; } - let total = count_repos(state, job_id).await; - update_job_counter(state, job_id, "total_repos", total).await; + running_total += page_count as i32; + update_job_counter(state, job_id, "total_repos", running_total).await; if is_cancelled(state, job_id).await { return Ok(()); @@ -413,8 +405,22 @@ // Reset processed_repos for the fetching phase update_job_counter(state, job_id, "processed_repos", already_completed).await; + // Seed total_records from DB so a resumed job doesn't lose its prior count + let existing_records: i32 = { + let sql = adapt_sql( + "SELECT total_records FROM backfill_jobs WHERE id = ?", + state.db_backend, + ); + sqlx::query_as::<_, (Option,)>(&sql) + .bind(job_id) + .fetch_one(&state.db) + .await + .map(|(c,)| c.unwrap_or(0)) + .unwrap_or(0) + }; + let processed_repos = Arc::new(AtomicI32::new(already_completed)); - let total_records = Arc::new(AtomicI32::new(0)); + let total_records = Arc::new(AtomicI32::new(existing_records)); let cancelled = Arc::new(AtomicBool::new(false)); let state = Arc::new(state.clone()); let collections = Arc::new(collections.to_vec()); -- tangled.sh