diff --git a/.beads/interactions.jsonl b/.beads/interactions.jsonl index 9911bb3..787814c 100644 --- a/.beads/interactions.jsonl +++ b/.beads/interactions.jsonl @@ -87,3 +87,4 @@ {"id":"int-d47b701c216df2d75a8351144aa4e648","kind":"field_change","created_at":"2026-08-05T13:28:54.108148Z","actor":"dawn","issue_id":"hydrant-0gg","extra":{"field":"status","new_value":"closed","old_value":"in_progress","reason":"Implemented and fully validated on the final working tree."}} {"id":"int-a476263d50fdc6414edb9ac747ba69b6","kind":"field_change","created_at":"2026-08-05T13:28:54.314066Z","actor":"dawn","issue_id":"hydrant-4hr","extra":{"field":"status","new_value":"closed","old_value":"in_progress","reason":"Implemented and fully validated on the final working tree."}} {"id":"int-ccc4d70168c95b110d3b7e29fcf40451","kind":"field_change","created_at":"2026-08-10T22:32:35.356361Z","actor":"dawn","issue_id":"hydrant-dom","extra":{"field":"status","new_value":"closed","old_value":"in_progress","reason":"Closed"}} +{"id":"int-1981f331f8b8daf93820df3650c16c6a","kind":"field_change","created_at":"2026-08-13T14:45:08.148441Z","actor":"dawn","issue_id":"hydrant-5dn","extra":{"field":"status","new_value":"closed","old_value":"in_progress","reason":"Persisted transition_repo state in the transaction; regression exercises BackfillFinished and rereads Synced from db.repos."}} diff --git a/src/ingest/indexer/shard.rs b/src/ingest/indexer/shard.rs index 4dafc2f..da001bf 100644 --- a/src/ingest/indexer/shard.rs +++ b/src/ingest/indexer/shard.rs @@ -625,7 +625,76 @@ mod tests { }; use super::*; - use crate::ingest::stream::{Datetime, RepoOp, RepoOpAction}; + use crate::{ + config::Config, + ingest::stream::{Datetime, RepoOp, RepoOpAction}, + }; + + #[tokio::test] + async fn backfill_finished_persists_synced_repo_status() -> 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(); + 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())?, + ); + batch.insert( + &state.db.repo_metadata, + keys::repo_metadata_key(&did), + ser_repo_meta(&RepoMetadata::backfilling(index_id))?, + ); + batch.insert(&state.db.indexer.pending, pending_key, keys::repo_key(&did)); + batch.commit().into_diagnostic()?; + + 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"))?; + + 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)?; + assert_eq!(repo_state.status, RepoStatus::Synced); + assert!(repo_state.active); + assert!( + state + .db + .indexer + .pending + .get(pending_key) + .into_diagnostic()? + .is_none() + ); + + Ok(()) + } #[test] fn resync_buffer_commit_msgpack_round_trip_preserves_path() { diff --git a/src/ops.rs b/src/ops.rs index 22a8c6f..cd2ac2e 100644 --- a/src/ops.rs +++ b/src/ops.rs @@ -201,6 +201,7 @@ pub fn transition_repo<'s>( repo_state.active = matches!(new_status, RepoStatus::Synced | RepoStatus::Error(_)); repo_state.status = new_status; repo_state.touch(); + batch.insert(&db.repos, repo_key, db::ser_repo_state(&repo_state)?); Ok(repo_state) }