very fast at protocol indexer with flexible filtering, xrpc queries, cursor-backed event stream, and more, built on fjall
rust fjall at-protocol atproto indexer
Something went wrong. Try again.
Rust
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476use std::ops::Not;use std::str::FromStr;use std::sync::Arc;use std::sync::atomic::{AtomicUsize, Ordering};use std::time::Duration;
use jacquard_common::Data;use jacquard_common::IntoStatic;use jacquard_common::types::crypto::PublicKey;use jacquard_common::types::ident::AtIdentifier;use jacquard_common::types::string::Did;use jacquard_common::types::string::Handle;use jacquard_identity::JacquardResolver;use jacquard_identity::resolver::{ IdentityError, IdentityErrorKind, IdentityResolver, PlcSource, ResolverOptions,};use miette::{Diagnostic, IntoDiagnostic};use reqwest::StatusCode;use scc::HashCache;use smol_str::SmolStr;use thiserror::Error;use url::Url;
use crate::net::{PublicHttpClient, public_http};use crate::util::url_to_fluent_uri;
// as per spec: https://web.plc.directory/spec/v0.1/did-plc// "As an anti-abuse mechanism, operations have a maximum size when encoded as// DAG-CBOR. The current limit is 7500 bytes."const MAX_DID_DOC_BYTES: usize = 16384; // this is generous i think but we deal with json so..
#[derive(Debug, Diagnostic, Error)]pub enum ResolverError { #[error("{0}")] Generic(miette::Report), #[error("too many requests")] Ratelimited, #[error("handle resolution exhausted: {0}")] HandleResolutionExhausted(miette::Report), #[error("transport error: {0}")] Transport(SmolStr), #[error("no PDS service found in DID Doc for {0}")] MissingPds(Did),}
impl From<IdentityError> for ResolverError { fn from(e: IdentityError) -> Self { match e.kind() { IdentityErrorKind::HttpStatus(reqwest::StatusCode::TOO_MANY_REQUESTS) => { Self::Ratelimited } IdentityErrorKind::Transport(msg) => Self::Transport(msg.clone()), // 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()), } }}
impl From<miette::Report> for ResolverError { fn from(report: miette::Report) -> Self { ResolverError::Generic(report) }}
#[derive(Clone)]pub struct MiniDoc { pub pds: Url, pub handle: Option<Handle>, pub key: Option<PublicKey<'static>>,}
struct ResolverInner { jacquards: Vec<JacquardResolver<PublicHttpClient>>, next_idx: AtomicUsize, cache: HashCache<Did, MiniDoc>,}
#[derive(Clone)]pub struct Resolver { inner: Arc<ResolverInner>,}
impl Resolver { pub fn new(plc_urls: Vec<Url>, identity_cache_size: u64) -> Self { let options = plc_urls .into_iter() .map(|url| ResolverOptions { plc_source: PlcSource::PlcDirectory { base: url_to_fluent_uri(&url), }, request_timeout: Some(Duration::from_secs(3)), ..Default::default() }) .collect(); Self::with_options(options, identity_cache_size) }
#[cfg(all(test, feature = "indexer_stream"))] pub(crate) fn slingshot_for_test( base: Url, identity_cache_size: u64, public: PublicHttpClient, ) -> Self { use jacquard_identity::resolver::{DidStep, HandleStep};
Self::with_options_and_client( vec![ResolverOptions { plc_source: PlcSource::Slingshot { base: url_to_fluent_uri(&base), }, handle_order: vec![HandleStep::PdsResolveHandle], did_order: vec![DidStep::PlcHttp], public_fallback_for_handle: false, request_timeout: Some(Duration::from_secs(3)), ..Default::default() }], identity_cache_size, public, ) }
fn with_options(options: Vec<ResolverOptions>, identity_cache_size: u64) -> Self { Self::with_options_and_client(options, identity_cache_size, public_http()) }
pub(crate) fn with_options_and_client( options: Vec<ResolverOptions>, identity_cache_size: u64, public: PublicHttpClient, ) -> Self { assert!(!options.is_empty(), "at least one PLC URL must be provided"); let jacquards = options .into_iter() .map(|options| JacquardResolver::new(public.clone(), options).with_system_dns()) .collect();
Self { inner: Arc::new(ResolverInner { jacquards, next_idx: AtomicUsize::new(0), cache: HashCache::with_capacity( std::cmp::min(1000, (identity_cache_size / 100) as usize), identity_cache_size as usize, ), }), } }
pub async fn invalidate(&self, did: &Did) { self.inner.cache.remove_async(did).await; }
pub fn invalidate_sync(&self, did: &Did) { self.inner.cache.remove_sync(did); }
async fn req<'r, T, Fut>( &'r self, is_plc: bool, f: impl Fn(&'r JacquardResolver<PublicHttpClient>) -> Fut, ) -> Result<T, ResolverError> where Fut: Future<Output = Result<T, IdentityError>>, { let mut idx = self.inner.next_idx.fetch_add(1, Ordering::Relaxed) % self.inner.jacquards.len(); let mut try_count = 0; loop { let res = f(&self.inner.jacquards[idx]).await; try_count += 1; // retry these with the different plc resolvers if is_plc { let is_retriable = matches!( res.as_ref().map_err(|e| e.kind()), Err(IdentityErrorKind::HttpStatus(StatusCode::TOO_MANY_REQUESTS) | IdentityErrorKind::Transport(_)) ); // check if retriable and we haven't gone through all the plc resolvers if is_retriable && try_count < self.inner.jacquards.len() { idx = (idx + 1) % self.inner.jacquards.len(); continue; } } return res.map_err(Into::into); } }
pub async fn resolve_did(&self, identifier: &AtIdentifier) -> Result<Did, ResolverError> { match identifier { AtIdentifier::Did(did) => Ok(did.clone()), AtIdentifier::Handle(handle) => { let did = self.req(false, |j| j.resolve_handle(handle)).await?; Ok(did) } } }
pub async fn resolve_doc(&self, did: &Did) -> Result<MiniDoc, ResolverError> { let did_static = did.clone(); if let Some(entry) = self.inner.cache.get_async(&did_static).await { 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<MiniDoc, ResolverError> { 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<MiniDoc, ResolverError> { let doc_resp = self .req(did.starts_with("did:plc:"), |j| j.resolve_did_doc(did)) .await?; let doc = doc_resp.parse()?;
let pds = doc .pds_endpoint() .ok_or_else(|| ResolverError::MissingPds(did.clone().into_static()))?;
let mut handles = doc.handles(); let handle = handles .is_empty() .not() .then(|| handles.remove(0).into_static()); let key = doc.atproto_public_key().ok().flatten();
Ok(MiniDoc { pds: Url::from_str(pds.as_str()).expect("that url is valid"), handle, key, }) }
pub async fn resolve_identity_info( &self, did: &Did, ) -> Result<(Url, Option<Handle>), ResolverError> { let mini = self.resolve_doc(did).await?; Ok((mini.pds, mini.handle)) }
pub async fn resolve_signing_key( &self, did: &Did, ) -> Result<PublicKey<'static>, ResolverError> { let did = did.clone(); let mini = self.resolve_doc(&did).await?; Ok(mini.key.ok_or(NoSigningKeyError(did)).into_diagnostic()?) }
/// resolves the full DID document as raw [`Data`] without caching, and /// extracts the first `alsoKnownAs` handle if present. pub async fn resolve_raw_doc( &self, did: &Did, ) -> Result<(Data, Option<Handle>), ResolverError> { let doc_resp = self .req(did.starts_with("did:plc:"), |j| j.resolve_did_doc(did)) .await?;
if doc_resp.buffer.len() > MAX_DID_DOC_BYTES { return Err(ResolverError::Generic(miette::miette!( "DID doc response exceeds {} byte limit", MAX_DID_DOC_BYTES ))); }
let doc = doc_resp.parse_validated()?; let mut handles = doc.handles(); let handle = handles .is_empty() .not() .then(|| handles.remove(0).into_static());
let data: Data = serde_json::from_slice(&doc_resp.buffer) .map_err(|e| ResolverError::Generic(miette::miette!("failed to parse DID doc: {e}")))?;
Ok((data, handle)) }
/// returns `true` if the given handle bi-directionally resolves to `did`. /// /// 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<bool, ResolverError> { 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::HandleResolutionExhausted(_)) => Ok(false), Err(e) => Err(e), } }}
#[derive(Debug, Diagnostic, Error)]#[error("no atproto signing key in DID doc for {0}")]pub struct NoSigningKeyError(Did);
#[cfg(test)]mod tests { use super::*; use axum::{Router, http::header, response::IntoResponse}; use std::sync::Arc; use tokio::sync::RwLock;
#[test] fn resolution_error_categories_are_preserved() { assert!(matches!( ResolverError::from(IdentityError::handle_resolution_exhausted()), ResolverError::HandleResolutionExhausted(_) )); assert!(matches!( ResolverError::from(IdentityError::timeout()), ResolverError::Transport(_) )); assert!(matches!( ResolverError::from(IdentityError::http_status( reqwest::StatusCode::TOO_MANY_REQUESTS )), ResolverError::Ratelimited )); }
async fn spawn_mutable_doc_server(body: String) -> (std::net::SocketAddr, Arc<RwLock<String>>) { 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(); }); (addr, current) }
async fn spawn_doc_server(body: String) -> std::net::SocketAddr { let app = Router::new().fallback(move || { let body = body.clone(); async move { ([(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(); }); addr }
fn test_public_client(address: std::net::SocketAddr) -> PublicHttpClient { let client = reqwest::Client::builder() .no_proxy() .resolve("public.example", address) .build() .unwrap(); PublicHttpClient::test(client) }
fn test_options() -> ResolverOptions { ResolverOptions { plc_source: PlcSource::PlcDirectory { base: url_to_fluent_uri(&Url::parse("http://public.example/").unwrap()), }, request_timeout: Some(Duration::from_secs(3)), ..Default::default() } }
#[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 (address, current) = spawn_mutable_doc_server(doc("https://pds-a.example")).await; let resolver = Resolver::with_options_and_client( vec![test_options()], 10, test_public_client(address), );
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(); let body = 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": "https://pds.example" }], "verificationMethod": [] }) .to_string(); let address = spawn_doc_server(body).await; let resolver = Resolver::with_options_and_client( vec![test_options()], 10, test_public_client(address), );
let (data, handle) = resolver.resolve_raw_doc(&did).await.unwrap(); let json = serde_json::to_value(data).unwrap(); assert_eq!(json["id"], did.as_str()); assert_eq!(json["service"][0]["serviceEndpoint"], "https://pds.example"); assert_eq!(handle.unwrap().as_str(), "alice.test");
let oversized = Resolver::with_options_and_client( vec![test_options()], 10, test_public_client(spawn_doc_server("x".repeat(MAX_DID_DOC_BYTES + 1)).await), ); let error = oversized.resolve_raw_doc(&did).await.unwrap_err(); assert!( error .to_string() .contains("DID doc response exceeds 16384 byte limit") ); }}