diff --git a/docs/api/repos.md b/docs/api/repos.md index 2f8379e..e570dec 100644 --- a/docs/api/repos.md +++ b/docs/api/repos.md @@ -22,7 +22,7 @@ both endpoints say whether hydrant still has work queued for the repository: `pe | `state` | fields | meaning | | :--- | :--- | :--- | -| `error` | `kind`, `retry_count`, `next_retry`, `host_failures` | the last backfill failed and is retried at `next_retry`. `kind` is `ratelimited` (the PDS asked to slow down, or hydrant skipped it because its PDS is throttled), `transport` (the PDS couldn't be reached) or `generic` (anything else). `host_failures` is how many transport failures in a row the repository's PDS has had since hydrant started, so a `ratelimited` repo on a PDS that keeps failing is waiting on a dead host rather than a rate limit | +| `error` | `kind`, `retry_count`, `next_retry`, `host_failures` | the last backfill failed and is retried at `next_retry`. `retry_count` is how many backfills failed in a row, across retries, until one succeeds, and the wait doubles with it from about 2 minutes up to an hour (a skip because the PDS is throttled waits a minute and doesn't count). `kind` is `ratelimited` (the PDS asked to slow down, or hydrant skipped it because its PDS is throttled), `transport` (the PDS couldn't be reached) or `generic` (anything else). `host_failures` is how many transport failures in a row the repository's PDS has had since hydrant started, so a `ratelimited` repo on a PDS that keeps failing is waiting on a dead host rather than a rate limit | | `gone` | `status` | the account is inactive (deactivated, takendown, ...) and gets backfilled again once it comes back | ## PUT /repos diff --git a/src/backfill/manager.rs b/src/backfill/manager.rs index 4ff28ba..171d3df 100644 --- a/src/backfill/manager.rs +++ b/src/backfill/manager.rs @@ -9,7 +9,7 @@ use std::sync::Arc; use std::time::Duration; use tracing::{debug, error, info}; -fn queue_resync_if( +pub(crate) fn queue_resync_if( state: &AppState, did: &Did, eligible: impl FnOnce(&ResyncState) -> bool, @@ -38,6 +38,10 @@ fn queue_resync_if( let mut metadata = crate::db::deser_repo_meta(&metadata_bytes)?; let old_pending = keys::pending_key(metadata.index_id); metadata.index_id = rand::random(); + // a database from before v14 keeps the count only on the entry + if let ResyncState::Error { retry_count, .. } = resync_state { + metadata.retry_count = metadata.retry_count.max(retry_count); + } txn.batch.remove(&state.db.indexer.resync, &did_key); let indexer = &state.db.indexer; @@ -247,6 +251,51 @@ mod tests { state.db.indexer.assert_pending_dids_match(); } + #[test] + fn requeue_keeps_the_retry_count_an_older_entry_holds() { + let (_tmp, state) = state(); + let did = did(); + let mut batch = state.db.inner.batch(); + batch.insert( + &state.db.repos, + keys::repo_key(&did), + crate::db::ser_repo_state(&RepoState::synced()).unwrap(), + ); + batch.insert( + &state.db.repo_metadata, + keys::repo_metadata_key(&did), + rmp_serde::to_vec(&crate::types::v4::RepoMetadata { + tracked: true, + index_id: 7, + }) + .unwrap(), + ); + batch.insert( + &state.db.indexer.resync, + keys::repo_key(&did), + rmp_serde::to_vec(&ResyncState::Error { + kind: crate::types::ResyncErrorKind::Generic, + retry_count: 5, + next_retry: 0, + }) + .unwrap(), + ); + batch.commit().unwrap(); + + assert!(queue_resync_if(&state, &did, |_| true).unwrap()); + let metadata = state + .db + .repo_metadata + .get(keys::repo_metadata_key(&did)) + .unwrap() + .unwrap(); + assert_eq!( + crate::db::deser_repo_meta(&metadata).unwrap().retry_count, + 5 + ); + state.db.indexer.assert_pending_dids_match(); + } + #[test] fn eligible_repo_moves_from_resync_to_pending() { let (_tmp, state) = state(); diff --git a/src/backfill/sparse.rs b/src/backfill/sparse.rs index 1a47b9f..904d7a4 100644 --- a/src/backfill/sparse.rs +++ b/src/backfill/sparse.rs @@ -675,6 +675,7 @@ async fn persist_sparse_backfill( .ok_or_else(|| miette::miette!("repo metadata not found for {}", did))?; let mut metadata = crate::db::deser_repo_meta(&metadata_bytes)?; metadata.tracked = true; + metadata.retry_count = 0; txn.batch.insert( &db.repo_metadata, &metadata_key, diff --git a/src/backfill/worker/process.rs b/src/backfill/worker/process.rs index f84412e..7e82f3b 100644 --- a/src/backfill/worker/process.rs +++ b/src/backfill/worker/process.rs @@ -494,6 +494,7 @@ pub(super) async fn process_did( .ok_or_else(|| miette::miette!("repo metadata not found for {}", did))?; let mut metadata = crate::db::deser_repo_meta(&metadata_bytes)?; metadata.tracked = true; + metadata.retry_count = 0; txn.batch.insert( &app_state.db.repo_metadata, &metadata_key, diff --git a/src/backfill/worker/task.rs b/src/backfill/worker/task.rs index 62a9ebf..a807cf7 100644 --- a/src/backfill/worker/task.rs +++ b/src/backfill/worker/task.rs @@ -284,7 +284,7 @@ pub(super) async fn did_task( rmp_serde::from_slice::(&bytes).into_diagnostic() }) .transpose()?; - let mut retry_count = match existing { + let entry_retries = match existing { Some(ResyncState::Error { retry_count, .. }) => retry_count, Some(ResyncState::Gone { .. }) => { db.indexer @@ -300,10 +300,21 @@ pub(super) async fn did_task( return Ok(Recorded::Retrying { applied: false }); } + let metadata_key = keys::repo_metadata_key(&did); + let mut metadata = db + .repo_metadata + .get(&metadata_key) + .into_diagnostic()? + .map(|bytes| crate::db::deser_repo_meta(&bytes)) + .transpose()?; + let mut retry_count = metadata + .as_ref() + .map_or(0, |metadata| metadata.retry_count) + .max(entry_retries); let next_retry = if throttled { chrono::Utc::now().timestamp() + 60 } else { - retry_count += 1; + retry_count = retry_count.saturating_add(1); ResyncState::next_backoff(retry_count) }; let resync_state = ResyncState::Error { @@ -324,6 +335,14 @@ pub(super) async fn did_task( &did_key, rmp_serde::to_vec(&resync_state).into_diagnostic()?, ); + if let Some(metadata) = metadata.as_mut() { + metadata.retry_count = retry_count; + txn.batch.insert( + &db.repo_metadata, + &metadata_key, + crate::db::ser_repo_meta(metadata)?, + ); + } if let Some(bytes) = db.repos.get(&did_key).into_diagnostic()? { let mut repo: RepoState = rmp_serde::from_slice(&bytes).into_diagnostic()?; @@ -1382,6 +1401,71 @@ mod tests { Ok(()) } + #[tokio::test] + async fn backoff_grows_across_requeues_until_a_fetch_succeeds() -> miette::Result<()> { + let fixture = Fixture::new().await?; + fixture.seed(&RepoState::backfilling().into_static(), 9, &[])?; + let metadata = || -> miette::Result { + let bytes = fixture + .state + .db + .repo_metadata + .get(keys::repo_metadata_key(&fixture.did)) + .into_diagnostic()? + .expect("metadata row"); + db::deser_repo_meta(&bytes) + }; + + for (failures, minutes) in [(1, 2), (2, 4), (3, 8), (4, 16), (5, 32), (6, 60), (7, 60)] { + let before = chrono::Utc::now().timestamp(); + let pending_key = keys::pending_key(metadata()?.index_id); + assert!(matches!( + fixture.run(pending_key).await, + Err(BackfillError::RepoNotFound) + )); + + let Some(ResyncState::Error { + retry_count, + next_retry, + .. + }) = fixture.resync()? + else { + panic!("expected a retryable error after failure {failures}"); + }; + assert_eq!(retry_count, failures); + assert_eq!(metadata()?.retry_count, failures); + let delay = next_retry - before; + let expected = minutes * 60; + assert!( + (expected * 9 / 10 - 1..=expected * 11 / 10 + 5).contains(&delay), + "failure {failures} waits {delay}s, expected about {expected}s" + ); + assert_eq!(fixture.state.db.get_count_sync("resync"), 1); + assert_eq!(fixture.state.db.get_count_sync("error_generic"), 1); + + // what the retry worker does once the entry is due + assert!(crate::backfill::manager::queue_resync_if( + &fixture.state, + &fixture.did, + |_| true + )?); + assert!(fixture.resync()?.is_none()); + assert_eq!(fixture.pending_keys()?.len(), 1); + assert_eq!(fixture.state.db.get_count_sync("resync"), 0); + assert_eq!(fixture.state.db.get_count_sync("error_generic"), 0); + } + + let body = post("a"); + let car = build_car("3jzfcijpj2z2b", &[("a", body.as_slice())]).await?; + fixture.pds.serve_car(car.bytes).await; + let pending_key = keys::pending_key(metadata()?.index_id); + assert_eq!(fixture.run(pending_key).await?, TaskDisposition::Finished); + assert_eq!(metadata()?.retry_count, 0); + assert!(fixture.resync()?.is_none()); + assert!(fixture.pending_keys()?.is_empty()); + Ok(()) + } + #[tokio::test] async fn repo_not_found_on_a_never_rooted_repo_uses_the_same_retry_policy() -> miette::Result<()> {