From 36c029c28d3fbeba028814f94f9f80bafb6ab73b Mon Sep 17 00:00:00 2001 From: Mia Date: Wed, 30 Jul 2025 17:37:15 +0000 Subject: [PATCH] feat(consumer): Faster Backfill --- Cargo.lock | 25 +++ consumer/Cargo.toml | 1 + consumer/src/backfill/downloader.rs | 334 ++++++++++++++++++++++++++++ consumer/src/backfill/mod.rs | 134 +++-------- consumer/src/backfill/repo.rs | 23 +- consumer/src/config.rs | 18 +- consumer/src/indexer/mod.rs | 11 +- consumer/src/main.rs | 1 + 8 files changed, 419 insertions(+), 128 deletions(-) create mode 100644 consumer/src/backfill/downloader.rs diff --git a/Cargo.lock b/Cargo.lock index 83fdcb20..2515b514 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -752,6 +752,7 @@ dependencies = [ "did-resolver", "eyre", "figment", + "flume", "foldhash", "futures", "ipld-core", @@ -1339,6 +1340,18 @@ version = "0.5.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1d674e81391d1e1ab681a28d99df07927c6d4aa5b027d7da16ba32d1d21ecd99" +[[package]] +name = "flume" +version = "0.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "da0e4dd2a88388a1f4ccc7c9ce104604dab68d9f408dc34cd45823d5a9069095" +dependencies = [ + "futures-core", + "futures-sink", + "nanorand", + "spin", +] + [[package]] name = "fnv" version = "1.0.7" @@ -2526,6 +2539,15 @@ version = "0.10.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "defc4c55412d89136f966bbb339008b474350e5e6e78d2714439c386b3137a03" +[[package]] +name = "nanorand" +version = "0.7.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6a51313c5820b0b02bd422f4b44776fbf47961755c74ce64afc73bfad10226c3" +dependencies = [ + "getrandom 0.2.15", +] + [[package]] name = "native-tls" version = "0.2.12" @@ -3972,6 +3994,9 @@ name = "spin" version = "0.9.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6980e8d7511241f8acf4aebddbb1ff938df5eebe98691418c4468d0b72a96a67" +dependencies = [ + "lock_api", +] [[package]] name = "spki" diff --git a/consumer/Cargo.toml b/consumer/Cargo.toml index 7869c089..bb655186 100644 --- a/consumer/Cargo.toml +++ b/consumer/Cargo.toml @@ -11,6 +11,7 @@ deadpool-postgres = { version = "0.14.1", features = ["serde"] } did-resolver = { path = "../did-resolver" } eyre = "0.6.12" figment = { version = "0.10.19", features = ["env", "toml"] } +flume = { version = "0.11", features = ["async"] } foldhash = "0.1.4" futures = "0.3.31" ipld-core = "0.4.1" diff --git a/consumer/src/backfill/downloader.rs b/consumer/src/backfill/downloader.rs new file mode 100644 index 00000000..35d37ee4 --- /dev/null +++ b/consumer/src/backfill/downloader.rs @@ -0,0 +1,334 @@ +use super::{DL_DONE_KEY, PDS_SERVICE_ID}; +use crate::db; +use chrono::prelude::*; +use deadpool_postgres::{Client as PgClient, Pool}; +use did_resolver::Resolver; +use futures::TryStreamExt; +use metrics::{counter, histogram}; +use parakeet_db::types::{ActorStatus, ActorSyncState}; +use redis::aio::MultiplexedConnection; +use redis::AsyncTypedCommands; +use reqwest::header::HeaderMap; +use reqwest::Client as HttpClient; +use std::path::{Path, PathBuf}; +use std::sync::Arc; +use tokio::sync::watch::Receiver as WatchReceiver; +use tokio::time::{Duration, Instant}; +use tokio_postgres::types::Type; +use tokio_util::io::StreamReader; +use tokio_util::task::TaskTracker; +use tracing::instrument; + +const BF_RESET_KEY: &str = "bf_download_ratelimit_reset"; +const BF_REM_KEY: &str = "bf_download_ratelimit_rem"; +const DL_DUP_KEY: &str = "bf_downloaded"; + +pub async fn downloader( + mut rc: MultiplexedConnection, + pool: Pool, + resolver: Arc, + tmp_dir: PathBuf, + concurrency: usize, + buffer: usize, + tracker: TaskTracker, + stop: WatchReceiver, +) { + let (tx, rx) = flume::bounded(64); + let mut conn = pool.get().await.unwrap(); + + let http = HttpClient::new(); + + for _ in 0..concurrency { + tracker.spawn(download_thread( + rc.clone(), + pool.clone(), + resolver.clone(), + http.clone(), + rx.clone(), + tmp_dir.clone(), + )); + } + + let status_stmt = conn.prepare_typed_cached( + "INSERT INTO actors (did, sync_state, last_indexed) VALUES ($1, 'processing', NOW()) ON CONFLICT (did) DO UPDATE SET sync_state = 'processing', last_indexed=NOW()", + &[Type::TEXT] + ).await.unwrap(); + + loop { + if stop.has_changed().unwrap_or(true) { + tracing::info!("stopping downloader"); + break; + } + + if let Ok(count) = rc.llen(DL_DONE_KEY).await { + if count > buffer { + tracing::info!("waiting due to full buffer"); + tokio::time::sleep(Duration::from_secs(5)).await; + continue; + } + } + + let did: String = match rc.lpop("backfill_queue", None).await { + Ok(Some(did)) => did, + Ok(None) => { + tokio::time::sleep(Duration::from_millis(250)).await; + continue; + } + Err(e) => { + tracing::error!("failed to get item from backfill queue: {e}"); + continue; + } + }; + + tracing::trace!("resolving repo {did}"); + + // has the repo already been downloaded? + if rc.sismember(DL_DUP_KEY, &did).await.unwrap_or_default() { + tracing::warn!("skipping duplicate repo {did}"); + continue; + } + + // check if they're already synced in DB too + match db::actor_get_statuses(&mut conn, &did).await { + Ok(Some((_, state))) => { + if state == ActorSyncState::Synced || state == ActorSyncState::Processing { + tracing::warn!("skipping duplicate repo {did}"); + continue; + } + } + Ok(None) => {} + Err(e) => { + tracing::error!(did, "failed to check current repo status: {e}"); + db::backfill_job_write(&mut conn, &did, "failed.resolve") + .await + .unwrap(); + } + } + + match resolver.resolve_did(&did).await { + Ok(Some(did_doc)) => { + let Some(service) = did_doc.find_service_by_id(PDS_SERVICE_ID) else { + tracing::warn!("bad DID doc for {did}"); + db::backfill_job_write(&mut conn, &did, "failed.resolve") + .await + .unwrap(); + continue; + }; + let service = service.service_endpoint.clone(); + + // set the repo to processing + if let Err(e) = conn.execute(&status_stmt, &[&did]).await { + tracing::error!("failed to update repo status for {did}: {e}"); + continue; + } + + let handle = did_doc + .also_known_as + .and_then(|akas| akas.first().map(|v| v[5..].to_owned())); + + tracing::trace!("resolved repo {did} {service}"); + if let Err(e) = tx.send_async((service, did, handle)).await { + tracing::error!("failed to send: {e}"); + } + } + Ok(None) => { + tracing::warn!(did, "bad DID doc"); + db::backfill_job_write(&mut conn, &did, "failed.resolve") + .await + .unwrap(); + } + Err(e) => { + tracing::error!(did, "failed to resolve DID doc: {e}"); + db::backfill_job_write(&mut conn, &did, "failed.resolve") + .await + .unwrap(); + } + } + } +} + +async fn download_thread( + mut rc: MultiplexedConnection, + pool: Pool, + resolver: Arc, + http: reqwest::Client, + rx: flume::Receiver<(String, String, Option)>, + tmp_dir: PathBuf, +) { + tracing::debug!("spawning thread"); + + // this will return Err(_) and exit when all senders (only held above) are dropped + while let Ok((pds, did, maybe_handle)) = rx.recv_async().await { + if let Err(e) = enforce_ratelimit(&mut rc, &pds).await { + tracing::error!("ratelimiter error: {e}"); + continue; + }; + + { + tracing::trace!("getting DB conn..."); + let mut conn = pool.get().await.unwrap(); + tracing::trace!("got DB conn..."); + match check_and_update_repo_status(&http, &mut conn, &pds, &did).await { + Ok(true) => {} + Ok(false) => continue, + Err(e) => { + tracing::error!(pds, did, "failed to check repo status: {e}"); + db::backfill_job_write(&mut conn, &did, "failed.resolve") + .await + .unwrap(); + continue; + } + } + + tracing::debug!("trying to resolve handle..."); + if let Some(handle) = maybe_handle { + if let Err(e) = resolve_and_set_handle(&conn, &resolver, &did, &handle).await { + tracing::error!(pds, did, "failed to resolve handle: {e}"); + db::backfill_job_write(&mut conn, &did, "failed.resolve") + .await + .unwrap(); + } + } + } + + let start = Instant::now(); + + tracing::trace!("downloading repo {did}"); + + match download_car(&http, &tmp_dir, &pds, &did).await { + Ok(Some((rem, reset))) => { + let _ = rc.zadd(BF_REM_KEY, &pds, rem).await; + let _ = rc.zadd(BF_RESET_KEY, &pds, reset).await; + } + Ok(_) => tracing::warn!(pds, "got response with no ratelimit headers."), + Err(e) => { + tracing::error!(pds, did, "failed to download repo: {e}"); + continue; + } + } + + histogram!("backfill_download_dur", "pds" => pds).record(start.elapsed().as_secs_f64()); + + let _ = rc.sadd(DL_DUP_KEY, &did).await; + if let Err(e) = rc.rpush(DL_DONE_KEY, &did).await { + tracing::error!(did, "failed to mark download complete: {e}"); + } else { + counter!("backfill_downloaded").increment(1); + } + } + + tracing::debug!("thread exiting"); +} + +async fn enforce_ratelimit(rc: &mut MultiplexedConnection, pds: &str) -> eyre::Result<()> { + let score = rc.zscore(BF_REM_KEY, pds).await?; + + if let Some(rem) = score { + if (rem as i32) < 100 { + // if we've got None for some reason, just hope that the next req will contain the reset header. + if let Some(at) = rc.zscore(BF_RESET_KEY, pds).await? { + tracing::debug!("rate limit for {pds} resets at {at}"); + let time = chrono::DateTime::from_timestamp(at as i64, 0).unwrap(); + let delta = (time - Utc::now()).num_milliseconds().max(0); + + tokio::time::sleep(Duration::from_millis(delta as u64)).await; + }; + } + } + + Ok(()) +} + +// you wouldn't... +#[instrument(skip(http, tmp_dir, pds))] +async fn download_car( + http: &HttpClient, + tmp_dir: &Path, + pds: &str, + did: &str, +) -> eyre::Result> { + let mut file = tokio::fs::File::create_new(tmp_dir.join(did)).await?; + + let res = http + .get(format!("{pds}/xrpc/com.atproto.sync.getRepo?did={did}")) + .send() + .await? + .error_for_status()?; + + let headers = res.headers(); + let ratelimit_rem = header_to_int(headers, "ratelimit-remaining"); + let ratelimit_reset = header_to_int(headers, "ratelimit-reset"); + + let strm = res.bytes_stream().map_err(std::io::Error::other); + let mut reader = StreamReader::new(strm); + + tokio::io::copy(&mut reader, &mut file).await?; + + Ok(ratelimit_rem.zip(ratelimit_reset)) +} + +// there's no ratelimit handling here because we pretty much always call download_car after. +#[instrument(skip(http, conn, pds))] +async fn check_and_update_repo_status( + http: &HttpClient, + conn: &mut PgClient, + pds: &str, + repo: &str, +) -> eyre::Result { + match super::check_pds_repo_status(http, pds, repo).await? { + Some(status) => { + if !status.active { + tracing::debug!("repo is inactive"); + + let status = status + .status + .unwrap_or(crate::firehose::AtpAccountStatus::Deleted); + conn.execute( + "UPDATE actors SET sync_state='dirty', status=$2 WHERE did=$1", + &[&repo, &ActorStatus::from(status)], + ) + .await?; + + Ok(false) + } else { + Ok(true) + } + } + None => { + // this repo can't be found - set dirty and assume deleted. + tracing::debug!("repo was deleted"); + conn.execute( + "UPDATE actors SET sync_state='dirty', status='deleted' WHERE did=$1", + &[&repo], + ) + .await?; + + Ok(false) + } + } +} + +async fn resolve_and_set_handle( + conn: &PgClient, + resolver: &Resolver, + did: &str, + handle: &str, +) -> eyre::Result<()> { + if let Some(handle_did) = resolver.resolve_handle(handle).await? { + if handle_did == did { + conn.execute("UPDATE actors SET handle=$2 WHERE did=$1", &[&did, &handle]) + .await?; + } else { + tracing::warn!("requested DID ({did}) doesn't match handle"); + } + } + + Ok(()) +} + +fn header_to_int(headers: &HeaderMap, name: &str) -> Option { + headers + .get(name) + .and_then(|v| v.to_str().ok()) + .and_then(|v| v.parse().ok()) +} diff --git a/consumer/src/backfill/mod.rs b/consumer/src/backfill/mod.rs index 53869521..b52a33c8 100644 --- a/consumer/src/backfill/mod.rs +++ b/consumer/src/backfill/mod.rs @@ -9,8 +9,9 @@ use ipld_core::cid::Cid; use metrics::counter; use parakeet_db::types::{ActorStatus, ActorSyncState}; use redis::aio::MultiplexedConnection; -use redis::{AsyncCommands, Direction}; +use redis::AsyncTypedCommands; use reqwest::{Client, StatusCode}; +use std::path::PathBuf; use std::str::FromStr; use std::sync::Arc; use tokio::sync::watch::Receiver as WatchReceiver; @@ -18,9 +19,11 @@ use tokio::sync::Semaphore; use tokio_util::task::TaskTracker; use tracing::instrument; +mod downloader; mod repo; mod types; +const DL_DONE_KEY: &str = "bf_download_complete"; const PDS_SERVICE_ID: &str = "#atproto_pds"; // There's a 4MiB limit on parakeet-index, so break delta batches up if there's loads. // this should be plenty low enough to not trigger the size limit. (59k did slightly) @@ -28,16 +31,16 @@ const DELTA_BATCH_SIZE: usize = 32 * 1024; #[derive(Clone)] pub struct BackfillManagerInner { - resolver: Arc, - client: Client, index_client: Option, - opts: BackfillConfig, + tmp_dir: PathBuf, } pub struct BackfillManager { pool: Pool, redis: MultiplexedConnection, + resolver: Arc, semaphore: Arc, + opts: BackfillConfig, inner: BackfillManagerInner, } @@ -49,41 +52,43 @@ impl BackfillManager { index_client: Option, opts: BackfillConfig, ) -> eyre::Result { - let client = Client::builder().brotli(true).build()?; let semaphore = Arc::new(Semaphore::new(opts.backfill_workers as usize)); Ok(BackfillManager { pool, redis, + resolver, semaphore, inner: BackfillManagerInner { - resolver, - client, index_client, - opts, + tmp_dir: PathBuf::from_str(&opts.download_tmp_dir)?, }, + opts, }) } pub async fn run(mut self, stop: WatchReceiver) -> eyre::Result<()> { let tracker = TaskTracker::new(); + tracker.spawn(downloader::downloader( + self.redis.clone(), + self.pool.clone(), + self.resolver, + self.inner.tmp_dir.clone(), + self.opts.download_workers, + self.opts.download_buffer, + tracker.clone(), + stop.clone(), + )); + loop { if stop.has_changed().unwrap_or(true) { tracker.close(); + tracing::info!("stopping backfiller"); break; } - let Some(job) = self - .redis - .lmove::<_, _, Option>( - "backfill_queue", - "backfill_processing", - Direction::Left, - Direction::Right, - ) - .await? - else { + let Some(job): Option = self.redis.lpop(DL_DONE_KEY, None).await? else { tokio::time::sleep(tokio::time::Duration::from_millis(100)).await; continue; }; @@ -92,7 +97,6 @@ impl BackfillManager { let mut inner = self.inner.clone(); let mut conn = self.pool.get().await?; - let mut redis = self.redis.clone(); tracker.spawn(async move { let _p = p; @@ -102,7 +106,7 @@ impl BackfillManager { tracing::error!(did = &job, "backfill failed: {e}"); counter!("backfill_failure").increment(1); - db::backfill_job_write(&mut conn, &job, "failed") + db::backfill_job_write(&mut conn, &job, "failed.write") .await .unwrap(); } else { @@ -113,10 +117,9 @@ impl BackfillManager { .unwrap(); } - redis - .lrem::<_, _, i32>("backfill_processing", 1, &job) - .await - .unwrap(); + if let Err(e) = tokio::fs::remove_file(inner.tmp_dir.join(&job)).await { + tracing::error!(did = &job, "failed to remove file: {e}"); + } }); } @@ -132,91 +135,12 @@ async fn backfill_actor( inner: &mut BackfillManagerInner, did: &str, ) -> eyre::Result<()> { - let Some((status, sync_state)) = db::actor_get_statuses(conn, did).await? else { - tracing::error!("skipping backfill on unknown repo"); - return Ok(()); - }; - - if sync_state != ActorSyncState::Dirty || status != ActorStatus::Active { - tracing::debug!("skipping non-dirty or inactive repo"); - return Ok(()); - } - - // resolve the did to a PDS (also validates the handle) - let Some(did_doc) = inner.resolver.resolve_did(did).await? else { - eyre::bail!("missing did doc"); - }; - - let Some(service) = did_doc.find_service_by_id(PDS_SERVICE_ID) else { - eyre::bail!("DID doc contained no service endpoint"); - }; - - let pds_url = service.service_endpoint.clone(); - - // check the repo status before we attempt to resolve the handle. There's a case where we can't - // resolve the handle in the DID doc because the acc is already deleted. - let Some(repo_status) = check_pds_repo_status(&inner.client, &pds_url, did).await? else { - // this repo can't be found - set dirty and assume deleted. - tracing::debug!("repo was deleted"); - db::actor_upsert( - conn, - did, - ActorStatus::Deleted, - ActorSyncState::Dirty, - Utc::now(), - ) - .await?; - return Ok(()); - }; - - if !repo_status.active { - tracing::debug!("repo is inactive"); - let status = repo_status - .status - .unwrap_or(crate::firehose::AtpAccountStatus::Deleted); - db::actor_upsert(conn, did, status.into(), ActorSyncState::Dirty, Utc::now()).await?; - return Ok(()); - } - - if !inner.opts.skip_handle_validation { - // at this point, the account will be active and we can attempt to resolve the handle. - let Some(handle) = did_doc - .also_known_as - .and_then(|aka| aka.first().cloned()) - .and_then(|handle| handle.strip_prefix("at://").map(String::from)) - else { - eyre::bail!("DID doc contained no handle"); - }; - - // in theory, we can use com.atproto.identity.resolveHandle against a PDS, but that seems - // like a way to end up with really sus handles. - if let Some(handle_did) = inner.resolver.resolve_handle(&handle).await? { - if handle_did != did { - tracing::warn!("requested DID doesn't match handle"); - } else { - // set the handle from above - db::actor_upsert_handle( - conn, - did, - ActorSyncState::Processing, - Some(handle), - Utc::now(), - ) - .await?; - } - } - } - - // now we can start actually backfilling - db::actor_set_sync_status(conn, did, ActorSyncState::Processing, Utc::now()).await?; - let mut t = conn.transaction().await?; t.execute("SET CONSTRAINTS ALL DEFERRED", &[]).await?; - tracing::trace!("pulling repo"); + tracing::trace!("loading repo"); - let (commit, mut deltas, copies) = - repo::stream_and_insert_repo(&mut t, &inner.client, did, &pds_url).await?; + let (commit, mut deltas, copies) = repo::insert_repo(&mut t, &inner.tmp_dir, did).await?; db::actor_set_repo_state(&mut t, did, &commit.rev, commit.data).await?; diff --git a/consumer/src/backfill/repo.rs b/consumer/src/backfill/repo.rs index d0e92069..60557dd2 100644 --- a/consumer/src/backfill/repo.rs +++ b/consumer/src/backfill/repo.rs @@ -6,36 +6,23 @@ use crate::indexer::records; use crate::indexer::types::{AggregateDeltaStore, RecordTypes}; use crate::{db, indexer}; use deadpool_postgres::Transaction; -use futures::TryStreamExt; use ipld_core::cid::Cid; use iroh_car::CarReader; use metrics::counter; use parakeet_index::AggregateType; -use reqwest::Client; use std::collections::HashMap; -use std::io::ErrorKind; +use std::path::Path; use tokio::io::BufReader; -use tokio_util::io::StreamReader; type BackfillDeltaStore = HashMap<(String, i32), i32>; -pub async fn stream_and_insert_repo( +pub async fn insert_repo( t: &mut Transaction<'_>, - client: &Client, + tmp_dir: &Path, repo: &str, - pds: &str, ) -> eyre::Result<(CarCommitEntry, BackfillDeltaStore, CopyStore)> { - let res = client - .get(format!("{pds}/xrpc/com.atproto.sync.getRepo?did={repo}")) - .send() - .await? - .error_for_status()?; - - let strm = res - .bytes_stream() - .map_err(|err| std::io::Error::new(ErrorKind::Other, err)); - let reader = StreamReader::new(strm); - let mut car_stream = CarReader::new(BufReader::new(reader)).await?; + let car = tokio::fs::File::open(tmp_dir.join(repo)).await?; + let mut car_stream = CarReader::new(BufReader::new(car)).await?; // the root should be the commit block let root = car_stream.header().roots().first().cloned().unwrap(); diff --git a/consumer/src/config.rs b/consumer/src/config.rs index 5c91d797..8a6a9c6d 100644 --- a/consumer/src/config.rs +++ b/consumer/src/config.rs @@ -40,6 +40,9 @@ pub struct IndexerConfig { /// You can use this to move handle resolution out of event handling and into another place. #[serde(default)] pub skip_handle_validation: bool, + /// Whether to submit backfill requests for new repos. (Only when history_mode == BackfillHistory). + #[serde(default)] + pub request_backfill: bool, } #[derive(Copy, Clone, Debug, PartialEq, PartialOrd, Deserialize)] @@ -57,8 +60,11 @@ pub struct BackfillConfig { pub backfill_workers: u8, #[serde(default)] pub skip_aggregation: bool, - #[serde(default)] - pub skip_handle_validation: bool, + #[serde(default = "default_download_workers")] + pub download_workers: usize, + #[serde(default = "default_download_buffer")] + pub download_buffer: usize, + pub download_tmp_dir: String, } fn default_backfill_workers() -> u8 { @@ -68,3 +74,11 @@ fn default_backfill_workers() -> u8 { fn default_indexer_workers() -> u8 { 4 } + +fn default_download_workers() -> usize { + 25 +} + +fn default_download_buffer() -> usize { + 25_000 +} diff --git a/consumer/src/indexer/mod.rs b/consumer/src/indexer/mod.rs index ca3b16c9..d20aa2aa 100644 --- a/consumer/src/indexer/mod.rs +++ b/consumer/src/indexer/mod.rs @@ -30,6 +30,7 @@ pub mod types; pub struct RelayIndexerOpts { pub history_mode: HistoryMode, pub skip_handle_validation: bool, + pub request_backfill: bool, } #[derive(Clone)] @@ -38,6 +39,7 @@ struct RelayIndexerState { resolver: Arc, do_backfill: bool, do_handle_res: bool, + req_backfill: bool, } pub struct RelayIndexer { @@ -66,6 +68,7 @@ impl RelayIndexer { state: RelayIndexerState { resolver, do_backfill: opts.history_mode == HistoryMode::BackfillHistory, + req_backfill: opts.request_backfill, do_handle_res: !opts.skip_handle_validation, idxc_tx, }, @@ -275,7 +278,7 @@ async fn index_account( .map(ActorStatus::from) .unwrap_or(ActorStatus::Active); - let trigger_bf = if state.do_backfill && status == ActorStatus::Active { + let trigger_bf = if state.do_backfill && state.req_backfill && status == ActorStatus::Active { // check old status - if they exist (Some(*)), AND were previously != Active but not Deleted, // AND have a rev == null, then trigger backfill. db::actor_get_status_and_rev(conn, &account.did) @@ -325,7 +328,7 @@ async fn index_commit( // backfill for them and they can be marked active and indexed normally. // TODO: bridgy doesn't implement since atm - we need a special case if commit.since.is_some() { - if state.do_backfill { + if state.do_backfill && state.req_backfill { rc.rpush::<_, _, i32>("backfill_queue", commit.repo).await?; } return Ok(()); @@ -356,7 +359,9 @@ async fn index_commit( .await?; if trigger_backfill { - rc.rpush::<_, _, i32>("backfill_queue", commit.repo).await?; + if state.req_backfill { + rc.rpush::<_, _, i32>("backfill_queue", commit.repo).await?; + } return Ok(()); } diff --git a/consumer/src/main.rs b/consumer/src/main.rs index f574f2cc..c1482b34 100644 --- a/consumer/src/main.rs +++ b/consumer/src/main.rs @@ -115,6 +115,7 @@ async fn main() -> eyre::Result<()> { let indexer_opts = indexer::RelayIndexerOpts { history_mode: indexer_cfg.history_mode, skip_handle_validation: indexer_cfg.skip_handle_validation, + request_backfill: indexer_cfg.request_backfill, }; let relay_indexer = indexer::RelayIndexer::new( -- 2.51.2