This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

bobbin/crates/{bobbin,ingest,resolver,xrpc}: have an identity cache that can be warmed by hydrant, use in xrpcs

Signed-off-by: dawn <dawn@tangled.org>

author
dawn
date (Aug 4, 2026, 3:15 PM +0300) commit 666a07cf parent 0c492ba6 change-id qzulsvqn
+786 -26
+3
Cargo.lock
··· 680 680 "bobbin-knot-ingest", 681 681 "bobbin-knot-proxy", 682 682 "bobbin-record-lru", 683 + "bobbin-resolver", 683 684 "bobbin-runtime", 684 685 "bobbin-search", 685 686 "bobbin-slingshot-client", ··· 810 811 "bobbin-types", 811 812 "jacquard-common", 812 813 "scc", 814 + "serde", 813 815 "serde_json", 816 + "thiserror 2.0.18", 814 817 "tokio", 815 818 "tracing", 816 819 "url",
+5
bobbin/crates/bobbin-sim/src/runtime.rs
··· 9 9 WarmingBuffer, WarmingShadowBuffer, run as run_ingest, 10 10 }; 11 11 use bobbin_record_lru::NoopRecordStore; 12 + use bobbin_resolver::IdentityResolver; 12 13 use bobbin_runtime::{ 13 14 Clock, DEFAULT_MEM_WS_CAPACITY, MemHttpTransport, MemWsTransport, RuntimeHasher, SeededEntropy, 14 15 SimClock, UnixMicros, ··· 115 116 search: Arc::new(NoopSearchSink), 116 117 records: records.clone() as Arc<dyn bobbin_record_lru::RecordStore>, 117 118 resolver: resolver.clone(), 119 + identity: Arc::new(IdentityResolver::detached( 120 + hasher.clone(), 121 + bobbin_resolver::DEFAULT_IDENTITY_CACHE_ENTRIES, 122 + )), 118 123 clock: clock.clone(), 119 124 entropy: entropy.clone(), 120 125 ws: mem_ws,
+1
bobbin/crates/bobbin/Cargo.toml
··· 15 15 bobbin-knot-ingest = { workspace = true } 16 16 bobbin-knot-proxy = { workspace = true } 17 17 bobbin-record-lru = { workspace = true } 18 + bobbin-resolver = { workspace = true } 18 19 bobbin-runtime = { workspace = true } 19 20 bobbin-search = { workspace = true } 20 21 bobbin-slingshot-client = { workspace = true }
+12
bobbin/crates/bobbin/src/config.rs
··· 26 26 "backpressure.reserved_index_bytes", 27 27 "slingshot.url", 28 28 "record_cache.lru_bytes", 29 + "identity_cache.max_entries", 29 30 "search.heap_bytes", 30 31 "knot.allow_private", 31 32 "knot.require_https", ··· 49 50 "BOBBIN_BACKPRESSURE_RESERVED_INDEX_BYTES", 50 51 "BOBBIN_SLINGSHOT_URL", 51 52 "BOBBIN_RECORD_LRU_BYTES", 53 + "BOBBIN_IDENTITY_CACHE_ENTRIES", 52 54 "BOBBIN_SEARCH_HEAP_BYTES", 53 55 "BOBBIN_KNOT_ALLOW_PRIVATE", 54 56 "BOBBIN_KNOT_REQUIRE_HTTPS", ··· 75 77 76 78 #[config(nested)] 77 79 pub record_cache: RecordCacheConfig, 80 + 81 + #[config(nested)] 82 + pub identity_cache: IdentityCacheConfig, 78 83 79 84 #[config(nested)] 80 85 pub search: SearchConfig, ··· 229 234 /// LRU policy keyed on URI plus payload length. 230 235 #[config(env = "BOBBIN_RECORD_LRU_BYTES", default = 67_108_864)] 231 236 pub lru_bytes: u64, 237 + } 238 + 239 + #[derive(Debug, Config)] 240 + pub struct IdentityCacheConfig { 241 + /// Maximum number of (mini) DID doc entries that the identity cache will hold. 242 + #[config(env = "BOBBIN_IDENTITY_CACHE_ENTRIES", default = 100_000)] 243 + pub max_entries: usize, 232 244 } 233 245 234 246 #[derive(Debug, Config)]
+9
bobbin/crates/bobbin/src/main.rs
··· 13 13 use bobbin_knot_ingest::{CapabilityGate, KnotClient, KnotRegistry, Orchestrator}; 14 14 use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig, classify_ip}; 15 15 use bobbin_record_lru::{CacheCapacity, LruRecordStore, RecordStore}; 16 + use bobbin_resolver::IdentityResolver; 16 17 use bobbin_runtime::{ 17 18 Clock, GuardedWs, MemoryBudget, NetworkError, OsEntropy, ReqwestHttp, RuntimeHasher, 18 19 SystemClock, TungsteniteWs, WsTransport, ··· 186 187 let records: Arc<dyn RecordStore> = 187 188 Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(lru_cap))); 188 189 let slingshot = SlingshotClient::with_default_http(cfg.slingshot.url.clone())?; 190 + let identity = Arc::new(IdentityResolver::with_slingshot( 191 + slingshot.clone(), 192 + hasher.clone(), 193 + cfg.identity_cache.max_entries, 194 + )); 189 195 let mut resolver_opts = ResolverOptions::default(); 190 196 // NOTE: see https://tangled.org/nonbinary.computer/jacquard/issues/39. 191 197 resolver_opts.did_order = vec![ ··· 273 279 search: search.clone(), 274 280 records: records.clone(), 275 281 resolver: resolver.clone(), 282 + identity: identity.clone(), 276 283 clock: clock.clone(), 277 284 entropy, 278 285 ws: ws.clone(), ··· 328 335 edges: edges.clone(), 329 336 search: search.clone(), 330 337 records: records.clone(), 338 + identity: identity.clone(), 331 339 issue_states: issue_states.clone(), 332 340 pull_statuses: pull_statuses.clone(), 333 341 }); ··· 343 351 resolver, 344 352 directory, 345 353 ) 354 + .with_identity(identity) 346 355 .with_limiter(limiter) 347 356 .with_proxies(trusted_proxies); 348 357 let app = router(state);
+20
bobbin/crates/bobbin/src/mem/report.rs
··· 6 6 use axum::{Json, Router}; 7 7 use bobbin_edge_index::{EdgeStore, IssueStateKind, PullStatusKind, StateIndex}; 8 8 use bobbin_record_lru::RecordStore; 9 + use bobbin_resolver::IdentityResolver; 9 10 use bobbin_search::SearchIndex; 10 11 use serde::Serialize; 11 12 use tikv_jemalloc_ctl::{epoch, stats}; ··· 15 16 pub edges: Arc<EdgeStore>, 16 17 pub search: Arc<SearchIndex>, 17 18 pub records: Arc<dyn RecordStore>, 19 + pub identity: Arc<IdentityResolver>, 18 20 pub issue_states: Arc<StateIndex<IssueStateKind>>, 19 21 pub pull_statuses: Arc<StateIndex<PullStatusKind>>, 20 22 } ··· 83 85 } 84 86 85 87 #[derive(Serialize)] 88 + struct Identity { 89 + entries: usize, 90 + capacity: usize, 91 + hits: u64, 92 + misses: u64, 93 + upstream_requests: u64, 94 + } 95 + 96 + #[derive(Serialize)] 86 97 struct Derived { 87 98 known_bytes: u64, 88 99 allocated_minus_known: u64, ··· 95 106 edges: Edges, 96 107 state: StateIdx, 97 108 lru: Lru, 109 + identity: Identity, 98 110 derived: Derived, 99 111 } 100 112 ··· 109 121 let su = p.search.space_usage(); 110 122 let er = p.edges.mem_report(); 111 123 let lru = p.records.cache_stats().unwrap_or_default(); 124 + let identity = p.identity.stats(); 112 125 113 126 let known_bytes = su.total_bytes 114 127 + er.source_interner_bytes ··· 162 175 weight: lru.weight, 163 176 len: lru.len, 164 177 capacity: lru.capacity, 178 + }, 179 + identity: Identity { 180 + entries: identity.entries, 181 + capacity: identity.capacity, 182 + hits: identity.hits, 183 + misses: identity.misses, 184 + upstream_requests: identity.upstream_requests, 165 185 }, 166 186 derived: Derived { 167 187 known_bytes,
+6 -1
bobbin/crates/ingest/examples/smoke.rs
··· 4 4 use bobbin_edge_index::{CoverageWatch, EdgeStore, StateIndex}; 5 5 use bobbin_ingest::{IngestConfig, IngestRuntime, RepoIdResolver, run}; 6 6 use bobbin_record_lru::{NoopRecordStore, RecordStore}; 7 + use bobbin_resolver::IdentityResolver; 7 8 use bobbin_runtime::{OsEntropy, RuntimeHasher, SystemClock, TungsteniteWs}; 8 9 use bobbin_types::search::NoopSearchSink; 9 10 use futures::stream::{self, StreamExt}; ··· 38 39 coverage: coverage.clone(), 39 40 search: Arc::new(NoopSearchSink), 40 41 records: Arc::new(NoopRecordStore) as Arc<dyn RecordStore>, 41 - resolver: Arc::new(RepoIdResolver::detached(hasher)), 42 + resolver: Arc::new(RepoIdResolver::detached(hasher.clone())), 43 + identity: Arc::new(IdentityResolver::detached( 44 + hasher, 45 + bobbin_resolver::DEFAULT_IDENTITY_CACHE_ENTRIES, 46 + )), 42 47 clock: Arc::new(SystemClock::new()), 43 48 entropy: Arc::new(OsEntropy), 44 49 ws: TungsteniteWs::shared(),
+10 -1
bobbin/crates/ingest/src/frame.rs
··· 2 2 use jacquard_common::types::did::Did; 3 3 use jacquard_common::types::nsid::Nsid; 4 4 use jacquard_common::types::recordkey::Rkey; 5 - use jacquard_common::types::string::Cid; 5 + use jacquard_common::types::string::{Cid, Handle}; 6 6 use jacquard_common::types::tid::Tid; 7 7 use serde::Deserialize; 8 8 use serde_json::value::RawValue; ··· 14 14 pub kind: FrameKind, 15 15 #[serde(default)] 16 16 pub record: Option<RecordFrame>, 17 + #[serde(default)] 18 + pub identity: Option<IdentityFrame>, 19 + } 20 + 21 + #[derive(Clone, Debug, Deserialize)] 22 + pub struct IdentityFrame { 23 + pub did: Did<DefaultStr>, 24 + pub handle: Handle<DefaultStr>, 25 + pub is_active: bool, 17 26 } 18 27 19 28 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
+121 -5
bobbin/crates/ingest/src/lib.rs
··· 9 9 }; 10 10 use bobbin_knot_ingest::{CapabilityGate, KnotRegistry}; 11 11 use bobbin_record_lru::RecordStore; 12 - use bobbin_resolver::{NormalizeRepoRefs, decode_canon_or_upgrade_bytes, synthesize_created_at}; 12 + #[cfg(test)] 13 + use bobbin_resolver::DEFAULT_IDENTITY_CACHE_ENTRIES; 14 + use bobbin_resolver::{ 15 + IdentityResolver, NormalizeRepoRefs, decode_canon_or_upgrade_bytes, synthesize_created_at, 16 + }; 13 17 use bobbin_runtime::{ 14 18 Clock, Entropy, NetworkError, RuntimeHasher, UnixMicros, WsConn, WsMessage, WsStream, 15 19 WsTransport, ··· 39 43 mod shadow; 40 44 mod warming; 41 45 use frame::HydrantStreamErrorFrame; 42 - pub use frame::{FrameKind, HydrantFrame, RecordAction, RecordFrame}; 46 + pub use frame::{FrameKind, HydrantFrame, IdentityFrame, RecordAction, RecordFrame}; 43 47 pub use resolver::{RepoIdResolver, Resolution}; 44 48 pub use shadow::{WarmingShadowBuffer, WarmingShadowSnapshot}; 45 49 pub use warming::{ParkedUpsert, WarmingBuffer, WarmingBufferSnapshot}; ··· 210 214 pub search: Arc<S>, 211 215 pub records: Arc<dyn RecordStore>, 212 216 pub resolver: Arc<RepoIdResolver>, 217 + pub identity: Arc<IdentityResolver>, 213 218 pub clock: Arc<dyn Clock>, 214 219 pub entropy: Arc<dyn Entropy>, 215 220 pub ws: Arc<dyn WsTransport>, ··· 231 236 search: self.search.clone(), 232 237 records: self.records.clone(), 233 238 resolver: self.resolver.clone(), 239 + identity: self.identity.clone(), 234 240 clock: self.clock.clone(), 235 241 entropy: self.entropy.clone(), 236 242 ws: self.ws.clone(), ··· 828 834 829 835 enum PendingOp { 830 836 Noop, 837 + Identity(Option<IdentityFrame>), 831 838 ClearCache { 832 839 source: AtUri<DefaultStr>, 833 840 }, ··· 867 874 PendingOp::Upsert { nsid, .. } => Some(nsid), 868 875 PendingOp::Delete { nsid, .. } => Some(nsid), 869 876 PendingOp::Parked { nsid, .. } => Some(nsid), 870 - PendingOp::Noop | PendingOp::ClearCache { .. } => None, 877 + PendingOp::Noop | PendingOp::Identity(_) | PendingOp::ClearCache { .. } => None, 871 878 } 872 879 } 873 880 ··· 937 944 &*rt.search, 938 945 &*rt.records, 939 946 &rt.resolver, 947 + &rt.identity, 940 948 ) 941 949 .await; 942 950 let commit_end = rt.clock.now_instant(); ··· 971 979 }; 972 980 let op = match frame.kind { 973 981 FrameKind::Record => prepare_record(frame.record, ctx).await, 974 - FrameKind::Identity | FrameKind::Account => PendingOp::Noop, 982 + FrameKind::Identity => PendingOp::Identity(frame.identity), 983 + FrameKind::Account => PendingOp::Noop, 975 984 FrameKind::Other => { 976 985 debug!(id = frame.id, "ignoring unknown hydrant frame kind"); 977 986 PendingOp::Noop ··· 1382 1391 search: &S, 1383 1392 records: &dyn RecordStore, 1384 1393 resolver: &RepoIdResolver, 1394 + identity: &IdentityResolver, 1385 1395 ) { 1386 1396 let Pending { 1387 1397 cursor, ··· 1391 1401 } = pending; 1392 1402 match op { 1393 1403 PendingOp::Noop | PendingOp::Parked { .. } => {} 1404 + PendingOp::Identity(observed) => { 1405 + if let Some(observed) = observed { 1406 + if observed.is_active { 1407 + identity.observe(observed.did, observed.handle); 1408 + } else { 1409 + identity.deactivate(observed.did); 1410 + } 1411 + } 1412 + } 1394 1413 PendingOp::ClearCache { source } => records.remove(&source), 1395 1414 PendingOp::Upsert { 1396 1415 source, ··· 1460 1479 }; 1461 1480 let pending = prepare_frame(frame, &ctx, now).await; 1462 1481 let pending = resolve_pending(pending, &ctx).await; 1482 + let identity = 1483 + IdentityResolver::detached(RuntimeHasher::default(), DEFAULT_IDENTITY_CACHE_ENTRIES); 1463 1484 commit_pending( 1464 1485 pending, 1465 1486 store, ··· 1469 1490 search, 1470 1491 records, 1471 1492 resolver, 1493 + &identity, 1472 1494 ) 1473 1495 .await; 1474 1496 } ··· 2493 2515 "type": "identity", 2494 2516 "identity": { 2495 2517 "did": "did:plc:olaren", 2496 - "handle": "olaren.dev" 2518 + "handle": "olaren.dev", 2519 + "is_active": true, 2520 + "status": "active" 2497 2521 } 2498 2522 })); 2499 2523 handle_frame( ··· 2515 2539 } 2516 2540 2517 2541 #[tokio::test] 2542 + async fn identity_frames_update_identity_resolver_lifecycle() { 2543 + let (store, issue_states, pull_statuses, coverage, resolver) = fresh(); 2544 + let identity = 2545 + IdentityResolver::detached(RuntimeHasher::default(), DEFAULT_IDENTITY_CACHE_ENTRIES); 2546 + let records = NoopRecordStore; 2547 + let search = NoopSearchSink; 2548 + let ctx = PipelineCtx { 2549 + resolver: &resolver, 2550 + store: &store, 2551 + issue_states: &issue_states, 2552 + pull_statuses: &pull_statuses, 2553 + coverage: &coverage, 2554 + records: &records, 2555 + search: &search, 2556 + shadow: None, 2557 + buffer: None, 2558 + knot_registry: None, 2559 + knot_gate: None, 2560 + }; 2561 + let frame: HydrantFrame = parse_frame(json!({ 2562 + "id": 5, 2563 + "type": "identity", 2564 + "identity": { 2565 + "did": "did:plc:olaren", 2566 + "handle": "olaren.dev", 2567 + "is_active": true, 2568 + "status": "active" 2569 + } 2570 + })); 2571 + let pending = prepare_frame(frame, &ctx, now()).await; 2572 + let pending = resolve_pending(pending, &ctx).await; 2573 + commit_pending( 2574 + pending, 2575 + &store, 2576 + &issue_states, 2577 + &pull_statuses, 2578 + &coverage, 2579 + &search, 2580 + &records, 2581 + &resolver, 2582 + &identity, 2583 + ) 2584 + .await; 2585 + 2586 + let doc = identity 2587 + .resolve_by_did(&Did::new_static("did:plc:olaren").unwrap()) 2588 + .await 2589 + .expect("identity frame should seed the resolver"); 2590 + assert_eq!(doc.handle.as_ref(), "olaren.dev"); 2591 + 2592 + let frame: HydrantFrame = parse_frame(json!({ 2593 + "id": 6, 2594 + "type": "identity", 2595 + "identity": { 2596 + "did": "did:plc:olaren", 2597 + "handle": "olaren.dev", 2598 + "is_active": false, 2599 + "status": "deactivated" 2600 + } 2601 + })); 2602 + let pending = prepare_frame(frame, &ctx, now()).await; 2603 + let pending = resolve_pending(pending, &ctx).await; 2604 + commit_pending( 2605 + pending, 2606 + &store, 2607 + &issue_states, 2608 + &pull_statuses, 2609 + &coverage, 2610 + &search, 2611 + &records, 2612 + &resolver, 2613 + &identity, 2614 + ) 2615 + .await; 2616 + 2617 + assert_eq!( 2618 + identity 2619 + .resolve_by_did(&Did::new_static("did:plc:olaren").unwrap()) 2620 + .await, 2621 + Err(bobbin_resolver::IdentityResolveError::NotFound) 2622 + ); 2623 + } 2624 + 2625 + #[tokio::test] 2518 2626 async fn account_frame_advances_cursor_only() { 2519 2627 let (store, issue_states, pull_statuses, cov, resolver) = fresh(); 2520 2628 let frame: HydrantFrame = parse_frame(json!({ ··· 2569 2677 search: Arc::new(NoopSearchSink), 2570 2678 records: Arc::new(NoopRecordStore) as Arc<dyn RecordStore>, 2571 2679 resolver: Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), 2680 + identity: Arc::new(IdentityResolver::detached( 2681 + RuntimeHasher::default(), 2682 + DEFAULT_IDENTITY_CACHE_ENTRIES, 2683 + )), 2572 2684 clock: Arc::new(SystemClock::new()), 2573 2685 entropy: Arc::new(OsEntropy), 2574 2686 ws: TungsteniteWs::shared(), ··· 3441 3553 search: Arc::new(NoopSearchSink), 3442 3554 records: capturing.clone() as Arc<dyn RecordStore>, 3443 3555 resolver, 3556 + identity: Arc::new(IdentityResolver::detached( 3557 + RuntimeHasher::default(), 3558 + DEFAULT_IDENTITY_CACHE_ENTRIES, 3559 + )), 3444 3560 clock, 3445 3561 entropy: Arc::new(OsEntropy), 3446 3562 ws: TungsteniteWs::shared(),
+2
bobbin/crates/resolver/Cargo.toml
··· 13 13 14 14 scc = { workspace = true } 15 15 serde_json = { workspace = true } 16 + serde = { workspace = true } 17 + thiserror = { workspace = true } 16 18 tokio = { workspace = true } 17 19 tracing = { workspace = true } 18 20
+554
bobbin/crates/resolver/src/identity.rs
··· 1 + use std::sync::Arc; 2 + use std::sync::atomic::{AtomicU64, Ordering}; 3 + 4 + use bobbin_runtime::RuntimeHasher; 5 + use bobbin_slingshot_client::{SlingshotClient, SlingshotError}; 6 + use jacquard_common::DefaultStr; 7 + use jacquard_common::types::did::Did; 8 + use jacquard_common::types::ident::AtIdentifier; 9 + use jacquard_common::types::string::Handle; 10 + use scc::hash_cache::Entry as CacheEntry; 11 + use scc::{HashCache as SccCache, HashMap as SccMap}; 12 + use serde::{Deserialize, Serialize}; 13 + use thiserror::Error; 14 + use tokio::sync::OnceCell; 15 + 16 + pub const DEFAULT_IDENTITY_CACHE_ENTRIES: usize = 100_000; 17 + 18 + #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] 19 + #[serde(rename_all = "camelCase")] 20 + pub struct MiniDoc { 21 + pub did: Did<DefaultStr>, 22 + pub handle: Handle<DefaultStr>, 23 + #[serde(skip_serializing_if = "Option::is_none")] 24 + pub pds: Option<String>, 25 + } 26 + 27 + #[derive(Clone, Debug, Error, Eq, PartialEq)] 28 + pub enum IdentityResolveError { 29 + #[error("identity not found")] 30 + NotFound, 31 + #[error("identity upstream: {0}")] 32 + Upstream(String), 33 + #[error("invalid identity response: {0}")] 34 + Decode(String), 35 + } 36 + 37 + impl From<SlingshotError> for IdentityResolveError { 38 + fn from(error: SlingshotError) -> Self { 39 + match error { 40 + SlingshotError::NotFound => Self::NotFound, 41 + other => Self::Upstream(other.to_string()), 42 + } 43 + } 44 + } 45 + 46 + // hydrant doesnt send pds so we have a separate states to have 47 + // public resolve call upgrade from a partial-cached to a full-cached 48 + #[derive(Clone)] 49 + enum IdentityState { 50 + // seen by bobbin from hydrant 51 + Observed(MiniDoc), 52 + // fetched by bobbin from slingshot 53 + Fetched(MiniDoc), 54 + Inactive, 55 + } 56 + 57 + impl IdentityState { 58 + fn doc(&self) -> Option<&MiniDoc> { 59 + match self { 60 + Self::Observed(doc) | Self::Fetched(doc) => Some(doc), 61 + Self::Inactive => None, 62 + } 63 + } 64 + } 65 + 66 + #[derive(Clone, Copy, Debug, Eq, PartialEq)] 67 + pub struct IdentityResolverStatsSnapshot { 68 + pub entries: usize, 69 + pub capacity: usize, 70 + pub hits: u64, 71 + pub misses: u64, 72 + pub upstream_requests: u64, 73 + } 74 + 75 + #[derive(Default)] 76 + struct IdentityResolverStats { 77 + hits: AtomicU64, 78 + misses: AtomicU64, 79 + upstream_requests: AtomicU64, 80 + } 81 + 82 + pub struct IdentityResolver { 83 + by_did: SccCache<Did<DefaultStr>, IdentityState, RuntimeHasher>, 84 + by_handle: SccMap<Handle<DefaultStr>, Did<DefaultStr>, RuntimeHasher>, 85 + in_flight: SccMap<String, Arc<OnceCell<Result<MiniDoc, IdentityResolveError>>>, RuntimeHasher>, 86 + slingshot: Option<SlingshotClient>, 87 + stats: IdentityResolverStats, 88 + } 89 + 90 + impl IdentityResolver { 91 + pub fn with_slingshot( 92 + slingshot: SlingshotClient, 93 + hasher: RuntimeHasher, 94 + capacity: usize, 95 + ) -> Self { 96 + Self::new(Some(slingshot), hasher, capacity) 97 + } 98 + 99 + pub fn detached(hasher: RuntimeHasher, capacity: usize) -> Self { 100 + Self::new(None, hasher, capacity) 101 + } 102 + 103 + fn new(slingshot: Option<SlingshotClient>, hasher: RuntimeHasher, capacity: usize) -> Self { 104 + Self { 105 + by_did: SccCache::with_capacity_and_hasher(0, capacity, hasher.clone()), 106 + by_handle: SccMap::with_hasher(hasher.clone()), 107 + in_flight: SccMap::with_hasher(hasher), 108 + slingshot, 109 + stats: IdentityResolverStats::default(), 110 + } 111 + } 112 + 113 + pub fn stats(&self) -> IdentityResolverStatsSnapshot { 114 + IdentityResolverStatsSnapshot { 115 + entries: self.by_did.len(), 116 + capacity: *self.by_did.capacity_range().end(), 117 + hits: self.stats.hits.load(Ordering::Relaxed), 118 + misses: self.stats.misses.load(Ordering::Relaxed), 119 + upstream_requests: self.stats.upstream_requests.load(Ordering::Relaxed), 120 + } 121 + } 122 + 123 + pub fn observe(&self, did: Did<DefaultStr>, handle: Handle<DefaultStr>) { 124 + let mut evicted = None; 125 + let mut previous_handle = None; 126 + match self.by_did.entry_sync(did.clone()) { 127 + CacheEntry::Occupied(mut occupied) => { 128 + let (pds, fetched) = match occupied.get() { 129 + IdentityState::Observed(previous) => { 130 + previous_handle = Some(previous.handle.clone()); 131 + (previous.pds.clone(), false) 132 + } 133 + IdentityState::Fetched(previous) => { 134 + previous_handle = Some(previous.handle.clone()); 135 + (previous.pds.clone(), previous.handle == handle) 136 + } 137 + IdentityState::Inactive => (None, false), 138 + }; 139 + let doc = MiniDoc { 140 + did: did.clone(), 141 + handle: handle.clone(), 142 + pds, 143 + }; 144 + occupied.put(if fetched { 145 + IdentityState::Fetched(doc) 146 + } else { 147 + IdentityState::Observed(doc) 148 + }); 149 + } 150 + CacheEntry::Vacant(vacant) => { 151 + let (removed, occupied) = vacant.put_entry(IdentityState::Observed(MiniDoc { 152 + did: did.clone(), 153 + handle: handle.clone(), 154 + pds: None, 155 + })); 156 + evicted = removed; 157 + drop(occupied); 158 + } 159 + } 160 + self.remove_by_handle_if_owned(&did, previous_handle.as_ref()); 161 + self.remove_by_handle_for_removed(evicted); 162 + self.insert_by_handle(did, handle); 163 + } 164 + 165 + pub fn deactivate(&self, did: Did<DefaultStr>) { 166 + let mut evicted = None; 167 + let mut previous_handle = None; 168 + match self.by_did.entry_sync(did.clone()) { 169 + CacheEntry::Occupied(mut occupied) => { 170 + previous_handle = occupied.get().doc().map(|doc| doc.handle.clone()); 171 + occupied.put(IdentityState::Inactive); 172 + } 173 + CacheEntry::Vacant(vacant) => { 174 + let (removed, occupied) = vacant.put_entry(IdentityState::Inactive); 175 + evicted = removed; 176 + drop(occupied); 177 + } 178 + } 179 + self.remove_by_handle_if_owned(&did, previous_handle.as_ref()); 180 + self.remove_by_handle_for_removed(evicted); 181 + } 182 + 183 + fn remove_by_handle_if_owned( 184 + &self, 185 + did: &Did<DefaultStr>, 186 + handle: Option<&Handle<DefaultStr>>, 187 + ) { 188 + let Some(handle) = handle else { 189 + return; 190 + }; 191 + self.by_handle.remove_if_sync(handle, |owner| owner == did); 192 + } 193 + 194 + fn remove_by_handle_for_removed(&self, removed: Option<(Did<DefaultStr>, IdentityState)>) { 195 + let Some((did, state)) = removed else { 196 + return; 197 + }; 198 + let Some(doc) = state.doc() else { 199 + return; 200 + }; 201 + self.remove_by_handle_if_owned(&did, Some(&doc.handle)); 202 + } 203 + 204 + fn insert_by_handle(&self, did: Did<DefaultStr>, handle: Handle<DefaultStr>) { 205 + if let Some(displaced_did) = self 206 + .by_handle 207 + .upsert_sync(handle.clone(), did.clone()) 208 + .filter(|displaced_did| displaced_did != &did) 209 + { 210 + let removed = self.by_did.remove_if_sync(&displaced_did, |state| { 211 + state.doc().is_some_and(|doc| doc.handle == handle) 212 + }); 213 + self.remove_by_handle_for_removed(removed); 214 + } 215 + 216 + if self.by_did_matches_handle(&did, &handle) { 217 + return; 218 + } 219 + self.remove_by_handle_if_owned(&did, Some(&handle)); 220 + } 221 + 222 + fn by_did_matches_handle(&self, did: &Did<DefaultStr>, handle: &Handle<DefaultStr>) -> bool { 223 + self.by_did 224 + .get_sync(did) 225 + .is_some_and(|state| state.get().doc().is_some_and(|doc| doc.handle == *handle)) 226 + } 227 + 228 + fn cached( 229 + &self, 230 + identifier: &AtIdentifier<DefaultStr>, 231 + require_fetched: bool, 232 + ) -> Result<Option<MiniDoc>, IdentityResolveError> { 233 + match identifier { 234 + AtIdentifier::Did(did) => match self.by_did.get_sync(did).as_deref() { 235 + Some(IdentityState::Observed(doc)) => Ok((!require_fetched).then(|| doc.clone())), 236 + Some(IdentityState::Fetched(doc)) => Ok(Some(doc.clone())), 237 + Some(IdentityState::Inactive) => Err(IdentityResolveError::NotFound), 238 + None => Ok(None), 239 + }, 240 + AtIdentifier::Handle(handle) => { 241 + let Some(did) = self 242 + .by_handle 243 + .get_sync(handle) 244 + .map(|entry| entry.get().clone()) 245 + else { 246 + return Ok(None); 247 + }; 248 + let state = self.by_did.get_sync(&did).map(|state| state.get().clone()); 249 + match state { 250 + Some(IdentityState::Observed(doc)) if doc.handle == *handle => { 251 + return Ok((!require_fetched).then_some(doc)); 252 + } 253 + Some(IdentityState::Fetched(doc)) if doc.handle == *handle => { 254 + return Ok(Some(doc)); 255 + } 256 + Some(IdentityState::Observed(_)) 257 + | Some(IdentityState::Fetched(_)) 258 + | Some(IdentityState::Inactive) 259 + | None => {} 260 + } 261 + self.remove_by_handle_if_owned(&did, Some(handle)); 262 + Ok(None) 263 + } 264 + } 265 + } 266 + 267 + /// Resolve a DID using Hydrant's partial identity data when available. 268 + pub async fn resolve_by_did( 269 + &self, 270 + did: &Did<DefaultStr>, 271 + ) -> Result<MiniDoc, IdentityResolveError> { 272 + self.resolve_with_cache(&AtIdentifier::Did(did.clone()), false) 273 + .await 274 + } 275 + 276 + /// Resolve a minidoc, fetching Hydrant-only observations upstream first. 277 + pub async fn resolve_minidoc( 278 + &self, 279 + identifier: &AtIdentifier<DefaultStr>, 280 + ) -> Result<MiniDoc, IdentityResolveError> { 281 + self.resolve_with_cache(identifier, true).await 282 + } 283 + 284 + async fn resolve_with_cache( 285 + &self, 286 + identifier: &AtIdentifier<DefaultStr>, 287 + require_fetched: bool, 288 + ) -> Result<MiniDoc, IdentityResolveError> { 289 + match self.cached(identifier, require_fetched) { 290 + Ok(Some(doc)) => { 291 + self.stats.hits.fetch_add(1, Ordering::Relaxed); 292 + return Ok(doc); 293 + } 294 + Err(error) => { 295 + self.stats.hits.fetch_add(1, Ordering::Relaxed); 296 + return Err(error); 297 + } 298 + Ok(None) => self.stats.misses.fetch_add(1, Ordering::Relaxed), 299 + }; 300 + 301 + let key = identifier.as_str().to_owned(); 302 + let cell = self 303 + .in_flight 304 + .entry_async(key.clone()) 305 + .await 306 + .or_insert_with(|| Arc::new(OnceCell::new())) 307 + .get() 308 + .clone(); 309 + let result = cell 310 + .get_or_init(|| async { self.fetch_minidoc(identifier).await }) 311 + .await 312 + .clone(); 313 + self.in_flight.remove_async(&key).await; 314 + result 315 + } 316 + 317 + async fn fetch_minidoc( 318 + &self, 319 + identifier: &AtIdentifier<DefaultStr>, 320 + ) -> Result<MiniDoc, IdentityResolveError> { 321 + let client = self 322 + .slingshot 323 + .as_ref() 324 + .ok_or(IdentityResolveError::NotFound)?; 325 + self.stats.upstream_requests.fetch_add(1, Ordering::Relaxed); 326 + let bytes = client 327 + .resolve_mini_doc(identifier) 328 + .await 329 + .map_err(IdentityResolveError::from)?; 330 + let doc = serde_json::from_slice::<MiniDoc>(&bytes) 331 + .map_err(|error| IdentityResolveError::Decode(error.to_string()))?; 332 + self.insert_fetched_by_did(doc) 333 + } 334 + 335 + fn insert_fetched_by_did(&self, doc: MiniDoc) -> Result<MiniDoc, IdentityResolveError> { 336 + let did = doc.did.clone(); 337 + let handle = doc.handle.clone(); 338 + let mut evicted = None; 339 + let mut previous_handle = None; 340 + let stored = match self.by_did.entry_sync(did.clone()) { 341 + CacheEntry::Occupied(mut occupied) => match occupied.get_mut() { 342 + IdentityState::Inactive => return Err(IdentityResolveError::NotFound), 343 + IdentityState::Observed(observed) if observed.handle != handle => { 344 + observed.pds = doc.pds; 345 + return Ok(observed.clone()); 346 + } 347 + IdentityState::Observed(previous) | IdentityState::Fetched(previous) => { 348 + previous_handle = Some(previous.handle.clone()); 349 + occupied.put(IdentityState::Fetched(doc.clone())); 350 + doc 351 + } 352 + }, 353 + CacheEntry::Vacant(vacant) => { 354 + let (removed, occupied) = vacant.put_entry(IdentityState::Fetched(doc.clone())); 355 + evicted = removed; 356 + drop(occupied); 357 + doc 358 + } 359 + }; 360 + self.remove_by_handle_if_owned(&did, previous_handle.as_ref()); 361 + self.remove_by_handle_for_removed(evicted); 362 + self.insert_by_handle(did, handle); 363 + Ok(stored) 364 + } 365 + } 366 + 367 + #[cfg(test)] 368 + mod tests { 369 + use super::*; 370 + use url::Url; 371 + use wiremock::matchers::{method, path, query_param}; 372 + use wiremock::{Mock, MockServer, ResponseTemplate}; 373 + 374 + const TEST_CAPACITY: usize = 64; 375 + 376 + fn hasher() -> RuntimeHasher { 377 + RuntimeHasher::from_seeds(1, 2, 3, 4) 378 + } 379 + 380 + fn resolver() -> IdentityResolver { 381 + IdentityResolver::detached(hasher(), TEST_CAPACITY) 382 + } 383 + 384 + fn did(value: &str) -> Did<DefaultStr> { 385 + Did::new_owned(value).unwrap() 386 + } 387 + 388 + fn handle(value: &str) -> Handle<DefaultStr> { 389 + Handle::new_owned(value).unwrap() 390 + } 391 + 392 + #[tokio::test] 393 + async fn observed_identity_resolves_by_did_without_upstream() { 394 + let resolver = resolver(); 395 + resolver.observe(did("did:plc:dawn"), handle("ptr.pet")); 396 + 397 + let doc = resolver.resolve_by_did(&did("did:plc:dawn")).await.unwrap(); 398 + assert_eq!(doc.handle, handle("ptr.pet")); 399 + assert_eq!(doc.pds, None); 400 + assert_eq!(resolver.stats().hits, 1); 401 + } 402 + 403 + #[tokio::test] 404 + async fn observed_identity_updates_existing_did_and_removes_old_handle() { 405 + let resolver = resolver(); 406 + let identity = did("did:plc:dawn"); 407 + resolver.observe(identity.clone(), handle("ptr.pet")); 408 + resolver.observe(identity.clone(), handle("new.ptr.pet")); 409 + 410 + assert!( 411 + resolver 412 + .cached(&AtIdentifier::Handle(handle("ptr.pet")), false) 413 + .unwrap() 414 + .is_none() 415 + ); 416 + let updated = resolver.resolve_by_did(&identity).await.unwrap(); 417 + assert_eq!(updated.handle, handle("new.ptr.pet")); 418 + } 419 + 420 + #[tokio::test] 421 + async fn handle_reassignment_keeps_the_new_owner() { 422 + let resolver = resolver(); 423 + let first = did("did:plc:first"); 424 + let second = did("did:plc:second"); 425 + let shared = handle("shared.example.com"); 426 + resolver.observe(first.clone(), shared.clone()); 427 + resolver.observe(second.clone(), shared.clone()); 428 + resolver.observe(first, handle("first.example.com")); 429 + 430 + let cached = resolver 431 + .cached(&AtIdentifier::Handle(shared), false) 432 + .unwrap() 433 + .unwrap(); 434 + assert_eq!(cached.did, second); 435 + } 436 + 437 + #[test] 438 + fn handle_lookup_discards_an_unvalidated_reverse_hint() { 439 + let resolver = resolver(); 440 + let identity = did("did:plc:dawn"); 441 + let stale = handle("stale.example.com"); 442 + resolver.observe(identity.clone(), handle("current.example.com")); 443 + resolver 444 + .by_handle 445 + .upsert_sync(stale.clone(), identity.clone()); 446 + 447 + assert!( 448 + resolver 449 + .cached(&AtIdentifier::Handle(stale.clone()), false) 450 + .unwrap() 451 + .is_none() 452 + ); 453 + assert!(resolver.by_handle.get_sync(&stale).is_none()); 454 + } 455 + 456 + #[test] 457 + fn minidoc_lookup_keeps_a_valid_observed_handle_hint() { 458 + let resolver = resolver(); 459 + let identity = did("did:plc:dawn"); 460 + let handle = handle("ptr.pet"); 461 + resolver.observe(identity.clone(), handle.clone()); 462 + 463 + assert!( 464 + resolver 465 + .cached(&AtIdentifier::Handle(handle.clone()), true) 466 + .unwrap() 467 + .is_none() 468 + ); 469 + assert_eq!( 470 + resolver 471 + .by_handle 472 + .get_sync(&handle) 473 + .map(|entry| entry.get().clone()), 474 + Some(identity) 475 + ); 476 + } 477 + 478 + #[tokio::test] 479 + async fn partial_observation_fetches_and_preserves_pds() { 480 + let server = MockServer::start().await; 481 + Mock::given(method("GET")) 482 + .and(path("/xrpc/com.bad-example.identity.resolveMiniDoc")) 483 + .and(query_param("identifier", "did:plc:dawn")) 484 + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ 485 + "did": "did:plc:dawn", 486 + "handle": "ptr.pet", 487 + "pds": "https://pds.example.com" 488 + }))) 489 + .expect(1) 490 + .mount(&server) 491 + .await; 492 + let client = 493 + SlingshotClient::with_default_http(Url::parse(&server.uri()).unwrap()).unwrap(); 494 + let resolver = IdentityResolver::with_slingshot(client, hasher(), TEST_CAPACITY); 495 + let identity = did("did:plc:dawn"); 496 + resolver.observe(identity.clone(), handle("ptr.pet")); 497 + 498 + let doc = resolver 499 + .resolve_minidoc(&AtIdentifier::Did(identity.clone())) 500 + .await 501 + .unwrap(); 502 + assert_eq!(doc.pds.as_deref(), Some("https://pds.example.com")); 503 + 504 + resolver.observe(identity.clone(), handle("ptr.pet")); 505 + let cached = resolver 506 + .resolve_minidoc(&AtIdentifier::Did(identity)) 507 + .await 508 + .unwrap(); 509 + assert_eq!(cached.pds.as_deref(), Some("https://pds.example.com")); 510 + assert_eq!(resolver.stats().upstream_requests, 1); 511 + } 512 + 513 + #[tokio::test] 514 + async fn inactive_identity_rejects_an_in_flight_result() { 515 + let resolver = resolver(); 516 + let identity = did("did:plc:dawn"); 517 + resolver.observe(identity.clone(), handle("ptr.pet")); 518 + resolver.deactivate(identity.clone()); 519 + 520 + assert_eq!( 521 + resolver.resolve_by_did(&identity).await, 522 + Err(IdentityResolveError::NotFound) 523 + ); 524 + assert_eq!( 525 + resolver.insert_fetched_by_did(MiniDoc { 526 + did: identity, 527 + handle: handle("ptr.pet"), 528 + pds: Some("https://pds.example.com".to_owned()), 529 + }), 530 + Err(IdentityResolveError::NotFound) 531 + ); 532 + assert!( 533 + resolver 534 + .cached(&AtIdentifier::Handle(handle("ptr.pet")), false) 535 + .unwrap() 536 + .is_none() 537 + ); 538 + } 539 + 540 + #[test] 541 + fn cache_capacity_bounds_forward_and_reverse_indexes() { 542 + let resolver = resolver(); 543 + for n in 0..512 { 544 + resolver.observe( 545 + did(&format!("did:plc:user{n}")), 546 + handle(&format!("user{n}.example.com")), 547 + ); 548 + } 549 + 550 + let stats = resolver.stats(); 551 + assert!(stats.entries <= stats.capacity); 552 + assert!(resolver.by_handle.len() <= stats.capacity); 553 + } 554 + }
+5
bobbin/crates/resolver/src/lib.rs
··· 1 + mod identity; 1 2 mod legacy_upgrade; 2 3 mod normalize; 3 4 5 + pub use identity::{ 6 + DEFAULT_IDENTITY_CACHE_ENTRIES, IdentityResolveError, IdentityResolver, 7 + IdentityResolverStatsSnapshot, MiniDoc, 8 + }; 4 9 pub use legacy_upgrade::{ 5 10 DecodedRecord, decode_canon_or_upgrade, decode_canon_or_upgrade_bytes, normalize_record_fields, 6 11 scrub_record_bytes, synthesize_created_at, upgrade, upgrade_wire_bytes,
+3 -3
bobbin/crates/xrpc/src/enrich.rs
··· 339 339 futures::stream::iter(targets) 340 340 .map(|(did, sources)| async move { 341 341 let doc = state 342 - .slingshot 343 - .resolve_mini_doc(&AtIdentifier::Did(did.clone())) 342 + .identity 343 + .resolve_by_did(&did) 344 344 .await 345 345 .ok() 346 - .and_then(|bytes| serde_json::from_slice::<Value>(&bytes).ok()); 346 + .and_then(|doc| serde_json::to_value(doc).ok()); 347 347 (did, sources, doc) 348 348 }) 349 349 .buffer_unordered(MINIDOC_CONCURRENCY)
+3 -10
bobbin/crates/xrpc/src/feed.rs
··· 15 15 use jacquard_common::DefaultStr; 16 16 use jacquard_common::types::string::{AtUri, Datetime, Did, Handle, UriValue}; 17 17 use jacquard_common::xrpc::XrpcResp; 18 - use jacquard_identity::resolver::IdentityResolver; 19 18 20 19 use crate::{ 21 20 AppState, XrpcError, XrpcQuery, fetch, json_stream, paged_tail, parse_cursor, parse_limit, ··· 350 349 351 350 async fn resolve_handle(state: &AppState, did: &Did<DefaultStr>) -> Handle<DefaultStr> { 352 351 state 353 - .directory 354 - .resolve_did_doc_owned(did) 352 + .identity 353 + .resolve_by_did(did) 355 354 .await 356 355 .ok() 357 - .and_then(|doc| { 358 - doc.handles() 359 - .into_iter() 360 - .next() 361 - .map(|h| h.as_str().to_owned()) 362 - }) 363 - .and_then(|s| Handle::new_owned(s).ok()) 356 + .map(|doc| doc.handle) 364 357 .unwrap_or_else(|| { 365 358 Handle::new_static("handle.invalid").expect("handle.invalid is a valid handle") 366 359 })
+24 -6
bobbin/crates/xrpc/src/lib.rs
··· 28 28 }; 29 29 use bobbin_knot_proxy::{KnotHost, KnotProxy, KnotProxyError, ProxyResponse, RepoSlug}; 30 30 use bobbin_record_lru::RecordStore; 31 - use bobbin_resolver::RepoIdResolver; 31 + use bobbin_resolver::{ 32 + DEFAULT_IDENTITY_CACHE_ENTRIES, IdentityResolveError, IdentityResolver, RepoIdResolver, 33 + }; 32 34 use bobbin_runtime::ReqwestHttp; 33 35 use bobbin_search::{ 34 36 SearchCursor, SearchError, SearchFilters, SearchHit, SearchOffset, SearchReader, ··· 134 136 pub knots: Arc<KnotProxy>, 135 137 pub search: Arc<dyn SearchReader>, 136 138 pub resolver: Arc<RepoIdResolver>, 139 + pub identity: Arc<IdentityResolver>, 137 140 pub directory: Arc<Directory>, 138 141 pub limiter: Option<Arc<HeavyLimiter>>, 139 142 pub client_address: Arc<ClientAddress>, ··· 154 157 resolver: Arc<RepoIdResolver>, 155 158 directory: Arc<Directory>, 156 159 ) -> Self { 160 + let identity = Arc::new(IdentityResolver::with_slingshot( 161 + slingshot.clone(), 162 + bobbin_runtime::RuntimeHasher::default(), 163 + DEFAULT_IDENTITY_CACHE_ENTRIES, 164 + )); 157 165 Self { 158 166 records, 159 167 slingshot, ··· 164 172 knots, 165 173 search, 166 174 resolver, 175 + identity, 167 176 directory, 168 177 limiter: None, 169 178 client_address: Arc::new(ClientAddress::default()), ··· 173 182 174 183 pub fn with_limiter(mut self, limiter: Option<Arc<HeavyLimiter>>) -> Self { 175 184 self.limiter = limiter; 185 + self 186 + } 187 + 188 + pub fn with_identity(mut self, identity: Arc<IdentityResolver>) -> Self { 189 + self.identity = identity; 176 190 self 177 191 } 178 192 ··· 2597 2611 State(state): State<AppState>, 2598 2612 XrpcQuery(q): XrpcQuery<ResolveMiniDocParams>, 2599 2613 ) -> Result<Response, XrpcError> { 2600 - let body = state 2601 - .slingshot 2602 - .resolve_mini_doc(&q.identifier) 2614 + let doc = state 2615 + .identity 2616 + .resolve_minidoc(&q.identifier) 2603 2617 .await 2604 - .map_err(map_slingshot)?; 2605 - Ok((StatusCode::OK, [(CONTENT_TYPE, "application/json")], body).into_response()) 2618 + .map_err(|error| match error { 2619 + IdentityResolveError::NotFound => XrpcError::NotFound, 2620 + IdentityResolveError::Upstream(message) => XrpcError::UpstreamUnavailable(message), 2621 + IdentityResolveError::Decode(message) => XrpcError::InvalidRecord(message), 2622 + })?; 2623 + Ok(Json(doc).into_response()) 2606 2624 } 2607 2625 2608 2626 async fn get_coverage(State(state): State<AppState>) -> Json<CoverageEnvelope> {
+8
bobbin/example.toml
··· 135 135 # Default value: 67108864 136 136 #lru_bytes = 67108864 137 137 138 + [identity_cache] 139 + # Maximum number of DID entries retained by the in-process identity cache. 140 + # 141 + # Can also be specified via environment variable `BOBBIN_IDENTITY_CACHE_ENTRIES`. 142 + # 143 + # Default value: 100000 144 + #max_entries = 100000 145 + 138 146 [search] 139 147 # The heap size in bytes for the in-mem tantivy writer. Larger values 140 148 # trade RAM for fewer segment merges - the index itself lives in