From 53b2b419487a75bccadce75440b88cd52987b75b Mon Sep 17 00:00:00 2001 From: Lewis Date: Sun, 10 May 2026 23:36:05 +0300 Subject: [PATCH] feat(xrpc): bulk fetch, xyzBy mirror routes, search filter Lewis: May this revision serve well! --- Cargo.lock | 1 + crates/xrpc/Cargo.toml | 2 + crates/xrpc/src/lib.rs | 1036 ++++++++++++++++++++++++++++-- crates/xrpc/tests/aggregation.rs | 2 +- 4 files changed, 993 insertions(+), 48 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 0ed5699..1c25922 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -493,6 +493,7 @@ dependencies = [ "bobbin-search", "bobbin-slingshot-client", "bobbin-types", + "chrono", "futures", "http", "jacquard-common", diff --git a/crates/xrpc/Cargo.toml b/crates/xrpc/Cargo.toml index 7778ea8..1d3eeb4 100644 --- a/crates/xrpc/Cargo.toml +++ b/crates/xrpc/Cargo.toml @@ -15,6 +15,7 @@ bobbin-knot-proxy = { workspace = true } jacquard-common = { workspace = true } axum = { workspace = true } +chrono = { workspace = true } futures = { workspace = true } http = { workspace = true } serde = { workspace = true } @@ -22,6 +23,7 @@ serde_json = { workspace = true } thiserror = { workspace = true } tower-http = { workspace = true } tracing = { workspace = true } +url = { workspace = true } [dev-dependencies] bobbin-runtime = { workspace = true } diff --git a/crates/xrpc/src/lib.rs b/crates/xrpc/src/lib.rs index 07df204..9ae1393 100644 --- a/crates/xrpc/src/lib.rs +++ b/crates/xrpc/src/lib.rs @@ -4,7 +4,7 @@ use std::sync::Arc; use axum::{ Router, body::Body, - extract::{FromRequestParts, Query, State, rejection::QueryRejection}, + extract::{FromRequestParts, Query, RawQuery, State, rejection::QueryRejection}, http::{ HeaderMap, HeaderName, StatusCode, header::{ @@ -22,17 +22,23 @@ use bobbin_edge_index::{ }; use bobbin_knot_proxy::{KnotHost, KnotProxy, KnotProxyError, ProxyResponse, RepoSlug}; use bobbin_record_lru::RecordStore; -use bobbin_search::{SearchCursor, SearchError, SearchHit, SearchOffset, SearchReader}; +use bobbin_search::{ + SearchCursor, SearchError, SearchFilters, SearchHit, SearchOffset, SearchReader, +}; use bobbin_slingshot_client::{SlingshotClient, SlingshotError}; use bobbin_types::ids::{EdgeKey, nsid_static}; use bobbin_types::record::RecordBody; use bobbin_types::search::SearchableRecord; use bobbin_types::sh_tangled::actor::profile::{Profile, ProfileGetRecordOutput, ProfileRecord}; +use bobbin_types::sh_tangled::feed::reaction::{Reaction, ReactionRecord}; use bobbin_types::sh_tangled::feed::star::{Star, StarRecord}; +use bobbin_types::sh_tangled::git::ref_update::{RefUpdate, RefUpdateRecord}; use bobbin_types::sh_tangled::graph::follow::{Follow, FollowRecord}; +use bobbin_types::sh_tangled::graph::vouch::{Vouch, VouchRecord}; use bobbin_types::sh_tangled::knot::member::{ Member as KnotMember, MemberRecord as KnotMemberRecord, }; +use bobbin_types::sh_tangled::knot::{Knot, KnotRecord}; use bobbin_types::sh_tangled::label::definition::{ Definition as LabelDefinition, DefinitionRecord as LabelDefinitionRecord, }; @@ -41,16 +47,28 @@ use bobbin_types::sh_tangled::pipeline::status::{ Status as PipelineStatus, StatusRecord as PipelineStatusRecord, }; use bobbin_types::sh_tangled::pipeline::{Pipeline, PipelineRecord}; +use bobbin_types::sh_tangled::public_key::{PublicKey, PublicKeyRecord}; use bobbin_types::sh_tangled::repo::artifact::{Artifact, ArtifactRecord}; +use bobbin_types::sh_tangled::repo::collaborator::{Collaborator, CollaboratorRecord}; use bobbin_types::sh_tangled::repo::issue::comment::{ Comment as IssueComment, CommentRecord as IssueCommentRecord, }; +use bobbin_types::sh_tangled::repo::issue::state::{ + State as IssueState, StateRecord as IssueStateRecord, +}; use bobbin_types::sh_tangled::repo::issue::{Issue, IssueGetRecordOutput, IssueRecord}; +use bobbin_types::sh_tangled::repo::pull::comment::{ + Comment as PullComment, CommentRecord as PullCommentRecord, +}; +use bobbin_types::sh_tangled::repo::pull::status::{ + Status as PullStatus, StatusRecord as PullStatusRecord, +}; use bobbin_types::sh_tangled::repo::pull::{Pull, PullGetRecordOutput, PullRecord}; use bobbin_types::sh_tangled::repo::{Repo, RepoGetRecordOutput, RepoRecord}; use bobbin_types::sh_tangled::spindle::member::{ Member as SpindleMember, MemberRecord as SpindleMemberRecord, }; +use bobbin_types::sh_tangled::spindle::{Spindle, SpindleRecord}; use bobbin_types::sh_tangled::string::{TangledString, TangledStringRecord}; use futures::stream::{self, StreamExt, TryStreamExt}; use jacquard_common::types::did::Did; @@ -62,6 +80,7 @@ use jacquard_common::{DefaultStr, IntoStatic}; use serde::{Deserialize, Serialize}; use std::time::Duration; use thiserror::Error; +use url::form_urlencoded; use tower_http::classify::ServerErrorsFailureClass; use tower_http::trace::{DefaultMakeSpan, OnFailure, OnResponse, TraceLayer}; @@ -103,9 +122,13 @@ impl AppState { pub fn router(state: AppState) -> Router { Router::new() .route("/xrpc/sh.tangled.repo.getRepo", get(get_repo)) + .route("/xrpc/sh.tangled.repo.getRepos", get(get_repos)) .route("/xrpc/sh.tangled.actor.getProfile", get(get_profile)) + .route("/xrpc/sh.tangled.actor.getProfiles", get(get_profiles)) .route("/xrpc/sh.tangled.repo.getIssue", get(get_issue)) + .route("/xrpc/sh.tangled.repo.getIssues", get(get_issues)) .route("/xrpc/sh.tangled.repo.getPull", get(get_pull)) + .route("/xrpc/sh.tangled.repo.getPulls", get(get_pulls)) .route("/xrpc/sh.tangled.feed.listStars", get(list_stars)) .route("/xrpc/sh.tangled.feed.countStars", get(count_stars)) .route("/xrpc/sh.tangled.graph.listFollows", get(list_follows)) @@ -122,6 +145,169 @@ pub fn router(state: AppState) -> Router { "/xrpc/sh.tangled.repo.issue.countComments", get(count_issue_comments), ) + .route( + "/xrpc/sh.tangled.repo.pull.listComments", + get(list_pull_comments), + ) + .route( + "/xrpc/sh.tangled.repo.pull.countComments", + get(count_pull_comments), + ) + .route("/xrpc/sh.tangled.feed.listReactions", get(list_reactions)) + .route("/xrpc/sh.tangled.feed.countReactions", get(count_reactions)) + .route("/xrpc/sh.tangled.git.listRefUpdates", get(list_ref_updates)) + .route( + "/xrpc/sh.tangled.git.countRefUpdates", + get(count_ref_updates), + ) + .route( + "/xrpc/sh.tangled.repo.listCollaborators", + get(list_collaborators), + ) + .route( + "/xrpc/sh.tangled.repo.countCollaborators", + get(count_collaborators), + ) + .route( + "/xrpc/sh.tangled.repo.issue.listStates", + get(list_issue_states), + ) + .route( + "/xrpc/sh.tangled.repo.issue.countStates", + get(count_issue_states), + ) + .route( + "/xrpc/sh.tangled.repo.pull.listStatuses", + get(list_pull_statuses), + ) + .route( + "/xrpc/sh.tangled.repo.pull.countStatuses", + get(count_pull_statuses), + ) + .route("/xrpc/sh.tangled.repo.listRepos", get(list_repos)) + .route("/xrpc/sh.tangled.repo.countRepos", get(count_repos)) + .route("/xrpc/sh.tangled.knot.listKnots", get(list_knots)) + .route("/xrpc/sh.tangled.knot.countKnots", get(count_knots)) + .route("/xrpc/sh.tangled.spindle.listSpindles", get(list_spindles)) + .route("/xrpc/sh.tangled.spindle.countSpindles", get(count_spindles)) + .route("/xrpc/sh.tangled.publicKey.listKeys", get(list_public_keys)) + .route("/xrpc/sh.tangled.publicKey.countKeys", get(count_public_keys)) + .route("/xrpc/sh.tangled.graph.listVouches", get(list_vouches)) + .route("/xrpc/sh.tangled.graph.countVouches", get(count_vouches)) + .route("/xrpc/sh.tangled.feed.listStarsBy", get(list_stars_by)) + .route("/xrpc/sh.tangled.feed.countStarsBy", get(count_stars_by)) + .route( + "/xrpc/sh.tangled.feed.listReactionsBy", + get(list_reactions_by), + ) + .route( + "/xrpc/sh.tangled.feed.countReactionsBy", + get(count_reactions_by), + ) + .route("/xrpc/sh.tangled.graph.listFollowsBy", get(list_follows_by)) + .route( + "/xrpc/sh.tangled.graph.countFollowsBy", + get(count_follows_by), + ) + .route("/xrpc/sh.tangled.graph.listVouchesBy", get(list_vouches_by)) + .route( + "/xrpc/sh.tangled.graph.countVouchesBy", + get(count_vouches_by), + ) + .route( + "/xrpc/sh.tangled.git.listRefUpdatesBy", + get(list_ref_updates_by), + ) + .route( + "/xrpc/sh.tangled.git.countRefUpdatesBy", + get(count_ref_updates_by), + ) + .route( + "/xrpc/sh.tangled.knot.listMembersBy", + get(list_knot_members_by), + ) + .route( + "/xrpc/sh.tangled.knot.countMembersBy", + get(count_knot_members_by), + ) + .route("/xrpc/sh.tangled.label.listOpsBy", get(list_label_ops_by)) + .route("/xrpc/sh.tangled.label.countOpsBy", get(count_label_ops_by)) + .route( + "/xrpc/sh.tangled.pipeline.listPipelinesBy", + get(list_pipelines_by), + ) + .route( + "/xrpc/sh.tangled.pipeline.countPipelinesBy", + get(count_pipelines_by), + ) + .route( + "/xrpc/sh.tangled.pipeline.listStatusesBy", + get(list_pipeline_statuses_by), + ) + .route( + "/xrpc/sh.tangled.pipeline.countStatusesBy", + get(count_pipeline_statuses_by), + ) + .route( + "/xrpc/sh.tangled.repo.listArtifactsBy", + get(list_artifacts_by), + ) + .route( + "/xrpc/sh.tangled.repo.countArtifactsBy", + get(count_artifacts_by), + ) + .route( + "/xrpc/sh.tangled.repo.listCollaboratorsBy", + get(list_collaborators_by), + ) + .route( + "/xrpc/sh.tangled.repo.countCollaboratorsBy", + get(count_collaborators_by), + ) + .route("/xrpc/sh.tangled.repo.listIssuesBy", get(list_issues_by)) + .route("/xrpc/sh.tangled.repo.countIssuesBy", get(count_issues_by)) + .route( + "/xrpc/sh.tangled.repo.issue.listCommentsBy", + get(list_issue_comments_by), + ) + .route( + "/xrpc/sh.tangled.repo.issue.countCommentsBy", + get(count_issue_comments_by), + ) + .route( + "/xrpc/sh.tangled.repo.issue.listStatesBy", + get(list_issue_states_by), + ) + .route( + "/xrpc/sh.tangled.repo.issue.countStatesBy", + get(count_issue_states_by), + ) + .route("/xrpc/sh.tangled.repo.listPullsBy", get(list_pulls_by)) + .route("/xrpc/sh.tangled.repo.countPullsBy", get(count_pulls_by)) + .route( + "/xrpc/sh.tangled.repo.pull.listCommentsBy", + get(list_pull_comments_by), + ) + .route( + "/xrpc/sh.tangled.repo.pull.countCommentsBy", + get(count_pull_comments_by), + ) + .route( + "/xrpc/sh.tangled.repo.pull.listStatusesBy", + get(list_pull_statuses_by), + ) + .route( + "/xrpc/sh.tangled.repo.pull.countStatusesBy", + get(count_pull_statuses_by), + ) + .route( + "/xrpc/sh.tangled.spindle.listMembersBy", + get(list_spindle_members_by), + ) + .route( + "/xrpc/sh.tangled.spindle.countMembersBy", + get(count_spindle_members_by), + ) .route( "/xrpc/sh.tangled.label.listDefinitions", get(list_label_definitions), @@ -337,6 +523,10 @@ struct CountQuery { struct SearchQueryParams { q: String, nsid: Option, + author: Option, + repo: Option, + since: Option, + until: Option, cursor: Option, limit: Option, } @@ -424,6 +614,13 @@ struct ListResponse { cursor: Option, } +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +struct BulkResponse { + coverage: CoverageEnvelope, + items: Vec>, +} + #[derive(Serialize)] #[serde(rename_all = "camelCase")] struct RecordView { @@ -479,18 +676,143 @@ fn parse_uri(raw: &str) -> Result, XrpcError> { AtUri::::new_owned(raw).map_err(|e| XrpcError::InvalidParams(format!("uri: {e}"))) } +fn parse_subject_uri(raw: &str) -> Result, XrpcError> { + if let Ok(did) = Did::::new_owned(raw) { + return AtUri::::new_owned(format!("at://{}", did.as_ref())) + .map_err(|e| XrpcError::InvalidParams(format!("uri: {e}"))); + } + parse_uri(raw) +} + #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub enum SubjectShape { BareDid, Collection(&'static str), BareDidOrOneOfCollections(&'static [&'static str]), OneOfCollections(&'static [&'static str]), + AnyAtUri, } pub trait HasSubject { const SHAPE: SubjectShape; } +pub trait MirrorOf { + type Record: XrpcResp; + const EDGE_KIND: &'static str; + 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 IssueBy; +pub struct IssueCommentBy; +pub struct IssueStateBy; +pub struct PullBy; +pub struct PullCommentBy; +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 IssueBy { + type Record = IssueRecord; + const EDGE_KIND: &'static str = "sh.tangled.repo.issue.by"; + const SHAPE: SubjectShape = SubjectShape::BareDid; +} +impl MirrorOf for IssueCommentBy { + type Record = IssueCommentRecord; + const EDGE_KIND: &'static str = "sh.tangled.repo.issue.comment.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 PullCommentBy { + type Record = PullCommentRecord; + const EDGE_KIND: &'static str = "sh.tangled.repo.pull.comment.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", "sh.tangled.repo"]); @@ -507,6 +829,9 @@ impl HasSubject for PullRecord { impl HasSubject for IssueCommentRecord { const SHAPE: SubjectShape = SubjectShape::Collection("sh.tangled.repo.issue"); } +impl HasSubject for PullCommentRecord { + const SHAPE: SubjectShape = SubjectShape::Collection("sh.tangled.repo.pull"); +} impl HasSubject for LabelDefinitionRecord { const SHAPE: SubjectShape = SubjectShape::BareDid; } @@ -532,9 +857,39 @@ impl HasSubject for SpindleMemberRecord { 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::BareDidOrOneOfCollections(&["sh.tangled.repo"]); +} +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; +} +impl HasSubject for VouchRecord { + const SHAPE: SubjectShape = SubjectShape::BareDid; +} fn parse_subject(raw: &RawAtUriParam, shape: SubjectShape) -> Result, XrpcError> { - let uri = parse_uri(raw.as_str())?; + let uri = parse_subject_uri(raw.as_str())?; if matches!(uri.authority(), AtIdentifier::Handle(_)) { return Err(XrpcError::InvalidParams( "subject authority must be a did, not a handle".into(), @@ -578,6 +933,7 @@ fn parse_subject(raw: &RawAtUriParam, shape: SubjectShape) -> Result// with nsid in [{}]", allowed.join(", "), ))), + (SubjectShape::AnyAtUri, _) => Ok(uri), } } @@ -737,6 +1093,98 @@ async fn get_pull( })) } +async fn get_repos( + State(state): State, + RawQuery(query): RawQuery, +) -> Result>>, XrpcError> { + let uris = collect_repeated(query.as_deref(), BULK_REPOS_KEY); + bulk_fetch::>(&state, uris).await.map(Json) +} + +async fn get_profiles( + State(state): State, + RawQuery(query): RawQuery, +) -> Result>>, XrpcError> { + let uris = collect_repeated(query.as_deref(), BULK_PROFILES_KEY); + bulk_fetch::>(&state, uris).await.map(Json) +} + +async fn get_issues( + State(state): State, + RawQuery(query): RawQuery, +) -> Result>>, XrpcError> { + let uris = collect_repeated(query.as_deref(), BULK_ISSUES_KEY); + bulk_fetch::>(&state, uris).await.map(Json) +} + +async fn get_pulls( + State(state): State, + RawQuery(query): RawQuery, +) -> Result>>, XrpcError> { + let uris = collect_repeated(query.as_deref(), BULK_PULLS_KEY); + bulk_fetch::>(&state, uris).await.map(Json) +} + +const BULK_REPOS_KEY: &str = "repos"; +const BULK_PROFILES_KEY: &str = "actors"; +const BULK_ISSUES_KEY: &str = "issues"; +const BULK_PULLS_KEY: &str = "pulls"; +const BULK_LIMIT: usize = 50; + +fn collect_repeated(query: Option<&str>, key: &str) -> Vec { + let Some(q) = query else { + return Vec::new(); + }; + form_urlencoded::parse(q.as_bytes()) + .filter_map(|(k, v)| (k == key).then(|| v.into_owned())) + .collect() +} + +async fn bulk_fetch(state: &AppState, uris: Vec) -> Result, XrpcError> +where + R: XrpcResp, + V: serde::de::DeserializeOwned + Serialize, +{ + if uris.is_empty() { + return Err(XrpcError::InvalidParams("at least one uri required".into())); + } + if uris.len() > BULK_LIMIT { + return Err(XrpcError::InvalidParams(format!( + "at most {BULK_LIMIT} uris per request" + ))); + } + let parsed: Vec> = uris + .iter() + .map(|s| parse_uri(s)) + .collect::>()?; + let coverage = state.coverage.snapshot(); + let items: Vec> = stream::iter(parsed) + .map(|uri| async move { + match resolve(state, ExpectedNsid::new(R::NSID), uri).await { + Ok((body, _)) => match serde_json::from_slice::(&body.value) { + Ok(value) => Ok(Some(RecordView { + uri: body.uri.clone(), + cid: Some(body.cid.clone()), + value, + })), + Err(_) => Ok(None), + }, + Err(XrpcError::NotFound | XrpcError::UpstreamGone(_) | XrpcError::InvalidRecord(_)) => { + Ok(None) + } + Err(other) => Err(other), + } + }) + .buffered(FETCH_CONCURRENCY) + .try_filter_map(|opt| async move { Ok(opt) }) + .try_collect() + .await?; + Ok(BulkResponse { + coverage: coverage.into(), + items, + }) +} + async fn list_records(state: &AppState, q: ListQuery) -> Result, XrpcError> where R: XrpcResp + HasSubject, @@ -783,6 +1231,52 @@ fn count_for( }) } +async fn list_mirror(state: &AppState, q: ListQuery) -> Result, XrpcError> +where + M: MirrorOf, + V: serde::de::DeserializeOwned + Serialize, +{ + let subject = parse_subject(&q.subject, M::SHAPE)?; + let cursor = parse_cursor(q.cursor.as_deref())?; + let limit = parse_limit(q.limit)?; + let key = EdgeKey::new(nsid_static(M::EDGE_KIND), subject); + let coverage = state.coverage.snapshot(); + let EdgePage { items, next } = state.edges.list(&key, cursor, limit); + let items = stream::iter(items) + .map(|uri| async move { + let body = resolve_for_view(state, ::NSID, uri).await?; + let value: V = serde_json::from_slice(&body.value) + .map_err(|e| XrpcError::InvalidRecord(e.to_string()))?; + Ok::<_, XrpcError>(RecordView { + uri: body.uri.clone(), + cid: Some(body.cid.clone()), + value, + }) + }) + .buffered(FETCH_CONCURRENCY) + .try_collect() + .await?; + Ok(ListResponse { + coverage: coverage.into(), + items, + cursor: next.map(SourceId::encode_token), + }) +} + +fn count_mirror( + state: &AppState, + q: CountQuery, +) -> Result { + let subject = parse_subject(&q.subject, M::SHAPE)?; + let key = EdgeKey::new(nsid_static(M::EDGE_KIND), subject); + let coverage = state.coverage.snapshot(); + Ok(CountResponse { + coverage: coverage.into(), + count: state.edges.count(&key), + distinct_authors: state.edges.count_distinct_authors(&key), + }) +} + async fn list_stars( State(state): State, XrpcQuery(q): XrpcQuery, @@ -855,103 +1349,499 @@ async fn count_issue_comments( count_for::(&state, q).map(Json) } -async fn list_label_definitions( +async fn list_pull_comments( State(state): State, XrpcQuery(q): XrpcQuery, -) -> Result>>, XrpcError> { - list_records::(&state, q) +) -> Result>>, XrpcError> { + list_records::(&state, q) .await .map(Json) } -async fn count_label_definitions( +async fn count_pull_comments( State(state): State, XrpcQuery(q): XrpcQuery, ) -> Result, XrpcError> { - count_for::(&state, q).map(Json) + count_for::(&state, q).map(Json) } -async fn list_label_ops( +async fn list_reactions( State(state): State, XrpcQuery(q): XrpcQuery, -) -> Result>>, XrpcError> { - list_records::(&state, q).await.map(Json) +) -> Result>>, XrpcError> { + list_records::(&state, q).await.map(Json) } -async fn count_label_ops( +async fn count_reactions( State(state): State, XrpcQuery(q): XrpcQuery, ) -> Result, XrpcError> { - count_for::(&state, q).map(Json) + count_for::(&state, q).map(Json) } -async fn list_pipelines( +async fn list_ref_updates( State(state): State, XrpcQuery(q): XrpcQuery, -) -> Result>>, XrpcError> { - list_records::(&state, q).await.map(Json) +) -> Result>>, XrpcError> { + list_records::(&state, q).await.map(Json) } -async fn count_pipelines( +async fn count_ref_updates( State(state): State, XrpcQuery(q): XrpcQuery, ) -> Result, XrpcError> { - count_for::(&state, q).map(Json) + count_for::(&state, q).map(Json) } -async fn list_pipeline_statuses( +async fn list_collaborators( State(state): State, XrpcQuery(q): XrpcQuery, -) -> Result>>, XrpcError> { - list_records::(&state, q) +) -> Result>>, XrpcError> { + list_records::(&state, q) .await .map(Json) } -async fn count_pipeline_statuses( +async fn count_collaborators( State(state): State, XrpcQuery(q): XrpcQuery, ) -> Result, XrpcError> { - count_for::(&state, q).map(Json) + count_for::(&state, q).map(Json) } -async fn list_artifacts( +async fn list_issue_states( State(state): State, XrpcQuery(q): XrpcQuery, -) -> Result>>, XrpcError> { - list_records::(&state, q).await.map(Json) +) -> Result>>, XrpcError> { + list_records::(&state, q) + .await + .map(Json) } -async fn count_artifacts( +async fn count_issue_states( State(state): State, XrpcQuery(q): XrpcQuery, ) -> Result, XrpcError> { - count_for::(&state, q).map(Json) + count_for::(&state, q).map(Json) } -async fn list_knot_members( +async fn list_pull_statuses( State(state): State, XrpcQuery(q): XrpcQuery, -) -> Result>>, XrpcError> { - list_records::(&state, q) +) -> Result>>, XrpcError> { + list_records::(&state, q) .await .map(Json) } -async fn count_knot_members( +async fn count_pull_statuses( State(state): State, XrpcQuery(q): XrpcQuery, ) -> Result, XrpcError> { - count_for::(&state, q).map(Json) + count_for::(&state, q).map(Json) } -async fn list_spindle_members( +async fn list_repos( State(state): State, XrpcQuery(q): XrpcQuery, -) -> Result>>, XrpcError> { - list_records::(&state, q) - .await - .map(Json) +) -> Result>>, XrpcError> { + list_records::(&state, q).await.map(Json) +} + +async fn count_repos( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_knots( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q).await.map(Json) +} + +async fn count_knots( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_spindles( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q).await.map(Json) +} + +async fn count_spindles( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_public_keys( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q).await.map(Json) +} + +async fn count_public_keys( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_vouches( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q).await.map(Json) +} + +async fn count_vouches( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_stars_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_stars_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_reactions_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_reactions_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_follows_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_follows_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_vouches_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_vouches_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_ref_updates_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_ref_updates_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_knot_members_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_knot_members_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_label_ops_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_label_ops_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_pipelines_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_pipelines_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_pipeline_statuses_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_pipeline_statuses_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_artifacts_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_artifacts_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_collaborators_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_collaborators_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_issues_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_issues_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_issue_comments_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_issue_comments_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_issue_states_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_issue_states_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_pulls_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_pulls_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_pull_comments_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_pull_comments_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_pull_statuses_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_pull_statuses_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_spindle_members_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_mirror::(&state, q).await.map(Json) +} +async fn count_spindle_members_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_mirror::(&state, q).map(Json) +} + +async fn list_label_definitions( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q) + .await + .map(Json) +} + +async fn count_label_definitions( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_label_ops( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q).await.map(Json) +} + +async fn count_label_ops( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_pipelines( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q).await.map(Json) +} + +async fn count_pipelines( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_pipeline_statuses( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q) + .await + .map(Json) +} + +async fn count_pipeline_statuses( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_artifacts( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q).await.map(Json) +} + +async fn count_artifacts( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_knot_members( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q) + .await + .map(Json) +} + +async fn count_knot_members( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + count_for::(&state, q).map(Json) +} + +async fn list_spindle_members( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result>>, XrpcError> { + list_records::(&state, q) + .await + .map(Json) } async fn count_spindle_members( @@ -1006,18 +1896,11 @@ async fn search_query( let cursor = SearchCursor::from_token(q.cursor.as_deref()) .map_err(|e| XrpcError::InvalidParams(format!("cursor: {e}")))?; let limit = parse_limit(q.limit)?; - let nsid_filter = q - .nsid - .as_deref() - .map(|s| { - Nsid::::new_owned(s) - .map_err(|e| XrpcError::InvalidParams(format!("nsid: {e}"))) - }) - .transpose()?; + let filters = build_search_filters(&q)?; let coverage = state.coverage.snapshot(); let page = state .search - .search(&q.q, nsid_filter.as_ref(), cursor, limit.get()) + .search(&q.q, filters, cursor, limit.get()) .await .map_err(map_search_err)?; let state_ref = &state; @@ -1034,6 +1917,65 @@ async fn search_query( })) } +fn build_search_filters(q: &SearchQueryParams) -> Result { + let nsid = q + .nsid + .as_deref() + .map(|s| { + Nsid::::new_owned(s) + .map_err(|e| XrpcError::InvalidParams(format!("nsid: {e}"))) + }) + .transpose()?; + let author = q + .author + .as_deref() + .map(|s| { + Did::::new_owned(s) + .map_err(|e| XrpcError::InvalidParams(format!("author: {e}"))) + }) + .transpose()?; + let repo = q + .repo + .as_deref() + .map(|s| { + Did::::new_owned(s) + .map_err(|e| XrpcError::InvalidParams(format!("repo: {e}"))) + }) + .transpose()?; + let since = q + .since + .as_deref() + .map(parse_rfc3339_seconds) + .transpose() + .map_err(|e| XrpcError::InvalidParams(format!("since: {e}")))?; + let until = q + .until + .as_deref() + .map(parse_rfc3339_seconds) + .transpose() + .map_err(|e| XrpcError::InvalidParams(format!("until: {e}")))?; + if let (Some(s), Some(u)) = (since, until) + && s > u + { + return Err(XrpcError::InvalidParams( + "since must be <= until".into(), + )); + } + Ok(SearchFilters { + nsid, + author, + repo, + since, + until, + }) +} + +fn parse_rfc3339_seconds(raw: &str) -> Result { + chrono::DateTime::parse_from_rfc3339(raw) + .map(|dt| dt.timestamp()) + .map_err(|e| format!("expected RFC3339, got {raw}: {e}")) +} + async fn hydrate_search_hit( state: &AppState, hit: SearchHit, diff --git a/crates/xrpc/tests/aggregation.rs b/crates/xrpc/tests/aggregation.rs index f0d019d..47da6d9 100644 --- a/crates/xrpc/tests/aggregation.rs +++ b/crates/xrpc/tests/aggregation.rs @@ -881,7 +881,7 @@ async fn repo_pointing_endpoints_accept_repo_uri_subject() { assert_eq!( status, StatusCode::OK, - "{endpoint} must accept a sh.tangled.repo path subject so NoRepoDid repos stay queryable", + "{endpoint} must accept a sh.tangled.repo path subject as input shape, even when the index keys edges on the repoDID", ); } }) -- 2.51.2