From ac86de738daddec88e2aedc6dbbc6970121c516b Mon Sep 17 00:00:00 2001 From: phil Date: Tue, 21 Jul 2026 15:40:15 -0400 Subject: [PATCH] resolution metrics + plc concurrency config --- .../src/crawl_strategy_upstream_repos.rs | 11 ++- hubble-sync/src/identity/did.rs | 9 +++ hubble-sync/src/identity/resolver.rs | 74 +++++++++++++------ hubble-sync/src/metrics.rs | 33 +++++++++ hubble-sync/src/pending_identity_scheduler.rs | 42 +++++++---- hubble/src/main.rs | 6 ++ 6 files changed, 137 insertions(+), 38 deletions(-) diff --git a/hubble-sync/src/crawl_strategy_upstream_repos.rs b/hubble-sync/src/crawl_strategy_upstream_repos.rs index 51765cd..55b93fb 100644 --- a/hubble-sync/src/crawl_strategy_upstream_repos.rs +++ b/hubble-sync/src/crawl_strategy_upstream_repos.rs @@ -33,10 +33,10 @@ pub async fn crawl_upstream_listrepos( discovery: RepoDiscovery, state: CrawlState, upstream: Arc, - recrawl: Duration, + recrawl_interval: Duration, cancel: CancellationToken, ) -> Result<(), CrawlError> { - let mut recrawl = interval(recrawl); + let mut recrawl = interval(recrawl_interval); recrawl.set_missed_tick_behavior(MissedTickBehavior::Delay); // startup: resume from a previous listRepos page, if any @@ -52,13 +52,20 @@ pub async fn crawl_upstream_listrepos( loop { // once per ~day (recrawl duration) + debug!( + ?upstream, + ?recrawl_interval, + "discover crawl outer loop about to tick" + ); let Some(_) = cancel.run(recrawl.tick()).await else { return Ok(()); }; + debug!(?upstream, ?recrawl_interval, "discover crawl proceeding"); // loop over pages for batches of repos let mut i = 1; 'paging: loop { + debug!(?upstream, ?cursor, "fetching page"); let Some(maybe_page) = cancel .run(list_repos_page(&upstream, cursor.as_deref())) .await diff --git a/hubble-sync/src/identity/did.rs b/hubble-sync/src/identity/did.rs index 7662a72..6bcde84 100644 --- a/hubble-sync/src/identity/did.rs +++ b/hubble-sync/src/identity/did.rs @@ -21,6 +21,15 @@ pub enum DidMethod { Web, } +impl DidMethod { + pub fn name(&self) -> &'static str { + match self { + Self::Plc => "plc", + Self::Web => "web", + } + } +} + #[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)] pub struct Did(Arc); diff --git a/hubble-sync/src/identity/resolver.rs b/hubble-sync/src/identity/resolver.rs index 22e4284..ba6125c 100644 --- a/hubble-sync/src/identity/resolver.rs +++ b/hubble-sync/src/identity/resolver.rs @@ -15,10 +15,15 @@ use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant, SystemTime}; +use metrics::{counter, histogram}; use reqwest::{StatusCode, Url, header}; use tracing::debug; use super::doc; +use crate::metrics::{ + IDENTITY_RESOLVE_DURATION_SECONDS, IDENTITY_RESOLVE_OUTCOMES_TOTAL, + IDENTITY_RESOLVE_REQUESTS_TOTAL, +}; use crate::storage::repo::RepoIdentity; use crate::{Did, DidMethod, HostRegistry}; @@ -117,6 +122,34 @@ impl HubbleSyncResolver { Some(Duration::from_secs(n)) } + /// metrics-wrapped actual fetch + async fn fetch_doc( + client: &reqwest::Client, + url: Url, + method: &DidMethod, + ) -> Result, ResolutionError> { + let start = Instant::now(); + let res = client.get(url).send().await; + + let (out, status) = match res { + Ok(r) => { + let status = r.status().as_str().to_string(); + (Self::get_bytes(r).await, status) + } + Err(err) => { + let o = Err(ResolutionError::Network(err.to_string())); + (o, "transport".to_string()) + } + }; + + counter!(IDENTITY_RESOLVE_REQUESTS_TOTAL, "method" => method.name(), "status" => status) + .increment(1); + histogram!(IDENTITY_RESOLVE_DURATION_SECONDS, "method" => method.name()) + .record(start.elapsed().as_secs_f64()); + + out + } + async fn get_bytes(mut resp: reqwest::Response) -> Result, ResolutionError> { let status = resp.status(); @@ -157,16 +190,8 @@ impl HubbleSyncResolver { .parse() .map_err(|e| ResolutionError::Unresolvable(format!("invalid in URL: {did:?}: {e}")))?; - let resp = self - .plc_client - .get(url) - .send() - .await - .map_err(|e| ResolutionError::Network(e.to_string()))?; - - let raw = Self::get_bytes(resp).await?; - let identity = doc::parse(raw, did, &self.hosts, now)?; - Ok(identity) + let raw = Self::fetch_doc(&self.plc_client, url, &did.method()).await?; + Ok(doc::parse(raw, did, &self.hosts, now)?) } async fn resolve_web(&self, did: &Did, now: SystemTime) -> ResolvedIdentity { @@ -203,16 +228,8 @@ impl HubbleSyncResolver { .parse() .map_err(|e| ResolutionError::Unresolvable(format!("invalid in URL: {did:?}: {e}")))?; - let resp = self - .web_client - .get(url) - .send() - .await - .map_err(|e| ResolutionError::Network(e.to_string()))?; - - let raw = Self::get_bytes(resp).await?; - let identity = doc::parse(raw, did, &self.hosts, now)?; - Ok(identity) + let raw = Self::fetch_doc(&self.web_client, url, &did.method()).await?; + Ok(doc::parse(raw, did, &self.hosts, now)?) } async fn plc_limit(&self) { @@ -241,7 +258,7 @@ impl HubbleSyncResolver { impl Resolve for HubbleSyncResolver { async fn resolve(&self, did: &Did, now: SystemTime) -> ResolvedIdentity { - match did.method() { + let res = match did.method() { DidMethod::Plc => { self.plc_limit().await; let res = self.resolve_plc(did, now).await; @@ -270,6 +287,19 @@ impl Resolve for HubbleSyncResolver { res } - } + }; + + use ResolutionError as RE; + let outcome: &'static str = match &res { + Ok(_) => "ok", + Err(RE::RateLimited { .. } | RE::TooSoon(_)) => "rate_limited", + Err(RE::NotFound) => "not_found", + Err(RE::Network(_) | RE::Body(_) | RE::DidDocError(_)) => "request", + Err(RE::Unresolvable(_)) => "unresolvable", + }; + counter!(IDENTITY_RESOLVE_OUTCOMES_TOTAL, "method" => did.method().name(), "outcome" => outcome) + .increment(1); + + res } } diff --git a/hubble-sync/src/metrics.rs b/hubble-sync/src/metrics.rs index dabe571..3133d80 100644 --- a/hubble-sync/src/metrics.rs +++ b/hubble-sync/src/metrics.rs @@ -51,6 +51,18 @@ pub(crate) const INITIAL_RESOLVE_OUTCOMES_TOTAL: &str = pub(crate) const IDENTITY_REFRESH_OUTCOMES_TOTAL: &str = "hubble_sync_identity_refresh_outcomes_total"; +/// actual requests made to resolve identity. `did_method` = `plc` | `web`; `status` = | `transport +pub(crate) const IDENTITY_RESOLVE_REQUESTS_TOTAL: &str = + "hubble_sync_identity_resolve_requests_total"; + +/// time to actually resolve identities by did method. `did_method` = `plc` | `web` +pub(crate) const IDENTITY_RESOLVE_DURATION_SECONDS: &str = + "hubble_sync_identity_resolve_duration_seconds"; + +/// identity resolution outcomes by did method and result. `did_method` = `plc` | `web`; `result` = `ok` | `not_found` | `rate_limited` | `request_error` | `unresolveable` +pub(crate) const IDENTITY_RESOLVE_OUTCOMES_TOTAL: &str = + "hubble_sync_identity_resolve_outcomes_total"; + // ///// host registry / host state /// hostname->Host lookups. `result` = `hit` | `miss` @@ -264,6 +276,18 @@ pub fn describe_metrics() { IDENTITY_REFRESH_OUTCOMES_TOTAL, "identity refreshes, outcome=refreshed|snoozed|failed|storage_error, trigger=wake|sig_retry|firehose_identity|resync_repo_missing|requested" ); + describe_counter!( + IDENTITY_RESOLVE_REQUESTS_TOTAL, + "actual requests made to resolve identity, did_method=plc|web, status=|transport" + ); + describe_histogram!( + IDENTITY_RESOLVE_DURATION_SECONDS, + "time to actually resolve identities by did method, did_method=plc|web" + ); + describe_counter!( + IDENTITY_RESOLVE_OUTCOMES_TOTAL, + "identity resolution outcomes by did method and result, did_method=plc|web; result=ok|not_found|rate_limited|request_error|unresolveable" + ); // host registry / host state describe_counter!( @@ -437,6 +461,11 @@ pub fn describe_metrics() { // ///// recommended histogram buckets +/// for [`IDENTITY_RESOLVE_DURATION_SECONDS`] +pub const IDENTITY_RESOLVE_DURATION_SECONDS_BUCKETS: &[f64] = &[ + 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, +]; + /// for [`COMMIT_PREVALIDATE_SECONDS`] pub const COMMIT_PREVALIDATE_SECONDS_BUCKETS: &[f64] = &[0.0001, 0.0005, 0.001, 0.005, 0.01, 0.05, 0.1, 0.5, 1.0]; @@ -495,6 +524,10 @@ pub const BIG_REPO_WAIT_SECONDS_BUCKETS: &[f64] = /// ``` pub fn histogram_buckets() -> &'static [(&'static str, &'static [f64])] { &[ + ( + IDENTITY_RESOLVE_DURATION_SECONDS, + IDENTITY_RESOLVE_DURATION_SECONDS_BUCKETS, + ), ( COMMIT_PREVALIDATE_SECONDS, COMMIT_PREVALIDATE_SECONDS_BUCKETS, diff --git a/hubble-sync/src/pending_identity_scheduler.rs b/hubble-sync/src/pending_identity_scheduler.rs index 7fd16ef..6082c64 100644 --- a/hubble-sync/src/pending_identity_scheduler.rs +++ b/hubble-sync/src/pending_identity_scheduler.rs @@ -38,8 +38,8 @@ use crate::metrics::{ use crate::repo_actor::{RepoMessage, RepoRegistry}; use crate::storage::repo::PendingIdentityQueueEntry; use crate::{ - CancelExt, Did, LoadError, PrefixedEngine, RepoSendError, Resolve, StorageBatch, StorageEngine, - SyncConsumer, + CancelExt, Did, DidMethod, LoadError, PrefixedEngine, RepoSendError, Resolve, StorageBatch, + StorageEngine, SyncConsumer, }; const IDLE_POLL: Duration = Duration::from_secs(15); @@ -319,7 +319,7 @@ pub fn bootstrap( /// TODO: probably goes on PendingScheduler impl pub async fn drive, R: Resolve>( scheduler: Arc, - resolve_limit: Arc, + plc_resolve_limit: Arc, storage: PrefixedEngine, repo_registry: Arc>, cancel: CancellationToken, @@ -342,19 +342,33 @@ pub async fn drive, R: Resolve>( } let now = SystemTime::now(); while let Some(did) = scheduler.next(now) { - // bound concurrency on resolves we dispatch - let Some(permit) = cancel - .run(resolve_limit.clone().acquire_owned()) - .await - .transpose() - .expect("semaphore not closed") - else { - return Ok(()); // cancelled + let resolve_permit = if matches!(did.method(), DidMethod::Plc) { + // concurrency bound on PLC throttles our requests to it + // + // we probably should have an actual rate-limit as well, since + // this couples backfill to plc RTT + // + // either way, we want to leave headroom for firehose-driven + // resolution, which isn't subject to the limit. + // + // and: for now, we just wait here. we could be blocking did:web + // resolutions from proceeding, but _for now_ that's a low- + // enough impact with acceptably small effect. + let Some(permit) = cancel + .run(plc_resolve_limit.clone().acquire_owned()) + .await + .transpose() + .expect("semaphore not closed") + else { + return Ok(()); // bail: cancelled + }; + Some(permit) + } else { + // no bound on did-web concurrency for now (+not shared with plc limit) + None }; - let message = RepoMessage::ScheduledIdentityResolve { - resolve_permit: Some(permit), - }; + let message = RepoMessage::ScheduledIdentityResolve { resolve_permit }; if let Err(e) = repo_registry.try_send(&did, message) { if matches!(e, RepoSendError::Draining(_)) { // seems we're shutting down diff --git a/hubble/src/main.rs b/hubble/src/main.rs index 4240072..9a6397d 100644 --- a/hubble/src/main.rs +++ b/hubble/src/main.rs @@ -1,4 +1,5 @@ use std::net::SocketAddr; +use std::num::NonZeroUsize; use std::path::{Path, PathBuf}; use std::sync::Arc; use std::time::Duration; @@ -106,6 +107,10 @@ struct Args { #[arg(long, action)] relay_bsky: bool, + /// how many scheduled plc resolution requests can be in flight + #[arg(long, default_value = "6")] + max_plc_scheduled_requests: NonZeroUsize, + /// http server listen #[arg(long, default_value = "0.0.0.0:3000")] bind: SocketAddr, @@ -178,6 +183,7 @@ async fn main() -> color_eyre::Result<()> { .max_big_repo_resync_concurrency .unwrap_or(SyncConfig::default().big_repo_resync_limit), huge_repo_spill_dir: resync_spill_dir, + scheduled_resolve_limit: args.max_plc_scheduled_requests, upstream: UpstreamConfig { hostname: args.upstream.clone(), ..Default::default() -- 2.51.2