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::{EdgeKey, SubjectRef, nsid_static, owner_did_from_aturi}; 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_axum::service_auth::ExtractOptionalServiceAuth; use jacquard_common::types::string::{AtUri, Datetime, Did, Handle, UriValue}; use jacquard_common::xrpc::XrpcResp; use jacquard_common::{DefaultStr, IntoStatic as _}; use jacquard_identity::resolver::IdentityResolver; use crate::{ AppState, XrpcError, XrpcQuery, fetch, json_stream, paged_tail, parse_cursor, parse_limit, }; 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 { #[serde(default)] following_only: bool, limit: Option, cursor: Option, } #[axum::debug_handler] pub(crate) async fn get_timeline( State(state): State, ExtractOptionalServiceAuth(auth): ExtractOptionalServiceAuth, 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 = auth.as_ref().map(|auth| auth.did().into_static()); // 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.as_ref().ok_or_else(|| { XrpcError::AuthRequired("followingOnly requires an authenticated viewer".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.as_ref(), 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") }) }