diff --git a/src/api/xrpc/resolve_mini_doc.rs b/src/api/xrpc/resolve_mini_doc.rs index 22c6677..cc187bb 100644 --- a/src/api/xrpc/resolve_mini_doc.rs +++ b/src/api/xrpc/resolve_mini_doc.rs @@ -84,10 +84,14 @@ pub(super) async fn resolve_mini_doc( .await .map_err(|e| bad_request(nsid, format!("can't resolve identifier: {e}")))?; + let requested_handle = match identifier { + AtIdentifier::Handle(handle) => Some(handle), + AtIdentifier::Did(_) => None, + }; let doc = hydrant .repos .get(&did) - .mini_doc() + .mini_doc_for(requested_handle) .await .map_err(|e| match e { MiniDocError::NotSynced => bad_request(nsid, "repo not synced"), diff --git a/src/backfill/error.rs b/src/backfill/error.rs index 36c6221..04b6e2c 100644 --- a/src/backfill/error.rs +++ b/src/backfill/error.rs @@ -72,7 +72,22 @@ impl From for BackfillError { match e { ResolverError::Ratelimited => Self::Ratelimited, ResolverError::Transport(s) => Self::Transport(s), - ResolverError::Generic(e) => Self::Generic(e), + ResolverError::HandleResolutionExhausted(e) | ResolverError::Generic(e) => { + Self::Generic(e) + } } } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn exhausted_handle_resolution_is_not_a_backfill_transport_failure() { + let error = BackfillError::from(ResolverError::HandleResolutionExhausted(miette::miette!( + "handle not found" + ))); + assert!(matches!(error, BackfillError::Generic(_))); + } +} diff --git a/src/control/repos/mod.rs b/src/control/repos/mod.rs index bbbf327..ccd640c 100644 --- a/src/control/repos/mod.rs +++ b/src/control/repos/mod.rs @@ -187,6 +187,61 @@ mod tests { use super::*; use miette::IntoDiagnostic; + #[test] + fn identity_reconcile_trigger_matrix_is_narrow() -> miette::Result<()> { + let mut complete = RepoState::synced(); + complete.pds = Some("https://pds.example".into()); + complete.handle = Some(Handle::new_static("alice.test").into_diagnostic()?); + complete.signing_key = Some(DidKey::from_did_key( + "did:key:zQ3shokFTS3brHcDQrn82RUDfCZESWL1ZdCEJwekUDPQiYBme", + )?); + let same_handle = complete.handle.as_ref(); + + assert!(!needs_identity_reconcile(&complete, None, false)); + assert!(!needs_identity_reconcile(&complete, same_handle, false)); + assert!(needs_identity_reconcile( + &complete, + Some(&Handle::new_static("carol.test").into_diagnostic()?), + false, + )); + + for missing in [ + { + let mut state = complete.clone(); + state.pds = None; + state + }, + { + let mut state = complete.clone(); + state.signing_key = None; + state + }, + ] { + assert!(needs_identity_reconcile(&missing, None, false)); + } + + let mut deactivated = complete.clone(); + deactivated.status = RepoStatus::Deactivated; + assert!(needs_identity_reconcile(&deactivated, None, false)); + + let mut inactive_synced = complete.clone(); + inactive_synced.active = false; + assert!(!needs_identity_reconcile(&inactive_synced, None, false)); + assert!(needs_identity_reconcile(&complete, None, true)); + + let mut deleted = complete; + deleted.status = RepoStatus::Deleted; + assert!(!needs_identity_reconcile(&deleted, None, true)); + deleted.pds = None; + deleted.signing_key = None; + assert!(!needs_identity_reconcile( + &deleted, + Some(&Handle::new_static("carol.test").into_diagnostic()?), + true, + )); + Ok(()) + } + #[test] fn repo_info_renders_message_time_from_millis() -> miette::Result<()> { let mut state = RepoState::backfilling(); @@ -304,6 +359,21 @@ pub struct MiniDoc<'i> { pub signing_key: DidKey<'i>, } +struct MiniDocSnapshot { + state: RepoState<'static>, + pending: bool, + gone: bool, +} + +struct MiniDocReconcile { + state: RepoState<'static>, + pending: bool, + accepted_doc: bool, + #[cfg_attr(not(feature = "indexer_stream"), allow(dead_code))] + identity_changed: bool, + queued_backfill: bool, +} + /// handle to access data related to this repository. #[derive(Clone)] pub struct RepoHandle { @@ -311,6 +381,76 @@ pub struct RepoHandle { pub did: Did, } +fn needs_identity_reconcile( + state: &RepoState<'_>, + requested_handle: Option<&Handle>, + gone: bool, +) -> bool { + if state.status == RepoStatus::Deleted { + return false; + } + + gone || state.status == RepoStatus::Deactivated + || requested_handle.is_some_and(|requested| state.handle.as_ref() != Some(requested)) + || state.pds.is_none() + || state.signing_key.is_none() +} + +fn read_mini_doc_snapshot(db: &Db, did: &Did) -> Result> { + let did_key = keys::repo_key(did); + let Some(state_bytes) = db.repos.get(&did_key).into_diagnostic()? else { + return Ok(None); + }; + let state = db::deser_repo_state(&state_bytes)?.into_static(); + + #[cfg(feature = "indexer")] + let pending = { + let metadata_key = keys::repo_metadata_key(did); + let metadata = db + .repo_metadata + .get(&metadata_key) + .into_diagnostic()? + .map(|bytes| db::deser_repo_meta(&bytes)) + .transpose()?; + metadata + .map(|metadata| { + db.indexer + .pending + .get(keys::pending_key(metadata.index_id)) + .into_diagnostic() + .map(|value| value.is_some()) + }) + .transpose()? + .unwrap_or(false) + }; + #[cfg(not(feature = "indexer"))] + let pending = false; + + #[cfg(feature = "indexer")] + let gone = db + .indexer + .resync + .get(&did_key) + .into_diagnostic()? + .map(|bytes| rmp_serde::from_slice::(&bytes).into_diagnostic()) + .transpose()? + .is_some_and(|resync| { + matches!( + resync, + crate::types::ResyncState::Gone { status } + if status != crate::types::RepoStatus::Deleted + ) + }); + #[cfg(not(feature = "indexer"))] + let gone = false; + + Ok(Some(MiniDocSnapshot { + state, + pending, + gone, + })) +} + impl RepoHandle { pub(crate) async fn state(&self) -> Result>> { let did_key = keys::repo_key(&self.did); @@ -396,42 +536,159 @@ impl RepoHandle { /// returns a bi-directionally validated mini doc. pub async fn mini_doc(&self) -> Result, MiniDocError> { - let Some(info) = self.info().await.map_err(MiniDocError::Other)? else { + self.mini_doc_for(None).await + } + + pub(crate) async fn mini_doc_for( + &self, + requested_handle: Option<&Handle>, + ) -> Result, MiniDocError> { + let requested_handle = requested_handle.cloned(); + let snapshot = self + .mini_doc_snapshot() + .await + .map_err(MiniDocError::Other)?; + let Some(snapshot) = snapshot else { return Err(MiniDocError::RepoNotFound); }; - // check if repo is still backfilling (in pending) - #[cfg(feature = "indexer")] - let is_pending = { - let metadata_key = keys::repo_metadata_key(&self.did); - self.state - .db - .run(move |db| { - let metadata_bytes = db.repo_metadata.get(&metadata_key).into_diagnostic()?; - let Some(metadata_bytes) = metadata_bytes else { - return Ok(false); - }; - - let metadata = crate::db::deser_repo_meta(metadata_bytes.as_ref())?; - Ok(db - .indexer - .pending - .get(crate::db::keys::pending_key(metadata.index_id)) - .into_diagnostic()? - .is_some()) - }) - .await - .map_err(MiniDocError::Other)? + if snapshot.pending { + return Err(MiniDocError::NotSynced); + } + + let needs_reconcile = + needs_identity_reconcile(&snapshot.state, requested_handle.as_ref(), snapshot.gone); + if !needs_reconcile { + return self.render_mini_doc(snapshot.state).await; + } + + // resolution is deliberately outside the repository write lock. the + // state is reloaded and the trigger is checked again before applying + // this result, so a slower concurrent resolver cannot overwrite a + // repair that already won the race. + let fresh = self + .state + .resolver + .resolve_doc_fresh(&self.did) + .await + .map_err(|_| MiniDocError::CouldNotResolveIdentity)?; + let fresh_for_db = fresh.clone(); + let did = self.did.clone().into_static(); + let did_for_db = did.clone(); + let initial_identity = snapshot.state.identity(); + let initial_pending = snapshot.pending; + let result = self + .state + .db + .run(move |db| { + #[cfg(feature = "indexer")] + let mut txn = crate::db::Txn::new(db); + #[cfg(feature = "indexer")] + txn.hold_repo_write_lock(&did_for_db); + #[cfg(not(feature = "indexer"))] + let mut batch = db.inner.batch(); + + let Some(MiniDocSnapshot { + mut state, + pending, + gone, + }) = read_mini_doc_snapshot(db, &did_for_db)? + else { + return Ok(None); + }; + let did_key = keys::repo_key(&did_for_db); + + let raced = state.identity() != initial_identity; + let current_trigger = + needs_identity_reconcile(&state, requested_handle.as_ref(), gone); + // a pending transition made while the network request was in + // flight belongs to another path. preserve the old read + // contract rather than applying a result to a backfill. + let pending_race = pending && !initial_pending; + let accepted_doc = !raced && !pending_race && current_trigger; + let identity_changed = if accepted_doc { + state.update_from_doc(fresh_for_db) + } else { + false + }; + + #[cfg(feature = "indexer")] + if identity_changed { + txn.batch + .insert(&db.repos, &did_key, crate::db::ser_repo_state(&state)?); + } + #[cfg(not(feature = "indexer"))] + if identity_changed { + batch.insert(&db.repos, &did_key, crate::db::ser_repo_state(&state)?); + } + + // a DID document never revives lifecycle state. gone recovery + // only moves the exact durable Gone row into the normal + // pending queue; BackfillFinished performs the eventual + // Synced transition. + #[cfg(feature = "indexer")] + let queued_backfill = if gone && !pending { + ReposControl::_resync(db, &did_for_db, &mut txn)? + } else { + false + }; + #[cfg(not(feature = "indexer"))] + let queued_backfill = false; + + if identity_changed || queued_backfill { + #[cfg(feature = "indexer")] + txn.commit()?; + #[cfg(not(feature = "indexer"))] + batch.commit().into_diagnostic()?; + } + + Ok(Some(MiniDocReconcile { + state, + pending, + accepted_doc, + identity_changed, + queued_backfill, + })) + }) + .await + .map_err(MiniDocError::Other)?; + let Some(result) = result else { + return Err(MiniDocError::RepoNotFound); }; - #[cfg(not(feature = "indexer"))] - let is_pending = false; - if is_pending { + if result.pending && !result.queued_backfill { return Err(MiniDocError::NotSynced); } + if result.accepted_doc { + self.state.resolver.cache_doc(did.clone(), fresh).await; + } + #[cfg(feature = "indexer_stream")] + if result.identity_changed { + crate::ops::announce_identity(&self.state.db, &did, result.state.handle.clone()); + } + #[cfg(feature = "indexer")] + if result.queued_backfill { + self.state.notify_backfill(); + } + self.render_mini_doc(result.state).await + } + + async fn mini_doc_snapshot(&self) -> Result> { + let did = self.did.clone().into_static(); + self.state + .db + .run(move |db| read_mini_doc_snapshot(db, &did)) + .await + } + + async fn render_mini_doc( + &self, + info: RepoState<'static>, + ) -> Result, MiniDocError> { let pds = info .pds + .and_then(|pds| Url::parse(pds.as_ref()).ok()) .ok_or_else(|| MiniDocError::CouldNotResolveIdentity)?; let signing_key = info .signing_key diff --git a/src/resolver.rs b/src/resolver.rs index ad78777..4e96d39 100644 --- a/src/resolver.rs +++ b/src/resolver.rs @@ -35,6 +35,8 @@ pub enum ResolverError { Generic(miette::Report), #[error("too many requests")] Ratelimited, + #[error("handle resolution exhausted: {0}")] + HandleResolutionExhausted(miette::Report), #[error("transport error: {0}")] Transport(SmolStr), } @@ -46,9 +48,12 @@ impl From for ResolverError { Self::Ratelimited } IdentityErrorKind::Transport(msg) => Self::Transport(msg.clone()), - // a timeout is the network failing to answer, not an answer. keeping it - // out of `Generic` is what lets callers read `Generic` as a verdict. + // a timeout is the network failing to answer, not an answer. keep it + // in the transport category rather than treating it as handle exhaustion. IdentityErrorKind::Timeout => Self::Transport("timed out".into()), + IdentityErrorKind::HandleResolutionExhausted => { + Self::HandleResolutionExhausted(e.into()) + } _ => Self::Generic(e.into()), } } @@ -170,6 +175,23 @@ impl Resolver { return Ok(entry.get().clone()); } + let mini = self.resolve_doc_uncached(did).await?; + let _ = self.inner.cache.put_async(did_static, mini.clone()).await; + Ok(mini) + } + + /// resolves the current DID document without reading or replacing the + /// cached mini document. callers that discard a raced result leave the + /// cache untouched. + pub async fn resolve_doc_fresh(&self, did: &Did) -> Result { + self.resolve_doc_uncached(did).await + } + + pub(crate) async fn cache_doc(&self, did: Did, doc: MiniDoc) { + let _ = self.inner.cache.entry_async(did).await.put_entry(doc); + } + + async fn resolve_doc_uncached(&self, did: &Did) -> Result { let doc_resp = self .req(did.starts_with("did:plc:"), |j| j.resolve_did_doc(did)) .await?; @@ -186,13 +208,11 @@ impl Resolver { .then(|| handles.remove(0).into_static()); let key = doc.atproto_public_key().ok().flatten(); - let mini = MiniDoc { + Ok(MiniDoc { pds: Url::from_str(pds.as_str()).expect("that url is valid"), handle, key, - }; - let _ = self.inner.cache.put_async(did_static, mini.clone()).await; - Ok(mini) + }) } pub async fn resolve_identity_info( @@ -244,16 +264,15 @@ impl Resolver { /// returns `true` if the given handle bi-directionally resolves to `did`. /// - /// a handle nothing resolves to is an unverified handle, not an outage of this - /// service, so it answers `false` instead of erroring: the caller's fallback is - /// `handle.invalid`, which is the spec's representation for exactly this - /// (atproto.com/specs/handle#invalid-handles). transient failures stay errors, - /// because an unanswered question is not a negative answer. + /// only an explicit exhausted-handle-resolution result answers `false`; all + /// other resolver failures remain errors. the caller's fallback is + /// `handle.invalid`, the representation for an unverified handle + /// (atproto.com/specs/handle#invalid-handles). pub async fn verify_handle(&self, did: &Did, handle: &Handle) -> Result { let id = AtIdentifier::Handle(handle.clone()); match self.resolve_did(&id).await { Ok(resolved_did) => Ok(resolved_did.as_str() == did.as_str()), - Err(ResolverError::Generic(_)) => Ok(false), + Err(ResolverError::HandleResolutionExhausted(_)) => Ok(false), Err(e) => Err(e), } } @@ -267,12 +286,14 @@ pub struct NoSigningKeyError(Did); mod tests { use super::*; use axum::{Router, http::header, response::IntoResponse}; + use std::sync::Arc; + use tokio::sync::RwLock; #[test] - fn resolution_verdicts_stay_apart_from_transient_failures() { + fn resolution_error_categories_are_preserved() { assert!(matches!( ResolverError::from(IdentityError::handle_resolution_exhausted()), - ResolverError::Generic(_) + ResolverError::HandleResolutionExhausted(_) )); assert!(matches!( ResolverError::from(IdentityError::timeout()), @@ -286,6 +307,24 @@ mod tests { )); } + async fn spawn_mutable_doc_server(body: String) -> (Url, Arc>) { + let current = Arc::new(RwLock::new(body)); + let served = current.clone(); + let app = Router::new().fallback(move || { + let served = served.clone(); + async move { + let body = served.read().await.clone(); + ([(header::CONTENT_TYPE, "application/json")], body).into_response() + } + }); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + (Url::parse(&format!("http://{addr}/")).unwrap(), current) + } + async fn spawn_doc_server(body: String) -> Url { let app = Router::new().fallback(move || { let body = body.clone(); @@ -299,6 +338,46 @@ mod tests { Url::parse(&format!("http://{addr}/")).unwrap() } + #[tokio::test] + async fn fresh_resolution_bypasses_cache_until_explicitly_installed() { + let did = Did::new_static("did:plc:ewvi7nxzyoun6zhxrhs64oiz").unwrap(); + let doc = |pds: &str| { + serde_json::json!({ + "@context": ["https://www.w3.org/ns/did/v1"], + "id": did.as_str(), + "alsoKnownAs": ["at://alice.test"], + "service": [{ + "id": "#atproto_pds", + "type": "AtprotoPersonalDataServer", + "serviceEndpoint": pds + }], + "verificationMethod": [] + }) + .to_string() + }; + let (url, current) = spawn_mutable_doc_server(doc("https://pds-a.example")).await; + let resolver = Resolver::new(vec![url], 10); + + assert_eq!( + resolver.resolve_doc(&did).await.unwrap().pds.as_str(), + "https://pds-a.example/" + ); + *current.write().await = doc("https://pds-b.example"); + + let fresh = resolver.resolve_doc_fresh(&did).await.unwrap(); + assert_eq!(fresh.pds.as_str(), "https://pds-b.example/"); + assert_eq!( + resolver.resolve_doc(&did).await.unwrap().pds.as_str(), + "https://pds-a.example/" + ); + + resolver.cache_doc(did.clone(), fresh).await; + assert_eq!( + resolver.resolve_doc(&did).await.unwrap().pds.as_str(), + "https://pds-b.example/" + ); + } + #[tokio::test] async fn raw_did_doc_json_contract_and_size_limit() { let did = Did::new_static("did:plc:ewvi7nxzyoun6zhxrhs64oiz").unwrap();