diff --git a/bobbin/crates/xrpc/Cargo.toml b/bobbin/crates/xrpc/Cargo.toml index 7f13094a..6ef8b957 100644 --- a/bobbin/crates/xrpc/Cargo.toml +++ b/bobbin/crates/xrpc/Cargo.toml @@ -23,6 +23,7 @@ http = { workspace = true } serde = { workspace = true } serde_json = { workspace = true } thiserror = { workspace = true } +tower = { workspace = true } tower-http = { workspace = true, features = ["trace"] } tracing = { workspace = true } url = { workspace = true } @@ -31,6 +32,5 @@ url = { workspace = true } bobbin-runtime = { workspace = true } http = { workspace = true } tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } -tower = { workspace = true } url = { workspace = true } wiremock = { workspace = true } diff --git a/bobbin/crates/xrpc/src/enrich.rs b/bobbin/crates/xrpc/src/enrich.rs new file mode 100644 index 00000000..82297ce9 --- /dev/null +++ b/bobbin/crates/xrpc/src/enrich.rs @@ -0,0 +1,372 @@ +use std::collections::HashSet; + +use axum::{ + Json, + body::{Body, to_bytes}, + extract::State, + http::{Request, StatusCode}, + response::Response, +}; +use bobbin_types::ids::{EdgeKey, SubjectRef, nsid_static}; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use jacquard_common::types::ident::AtIdentifier; +use jacquard_common::types::string::AtUri; +use serde::Deserialize; +use serde_json::{Map, Value, json}; +use tower::ServiceExt; + +use crate::recordpath::{parse_record_path, walk_path}; +use crate::{AppState, SubjectShape, XrpcError, mirror_kind, subject_shape}; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)] +#[serde(rename_all = "camelCase")] +pub enum Aggregation { + Count, + DistinctAuthors, + Viewer, +} + +impl Aggregation { + fn key(self) -> &'static str { + match self { + Aggregation::Count => "count", + Aggregation::DistinctAuthors => "distinctAuthors", + Aggregation::Viewer => "viewer", + } + } +} + +#[derive(Debug, Deserialize)] +pub struct LinkDescriptor { + /// constellation link source + source: String, + #[serde(rename = "type")] + aggregation: Aggregation, +} + +impl LinkDescriptor { + fn parts(&self) -> Result<(&str, &str, Aggregation), XrpcError> { + match self.source.split_once(':') { + Some((collection, path)) if !collection.is_empty() && !path.is_empty() => { + Ok((collection, path, self.aggregation)) + } + _ => Err(XrpcError::InvalidParams(format!( + "enrich {:?}: expected \"collection:path\"", + self.source + ))), + } + } +} + +#[derive(Debug, Deserialize)] +pub struct EnrichInput { + xrpc: String, + #[serde(default)] + params: Option>, + enrich: Vec, + #[serde(default)] + sources: Option>, + #[serde(default)] + viewer: Option, +} + +pub async fn enrich( + State(state): State, + Json(input): Json, +) -> Result, XrpcError> { + for descriptor in &input.enrich { + validate_descriptor(descriptor)?; + if descriptor.aggregation == Aggregation::Viewer && input.viewer.is_none() { + return Err(descriptor_error( + descriptor, + "viewer aggregation requires a viewer param", + )); + } + } + + let inner = run_inner(&state, &input.xrpc, input.params.unwrap_or_default()).await?; + + let mut refs: Vec = Vec::new(); + let mut seen: HashSet = HashSet::new(); + match &input.sources { + Some(sources) => { + for (i, path) in sources.iter().enumerate() { + let segs = parse_record_path(path) + .map_err(|e| XrpcError::InvalidParams(format!("sources[{i}] {path:?}: {e}")))?; + for node in walk_path(&segs, [&inner]) { + collect_ref(node, &mut refs, &mut seen); + } + } + } + None => discover_refs(&inner, &mut refs, &mut seen), + } + + let mut stats = Map::new(); + for reference in &refs { + let mut per_ref = Map::new(); + for descriptor in &input.enrich { + let (.., aggregation) = descriptor.parts()?; + let Some(kind) = resolve_descriptor(descriptor, reference)? else { + continue; + }; + let key = EdgeKey::new(kind, reference.clone()); + let result = match aggregation { + Aggregation::Count => Value::from(state.edges.count(&key)), + Aggregation::DistinctAuthors => { + Value::from(state.edges.count_distinct_authors(&key)) + } + Aggregation::Viewer => { + // viewer presence is validated up front + let viewer = input.viewer.as_ref().expect("viewer param present"); + match state.edges.viewer_source(&key, viewer.as_str()) { + Some(uri) => Value::String(uri.to_string()), + None => Value::Null, + } + } + }; + let entry = per_ref + .entry(descriptor.source.clone()) + .or_insert_with(|| Value::Object(Map::new())); + if let Value::Object(counts) = entry { + counts.insert(aggregation.key().to_owned(), result); + } + } + if !per_ref.is_empty() { + stats.insert(subject_string(reference), Value::Object(per_ref)); + } + } + + Ok(Json(json!({ "output": inner, "stats": stats }))) +} + +fn descriptor_error(descriptor: &LinkDescriptor, msg: &str) -> XrpcError { + XrpcError::InvalidParams(format!("enrich {}: {msg}", descriptor.source)) +} + +fn validate_descriptor(descriptor: &LinkDescriptor) -> Result<(), XrpcError> { + let (collection, path, _) = descriptor.parts()?; + match path { + "subject" => subject_shape(collection) + .map(|_| ()) + .ok_or_else(|| descriptor_error(descriptor, "unknown collection")), + ".repo" => mirror_kind(collection) + .map(|_| ()) + .ok_or_else(|| descriptor_error(descriptor, "collection has no author index")), + _ if path.starts_with('.') => Err(descriptor_error( + descriptor, + "envelope field not index-backed; try .repo", + )), + _ => Err(descriptor_error( + descriptor, + "path not index-backed; only `subject` and `.repo` are supported", + )), + } +} + +/// the edge kind to look up for one reference, or None if the descriptor's shape +/// doesn't apply to this ref kind, eg. a did-only descriptor asked about an at-uri +fn resolve_descriptor( + descriptor: &LinkDescriptor, + reference: &SubjectRef, +) -> Result>, XrpcError> { + let (collection, path, _) = descriptor.parts()?; + match path { + "subject" => { + let Some((nsid, shape)) = subject_shape(collection) else { + return Err(descriptor_error(descriptor, "unknown collection")); + }; + Ok(shape_accepts(shape, reference).then(|| nsid_static(nsid))) + } + ".repo" => { + let Some(kind) = mirror_kind(collection) else { + return Err(descriptor_error( + descriptor, + "collection has no author index", + )); + }; + Ok(matches!(reference, SubjectRef::Did(_)).then(|| nsid_static(kind))) + } + _ => unreachable!("validated up front"), + } +} + +/// mismatch means skip, not reject +fn shape_accepts(shape: SubjectShape, reference: &SubjectRef) -> bool { + match (shape, reference) { + (SubjectShape::BareDid, SubjectRef::Did(_)) => true, + (SubjectShape::Collection(expected), SubjectRef::Uri(uri)) => { + uri.collection().is_some_and(|c| c.as_ref() == expected) + } + (SubjectShape::OneOfCollections(allowed), SubjectRef::Uri(uri)) + | (SubjectShape::BareDidOrOneOfCollections(allowed), SubjectRef::Uri(uri)) => uri + .collection() + .is_some_and(|c| allowed.contains(&c.as_ref())), + (SubjectShape::BareDidOrOneOfCollections(_), SubjectRef::Did(_)) => true, + (SubjectShape::AnyAtUri, SubjectRef::Uri(_)) => true, + _ => false, + } +} + +fn subject_string(reference: &SubjectRef) -> String { + match reference { + SubjectRef::Did(d) => d.as_ref().to_owned(), + SubjectRef::Uri(u) => u.as_ref().to_owned(), + } +} + +/// a value counts as a reference if it is a did string or an at-uri string +/// with collection and rkey, or a {uri, cid} strong ref +fn collect_ref(value: &Value, refs: &mut Vec, seen: &mut HashSet) { + let candidate = match value { + Value::String(s) => Some(s.as_str()), + Value::Object(o) => o + .get("uri") + .and_then(Value::as_str) + .filter(|_| o.contains_key("cid")), + _ => None, + }; + let Some(candidate) = candidate else { return }; + let reference = if candidate.starts_with("did:") { + Did::::new_owned(candidate) + .ok() + .map(SubjectRef::Did) + } else if candidate.starts_with("at://") { + AtUri::::new_owned(candidate) + .ok() + .filter(|u| { + matches!(u.authority(), AtIdentifier::Did(_)) + && u.collection().is_some() + && u.rkey().is_some() + }) + .map(SubjectRef::Uri) + } else { + None + }; + if let Some(reference) = reference + && seen.insert(reference.clone()) + { + refs.push(reference); + } +} + +fn discover_refs(value: &Value, refs: &mut Vec, seen: &mut HashSet) { + collect_ref(value, refs, seen); + match value { + Value::Array(items) => { + for item in items { + discover_refs(item, refs, seen); + } + } + Value::Object(map) => { + for v in map.values() { + discover_refs(v, refs, seen); + } + } + _ => {} + } +} + +async fn run_inner( + state: &AppState, + nsid: &str, + params: Map, +) -> Result { + let qs = encode_params(¶ms); + let uri = format!("/xrpc/{nsid}?{qs}"); + let request = Request::builder() + .method("GET") + .uri(&uri) + .body(Body::empty()) + .map_err(|e| XrpcError::Internal(format!("inner request: {e}")))?; + let response = state + .self_router() + .oneshot(request) + .await + .map_err(|e| XrpcError::Internal(format!("inner dispatch: {e}")))?; + if response.status() == StatusCode::NOT_FOUND { + let bytes = to_bytes(response.into_body(), usize::MAX) + .await + .map_err(|e| XrpcError::Internal(format!("inner response: {e}")))?; + return Err(bytes + .is_empty() + .then(|| XrpcError::InvalidParams(format!("unknown or unenrichable query: {nsid}"))) + .unwrap_or(XrpcError::NotFound)); + } + finish(response).await +} + +fn encode_params(params: &Map) -> String { + let mut out = url::form_urlencoded::Serializer::new(String::new()); + for (key, value) in params { + match value { + Value::Array(items) => { + for item in items { + if let Some(scalar) = scalar_str(item) { + out.append_pair(key, &scalar); + } + } + } + _ => { + if let Some(scalar) = scalar_str(value) { + out.append_pair(key, &scalar); + } + } + } + } + out.finish() +} + +fn scalar_str(value: &Value) -> Option { + match value { + Value::String(s) => Some(s.clone()), + Value::Number(n) => Some(n.to_string()), + Value::Bool(b) => Some(b.to_string()), + _ => None, + } +} + +async fn finish(resp: Response) -> Result { + let status = resp.status(); + let bytes = to_bytes(resp.into_body(), usize::MAX) + .await + .map_err(|e| XrpcError::Internal(format!("inner response: {e}")))?; + if !status.is_success() { + let msg = String::from_utf8_lossy(&bytes).into_owned(); + return Err(match status { + StatusCode::BAD_REQUEST => XrpcError::InvalidParams(msg), + StatusCode::NOT_FOUND => XrpcError::NotFound, + StatusCode::SERVICE_UNAVAILABLE => XrpcError::Overloaded, + _ => XrpcError::UpstreamUnavailable(format!("inner query ({status}): {msg}")), + }); + } + serde_json::from_slice(&bytes) + .map_err(|e| XrpcError::Internal(format!("inner response decode: {e}"))) +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + + #[test] + fn collects_dids_uris_and_strong_refs() { + let mut refs = Vec::new(); + let mut seen = HashSet::new(); + collect_ref(&json!("did:plc:abc"), &mut refs, &mut seen); + collect_ref( + &json!("at://did:plc:abc/sh.tangled.repo/x"), + &mut refs, + &mut seen, + ); + collect_ref( + &json!({"uri": "at://did:plc:abc/sh.tangled.repo/y", "cid": "bafy"}), + &mut refs, + &mut seen, + ); + collect_ref(&json!("at://did:plc:abc"), &mut refs, &mut seen); + collect_ref(&json!("oppi.li"), &mut refs, &mut seen); + collect_ref(&json!("did:plc:abc"), &mut refs, &mut seen); + assert_eq!(refs.len(), 3); + } +} diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index 28612da2..99d09ead 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -95,7 +95,9 @@ use tower_http::trace::{DefaultMakeSpan, OnFailure, OnResponse, TraceLayer}; use tracing::{Level, Span}; mod backpressure; +mod enrich; mod filter; +mod recordpath; pub use backpressure::{ HeavyLimiter, HeavyPermit, MaxInFlight, PerRequestAnonBytes, PressureVerdict, ReservedFloor, @@ -117,6 +119,7 @@ pub struct AppState { pub search: Arc, pub resolver: Arc, pub limiter: Option>, + enrich_router: Arc>, } impl AppState { @@ -143,6 +146,7 @@ impl AppState { search, resolver, limiter: None, + enrich_router: Arc::new(std::sync::OnceLock::new()), } } @@ -151,6 +155,13 @@ impl AppState { self } + /// for internal xrpc dispatch [`enrich`] + pub fn self_router(&self) -> Router { + self.enrich_router + .get_or_init(|| router(self.clone())) + .clone() + } + fn heavy_permit(&self) -> Result, XrpcError> { self.limiter.as_ref().map(|l| l.try_enter()).transpose() } @@ -385,6 +396,10 @@ pub fn router(state: AppState) -> Router { .route("/xrpc/sh.tangled.string.listStrings", get(list_strings)) .route("/xrpc/sh.tangled.string.countStrings", get(count_strings)) .route("/xrpc/sh.tangled.search.query", get(search_query)) + .route( + "/xrpc/sh.tangled.query.enrichResponse", + axum::routing::post(enrich::enrich), + ) .route("/xrpc/sh.tangled.bobbin.getCoverage", get(get_coverage)) .route( "/xrpc/com.bad-example.identity.resolveMiniDoc", @@ -916,183 +931,62 @@ pub trait MirrorOf { const SHAPE: SubjectShape; } -pub struct StarBy; -pub struct ReactionBy; -pub struct FollowBy; -pub struct VouchBy; -pub struct RefUpdateBy; -pub struct KnotMemberBy; -pub struct LabelOpBy; -pub struct PipelineBy; -pub struct PipelineStatusBy; -pub struct ArtifactBy; -pub struct CollaboratorBy; -pub struct FeedCommentBy; -pub struct IssueBy; -pub struct IssueStateBy; -pub struct PullBy; -pub struct PullStatusBy; -pub struct SpindleMemberBy; - -impl MirrorOf for StarBy { - type Record = StarRecord; - const EDGE_KIND: &'static str = "sh.tangled.feed.star.by"; - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl MirrorOf for ReactionBy { - type Record = ReactionRecord; - const EDGE_KIND: &'static str = "sh.tangled.feed.reaction.by"; - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl MirrorOf for FollowBy { - type Record = FollowRecord; - const EDGE_KIND: &'static str = "sh.tangled.graph.follow.by"; - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl MirrorOf for VouchBy { - type Record = VouchRecord; - const EDGE_KIND: &'static str = "sh.tangled.graph.vouch.by"; - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl MirrorOf for RefUpdateBy { - type Record = RefUpdateRecord; - const EDGE_KIND: &'static str = "sh.tangled.git.refUpdate.by"; - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl MirrorOf for KnotMemberBy { - type Record = KnotMemberRecord; - const EDGE_KIND: &'static str = "sh.tangled.knot.member.by"; - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl MirrorOf for LabelOpBy { - type Record = LabelOpRecord; - const EDGE_KIND: &'static str = "sh.tangled.label.op.by"; - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl MirrorOf for PipelineBy { - type Record = PipelineRecord; - const EDGE_KIND: &'static str = "sh.tangled.pipeline.by"; - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl MirrorOf for PipelineStatusBy { - type Record = PipelineStatusRecord; - const EDGE_KIND: &'static str = "sh.tangled.pipeline.status.by"; - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl MirrorOf for ArtifactBy { - type Record = ArtifactRecord; - const EDGE_KIND: &'static str = "sh.tangled.repo.artifact.by"; - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl MirrorOf for CollaboratorBy { - type Record = CollaboratorRecord; - const EDGE_KIND: &'static str = "sh.tangled.repo.collaborator.by"; - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl MirrorOf for FeedCommentBy { - type Record = FeedCommentRecord; - const EDGE_KIND: &'static str = "sh.tangled.feed.comment.by"; - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl MirrorOf for IssueBy { - type Record = IssueRecord; - const EDGE_KIND: &'static str = "sh.tangled.repo.issue.by"; - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl MirrorOf for IssueStateBy { - type Record = IssueStateRecord; - const EDGE_KIND: &'static str = "sh.tangled.repo.issue.state.by"; - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl MirrorOf for PullBy { - type Record = PullRecord; - const EDGE_KIND: &'static str = "sh.tangled.repo.pull.by"; - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl MirrorOf for PullStatusBy { - type Record = PullStatusRecord; - const EDGE_KIND: &'static str = "sh.tangled.repo.pull.status.by"; - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl MirrorOf for SpindleMemberBy { - type Record = SpindleMemberRecord; - const EDGE_KIND: &'static str = "sh.tangled.spindle.member.by"; - const SHAPE: SubjectShape = SubjectShape::BareDid; -} - -impl HasSubject for StarRecord { - const SHAPE: SubjectShape = SubjectShape::BareDidOrOneOfCollections(&["sh.tangled.string"]); -} -impl HasSubject for FollowRecord { - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl HasSubject for IssueRecord { - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl HasSubject for PullRecord { - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl HasSubject for FeedCommentRecord { - const SHAPE: SubjectShape = SubjectShape::OneOfCollections(&[ - "sh.tangled.repo.issue", - "sh.tangled.repo.pull", - "sh.tangled.string", - ]); -} -impl HasSubject for LabelDefinitionRecord { - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl HasSubject for LabelOpRecord { - const SHAPE: SubjectShape = - SubjectShape::OneOfCollections(&["sh.tangled.repo.issue", "sh.tangled.repo.pull"]); -} -impl HasSubject for PipelineRecord { - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl HasSubject for PipelineStatusRecord { - const SHAPE: SubjectShape = SubjectShape::Collection("sh.tangled.pipeline"); -} -impl HasSubject for ArtifactRecord { - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl HasSubject for KnotMemberRecord { - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl HasSubject for SpindleMemberRecord { - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl HasSubject for TangledStringRecord { - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl HasSubject for ReactionRecord { - const SHAPE: SubjectShape = SubjectShape::AnyAtUri; -} -impl HasSubject for RefUpdateRecord { - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl HasSubject for CollaboratorRecord { - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl HasSubject for IssueStateRecord { - const SHAPE: SubjectShape = SubjectShape::Collection("sh.tangled.repo.issue"); -} -impl HasSubject for PullStatusRecord { - const SHAPE: SubjectShape = SubjectShape::Collection("sh.tangled.repo.pull"); -} -impl HasSubject for KnotRecord { - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl HasSubject for SpindleRecord { - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl HasSubject for PublicKeyRecord { - const SHAPE: SubjectShape = SubjectShape::BareDid; -} -impl HasSubject for RepoRecord { - const SHAPE: SubjectShape = SubjectShape::BareDid; +macro_rules! edge_kinds { + ($($collection:literal => $record:ty, $shape:expr $(, mirror $by:ident)? ;)*) => { + $( + impl HasSubject for $record { + const SHAPE: SubjectShape = $shape; + } + )* + $($( + pub struct $by; + impl MirrorOf for $by { + type Record = $record; + const EDGE_KIND: &'static str = concat!($collection, ".by"); + const SHAPE: SubjectShape = SubjectShape::BareDid; + } + )?)* + + pub(crate) fn subject_shape(collection: &str) -> Option<(&'static str, SubjectShape)> { + Some(match collection { + $($collection => ($collection, <$record as HasSubject>::SHAPE),)* + _ => return None, + }) + } + + pub(crate) fn mirror_kind(collection: &str) -> Option<&'static str> { + Some(match collection { + $($($collection => <$by as MirrorOf>::EDGE_KIND,)?)* + _ => return None, + }) + } + }; } -impl HasSubject for VouchRecord { - const SHAPE: SubjectShape = SubjectShape::BareDid; + +edge_kinds! { + "sh.tangled.feed.star" => StarRecord, SubjectShape::BareDidOrOneOfCollections(&["sh.tangled.string"]), mirror StarBy; + "sh.tangled.feed.comment" => FeedCommentRecord, SubjectShape::OneOfCollections(&["sh.tangled.repo.issue", "sh.tangled.repo.pull", "sh.tangled.string"]), mirror FeedCommentBy; + "sh.tangled.feed.reaction" => ReactionRecord, SubjectShape::AnyAtUri, mirror ReactionBy; + "sh.tangled.graph.follow" => FollowRecord, SubjectShape::BareDid, mirror FollowBy; + "sh.tangled.graph.vouch" => VouchRecord, SubjectShape::BareDid, mirror VouchBy; + "sh.tangled.git.refUpdate" => RefUpdateRecord, SubjectShape::BareDid, mirror RefUpdateBy; + "sh.tangled.knot" => KnotRecord, SubjectShape::BareDid; + "sh.tangled.knot.member" => KnotMemberRecord, SubjectShape::BareDid, mirror KnotMemberBy; + "sh.tangled.label.definition" => LabelDefinitionRecord, SubjectShape::BareDid; + "sh.tangled.label.op" => LabelOpRecord, SubjectShape::OneOfCollections(&["sh.tangled.repo.issue", "sh.tangled.repo.pull"]), mirror LabelOpBy; + "sh.tangled.pipeline" => PipelineRecord, SubjectShape::BareDid, mirror PipelineBy; + "sh.tangled.pipeline.status" => PipelineStatusRecord, SubjectShape::Collection("sh.tangled.pipeline"), mirror PipelineStatusBy; + "sh.tangled.publicKey" => PublicKeyRecord, SubjectShape::BareDid; + "sh.tangled.repo" => RepoRecord, SubjectShape::BareDid; + "sh.tangled.repo.artifact" => ArtifactRecord, SubjectShape::BareDid, mirror ArtifactBy; + "sh.tangled.repo.collaborator" => CollaboratorRecord, SubjectShape::BareDid, mirror CollaboratorBy; + "sh.tangled.repo.issue" => IssueRecord, SubjectShape::BareDid, mirror IssueBy; + "sh.tangled.repo.issue.state" => IssueStateRecord, SubjectShape::Collection("sh.tangled.repo.issue"), mirror IssueStateBy; + "sh.tangled.repo.pull" => PullRecord, SubjectShape::BareDid, mirror PullBy; + "sh.tangled.repo.pull.status" => PullStatusRecord, SubjectShape::Collection("sh.tangled.repo.pull"), mirror PullStatusBy; + "sh.tangled.spindle" => SpindleRecord, SubjectShape::BareDid; + "sh.tangled.spindle.member" => SpindleMemberRecord, SubjectShape::BareDid, mirror SpindleMemberBy; + "sh.tangled.string" => TangledStringRecord, SubjectShape::BareDid; } fn parse_subject(raw: &SubjectQuery, shape: SubjectShape) -> Result { diff --git a/bobbin/crates/xrpc/src/recordpath.rs b/bobbin/crates/xrpc/src/recordpath.rs new file mode 100644 index 00000000..3b6305ec --- /dev/null +++ b/bobbin/crates/xrpc/src/recordpath.rs @@ -0,0 +1,301 @@ +//! subset of microcosm.blue/RecordPath + +use serde_json::Value; + +#[derive(Debug, Clone, PartialEq)] +pub(crate) enum Modifier { + /// `[]` + Elements, + /// `[nsid]` + UnionElements(String), + /// `{nsid}` + Union(String), +} + +#[derive(Debug, Clone, PartialEq)] +pub(crate) struct Segment { + field: String, + modifier: Option, +} + +pub(crate) fn parse_record_path(path: &str) -> Result, String> { + let chars: Vec = path.chars().collect(); + let mut segments = Vec::new(); + let mut field = String::new(); + let mut i = 0; + let mut saw_any = false; + + let unescape = |chars: &[char], i: &mut usize| -> Result { + *i += 1; + match chars.get(*i) { + Some(&c @ ('.' | '[' | ']' | '{' | '}' | '!')) => Ok(c), + Some(c) => Err(format!("invalid escape !{c}")), + None => Err("trailing ! escape".to_owned()), + } + }; + + while i < chars.len() { + let c = chars[i]; + match c { + '.' => { + if !saw_any { + return Err("empty path segment".to_owned()); + } + segments.push(Segment { + field: std::mem::take(&mut field), + modifier: None, + }); + saw_any = false; + i += 1; + } + '!' => { + field.push(unescape(&chars, &mut i)?); + saw_any = true; + i += 1; + } + '[' | '{' => { + let (close, is_union_brace) = if c == '[' { (']', false) } else { ('}', true) }; + let mut inner = String::new(); + i += 1; + loop { + match chars.get(i) { + None => return Err(format!("unclosed {c}")), + Some(&close_c) if close_c == close => break, + Some(&'!') => inner.push(unescape(&chars, &mut i)?), + Some(&ch) => inner.push(ch), + } + i += 1; + } + i += 1; // consume closer + let modifier = match (is_union_brace, inner.is_empty()) { + (false, true) => Modifier::Elements, + (false, false) => Modifier::UnionElements(inner), + (true, true) => return Err("empty {} union ref".to_owned()), + (true, false) => Modifier::Union(inner), + }; + if field.is_empty() && !saw_any { + return Err("modifier without a field".to_owned()); + } + segments.push(Segment { + field: std::mem::take(&mut field), + modifier: Some(modifier), + }); + saw_any = false; + // a modified segment must end the path or be followed by '.' + match chars.get(i) { + None => break, + Some('.') => { + i += 1; + } + Some(other) => { + return Err(format!("expected '.' after {c}{close}, got {other:?}")); + } + } + } + _ => { + field.push(c); + saw_any = true; + i += 1; + } + } + } + if saw_any { + segments.push(Segment { + field, + modifier: None, + }); + } + if segments.is_empty() { + return Err("empty path".to_owned()); + } + Ok(segments) +} + +pub(crate) fn type_tag(value: &Value) -> Option<&str> { + value.get("$type").and_then(Value::as_str) +} + +pub(crate) fn walk_path<'a>( + segments: &[Segment], + roots: impl IntoIterator, +) -> Vec<&'a Value> { + let mut nodes: Vec<&Value> = roots.into_iter().collect(); + for segment in segments { + let mut next: Vec<&Value> = nodes + .into_iter() + .filter_map(|node| node.get(&segment.field)) + .collect(); + if let Some(modifier) = &segment.modifier { + next = match modifier { + Modifier::Elements => next + .into_iter() + .flat_map(|node| node.as_array().into_iter().flatten()) + .collect(), + Modifier::UnionElements(nsid) => next + .into_iter() + .flat_map(|node| node.as_array().into_iter().flatten()) + .filter(|item| type_tag(item) == Some(nsid.as_str())) + .collect(), + Modifier::Union(nsid) => next + .into_iter() + .filter(|node| type_tag(node) == Some(nsid.as_str())) + .collect(), + }; + } + nodes = next; + } + nodes +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + + fn parse_ok(path: &str) -> Vec { + parse_record_path(path).unwrap_or_else(|e| panic!("{path}: {e}")) + } + + #[test] + fn parses_field_paths() { + assert_eq!( + parse_ok("subject.uri"), + vec![ + Segment { + field: "subject".into(), + modifier: None + }, + Segment { + field: "uri".into(), + modifier: None + }, + ] + ); + } + + #[test] + fn parses_array_descents() { + assert_eq!( + parse_ok("repos[]"), + vec![Segment { + field: "repos".into(), + modifier: Some(Modifier::Elements) + }] + ); + assert_eq!( + parse_ok("facets[].features[app.bsky.richtext.facet#mention].did"), + vec![ + Segment { + field: "facets".into(), + modifier: Some(Modifier::Elements) + }, + Segment { + field: "features".into(), + modifier: Some(Modifier::UnionElements( + "app.bsky.richtext.facet#mention".into() + )) + }, + Segment { + field: "did".into(), + modifier: None + }, + ] + ); + assert_eq!( + parse_ok("embed{app.bsky.embed.record}.record.uri"), + vec![ + Segment { + field: "embed".into(), + modifier: Some(Modifier::Union("app.bsky.embed.record".into())) + }, + Segment { + field: "record".into(), + modifier: None + }, + Segment { + field: "uri".into(), + modifier: None + }, + ] + ); + } + + #[test] + fn parses_escaped_field_names() { + assert_eq!( + parse_ok("meta.dot!.name"), + vec![Segment { + field: "meta".into(), + modifier: None + }] + .into_iter() + .chain([Segment { + field: "dot.name".into(), + modifier: None + }]) + .collect::>() + ); + assert_eq!( + parse_ok("meta.a!!b"), + vec![ + Segment { + field: "meta".into(), + modifier: None + }, + Segment { + field: "a!b".into(), + modifier: None + }, + ] + ); + assert_eq!( + parse_ok("meta.$unknown"), + vec![ + Segment { + field: "meta".into(), + modifier: None + }, + Segment { + field: "$unknown".into(), + modifier: None + }, + ] + ); + } + + #[test] + fn rejects_malformed_paths() { + for bad in [ + "", ".", "a..b", "a!", "a!x", "a[", "a[]b", "a{}", "[did]", "a[nsid]b", + ] { + assert!(parse_record_path(bad).is_err(), "{bad} should fail"); + } + } + + #[test] + fn walks_vector_matches() { + let doc = json!({ + "repos": [ + {"uri": "at://did:plc:a/sh.tangled.repo/x"}, + {"uri": "at://did:plc:b/sh.tangled.repo/y"}, + ], + "owner": "did:plc:z", + }); + let segs = parse_ok("repos[].uri"); + let found = walk_path(&segs, [&doc]); + assert_eq!(found.len(), 2); + } + + #[test] + fn walks_union_filters() { + let doc = json!({ + "items": [ + {"$type": "sh.tangled.repo", "uri": "at://did:plc:a/sh.tangled.repo/x"}, + {"$type": "sh.tangled.actor.profile", "did": "did:plc:b"}, + ] + }); + let segs = parse_ok("items[sh.tangled.repo].uri"); + let found = walk_path(&segs, [&doc]); + assert_eq!(found, vec![&json!("at://did:plc:a/sh.tangled.repo/x")]); + } +} diff --git a/bobbin/crates/xrpc/tests/enrich.rs b/bobbin/crates/xrpc/tests/enrich.rs new file mode 100644 index 00000000..3f4daac2 --- /dev/null +++ b/bobbin/crates/xrpc/tests/enrich.rs @@ -0,0 +1,441 @@ +use std::sync::Arc; + +use axum::body::{Body, to_bytes}; +use bobbin_edge_index::{CoverageWatch, EdgeStore, IssueStateKind, PullStatusKind, StateIndex}; +use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig}; +use bobbin_record_lru::{CacheCapacity, LruRecordStore}; +use bobbin_resolver::RepoIdResolver; +use bobbin_runtime::{RuntimeHasher, SystemClock}; +use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex, SearchReader}; +use bobbin_slingshot_client::SlingshotClient; +use bobbin_types::edges::Edge; +use bobbin_types::ids::SubjectRef; +use bobbin_xrpc::{AppState, router}; +use http::{Request, StatusCode}; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use jacquard_common::types::nsid::Nsid; +use jacquard_common::types::string::AtUri; +use serde_json::{Value, json}; +use tower::ServiceExt; +use url::Url; +use wiremock::matchers::{method, path, query_param}; +use wiremock::{Mock, MockServer, ResponseTemplate}; + +const CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; + +fn at(s: &str) -> AtUri { + AtUri::new_owned(s).unwrap() +} + +fn did(s: &str) -> Did { + Did::new_owned(s).unwrap() +} + +fn nsid(s: &'static str) -> Nsid { + Nsid::new_static(s).unwrap() +} + +static EDGE_COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1); + +fn next_sort_micros() -> u64 { + EDGE_COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed) +} + +struct Harness { + server: MockServer, + edges: Arc, + state: AppState, +} + +impl Harness { + async fn new() -> Self { + let server = MockServer::start().await; + let edges = Arc::new(EdgeStore::new(RuntimeHasher::default())); + let coverage = Arc::new(CoverageWatch::new()); + let state = AppState::new( + Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024))), + SlingshotClient::with_default_http(Url::parse(&server.uri()).unwrap()).unwrap(), + edges.clone(), + Arc::new(StateIndex::::new(RuntimeHasher::default())), + Arc::new(StateIndex::::new(RuntimeHasher::default())), + coverage.clone(), + Arc::new( + KnotProxy::new( + KnotProxyConfig::default(), + KnotHttpConfig::default(), + Arc::new(SystemClock::new()), + RuntimeHasher::default(), + ) + .unwrap(), + ), + Arc::new( + SearchIndex::new(DEFAULT_WRITER_HEAP_BYTES, Arc::new(SystemClock::new())).unwrap(), + ) as Arc, + Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), + ); + Self { + server, + edges, + state, + } + } + + fn add_edge(&self, kind: &'static str, subject: SubjectRef, source: &AtUri) { + self.edges.add(Edge { + kind: nsid(kind), + subject, + source: source.clone(), + sort_micros: next_sort_micros(), + }); + } + + async fn mount(&self, did: &Did, collection: &str, rkey: &str, value: Value) { + let uri = format!("at://{}/{}/{}", did.as_ref(), collection, rkey); + let body = json!({ "uri": uri, "cid": CID, "value": value }); + Mock::given(method("GET")) + .and(path("/xrpc/com.atproto.repo.getRecord")) + .and(query_param("repo", did.as_ref())) + .and(query_param("collection", collection)) + .and(query_param("rkey", rkey)) + .respond_with(ResponseTemplate::new(200).set_body_json(body)) + .mount(&self.server) + .await; + } +} + +fn enrich_request(body: Value) -> Request { + Request::builder() + .method("POST") + .uri("/xrpc/sh.tangled.query.enrichResponse") + .header("content-type", "application/json") + .body(Body::from(serde_json::to_vec(&body).unwrap())) + .unwrap() +} + +async fn json_response(resp: axum::response::Response) -> (StatusCode, Value) { + let status = resp.status(); + let bytes = to_bytes(resp.into_body(), 1 << 20).await.unwrap(); + let parsed: Value = serde_json::from_slice(&bytes).expect("JSON body"); + (status, parsed) +} + +fn repo_body(name: &str, repo_did: &Did) -> Value { + json!({ + "$type": "sh.tangled.repo", + "name": name, + "knot": "oyster.cafe", + "repoDid": repo_did.as_ref(), + "createdAt": "2026-05-01T00:00:00Z" + }) +} + +fn follow_body(subject: &Did) -> Value { + json!({ + "$type": "sh.tangled.graph.follow", + "createdAt": "2026-05-01T00:00:00Z", + "subject": subject.as_ref() + }) +} + +/// one repo owned by `owner`, with `stars`/`issues` counts against its repo did +async fn repo_fixture(h: &Harness, owner: &Did, repo_did: &Did) { + let repo_uri = at(&format!("at://{}/sh.tangled.repo/reef", owner.as_ref())); + h.add_edge("sh.tangled.repo", SubjectRef::Did(owner.clone()), &repo_uri); + h.mount( + owner, + "sh.tangled.repo", + "reef", + repo_body("reef", repo_did), + ) + .await; + for (i, stargazer) in ["did:plc:a", "did:plc:b", "did:plc:a"].iter().enumerate() { + h.add_edge( + "sh.tangled.feed.star", + SubjectRef::Did(repo_did.clone()), + &at(&format!("at://{stargazer}/sh.tangled.feed.star/s{i}")), + ); + } + h.add_edge( + "sh.tangled.repo.issue", + SubjectRef::Did(repo_did.clone()), + &at("at://did:plc:a/sh.tangled.repo.issue/i0"), + ); +} + +#[tokio::test] +async fn zero_config_counts_stars_and_issues_for_repo_did() { + let h = Harness::new().await; + let owner = did("did:plc:nel"); + let repo_did = did("did:plc:limpet"); + repo_fixture(&h, &owner, &repo_did).await; + + let app = router(h.state.clone()); + let (status, body) = json_response( + app.oneshot(enrich_request(json!({ + "xrpc": "sh.tangled.repo.listRepos", + "params": { "subject": owner.as_ref() }, + "enrich": [ + { "source": "sh.tangled.feed.star:subject", "type": "count" }, + { "source": "sh.tangled.feed.star:subject", "type": "distinctAuthors" }, + { "source": "sh.tangled.repo.issue:subject", "type": "count" } + ] + }))) + .await + .unwrap(), + ) + .await; + + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(body["output"]["items"].as_array().unwrap().len(), 1); + let stats = &body["stats"]; + assert_eq!( + stats["did:plc:limpet"]["sh.tangled.feed.star:subject"]["count"], + json!(3) + ); + assert_eq!( + stats["did:plc:limpet"]["sh.tangled.feed.star:subject"]["distinctAuthors"], + json!(2) + ); + assert_eq!( + stats["did:plc:limpet"]["sh.tangled.repo.issue:subject"]["count"], + json!(1) + ); + assert!(stats["at://did:plc:nel/sh.tangled.repo/reef"].is_null()); +} + +#[tokio::test] +async fn follow_counts_cover_both_directions() { + let h = Harness::new().await; + let owner = did("did:plc:nel"); + // followers, edges pointing at owner + for (i, fan) in ["did:plc:a", "did:plc:b"].iter().enumerate() { + h.add_edge( + "sh.tangled.graph.follow", + SubjectRef::Did(owner.clone()), + &at(&format!("at://{fan}/sh.tangled.graph.follow/f{i}")), + ); + h.mount( + &did(fan), + "sh.tangled.graph.follow", + &format!("f{i}"), + follow_body(&owner), + ) + .await; + } + // following, via the .by mirror edge since owner is the author here + h.add_edge( + "sh.tangled.graph.follow.by", + SubjectRef::Did(owner.clone()), + &at("at://did:plc:nel/sh.tangled.graph.follow/f0"), + ); + + let app = router(h.state.clone()); + let (status, body) = json_response( + app.oneshot(enrich_request(json!({ + "xrpc": "sh.tangled.graph.listFollows", + "params": { "subject": owner.as_ref() }, + "enrich": [{ "source": "sh.tangled.graph.follow:subject", "type": "count" }, { "source": "sh.tangled.graph.follow:.repo", "type": "count" }] + }))) + .await + .unwrap(), + ) + .await; + + assert_eq!(status, StatusCode::OK, "{body}"); + let nel = &body["stats"]["did:plc:nel"]; + assert_eq!( + nel["sh.tangled.graph.follow:subject"]["count"], + json!(2), + "{body}" + ); + assert_eq!( + nel["sh.tangled.graph.follow:.repo"]["count"], + json!(1), + "{body}" + ); +} + +#[tokio::test] +async fn sources_scope_which_refs_get_enriched() { + let h = Harness::new().await; + let owner = did("did:plc:nel"); + let repo_did = did("did:plc:limpet"); + repo_fixture(&h, &owner, &repo_did).await; + + let app = router(h.state.clone()); + let (status, body) = json_response( + app.oneshot(enrich_request(json!({ + "xrpc": "sh.tangled.repo.listRepos", + "params": { "subject": owner.as_ref() }, + "enrich": [{ "source": "sh.tangled.feed.star:subject", "type": "count" }], + "sources": ["items[].value.repoDid"] + }))) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::OK, "{body}"); + let stats = &body["stats"]; + assert_eq!( + stats["did:plc:limpet"]["sh.tangled.feed.star:subject"]["count"], + json!(3) + ); + assert_eq!(stats.as_object().unwrap().len(), 1, "{body}"); + + // a path matching nothing is empty stats, not an error, since selection is vector-matched + let app = router(h.state.clone()); + let (status, body) = json_response( + app.oneshot(enrich_request(json!({ + "xrpc": "sh.tangled.repo.listRepos", + "params": { "subject": owner.as_ref() }, + "enrich": [{ "source": "sh.tangled.feed.star:subject", "type": "count" }], + "sources": ["items[].value.nope"] + }))) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(body["stats"], json!({})); +} + +#[tokio::test] +async fn inner_record_miss_passes_through_as_404() { + let h = Harness::new().await; + let app = router(h.state.clone()); + let (status, body) = json_response( + app.oneshot(enrich_request(json!({ + "xrpc": "sh.tangled.repo.getRepo", + "params": { "repo": "at://did:plc:nel/sh.tangled.repo/absent" }, + "enrich": [{ "source": "sh.tangled.feed.star:subject", "type": "count" }] + }))) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::NOT_FOUND, "{body}"); + assert_eq!(body["error"], json!("RecordNotFound")); +} + +#[tokio::test] +async fn rejects_bad_requests() { + let h = Harness::new().await; + let cases = [ + json!({ "xrpc": "sh.tangled.nope.nope", "enrich": [] }), + json!({ + "xrpc": "sh.tangled.repo.countRepos", + "params": { "subject": "did:plc:nel" }, + "enrich": [{ "source": "sh.tangled.nope:subject", "type": "count" }] + }), + json!({ + "xrpc": "sh.tangled.repo.countRepos", + "params": { "subject": "did:plc:nel" }, + "enrich": [{ "source": "sh.tangled.feed.star:subject.uri", "type": "count" }] + }), + json!({ + "xrpc": "sh.tangled.repo.countRepos", + "params": { "subject": "did:plc:nel" }, + "enrich": [{ "source": "sh.tangled.feed.star:.rkey", "type": "count" }] + }), + json!({ + "xrpc": "sh.tangled.repo.countRepos", + "params": { "subject": "did:plc:nel" }, + "enrich": [{ "source": "sh.tangled.feed.star:subject", "type": "count" }], + "sources": ["items["] + }), + // no defaults, a source without a colon is rejected + json!({ + "xrpc": "sh.tangled.repo.countRepos", + "params": { "subject": "did:plc:nel" }, + "enrich": [{ "source": "sh.tangled.feed.star", "type": "count" }] + }), + ]; + for case in cases { + let app = router(h.state.clone()); + let (status, body) = json_response(app.oneshot(enrich_request(case)).await.unwrap()).await; + assert_eq!(status, StatusCode::BAD_REQUEST, "{body}"); + assert_eq!(body["error"], json!("InvalidRequest"), "{body}"); + } + + // serde rejects a missing type before the handler sees it + // that returns plain-text 422, not our 400 json body + let app = router(h.state.clone()); + let response = app + .oneshot(enrich_request(json!({ + "xrpc": "sh.tangled.repo.countRepos", + "params": { "subject": "did:plc:nel" }, + "enrich": [{ "source": "sh.tangled.feed.star:subject" }] + }))) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::UNPROCESSABLE_ENTITY); +} + +#[tokio::test] +async fn viewer_aggregation_uses_explicit_viewer_param() { + let h = Harness::new().await; + let owner = did("did:plc:abc"); + let repo_did = did("did:plc:limpet"); + repo_fixture(&h, &owner, &repo_did).await; + + // the viewer already starred this repo, for the checks below + let subject = SubjectRef::Did(repo_did.clone()); + h.state.edges.add(Edge { + kind: nsid("sh.tangled.feed.star"), + subject, + source: at("at://did:plc:nel/sh.tangled.feed.star/r99"), + sort_micros: 99, + }); + + let app = router(h.state.clone()); + + // missing viewer param is a 400, viewer descriptors require it + let no_viewer = json!({ + "xrpc": "sh.tangled.repo.listRepos", + "params": { "subject": owner.as_ref() }, + "enrich": [{ "source": "sh.tangled.feed.star:subject", "type": "viewer" }] + }); + let (status, _) = json_response( + app.clone() + .oneshot(enrich_request(no_viewer)) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::BAD_REQUEST); + + // viewer who starred it gets their own star uri back + let starred_viewer = json!({ + "xrpc": "sh.tangled.repo.listRepos", + "params": { "subject": owner.as_ref() }, + "enrich": [{ "source": "sh.tangled.feed.star:subject", "type": "viewer" }], + "viewer": "did:plc:nel" + }); + let (status, resp) = json_response( + app.clone() + .oneshot(enrich_request(starred_viewer)) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::OK); + let stats = &resp["stats"][repo_did.as_str()]["sh.tangled.feed.star:subject"]; + assert_eq!( + stats["viewer"], + json!("at://did:plc:nel/sh.tangled.feed.star/r99") + ); + + // a viewer who never starred it gets an explicit null, not absent + let other_viewer = json!({ + "xrpc": "sh.tangled.repo.listRepos", + "params": { "subject": owner.as_ref() }, + "enrich": [{ "source": "sh.tangled.feed.star:subject", "type": "viewer" }], + "viewer": "did:plc:someoneelse" + }); + let (status, resp) = + json_response(app.oneshot(enrich_request(other_viewer)).await.unwrap()).await; + assert_eq!(status, StatusCode::OK); + let stats = &resp["stats"][repo_did.as_str()]["sh.tangled.feed.star:subject"]; + assert_eq!(stats["viewer"], Value::Null); +} diff --git a/lexicons/query/enrichResponse.json b/lexicons/query/enrichResponse.json new file mode 100644 index 00000000..c8a9a8f3 --- /dev/null +++ b/lexicons/query/enrichResponse.json @@ -0,0 +1,104 @@ +{ + "lexicon": 1, + "id": "sh.tangled.query.enrichResponse", + "defs": { + "main": { + "type": "procedure", + "description": "like calling an xrpc query directly, but the response comes back with a stats sidecar: aggregate counts (stars, follows, issues...) for every did / at-uri / strong-ref found in the response. counts are computed synchronously from the server's local index, so there are no pending states to resolve later. any GET query the server implements can be hydrated, including queries it proxies elsewhere.", + "input": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": [ + "xrpc", + "enrich" + ], + "properties": { + "xrpc": { + "type": "string", + "format": "nsid", + "description": "NSID of the inner query to run, e.g. sh.tangled.repo.listRepos." + }, + "params": { + "type": "unknown", + "description": "Parameters for the inner query, exactly as it declares them." + }, + "enrich": { + "type": "array", + "items": { + "type": "ref", + "ref": "#linkDescriptor" + }, + "description": "Counts to compute for each found reference. A descriptor only applies to references it can meaningfully point at: a did can carry author-side counts (.repo), an at-uri can carry subject-side counts." + }, + "sources": { + "type": "array", + "items": { + "type": "string" + }, + "description": "RecordPaths into the inner response restricting which references get counted. When omitted, every did / at-uri / strong-ref in the response is counted. Paths that match nothing yield empty stats rather than an error." + }, + "viewer": { + "type": "string", + "format": "did", + "description": "DID whose own relation to each reference is looked up for viewer-type descriptors, e.g. whether this did starred the found repo and with which record. Required when any enrich descriptor uses type viewer." + } + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": [ + "output", + "stats" + ], + "properties": { + "output": { + "type": "unknown", + "description": "The inner query's response, unmodified." + }, + "stats": { + "type": "unknown", + "description": "Reference value (did or at-uri) -> link source (echoed verbatim from the request) -> aggregation results. Example: stats[\"did:plc:...\"][\"sh.tangled.graph.follow:subject\"] = {\"count\": 12, \"distinctAuthors\": 11}." + } + } + } + }, + "errors": [ + { + "name": "UnknownQuery", + "description": "The xrpc parameter does not name a query this server can execute." + }, + { + "name": "InvalidSourcePath", + "description": "A sources entry is not valid RecordPath syntax." + }, + { + "name": "InvalidLinkDescriptor", + "description": "An enrich entry names an unknown collection, or a path the server's index does not support." + } + ] + }, + "linkDescriptor": { + "type": "object", + "required": [ + "source", + "type" + ], + "description": "one aggregate count: records whose reference field equals the found reference, addressed by a constellation-style link source (\"collection:recordpath\"). envelope paths (.repo etc.) address the record's own metadata instead of its contents.", + "properties": { + "source": { + "type": "string", + "description": "Link source: the collection whose records are counted, a colon, then a RecordPath naming the reference field (e.g. sh.tangled.feed.star:subject, sh.tangled.feed.star:subject.uri), or an envelope field: .repo (authoring repo), .collection, .rkey, or bare . (the record's own at-uri)." + }, + "type": { + "type": "string", + "knownValues": ["count", "distinctAuthors", "viewer"], + "description": "What to compute over matching records: how many records (count), how many distinct authors (distinctAuthors), or the viewer's own record among the matches (viewer, requires the top-level viewer param). Results land in the stats sidecar at stats[ref][source][type]." + } + } + } + } +} diff --git a/web/lex.config.ts b/web/lex.config.ts index e5851854..2ef052c9 100644 --- a/web/lex.config.ts +++ b/web/lex.config.ts @@ -20,6 +20,7 @@ export default defineLexiconConfig({ "../lexicons/markup/**/*.json", "../lexicons/pipeline/**/*.json", "../lexicons/pulls/**/*.json", + "../lexicons/query/**/*.json", "../lexicons/repo/**/*.json", "../lexicons/spindle/**/*.json", "../lexicons/string/**/*.json", diff --git a/web/src/lib/api/lexicons/index.ts b/web/src/lib/api/lexicons/index.ts index 2fa56575..a59f2eee 100644 --- a/web/src/lib/api/lexicons/index.ts +++ b/web/src/lib/api/lexicons/index.ts @@ -68,6 +68,7 @@ export * as ShTangledPipelineListStatusesBy from "./types/sh/tangled/pipeline/li export * as ShTangledPipelineStatus from "./types/sh/tangled/pipeline/status.js"; export * as ShTangledPublicKey from "./types/sh/tangled/publicKey.js"; export * as ShTangledPublicKeyListKeys from "./types/sh/tangled/publicKey/listKeys.js"; +export * as ShTangledQueryEnrichResponse from "./types/sh/tangled/query/enrichResponse.js"; export * as ShTangledRepo from "./types/sh/tangled/repo.js"; export * as ShTangledRepoAddCollaborator from "./types/sh/tangled/repo/addCollaborator.js"; export * as ShTangledRepoAddSecret from "./types/sh/tangled/repo/addSecret.js"; diff --git a/web/src/lib/api/lexicons/types/sh/tangled/query/enrichResponse.ts b/web/src/lib/api/lexicons/types/sh/tangled/query/enrichResponse.ts new file mode 100644 index 00000000..81196fa8 --- /dev/null +++ b/web/src/lib/api/lexicons/types/sh/tangled/query/enrichResponse.ts @@ -0,0 +1,91 @@ +import type {} from "@atcute/lexicons"; +import * as v from "@atcute/lexicons/validations"; +import type {} from "@atcute/lexicons/ambient"; + +const _linkDescriptorSchema = /*#__PURE__*/ v.object({ + $type: /*#__PURE__*/ v.optional( + /*#__PURE__*/ v.literal("sh.tangled.query.enrichResponse#linkDescriptor"), + ), + /** + * Link source: the collection whose records are counted, a colon, then a RecordPath naming the reference field (e.g. sh.tangled.feed.star:subject, sh.tangled.feed.star:subject.uri), or an envelope field: .repo (authoring repo), .collection, .rkey, or bare . (the record's own at-uri). + */ + source: /*#__PURE__*/ v.string(), + /** + * What to compute over matching records: how many records (count), how many distinct authors (distinctAuthors), or the viewer's own record among the matches (viewer, requires the top-level viewer param). Results land in the stats sidecar at stats[ref][source][type]. + */ + type: /*#__PURE__*/ v.string< + "count" | "distinctAuthors" | "viewer" | (string & {}) + >(), +}); +const _mainSchema = /*#__PURE__*/ v.procedure( + "sh.tangled.query.enrichResponse", + { + params: null, + input: { + type: "lex", + schema: /*#__PURE__*/ v.object({ + /** + * Counts to compute for each found reference. A descriptor only applies to references it can meaningfully point at: a did can carry author-side counts (.repo), an at-uri can carry subject-side counts. + */ + get enrich() { + return /*#__PURE__*/ v.array(linkDescriptorSchema); + }, + /** + * Parameters for the inner query, exactly as it declares them. + */ + params: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.unknown()), + /** + * RecordPaths into the inner response restricting which references get counted. When omitted, every did / at-uri / strong-ref in the response is counted. Paths that match nothing yield empty stats rather than an error. + */ + sources: /*#__PURE__*/ v.optional( + /*#__PURE__*/ v.array(/*#__PURE__*/ v.string()), + ), + /** + * DID whose own relation to each reference is looked up for viewer-type descriptors, e.g. whether this did starred the found repo and with which record. Required when any enrich descriptor uses type viewer. + */ + viewer: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.didString()), + /** + * NSID of the inner query to run, e.g. sh.tangled.repo.listRepos. + */ + xrpc: /*#__PURE__*/ v.nsidString(), + }), + }, + output: { + type: "lex", + schema: /*#__PURE__*/ v.object({ + /** + * The inner query's response, unmodified. + */ + output: /*#__PURE__*/ v.unknown(), + /** + * Reference value (did or at-uri) -> link source (echoed verbatim from the request) -> aggregation results. Example: stats["did:plc:..."]["sh.tangled.graph.follow:subject"] = {"count": 12, "distinctAuthors": 11}. + */ + stats: /*#__PURE__*/ v.unknown(), + }), + }, + }, +); + +type linkDescriptor$schematype = typeof _linkDescriptorSchema; +type main$schematype = typeof _mainSchema; + +export interface linkDescriptorSchema extends linkDescriptor$schematype {} +export interface mainSchema extends main$schematype {} + +export const linkDescriptorSchema = + _linkDescriptorSchema as linkDescriptorSchema; +export const mainSchema = _mainSchema as mainSchema; + +export interface LinkDescriptor extends v.InferInput< + typeof linkDescriptorSchema +> {} + +export interface $params {} +export interface $input extends v.InferXRPCBodyInput {} +export interface $output extends v.InferXRPCBodyInput {} + +declare module "@atcute/lexicons/ambient" { + interface XRPCProcedures { + "sh.tangled.query.enrichResponse": mainSchema; + } +}