From 0a11ed5ba7e371eb2fc48e656c63db647d91d6ed Mon Sep 17 00:00:00 2001 From: Lewis Date: Thu, 07 May 2026 20:05:42 +0000 Subject: [PATCH] feat(ingest): ResolverStats, HydrantStreamErrorFrame, profiling Lewis: May this revision serve well! --- Cargo.lock | 12 ++++++++++++ Cargo.toml | 7 +++++++ containerfiles/bobbin.Containerfile | 8 +++++--- crates/ingest/Cargo.toml | 1 + crates/ingest/src/frame.rs | 7 +++++++ crates/ingest/src/resolver.rs | 224 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------------ crates/runtime/src/entropy.rs | 2 +- crates/xrpc/tests/aggregation.rs | 2 +- crates/xrpc/tests/extended.rs | 6 +++--- crates/xrpc/tests/knot_proxy.rs | 2 +- 10 file(s) changed, 250 insertion(s)(+), 21 deletion(s)(-) diff --git a/Cargo.lock b/Cargo.lock --- a/Cargo.lock +++ b/Cargo.lock @@ -338,6 +338,7 @@ "serde_json", "thiserror 2.0.18", "tokio", + "tokio-stream", "tokio-util", "tracing", "tracing-subscriber", @@ -3747,6 +3748,17 @@ checksum = "1729aa945f29d91ba541258c8df89027d5792d85a8841fb65e8bf0f4ede4ef61" dependencies = [ "rustls", + "tokio", +] + +[[package]] +name = "tokio-stream" +version = "0.1.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "32da49809aab5c3bc678af03902d4ccddea2a87d028d86392a4b1560c6906c70" +dependencies = [ + "futures-core", + "pin-project-lite", "tokio", ] diff --git a/Cargo.toml b/Cargo.toml --- a/Cargo.toml +++ b/Cargo.toml @@ -40,6 +40,7 @@ tokio = { version = "1.52", features = ["macros", "rt-multi-thread", "time", "signal", "io-util", "sync"] } tokio-util = "0.7" +tokio-stream = "0.1" tokio-tungstenite = { version = "0.29", features = ["rustls-tls-webpki-roots"] } futures = "0.3" @@ -87,3 +88,9 @@ [profile.bench] debug = 1 strip = false + +[profile.profiling] +inherits = "release" +debug = 1 +strip = false +lto = "thin" diff --git a/containerfiles/bobbin.Containerfile b/containerfiles/bobbin.Containerfile --- a/containerfiles/bobbin.Containerfile +++ b/containerfiles/bobbin.Containerfile @@ -1,14 +1,16 @@ FROM docker.io/library/rust:1-alpine3.23 AS builder RUN apk add --no-cache build-base musl-dev cmake perl pkgconfig +ARG BOBBIN_PROFILE=release WORKDIR /src COPY . ./ RUN rm -f .cargo/config.toml -RUN cargo build --release --bin bobbin --package bobbin -RUN strip target/release/bobbin +RUN cargo build --profile ${BOBBIN_PROFILE} --bin bobbin --package bobbin +RUN if [ "${BOBBIN_PROFILE}" = "release" ]; then strip target/${BOBBIN_PROFILE}/bobbin; fi FROM docker.io/library/alpine:3.23 +ARG BOBBIN_PROFILE=release RUN apk add --no-cache ca-certificates -COPY --from=builder /src/target/release/bobbin /usr/local/bin/bobbin +COPY --from=builder /src/target/${BOBBIN_PROFILE}/bobbin /usr/local/bin/bobbin ENV BOBBIN_BIND=0.0.0.0:8090 EXPOSE 8090 ENTRYPOINT ["/usr/local/bin/bobbin"] diff --git a/crates/ingest/Cargo.toml b/crates/ingest/Cargo.toml --- a/crates/ingest/Cargo.toml +++ b/crates/ingest/Cargo.toml @@ -20,6 +20,7 @@ serde_json = { workspace = true } thiserror = { workspace = true } tokio = { workspace = true } +tokio-stream = { workspace = true } tokio-util = { workspace = true } tracing = { workspace = true } url = { workspace = true } diff --git a/crates/ingest/src/frame.rs b/crates/ingest/src/frame.rs --- a/crates/ingest/src/frame.rs +++ b/crates/ingest/src/frame.rs @@ -39,6 +39,13 @@ } #[derive(Clone, Debug, Deserialize)] +pub struct HydrantStreamErrorFrame { + pub error: String, + #[serde(default)] + pub message: Option, +} + +#[derive(Clone, Debug, Deserialize)] pub struct RecordFrame { pub live: bool, pub did: Did, diff --git a/crates/ingest/src/resolver.rs b/crates/ingest/src/resolver.rs --- a/crates/ingest/src/resolver.rs +++ b/crates/ingest/src/resolver.rs @@ -1,4 +1,8 @@ -use bobbin_runtime::RuntimeHasher; +use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Duration; + +use bobbin_runtime::{Clock, RuntimeHasher}; use bobbin_slingshot_client::{SlingshotClient, SlingshotError}; use bobbin_types::edges::{ExtractError, Record}; use bobbin_types::ids::nsid_static; @@ -61,24 +65,129 @@ } } +#[derive(Default)] +pub struct ResolverStats { + hits: AtomicU64, + misses_mapped: AtomicU64, + misses_no_repo_did: AtomicU64, + misses_unresolvable: AtomicU64, + misses_transient: AtomicU64, + misses_no_client: AtomicU64, + miss_latency_micros_sum: AtomicU64, + miss_latency_micros_max: AtomicU64, +} + +#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] +pub struct ResolverStatsSnapshot { + pub hits: u64, + pub misses_mapped: u64, + pub misses_no_repo_did: u64, + pub misses_unresolvable: u64, + pub misses_transient: u64, + pub misses_no_client: u64, + pub miss_latency_micros_sum: u64, + pub miss_latency_micros_max: u64, +} + +impl ResolverStatsSnapshot { + pub fn miss_count(&self) -> u64 { + self.misses_mapped + + self.misses_no_repo_did + + self.misses_unresolvable + + self.misses_transient + + self.misses_no_client + } + + pub fn total(&self) -> u64 { + self.hits + self.miss_count() + } + + pub fn miss_latency_micros_avg(&self) -> Option { + let misses = self.miss_count() - self.misses_no_client; + (misses > 0).then(|| self.miss_latency_micros_sum / misses) + } +} + +#[derive(Clone, Copy)] +enum MissKind { + Mapped, + NoRepoDid, + Unresolvable, + Transient, + NoClient, +} + +impl ResolverStats { + fn record_hit(&self) { + self.hits.fetch_add(1, Ordering::Relaxed); + } + + fn record_miss(&self, kind: MissKind, latency: Option) { + let counter = match kind { + MissKind::Mapped => &self.misses_mapped, + MissKind::NoRepoDid => &self.misses_no_repo_did, + MissKind::Unresolvable => &self.misses_unresolvable, + MissKind::Transient => &self.misses_transient, + MissKind::NoClient => &self.misses_no_client, + }; + counter.fetch_add(1, Ordering::Relaxed); + if let Some(latency) = latency { + let micros = u64::try_from(latency.as_micros()).unwrap_or(u64::MAX); + self.miss_latency_micros_sum + .fetch_add(micros, Ordering::Relaxed); + self.miss_latency_micros_max + .fetch_max(micros, Ordering::Relaxed); + } + } + + pub fn snapshot(&self) -> ResolverStatsSnapshot { + ResolverStatsSnapshot { + hits: self.hits.load(Ordering::Relaxed), + misses_mapped: self.misses_mapped.load(Ordering::Relaxed), + misses_no_repo_did: self.misses_no_repo_did.load(Ordering::Relaxed), + misses_unresolvable: self.misses_unresolvable.load(Ordering::Relaxed), + misses_transient: self.misses_transient.load(Ordering::Relaxed), + misses_no_client: self.misses_no_client.load(Ordering::Relaxed), + miss_latency_micros_sum: self.miss_latency_micros_sum.load(Ordering::Relaxed), + miss_latency_micros_max: self.miss_latency_micros_max.load(Ordering::Relaxed), + } + } +} + +struct SlingshotProbe { + client: SlingshotClient, + clock: Arc, +} + pub struct RepoIdResolver { cache: SccMap, - client: Option, + probe: Option, + stats: ResolverStats, } impl RepoIdResolver { - pub fn with_slingshot(client: SlingshotClient, hasher: RuntimeHasher) -> Self { + pub fn with_slingshot( + client: SlingshotClient, + clock: Arc, + hasher: RuntimeHasher, + ) -> Self { Self { cache: SccMap::with_hasher(hasher), - client: Some(client), + probe: Some(SlingshotProbe { client, clock }), + stats: ResolverStats::default(), } } pub fn detached(hasher: RuntimeHasher) -> Self { Self { cache: SccMap::with_hasher(hasher), - client: None, + probe: None, + stats: ResolverStats::default(), } + } + + pub fn stats(&self) -> ResolverStatsSnapshot { + self.stats.snapshot() } pub async fn observe( @@ -112,13 +221,16 @@ pub async fn resolve(&self, owner: &Did, rkey: &Rkey) -> Resolution { let key = RepoRef::new(owner.clone(), rkey.clone()); if let Some(entry) = self.cache.get_async(&key).await { + self.stats.record_hit(); return entry.get().clone().into_resolution(); } - let Some(client) = self.client.as_ref() else { + let Some(probe) = self.probe.as_ref() else { + self.stats.record_miss(MissKind::NoClient, None); return Resolution::Unresolvable; }; + let started = probe.clock.now_instant(); let nsid: Nsid = nsid_static(REPO_COLLECTION); - let provisional = match client.get_record(owner, &nsid, rkey).await { + let provisional = match probe.client.get_record(owner, &nsid, rkey).await { Ok(body) => match repo_did_from_body(&body.value) { Ok(Some(did)) => Resolution::Mapped(did), Ok(None) => Resolution::NoRepoDid, @@ -156,9 +268,19 @@ rkey = rkey.as_ref(), "slingshot transient failure during repoDID lookup, will retry", ); + let elapsed = probe.clock.now_instant().duration_since(started); + self.stats + .record_miss(MissKind::Transient, Some(elapsed)); return Resolution::Unresolvable; } }; + let elapsed = probe.clock.now_instant().duration_since(started); + let kind = match &provisional { + Resolution::Mapped(_) => MissKind::Mapped, + Resolution::NoRepoDid => MissKind::NoRepoDid, + Resolution::Unresolvable => MissKind::Unresolvable, + }; + self.stats.record_miss(kind, Some(elapsed)); self.fill_provisional(key, provisional.clone()).await; provisional } @@ -185,6 +307,7 @@ #[cfg(test)] mod tests { use super::*; + use bobbin_runtime::SystemClock; use jacquard_common::types::did::Did; use jacquard_common::types::recordkey::Rkey; @@ -194,6 +317,10 @@ fn rkey(s: &str) -> Rkey { Rkey::new_owned(s).unwrap() + } + + fn test_clock() -> Arc { + Arc::new(SystemClock::new()) } #[tokio::test] @@ -314,7 +441,7 @@ let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); - let resolver = RepoIdResolver::with_slingshot(client, RuntimeHasher::default()); + let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); @@ -340,7 +467,7 @@ let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); - let resolver = RepoIdResolver::with_slingshot(client, RuntimeHasher::default()); + let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); @@ -371,7 +498,7 @@ let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); - let resolver = RepoIdResolver::with_slingshot(client, RuntimeHasher::default()); + let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); @@ -398,7 +525,7 @@ let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); - let resolver = RepoIdResolver::with_slingshot(client, RuntimeHasher::default()); + let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); @@ -424,7 +551,7 @@ let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); - let resolver = RepoIdResolver::with_slingshot(client, RuntimeHasher::default()); + let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); @@ -432,6 +559,79 @@ let second = resolver.resolve(&owner, &key).await; assert_eq!(first, Resolution::Unresolvable); assert_eq!(second, Resolution::Unresolvable); + } + + #[tokio::test] + async fn slingshot_transient_recorded_separately_from_unresolvable() { + let server = wiremock::MockServer::start().await; + wiremock::Mock::given(wiremock::matchers::method("GET")) + .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) + .respond_with(wiremock::ResponseTemplate::new(503)) + .expect(2) + .mount(&server) + .await; + + let client = + SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); + let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); + + let owner = did("did:plc:nel"); + let key = rkey("abcabcabcabcz"); + resolver.resolve(&owner, &key).await; + resolver.resolve(&owner, &key).await; + + let snap = resolver.stats(); + assert_eq!( + snap.misses_transient, 2, + "transport-error retries must count as transient, not as canonical unresolvable", + ); + assert_eq!( + snap.misses_unresolvable, 0, + "canonical unresolvable counter is reserved for cached terminal answers", + ); + assert!( + snap.miss_latency_micros_sum > 0, + "transient misses still have latency contributions", + ); + assert_eq!(snap.miss_count(), 2); + } + + #[tokio::test] + async fn stats_count_hits_misses_and_latency() { + let server = wiremock::MockServer::start().await; + wiremock::Mock::given(wiremock::matchers::method("GET")) + .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) + .respond_with(wiremock::ResponseTemplate::new(404)) + .mount(&server) + .await; + let client = + SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); + let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); + + let owner = did("did:plc:nel"); + let key = rkey("abcabcabcabcz"); + resolver.resolve(&owner, &key).await; + resolver.resolve(&owner, &key).await; + + let snap = resolver.stats(); + assert_eq!(snap.misses_unresolvable, 1, "first call is the slingshot miss"); + assert_eq!(snap.hits, 1, "second call hits the unresolvable cache"); + assert_eq!(snap.miss_count(), 1); + assert_eq!(snap.total(), 2); + assert!(snap.miss_latency_micros_sum > 0, "latency recorded for slingshot miss"); + assert!(snap.miss_latency_micros_avg().unwrap() > 0); + } + + #[tokio::test] + async fn stats_no_client_miss_recorded_without_latency() { + let resolver = RepoIdResolver::detached(RuntimeHasher::default()); + resolver + .resolve(&did("did:plc:nel"), &rkey("abcabcabcabcz")) + .await; + let snap = resolver.stats(); + assert_eq!(snap.misses_no_client, 1); + assert_eq!(snap.miss_latency_micros_sum, 0); + assert_eq!(snap.miss_latency_micros_avg(), None); } #[tokio::test] diff --git a/crates/runtime/src/entropy.rs b/crates/runtime/src/entropy.rs --- a/crates/runtime/src/entropy.rs +++ b/crates/runtime/src/entropy.rs @@ -41,7 +41,7 @@ let b = e.next_u64(); assert_ne!( a, b, - "back-to-back os entropy collided; getrandom likely not actually wired up", + "back-to-back os entropy collided, getrandom likely not actually wired up", ); } } diff --git a/crates/xrpc/tests/aggregation.rs b/crates/xrpc/tests/aggregation.rs --- a/crates/xrpc/tests/aggregation.rs +++ b/crates/xrpc/tests/aggregation.rs @@ -1188,7 +1188,7 @@ assert_eq!( status, StatusCode::OK, - "extractor key must match handler subject; body: {json}", + "extractor key must match handler subject, body was {json}", ); let items = json["items"].as_array().unwrap(); assert_eq!(items.len(), 1, "expected exactly one star edge"); diff --git a/crates/xrpc/tests/extended.rs b/crates/xrpc/tests/extended.rs --- a/crates/xrpc/tests/extended.rs +++ b/crates/xrpc/tests/extended.rs @@ -719,7 +719,7 @@ assert_eq!( status, StatusCode::OK, - "extractor key must match handler subject; body: {json}" + "extractor key must match handler subject, body was {json}" ); let items = json["items"].as_array().unwrap(); assert_eq!(items.len(), 1, "expected exactly one pipeline edge"); @@ -758,7 +758,7 @@ assert_eq!( status, StatusCode::OK, - "owner-DID fallback must hydrate; body: {json}" + "owner-DID fallback must hydrate, body was {json}" ); let items = json["items"].as_array().unwrap(); assert_eq!(items.len(), 1); @@ -770,7 +770,7 @@ items[0]["value"]["triggerMetadata"]["repo"] .get("repoDid") .is_none(), - "fixture must omit repoDid; got {}", + "fixture must omit repoDid but got {}", items[0]["value"]["triggerMetadata"]["repo"] ); } diff --git a/crates/xrpc/tests/knot_proxy.rs b/crates/xrpc/tests/knot_proxy.rs --- a/crates/xrpc/tests/knot_proxy.rs +++ b/crates/xrpc/tests/knot_proxy.rs @@ -392,7 +392,7 @@ .count(); assert_eq!( getrecord, 1, - "slingshot must be hit exactly once; LRU serves the second proxy call", + "slingshot must be hit exactly once because the LRU serves the second proxy call", ); } -- tangled.sh