diff --git a/Cargo.lock b/Cargo.lock index 821778b86..56f6af71f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -995,6 +995,7 @@ dependencies = [ "chrono", "futures", "http", + "jacquard-axum", "jacquard-common", "jacquard-identity", "reqwest 0.13.1", diff --git a/bobbin/crates/xrpc/Cargo.toml b/bobbin/crates/xrpc/Cargo.toml index b94b36158..539f069c4 100644 --- a/bobbin/crates/xrpc/Cargo.toml +++ b/bobbin/crates/xrpc/Cargo.toml @@ -14,6 +14,7 @@ bobbin-runtime = { workspace = true } bobbin-search = { workspace = true } bobbin-slingshot-client = { workspace = true } bobbin-knot-proxy = { workspace = true } +jacquard-axum = { workspace = true } jacquard-common = { workspace = true } jacquard-identity = { workspace = true } diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index 05c9f3e2c..334232dc2 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -108,6 +108,8 @@ mod enrich; mod feed; mod filter; mod recordpath; +mod repo; +mod view; pub use backpressure::{ HeavyLimiter, HeavyPermit, MaxInFlight, PerRequestAnonBytes, PressureVerdict, ReservedFloor, @@ -410,6 +412,8 @@ pub fn router(state: AppState) -> Router { "/xrpc/sh.tangled.repo.pull.countStatusesBy", get(count_pull_statuses_by), ) + .route("/xrpc/sh.tangled.pull.getPullView", get(repo::get_pull_view)) + // .route("/xrpc/sh.tangled.pull.listPullViews", get(repo::list_pull_views)) .route( "/xrpc/sh.tangled.spindle.listMembersBy", get(list_spindle_members_by), diff --git a/bobbin/crates/xrpc/src/repo.rs b/bobbin/crates/xrpc/src/repo.rs new file mode 100644 index 000000000..322209087 --- /dev/null +++ b/bobbin/crates/xrpc/src/repo.rs @@ -0,0 +1,139 @@ +#![allow(unused)] +use crate::{ + AppState, SubjectQuery, TypedListQuery, XrpcError, accept_state_source, fetch, filter::ListFilter, record_edge_page, view::{build_profile_basic, build_repo_view_basic} +}; +use axum::extract::State; +use bobbin_edge_index::{EdgeItem, StateKind}; +use bobbin_types::{ + ids::owner_did_from_aturi, + sh_tangled::{ + markup::markdown::Markdown, + pull::{self, PatchViewBuilder, PullView, get_pull_view, list_pull_views}, + repo::pull::{Pull, PullRecord}, + }, +}; +use futures::{stream, StreamExt as _}; +use jacquard_axum::{ExtractXrpc, XrpcResponse}; + +const HYDRATE_CONCURRENCY: usize = 8; + +// impl ListFilter for list_pull_views::ListPullViews { +// fn predicate( +// &self, +// state: &AppState, +// subject: &bobbin_types::ids::SubjectRef, +// ) -> crate::filter::FilterPredicate { +// // compose_state_filter::( +// // self.author.clone(), +// // self.status.map(Into::into), +// // state.pull_statuses.clone(), +// // subject.as_did().cloned(), +// // ) +// todo!() +// } +// +// fn is_identity(&self) -> bool { +// self.state.is_none() +// } +// } + +pub(crate) async fn get_pull_view( + State(state): State, + ExtractXrpc(params): ExtractXrpc, +) -> Result, XrpcError> { + let (body, pull) = fetch::(&state, ¶ms.uri).await?; + + let viewer = None; + + let author_did = owner_did_from_aturi(¶ms.uri) + .ok_or_else(|| XrpcError::InvalidParams("pull uri authority must be a did".into()))?; + let author = build_profile_basic(&state, viewer, &author_did).await; + + let repo = build_repo_view_basic(&state, viewer, &pull.target.repo).await?; + + // TODO(cob): state will be embedded to ticket record. + let pull_state = state + .pull_statuses + .latest_by(¶ms.uri, |src| { + accept_state_source(src, Some(&author_did), &pull.target.repo) + }) + .map(|(kind, _)| kind.wire().into()); + + Ok(XrpcResponse(get_pull_view::GetPullViewOutput { + value: pull::PullViewDetailed { + uri: body.uri.clone(), + cid: body.cid.clone(), + repo, + author, + title: pull.title, + body: pull.body.map(|body| Markdown { + text: body, + ..Default::default() + }), + created_at: pull.created_at, + source: pull.source.map(|source| pull::Source { + branch: source.branch, + repo: None, + extra_data: Default::default(), + }), + target_branch: pull.target.branch, + versions: pull.versions.into_iter().map(|version| { + PatchViewBuilder::new() + .base(version.base) + .head(version.head) + .created_at(version.created_at) + .build() + }).collect(), + state: pull_state, + extra_data: Default::default(), + }, + extra_data: Default::default(), + })) +} + +// pub(crate) async fn list_pull_views( +// State(state): State, +// ExtractXrpc(params): ExtractXrpc, +// ) -> Result, XrpcError> { +// let q = TypedListQuery { +// subject: SubjectQuery::Did(params.subject.clone()), +// cursor: params.cursor.as_ref().map(|cursor| cursor.to_string()), +// limit: params.limit.as_ref().map(|limit| *limit as u32), +// order: Default::default(), +// filter: params, +// }; +// +// let (page, _nsid) = record_edge_page::(&state, &q)?; +// +// // Hydrate skeleton +// let owned = state.clone(); +// let hydrated = stream::iter( +// page.items +// .iter() +// .map(|it| hydrate_item(&owned, it)) +// .collect::>(), +// ) +// .buffered(HYDRATE_CONCURRENCY) +// .collect::>() +// .await; +// +// let mut pulls: Vec = Vec::with_capacity(hydrated.len()); +// for (item, r) in page.items.iter().zip(hydrated) { +// match r { +// Ok(item) => pulls.push(item), +// Err(e) => { +// tracing::warn!(uri = %item.uri, error = %e, "skipping pullView item, hydration failed") +// } +// } +// } +// +// Ok(XrpcResponse(list_pull_views::ListPullViewsOutput { +// items: pulls, +// cursor: None, +// extra_data: Default::default(), +// })) +// } +// +// async fn hydrate_item(state: &AppState, item: &EdgeItem) -> Result { +// todo!("Edge -> PullView") +// } diff --git a/bobbin/crates/xrpc/src/view.rs b/bobbin/crates/xrpc/src/view.rs new file mode 100644 index 000000000..31136a799 --- /dev/null +++ b/bobbin/crates/xrpc/src/view.rs @@ -0,0 +1,125 @@ +use bobbin_edge_index::EdgeStore; +use bobbin_types::{ + ids::{nsid_static, owner_did_from_aturi, EdgeKey, SubjectRef}, + sh_tangled::{ + actor, + feed::star::StarRecord, + graph::follow::FollowRecord, + repo::{self, Repo, RepoRecord}, + }, +}; +use jacquard_common::{ + types::{aturi::AtUri, did::Did, string::Handle}, + xrpc::XrpcResp as _, + DefaultStr, +}; +use tracing::error; + +use crate::{fetch, AppState, XrpcError}; + +pub(crate) async fn build_profile_basic( + state: &AppState, + viewer: Option<&Did>, + did: &Did, +) -> actor::ProfileViewBasic { + let handle = resolve_handle(state, did); + let avatar = None; // state.avatar.as_ref().and_then(|s| s.url(did)); + let viewer_state = viewer + .filter(|viewer| *viewer != did) + .map(|viewer| build_actor_viewer_state(&state.edges, viewer, did)); + actor::ProfileViewBasic::new() + .did(did.clone()) + .handle(handle) + .maybe_avatar(avatar) + .maybe_viewer(viewer_state) + .build() +} + +fn build_actor_viewer_state(edges: &EdgeStore, viewer: &Did, subject: &Did) -> actor::ViewerState { + let following = edges.viewer_source( + &EdgeKey::new( + nsid_static(FollowRecord::NSID), + SubjectRef::Did(subject.clone()), + ), + viewer.as_str(), + ); + let followed_by = edges.viewer_source( + &EdgeKey::new( + nsid_static(FollowRecord::NSID), + SubjectRef::Did(viewer.clone()), + ), + subject.as_str(), + ); + actor::ViewerState { + following, + followed_by, + ..Default::default() + } +} + +pub(crate) async fn build_repo_view_basic( + state: &AppState, + viewer: Option<&Did>, + did: &Did, +) -> Result { + let ident = state + .resolver + .lookup_by_repo_did(did) + .await + .ok_or(XrpcError::NotFound)?; + let uri = AtUri::::from_parts_owned( + ident.owner.as_str(), + RepoRecord::NSID, + ident.rkey.as_str(), + ) + .map_err(|e| XrpcError::InvalidRecord(format!("repo uri assembly: {e}")))?; + let (_, repo) = fetch::>(state, &uri).await?; + build_repo_view_basic_inner(state, viewer, &uri, &repo, did.clone()).await +} + +async fn build_repo_view_basic_inner( + state: &AppState, + viewer: Option<&Did>, + repo_uri: &AtUri, + repo_record: &Repo, + repo_did: Did, +) -> Result { + let owner_did = owner_did_from_aturi(repo_uri) + .ok_or_else(|| XrpcError::InvalidParams("repo uri authority must be a did".into()))?; + let owner = build_profile_basic(state, viewer, &owner_did).await; + let slug = repo_record.name.clone().unwrap_or_else(|| { + repo_uri + .rkey() + .map(|r| DefaultStr::from(r.as_ref())) + .unwrap_or_default() + }); + let star_key = EdgeKey::new( + nsid_static(StarRecord::NSID), + SubjectRef::Did(repo_did.clone()), + ); + let star_count = state.edges.count(&star_key) as i64; + let viewer = viewer.map(|v| { + let mut viewer = repo::ViewerState::default(); + viewer.star = state.edges.viewer_source(&star_key, v.as_str()); + viewer + }); + Ok(repo::RepoViewBasic::new() + .did(repo_did) + .owner(owner) + .slug(slug) + .created_at(repo_record.created_at.clone()) + .description(repo_record.description.clone()) + .star_count(star_count) + .viewer(viewer) + .build()) +} + +fn resolve_handle(state: &AppState, did: &Did) -> Handle { + match state.identity.get_by_did(did) { + Ok(doc) => doc.handle, + Err(err) => { + error!(error = %err, did = %did, "failed to resolve did"); + Handle::new_static("handle.invalid").expect("handle.invalid is a valid handle") + } + } +}