diff --git a/README.md b/README.md index 35c3285..e2b996c 100644 --- a/README.md +++ b/README.md @@ -4,7 +4,7 @@ -> [vs tap](#vs-tap) | [stream](#stream-behavior) | [multi-relay](#multiple-relay-support) | [crawler sources](#crawler-sources)
-> [configuration](#configuration)
-> [rest api](#rest-api) | [filter](#filter-management) | [ingestion](#ingestion-control) | [crawler](#crawler-management) | [firehose](#firehose-management) | [repos](#repository-management)
--> [xrpc api](#data-access-xrpc) | [backlinks](#bluemicrocosmlinks) | [atproto](#comatproto) | [custom](#systemsgazehydrant) +-> [xrpc api](#data-access-xrpc) | [backlinks](#bluemicrocosmlinks) | [identity](#bluemicrocosmidentity) | [atproto](#comatproto) | [custom](#systemsgazehydrant) # hydrant @@ -328,6 +328,19 @@ return the total number of stored records in a collection. returns `{ count }`. +#### systems.gaze.hydrant.describeRepo + +return account and identity information about this repo. +this is equal to `com.atproto.repo.describeRepo`, except we don't return the full DID document. +the handle is bi-directionally verified, if its invalid or the handle does not exist we return +"handle.invalid". + +| param | required | description | +| :--- | :--- | :--- | +| `identifier` | yes | DID or handle of the repository. | + +returns `{ did, handle, pds, collections }`. + ### blue.microcosm.links.* [<- back to toc](#table-of-contents) @@ -365,3 +378,11 @@ return the number of records that link to a given subject. | `source` | no | filter by source collection (same format as `getBacklinks`). | returns `{ count }`. + +### blue.microcosm.identity.* + +[<- back to toc](#table-of-contents) + +#### blue.microcosm.identity.resolveMiniDoc + +see [here](https://slingshot.microcosm.blue/#tag/slingshot-specific-queries/GET/xrpc/blue.microcosm.identity.resolveMiniDoc) for this XRPC's documentation. diff --git a/src/api/xrpc/describe_repo.rs b/src/api/xrpc/describe_repo.rs new file mode 100644 index 0000000..021f9ef --- /dev/null +++ b/src/api/xrpc/describe_repo.rs @@ -0,0 +1,85 @@ +use std::collections::HashSet; + +use futures::TryFutureExt; +use jacquard_common::types::{did::Did, nsid::Nsid, string::Handle}; +use smol_str::SmolStr; + +use crate::control::repos::MiniDocError; +use crate::db::types::DidKey; + +use super::*; + +#[derive(Serialize, Deserialize, jacquard_derive::IntoStatic)] +pub struct DescribeRepoOutput<'d> { + #[serde(borrow)] + pub did: Did<'d>, + #[serde(borrow)] + pub handle: Handle<'d>, + #[serde(serialize_with = "crate::util::did_key_serialize_str")] + #[serde(borrow)] + pub signing_key: DidKey<'d>, + pub pds: SmolStr, + #[serde(borrow)] + pub collections: HashSet>, +} + +pub struct DescribeRepoResponse; +impl jacquard_common::xrpc::XrpcResp for DescribeRepoResponse { + const NSID: &'static str = "systems.gaze.hydrant.describeRepo"; + const ENCODING: &'static str = "application/json"; + type Output<'de> = DescribeRepoOutput<'de>; + type Err<'de> = GenericXrpcError; +} + +#[derive(Serialize, Deserialize, jacquard_derive::IntoStatic)] +pub struct DescribeRepoRequestData<'i> { + #[serde(borrow)] + pub identifier: AtIdentifier<'i>, +} + +impl<'a> jacquard_common::xrpc::XrpcRequest for DescribeRepoRequestData<'a> { + type Response = DescribeRepoResponse; + const NSID: &'static str = Self::Response::NSID; + const METHOD: jacquard_common::xrpc::XrpcMethod = jacquard_common::xrpc::XrpcMethod::Query; +} + +pub struct DescribeRepo; +impl jacquard_common::xrpc::XrpcEndpoint for DescribeRepo { + const PATH: &'static str = "/xrpc/systems.gaze.hydrant.describeRepo"; + const METHOD: jacquard_common::xrpc::XrpcMethod = jacquard_common::xrpc::XrpcMethod::Query; + type Request<'de> = DescribeRepoRequestData<'de>; + type Response = DescribeRepoResponse; +} + +pub async fn handle( + State(hydrant): State, + ExtractXrpc(req): ExtractXrpc, +) -> XrpcResult>> { + let nsid = DescribeRepoResponse::NSID; + let did = hydrant + .state + .resolver + .resolve_did(&req.identifier) + .await + .map_err(|e| internal_error(nsid, format!("can't resolve identifier: {e}")))?; + + let repo = hydrant.repos.get(&did); + let doc = repo.mini_doc().map_err(|e| match e { + MiniDocError::NotSynced => bad_request(nsid, "repo not synced"), + MiniDocError::RepoNotFound => bad_request(nsid, "repo not found"), + MiniDocError::CouldNotResolveIdentity => { + upstream_error(nsid, "identity could not be resolved") + } + MiniDocError::Other(e) => internal_error(nsid, e), + }); + let collections = repo.collections().map_err(|e| internal_error(nsid, e)); + let (doc, collections) = tokio::try_join!(doc, collections)?; + + Ok(Json(DescribeRepoOutput { + did: doc.did, + handle: doc.handle, + pds: doc.pds.to_smolstr(), + signing_key: doc.signing_key, + collections: collections.into_iter().map(|(k, _)| k).collect(), + })) +} diff --git a/src/api/xrpc/mod.rs b/src/api/xrpc/mod.rs index 32cc925..2208717 100644 --- a/src/api/xrpc/mod.rs +++ b/src/api/xrpc/mod.rs @@ -1,4 +1,5 @@ use crate::api::xrpc::count_records::CountRecords; +use crate::api::xrpc::describe_repo::DescribeRepo; use crate::control::Hydrant; use axum::extract::FromRequest; use axum::response::IntoResponse; @@ -9,6 +10,7 @@ use jacquard_api::com_atproto::repo::{ list_records::{ListRecordsOutput, ListRecordsRequest, Record as RepoRecord}, }; use jacquard_common::types::ident::AtIdentifier; +use jacquard_common::xrpc::XrpcResp; use jacquard_common::xrpc::{XrpcEndpoint, XrpcMethod}; use jacquard_common::{IntoStatic, xrpc::XrpcRequest}; use jacquard_common::{ @@ -20,6 +22,7 @@ use smol_str::ToSmolStr; use std::fmt::Display; mod count_records; +mod describe_repo; mod get_record; mod list_records; @@ -28,6 +31,7 @@ pub fn router() -> Router { .route(GetRecordRequest::PATH, get(get_record::handle)) .route(ListRecordsRequest::PATH, get(list_records::handle)) .route(CountRecords::PATH, get(count_records::handle)) + .route(DescribeRepo::PATH, get(describe_repo::handle)) } #[derive(Debug)] @@ -110,3 +114,19 @@ fn bad_request( }), } } + +fn upstream_error( + nsid: &'static str, + message: impl Display, +) -> XrpcErrorResponse { + XrpcErrorResponse { + status: StatusCode::BAD_GATEWAY, + error: XrpcError::Generic(GenericXrpcError { + error: "UpstreamError".into(), + message: Some(message.to_smolstr()), + nsid, + method: "GET", + http_status: StatusCode::BAD_GATEWAY, + }), + } +} diff --git a/src/control/mod.rs b/src/control/mod.rs index 3f92b1b..c03b114 100644 --- a/src/control/mod.rs +++ b/src/control/mod.rs @@ -1,8 +1,8 @@ -mod crawler; -mod filter; -mod firehose; -mod repos; -mod stream; +pub(crate) mod crawler; +pub(crate) mod filter; +pub(crate) mod firehose; +pub(crate) mod repos; +pub(crate) mod stream; pub use crawler::{CrawlerHandle, CrawlerSourceInfo}; pub use filter::{FilterControl, FilterPatch, FilterSnapshot}; diff --git a/src/control/repos.rs b/src/control/repos.rs index 7cb71e3..7275398 100644 --- a/src/control/repos.rs +++ b/src/control/repos.rs @@ -1,3 +1,4 @@ +use std::collections::HashMap; use std::sync::Arc; use chrono::{DateTime, Utc}; @@ -5,6 +6,7 @@ use fjall::OwnedWriteBatch; use jacquard_common::cowstr::ToCowStr; use jacquard_common::types::cid::{Cid, IpldCid}; use jacquard_common::types::ident::AtIdentifier; +use jacquard_common::types::nsid::Nsid; use jacquard_common::types::string::{Did, Handle, Rkey}; use jacquard_common::types::tid::Tid; use jacquard_common::{CowStr, Data, IntoStatic}; @@ -13,7 +15,7 @@ use rand::Rng; use smol_str::ToSmolStr; use url::Url; -use crate::db::types::{DbRkey, TrimmedDid}; +use crate::db::types::{DbRkey, DidKey, TrimmedDid}; use crate::db::{self, Db, keys, ser_repo_state}; use crate::state::AppState; use crate::types::{GaugeState, RepoState, RepoStatus}; @@ -37,14 +39,17 @@ pub struct RepoInfo { #[serde(skip_serializing_if = "Option::is_none")] pub data: Option, /// the handle for the DID of this repository. + /// + /// note that this handle is not bi-directionally verified. #[serde(skip_serializing_if = "Option::is_none")] pub handle: Option>, /// the URL for the PDS in which this repository is hosted on. #[serde(skip_serializing_if = "Option::is_none")] pub pds: Option, /// ATProto signing key of this repository. + #[serde(serialize_with = "crate::util::opt_did_key_serialize_str")] #[serde(skip_serializing_if = "Option::is_none")] - pub signing_key: Option, + pub signing_key: Option>, /// when this repository was last touched (status update, commit ingested, etc.). #[serde(skip_serializing_if = "Option::is_none")] pub last_updated_at: Option>, @@ -154,11 +159,11 @@ impl ReposControl { } /// gets a handle for a repository to read from it. - pub fn get<'i>(&self, did: &Did<'i>) -> Result> { - Ok(RepoHandle { + pub fn get<'i>(&self, did: &Did<'i>) -> RepoHandle<'i> { + RepoHandle { state: self.0.clone(), did: did.clone(), - }) + } } /// same as [`ReposControl::get`] but allows you to pass in an identifier that can be @@ -171,10 +176,10 @@ impl ReposControl { }) } - /// fetch the current state of repository. + /// fetch the current state of a repository. /// returns `None` if hydrant has never seen this repository. pub async fn info(&self, did: &Did<'_>) -> Result> { - self.get(did)?.info().await + self.get(did).info().await } fn _resync( @@ -384,7 +389,7 @@ pub(crate) fn repo_state_to_info(did: Did<'static>, s: RepoState<'_>) -> RepoInf data: s.data, handle: s.handle.map(|h| h.into_static()), pds: s.pds.and_then(|p| p.parse().ok()), - signing_key: s.signing_key.map(|k| k.encode()), + signing_key: s.signing_key.map(|k| k.into_static()), last_updated_at: DateTime::from_timestamp_secs(s.last_updated_at), last_message_at: s.last_message_time.and_then(DateTime::from_timestamp_secs), } @@ -407,6 +412,31 @@ pub struct RecordList { pub cursor: Option>, } +#[derive(Debug, thiserror::Error)] +pub enum MiniDocError { + #[error("repo is not synced yet")] + NotSynced, + #[error("repo not found")] + RepoNotFound, + #[error("could not resolve identity")] + CouldNotResolveIdentity, + #[error("{0}")] + Other(miette::Error), +} + +/// a mini doc with a bi-directionally verified handle. +pub struct MiniDoc<'i> { + /// the did. + pub did: Did<'i>, + /// the handle. if verification fails or no handle is found, + /// this will be "handle.invalid". + pub handle: Handle<'i>, + /// the url of the PDS of this repo. + pub pds: Url, + /// the atproto signing key of this repo. + pub signing_key: DidKey<'i>, +} + /// handle to access data related to this repository. #[derive(Clone)] pub struct RepoHandle<'i> { @@ -415,6 +445,8 @@ pub struct RepoHandle<'i> { } impl<'i> RepoHandle<'i> { + /// fetch the current state of this repository. + /// returns `None` if hydrant has never seen this repository. pub async fn info(&self) -> Result> { let did_key = keys::repo_key(&self.did); let state = self.state.clone(); @@ -429,6 +461,87 @@ impl<'i> RepoHandle<'i> { .into_diagnostic()? } + /// returns the collections of this repository and the number of records it has in each. + pub async fn collections(&self) -> Result, u64>> { + let did = self.did.clone().into_static(); + let state = self.state.clone(); + + tokio::task::spawn_blocking(move || { + let prefix = keys::did_collection_prefix(&did); + let mut res = HashMap::new(); + for item in state.db.counts.prefix(&prefix) { + let (k, v) = item.into_inner().into_diagnostic()?; + let col = k + .strip_prefix(prefix.as_slice()) + .ok_or_else(|| miette::miette!("invalid collection count key: {k:?}")) + .and_then(|r| std::str::from_utf8(r).into_diagnostic()) + .and_then(|n| Nsid::new(n).into_diagnostic())? + .into_static(); + let count = u64::from_be_bytes( + v.as_ref() + .try_into() + .into_diagnostic() + .wrap_err("expected to be count (8 bytes)")?, + ); + res.insert(col, count); + } + Ok(res) + }) + .await + .into_diagnostic()? + } + + /// returns a bi-directionally validated mini doc. + pub async fn mini_doc(&self) -> Result, MiniDocError> { + fn invalid_handle() -> Handle<'static> { + unsafe { Handle::unchecked("handle.invalid") } + } + + let Some(info) = self.info().await.map_err(MiniDocError::Other)? else { + return Err(MiniDocError::RepoNotFound); + }; + + if info.status == RepoStatus::Backfilling { + return Err(MiniDocError::NotSynced); + } + + let pds = info + .pds + .ok_or_else(|| MiniDocError::CouldNotResolveIdentity)?; + let signing_key = info + .signing_key + .ok_or_else(|| MiniDocError::CouldNotResolveIdentity)? + .into_static(); + + let handle = if let Some(handle_unverified) = info.handle { + let id = AtIdentifier::Handle(handle_unverified); + let handle_did = self + .state + .resolver + .resolve_did(&id) + .await + .into_diagnostic() + .map_err(MiniDocError::Other)?; + + (handle_did == self.did) + .then(|| match id { + AtIdentifier::Handle(h) => h, + _ => unreachable!("can only be handle"), + }) + .unwrap_or_else(invalid_handle) + } else { + invalid_handle() + }; + + Ok(MiniDoc { + did: self.did.clone().into_static(), + handle, + pds, + signing_key, + }) + } + + /// gets a record from this repository. pub async fn get_record(&self, collection: &str, rkey: &str) -> Result> { let did = self.did.clone().into_static(); let db_key = keys::record_key(&did, collection, &DbRkey::new(rkey)); @@ -464,6 +577,7 @@ impl<'i> RepoHandle<'i> { .into_diagnostic()? } + /// lists records from this repository. pub async fn list_records( &self, collection: &str, @@ -559,6 +673,7 @@ impl<'i> RepoHandle<'i> { }) } + /// gets how many records of a collection this repository has. pub async fn count_records(&self, collection: &str) -> Result { let did = self.did.clone().into_static(); let state = self.state.clone(); diff --git a/src/db/keys.rs b/src/db/keys.rs index 3edc61d..4030f8c 100644 --- a/src/db/keys.rs +++ b/src/db/keys.rs @@ -110,14 +110,18 @@ pub fn count_keyspace_key(name: &str) -> Vec { pub const COUNT_COLLECTION_PREFIX: &[u8] = &[b'r', SEP]; -// key format: r|{DID}|{collection} (DID trimmed) -pub fn count_collection_key(did: &Did, collection: &str) -> Vec { +pub fn did_collection_prefix(did: &Did) -> Vec { let repo = TrimmedDid::from(did); - let mut key = - Vec::with_capacity(COUNT_COLLECTION_PREFIX.len() + repo.len() + 1 + collection.len()); + let mut key = Vec::with_capacity(COUNT_COLLECTION_PREFIX.len() + repo.len() + 1); key.extend_from_slice(COUNT_COLLECTION_PREFIX); repo.write_to_vec(&mut key); key.push(SEP); + key +} + +// key format: r|{DID}|{collection} (DID trimmed) +pub fn count_collection_key(did: &Did, collection: &str) -> Vec { + let mut key = did_collection_prefix(did); key.extend_from_slice(collection.as_bytes()); key } diff --git a/src/util.rs b/src/util.rs index 7465c78..7596a86 100644 --- a/src/util.rs +++ b/src/util.rs @@ -8,7 +8,7 @@ use tokio::sync::watch; use tracing::info; use url::Url; -use crate::types::RepoStatus; +use crate::{db::types::DidKey, types::RepoStatus}; /// outcome of [`RetryWithBackoff::retry`] when the operation does not succeed. pub enum RetryOutcome { @@ -153,6 +153,20 @@ pub fn opt_cid_serialize_str(v: &Option, s: S) -> Resul } } +pub fn did_key_serialize_str(v: &DidKey<'_>, s: S) -> Result { + s.serialize_str(&v.encode()) +} + +pub fn opt_did_key_serialize_str( + v: &Option>, + s: S, +) -> Result { + match v { + Some(k) => s.serialize_some(k.encode().as_str()), + None => s.serialize_none(), + } +} + pub fn repo_status_serialize_str(v: &RepoStatus, s: S) -> Result { s.serialize_str(&v.to_string()) }