From b76c0c484533babb8a7ec13a3ce6dc43d6e9f3e7 Mon Sep 17 00:00:00 2001 From: dawn Date: Sat, 12 Sep 2026 07:24:41 +0300 Subject: [PATCH] bobbin/{resolver,ingest}: re-verify deactivated identities on read Signed-off-by: dawn --- bobbin/crates/ingest/src/lib.rs | 111 +++++++++++++++++++++ bobbin/crates/resolver/src/identity.rs | 133 ++++++++++++++++++++++++- 2 files changed, 242 insertions(+), 2 deletions(-) diff --git a/bobbin/crates/ingest/src/lib.rs b/bobbin/crates/ingest/src/lib.rs index fdc4de479..15484f0f6 100644 --- a/bobbin/crates/ingest/src/lib.rs +++ b/bobbin/crates/ingest/src/lib.rs @@ -1800,12 +1800,15 @@ mod tests { use bobbin_edge_index::Coverage; use bobbin_record_lru::{CacheCapacity, LruRecordStore, NoopRecordStore, RecordStore}; use bobbin_runtime::{OsEntropy, RuntimeHasher, SystemClock, WsConnectFuture, WsTransport}; + use bobbin_slingshot_client::SlingshotClient; use bobbin_types::search::NoopSearchSink; use jacquard_common::types::datetime::Datetime; use jacquard_common::types::handle::Handle; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::tid::Tid; use serde_json::json; + use wiremock::matchers::{method, path, query_param}; + use wiremock::{Mock, MockServer, ResponseTemplate}; const VALID_CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; @@ -1908,6 +1911,21 @@ mod tests { IdentityResolver::detached(RuntimeHasher::default()) } + async fn drain_identity_warming(identity: &Arc) { + let settled_before = identity.stats().warm_resolved + identity.stats().warm_failed; + let cancel = CancellationToken::new(); + let warmer = tokio::spawn(identity.clone().run_warming(cancel.clone())); + tokio::time::timeout(Duration::from_secs(5), async { + while identity.stats().warm_resolved + identity.stats().warm_failed == settled_before { + tokio::task::yield_now().await; + } + }) + .await + .expect("identity warmer never settled"); + cancel.cancel(); + warmer.await.unwrap(); + } + fn parse_frame(value: serde_json::Value) -> HydrantFrame { let text = serde_json::to_string(&value).expect("serialize fixture"); serde_json::from_str(&text).expect("deserialize fixture") @@ -2822,6 +2840,99 @@ mod tests { assert!(!cov.snapshot().is_ready()); } + #[tokio::test] + async fn identity_frames_refresh_incomplete_docs_and_observe_complete_ones() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/blue.microcosm.identity.resolveMiniDoc")) + .and(query_param("identifier", "did:plc:cached")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + "did": "did:plc:cached", + "handle": "cached.example.com", + "pds": "https://cached-pds.example.com" + }))) + .expect(1) + .mount(&server) + .await; + Mock::given(method("GET")) + .and(path("/xrpc/blue.microcosm.identity.resolveMiniDoc")) + .and(query_param("identifier", "did:plc:inactive")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + "did": "did:plc:inactive", + "handle": "inactive.example.com", + "pds": "https://inactive-pds.example.com" + }))) + .expect(1) + .mount(&server) + .await; + + let (store, issue_states, pull_statuses, coverage, resolver) = fresh(); + let identity = Arc::new(IdentityResolver::with_slingshot( + SlingshotClient::with_default_http(Url::parse(&server.uri()).unwrap()).unwrap(), + Arc::new(SystemClock::new()), + RuntimeHasher::default(), + )); + let records = NoopRecordStore; + let search = NoopSearchSink; + let ctx = PipelineCtx { + resolver: &resolver, + identity: &identity, + 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, + settlements: None, + }; + let apply = async |frame: HydrantFrame| { + let pending = prepare_frame(frame, &ctx, now()).await; + let pending = resolve_pending(pending, &ctx).await; + commit_pending(pending, &ctx).await; + }; + + let cached = Did::new_static("did:plc:cached").unwrap(); + identity + .resolve_minidoc(&AtIdentifier::Did(cached.clone())) + .await + .unwrap(); + apply(parse_frame(json!({ + "id": 1, + "type": "identity", + "identity": {"did": "did:plc:cached", "handle": "renamed.example.com"} + }))) + .await; + let cached_doc = identity.get_by_did(&cached).unwrap(); + assert_eq!(cached_doc.handle.as_ref(), "renamed.example.com"); + assert_eq!( + cached_doc.pds, + Some(Url::parse("https://cached-pds.example.com").unwrap()) + ); + assert_eq!(identity.stats().warm_queued, 0); + + let inactive = Did::new_static("did:plc:inactive").unwrap(); + identity.deactivate(inactive.clone()); + apply(parse_frame(json!({ + "id": 2, + "type": "identity", + "identity": {"did": "did:plc:inactive", "handle": "frame.example.com"} + }))) + .await; + assert_eq!(identity.stats().warm_queued, 1); + drain_identity_warming(&identity).await; + + let rehydrated = identity.get_by_did(&inactive).unwrap(); + assert_eq!(rehydrated.handle.as_ref(), "inactive.example.com"); + assert_eq!( + rehydrated.pds, + Some(Url::parse("https://inactive-pds.example.com").unwrap()) + ); + } + #[tokio::test] async fn identity_and_account_frames_drive_the_identity_resolver() { let (store, issue_states, pull_statuses, coverage, resolver) = fresh(); diff --git a/bobbin/crates/resolver/src/identity.rs b/bobbin/crates/resolver/src/identity.rs index c18768897..2979e38da 100644 --- a/bobbin/crates/resolver/src/identity.rs +++ b/bobbin/crates/resolver/src/identity.rs @@ -329,6 +329,20 @@ impl IdentityResolver { } pub fn observe(&self, did: Did, handle: Handle) { + let (inactive, incomplete) = self + .by_did + .read_sync(&did, |_, state| match state { + IdentityState::Inactive => (true, false), + IdentityState::Cached(doc) => (false, doc.pds.is_none()), + }) + .unwrap_or((false, false)); + if inactive { + self.refresh(&did); + return; + } + if incomplete { + self.refresh(&did); + } self.observe_doc(MiniDoc { did, handle, @@ -443,7 +457,10 @@ impl IdentityResolver { match identifier { AtIdentifier::Did(did) => match self.by_did.get_sync(did).as_deref() { Some(IdentityState::Cached(doc)) => Ok(Some(doc.clone())), - Some(IdentityState::Inactive) => Err(IdentityResolveError::NotFound), + Some(IdentityState::Inactive) => { + self.refresh(did); + Ok(None) + } None => Ok(None), }, AtIdentifier::Handle(handle) => { @@ -544,7 +561,10 @@ impl IdentityResolver { .expect("identity mutation mutex poisoned"); let (did, handle) = match self.by_did.entry_sync(doc.did.clone()) { MapEntry::Occupied(occupied) => match occupied.get() { - IdentityState::Inactive => return Err(IdentityResolveError::NotFound), + IdentityState::Inactive => { + self.refresh(&doc.did); + return Err(IdentityResolveError::NotFound); + } IdentityState::Cached(current) => return Ok(current.clone()), }, MapEntry::Vacant(vacant) => { @@ -1006,6 +1026,115 @@ mod tests { assert_eq!(resolver.stats().warm_failed, 1); } + #[tokio::test] + async fn a_deactivated_did_reverifies_on_read() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/blue.microcosm.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 resolver = hydrant_backed(&server); + let identity = did("did:plc:dawn"); + resolver.deactivate(identity.clone()); + + assert_eq!( + resolver.get_by_did(&identity), + Err(IdentityResolveError::NotFound), + "the read queues a revalidation instead of waiting on it" + ); + assert_eq!(resolver.stats().warm_queued, 1); + + drain_warming(&resolver, &identity).await; + + let doc = resolver.get_by_did(&identity).unwrap(); + assert_eq!(doc.handle, handle("ptr.pet")); + assert_eq!( + doc.pds, + Some(Url::parse("https://pds.example.com").unwrap()) + ); + } + + #[tokio::test] + async fn a_deactivated_did_that_stays_dead_is_not_refetched_per_read() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/blue.microcosm.identity.resolveMiniDoc")) + .respond_with(ResponseTemplate::new(404)) + .expect(1) + .mount(&server) + .await; + + let resolver = hydrant_backed(&server); + let identity = did("did:plc:gone"); + resolver.deactivate(identity.clone()); + + for _ in 0..3 { + assert_eq!( + resolver.get_by_did(&identity), + Err(IdentityResolveError::NotFound) + ); + } + assert_eq!(resolver.stats().warm_queued, 1); + drain_warming(&resolver, &identity).await; + + for _ in 0..3 { + assert_eq!( + resolver.get_by_did(&identity), + Err(IdentityResolveError::NotFound) + ); + } + let stats = resolver.stats(); + assert_eq!(stats.warm_queued, 1); + assert_eq!(stats.warm_failed, 1); + } + + #[tokio::test] + async fn a_handle_lookup_for_a_deactivated_did_queues_a_revalidation() { + let server = MockServer::start().await; + for identifier in ["dawn.example.com", "did:plc:dawn"] { + Mock::given(method("GET")) + .and(path("/xrpc/blue.microcosm.identity.resolveMiniDoc")) + .and(query_param("identifier", identifier)) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "did": "did:plc:dawn", + "handle": "dawn.example.com", + "pds": "https://pds.example.com" + }))) + .expect(1) + .mount(&server) + .await; + } + + let resolver = hydrant_backed(&server); + let identity = did("did:plc:dawn"); + resolver.observe(identity.clone(), handle("dawn.example.com")); + resolver.deactivate(identity.clone()); + + assert_eq!( + resolver + .resolve_minidoc(&AtIdentifier::Handle(handle("dawn.example.com"))) + .await, + Err(IdentityResolveError::NotFound) + ); + assert_eq!(resolver.stats().warm_queued, 1); + + drain_warming(&resolver, &identity).await; + + let doc = resolver + .resolve_minidoc(&AtIdentifier::Handle(handle("dawn.example.com"))) + .await + .unwrap(); + assert_eq!(doc.handle, handle("dawn.example.com")); + } + #[tokio::test] async fn warming_a_dead_repo_caches_no_handle() { let server = MockServer::start().await; -- 2.51.2