diff --git a/hubble-sync/src/lib.rs b/hubble-sync/src/lib.rs index b42b75b..d132491 100644 --- a/hubble-sync/src/lib.rs +++ b/hubble-sync/src/lib.rs @@ -36,7 +36,8 @@ pub use hubble_sync::{HubbleSync, HubbleSyncError}; pub use identity::{Did, DidMethod, HubbleSyncResolver, Resolve, SigningKey, did}; pub use metrics::{ BIG_REPO_WAIT_SECONDS_BUCKETS, REPO_ACTOR_ALIVE_SECONDS_BUCKETS, RESYNC_COUNT_BUCKETS, - RESYNC_DURATION_SECONDS_BUCKETS, RESYNC_SIZE_BUCKETS, describe_metrics, histogram_buckets, + RESYNC_DURATION_SECONDS_BUCKETS, RESYNC_PHASE_SECONDS_BUCKETS, RESYNC_SIZE_BUCKETS, + describe_metrics, histogram_buckets, }; pub use pending_identity_scheduler::PendingScheduler; pub use repo_actor::{ diff --git a/hubble-sync/src/metrics.rs b/hubble-sync/src/metrics.rs index 930b9d4..f68f01b 100644 --- a/hubble-sync/src/metrics.rs +++ b/hubble-sync/src/metrics.rs @@ -138,6 +138,12 @@ pub(crate) const RESYNC_COUNT: &str = "hubble_sync_resync_count"; /// wall-clock seconds for a completed resync (fetch + load + apply) pub(crate) const RESYNC_DURATION_SECONDS: &str = "hubble_sync_resync_duration_seconds"; +/// wall-clock seconds for one phase of a resync fetch, to separate host-permit +/// queuing from stream draining. +/// `phase` = `request` (getRepo: host-permit acquire/wait + send + response +/// headers) | `drain` (streaming the response body into a loaded repo) +pub(crate) const RESYNC_PHASE_SECONDS: &str = "hubble_sync_resync_phase_seconds"; + /// seconds spent waiting to acquire a big-repo permit during resync. /// `path` = `preemptive` (predicted big up front) | `reactive` (small load overflowed) pub(crate) const BIG_REPO_WAIT_SECONDS: &str = "hubble_sync_big_repo_wait_seconds"; @@ -440,6 +446,10 @@ pub fn describe_metrics() { RESYNC_DURATION_SECONDS, "seconds for a completed resync (fetch+load+apply)" ); + describe_histogram!( + RESYNC_PHASE_SECONDS, + "seconds for a resync fetch phase, phase=request|drain" + ); describe_histogram!( BIG_REPO_WAIT_SECONDS, "seconds waiting for a big-repo permit, path=preemptive|reactive" @@ -510,6 +520,12 @@ pub const RESYNC_DURATION_SECONDS_BUCKETS: &[f64] = &[ 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, 30.0, 60.0, 120.0, 300.0, ]; +/// for [`RESYNC_PHASE_SECONDS`] -- fine around the 0.3s (fast drain) vs ~1s +/// (queuing/slow drain) split we're trying to distinguish +pub const RESYNC_PHASE_SECONDS_BUCKETS: &[f64] = &[ + 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, 30.0, 60.0, 300.0, +]; + /// for [`REPO_ACTOR_ALIVE_SECONDS`] pub const REPO_ACTOR_ALIVE_SECONDS_BUCKETS: &[f64] = &[ 1.0, 5.0, 30.0, 60.0, 300.0, 1_800.0, 3_600.0, 21_600.0, 86_400.0, @@ -542,6 +558,7 @@ pub fn histogram_buckets() -> &'static [(&'static str, &'static [f64])] { (RESYNC_SIZE, RESYNC_SIZE_BUCKETS), (RESYNC_COUNT, RESYNC_COUNT_BUCKETS), (RESYNC_DURATION_SECONDS, RESYNC_DURATION_SECONDS_BUCKETS), + (RESYNC_PHASE_SECONDS, RESYNC_PHASE_SECONDS_BUCKETS), (REPO_ACTOR_ALIVE_SECONDS, REPO_ACTOR_ALIVE_SECONDS_BUCKETS), (BIG_REPO_WAIT_SECONDS, BIG_REPO_WAIT_SECONDS_BUCKETS), ] diff --git a/hubble-sync/src/repo_actor/task_processor/mod.rs b/hubble-sync/src/repo_actor/task_processor/mod.rs index c6ed1ec..e2caf0a 100644 --- a/hubble-sync/src/repo_actor/task_processor/mod.rs +++ b/hubble-sync/src/repo_actor/task_processor/mod.rs @@ -22,7 +22,8 @@ use crate::host::{HostRequestError, ResolvedHost}; use crate::identity::{ResolutionError, Resolve, ResolvedIdentity, Validity}; use crate::metrics::{ ACCOUNT_OUTCOMES_TOTAL, COMMIT_OUTCOMES_TOTAL, IDENTITY_REFRESH_OUTCOMES_TOTAL, RESYNC_COUNT, - RESYNC_DURATION_SECONDS, RESYNC_OUTCOMES_TOTAL, RESYNC_SIZE, SYNC_OUTCOMES_TOTAL, + RESYNC_DURATION_SECONDS, RESYNC_OUTCOMES_TOTAL, RESYNC_PHASE_SECONDS, RESYNC_SIZE, + SYNC_OUTCOMES_TOTAL, }; use crate::resync::{ BigRepoPermits, ResyncData, ResyncError, Resyncable, TransientResyncError, load_repo, @@ -173,6 +174,7 @@ impl<'a> ResyncContext<'a> { }; // set up the initial request + let request_start = Instant::now(); let (reader, extensions) = target .get_repo(self.repo.did()) @@ -204,6 +206,8 @@ impl<'a> ResyncContext<'a> { TransientResyncError::Request(e.to_string()).into() } })?; + histogram!(RESYNC_PHASE_SECONDS, "phase" => "request") + .record(request_start.elapsed().as_secs_f64()); self.resolved = extensions.get::().map(|r| r.0.clone()); @@ -211,7 +215,9 @@ impl<'a> ResyncContext<'a> { let permit = self.permit.take(); // actually load the repo - self.cancel + let drain_start = Instant::now(); + let loaded = self + .cancel .timeout( STREAM_CAR_TIMEOUT, load_repo( @@ -224,7 +230,10 @@ impl<'a> ResyncContext<'a> { ) .await .ok_or(ResyncError::Cancelled)? - .map_err(TransientResyncError::Timeout)? + .map_err(TransientResyncError::Timeout)?; + histogram!(RESYNC_PHASE_SECONDS, "phase" => "drain") + .record(drain_start.elapsed().as_secs_f64()); + loaded } }