From 35731cdbed32407a83f9cdbf40ddd244526ef290 Mon Sep 17 00:00:00 2001 From: phil Date: Tue, 17 Mar 2026 18:50:09 -0400 Subject: [PATCH] resync queue timing metrics --- src/storage/resync_queue.rs | 21 +++++++++++++++++++++ 1 file changed, 21 insertions(+) diff --git a/src/storage/resync_queue.rs b/src/storage/resync_queue.rs index c9bcb85..54518e3 100644 --- a/src/storage/resync_queue.rs +++ b/src/storage/resync_queue.rs @@ -208,6 +208,7 @@ pub fn dequeue_ready( now: SystemTime, since: Option>, ) -> StorageResult)>> { + let t0 = std::time::Instant::now(); let now_ms = crate::util::to_millis(now); let prefix = key_prefix_all(); @@ -219,6 +220,8 @@ pub fn dequeue_ready( .range(prefixed_range(&prefix, lower_suffix..upper_suffix)) .next() else { + metrics::histogram!("lightrail_resync_queue_op_seconds", "op" => "dequeue", "outcome" => "empty") + .record(t0.elapsed().as_secs_f64()); return Ok(None); }; @@ -239,6 +242,8 @@ pub fn dequeue_ready( let next_since = key_bytes .get(prefix.len()..) .expect("a resync queue key must start with the resync queue prefix"); + metrics::histogram!("lightrail_resync_queue_op_seconds", "op" => "dequeue", "outcome" => "found") + .record(t0.elapsed().as_secs_f64()); Ok(Some((item, next_since.to_vec()))) } @@ -257,15 +262,19 @@ pub fn claim_resync( since: Option>, busy: &HashSet>, ) -> StorageResult)>> { + let t0 = std::time::Instant::now(); let now_ms = crate::util::to_millis(now); let prefix = key_prefix_all(); let lower_suffix = since.unwrap_or_default(); let upper_suffix = key_ts_midfix(now_ms); + let mut scanned: u64 = 0; + for guard in db .ks .range(prefixed_range(&prefix, lower_suffix..upper_suffix)) { + scanned += 1; let (key_slice, val_slice) = guard.into_inner()?; let key_bytes = key_slice.as_ref(); let (_, did) = key_parse(key_bytes)?; @@ -275,6 +284,8 @@ pub fn claim_resync( continue; } + let scan_elapsed = t0.elapsed(); + let key_str = String::from_utf8_lossy(key_bytes).into_owned(); let item = decode(val_slice.as_ref(), &key_str, did.clone())?; let next_since = key_bytes[prefix.len()..].to_vec(); @@ -303,6 +314,7 @@ pub fn claim_resync( }; // Atomically: remove from queue + write state=Resyncing. + let t_commit = std::time::Instant::now(); let mut batch = db.database.batch(); batch.remove_weak(&db.ks, key_bytes); batch.insert(&db.ks, &repo_key, repo::encode_repo_info(&new_info)); @@ -315,9 +327,18 @@ pub fn claim_resync( retry = item.retry_count, "claimed resync job" ); + + metrics::histogram!("lightrail_resync_queue_op_seconds", "op" => "claim", "phase" => "scan") + .record(scan_elapsed.as_secs_f64()); + metrics::histogram!("lightrail_resync_queue_op_seconds", "op" => "claim", "phase" => "commit") + .record(t_commit.elapsed().as_secs_f64()); + metrics::histogram!("lightrail_resync_queue_claim_scanned").record(scanned as f64); return Ok(Some((item, next_since))); } + metrics::histogram!("lightrail_resync_queue_op_seconds", "op" => "claim", "phase" => "scan") + .record(t0.elapsed().as_secs_f64()); + metrics::histogram!("lightrail_resync_queue_claim_scanned").record(scanned as f64); Ok(None) } -- 2.51.2