diff --git a/src/backfill/worker/task.rs b/src/backfill/worker/task.rs index 7db0926..277139f 100644 --- a/src/backfill/worker/task.rs +++ b/src/backfill/worker/task.rs @@ -38,35 +38,46 @@ pub(crate) async fn did_task( .await { Ok(Some(_repo_state)) => { - let applied = state.db.run({ - let did = did.clone(); - let pending_key = pending_key.clone(); - move |db| { - let did_key = keys::repo_key(&did); - let mut txn = DbTxn::new(db); - let applied = - txn.transition_pending_key(&did, pending_key.as_ref(), GaugeState::Synced)?; - txn.batch.remove(&db.indexer.pending, pending_key.clone()); - if applied { - txn.batch.remove(&db.indexer.resync, &did_key); + let applied = state + .db + .run({ + let did = did.clone(); + let pending_key = pending_key.clone(); + move |db| { + let did_key = keys::repo_key(&did); + let mut txn = DbTxn::new(db); + let applied = txn.transition_pending_key( + &did, + pending_key.as_ref(), + GaugeState::Synced, + )?; + txn.batch.remove(&db.indexer.pending, pending_key.clone()); + if applied { + txn.batch.remove(&db.indexer.resync, &did_key); + } + txn.commit()?; + Ok::<_, miette::Report>(applied) } - txn.commit()?; - Ok::<_, miette::Report>(applied) - } - }) - .await?; + }) + .await?; if !applied { return Ok(()); } + state + .backfill_drained + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); + let state = state.clone(); - state.db.run(move |db| { - db.inner - .persist(fjall::PersistMode::Buffer) - .into_diagnostic() - }) - .await?; + state + .db + .run(move |db| { + db.inner + .persist(fjall::PersistMode::Buffer) + .into_diagnostic() + }) + .await?; if let Err(e) = buffer_tx .send(IndexerMessage::BackfillFinished(did.clone())) @@ -79,18 +90,23 @@ pub(crate) async fn did_task( Ok(None) => Ok(()), Err(BackfillError::Deleted) => { warn!("orphaned pending entry, cleaning up"); - state.db.run({ - let did = did.clone(); - let pending_key = pending_key.clone(); - move |db| { - let mut txn = DbTxn::new(db); - txn.transition_pending_key(&did, pending_key.as_ref(), GaugeState::Synced)?; - txn.batch.remove(&db.indexer.pending, pending_key); - txn.commit()?; - Ok::<_, miette::Report>(()) - } - }) - .await?; + state + .db + .run({ + let did = did.clone(); + let pending_key = pending_key.clone(); + move |db| { + let mut txn = DbTxn::new(db); + txn.transition_pending_key(&did, pending_key.as_ref(), GaugeState::Synced)?; + txn.batch.remove(&db.indexer.pending, pending_key); + txn.commit()?; + Ok::<_, miette::Report>(()) + } + }) + .await?; + state + .backfill_drained + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); Ok(()) } Err(e) => { @@ -154,46 +170,53 @@ pub(crate) async fn did_task( }; let error_string = e.to_string(); - state.db.run({ - let did_key = did_key.into_static(); - let did = did.clone(); - let pending_key = pending_key.clone(); - move |db| { - // 3. save to resync - let serialized_resync_state = - rmp_serde::to_vec(&resync_state).into_diagnostic()?; - - // 4. and update the main repo state - let serialized_repo_state = if let Some(state_bytes) = - db.repos.get(&did_key).into_diagnostic()? - { - let mut state: RepoState = - rmp_serde::from_slice(&state_bytes).into_diagnostic()?; - state.active = true; - state.status = RepoStatus::Error(error_string.into()); - Some(rmp_serde::to_vec(&state).into_diagnostic()?) - } else { - None - }; - let mut txn = DbTxn::new(db); - let applied = txn.transition_pending_key( - &did, - pending_key.as_ref(), - GaugeState::Resync(Some(error_kind)), - )?; - txn.batch.remove(&db.indexer.pending, pending_key.clone()); - if applied { - txn.batch - .insert(&db.indexer.resync, &did_key, serialized_resync_state); - if let Some(state_bytes) = serialized_repo_state { - txn.batch.insert(&db.repos, &did_key, state_bytes); + let applied = state + .db + .run({ + let did_key = did_key.into_static(); + let did = did.clone(); + let pending_key = pending_key.clone(); + move |db| { + // 3. save to resync + let serialized_resync_state = + rmp_serde::to_vec(&resync_state).into_diagnostic()?; + + // 4. and update the main repo state + let serialized_repo_state = + if let Some(state_bytes) = db.repos.get(&did_key).into_diagnostic()? { + let mut state: RepoState = + rmp_serde::from_slice(&state_bytes).into_diagnostic()?; + state.active = true; + state.status = RepoStatus::Error(error_string.into()); + Some(rmp_serde::to_vec(&state).into_diagnostic()?) + } else { + None + }; + let mut txn = DbTxn::new(db); + let applied = txn.transition_pending_key( + &did, + pending_key.as_ref(), + GaugeState::Resync(Some(error_kind)), + )?; + txn.batch.remove(&db.indexer.pending, pending_key.clone()); + if applied { + txn.batch + .insert(&db.indexer.resync, &did_key, serialized_resync_state); + if let Some(state_bytes) = serialized_repo_state { + txn.batch.insert(&db.repos, &did_key, state_bytes); + } } + txn.commit()?; + Ok::<_, miette::Report>(applied) } - txn.commit()?; - Ok::<_, miette::Report>(()) - } - }) - .await?; + }) + .await?; + + if applied { + state + .backfill_drained + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); + } Err(e) } diff --git a/src/control/hydrant/run.rs b/src/control/hydrant/run.rs index c3b6f9e..4fd9f64 100644 --- a/src/control/hydrant/run.rs +++ b/src/control/hydrant/run.rs @@ -214,6 +214,39 @@ impl Hydrant { } }); + // 9b. backfill stats ticker + #[cfg(feature = "indexer")] + tokio::spawn({ + let state = state.clone(); + let mut last_time = std::time::Instant::now(); + let mut interval = tokio::time::interval(std::time::Duration::from_secs(60)); + async move { + loop { + interval.tick().await; + + let pending_len = state.db.indexer.approximate_pending_count(); + let drained = state.backfill_drained.swap(0, Ordering::Relaxed); + let current_time = std::time::Instant::now(); + let elapsed = current_time.duration_since(last_time).as_secs_f64(); + last_time = current_time; + + if drained == 0 && pending_len == 0 { + continue; + } + + let rate = if elapsed > 0.0 { + drained as f64 / elapsed + } else { + 0.0 + }; + + info!( + "backfill: {rate:.2} repos/s drained ({drained} repos in {elapsed:.1}s, {pending_len} pending)" + ); + } + } + }); + let (fatal_tx_inner, mut fatal_rx) = watch::channel(None); let fatal_tx = Arc::new(fatal_tx_inner); diff --git a/src/db/open.rs b/src/db/open.rs index 24a54ed..db62528 100644 --- a/src/db/open.rs +++ b/src/db/open.rs @@ -18,6 +18,8 @@ use super::{Db, migration, registry, schema}; impl Db { pub fn open(cfg: &Config) -> Result { + tracing::info!("opening database at {}...", cfg.database_path.display()); + let start = std::time::Instant::now(); let (db, count_delta_gc_watermark) = Self::open_database(cfg)?; let dicts = Self::load_dicts(cfg); @@ -49,6 +51,8 @@ impl Db { this.restore_count_deltas_and_next_id()?; + tracing::info!("database ready (recovered in {:?})", start.elapsed()); + Ok(this) } diff --git a/src/ingest/indexer/shard.rs b/src/ingest/indexer/shard.rs index 457335e..555f7ec 100644 --- a/src/ingest/indexer/shard.rs +++ b/src/ingest/indexer/shard.rs @@ -313,7 +313,7 @@ impl FirehoseWorker { }; if chain_break { - warn!("chain break detected, triggering backfill"); + debug!("chain break detected, triggering backfill"); Self::trigger_backfill(ctx, did, repo_state)?; return Ok(RepoProcessResult::NeedsBackfill(Some(commit))); } diff --git a/src/state.rs b/src/state.rs index 8f4a0ff..7e0a627 100644 --- a/src/state.rs +++ b/src/state.rs @@ -40,6 +40,8 @@ pub struct AppState { pub crawler_enabled: watch::Sender, #[cfg(feature = "indexer")] pub backfill_enabled: watch::Sender, + #[cfg(feature = "indexer")] + pub backfill_drained: std::sync::atomic::AtomicUsize, pub ephemeral: bool, pub ephemeral_ttl: Duration, pub only_index_links: bool, @@ -121,6 +123,8 @@ impl AppState { firehose_enabled, #[cfg(feature = "indexer")] backfill_enabled, + #[cfg(feature = "indexer")] + backfill_drained: std::sync::atomic::AtomicUsize::new(0), ephemeral: config.ephemeral, ephemeral_ttl: config.ephemeral_ttl, only_index_links: config.only_index_links,