diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -160,6 +160,182 @@ - name: Clippy run: cargo clippy --all-targets -- -D warnings # --------------------------------------------------------------------------- + # PR binaries — build on every PR push so reviewers can test + # --------------------------------------------------------------------------- + pr-binary: + needs: changes + if: >- + github.event_name == 'pull_request' + && needs.changes.outputs.server == 'true' + runs-on: ${{ matrix.runs-on }} + strategy: + fail-fast: false + matrix: + include: + - runs-on: depot-ubuntu-24.04 + arch: amd64 + - runs-on: depot-ubuntu-24.04-arm + arch: arm64 + steps: + - name: Checkout repository + uses: actions/checkout@v6 + + - name: Cache cargo + uses: actions/cache@v5 + with: + path: | + ~/.cargo/registry + ~/.cargo/git + key: cargo-release-${{ matrix.arch }}-${{ hashFiles('Cargo.lock') }} + restore-keys: | + cargo-release-${{ matrix.arch }}- + + - name: Build release binary + run: | + docker run --rm \ + -v "$PWD:/app" \ + -v "$HOME/.cargo/registry:/usr/local/cargo/registry" \ + -v "$HOME/.cargo/git:/usr/local/cargo/git" \ + -w /app \ + -e SQLX_OFFLINE=true \ + rust:1.93-bookworm \ + cargo build --release + + - name: Upload artifact + uses: actions/upload-artifact@v4 + with: + name: happyview-linux-${{ matrix.arch }} + path: target/release/happyview + + pr-docker: + needs: pr-binary + runs-on: ${{ matrix.runs-on }} + permissions: + contents: read + packages: write + strategy: + fail-fast: false + matrix: + include: + - platform: linux/amd64 + runs-on: depot-ubuntu-24.04 + arch: amd64 + - platform: linux/arm64 + runs-on: depot-ubuntu-24.04-arm + arch: arm64 + env: + GHCR_IMAGE: ghcr.io/${{ github.repository }} + steps: + - name: Checkout repository + uses: actions/checkout@v6 + + - name: Download pre-built binary + uses: actions/download-artifact@v4 + with: + name: happyview-linux-${{ matrix.arch }} + path: .binary + + - name: Prepare binary override + run: | + mkdir -p .builder-override/app/target/release + cp .binary/happyview .builder-override/app/target/release/happyview + + - name: Set up Docker Buildx + uses: docker/setup-buildx-action@v3 + + - name: Log in to GitHub Container Registry + uses: docker/login-action@v3 + with: + registry: ghcr.io + username: ${{ github.actor }} + password: ${{ secrets.GITHUB_TOKEN }} + + - name: Extract metadata + id: meta + uses: docker/metadata-action@v5 + with: + images: ${{ env.GHCR_IMAGE }} + tags: | + type=raw,value=pr-${{ github.event.pull_request.number }} + + - name: Build and push by digest + id: build + uses: docker/build-push-action@v5 + with: + context: . + platforms: ${{ matrix.platform }} + labels: ${{ steps.meta.outputs.labels }} + build-contexts: | + builder=.builder-override + cache-from: type=gha,scope=pr-${{ matrix.platform }} + cache-to: type=gha,scope=pr-${{ matrix.platform }},mode=max,ignore-error=true + outputs: type=image,"name=${{ env.GHCR_IMAGE }}",push-by-digest=true,name-canonical=true,push=true + + - name: Create manifest + run: | + TAGS=$(jq -cr '.tags | map("-t " + .) | join(" ")' <<< '${{ steps.meta.outputs.json }}') + docker buildx imagetools create --append $TAGS \ + ${{ env.GHCR_IMAGE }}@${{ steps.build.outputs.digest }} 2>/dev/null || \ + docker buildx imagetools create $TAGS \ + ${{ env.GHCR_IMAGE }}@${{ steps.build.outputs.digest }} + + pr-build-comment: + needs: pr-docker + runs-on: depot-ubuntu-24.04 + permissions: + pull-requests: write + steps: + - name: Comment on PR + uses: actions/github-script@v7 + with: + script: | + const sha = context.payload.pull_request.head.sha.substring(0, 7); + const prNumber = context.payload.pull_request.number; + const runUrl = `${context.serverUrl}/${context.repo.owner}/${context.repo.repo}/actions/runs/${context.runId}`; + const image = `ghcr.io/${context.repo.owner}/${context.repo.repo}`; + const body = [ + '### PR Build', + '', + `Builds for \`${sha}\`:`, + '', + '**Docker**', + '```sh', + `docker pull ${image}:pr-${prNumber}`, + '```', + '', + '**Binaries**', + '```sh', + `gh run download ${context.runId} -n happyview-linux-amd64`, + `gh run download ${context.runId} -n happyview-linux-arm64`, + '```', + '', + `Or download binaries from the [workflow run](${runUrl}).`, + ].join('\n'); + + const { data: comments } = await github.rest.issues.listComments({ + owner: context.repo.owner, + repo: context.repo.repo, + issue_number: context.issue.number, + }); + const existing = comments.find(c => c.body.startsWith('### PR Build')); + + if (existing) { + await github.rest.issues.updateComment({ + owner: context.repo.owner, + repo: context.repo.repo, + comment_id: existing.id, + body, + }); + } else { + await github.rest.issues.createComment({ + owner: context.repo.owner, + repo: context.repo.repo, + issue_number: context.issue.number, + body, + }); + } + + # --------------------------------------------------------------------------- # Server release # --------------------------------------------------------------------------- release: diff --git a/migrations/postgres/20260520000000_add_resolved_repos.sql b/migrations/postgres/20260520000000_add_resolved_repos.sql new file mode 100644 --- /dev/null +++ b/migrations/postgres/20260520000000_add_resolved_repos.sql @@ -0,0 +1,1 @@ +ALTER TABLE backfill_jobs ADD COLUMN resolved_repos INTEGER; diff --git a/migrations/sqlite/20260520000000_add_resolved_repos.sql b/migrations/sqlite/20260520000000_add_resolved_repos.sql new file mode 100644 --- /dev/null +++ b/migrations/sqlite/20260520000000_add_resolved_repos.sql @@ -0,0 +1,1 @@ +ALTER TABLE backfill_jobs ADD COLUMN resolved_repos INTEGER; diff --git a/src/admin/backfill.rs b/src/admin/backfill.rs --- a/src/admin/backfill.rs +++ b/src/admin/backfill.rs @@ -5,9 +5,11 @@ use axum::Json; use axum::extract::{Path, State}; use axum::http::StatusCode; -use futures_util::stream::{self, StreamExt}; +use futures_util::FutureExt; +use futures_util::stream::{self, FuturesUnordered, StreamExt}; use serde::Deserialize; use serde_json::Value; +use tokio::sync::mpsc; use uuid::Uuid; use crate::AppState; @@ -69,6 +71,7 @@ async fn update_job_counter(state: &AppState, job_id: &str, column: &str, value: i32) { let query = match column { "total_repos" => "UPDATE backfill_jobs SET total_repos = ? WHERE id = ?", + "resolved_repos" => "UPDATE backfill_jobs SET resolved_repos = ? WHERE id = ?", "processed_repos" => "UPDATE backfill_jobs SET processed_repos = ? WHERE id = ?", "total_records" => "UPDATE backfill_jobs SET total_records = ? WHERE id = ?", other => { @@ -317,69 +320,448 @@ Ok(()) } // --------------------------------------------------------------------------- -// Phase 2: Resolve PDS endpoints +// Pipelined Phase 2+3: Resolve PDS endpoints and fetch records concurrently // --------------------------------------------------------------------------- -async fn run_resolution_phase(state: &AppState, job_id: &str) { - set_stage(state, job_id, "resolving_pds").await; +async fn run_pipelined_resolve_and_fetch( + state: &AppState, + job_id: &str, + collections: &[String], +) -> (i32, i32) { + set_stage(state, job_id, "resolving_and_fetching").await; + + // Count already-resolved and already-completed repos for accurate progress + let already_resolved: i32 = { + let sql = adapt_sql( + "SELECT COUNT(*) FROM backfill_repos WHERE job_id = ? AND pds_endpoint IS NOT NULL", + state.db_backend, + ); + sqlx::query_as::<_, (i32,)>(&sql) + .bind(job_id) + .fetch_one(&state.db) + .await + .map(|(c,)| c) + .unwrap_or(0) + }; + + let already_completed: i32 = { + let sql = adapt_sql( + "SELECT COUNT(*) FROM backfill_repos WHERE job_id = ? AND status = 'completed'", + state.db_backend, + ); + sqlx::query_as::<_, (i32,)>(&sql) + .bind(job_id) + .fetch_one(&state.db) + .await + .map(|(c,)| c) + .unwrap_or(0) + }; + + update_job_counter(state, job_id, "resolved_repos", already_resolved).await; + update_job_counter(state, job_id, "processed_repos", already_completed).await; + + 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) + }; + + // Shared atomics for lock-free counter updates + let resolved_repos = Arc::new(AtomicI32::new(already_resolved)); + let processed_repos = Arc::new(AtomicI32::new(already_completed)); + let total_records = Arc::new(AtomicI32::new(existing_records)); + let cancelled = Arc::new(AtomicBool::new(false)); + + let (tx, mut rx) = mpsc::channel::<(String, String)>(256); + let tx_resolver = tx.clone(); + let tx_backlog = tx.clone(); + + // --- Resolver task --- + let resolver_state = state.clone(); + let resolver_job_id = job_id.to_string(); + let resolver_resolved = Arc::clone(&resolved_repos); + let resolver_cancelled = Arc::clone(&cancelled); + + let resolver_handle = tokio::spawn(async move { + let sql = adapt_sql( + "SELECT did FROM backfill_repos WHERE job_id = ? AND pds_endpoint IS NULL", + resolver_state.db_backend, + ); + let unresolved: Vec<(String,)> = sqlx::query_as(&sql) + .bind(&resolver_job_id) + .fetch_all(&resolver_state.db) + .await + .unwrap_or_default(); + + let mut attempted: i32 = 0; + for (did,) in &unresolved { + if resolver_cancelled.load(Ordering::Relaxed) { + break; + } + + match profile::resolve_pds_endpoint( + &resolver_state.http, + &resolver_state.config.plc_url, + did, + ) + .await + { + Ok(pds) => { + let sql = adapt_sql( + "UPDATE backfill_repos SET pds_endpoint = ? WHERE job_id = ? AND did = ?", + resolver_state.db_backend, + ); + let _ = sqlx::query(&sql) + .bind(&pds) + .bind(&resolver_job_id) + .bind(did) + .execute(&resolver_state.db) + .await; + + let count = resolver_resolved.fetch_add(1, Ordering::Relaxed) + 1; + if count % 100 == 0 { + update_job_counter( + &resolver_state, + &resolver_job_id, + "resolved_repos", + count, + ) + .await; + } + + if tx_resolver.send((did.clone(), pds)).await.is_err() { + break; + } + } + Err(e) => { + tracing::warn!(did, error = %e, "failed to resolve PDS endpoint, skipping DID"); + } + } + + attempted += 1; + if attempted % 100 == 0 && is_cancelled(&resolver_state, &resolver_job_id).await { + resolver_cancelled.store(true, Ordering::Relaxed); + break; + } + } + + // Persist final resolved count + let final_resolved = resolver_resolved.load(Ordering::Relaxed); + update_job_counter( + &resolver_state, + &resolver_job_id, + "resolved_repos", + final_resolved, + ) + .await; + // tx is dropped here, signalling the fetcher that no more DIDs are coming + }); - let sql = adapt_sql( - "SELECT did FROM backfill_repos WHERE job_id = ? AND pds_endpoint IS NULL", + // --- Also send already-resolved-but-unfetched DIDs to the fetcher --- + let pending_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 unresolved: Vec<(String,)> = sqlx::query_as(&sql) + let pending_rows: Vec<(String, String)> = sqlx::query_as(&pending_sql) .bind(job_id) .fetch_all(&state.db) .await .unwrap_or_default(); + let backlog_cancelled = Arc::clone(&cancelled); + let backlog_handle = tokio::spawn(async move { + for (did, pds) in pending_rows { + if backlog_cancelled.load(Ordering::Relaxed) { + break; + } + if tx_backlog.send((did, pds)).await.is_err() { + break; + } + } + }); + + // Drop our copy of tx so the channel closes when both senders finish + drop(tx); + + // --- Fetcher: receive (did, pds) pairs and dispatch to PDS workers --- + // Each PDS gets its own worker with a DID channel. Workers acquire a + // semaphore permit before starting, limiting concurrent PDS connections. + // We never hold the workers lock across an `.await` — use `try_send` to + // avoid blocking when a worker's channel is full (overflow goes to a + // retry queue drained on each iteration). + 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_semaphore = Arc::new(tokio::sync::Semaphore::new(10)); + let mut pds_workers: HashMap> = HashMap::new(); + let mut worker_handles = FuturesUnordered::new(); + let mut overflow: Vec<(String, String)> = Vec::new(); + + while let Some((did, pds_endpoint)) = rx.recv().await { + // Also drain any overflow from previous iterations + overflow.push((did, pds_endpoint)); + + let mut still_pending = Vec::new(); + for (did, pds_endpoint) in overflow.drain(..) { + if cancelled.load(Ordering::Relaxed) { + break; + } + + // Try to send to an existing PDS worker + if let Some(pds_tx) = pds_workers.get(&pds_endpoint) { + match pds_tx.try_send(did.clone()) { + Ok(()) => continue, + Err(mpsc::error::TrySendError::Full(_)) => { + still_pending.push((did, pds_endpoint)); + continue; + } + Err(mpsc::error::TrySendError::Closed(_)) => { + // Worker finished, will be removed below + } + } + } + + // Remove stale workers whose channels have closed + pds_workers.retain(|_, tx| !tx.is_closed()); + + // Spawn a new PDS worker + let permit = Arc::clone(&pds_semaphore); + let (pds_tx, pds_rx) = mpsc::channel::(64); + let _ = pds_tx.try_send(did); + pds_workers.insert(pds_endpoint.clone(), pds_tx); + + let ctx = FetchContext { + state: Arc::clone(&state), + job_id: Arc::clone(&job_id_arc), + collections: Arc::clone(&collections), + processed_repos: Arc::clone(&processed_repos), + total_records: Arc::clone(&total_records), + cancelled: Arc::clone(&cancelled), + }; + + worker_handles.push(tokio::spawn(async move { + let _permit = permit + .acquire() + .await + .expect("semaphore should not be closed"); + + run_pds_worker(ctx, pds_endpoint, pds_rx).await; + })); + } + overflow = still_pending; + + // Drain any completed worker handles to avoid unbounded accumulation + while let Some(result) = worker_handles.next().now_or_never() { + if let Some(Err(e)) = result { + tracing::warn!(error = %e, "PDS worker task panicked"); + } + } + } + + // Drain remaining overflow after channel closes + for (did, pds_endpoint) in overflow.drain(..) { + if cancelled.load(Ordering::Relaxed) { + break; + } + + // Remove stale workers + pds_workers.retain(|_, tx| !tx.is_closed()); + + if let Some(pds_tx) = pds_workers.get(&pds_endpoint) { + // Channel is bounded; this can block, but all senders are done so it's fine + let _ = pds_tx.send(did).await; + continue; + } + + let permit = Arc::clone(&pds_semaphore); + let (pds_tx, pds_rx) = mpsc::channel::(64); + let _ = pds_tx.try_send(did); + pds_workers.insert(pds_endpoint.clone(), pds_tx); + + let ctx = FetchContext { + state: Arc::clone(&state), + job_id: Arc::clone(&job_id_arc), + collections: Arc::clone(&collections), + processed_repos: Arc::clone(&processed_repos), + total_records: Arc::clone(&total_records), + cancelled: Arc::clone(&cancelled), + }; + + worker_handles.push(tokio::spawn(async move { + let _permit = permit + .acquire() + .await + .expect("semaphore should not be closed"); + + run_pds_worker(ctx, pds_endpoint.clone(), pds_rx).await; + })); + } + + // Drop all PDS senders so workers know no more DIDs are coming + drop(pds_workers); + + // Wait for all PDS workers to finish + while let Some(result) = worker_handles.next().await { + if let Err(e) = result { + tracing::warn!(error = %e, "PDS worker task panicked"); + } + } + + // Wait for resolver and backlog tasks + let _ = resolver_handle.await; + let _ = backlog_handle.await; + + let final_repos = processed_repos.load(Ordering::Relaxed); + let final_records = total_records.load(Ordering::Relaxed); + + // Persist final counts let sql = adapt_sql( - "SELECT COUNT(*) FROM backfill_repos WHERE job_id = ? AND pds_endpoint IS NOT NULL", + "UPDATE backfill_jobs SET processed_repos = ?, total_records = ? WHERE id = ?", state.db_backend, ); - let already_resolved: i32 = sqlx::query_as::<_, (i32,)>(&sql) + let _ = sqlx::query(&sql) + .bind(final_repos) + .bind(final_records) .bind(job_id) - .fetch_one(&state.db) - .await - .map(|(c,)| c) - .unwrap_or(0); + .execute(&state.db) + .await; + + (final_repos, final_records) +} + +struct FetchContext { + state: Arc, + job_id: Arc, + collections: Arc>, + processed_repos: Arc, + total_records: Arc, + cancelled: Arc, +} + +async fn run_pds_worker(ctx: FetchContext, pds_endpoint: String, mut rx: mpsc::Receiver) { + let FetchContext { + state, + job_id, + collections, + processed_repos, + total_records, + cancelled, + } = ctx; + let mut fetches = FuturesUnordered::new(); + let mut rx_open = true; + + loop { + tokio::select! { + biased; - let mut resolved_count = already_resolved; + Some(result) = fetches.next(), if !fetches.is_empty() => { + let (did, records): (String, i32) = result; + total_records.fetch_add(records, Ordering::Relaxed); - let mut attempted = already_resolved; - for (did,) in &unresolved { - match profile::resolve_pds_endpoint(&state.http, &state.config.plc_url, did).await { - Ok(pds) => { + // Mark DID as completed let sql = adapt_sql( - "UPDATE backfill_repos SET pds_endpoint = ? WHERE job_id = ? AND did = ?", + "UPDATE backfill_repos SET status = 'completed' WHERE job_id = ? AND did = ?", state.db_backend, ); let _ = sqlx::query(&sql) - .bind(&pds) - .bind(job_id) - .bind(did) + .bind(job_id.as_str()) + .bind(&did) .execute(&state.db) .await; - resolved_count += 1; + + let repos = processed_repos.fetch_add(1, Ordering::Relaxed) + 1; + if repos % 10 == 0 { + let records = total_records.load(Ordering::Relaxed); + let sql = adapt_sql( + "UPDATE backfill_jobs SET processed_repos = ?, total_records = ? WHERE id = ?", + state.db_backend, + ); + let _ = sqlx::query(&sql) + .bind(repos) + .bind(records) + .bind(job_id.as_str()) + .execute(&state.db) + .await; + + if is_cancelled(&state, job_id.as_str()).await { + cancelled.store(true, Ordering::Relaxed); + break; + } + } } - Err(e) => { - tracing::warn!(did, error = %e, "failed to resolve PDS endpoint, skipping DID"); - } - } - attempted += 1; - if attempted % 100 == 0 { - update_job_counter(state, job_id, "processed_repos", resolved_count).await; - if is_cancelled(state, job_id).await { - return; + + did = rx.recv(), if rx_open && fetches.len() < 3 => { + match did { + Some(did) if !cancelled.load(Ordering::Relaxed) => { + let state = Arc::clone(&state); + let collections = collections.clone(); + let pds_endpoint = pds_endpoint.clone(); + + fetches.push(async move { + let mut count: i32 = 0; + for collection in collections.iter() { + match fetch_records_from_pds( + &state, + &pds_endpoint, + &did, + collection, + ) + .await + { + Ok(c) => count += c as i32, + Err(e) => { + tracing::warn!( + did, + collection, + pds = %pds_endpoint, + error = %e, + "failed to fetch records from PDS" + ); + } + } + } + (did, count) + }); + } + _ => { + rx_open = false; + } + } } + + else => break, } } - update_job_counter(state, job_id, "processed_repos", resolved_count).await; + // Drain any remaining fetches + while let Some(result) = fetches.next().await { + let (did, records): (String, i32) = result; + total_records.fetch_add(records, Ordering::Relaxed); + + 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; + + processed_repos.fetch_add(1, Ordering::Relaxed); + } } // --------------------------------------------------------------------------- -// Phase 3: Fetch records from PDS instances +// Phase 3: Fetch records from PDS instances (legacy, for resumed jobs) // --------------------------------------------------------------------------- async fn run_fetching_phase(state: &AppState, job_id: &str, collections: &[String]) -> (i32, i32) { @@ -713,20 +1095,15 @@ return; } } - if matches!( + let (final_processed, final_records) = if matches!( stage.as_str(), - "pending" | "discovering_repos" | "resolving_pds" + "pending" | "discovering_repos" | "resolving_pds" | "resolving_and_fetching" ) { - 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; - } - } - - let (final_processed, final_records) = run_fetching_phase(&state, &job_id, &collections).await; + run_pipelined_resolve_and_fetch(&state, &job_id, &collections).await + } else { + // stage == "fetching_records": resolution already done (legacy or resumed) + run_fetching_phase(&state, &job_id, &collections).await + }; if is_cancelled(&state, &job_id).await { tracing::info!(job_id, "backfill job cancelled"); @@ -871,7 +1248,7 @@ auth.require(Permission::BackfillRead).await?; let backend = state.db_backend; let sql = adapt_sql( - "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", + "SELECT id, collection, did, status, stage, total_repos, resolved_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)] @@ -881,6 +1258,7 @@ Option, Option, String, String, + Option, Option, Option, Option, @@ -903,6 +1281,7 @@ did, status, stage, total_repos, + resolved_repos, processed_repos, total_records, error, @@ -917,6 +1296,7 @@ did, status, stage, total_repos, + resolved_repos, processed_repos, total_records, error, diff --git a/src/admin/types.rs b/src/admin/types.rs --- a/src/admin/types.rs +++ b/src/admin/types.rs @@ -76,6 +76,7 @@ pub(crate) did: Option, pub(crate) status: String, pub(crate) stage: String, pub(crate) total_repos: Option, + pub(crate) resolved_repos: Option, pub(crate) processed_repos: Option, pub(crate) total_records: Option, pub(crate) error: Option, diff --git a/web/src/app/dashboard/backfill/page.tsx b/web/src/app/dashboard/backfill/page.tsx --- a/web/src/app/dashboard/backfill/page.tsx +++ b/web/src/app/dashboard/backfill/page.tsx @@ -90,6 +90,7 @@ } } function phaseIndex(stage: string): number { + if (stage === "resolving_and_fetching") return 1; const idx = PROGRESS_PHASES.indexOf( stage as (typeof PROGRESS_PHASES)[number], ); @@ -229,6 +230,9 @@ const isActive = job.status === "running" || job.status === "cancelling"; function hasReached(phase: (typeof PROGRESS_PHASES)[number]): boolean { if (allDone) return true; + if (job.stage === "resolving_and_fetching") { + return phase === "discovering_repos" || phase === "resolving_pds" || phase === "fetching_records"; + } return current >= phaseIndex(phase); } @@ -307,26 +311,36 @@ suffix="repos found" />