From 2b4840c6b4f7becb2aa045cbec3adf96ed3549a6 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Mon, 5 Oct 2026 22:19:43 +0300 Subject: [PATCH] [backfill] keep a repo queued until the shard has drained its buffer the task took the repo off the queue as soon as its records were written and sent BackfillFinished without waiting. a crash before the shard drained, a failed send or a drain error left the commits buffered during the fetch in resync_buffer with nothing queued to apply them, and a live commit landing in that gap was applied directly and then rolled back by the drain. the queue entry now stays until the shard drains and drops it in the same batch, and the task waits for that answer. a failed drain goes through the usual retry, and a shutdown leaves the repo queued for the next start. the shard only drains while the task's key is still the one the repo is queued under, and it also drops that key when metadata names another. the drain applies buffered commits through apply_checked_commit, since handle_commit would buffer them again while the queue entry is there, and RepoProcessResult drops the state nothing reads anymore. --- src/api/xrpc/resolve_mini_doc.rs | 26 ++- src/backfill/worker/task.rs | 273 ++++++++++++++++++------ src/ingest/indexer/message.rs | 14 +- src/ingest/indexer/shard.rs | 342 ++++++++++++++++++------------- src/ingest/indexer/worker.rs | 23 ++- 5 files changed, 450 insertions(+), 228 deletions(-) diff --git a/src/api/xrpc/resolve_mini_doc.rs b/src/api/xrpc/resolve_mini_doc.rs index 8e7893e..2713069 100644 --- a/src/api/xrpc/resolve_mini_doc.rs +++ b/src/api/xrpc/resolve_mini_doc.rs @@ -512,15 +512,19 @@ mod tests { } async fn finish_backfill(fixture: &Fixture) -> miette::Result<()> { - let state = fixture.hydrant.state.clone(); - let (tx, mut worker) = FirehoseWorker::new(state.clone(), 1); - let rx = worker.rxs.pop().expect("single shard receiver"); - let runtime = tokio::runtime::Handle::current(); - let shard = std::thread::spawn(move || FirehoseWorker::shard(0, rx, state, runtime)); - - tx.send(IndexerMessage::BackfillFinished(fixture.did.clone())) - .await - .into_diagnostic()?; + let [(pending_key, _)] = fixture + .pending_rows()? + .try_into() + .expect("one queued backfill"); + let (tx, shard) = FirehoseWorker::spawn_one_shard(fixture.hydrant.state.clone()); + let (done, drained) = tokio::sync::oneshot::channel(); + tx.send(IndexerMessage::BackfillFinished { + did: fixture.did.clone(), + pending_key: pending_key.into(), + done, + }) + .await + .into_diagnostic()?; drop(tx); timeout( Duration::from_secs(5), @@ -530,6 +534,10 @@ mod tests { .expect("indexer shard should stop") .into_diagnostic()? .map_err(|_| miette::miette!("indexer shard panicked"))?; + assert!( + drained.await.into_diagnostic()?, + "the queued repo wasn't drained" + ); Ok(()) } diff --git a/src/backfill/worker/task.rs b/src/backfill/worker/task.rs index 429c17c..cbbce00 100644 --- a/src/backfill/worker/task.rs +++ b/src/backfill/worker/task.rs @@ -14,7 +14,7 @@ use crate::ingest::indexer::{IndexerMessage, IndexerTx}; use crate::state::AppState; use crate::types::{GaugeState, RepoState, RepoStatus, ResyncErrorKind, ResyncState}; -use super::process::process_did; +use super::process::{BackfillImported, process_did}; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub(super) enum TaskDisposition { @@ -84,6 +84,60 @@ async fn remove_pending(state: &AppState, did: &Did, pending_key: Slice) -> Resu .await } +/// what the shard did with a backfill that committed its records +enum Handoff { + /// applied what it buffered during the fetch and took the repo off the queue + Drained, + /// left it alone, because the queue entry went away or the repo is excluded or gone + Skipped, + /// hydrant is stopping, so the queue entry stays for the next start to run again + Stopping, +} + +/// gives a backfill that committed its records to its shard and waits for the drain. the queue +/// entry stays until the drain commits, so a crash or a failed drain in between leaves the repo +/// queued instead of quietly missing what was buffered during the fetch +async fn hand_off( + state: &Arc, + buffer_tx: &IndexerTx, + did: &Did, + pending_key: &Slice, + imported: BackfillImported, +) -> Result { + state + .db + .run(move |db| { + db.inner + .persist(fjall::PersistMode::Buffer) + .into_diagnostic() + }) + .await?; + #[cfg(feature = "indexer_stream")] + imported.emit(state, did).await; + #[cfg(not(feature = "indexer_stream"))] + drop(imported); + + let (done, drained) = tokio::sync::oneshot::channel(); + let finished = IndexerMessage::BackfillFinished { + did: did.clone(), + pending_key: pending_key.clone(), + done, + }; + let drained = match buffer_tx.send(finished).await { + Ok(()) => drained.await.ok(), + Err(_) => None, + }; + match drained { + Some(true) => Ok(Handoff::Drained), + Some(false) => Ok(Handoff::Skipped), + None if state.stop.is_set() => Ok(Handoff::Stopping), + None => Err(miette::miette!( + "the indexer couldn't drain what was buffered during the fetch" + ) + .into()), + } +} + pub(super) async fn did_task( state: &Arc, http: ThrottledHttpClient, @@ -108,7 +162,7 @@ pub(super) async fn did_task( return Ok(TaskDisposition::Finished); } - match process_did( + let fetched = process_did( state, &http, did, @@ -117,69 +171,26 @@ pub(super) async fn did_task( verify_signatures, strategy, ) - .await - { - Ok(Some(imported)) => { - 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); - if stage_excluded_pending_cleanup(&mut txn, db, &did, pending_key.as_ref())? - { - txn.commit()?; - return Ok::<_, miette::Report>(false); - } - let applied = txn.transition_pending_key( - &did, - pending_key.as_ref(), - GaugeState::Synced, - )?; - db.indexer - .stage_pending_remove(&mut txn.batch, &did, &pending_key); - if applied { - txn.batch.remove(&db.indexer.resync, &did_key); - } - txn.commit()?; - Ok::<_, miette::Report>(applied) - } - }) - .await?; - - if !applied { - return Ok(TaskDisposition::Finished); - } - + .await; + let result = match fetched { + Ok(Some(imported)) => hand_off(state, &buffer_tx, did, &pending_key, imported) + .await + .map(Some), + Ok(None) => Ok(None), + Err(e) => Err(e), + }; + match result { + Ok(Some(Handoff::Drained)) => { 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?; - - #[cfg(feature = "indexer_stream")] - imported.emit(&state, did).await; - #[cfg(not(feature = "indexer_stream"))] - drop(imported); - - if let Err(e) = buffer_tx - .send(IndexerMessage::BackfillFinished(did.clone())) - .await - { - error!(err = %e, "failed to send BackfillFinished"); - } Ok(TaskDisposition::Finished) } + Ok(Some(Handoff::Skipped)) => { + remove_pending_if_excluded(state, did, pending_key).await?; + Ok(TaskDisposition::Finished) + } + Ok(Some(Handoff::Stopping)) => Ok(TaskDisposition::Finished), Ok(None) => { remove_pending(state, did, pending_key).await?; Ok(TaskDisposition::Finished) @@ -358,12 +369,13 @@ mod tests { use crate::db; use crate::db::types::DbRkey; use crate::ingest::indexer::FirehoseWorker; + use crate::ingest::stream::{Commit as StreamCommit, Datetime, RepoOp, RepoOpAction}; use crate::types::RepoMetadata; use axum::Router; use axum::routing::get; use bytes::Bytes; use jacquard_common::CowStr; - use jacquard_common::types::cid::IpldCid; + use jacquard_common::types::cid::{CidLink, IpldCid}; use jacquard_common::types::string::Handle; use jacquard_common::types::tid::Tid; use jacquard_repo::commit::Commit as AtpCommit; @@ -784,7 +796,22 @@ mod tests { } async fn run(&self, pending_key: [u8; 8]) -> Result { - let (tx, _worker) = FirehoseWorker::new(self.state.clone(), 1); + let (tx, shard) = FirehoseWorker::spawn_one_shard(self.state.clone()); + let result = self.run_with(tx, pending_key).await; + // the task held the shard's only sender, so the shard stops now + tokio::task::spawn_blocking(move || shard.join()) + .await + .into_diagnostic()? + .map_err(|_| miette::miette!("indexer shard panicked"))?; + result + } + + /// `run`, handing the finished backfill to `tx`, which may have no shard behind it + async fn run_with( + &self, + tx: IndexerTx, + pending_key: [u8; 8], + ) -> Result { let client = reqwest::Client::builder() .resolve("pds.example", self.pds.address) .build() @@ -808,6 +835,63 @@ mod tests { .await } + /// a live commit creating `rkey` out of `car`, buffered the way the shard does while + /// the repo is queued + fn buffer(&self, car: &BuiltCar, rkey: &str, body: &[u8]) -> miette::Result<()> { + let root = crate::car::parse_car(car.bytes.clone(), false)?.root; + let cid = jacquard_repo::mst::util::compute_cid(body).into_diagnostic()?; + let commit = StreamCommit { + blobs: Vec::new(), + blocks: car.bytes.clone(), + commit: CidLink::from(root), + ops: vec![RepoOp { + action: RepoOpAction::Create, + cid: Some(CidLink::from(cid)), + path: CowStr::Owned(format!("{COLLECTION}/{rkey}").into()), + prev: None, + }], + prev_data: None, + rebase: false, + repo: self.did.clone(), + rev: car.rev.clone(), + seq: 1, + since: None, + time: Datetime::try_from("2026-10-05T12:00:00Z".to_owned()).into_diagnostic()?, + too_big: false, + }; + let mut txn = DbTxn::new(&self.state.db); + crate::ops::persist_to_resync_buffer(&mut txn, &self.state.db, &self.did, &commit)?; + txn.commit() + } + + /// queues the repo under 7 and serves a car with record `a`, after the shard buffered a + /// live commit adding `b` during the fetch. returns that commit's car + async fn fetch_with_a_buffered_commit(&self) -> miette::Result { + self.seed(&rooted("3jzfcijpj2z2a")?, 7, &[])?; + let a = post("a"); + let b = post("b"); + let fetched = build_car("3jzfcijpj2z2b", &[("a", a.as_slice())]).await?; + let newer = + build_car("3jzfcijpj2z2c", &[("a", a.as_slice()), ("b", b.as_slice())]).await?; + self.buffer(&newer, "b", &b)?; + self.pds.serve_car(fetched.bytes).await; + Ok(newer) + } + + /// a channel with no shard behind it, so nobody answers the handoff + fn unanswered(&self) -> IndexerTx { + FirehoseWorker::new(self.state.clone(), 1).0 + } + + fn buffered(&self) -> usize { + self.state + .db + .indexer + .resync_buffer + .prefix(keys::resync_buffer_prefix(&self.did)) + .count() + } + /// runs the task for `pending_key` and applies `change` while the task waits on getRepo, /// the way another writer can move the queue under a task that's already in flight async fn run_changing( @@ -1503,7 +1587,7 @@ mod tests { ); let after = fixture.repo()?; assert!(after.active); - assert_eq!(after.status, RepoStatus::Desynchronized); + assert_eq!(after.status, RepoStatus::Synced); assert_eq!(fixture.record_rkeys()?, ["a"]); Ok(()) } @@ -1537,6 +1621,69 @@ mod tests { Ok(()) } + #[tokio::test] + async fn a_commit_buffered_during_the_fetch_is_applied_before_the_repo_leaves_the_queue() + -> miette::Result<()> { + let fixture = Fixture::new().await?; + let newer = fixture.fetch_with_a_buffered_commit().await?; + + assert_eq!( + fixture.run(keys::pending_key(7)).await?, + TaskDisposition::Finished + ); + assert_eq!(fixture.record_rkeys()?, ["a", "b"]); + let after = fixture.repo()?; + assert_eq!( + after.root.map(|root| root.rev), + Some(crate::db::types::DbTid::from(&newer.rev)) + ); + assert_eq!(after.status, RepoStatus::Synced); + assert!(fixture.pending_keys()?.is_empty()); + assert_eq!(fixture.buffered(), 0); + Ok(()) + } + + #[tokio::test] + async fn a_backfill_whose_drain_fails_is_retried_with_its_buffer_kept() -> miette::Result<()> { + let fixture = Fixture::new().await?; + fixture.fetch_with_a_buffered_commit().await?; + let result = fixture + .run_with(fixture.unanswered(), keys::pending_key(7)) + .await; + + assert!( + matches!(result, Err(BackfillError::Generic(_))), + "{result:?}" + ); + assert!(matches!( + fixture.resync()?, + Some(ResyncState::Error { retry_count: 1, .. }) + )); + assert!(fixture.pending_keys()?.is_empty()); + assert_eq!(fixture.buffered(), 1); + Ok(()) + } + + #[tokio::test] + async fn a_backfill_stopped_before_its_drain_stays_queued() -> miette::Result<()> { + let fixture = Fixture::new().await?; + fixture.seed(&rooted("3jzfcijpj2z2a")?, 7, &[])?; + let a = post("a"); + let fetched = build_car("3jzfcijpj2z2b", &[("a", a.as_slice())]).await?; + fixture.pds.serve_car(fetched.bytes).await; + fixture.state.stop.trigger(); + + assert_eq!( + fixture + .run_with(fixture.unanswered(), keys::pending_key(7)) + .await?, + TaskDisposition::Finished + ); + assert_eq!(fixture.pending_keys()?, [keys::pending_key(7).to_vec()]); + assert!(fixture.resync()?.is_none()); + Ok(()) + } + #[tokio::test] async fn a_newer_event_repeating_the_old_value_beats_what_the_fetch_saw() -> miette::Result<()> { diff --git a/src/ingest/indexer/message.rs b/src/ingest/indexer/message.rs index 113ac6b..790ba83 100644 --- a/src/ingest/indexer/message.rs +++ b/src/ingest/indexer/message.rs @@ -1,7 +1,9 @@ use crate::ingest::mailbox::{ShardedMessage, ShardedReceiver, ShardedSender}; use crate::ingest::stream; +use fjall::Slice; use jacquard_common::types::did::Did; use miette::Result; +use tokio::sync::oneshot; use url::Url; #[derive(Debug)] @@ -56,8 +58,14 @@ pub enum IndexerMessage { Event(Box), /// a new repo was discovered and needs backfill. NewRepo(Did), - /// backfill for this DID has completed; drain the resync buffer. - BackfillFinished(Did), + /// a backfill committed its records. the shard applies what it buffered during the fetch + /// and takes the repo off the queue in one batch, if `pending_key` is still what the repo + /// is queued under, then says on `done` whether it did. a dropped `done` means it failed + BackfillFinished { + did: Did, + pending_key: Slice, + done: oneshot::Sender, + }, } #[derive(Clone, Debug)] @@ -94,7 +102,7 @@ impl ShardedMessage for IndexerMessage { IndexerEventData::Sync(did) => did, }, IndexerMessage::NewRepo(did) => did, - IndexerMessage::BackfillFinished(did) => did, + IndexerMessage::BackfillFinished { did, .. } => did, }; (crate::util::hash(did) as usize) % num_shards diff --git a/src/ingest/indexer/shard.rs b/src/ingest/indexer/shard.rs index 7804363..04a8c6a 100644 --- a/src/ingest/indexer/shard.rs +++ b/src/ingest/indexer/shard.rs @@ -58,56 +58,19 @@ impl FirehoseWorker { lifecycle_transitions: &mut lifecycle_transitions, notify_backfill: false, }; + let mut finished = None; match msg { - IndexerMessage::BackfillFinished(did) => { + IndexerMessage::BackfillFinished { + did, + pending_key, + done, + } => { let _span = tracing::info_span!("ingest", did = %did).entered(); debug!("backfill finished, verifying state and draining buffer"); - match ctx.txn.lock_repo_and_is_excluded(&did) { - Ok(false) => {} - Ok(true) => continue, - Err(e) => { - error!(err = %e, "failed to check repository exclusion"); - continue; - } - } - - let repo_key = keys::repo_key(&did); - if let Ok(Some(state_bytes)) = state.db.repos.get(&repo_key).into_diagnostic() { - match crate::db::deser_repo_state(&state_bytes) { - Ok(repo_state) => { - let repo_state = repo_state.into_static(); - - match Self::drain_resync_buffer(&mut ctx, &did, repo_state) { - Ok(RepoProcessResult::Ok(s)) => { - // a repo that went inactive while it backfilled keeps - // its status, the backfill only caught its records up - let status = if s.active { - RepoStatus::Synced - } else { - s.status.clone() - }; - let res = ops::transition_repo( - &mut ctx.txn.batch, - &state.db, - ctx.lifecycle_transitions, - &did, - s, - status, - ); - if let Err(e) = res { - error!(err = %e, "failed to transition the backfilled repo"); - } - } - Ok(RepoProcessResult::NeedsBackfill(_)) => {} - Ok(RepoProcessResult::Deleted) => {} - Err(e) => { - error!(err = %e, "failed to drain resync buffer") - } - }; - } - Err(e) => error!(err = %e, "failed to deser repo state"), - } + match Self::finish_backfill(&mut ctx, &did, &pending_key) { + Ok(drained) => finished = Some((done, drained)), + Err(e) => error!(err = %e, "failed to drain the resync buffer"), } } IndexerMessage::NewRepo(did) => { @@ -220,14 +183,13 @@ impl FirehoseWorker { chain_break, parsed_blocks, ) { - Ok(RepoProcessResult::Ok(_)) => {} + Ok(RepoProcessResult::Ok) => {} Ok(RepoProcessResult::Deleted) => { ctx.txn.counts.add_repos(-1); } - Ok(RepoProcessResult::NeedsBackfill(Some(commit))) => { + Ok(RepoProcessResult::NeedsBackfill(commit)) => { try_persist(&mut ctx, commit); } - Ok(RepoProcessResult::NeedsBackfill(None)) => {} Err(e) => { if let IngestError::Generic(ref r) = e { state.db.poison.check_report(r); @@ -282,6 +244,11 @@ impl FirehoseWorker { error!(shard = id, err = %e, "failed to commit transaction"); continue; } + // only after the commit, because the backfill takes the answer to mean its repo is + // off the queue + if let Some((done, drained)) = finished { + let _ = done.send(drained); + } if ctx.notify_backfill { state.notify_backfill(); } @@ -315,7 +282,7 @@ impl FirehoseWorker { commit: &'c Commit<'c>, chain_break: bool, parsed_blocks: crate::car::CommitCar, - ) -> Result, IngestError> { + ) -> Result, IngestError> { let db = &ctx.state.db; let did = &commit.repo; repo_state.advance_message_time(commit.time.0.timestamp_millis()); @@ -334,13 +301,25 @@ impl FirehoseWorker { if chain_break { debug!("chain break detected, triggering backfill"); Self::trigger_backfill(ctx, did, repo_state)?; - return Ok(RepoProcessResult::NeedsBackfill(Some(commit))); + return Ok(RepoProcessResult::NeedsBackfill(commit)); } if is_backfilling { - return Ok(RepoProcessResult::NeedsBackfill(Some(commit))); + return Ok(RepoProcessResult::NeedsBackfill(commit)); } + Self::apply_checked_commit(ctx, repo_state, commit, parsed_blocks) + .map(|_| RepoProcessResult::Ok) + } + + /// applies a commit that got past [`Self::handle_commit`]'s checks. the drain comes in here + /// directly, because its own backfill still holds the queue entry those checks buffer behind + fn apply_checked_commit<'s>( + ctx: &mut WorkerContext, + repo_state: RepoState<'s>, + commit: &Commit<'_>, + parsed_blocks: crate::car::CommitCar, + ) -> Result, IngestError> { let root_bytes = parsed_blocks .blocks .get(&parsed_blocks.root) @@ -364,7 +343,7 @@ impl FirehoseWorker { validated, &ctx.state.filter.load(), )?; - Ok(RepoProcessResult::Ok(res.repo_state)) + Ok(res.repo_state) } #[cfg_attr(not(feature = "indexer_stream"), allow(unused_variables))] @@ -373,7 +352,7 @@ impl FirehoseWorker { repo_state: RepoState<'s>, identity: &Identity, changed: bool, - ) -> Result, IngestError> { + ) -> Result, IngestError> { #[cfg(feature = "indexer_stream")] { let did = &identity.did; @@ -396,7 +375,7 @@ impl FirehoseWorker { } } - Ok(RepoProcessResult::Ok(repo_state)) + Ok(RepoProcessResult::Ok) } fn handle_account<'s, 'c>( @@ -405,7 +384,7 @@ impl FirehoseWorker { changed: bool, account: &'c Account<'c>, was_active: bool, - ) -> Result, IngestError> { + ) -> Result, IngestError> { let db = &ctx.state.db; let did = &account.did; let is_inactive = !account.active; @@ -458,14 +437,61 @@ impl FirehoseWorker { ctx.txn.outbox.stream.push_account(evt); } - Ok(RepoProcessResult::Ok(repo_state)) + Ok(RepoProcessResult::Ok) + } + + /// applies what a finished backfill's repo buffered during the fetch and takes the repo off + /// the queue in the same batch, while `pending_key` is still what it's queued under. `false` + /// when it isn't, or the repo is excluded or gone, so nothing was drained + fn finish_backfill( + ctx: &mut WorkerContext, + did: &Did, + pending_key: &[u8], + ) -> Result { + if ctx.txn.lock_repo_and_is_excluded(did)? + || !ctx.txn.lock_repo_and_is_pending(did, pending_key)? + { + return Ok(false); + } + let Some(bytes) = ctx + .state + .db + .repos + .get(keys::repo_key(did)) + .into_diagnostic()? + else { + return Ok(false); + }; + let repo_state = crate::db::deser_repo_state(&bytes)?.into_static(); + let repo_state = Self::drain_resync_buffer(ctx, did, repo_state)?; + // a repo that went inactive while it backfilled keeps its status, the backfill only + // caught its records up + let status = if repo_state.active { + RepoStatus::Synced + } else { + repo_state.status.clone() + }; + ops::transition_repo( + &mut ctx.txn.batch, + &ctx.state.db, + ctx.lifecycle_transitions, + did, + repo_state, + status, + )?; + // transition_repo drops the key metadata names, which an orphaned key isn't + ctx.state + .db + .indexer + .stage_pending_remove(&mut ctx.txn.batch, did, pending_key); + Ok(true) } fn drain_resync_buffer<'s>( ctx: &mut WorkerContext, did: &Did, mut repo_state: RepoState<'s>, - ) -> Result, IngestError> { + ) -> Result, IngestError> { let db = &ctx.state.db; let prefix = keys::resync_buffer_prefix(did); @@ -478,33 +504,22 @@ impl FirehoseWorker { // buffer, so each record block is hashed once more as it is stored let parsed_blocks = crate::car::parse_commit_car(commit.blocks.clone(), false) .map_err(|e| IngestError::Generic(miette::miette!("malformed CAR: {e}")))?; - let res = Self::handle_commit(ctx, repo_state, &commit, false, parsed_blocks); - let res = match res { - Ok(r) => r, + repo_state.advance_message_time(commit.time.0.timestamp_millis()); + match Self::apply_checked_commit(ctx, repo_state, &commit, parsed_blocks) { + Ok(rs) => { + ctx.txn.batch.remove(&db.indexer.resync_buffer, key); + repo_state = rs; + } Err(e) => { if !Self::check_if_retriable_failure(&e) { ctx.txn.batch.remove(&db.indexer.resync_buffer, key); } return Err(e); } - }; - match res { - RepoProcessResult::Ok(rs) => { - ctx.txn.batch.remove(&db.indexer.resync_buffer, key); - repo_state = rs; - } - RepoProcessResult::NeedsBackfill(_) => { - // commit is already in the buffer, leave it there for the next backfill - return Ok(RepoProcessResult::NeedsBackfill(None)); - } - RepoProcessResult::Deleted => { - ctx.txn.batch.remove(&db.indexer.resync_buffer, key); - return Ok(RepoProcessResult::Deleted); - } } } - Ok(RepoProcessResult::Ok(repo_state)) + Ok(repo_state) } fn trigger_backfill<'s>( @@ -574,8 +589,12 @@ mod tests { ingest::stream::{Datetime, RepoOp, RepoOpAction}, }; - #[tokio::test] - async fn backfill_finished_persists_synced_repo_status() -> Result<()> { + /// a repo row whose metadata names index `named` while the queue holds the repo under `queued` + fn queued_repo( + row: &RepoState<'_>, + named: u64, + queued: u64, + ) -> Result<(tempfile::TempDir, Arc, Did)> { let tmp = tempfile::tempdir().into_diagnostic()?; let config = Config { database_path: tmp.path().to_path_buf(), @@ -585,56 +604,71 @@ mod tests { let did = Did::new("did:plc:backfillfinishedstatus") .into_diagnostic()? .into_static(); - let index_id = 7; - let pending_key = keys::pending_key(index_id); - let mut batch = state.db.inner.batch(); batch.insert( &state.db.repos, keys::repo_key(&did), - db::ser_repo_state(&RepoState::backfilling())?, + db::ser_repo_state(row)?, ); batch.insert( &state.db.repo_metadata, keys::repo_metadata_key(&did), - ser_repo_meta(&RepoMetadata::backfilling(index_id))?, + ser_repo_meta(&RepoMetadata::backfilling(named))?, ); state .db .indexer - .stage_pending_insert(&mut batch, &did, &pending_key); + .stage_pending_insert(&mut batch, &did, &keys::pending_key(queued)); batch.commit().into_diagnostic()?; + Ok((tmp, state, did)) + } - let (tx, mut worker) = FirehoseWorker::new(state.clone(), 1); - let rx = worker.rxs.pop().expect("single shard receiver"); - let shard_state = state.clone(); - let runtime = TokioHandle::current(); - let shard = std::thread::spawn(move || { - FirehoseWorker::shard(0, rx, shard_state, runtime); - }); + fn repo_row(state: &AppState, did: &Did) -> Result> { + let bytes = state + .db + .repos + .get(keys::repo_key(did)) + .into_diagnostic()? + .expect("repository state"); + Ok(db::deser_repo_state(&bytes)?.into_static()) + } - tx.send(IndexerMessage::BackfillFinished(did.clone())) - .await - .into_diagnostic()?; + /// hands `did`'s finished backfill under `pending_key` to a shard and says whether it drained + async fn finish_on_a_shard( + state: &Arc, + did: &Did, + pending_key: [u8; 8], + ) -> Result { + let (tx, shard) = FirehoseWorker::spawn_one_shard(state.clone()); + let (done, drained) = tokio::sync::oneshot::channel(); + tx.send(IndexerMessage::BackfillFinished { + did: did.clone(), + pending_key: pending_key.into(), + done, + }) + .await + .into_diagnostic()?; drop(tx); shard .join() .map_err(|_| miette::miette!("indexer shard panicked"))?; + drained.await.into_diagnostic() + } - let state_bytes = state - .db - .repos - .get(keys::repo_key(&did)) - .into_diagnostic()? - .expect("repository state"); - let repo_state = db::deser_repo_state(&state_bytes)?; + #[tokio::test] + async fn backfill_finished_persists_synced_repo_status() -> Result<()> { + let (_tmp, state, did) = queued_repo(&RepoState::backfilling(), 7, 7)?; + + assert!(finish_on_a_shard(&state, &did, keys::pending_key(7)).await?); + + let repo_state = repo_row(&state, &did)?; assert_eq!(repo_state.status, RepoStatus::Synced); assert!(repo_state.active); assert!( !state .db .indexer - .is_pending(&pending_key) + .is_pending(&keys::pending_key(7)) .into_diagnostic()? ); state.db.indexer.assert_pending_dids_match(); @@ -644,55 +678,15 @@ mod tests { #[tokio::test] async fn backfill_finished_leaves_a_repo_that_went_inactive_inactive() -> Result<()> { - let tmp = tempfile::tempdir().into_diagnostic()?; - let config = Config { - database_path: tmp.path().to_path_buf(), - ..Default::default() - }; - let state = Arc::new(AppState::new(&config)?); - let did = Did::new("did:plc:backfillfinishedstatus") - .into_diagnostic()? - .into_static(); - - // deactivated while its getRepo streamed, so the task already dropped its queue entry + // deactivated while its getRepo streamed let mut deactivated = RepoState::backfilling(); deactivated.active = false; deactivated.status = RepoStatus::Deactivated; - let mut batch = state.db.inner.batch(); - batch.insert( - &state.db.repos, - keys::repo_key(&did), - db::ser_repo_state(&deactivated)?, - ); - batch.insert( - &state.db.repo_metadata, - keys::repo_metadata_key(&did), - ser_repo_meta(&RepoMetadata::backfilling(7))?, - ); - batch.commit().into_diagnostic()?; + let (_tmp, state, did) = queued_repo(&deactivated, 7, 7)?; - let (tx, mut worker) = FirehoseWorker::new(state.clone(), 1); - let rx = worker.rxs.pop().expect("single shard receiver"); - let shard_state = state.clone(); - let runtime = TokioHandle::current(); - let shard = std::thread::spawn(move || { - FirehoseWorker::shard(0, rx, shard_state, runtime); - }); - tx.send(IndexerMessage::BackfillFinished(did.clone())) - .await - .into_diagnostic()?; - drop(tx); - shard - .join() - .map_err(|_| miette::miette!("indexer shard panicked"))?; + assert!(finish_on_a_shard(&state, &did, keys::pending_key(7)).await?); - let state_bytes = state - .db - .repos - .get(keys::repo_key(&did)) - .into_diagnostic()? - .expect("repository state"); - let repo_state = db::deser_repo_state(&state_bytes)?; + let repo_state = repo_row(&state, &did)?; assert_eq!(repo_state.status, RepoStatus::Deactivated); assert!(!repo_state.active); let resync = state @@ -711,6 +705,62 @@ mod tests { Ok(()) } + #[tokio::test] + async fn a_backfill_whose_repo_was_requeued_drains_nothing() -> Result<()> { + // requeued under 8 while the backfill queued under 7 fetched + let (_tmp, state, did) = queued_repo(&RepoState::backfilling(), 8, 8)?; + let buffered = keys::resync_buffer_key( + &did, + crate::db::types::DbTid::new_from_bytes(4_u64.to_be_bytes()), + ); + state + .db + .indexer + .resync_buffer + .insert(&buffered, b"for the next backfill") + .into_diagnostic()?; + + assert!(!finish_on_a_shard(&state, &did, keys::pending_key(7)).await?); + assert!( + state + .db + .indexer + .is_pending(&keys::pending_key(8)) + .into_diagnostic()? + ); + assert!( + state + .db + .indexer + .resync_buffer + .get(&buffered) + .into_diagnostic()? + .is_some() + ); + assert_eq!( + repo_row(&state, &did)?.status, + RepoState::backfilling().status + ); + state.db.indexer.assert_pending_dids_match(); + Ok(()) + } + + #[tokio::test] + async fn a_drain_also_drops_a_queue_entry_its_metadata_doesnt_name() -> Result<()> { + let (_tmp, state, did) = queued_repo(&RepoState::backfilling(), 8, 7)?; + + 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()); diff --git a/src/ingest/indexer/worker.rs b/src/ingest/indexer/worker.rs index 6e512ec..8fe02f3 100644 --- a/src/ingest/indexer/worker.rs +++ b/src/ingest/indexer/worker.rs @@ -9,7 +9,6 @@ use crate::ingest::mailbox::ShardedSender; use crate::ingest::stream::Commit; use crate::resolver::{NoSigningKeyError, ResolverError}; use crate::state::AppState; -use crate::types::RepoState; use jacquard_repo::error::CommitError; #[derive(Debug, Diagnostic, Error)] @@ -37,14 +36,12 @@ impl From for IngestError { } #[derive(Debug)] -#[allow(clippy::large_enum_variant)] -pub(crate) enum RepoProcessResult<'s, 'c> { - // message processed successfully, here is the (possibly updated) state - Ok(RepoState<'s>), +pub(crate) enum RepoProcessResult<'c> { + Ok, // repo was deleted as part of processing Deleted, - // needs backfill; carries the triggering commit to buffer (None when already in the buffer) - NeedsBackfill(Option<&'c Commit<'c>>), + // needs backfill and carries the triggering commit to buffer + NeedsBackfill(&'c Commit<'c>), } pub struct FirehoseWorker { @@ -58,6 +55,18 @@ impl FirehoseWorker { (IndexerTx { inner }, Self { state, rxs }) } + /// one shard on its own thread, which stops once every sender is dropped + #[cfg(test)] + pub(crate) fn spawn_one_shard( + state: Arc, + ) -> (IndexerTx, std::thread::JoinHandle<()>) { + let (tx, mut worker) = Self::new(state.clone(), 1); + let rx = worker.rxs.pop().expect("single shard receiver"); + let runtime = TokioHandle::current(); + let shard = std::thread::spawn(move || Self::shard(0, rx, state, runtime)); + (tx, shard) + } + pub fn run(self, handle: TokioHandle) -> Result<()> { let num_shards = self.rxs.len(); let (exit_tx, exit_rx) = std::sync::mpsc::channel(); -- 2.51.2