diff --git a/hubble-sync/src/resync_scheduler.rs b/hubble-sync/src/resync_scheduler.rs index 0d6ad3c..76e2029 100644 --- a/hubble-sync/src/resync_scheduler.rs +++ b/hubble-sync/src/resync_scheduler.rs @@ -75,6 +75,8 @@ const DISPATCH_LAG_MIN: Duration = Duration::from_millis(4); const CURSOR_DROP_INTERVAL: Duration = Duration::from_secs(4); /// all cursors dropped if this is somehow reached between normal drop intervals const MAX_CURSORS: usize = 1_000_000; +/// queued resyncs per host to prefetch at once +const PREFETCH_SIZE: usize = 64; #[derive(Debug, thiserror::Error)] pub enum ResyncDriveError { @@ -381,13 +383,16 @@ pub async fn drive, R: Resolve>( RateLimiter::direct(Quota::per_second(dispatch_qps)); let concurrency_limit = Arc::new(Semaphore::new(concurrency_limit)); + let mut prefetch_buffer: HashMap, VecDeque> = HashMap::new(); + let mut last_rebootstrap = Instant::now(); let mut last_cursors_drop = Instant::now(); while !cancel.is_cancelled() { if last_cursors_drop.elapsed() >= CURSOR_DROP_INTERVAL { - trace!("droping cursors to catch stragglers"); + trace!("droping cursors and prefetch buffers to catch stragglers"); scheduler.drop_cursors(); + prefetch_buffer.clear(); last_cursors_drop = Instant::now(); } rebootstrap_if_due(&scheduler, &storage, &host_registry, &mut last_rebootstrap).await?; @@ -464,16 +469,37 @@ pub async fn drive, R: Resolve>( // yay. update host cursor, get its next thing or unschedule it let host_label = if mhost.is_some() { "known" } else { "unknown" }; counter!(RESYNC_SCHEDULER_DISPATCHED_TOTAL, "host" => host_label).increment(1); - let peek_storage = storage.clone(); scheduler.record_cursor(&mhost, item.due, item.did.clone()); - if let Some(host_next) = - tokio::task::spawn_blocking(move || item.peek_host_next(&peek_storage)) - .await - .expect("storage not to panic")? + + let host_key = mhost.as_ref().map(|h| h.name().clone()); + + // get the next task for this host-- + // if there's nothing prefetched, refresh its prefetch buffer first + if prefetch_buffer.get(&host_key).is_none_or(|b| b.is_empty()) { + let peek_storage = storage.clone(); + let host_batch = tokio::task::spawn_blocking(move || { + item.peek_host_next_batch(PREFETCH_SIZE, &peek_storage) + }) + .await + .expect("storage not to panic")?; + if !host_batch.is_empty() { + prefetch_buffer + .entry(host_key.clone()) + .or_default() + .extend(host_batch); + } + } + if let Some(next) = prefetch_buffer + .get_mut(&host_key) + .and_then(|b| b.pop_front()) { - scheduler.insert(host_next); + scheduler.insert(next) } else { - trace!(host_label, "no next due for host"); + trace!( + host_label, + "no next due for host, unscheduling and cleaning up prefetch" + ); + prefetch_buffer.remove(&host_key); scheduler.unschedule_host(&mhost); } } diff --git a/hubble-sync/src/storage/repo/info_idx_resync.rs b/hubble-sync/src/storage/repo/info_idx_resync.rs index 12848bd..90ea44c 100644 --- a/hubble-sync/src/storage/repo/info_idx_resync.rs +++ b/hubble-sync/src/storage/repo/info_idx_resync.rs @@ -235,6 +235,23 @@ impl QueuedResync { let entry = decode_host_suffix(self.host.clone(), &k)?; Ok(Some(entry)) } + + /// Peek the host's next n queued resyncs strictly after this one + pub fn peek_host_next_batch( + &self, + n: usize, + storage: &S, + ) -> Result, LoadError> { + let prefix = build_host_prefix(self.host.as_deref()); + let suffix = build_host_next_suffix(&self.did, self.due()); + let mut out = Vec::with_capacity(n); + for item in storage.scan_from_queue(&prefix, &suffix).take(n) { + let (k, _v) = item.map_err(LoadError::Storage)?; + let entry = decode_host_suffix(self.host.clone(), &k)?; + out.push(entry); + } + Ok(out) + } } /// Iterator to get one QueuedResync per host with items in its queue