A lexicon-driven AppView for ATProto.
0

Configure Feed

Select the types of activity you want to include in your feed.

fix: resolve backfill fetcher deadlock when PDS hosts exceed concurrency limit

+107 -514
+107 -514
src/admin/backfill.rs
··· 5 5 use axum::Json; 6 6 use axum::extract::{Path, State}; 7 7 use axum::http::StatusCode; 8 - use futures_util::FutureExt; 9 - use futures_util::stream::{self, FuturesUnordered, StreamExt}; 8 + use futures_util::stream::{self, StreamExt}; 10 9 use serde::Deserialize; 11 10 use serde_json::Value; 12 - use tokio::sync::mpsc; 13 11 use uuid::Uuid; 14 12 15 13 use rand::RngExt; ··· 459 457 } 460 458 461 459 // --------------------------------------------------------------------------- 462 - // Pipelined Phase 2+3: Resolve PDS endpoints and fetch records concurrently 460 + // Phase 2: Resolve PDS endpoints for all discovered repos 463 461 // --------------------------------------------------------------------------- 464 462 465 - async fn run_pipelined_resolve_and_fetch( 466 - state: &AppState, 467 - job_id: &str, 468 - collections: &[String], 469 - concurrency: &BackfillConcurrency, 470 - ) -> (i32, i32) { 471 - set_stage(state, job_id, "resolving_and_fetching").await; 463 + async fn run_resolution_phase(state: &AppState, job_id: &str, concurrency: &BackfillConcurrency) { 464 + set_stage(state, job_id, "resolving_pds").await; 472 465 473 - // Count already-resolved and already-completed repos for accurate progress 466 + // Count already-resolved repos so a resumed job reports accurate progress. 474 467 let already_resolved: i32 = { 475 468 let sql = adapt_sql( 476 469 "SELECT COUNT(*) FROM happyview_backfill_repos WHERE job_id = ? AND pds_endpoint IS NOT NULL", ··· 483 476 .map(|(c,)| c) 484 477 .unwrap_or(0) 485 478 }; 486 - 487 - let already_completed: i32 = { 488 - let sql = adapt_sql( 489 - "SELECT COUNT(*) FROM happyview_backfill_repos WHERE job_id = ? AND status = 'completed'", 490 - state.db_backend, 491 - ); 492 - crate::db::query_as::<(i32,)>(&sql) 493 - .bind(job_id) 494 - .fetch_one(&state.backfill_db) 495 - .await 496 - .map(|(c,)| c) 497 - .unwrap_or(0) 498 - }; 499 - 500 479 update_job_counter(state, job_id, "resolved_repos", already_resolved).await; 501 - update_job_counter(state, job_id, "processed_repos", already_completed).await; 502 480 503 - let existing_records: i32 = { 504 - let sql = adapt_sql( 505 - "SELECT total_records FROM happyview_backfill_jobs WHERE id = ?", 506 - state.db_backend, 507 - ); 508 - crate::db::query_as::<(Option<i32>,)>(&sql) 509 - .bind(job_id) 510 - .fetch_one(&state.backfill_db) 511 - .await 512 - .map(|(c,)| c.unwrap_or(0)) 513 - .unwrap_or(0) 514 - }; 515 - 516 - // Shared atomics for lock-free counter updates 517 - let resolved_repos = Arc::new(AtomicI32::new(already_resolved)); 518 - let processed_repos = Arc::new(AtomicI32::new(already_completed)); 519 - let total_records = Arc::new(AtomicI32::new(existing_records)); 520 - let cancelled = Arc::new(AtomicBool::new(false)); 521 - 522 - let (tx, mut rx) = mpsc::channel::<(String, String)>(256); 523 - let tx_resolver = tx.clone(); 524 - let tx_backlog = tx.clone(); 525 - 526 - // --- Resolver task --- 527 - let resolution_concurrency = concurrency.resolution; 528 - let resolver_state = state.clone(); 529 - let resolver_job_id = job_id.to_string(); 530 - let resolver_resolved = Arc::clone(&resolved_repos); 531 - let resolver_cancelled = Arc::clone(&cancelled); 532 - 533 - let resolver_handle = tokio::spawn(async move { 534 - let sql = adapt_sql( 535 - "SELECT did FROM happyview_backfill_repos WHERE job_id = ? AND pds_endpoint IS NULL", 536 - resolver_state.db_backend, 537 - ); 538 - let unresolved: Vec<(String,)> = crate::db::query_as(&sql) 539 - .bind(&resolver_job_id) 540 - .fetch_all(&resolver_state.backfill_db) 541 - .await 542 - .unwrap_or_default(); 543 - 544 - let mut attempted: i32 = 0; 545 - let mut next_flush = random_batch_threshold(100); 546 - let mut next_cancel_check = random_batch_threshold(10); 547 - 548 - let stream_state = resolver_state.clone(); 549 - let stream_cancelled = Arc::clone(&resolver_cancelled); 550 - let mut results = stream::iter(unresolved) 551 - .map(move |(did,)| { 552 - let state = stream_state.clone(); 553 - let cancelled = Arc::clone(&stream_cancelled); 554 - async move { 555 - if cancelled.load(Ordering::Relaxed) { 556 - return None; 557 - } 558 - // Bound the entire resolution of one DID (DNS, connect, and 559 - // any rate-limit retry loop) so a single stuck DID can never 560 - // hang the resolver stream. On expiry the DID is skipped. 561 - let result = match tokio::time::timeout( 562 - crate::http_retry::RESOLVE_DEADLINE, 563 - profile::resolve_pds_endpoint(&state.http, &state.config.plc_url, &did), 564 - ) 565 - .await 566 - { 567 - Ok(result) => result, 568 - Err(_) => Err(AppError::Internal(format!( 569 - "PDS resolution timed out after {}s", 570 - crate::http_retry::RESOLVE_DEADLINE.as_secs() 571 - ))), 572 - }; 573 - Some((did, result)) 574 - } 575 - }) 576 - .buffer_unordered(resolution_concurrency); 577 - 578 - while let Some(item) = results.next().await { 579 - let Some((did, result)) = item else { 580 - break; 581 - }; 582 - 583 - match result { 584 - Ok(pds) => { 585 - let sql = adapt_sql( 586 - "UPDATE happyview_backfill_repos SET pds_endpoint = ? WHERE job_id = ? AND did = ?", 587 - resolver_state.db_backend, 588 - ); 589 - let _ = crate::db::query(&sql) 590 - .bind(&pds) 591 - .bind(&resolver_job_id) 592 - .bind(&did) 593 - .execute(&resolver_state.backfill_db) 594 - .await; 595 - 596 - publish_event( 597 - &resolver_state, 598 - super::types::BackfillEvent::RepoResolved { 599 - job_id: resolver_job_id.clone(), 600 - did: did.clone(), 601 - pds_endpoint: pds.clone(), 602 - }, 603 - ); 604 - 605 - let count = resolver_resolved.fetch_add(1, Ordering::Relaxed) + 1; 606 - if count >= next_flush { 607 - update_job_counter( 608 - &resolver_state, 609 - &resolver_job_id, 610 - "resolved_repos", 611 - count, 612 - ) 613 - .await; 614 - next_flush = count + random_batch_threshold(100); 615 - } 616 - publish_event( 617 - &resolver_state, 618 - super::types::BackfillEvent::JobCounters { 619 - job_id: resolver_job_id.clone(), 620 - total_repos: None, 621 - resolved_repos: Some(count), 622 - processed_repos: None, 623 - total_records: None, 624 - }, 625 - ); 626 - 627 - if tx_resolver.send((did, pds)).await.is_err() { 628 - break; 629 - } 630 - } 631 - Err(e) => { 632 - tracing::warn!(did, error = %e, "failed to resolve PDS endpoint, skipping DID"); 633 - } 634 - } 635 - 636 - attempted += 1; 637 - if attempted >= next_cancel_check { 638 - if should_stop_worker(&resolver_state, &resolver_job_id).await { 639 - resolver_cancelled.store(true, Ordering::Relaxed); 640 - break; 641 - } 642 - next_cancel_check = attempted + random_batch_threshold(10); 643 - } 644 - } 645 - 646 - // Persist final resolved count 647 - let final_resolved = resolver_resolved.load(Ordering::Relaxed); 648 - update_job_counter( 649 - &resolver_state, 650 - &resolver_job_id, 651 - "resolved_repos", 652 - final_resolved, 653 - ) 654 - .await; 655 - // tx is dropped here, signalling the fetcher that no more DIDs are coming 656 - }); 657 - 658 - // --- Also send already-resolved-but-unfetched DIDs to the fetcher --- 659 - let pending_sql = adapt_sql( 660 - "SELECT did, pds_endpoint FROM happyview_backfill_repos WHERE job_id = ? AND status = 'pending' AND pds_endpoint IS NOT NULL", 481 + let sql = adapt_sql( 482 + "SELECT did FROM happyview_backfill_repos WHERE job_id = ? AND pds_endpoint IS NULL", 661 483 state.db_backend, 662 484 ); 663 - let pending_rows: Vec<(String, String)> = crate::db::query_as(&pending_sql) 485 + let unresolved: Vec<(String,)> = crate::db::query_as(&sql) 664 486 .bind(job_id) 665 487 .fetch_all(&state.backfill_db) 666 488 .await 667 489 .unwrap_or_default(); 668 490 669 - let backlog_cancelled = Arc::clone(&cancelled); 670 - let backlog_handle = tokio::spawn(async move { 671 - for (did, pds) in pending_rows { 672 - if backlog_cancelled.load(Ordering::Relaxed) { 673 - break; 674 - } 675 - if tx_backlog.send((did, pds)).await.is_err() { 676 - break; 677 - } 678 - } 679 - }); 491 + let resolved = Arc::new(AtomicI32::new(already_resolved)); 492 + let cancelled = Arc::new(AtomicBool::new(false)); 680 493 681 - // Drop our copy of tx so the channel closes when both senders finish 682 - drop(tx); 683 - 684 - // --- Fetcher: receive (did, pds) pairs and dispatch to PDS workers --- 685 - // Each PDS gets its own worker with a DID channel. Workers acquire a 686 - // semaphore permit before starting, limiting concurrent PDS connections. 687 - // We never hold the workers lock across an `.await` — use `try_send` to 688 - // avoid blocking when a worker's channel is full (overflow goes to a 689 - // retry queue drained on each iteration). 690 - let state = Arc::new(state.clone()); 691 - let collections = Arc::new(collections.to_vec()); 692 - let job_id_arc = Arc::new(job_id.to_string()); 494 + let mut attempted: i32 = 0; 495 + let mut next_flush = random_batch_threshold(100); 496 + let mut next_cancel_check = random_batch_threshold(10); 693 497 694 - let pds_semaphore = Arc::new(tokio::sync::Semaphore::new(concurrency.pds)); 695 - let mut pds_workers: HashMap<String, mpsc::Sender<String>> = HashMap::new(); 696 - let mut worker_handles = FuturesUnordered::new(); 697 - let mut overflow: Vec<(String, String)> = Vec::new(); 698 - 699 - loop { 700 - if cancelled.load(Ordering::Relaxed) { 701 - break; 702 - } 703 - 704 - let poll_state = Arc::clone(&state); 705 - let poll_job_id = Arc::clone(&job_id_arc); 706 - let poll_cancelled = Arc::clone(&cancelled); 707 - let pair = tokio::select! { 708 - result = rx.recv() => result, 709 - _ = async move { 710 - loop { 711 - tokio::time::sleep(std::time::Duration::from_millis(500)).await; 712 - if poll_cancelled.load(Ordering::Relaxed) || should_stop_worker(&poll_state, &poll_job_id).await { 713 - poll_cancelled.store(true, Ordering::Relaxed); 714 - return; 715 - } 498 + let stream_state = state.clone(); 499 + let stream_cancelled = Arc::clone(&cancelled); 500 + let mut results = stream::iter(unresolved) 501 + .map(move |(did,)| { 502 + let state = stream_state.clone(); 503 + let cancelled = Arc::clone(&stream_cancelled); 504 + async move { 505 + if cancelled.load(Ordering::Relaxed) { 506 + return None; 716 507 } 717 - } => None, 718 - }; 719 - let Some((did, pds_endpoint)) = pair else { 720 - break; 721 - }; 722 - 723 - // Also drain any overflow from previous iterations 724 - overflow.push((did, pds_endpoint)); 725 - 726 - let mut still_pending = Vec::new(); 727 - for (did, pds_endpoint) in overflow.drain(..) { 728 - if cancelled.load(Ordering::Relaxed) { 729 - break; 730 - } 731 - 732 - // Try to send to an existing PDS worker 733 - if let Some(pds_tx) = pds_workers.get(&pds_endpoint) { 734 - match pds_tx.try_send(did.clone()) { 735 - Ok(()) => continue, 736 - Err(mpsc::error::TrySendError::Full(_)) => { 737 - still_pending.push((did, pds_endpoint)); 738 - continue; 739 - } 740 - Err(mpsc::error::TrySendError::Closed(_)) => { 741 - // Worker finished, will be removed below 742 - } 743 - } 508 + // Bound the entire resolution of one DID (DNS, connect, and any 509 + // rate-limit retry loop) so a single stuck DID can never hang 510 + // the resolver stream. On expiry the DID is skipped. 511 + let result = match tokio::time::timeout( 512 + crate::http_retry::RESOLVE_DEADLINE, 513 + profile::resolve_pds_endpoint(&state.http, &state.config.plc_url, &did), 514 + ) 515 + .await 516 + { 517 + Ok(result) => result, 518 + Err(_) => Err(AppError::Internal(format!( 519 + "PDS resolution timed out after {}s", 520 + crate::http_retry::RESOLVE_DEADLINE.as_secs() 521 + ))), 522 + }; 523 + Some((did, result)) 744 524 } 525 + }) 526 + .buffer_unordered(concurrency.resolution); 745 527 746 - // Remove stale workers whose channels have closed 747 - pds_workers.retain(|_, tx| !tx.is_closed()); 748 - 749 - // Spawn a new PDS worker 750 - let permit = Arc::clone(&pds_semaphore); 751 - let (pds_tx, pds_rx) = mpsc::channel::<String>(64); 752 - let _ = pds_tx.try_send(did); 753 - pds_workers.insert(pds_endpoint.clone(), pds_tx); 754 - 755 - let ctx = FetchContext { 756 - state: Arc::clone(&state), 757 - job_id: Arc::clone(&job_id_arc), 758 - collections: Arc::clone(&collections), 759 - processed_repos: Arc::clone(&processed_repos), 760 - total_records: Arc::clone(&total_records), 761 - cancelled: Arc::clone(&cancelled), 762 - dids_per_pds: concurrency.dids_per_pds, 763 - }; 764 - 765 - worker_handles.push(tokio::spawn(async move { 766 - let _permit = permit 767 - .acquire() 768 - .await 769 - .expect("semaphore should not be closed"); 770 - 771 - run_pds_worker(ctx, pds_endpoint, pds_rx).await; 772 - })); 773 - } 774 - overflow = still_pending; 775 - 776 - // Drain any completed worker handles to avoid unbounded accumulation 777 - while let Some(result) = worker_handles.next().now_or_never() { 778 - if let Some(Err(e)) = result { 779 - tracing::warn!(error = %e, "PDS worker task panicked"); 780 - } 781 - } 782 - } 783 - 784 - // Drain remaining overflow after channel closes 785 - for (did, pds_endpoint) in overflow.drain(..) { 786 - if cancelled.load(Ordering::Relaxed) { 528 + while let Some(item) = results.next().await { 529 + let Some((did, result)) = item else { 787 530 break; 788 - } 789 - 790 - // Remove stale workers 791 - pds_workers.retain(|_, tx| !tx.is_closed()); 792 - 793 - if let Some(pds_tx) = pds_workers.get(&pds_endpoint) { 794 - // Channel is bounded; this can block, but all senders are done so it's fine 795 - let _ = pds_tx.send(did).await; 796 - continue; 797 - } 798 - 799 - let permit = Arc::clone(&pds_semaphore); 800 - let (pds_tx, pds_rx) = mpsc::channel::<String>(64); 801 - let _ = pds_tx.try_send(did); 802 - pds_workers.insert(pds_endpoint.clone(), pds_tx); 803 - 804 - let ctx = FetchContext { 805 - state: Arc::clone(&state), 806 - job_id: Arc::clone(&job_id_arc), 807 - collections: Arc::clone(&collections), 808 - processed_repos: Arc::clone(&processed_repos), 809 - total_records: Arc::clone(&total_records), 810 - cancelled: Arc::clone(&cancelled), 811 - dids_per_pds: concurrency.dids_per_pds, 812 531 }; 813 532 814 - worker_handles.push(tokio::spawn(async move { 815 - let _permit = permit 816 - .acquire() 817 - .await 818 - .expect("semaphore should not be closed"); 819 - 820 - run_pds_worker(ctx, pds_endpoint.clone(), pds_rx).await; 821 - })); 822 - } 823 - 824 - // Drop all PDS senders so workers know no more DIDs are coming 825 - drop(pds_workers); 826 - 827 - // Wait for all PDS workers to finish 828 - while let Some(result) = worker_handles.next().await { 829 - if let Err(e) = result { 830 - tracing::warn!(error = %e, "PDS worker task panicked"); 831 - } 832 - } 833 - 834 - // Wait for resolver and backlog tasks 835 - let _ = resolver_handle.await; 836 - let _ = backlog_handle.await; 837 - 838 - let final_repos = processed_repos.load(Ordering::Relaxed); 839 - let final_records = total_records.load(Ordering::Relaxed); 840 - 841 - // Persist final counts 842 - let sql = adapt_sql( 843 - "UPDATE happyview_backfill_jobs SET processed_repos = ?, total_records = ? WHERE id = ?", 844 - state.db_backend, 845 - ); 846 - let _ = crate::db::query(&sql) 847 - .bind(final_repos) 848 - .bind(final_records) 849 - .bind(job_id) 850 - .execute(&state.backfill_db) 851 - .await; 852 - 853 - (final_repos, final_records) 854 - } 855 - 856 - struct FetchContext { 857 - state: Arc<AppState>, 858 - job_id: Arc<String>, 859 - collections: Arc<Vec<String>>, 860 - processed_repos: Arc<AtomicI32>, 861 - total_records: Arc<AtomicI32>, 862 - cancelled: Arc<AtomicBool>, 863 - dids_per_pds: usize, 864 - } 865 - 866 - async fn run_pds_worker(ctx: FetchContext, pds_endpoint: String, mut rx: mpsc::Receiver<String>) { 867 - let FetchContext { 868 - state, 869 - job_id, 870 - collections, 871 - processed_repos, 872 - total_records, 873 - cancelled, 874 - dids_per_pds, 875 - } = ctx; 876 - let mut fetches = FuturesUnordered::new(); 877 - let mut rx_open = true; 878 - let mut next_flush = random_batch_threshold(10); 879 - 880 - loop { 881 - tokio::select! { 882 - biased; 883 - 884 - Some(result) = fetches.next(), if !fetches.is_empty() => { 885 - let (did, records): (String, i32) = result; 886 - total_records.fetch_add(records, Ordering::Relaxed); 887 - 888 - // Mark DID as completed 533 + match result { 534 + Ok(pds) => { 889 535 let sql = adapt_sql( 890 - "UPDATE happyview_backfill_repos SET status = 'completed', records_fetched = ? WHERE job_id = ? AND did = ?", 536 + "UPDATE happyview_backfill_repos SET pds_endpoint = ? WHERE job_id = ? AND did = ?", 891 537 state.db_backend, 892 538 ); 893 539 let _ = crate::db::query(&sql) 894 - .bind(records) 895 - .bind(job_id.as_str()) 540 + .bind(&pds) 541 + .bind(job_id) 896 542 .bind(&did) 897 543 .execute(&state.backfill_db) 898 544 .await; 899 545 900 - publish_event(&state, super::types::BackfillEvent::RepoFetched { 901 - job_id: job_id.to_string(), 902 - did: did.clone(), 903 - pds_endpoint: pds_endpoint.clone(), 904 - records_fetched: records, 905 - }); 546 + publish_event( 547 + state, 548 + super::types::BackfillEvent::RepoResolved { 549 + job_id: job_id.to_string(), 550 + did: did.clone(), 551 + pds_endpoint: pds.clone(), 552 + }, 553 + ); 906 554 907 - let repos = processed_repos.fetch_add(1, Ordering::Relaxed) + 1; 908 - let records = total_records.load(Ordering::Relaxed); 909 - if repos >= next_flush { 910 - let sql = adapt_sql( 911 - "UPDATE happyview_backfill_jobs SET processed_repos = ?, total_records = ? WHERE id = ?", 912 - state.db_backend, 913 - ); 914 - let _ = crate::db::query(&sql) 915 - .bind(repos) 916 - .bind(records) 917 - .bind(job_id.as_str()) 918 - .execute(&state.backfill_db) 919 - .await; 920 - next_flush = repos + random_batch_threshold(10); 555 + let count = resolved.fetch_add(1, Ordering::Relaxed) + 1; 556 + if count >= next_flush { 557 + update_job_counter(state, job_id, "resolved_repos", count).await; 558 + next_flush = count + random_batch_threshold(100); 921 559 } 922 - if cancelled.load(Ordering::Relaxed) || should_stop_worker(&state, job_id.as_str()).await { 923 - cancelled.store(true, Ordering::Relaxed); 924 - break; 925 - } 926 - publish_event(&state, super::types::BackfillEvent::JobCounters { 927 - job_id: job_id.to_string(), 928 - total_repos: None, 929 - resolved_repos: None, 930 - processed_repos: Some(repos), 931 - total_records: Some(records), 932 - }); 560 + publish_event( 561 + state, 562 + super::types::BackfillEvent::JobCounters { 563 + job_id: job_id.to_string(), 564 + total_repos: None, 565 + resolved_repos: Some(count), 566 + processed_repos: None, 567 + total_records: None, 568 + }, 569 + ); 933 570 } 934 - 935 - did = rx.recv(), if rx_open && fetches.len() < dids_per_pds => { 936 - match did { 937 - Some(did) if !cancelled.load(Ordering::Relaxed) => { 938 - let state = Arc::clone(&state); 939 - let collections = collections.clone(); 940 - let pds_endpoint = pds_endpoint.clone(); 941 - let cancelled = Arc::clone(&cancelled); 942 - 943 - fetches.push(async move { 944 - let mut count: i32 = 0; 945 - for collection in collections.iter() { 946 - if cancelled.load(Ordering::Relaxed) { 947 - break; 948 - } 949 - match fetch_records_from_pds( 950 - &state, 951 - &pds_endpoint, 952 - &did, 953 - collection, 954 - &cancelled, 955 - ) 956 - .await 957 - { 958 - Ok(c) => count += c as i32, 959 - Err(e) => { 960 - tracing::warn!( 961 - did, 962 - collection, 963 - pds = %pds_endpoint, 964 - error = %e, 965 - "failed to fetch records from PDS" 966 - ); 967 - } 968 - } 969 - } 970 - (did, count) 971 - }); 972 - } 973 - _ => { 974 - rx_open = false; 975 - } 976 - } 571 + Err(e) => { 572 + tracing::warn!(did, error = %e, "failed to resolve PDS endpoint, skipping DID"); 977 573 } 574 + } 978 575 979 - else => break, 576 + attempted += 1; 577 + if attempted >= next_cancel_check { 578 + if should_stop_worker(state, job_id).await { 579 + cancelled.store(true, Ordering::Relaxed); 580 + break; 581 + } 582 + next_cancel_check = attempted + random_batch_threshold(10); 980 583 } 981 584 } 982 585 983 - // Drain any remaining fetches 984 - while let Some(result) = fetches.next().await { 985 - let (did, records): (String, i32) = result; 986 - total_records.fetch_add(records, Ordering::Relaxed); 987 - 988 - let sql = adapt_sql( 989 - "UPDATE happyview_backfill_repos SET status = 'completed', records_fetched = ? WHERE job_id = ? AND did = ?", 990 - state.db_backend, 991 - ); 992 - let _ = crate::db::query(&sql) 993 - .bind(records) 994 - .bind(job_id.as_str()) 995 - .bind(&did) 996 - .execute(&state.backfill_db) 997 - .await; 998 - 999 - publish_event( 1000 - &state, 1001 - super::types::BackfillEvent::RepoFetched { 1002 - job_id: job_id.to_string(), 1003 - did: did.clone(), 1004 - pds_endpoint: pds_endpoint.clone(), 1005 - records_fetched: records, 1006 - }, 1007 - ); 1008 - 1009 - processed_repos.fetch_add(1, Ordering::Relaxed); 1010 - } 586 + // Persist the final resolved count regardless of the last flush threshold. 587 + let final_resolved = resolved.load(Ordering::Relaxed); 588 + update_job_counter(state, job_id, "resolved_repos", final_resolved).await; 1011 589 } 1012 590 1013 591 // --------------------------------------------------------------------------- 1014 - // Phase 3: Fetch records from PDS instances (legacy, for resumed jobs) 592 + // Phase 3: Fetch records from PDS instances, grouped by PDS 1015 593 // --------------------------------------------------------------------------- 1016 594 1017 595 async fn run_fetching_phase( ··· 1528 1106 } 1529 1107 1530 1108 let concurrency = load_concurrency(&state).await; 1531 - let (final_processed, final_records) = if matches!( 1532 - stage.as_str(), 1533 - "pending" | "discovering_repos" | "resolving_pds" | "resolving_and_fetching" 1534 - ) { 1535 - run_pipelined_resolve_and_fetch(&state, &job_id, &collections, &concurrency).await 1536 - } else { 1537 - // stage == "fetching_records": resolution already done (legacy or resumed) 1538 - run_fetching_phase(&state, &job_id, &collections, &concurrency).await 1539 - }; 1109 + 1110 + // Resolve PDS endpoints for all discovered repos. Skipped only when a 1111 + // resumed job has already advanced to the fetching stage; resolution is 1112 + // idempotent (repos with a pds_endpoint already set are left untouched). 1113 + if stage.as_str() != "fetching_records" { 1114 + run_resolution_phase(&state, &job_id, &concurrency).await; 1115 + 1116 + match should_stop(&state, &job_id).await { 1117 + Some("cancelling") => { 1118 + tracing::info!(job_id, "backfill job cancelled"); 1119 + finalise_cancel(&state, &job_id).await; 1120 + return; 1121 + } 1122 + Some("pausing") => { 1123 + tracing::info!(job_id, "backfill job paused"); 1124 + finalise_pause(&state, &job_id).await; 1125 + return; 1126 + } 1127 + _ => {} 1128 + } 1129 + } 1130 + 1131 + let (final_processed, final_records) = 1132 + run_fetching_phase(&state, &job_id, &collections, &concurrency).await; 1540 1133 1541 1134 match should_stop(&state, &job_id).await { 1542 1135 Some("cancelling") => {