diff --git a/bobbin/crates/xrpc/src/feed.rs b/bobbin/crates/xrpc/src/feed.rs new file mode 100644 index 00000000..4a78c416 --- /dev/null +++ b/bobbin/crates/xrpc/src/feed.rs @@ -0,0 +1,367 @@ +use axum::extract::State; +use axum::response::Response; +use futures::stream::{self, StreamExt}; +use serde::Deserialize; + +use bobbin_edge_index::{EdgeItem, EdgePage, EdgeStore, PageCursor, PageLimit, PageToken, SortDir}; +use bobbin_types::ids::{nsid_static, owner_did_from_aturi, EdgeKey, SubjectRef}; +use bobbin_types::sh_tangled::actor::{ProfileViewBasic, ProfileViewDetailed, ViewerState}; +use bobbin_types::sh_tangled::feed::get_timeline::{ + FollowEvent, RepoEvent, StarEvent, TimelineItem, TimelineItemEvent, +}; +use bobbin_types::sh_tangled::feed::star::{Star, StarRecord, StarSubject}; +use bobbin_types::sh_tangled::graph::follow::{Follow, FollowRecord}; +use bobbin_types::sh_tangled::repo::{self, Repo, RepoRecord, RepoViewBasic}; +use jacquard_common::types::string::{AtUri, Datetime, Did, Handle, UriValue}; +use jacquard_common::xrpc::XrpcResp; +use jacquard_common::DefaultStr; +use jacquard_identity::resolver::IdentityResolver; + +use crate::{ + fetch, json_stream, paged_tail, parse_cursor, parse_limit, AppState, XrpcError, XrpcQuery, +}; + +const REPO_NSID: &str = "sh.tangled.repo"; +const STAR_NSID: &str = "sh.tangled.feed.star"; +const FOLLOW_NSID: &str = "sh.tangled.graph.follow"; +const STAR_BY_NSID: &str = "sh.tangled.feed.star.by"; +const FOLLOW_BY_NSID: &str = "sh.tangled.graph.follow.by"; + +const HYDRATE_CONCURRENCY: usize = 8; + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct GetTimelineQuery { + viewer: Option>, + #[serde(default)] + following_only: bool, + limit: Option, + cursor: Option, +} + +pub(crate) async fn get_timeline( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result { + let limit = parse_limit(q.limit)?; + let cursor = parse_cursor(q.cursor.as_deref())?; + let permit = state.heavy_permit()?; + let viewer = q.viewer.as_ref(); + + // Prepare timeline skeleton + let page = if q.following_only { + // Following feed: the viewer's followed set, then a time-ordered k-way merge. + let viewer = viewer.ok_or_else(|| { + XrpcError::InvalidParams("followingOnly requires a viewer did".into()) + })?; + let followed = followed_dids(&state, viewer).await; + following_skeleton(&state, &followed, cursor, limit) + } else { + global_skeleton(&state, cursor, limit) + }; + + // Hydrate skeleton + let hydrated = stream::iter( + page.items + .iter() + .map(|it| hydrate_item(&state, viewer, it)) + .collect::>(), + ) + .buffered(HYDRATE_CONCURRENCY) + .collect::>() + .await; + + let mut feed: Vec> = Vec::with_capacity(hydrated.len()); + for (item, r) in page.items.iter().zip(hydrated) { + match r { + Ok(item) => feed.push(item), + Err(e) => { + tracing::warn!(uri = %item.uri, error = %e, "skipping timeline item, hydration failed") + } + } + } + + let items = stream::iter(feed.into_iter().map(Ok::<_, XrpcError>)); + Ok(json_stream::, _>( + "feed", + items, + paged_tail(page.next.map(PageToken::encode_token)), + permit, + )) +} + +async fn followed_dids(state: &AppState, viewer: &Did) -> Vec> { + let key = EdgeKey::new(nsid_static(FOLLOW_BY_NSID), SubjectRef::Did(viewer.clone())); + let uris = state.edges.sources_for(&key); + // TODO: store follow actor information in edge + stream::iter(uris.into_iter().map(|uri| async move { + match fetch::>(state, &uri).await { + Ok((_, follow)) => Some(follow.subject), + Err(e) => { + tracing::warn!(uri = %uri, error = %e, "skipping followed did, follow record unresolved"); + None + } + } + })) + .buffered(HYDRATE_CONCURRENCY) + .collect::>() + .await + .into_iter() + .flatten() + .collect() +} + +fn following_skeleton( + state: &AppState, + followed: &[Did], + cursor: PageCursor, + limit: PageLimit, +) -> EdgePage { + let keys: Vec = followed + .iter() + .flat_map(|u| { + [REPO_NSID, STAR_BY_NSID, FOLLOW_BY_NSID] + .map(|kind| EdgeKey::new(nsid_static(kind), SubjectRef::Did(u.clone()))) + }) + .collect(); + state.edges.list_multi(&keys, cursor, limit, SortDir::Desc) +} + +fn global_skeleton(state: &AppState, cursor: PageCursor, limit: PageLimit) -> EdgePage { + let keys = [REPO_NSID, STAR_NSID, FOLLOW_NSID] + .map(|kind| EdgeKey::new(nsid_static(kind), SubjectRef::Global)); + state.edges.list_multi(&keys, cursor, limit, SortDir::Desc) +} + +async fn hydrate_item( + state: &AppState, + viewer: Option<&Did>, + item: &EdgeItem, +) -> Result, XrpcError> { + let collection = item + .uri + .collection() + .ok_or_else(|| XrpcError::InvalidParams("timeline uri missing collection".into()))?; + let (event, event_at) = match collection.as_ref() { + REPO_NSID => hydrate_repo_event(state, viewer, &item.uri).await?, + STAR_NSID => hydrate_star_event(state, viewer, &item.uri).await?, + FOLLOW_NSID => hydrate_follow_event(state, viewer, &item.uri).await?, + other => { + return Err(XrpcError::InvalidRecord(format!( + "unexpected timeline collection: {other}" + ))); + } + }; + Ok(TimelineItem::new().event(event).event_at(event_at).build()) +} + +async fn hydrate_repo_event( + state: &AppState, + viewer: Option<&Did>, + uri: &AtUri, +) -> Result<(TimelineItemEvent, Datetime), XrpcError> { + let (body, repo) = fetch::>(state, uri).await?; + let repo_did = repo.repo_did.clone().ok_or_else(|| { + XrpcError::InvalidRecord("repo has no repo_did; cannot build repoViewBasic".into()) + })?; + let view = build_repo_view_basic(state, viewer, uri, &repo, repo_did).await?; + let source = match &repo.source { + Some(UriValue::Did(did)) => build_repo_view_from_did(state, viewer, did).await.ok(), + Some(UriValue::At(_aturi)) => None, // we drop the legacy support + _ => None, + }; + let event_at = repo.created_at.clone(); + let ev = RepoEvent::new() + .uri(body.uri.clone()) + .cid(body.cid.clone()) + .repo(view) + .source(source) + .build(); + Ok((TimelineItemEvent::RepoEvent(Box::new(ev)), event_at)) +} + +async fn hydrate_star_event( + state: &AppState, + viewer: Option<&Did>, + uri: &AtUri, +) -> Result<(TimelineItemEvent, Datetime), XrpcError> { + let (body, star) = fetch::>(state, uri).await?; + let starrer = owner_did_from_aturi(uri) + .ok_or_else(|| XrpcError::InvalidParams("star uri authority must be a did".into()))?; + let actor = build_profile_basic(state, viewer, &starrer).await; + let repo_did = match &star.subject { + StarSubject::Repo(r) => r.did.clone(), + StarSubject::String(_) => { + return Err(XrpcError::InvalidRecord( + "timeline star subject is not a repo".into(), + )); + } + }; + let repo = build_repo_view_from_did(state, viewer, &repo_did).await?; + let event_at = star.created_at.clone(); + let ev = StarEvent::new() + .uri(body.uri.clone()) + .cid(body.cid.clone()) + .actor(actor) + .repo(repo) + .build(); + Ok((TimelineItemEvent::StarEvent(Box::new(ev)), event_at)) +} + +async fn hydrate_follow_event( + state: &AppState, + viewer: Option<&Did>, + uri: &AtUri, +) -> Result<(TimelineItemEvent, Datetime), XrpcError> { + let (body, follow) = fetch::>(state, uri).await?; + let follower = owner_did_from_aturi(uri) + .ok_or_else(|| XrpcError::InvalidParams("follow uri authority must be a did".into()))?; + let actor = build_profile_basic(state, viewer, &follower).await; + let subject = build_profile_detailed(state, viewer, &follow.subject).await; + let event_at = follow.created_at.clone(); + let ev = FollowEvent::new() + .uri(body.uri.clone()) + .cid(body.cid.clone()) + .actor(actor) + .subject(subject) + .build(); + Ok((TimelineItemEvent::FollowEvent(Box::new(ev)), event_at)) +} + +async fn build_repo_view_from_did( + state: &AppState, + viewer: Option<&Did>, + repo_did: &Did, +) -> Result, XrpcError> { + let ident = state + .resolver + .lookup_by_repo_did(repo_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(state, viewer, &uri, &repo, repo_did.clone()).await +} + +async fn build_repo_view_basic( + state: &AppState, + viewer: Option<&Did>, + repo_uri: &AtUri, + repo_record: &Repo, + repo_did: Did, +) -> Result, XrpcError> { + 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(STAR_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(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()) +} + +async fn build_profile_basic( + state: &AppState, + viewer: Option<&Did>, + did: &Did, +) -> ProfileViewBasic { + let handle = resolve_handle(state, did).await; + let avatar = None; // state.avatar.as_ref().and_then(|s| s.url(did)); + let viewer_state = build_actor_viewer_state(&state.edges, viewer, did); + ProfileViewBasic::new() + .did(did.clone()) + .handle(handle) + .maybe_avatar(avatar) + .maybe_viewer(viewer_state) + .build() +} + +async fn build_profile_detailed( + state: &AppState, + viewer: Option<&Did>, + did: &Did, +) -> ProfileViewDetailed { + let handle = resolve_handle(state, did).await; + let avatar = None; // state.avatar.as_ref().and_then(|s| s.url(did)); + let viewer_state = build_actor_viewer_state(&state.edges, viewer, did); + let followers = state.edges.count(&EdgeKey::new( + nsid_static(FOLLOW_NSID), + SubjectRef::Did(did.clone()), + )) as i64; + let follows = state.edges.count(&EdgeKey::new( + nsid_static(FOLLOW_BY_NSID), + SubjectRef::Did(did.clone()), + )) as i64; + ProfileViewDetailed::new() + .did(did.clone()) + .handle(handle) + .followers_count(followers) + .follows_count(follows) + .maybe_avatar(avatar) + .maybe_viewer(viewer_state) + .build() +} + +fn build_actor_viewer_state( + edges: &EdgeStore, + viewer: Option<&Did>, + subject: &Did, +) -> Option> { + let viewer = viewer?; + if viewer == subject { + return None; + } + let following = edges.viewer_source( + &EdgeKey::new(nsid_static(FOLLOW_NSID), SubjectRef::Did(subject.clone())), + viewer.as_str(), + ); + let followed_by = edges.viewer_source( + &EdgeKey::new(nsid_static(FOLLOW_NSID), SubjectRef::Did(viewer.clone())), + subject.as_str(), + ); + (following.is_some() || followed_by.is_some()).then(|| ViewerState { + following, + followed_by, + ..Default::default() + }) +} + +async fn resolve_handle(state: &AppState, did: &Did) -> Handle { + state + .directory + .resolve_did_doc_owned(did) + .await + .ok() + .and_then(|doc| { + doc.handles() + .into_iter() + .next() + .map(|h| h.as_str().to_owned()) + }) + .and_then(|s| Handle::new_owned(s).ok()) + .unwrap_or_else(|| { + Handle::new_static("handle.invalid").expect("handle.invalid is a valid handle") + }) +} diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index 9839d513..db1e9643 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -101,6 +101,7 @@ use tracing::{Level, Span}; mod backpressure; mod enrich; +mod feed; mod filter; mod recordpath; @@ -176,7 +177,7 @@ impl AppState { .clone() } - fn heavy_permit(&self) -> Result, XrpcError> { + pub(crate) fn heavy_permit(&self) -> Result, XrpcError> { self.limiter.as_ref().map(|l| l.try_enter()).transpose() } } @@ -270,6 +271,7 @@ pub fn router(state: AppState) -> Router { ) .route("/xrpc/sh.tangled.graph.listVouches", get(list_vouches)) .route("/xrpc/sh.tangled.graph.countVouches", get(count_vouches)) + .route("/xrpc/sh.tangled.feed.getTimeline", get(feed::get_timeline)) .route("/xrpc/sh.tangled.feed.listStarsBy", get(list_stars_by)) .route("/xrpc/sh.tangled.feed.countStarsBy", get(count_stars_by)) .route( @@ -1112,12 +1114,12 @@ fn require_rkey(uri: &AtUri, expected: &str) -> Result<(), XrpcError }) } -fn parse_cursor(raw: Option<&str>) -> Result { +pub(crate) fn parse_cursor(raw: Option<&str>) -> Result { PageCursor::from_token(raw) .map_err(|e: CursorParseError| XrpcError::InvalidParams(format!("cursor: {e}"))) } -fn parse_limit(raw: Option) -> Result { +pub(crate) fn parse_limit(raw: Option) -> Result { PageLimit::new(raw.unwrap_or(DEFAULT_LIMIT)) .map_err(|e| XrpcError::InvalidParams(format!("limit: {e}"))) } @@ -1275,7 +1277,7 @@ where Ok((body, value)) } -async fn fetch( +pub(crate) async fn fetch( state: &AppState, uri: &AtUri, ) -> Result<(Arc, V), XrpcError> @@ -1682,7 +1684,7 @@ struct PageState { permit: Option, } -fn paged_tail(cursor: Option) -> Vec { +pub(crate) fn paged_tail(cursor: Option) -> Vec { let encoded = serde_json::to_string(&cursor).unwrap_or_else(|_| "null".to_owned()); format!("],\"cursor\":{encoded}}}").into_bytes() } @@ -1691,7 +1693,7 @@ fn unpaged_tail() -> Vec { b"]}".to_vec() } -fn json_stream( +pub(crate) fn json_stream( array_key: &'static str, items: S, tail: Vec,