diff --git a/src/backfill/manager.rs b/src/backfill/manager.rs index faced97..bbf8011 100644 --- a/src/backfill/manager.rs +++ b/src/backfill/manager.rs @@ -1,5 +1,5 @@ use crate::db::types::TrimmedDid; -use crate::db::{self, deser_repo_state, ser_repo_state}; +use crate::db::{self, deser_repo_state, keys, ser_repo_state}; use crate::state::AppState; use crate::types::{RepoStatus, ResyncState}; use miette::{IntoDiagnostic, Result}; @@ -28,18 +28,13 @@ pub fn queue_gone_backfills(state: &Arc) -> Result<()> { // move back to pending let mut batch = state.db.inner.batch(); batch.remove(&state.db.resync, key.clone()); - batch.insert( - &state.db.pending, - crate::db::keys::pending_key(&did), - Vec::new(), - ); + batch.insert(&state.db.pending, key.clone(), Vec::new()); // update repo state back to backfilling - let repo_key = crate::db::keys::repo_key(&did); - if let Some(state_bytes) = state.db.repos.get(&repo_key).into_diagnostic()? { + if let Some(state_bytes) = state.db.repos.get(&key).into_diagnostic()? { let mut repo_state = deser_repo_state(&state_bytes)?; repo_state.status = RepoStatus::Backfilling; - batch.insert(&state.db.repos, &repo_key, ser_repo_state(&repo_state)?); + batch.insert(&state.db.repos, key, ser_repo_state(&repo_state)?); } state.db.update_count("resync", -1); @@ -91,10 +86,7 @@ pub fn retry_worker(state: Arc) { // move back to pending state.db.update_count("pending", 1); - if let Err(e) = db - .pending - .insert(crate::db::keys::pending_key(&did), Vec::new()) - { + if let Err(e) = db.pending.insert(keys::repo_key(&did), Vec::new()) { error!("failed to move {did} to pending: {e}"); db::check_poisoned(&e); continue; diff --git a/src/backfill/mod.rs b/src/backfill/mod.rs index e8218ba..508f589 100644 --- a/src/backfill/mod.rs +++ b/src/backfill/mod.rs @@ -126,12 +126,7 @@ impl BackfillWorker { let mut spawned = 0; - // limit the number of active tasks based on adaptive limit - // we iterate in reverse to prioritize newer items (LIFO) - // effective key comparison: {timestamp}|{did} - // older timestamps are smaller, newer are larger. - // rev() starts from largest (newest). - for guard in self.state.db.pending.iter().rev() { + for guard in self.state.db.pending.iter() { if self.in_flight.len() >= limiter.current_limit { break; } @@ -145,17 +140,12 @@ impl BackfillWorker { } }; - let did = if key.len() > 9 && key[8] == keys::SEP { - match TrimmedDid::try_from(&key[9..]) { - Ok(d) => d.to_did(), - Err(e) => { - error!("invalid did '{key:?}' in pending: {e}"); - continue; - } + let did = match TrimmedDid::try_from(key.as_ref()) { + Ok(d) => d.to_did(), + Err(e) => { + error!("invalid did '{key:?}' in pending: {e}"); + continue; } - } else { - error!("invalid did '{key:?}' in pending"); - continue; }; if self.in_flight.contains_sync(&did) { diff --git a/src/db/keys.rs b/src/db/keys.rs index 48a3f17..916efef 100644 --- a/src/db/keys.rs +++ b/src/db/keys.rs @@ -114,15 +114,3 @@ pub fn resync_buffer_prefix(did: &Did) -> Vec { prefix.push(SEP); prefix } - -// key format: {timestamp}|{DID} (DID trimmed) -// timestamp is big-endian u64 micros -pub fn pending_key(did: &Did) -> Vec { - let repo = TrimmedDid::from(did); - let mut key = Vec::with_capacity(8 + 1 + repo.len()); - let ts = chrono::Utc::now().timestamp_micros() as u64; - key.extend_from_slice(&ts.to_be_bytes()); - key.push(SEP); - repo.write_to_vec(&mut key); - key -} diff --git a/src/ingest/worker.rs b/src/ingest/worker.rs index a279b4b..64a0d29 100644 --- a/src/ingest/worker.rs +++ b/src/ingest/worker.rs @@ -360,7 +360,7 @@ impl FirehoseWorker { RepoStatus::Backfilling, )?; ctx.state.db.update_count("pending", 1); - batch.insert(&ctx.state.db.pending, keys::pending_key(did), &[]); + batch.insert(&ctx.state.db.pending, keys::repo_key(did), &[]); batch.commit().into_diagnostic()?; ctx.state.notify_backfill(); return Ok(RepoProcessResult::Ok(repo_state)); @@ -506,7 +506,7 @@ impl FirehoseWorker { RepoStatus::Backfilling, )?; ctx.state.db.update_count("pending", 1); - batch.insert(&ctx.state.db.pending, keys::pending_key(did), &[]); + batch.insert(&ctx.state.db.pending, keys::repo_key(did), &[]); batch.commit().into_diagnostic()?; ctx.repo_cache .insert(did.clone().into_static(), repo_state.clone().into_static()); @@ -567,7 +567,7 @@ impl FirehoseWorker { &repo_key, crate::db::ser_repo_state(&new_state)?, ); - batch.insert(&ctx.state.db.pending, keys::pending_key(did), &[]); + batch.insert(&ctx.state.db.pending, repo_key, &[]); ctx.state.db.update_count("repos", 1); ctx.state.db.update_count("pending", 1); batch.commit().into_diagnostic()?;