diff --git a/Cargo.lock b/Cargo.lock index ec5a09f2..de0b1e25 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -680,6 +680,7 @@ dependencies = [ "bobbin-knot-ingest", "bobbin-knot-proxy", "bobbin-record-lru", + "bobbin-resolver", "bobbin-runtime", "bobbin-search", "bobbin-slingshot-client", @@ -810,7 +811,9 @@ dependencies = [ "bobbin-types", "jacquard-common", "scc", + "serde", "serde_json", + "thiserror 2.0.18", "tokio", "tracing", "url", diff --git a/bobbin/crates/bobbin-sim/src/runtime.rs b/bobbin/crates/bobbin-sim/src/runtime.rs index fd229047..9967cb1f 100644 --- a/bobbin/crates/bobbin-sim/src/runtime.rs +++ b/bobbin/crates/bobbin-sim/src/runtime.rs @@ -9,6 +9,7 @@ use bobbin_ingest::{ WarmingBuffer, WarmingShadowBuffer, run as run_ingest, }; use bobbin_record_lru::NoopRecordStore; +use bobbin_resolver::IdentityResolver; use bobbin_runtime::{ Clock, DEFAULT_MEM_WS_CAPACITY, MemHttpTransport, MemWsTransport, RuntimeHasher, SeededEntropy, SimClock, UnixMicros, @@ -115,6 +116,10 @@ impl Sim { search: Arc::new(NoopSearchSink), records: records.clone() as Arc, resolver: resolver.clone(), + identity: Arc::new(IdentityResolver::detached( + hasher.clone(), + bobbin_resolver::DEFAULT_IDENTITY_CACHE_ENTRIES, + )), clock: clock.clone(), entropy: entropy.clone(), ws: mem_ws, diff --git a/bobbin/crates/bobbin/Cargo.toml b/bobbin/crates/bobbin/Cargo.toml index 4c3f959a..5fffa1d5 100644 --- a/bobbin/crates/bobbin/Cargo.toml +++ b/bobbin/crates/bobbin/Cargo.toml @@ -15,6 +15,7 @@ bobbin-ingest = { workspace = true } bobbin-knot-ingest = { workspace = true } bobbin-knot-proxy = { workspace = true } bobbin-record-lru = { workspace = true } +bobbin-resolver = { workspace = true } bobbin-runtime = { workspace = true } bobbin-search = { workspace = true } bobbin-slingshot-client = { workspace = true } diff --git a/bobbin/crates/bobbin/src/config.rs b/bobbin/crates/bobbin/src/config.rs index 137eff33..57646baa 100644 --- a/bobbin/crates/bobbin/src/config.rs +++ b/bobbin/crates/bobbin/src/config.rs @@ -26,6 +26,7 @@ const KNOWN_KEYS: &[&str] = &[ "backpressure.reserved_index_bytes", "slingshot.url", "record_cache.lru_bytes", + "identity_cache.max_entries", "search.heap_bytes", "knot.allow_private", "knot.require_https", @@ -50,6 +51,7 @@ const KNOWN_ENVS: &[&str] = &[ "BOBBIN_BACKPRESSURE_RESERVED_INDEX_BYTES", "BOBBIN_SLINGSHOT_URL", "BOBBIN_RECORD_LRU_BYTES", + "BOBBIN_IDENTITY_CACHE_ENTRIES", "BOBBIN_SEARCH_HEAP_BYTES", "BOBBIN_KNOT_ALLOW_PRIVATE", "BOBBIN_KNOT_REQUIRE_HTTPS", @@ -78,6 +80,9 @@ pub struct BobbinConfig { #[config(nested)] pub record_cache: RecordCacheConfig, + #[config(nested)] + pub identity_cache: IdentityCacheConfig, + #[config(nested)] pub search: SearchConfig, @@ -236,6 +241,13 @@ pub struct RecordCacheConfig { pub lru_bytes: u64, } +#[derive(Debug, Config)] +pub struct IdentityCacheConfig { + /// Maximum number of (mini) DID doc entries that the identity cache will hold. + #[config(env = "BOBBIN_IDENTITY_CACHE_ENTRIES", default = 100_000)] + pub max_entries: usize, +} + #[derive(Debug, Config)] pub struct SearchConfig { /// The heap size in bytes for the in-mem tantivy writer. Larger values diff --git a/bobbin/crates/bobbin/src/main.rs b/bobbin/crates/bobbin/src/main.rs index 56de2572..789d7c7d 100644 --- a/bobbin/crates/bobbin/src/main.rs +++ b/bobbin/crates/bobbin/src/main.rs @@ -13,6 +13,7 @@ use bobbin_ingest::{ use bobbin_knot_ingest::{CapabilityGate, KnotClient, KnotRegistry, Orchestrator}; use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig, MirrorProxy, classify_ip}; use bobbin_record_lru::{CacheCapacity, LruRecordStore, RecordStore}; +use bobbin_resolver::IdentityResolver; use bobbin_runtime::{ Clock, GuardedWs, MemoryBudget, NetworkError, OsEntropy, ReqwestHttp, RuntimeHasher, SystemClock, TungsteniteWs, WsTransport, @@ -186,6 +187,11 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { let records: Arc = Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(lru_cap))); let slingshot = SlingshotClient::with_default_http(cfg.slingshot.url.clone())?; + let identity = Arc::new(IdentityResolver::with_slingshot( + slingshot.clone(), + hasher.clone(), + cfg.identity_cache.max_entries, + )); let mut resolver_opts = ResolverOptions::default(); // NOTE: see https://tangled.org/nonbinary.computer/jacquard/issues/39. resolver_opts.did_order = vec![ @@ -287,6 +293,7 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { search: search.clone(), records: records.clone(), resolver: resolver.clone(), + identity: identity.clone(), clock: clock.clone(), entropy, ws: ws.clone(), @@ -342,6 +349,7 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { edges: edges.clone(), search: search.clone(), records: records.clone(), + identity: identity.clone(), issue_states: issue_states.clone(), pull_statuses: pull_statuses.clone(), }); @@ -357,6 +365,7 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { resolver, directory, ) + .with_identity(identity) .with_limiter(limiter) .with_mirror(mirror) .with_proxies(trusted_proxies); diff --git a/bobbin/crates/bobbin/src/mem/report.rs b/bobbin/crates/bobbin/src/mem/report.rs index ed4f57d2..aae5f53b 100644 --- a/bobbin/crates/bobbin/src/mem/report.rs +++ b/bobbin/crates/bobbin/src/mem/report.rs @@ -6,6 +6,7 @@ use axum::routing::get; use axum::{Json, Router}; use bobbin_edge_index::{EdgeStore, IssueStateKind, PullStatusKind, StateIndex}; use bobbin_record_lru::RecordStore; +use bobbin_resolver::IdentityResolver; use bobbin_search::SearchIndex; use serde::Serialize; use tikv_jemalloc_ctl::{epoch, stats}; @@ -15,6 +16,7 @@ pub struct MemProbe { pub edges: Arc, pub search: Arc, pub records: Arc, + pub identity: Arc, pub issue_states: Arc>, pub pull_statuses: Arc>, } @@ -82,6 +84,15 @@ struct Lru { capacity: u64, } +#[derive(Serialize)] +struct Identity { + entries: usize, + capacity: usize, + hits: u64, + misses: u64, + upstream_requests: u64, +} + #[derive(Serialize)] struct Derived { known_bytes: u64, @@ -95,6 +106,7 @@ struct MemSnapshot { edges: Edges, state: StateIdx, lru: Lru, + identity: Identity, derived: Derived, } @@ -109,6 +121,7 @@ async fn mem_report(State(p): State) -> Json { let su = p.search.space_usage(); let er = p.edges.mem_report(); let lru = p.records.cache_stats().unwrap_or_default(); + let identity = p.identity.stats(); let known_bytes = su.total_bytes + er.source_interner_bytes @@ -163,6 +176,13 @@ async fn mem_report(State(p): State) -> Json { len: lru.len, capacity: lru.capacity, }, + identity: Identity { + entries: identity.entries, + capacity: identity.capacity, + hits: identity.hits, + misses: identity.misses, + upstream_requests: identity.upstream_requests, + }, derived: Derived { known_bytes, allocated_minus_known: allocated.saturating_sub(known_bytes), diff --git a/bobbin/crates/ingest/examples/smoke.rs b/bobbin/crates/ingest/examples/smoke.rs index 7a105c7f..4a9a0c81 100644 --- a/bobbin/crates/ingest/examples/smoke.rs +++ b/bobbin/crates/ingest/examples/smoke.rs @@ -4,6 +4,7 @@ use std::time::Duration; use bobbin_edge_index::{CoverageWatch, EdgeStore, StateIndex}; use bobbin_ingest::{IngestConfig, IngestRuntime, RepoIdResolver, run}; use bobbin_record_lru::{NoopRecordStore, RecordStore}; +use bobbin_resolver::IdentityResolver; use bobbin_runtime::{OsEntropy, RuntimeHasher, SystemClock, TungsteniteWs}; use bobbin_types::search::NoopSearchSink; use futures::stream::{self, StreamExt}; @@ -38,7 +39,11 @@ async fn main() { coverage: coverage.clone(), search: Arc::new(NoopSearchSink), records: Arc::new(NoopRecordStore) as Arc, - resolver: Arc::new(RepoIdResolver::detached(hasher)), + resolver: Arc::new(RepoIdResolver::detached(hasher.clone())), + identity: Arc::new(IdentityResolver::detached( + hasher, + bobbin_resolver::DEFAULT_IDENTITY_CACHE_ENTRIES, + )), clock: Arc::new(SystemClock::new()), entropy: Arc::new(OsEntropy), ws: TungsteniteWs::shared(), diff --git a/bobbin/crates/ingest/src/frame.rs b/bobbin/crates/ingest/src/frame.rs index 21f39932..b96d7f3b 100644 --- a/bobbin/crates/ingest/src/frame.rs +++ b/bobbin/crates/ingest/src/frame.rs @@ -2,7 +2,7 @@ use jacquard_common::DefaultStr; use jacquard_common::types::did::Did; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::recordkey::Rkey; -use jacquard_common::types::string::Cid; +use jacquard_common::types::string::{Cid, Handle}; use jacquard_common::types::tid::Tid; use serde::Deserialize; use serde_json::value::RawValue; @@ -14,6 +14,15 @@ pub struct HydrantFrame { pub kind: FrameKind, #[serde(default)] pub record: Option, + #[serde(default)] + pub identity: Option, +} + +#[derive(Clone, Debug, Deserialize)] +pub struct IdentityFrame { + pub did: Did, + pub handle: Handle, + pub is_active: bool, } #[derive(Clone, Copy, Debug, Eq, PartialEq)] diff --git a/bobbin/crates/ingest/src/lib.rs b/bobbin/crates/ingest/src/lib.rs index 3d9113a7..ba8fcda5 100644 --- a/bobbin/crates/ingest/src/lib.rs +++ b/bobbin/crates/ingest/src/lib.rs @@ -9,7 +9,11 @@ use bobbin_edge_index::{ }; use bobbin_knot_ingest::{CapabilityGate, KnotRegistry}; use bobbin_record_lru::RecordStore; -use bobbin_resolver::{NormalizeRepoRefs, decode_canon_or_upgrade_bytes, synthesize_created_at}; +#[cfg(test)] +use bobbin_resolver::DEFAULT_IDENTITY_CACHE_ENTRIES; +use bobbin_resolver::{ + IdentityResolver, NormalizeRepoRefs, decode_canon_or_upgrade_bytes, synthesize_created_at, +}; use bobbin_runtime::{ Clock, Entropy, NetworkError, RuntimeHasher, UnixMicros, WsConn, WsMessage, WsStream, WsTransport, @@ -40,7 +44,7 @@ mod resolver; mod shadow; mod warming; use frame::HydrantStreamErrorFrame; -pub use frame::{FrameKind, HydrantFrame, RecordAction, RecordFrame}; +pub use frame::{FrameKind, HydrantFrame, IdentityFrame, RecordAction, RecordFrame}; pub use resolver::{RepoIdResolver, Resolution}; pub use shadow::{WarmingShadowBuffer, WarmingShadowSnapshot}; pub use warming::{ParkedUpsert, WarmingBuffer, WarmingBufferSnapshot}; @@ -211,6 +215,7 @@ pub struct IngestRuntime { pub search: Arc, pub records: Arc, pub resolver: Arc, + pub identity: Arc, pub clock: Arc, pub entropy: Arc, pub ws: Arc, @@ -232,6 +237,7 @@ impl Clone for IngestRuntime { search: self.search.clone(), records: self.records.clone(), resolver: self.resolver.clone(), + identity: self.identity.clone(), clock: self.clock.clone(), entropy: self.entropy.clone(), ws: self.ws.clone(), @@ -832,6 +838,7 @@ struct Pending { enum PendingOp { Noop, + Identity(Option), ClearCache { source: AtUri, }, @@ -864,7 +871,7 @@ fn pending_nsid(op: &PendingOp) -> Option<&Nsid> { PendingOp::Upsert(pieces) => Some(&pieces.nsid), PendingOp::Delete { nsid, .. } => Some(nsid), PendingOp::Parked { nsid, .. } => Some(nsid), - PendingOp::Noop | PendingOp::ClearCache { .. } => None, + PendingOp::Noop | PendingOp::Identity(_) | PendingOp::ClearCache { .. } => None, } } @@ -1010,6 +1017,7 @@ async fn commit_stage( &*rt.search, &*rt.records, &rt.resolver, + &rt.identity, ) .await; let commit_end = rt.clock.now_instant(); @@ -1044,7 +1052,8 @@ async fn prepare_frame( }; let op = match frame.kind { FrameKind::Record => prepare_record(frame.record, ctx).await, - FrameKind::Identity | FrameKind::Account => PendingOp::Noop, + FrameKind::Identity => PendingOp::Identity(frame.identity), + FrameKind::Account => PendingOp::Noop, FrameKind::Other => { debug!(id = frame.id, "ignoring unknown hydrant frame kind"); PendingOp::Noop @@ -1402,6 +1411,7 @@ async fn commit_pending( search: &S, records: &dyn RecordStore, resolver: &RepoIdResolver, + identity: &IdentityResolver, ) { let Pending { cursor, @@ -1411,6 +1421,15 @@ async fn commit_pending( } = pending; match op { PendingOp::Noop | PendingOp::Parked { .. } => {} + PendingOp::Identity(observed) => { + if let Some(observed) = observed { + if observed.is_active { + identity.observe(observed.did, observed.handle); + } else { + identity.deactivate(observed.did); + } + } + } PendingOp::ClearCache { source } => records.remove(&source), PendingOp::Upsert(pieces) => { let UpsertPieces { @@ -1483,6 +1502,8 @@ async fn handle_frame( }; let pending = claim_pending(prepare_frame(frame, &ctx, now).await, &ctx).await; let pending = resolve_pending(pending, &ctx).await; + let identity = + IdentityResolver::detached(RuntimeHasher::default(), DEFAULT_IDENTITY_CACHE_ENTRIES); commit_pending( pending, store, @@ -1492,6 +1513,7 @@ async fn handle_frame( search, records, resolver, + &identity, ) .await; } @@ -2524,7 +2546,9 @@ mod tests { "type": "identity", "identity": { "did": "did:plc:olaren", - "handle": "olaren.dev" + "handle": "olaren.dev", + "is_active": true, + "status": "active" } })); handle_frame( @@ -2545,6 +2569,90 @@ mod tests { assert!(!cov.snapshot().is_ready()); } + #[tokio::test] + async fn identity_frames_update_identity_resolver_lifecycle() { + let (store, issue_states, pull_statuses, coverage, resolver) = fresh(); + let identity = + IdentityResolver::detached(RuntimeHasher::default(), DEFAULT_IDENTITY_CACHE_ENTRIES); + let records = NoopRecordStore; + let search = NoopSearchSink; + let ctx = PipelineCtx { + resolver: &resolver, + store: &store, + issue_states: &issue_states, + pull_statuses: &pull_statuses, + coverage: &coverage, + records: &records, + search: &search, + shadow: None, + buffer: None, + knot_registry: None, + knot_gate: None, + }; + let frame: HydrantFrame = parse_frame(json!({ + "id": 5, + "type": "identity", + "identity": { + "did": "did:plc:olaren", + "handle": "olaren.dev", + "is_active": true, + "status": "active" + } + })); + let pending = prepare_frame(frame, &ctx, now()).await; + let pending = resolve_pending(pending, &ctx).await; + commit_pending( + pending, + &store, + &issue_states, + &pull_statuses, + &coverage, + &search, + &records, + &resolver, + &identity, + ) + .await; + + let doc = identity + .resolve_by_did(&Did::new_static("did:plc:olaren").unwrap()) + .await + .expect("identity frame should seed the resolver"); + assert_eq!(doc.handle.as_ref(), "olaren.dev"); + + let frame: HydrantFrame = parse_frame(json!({ + "id": 6, + "type": "identity", + "identity": { + "did": "did:plc:olaren", + "handle": "olaren.dev", + "is_active": false, + "status": "deactivated" + } + })); + let pending = prepare_frame(frame, &ctx, now()).await; + let pending = resolve_pending(pending, &ctx).await; + commit_pending( + pending, + &store, + &issue_states, + &pull_statuses, + &coverage, + &search, + &records, + &resolver, + &identity, + ) + .await; + + assert_eq!( + identity + .resolve_by_did(&Did::new_static("did:plc:olaren").unwrap()) + .await, + Err(bobbin_resolver::IdentityResolveError::NotFound) + ); + } + #[tokio::test] async fn account_frame_advances_cursor_only() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); @@ -2600,6 +2708,10 @@ mod tests { search: Arc::new(NoopSearchSink), records: Arc::new(NoopRecordStore) as Arc, resolver: Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), + identity: Arc::new(IdentityResolver::detached( + RuntimeHasher::default(), + DEFAULT_IDENTITY_CACHE_ENTRIES, + )), clock: Arc::new(SystemClock::new()), entropy: Arc::new(OsEntropy), ws: TungsteniteWs::shared(), @@ -3560,6 +3672,10 @@ mod tests { search: Arc::new(NoopSearchSink), records: capturing.clone() as Arc, resolver, + identity: Arc::new(IdentityResolver::detached( + RuntimeHasher::default(), + DEFAULT_IDENTITY_CACHE_ENTRIES, + )), clock, entropy: Arc::new(OsEntropy), ws: TungsteniteWs::shared(), diff --git a/bobbin/crates/resolver/Cargo.toml b/bobbin/crates/resolver/Cargo.toml index f0e2fbf3..b59645f6 100644 --- a/bobbin/crates/resolver/Cargo.toml +++ b/bobbin/crates/resolver/Cargo.toml @@ -13,6 +13,8 @@ jacquard-common = { workspace = true } scc = { workspace = true } serde_json = { workspace = true } +serde = { workspace = true } +thiserror = { workspace = true } tokio = { workspace = true } tracing = { workspace = true } diff --git a/bobbin/crates/resolver/src/identity.rs b/bobbin/crates/resolver/src/identity.rs new file mode 100644 index 00000000..e8187f24 --- /dev/null +++ b/bobbin/crates/resolver/src/identity.rs @@ -0,0 +1,554 @@ +use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; + +use bobbin_runtime::RuntimeHasher; +use bobbin_slingshot_client::{SlingshotClient, SlingshotError}; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use jacquard_common::types::ident::AtIdentifier; +use jacquard_common::types::string::Handle; +use scc::hash_cache::Entry as CacheEntry; +use scc::{HashCache as SccCache, HashMap as SccMap}; +use serde::{Deserialize, Serialize}; +use thiserror::Error; +use tokio::sync::OnceCell; + +pub const DEFAULT_IDENTITY_CACHE_ENTRIES: usize = 100_000; + +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct MiniDoc { + pub did: Did, + pub handle: Handle, + #[serde(skip_serializing_if = "Option::is_none")] + pub pds: Option, +} + +#[derive(Clone, Debug, Error, Eq, PartialEq)] +pub enum IdentityResolveError { + #[error("identity not found")] + NotFound, + #[error("identity upstream: {0}")] + Upstream(String), + #[error("invalid identity response: {0}")] + Decode(String), +} + +impl From for IdentityResolveError { + fn from(error: SlingshotError) -> Self { + match error { + SlingshotError::NotFound => Self::NotFound, + other => Self::Upstream(other.to_string()), + } + } +} + +// hydrant doesnt send pds so we have a separate states to have +// public resolve call upgrade from a partial-cached to a full-cached +#[derive(Clone)] +enum IdentityState { + // seen by bobbin from hydrant + Observed(MiniDoc), + // fetched by bobbin from slingshot + Fetched(MiniDoc), + Inactive, +} + +impl IdentityState { + fn doc(&self) -> Option<&MiniDoc> { + match self { + Self::Observed(doc) | Self::Fetched(doc) => Some(doc), + Self::Inactive => None, + } + } +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct IdentityResolverStatsSnapshot { + pub entries: usize, + pub capacity: usize, + pub hits: u64, + pub misses: u64, + pub upstream_requests: u64, +} + +#[derive(Default)] +struct IdentityResolverStats { + hits: AtomicU64, + misses: AtomicU64, + upstream_requests: AtomicU64, +} + +pub struct IdentityResolver { + by_did: SccCache, IdentityState, RuntimeHasher>, + by_handle: SccMap, Did, RuntimeHasher>, + in_flight: SccMap>>, RuntimeHasher>, + slingshot: Option, + stats: IdentityResolverStats, +} + +impl IdentityResolver { + pub fn with_slingshot( + slingshot: SlingshotClient, + hasher: RuntimeHasher, + capacity: usize, + ) -> Self { + Self::new(Some(slingshot), hasher, capacity) + } + + pub fn detached(hasher: RuntimeHasher, capacity: usize) -> Self { + Self::new(None, hasher, capacity) + } + + fn new(slingshot: Option, hasher: RuntimeHasher, capacity: usize) -> Self { + Self { + by_did: SccCache::with_capacity_and_hasher(0, capacity, hasher.clone()), + by_handle: SccMap::with_hasher(hasher.clone()), + in_flight: SccMap::with_hasher(hasher), + slingshot, + stats: IdentityResolverStats::default(), + } + } + + pub fn stats(&self) -> IdentityResolverStatsSnapshot { + IdentityResolverStatsSnapshot { + entries: self.by_did.len(), + capacity: *self.by_did.capacity_range().end(), + hits: self.stats.hits.load(Ordering::Relaxed), + misses: self.stats.misses.load(Ordering::Relaxed), + upstream_requests: self.stats.upstream_requests.load(Ordering::Relaxed), + } + } + + pub fn observe(&self, did: Did, handle: Handle) { + let mut evicted = None; + let mut previous_handle = None; + match self.by_did.entry_sync(did.clone()) { + CacheEntry::Occupied(mut occupied) => { + let (pds, fetched) = match occupied.get() { + IdentityState::Observed(previous) => { + previous_handle = Some(previous.handle.clone()); + (previous.pds.clone(), false) + } + IdentityState::Fetched(previous) => { + previous_handle = Some(previous.handle.clone()); + (previous.pds.clone(), previous.handle == handle) + } + IdentityState::Inactive => (None, false), + }; + let doc = MiniDoc { + did: did.clone(), + handle: handle.clone(), + pds, + }; + occupied.put(if fetched { + IdentityState::Fetched(doc) + } else { + IdentityState::Observed(doc) + }); + } + CacheEntry::Vacant(vacant) => { + let (removed, occupied) = vacant.put_entry(IdentityState::Observed(MiniDoc { + did: did.clone(), + handle: handle.clone(), + pds: None, + })); + evicted = removed; + drop(occupied); + } + } + self.remove_by_handle_if_owned(&did, previous_handle.as_ref()); + self.remove_by_handle_for_removed(evicted); + self.insert_by_handle(did, handle); + } + + pub fn deactivate(&self, did: Did) { + let mut evicted = None; + let mut previous_handle = None; + match self.by_did.entry_sync(did.clone()) { + CacheEntry::Occupied(mut occupied) => { + previous_handle = occupied.get().doc().map(|doc| doc.handle.clone()); + occupied.put(IdentityState::Inactive); + } + CacheEntry::Vacant(vacant) => { + let (removed, occupied) = vacant.put_entry(IdentityState::Inactive); + evicted = removed; + drop(occupied); + } + } + self.remove_by_handle_if_owned(&did, previous_handle.as_ref()); + self.remove_by_handle_for_removed(evicted); + } + + fn remove_by_handle_if_owned( + &self, + did: &Did, + handle: Option<&Handle>, + ) { + let Some(handle) = handle else { + return; + }; + self.by_handle.remove_if_sync(handle, |owner| owner == did); + } + + fn remove_by_handle_for_removed(&self, removed: Option<(Did, IdentityState)>) { + let Some((did, state)) = removed else { + return; + }; + let Some(doc) = state.doc() else { + return; + }; + self.remove_by_handle_if_owned(&did, Some(&doc.handle)); + } + + fn insert_by_handle(&self, did: Did, handle: Handle) { + if let Some(displaced_did) = self + .by_handle + .upsert_sync(handle.clone(), did.clone()) + .filter(|displaced_did| displaced_did != &did) + { + let removed = self.by_did.remove_if_sync(&displaced_did, |state| { + state.doc().is_some_and(|doc| doc.handle == handle) + }); + self.remove_by_handle_for_removed(removed); + } + + if self.by_did_matches_handle(&did, &handle) { + return; + } + self.remove_by_handle_if_owned(&did, Some(&handle)); + } + + fn by_did_matches_handle(&self, did: &Did, handle: &Handle) -> bool { + self.by_did + .get_sync(did) + .is_some_and(|state| state.get().doc().is_some_and(|doc| doc.handle == *handle)) + } + + fn cached( + &self, + identifier: &AtIdentifier, + require_fetched: bool, + ) -> Result, IdentityResolveError> { + match identifier { + AtIdentifier::Did(did) => match self.by_did.get_sync(did).as_deref() { + Some(IdentityState::Observed(doc)) => Ok((!require_fetched).then(|| doc.clone())), + Some(IdentityState::Fetched(doc)) => Ok(Some(doc.clone())), + Some(IdentityState::Inactive) => Err(IdentityResolveError::NotFound), + None => Ok(None), + }, + AtIdentifier::Handle(handle) => { + let Some(did) = self + .by_handle + .get_sync(handle) + .map(|entry| entry.get().clone()) + else { + return Ok(None); + }; + let state = self.by_did.get_sync(&did).map(|state| state.get().clone()); + match state { + Some(IdentityState::Observed(doc)) if doc.handle == *handle => { + return Ok((!require_fetched).then_some(doc)); + } + Some(IdentityState::Fetched(doc)) if doc.handle == *handle => { + return Ok(Some(doc)); + } + Some(IdentityState::Observed(_)) + | Some(IdentityState::Fetched(_)) + | Some(IdentityState::Inactive) + | None => {} + } + self.remove_by_handle_if_owned(&did, Some(handle)); + Ok(None) + } + } + } + + /// Resolve a DID using Hydrant's partial identity data when available. + pub async fn resolve_by_did( + &self, + did: &Did, + ) -> Result { + self.resolve_with_cache(&AtIdentifier::Did(did.clone()), false) + .await + } + + /// Resolve a minidoc, fetching Hydrant-only observations upstream first. + pub async fn resolve_minidoc( + &self, + identifier: &AtIdentifier, + ) -> Result { + self.resolve_with_cache(identifier, true).await + } + + async fn resolve_with_cache( + &self, + identifier: &AtIdentifier, + require_fetched: bool, + ) -> Result { + match self.cached(identifier, require_fetched) { + Ok(Some(doc)) => { + self.stats.hits.fetch_add(1, Ordering::Relaxed); + return Ok(doc); + } + Err(error) => { + self.stats.hits.fetch_add(1, Ordering::Relaxed); + return Err(error); + } + Ok(None) => self.stats.misses.fetch_add(1, Ordering::Relaxed), + }; + + let key = identifier.as_str().to_owned(); + let cell = self + .in_flight + .entry_async(key.clone()) + .await + .or_insert_with(|| Arc::new(OnceCell::new())) + .get() + .clone(); + let result = cell + .get_or_init(|| async { self.fetch_minidoc(identifier).await }) + .await + .clone(); + self.in_flight.remove_async(&key).await; + result + } + + async fn fetch_minidoc( + &self, + identifier: &AtIdentifier, + ) -> Result { + let client = self + .slingshot + .as_ref() + .ok_or(IdentityResolveError::NotFound)?; + self.stats.upstream_requests.fetch_add(1, Ordering::Relaxed); + let bytes = client + .resolve_mini_doc(identifier) + .await + .map_err(IdentityResolveError::from)?; + let doc = serde_json::from_slice::(&bytes) + .map_err(|error| IdentityResolveError::Decode(error.to_string()))?; + self.insert_fetched_by_did(doc) + } + + fn insert_fetched_by_did(&self, doc: MiniDoc) -> Result { + let did = doc.did.clone(); + let handle = doc.handle.clone(); + let mut evicted = None; + let mut previous_handle = None; + let stored = match self.by_did.entry_sync(did.clone()) { + CacheEntry::Occupied(mut occupied) => match occupied.get_mut() { + IdentityState::Inactive => return Err(IdentityResolveError::NotFound), + IdentityState::Observed(observed) if observed.handle != handle => { + observed.pds = doc.pds; + return Ok(observed.clone()); + } + IdentityState::Observed(previous) | IdentityState::Fetched(previous) => { + previous_handle = Some(previous.handle.clone()); + occupied.put(IdentityState::Fetched(doc.clone())); + doc + } + }, + CacheEntry::Vacant(vacant) => { + let (removed, occupied) = vacant.put_entry(IdentityState::Fetched(doc.clone())); + evicted = removed; + drop(occupied); + doc + } + }; + self.remove_by_handle_if_owned(&did, previous_handle.as_ref()); + self.remove_by_handle_for_removed(evicted); + self.insert_by_handle(did, handle); + Ok(stored) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use url::Url; + use wiremock::matchers::{method, path, query_param}; + use wiremock::{Mock, MockServer, ResponseTemplate}; + + const TEST_CAPACITY: usize = 64; + + fn hasher() -> RuntimeHasher { + RuntimeHasher::from_seeds(1, 2, 3, 4) + } + + fn resolver() -> IdentityResolver { + IdentityResolver::detached(hasher(), TEST_CAPACITY) + } + + fn did(value: &str) -> Did { + Did::new_owned(value).unwrap() + } + + fn handle(value: &str) -> Handle { + Handle::new_owned(value).unwrap() + } + + #[tokio::test] + async fn observed_identity_resolves_by_did_without_upstream() { + let resolver = resolver(); + resolver.observe(did("did:plc:dawn"), handle("ptr.pet")); + + let doc = resolver.resolve_by_did(&did("did:plc:dawn")).await.unwrap(); + assert_eq!(doc.handle, handle("ptr.pet")); + assert_eq!(doc.pds, None); + assert_eq!(resolver.stats().hits, 1); + } + + #[tokio::test] + async fn observed_identity_updates_existing_did_and_removes_old_handle() { + let resolver = resolver(); + let identity = did("did:plc:dawn"); + resolver.observe(identity.clone(), handle("ptr.pet")); + resolver.observe(identity.clone(), handle("new.ptr.pet")); + + assert!( + resolver + .cached(&AtIdentifier::Handle(handle("ptr.pet")), false) + .unwrap() + .is_none() + ); + let updated = resolver.resolve_by_did(&identity).await.unwrap(); + assert_eq!(updated.handle, handle("new.ptr.pet")); + } + + #[tokio::test] + async fn handle_reassignment_keeps_the_new_owner() { + let resolver = resolver(); + let first = did("did:plc:first"); + let second = did("did:plc:second"); + let shared = handle("shared.example.com"); + resolver.observe(first.clone(), shared.clone()); + resolver.observe(second.clone(), shared.clone()); + resolver.observe(first, handle("first.example.com")); + + let cached = resolver + .cached(&AtIdentifier::Handle(shared), false) + .unwrap() + .unwrap(); + assert_eq!(cached.did, second); + } + + #[test] + fn handle_lookup_discards_an_unvalidated_reverse_hint() { + let resolver = resolver(); + let identity = did("did:plc:dawn"); + let stale = handle("stale.example.com"); + resolver.observe(identity.clone(), handle("current.example.com")); + resolver + .by_handle + .upsert_sync(stale.clone(), identity.clone()); + + assert!( + resolver + .cached(&AtIdentifier::Handle(stale.clone()), false) + .unwrap() + .is_none() + ); + assert!(resolver.by_handle.get_sync(&stale).is_none()); + } + + #[test] + fn minidoc_lookup_keeps_a_valid_observed_handle_hint() { + let resolver = resolver(); + let identity = did("did:plc:dawn"); + let handle = handle("ptr.pet"); + resolver.observe(identity.clone(), handle.clone()); + + assert!( + resolver + .cached(&AtIdentifier::Handle(handle.clone()), true) + .unwrap() + .is_none() + ); + assert_eq!( + resolver + .by_handle + .get_sync(&handle) + .map(|entry| entry.get().clone()), + Some(identity) + ); + } + + #[tokio::test] + async fn partial_observation_fetches_and_preserves_pds() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/com.bad-example.identity.resolveMiniDoc")) + .and(query_param("identifier", "did:plc:dawn")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "did": "did:plc:dawn", + "handle": "ptr.pet", + "pds": "https://pds.example.com" + }))) + .expect(1) + .mount(&server) + .await; + let client = + SlingshotClient::with_default_http(Url::parse(&server.uri()).unwrap()).unwrap(); + let resolver = IdentityResolver::with_slingshot(client, hasher(), TEST_CAPACITY); + let identity = did("did:plc:dawn"); + resolver.observe(identity.clone(), handle("ptr.pet")); + + let doc = resolver + .resolve_minidoc(&AtIdentifier::Did(identity.clone())) + .await + .unwrap(); + assert_eq!(doc.pds.as_deref(), Some("https://pds.example.com")); + + resolver.observe(identity.clone(), handle("ptr.pet")); + let cached = resolver + .resolve_minidoc(&AtIdentifier::Did(identity)) + .await + .unwrap(); + assert_eq!(cached.pds.as_deref(), Some("https://pds.example.com")); + assert_eq!(resolver.stats().upstream_requests, 1); + } + + #[tokio::test] + async fn inactive_identity_rejects_an_in_flight_result() { + let resolver = resolver(); + let identity = did("did:plc:dawn"); + resolver.observe(identity.clone(), handle("ptr.pet")); + resolver.deactivate(identity.clone()); + + assert_eq!( + resolver.resolve_by_did(&identity).await, + Err(IdentityResolveError::NotFound) + ); + assert_eq!( + resolver.insert_fetched_by_did(MiniDoc { + did: identity, + handle: handle("ptr.pet"), + pds: Some("https://pds.example.com".to_owned()), + }), + Err(IdentityResolveError::NotFound) + ); + assert!( + resolver + .cached(&AtIdentifier::Handle(handle("ptr.pet")), false) + .unwrap() + .is_none() + ); + } + + #[test] + fn cache_capacity_bounds_forward_and_reverse_indexes() { + let resolver = resolver(); + for n in 0..512 { + resolver.observe( + did(&format!("did:plc:user{n}")), + handle(&format!("user{n}.example.com")), + ); + } + + let stats = resolver.stats(); + assert!(stats.entries <= stats.capacity); + assert!(resolver.by_handle.len() <= stats.capacity); + } +} diff --git a/bobbin/crates/resolver/src/lib.rs b/bobbin/crates/resolver/src/lib.rs index cee8a4fc..231b58be 100644 --- a/bobbin/crates/resolver/src/lib.rs +++ b/bobbin/crates/resolver/src/lib.rs @@ -1,6 +1,11 @@ +mod identity; mod legacy_upgrade; mod normalize; +pub use identity::{ + DEFAULT_IDENTITY_CACHE_ENTRIES, IdentityResolveError, IdentityResolver, + IdentityResolverStatsSnapshot, MiniDoc, +}; pub use legacy_upgrade::{ DecodedRecord, decode_canon_or_upgrade, decode_canon_or_upgrade_bytes, normalize_record_fields, scrub_record_bytes, synthesize_created_at, upgrade, upgrade_wire_bytes, diff --git a/bobbin/crates/xrpc/src/enrich.rs b/bobbin/crates/xrpc/src/enrich.rs index 7cdc02d8..6cea9e73 100644 --- a/bobbin/crates/xrpc/src/enrich.rs +++ b/bobbin/crates/xrpc/src/enrich.rs @@ -339,11 +339,11 @@ async fn resolve_minidocs( futures::stream::iter(targets) .map(|(did, sources)| async move { let doc = state - .slingshot - .resolve_mini_doc(&AtIdentifier::Did(did.clone())) + .identity + .resolve_by_did(&did) .await .ok() - .and_then(|bytes| serde_json::from_slice::(&bytes).ok()); + .and_then(|doc| serde_json::to_value(doc).ok()); (did, sources, doc) }) .buffer_unordered(MINIDOC_CONCURRENCY) diff --git a/bobbin/crates/xrpc/src/feed.rs b/bobbin/crates/xrpc/src/feed.rs index 27049e97..2f6c274f 100644 --- a/bobbin/crates/xrpc/src/feed.rs +++ b/bobbin/crates/xrpc/src/feed.rs @@ -15,7 +15,6 @@ use bobbin_types::sh_tangled::repo::{self, Repo, RepoRecord, RepoViewBasic}; use jacquard_common::DefaultStr; use jacquard_common::types::string::{AtUri, Datetime, Did, Handle, UriValue}; use jacquard_common::xrpc::XrpcResp; -use jacquard_identity::resolver::IdentityResolver; use crate::{ AppState, XrpcError, XrpcQuery, fetch, json_stream, paged_tail, parse_cursor, parse_limit, @@ -350,17 +349,11 @@ fn build_actor_viewer_state( async fn resolve_handle(state: &AppState, did: &Did) -> Handle { state - .directory - .resolve_did_doc_owned(did) + .identity + .resolve_by_did(did) .await .ok() - .and_then(|doc| { - doc.handles() - .into_iter() - .next() - .map(|h| h.as_str().to_owned()) - }) - .and_then(|s| Handle::new_owned(s).ok()) + .map(|doc| doc.handle) .unwrap_or_else(|| { Handle::new_static("handle.invalid").expect("handle.invalid is a valid handle") }) diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index 50d51ddf..3515c2fa 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -30,7 +30,9 @@ use bobbin_knot_proxy::{ KnotHost, KnotProxy, KnotProxyError, MirrorNsid, MirrorProxy, ProxyResponse, RepoSlug, }; use bobbin_record_lru::RecordStore; -use bobbin_resolver::RepoIdResolver; +use bobbin_resolver::{ + DEFAULT_IDENTITY_CACHE_ENTRIES, IdentityResolveError, IdentityResolver, RepoIdResolver, +}; use bobbin_runtime::ReqwestHttp; use bobbin_search::{ SearchCursor, SearchError, SearchFilters, SearchHit, SearchOffset, SearchReader, @@ -137,6 +139,7 @@ pub struct AppState { pub mirror: Option>, pub search: Arc, pub resolver: Arc, + pub identity: Arc, pub directory: Arc, pub limiter: Option>, pub client_address: Arc, @@ -157,6 +160,11 @@ impl AppState { resolver: Arc, directory: Arc, ) -> Self { + let identity = Arc::new(IdentityResolver::with_slingshot( + slingshot.clone(), + bobbin_runtime::RuntimeHasher::default(), + DEFAULT_IDENTITY_CACHE_ENTRIES, + )); Self { records, slingshot, @@ -168,6 +176,7 @@ impl AppState { mirror: None, search, resolver, + identity, directory, limiter: None, client_address: Arc::new(ClientAddress::default()), @@ -185,6 +194,11 @@ impl AppState { self } + pub fn with_identity(mut self, identity: Arc) -> Self { + self.identity = identity; + self + } + pub fn with_proxies(mut self, proxies: TrustedProxies) -> Self { self.client_address = Arc::new(ClientAddress::new(proxies)); self @@ -2608,12 +2622,16 @@ async fn resolve_mini_doc( State(state): State, XrpcQuery(q): XrpcQuery, ) -> Result { - let body = state - .slingshot - .resolve_mini_doc(&q.identifier) + let doc = state + .identity + .resolve_minidoc(&q.identifier) .await - .map_err(map_slingshot)?; - Ok((StatusCode::OK, [(CONTENT_TYPE, "application/json")], body).into_response()) + .map_err(|error| match error { + IdentityResolveError::NotFound => XrpcError::NotFound, + IdentityResolveError::Upstream(message) => XrpcError::UpstreamUnavailable(message), + IdentityResolveError::Decode(message) => XrpcError::InvalidRecord(message), + })?; + Ok(Json(doc).into_response()) } async fn get_coverage(State(state): State) -> Json { diff --git a/bobbin/example.toml b/bobbin/example.toml index 081dc5c0..17a9bac7 100644 --- a/bobbin/example.toml +++ b/bobbin/example.toml @@ -135,6 +135,14 @@ # Default value: 67108864 #lru_bytes = 67108864 +[identity_cache] +# Maximum number of DID entries retained by the in-process identity cache. +# +# Can also be specified via environment variable `BOBBIN_IDENTITY_CACHE_ENTRIES`. +# +# Default value: 100000 +#max_entries = 100000 + [search] # The heap size in bytes for the in-mem tantivy writer. Larger values # trade RAM for fewer segment merges - the index itself lives in