diff --git a/src/ingest/indexer/shard.rs b/src/ingest/indexer/shard.rs index a40e8bf..c2548cf 100644 --- a/src/ingest/indexer/shard.rs +++ b/src/ingest/indexer/shard.rs @@ -451,6 +451,14 @@ impl FirehoseWorker { .get(keys::repo_key(did)) .into_diagnostic()? else { + // whatever erased the row is the last word, so the key it left behind goes now + // instead of on its next run + ctx.txn + .transition_pending_key(did, pending_key, GaugeState::Synced)?; + ctx.state + .db + .indexer + .stage_pending_remove(&mut ctx.txn.batch, did, pending_key); return Ok(false); }; let repo_state = crate::db::deser_repo_state(&bytes)?.into_static(); @@ -777,6 +785,27 @@ mod tests { Ok(()) } + #[tokio::test] + async fn a_backfill_whose_repo_row_is_gone_comes_off_the_queue() -> Result<()> { + let (_tmp, state, did) = queued_repo(&RepoState::backfilling(), 7, 7)?; + state + .db + .repos + .remove(keys::repo_key(&did)) + .into_diagnostic()?; + + assert!(!finish_on_a_shard(&state, &did, keys::pending_key(7)).await?); + assert!( + !state + .db + .indexer + .is_pending(&keys::pending_key(7)) + .into_diagnostic()? + ); + state.db.indexer.assert_pending_dids_match(); + Ok(()) + } + #[test] fn resync_buffer_commit_msgpack_round_trip_preserves_path() { let cid = CidLink::from(jacquard_repo::mst::util::compute_cid(b"commit").unwrap());