This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / bobbin / crates / xrpc / src / feed.rs
13 kB 369 lines
1use axum::extract::State; 2use axum::response::Response; 3use futures::stream::{self, StreamExt}; 4use serde::Deserialize; 5 6use bobbin_edge_index::{EdgeItem, EdgePage, EdgeStore, PageCursor, PageLimit, PageToken, SortDir}; 7use bobbin_types::ids::{EdgeKey, SubjectRef, nsid_static, owner_did_from_aturi}; 8use bobbin_types::sh_tangled::actor::{ProfileViewBasic, ProfileViewDetailed, ViewerState}; 9use bobbin_types::sh_tangled::feed::get_timeline::{ 10 FollowEvent, RepoEvent, StarEvent, TimelineItem, TimelineItemEvent, 11}; 12use bobbin_types::sh_tangled::feed::star::{Star, StarRecord, StarSubject}; 13use bobbin_types::sh_tangled::graph::follow::{Follow, FollowRecord}; 14use bobbin_types::sh_tangled::repo::{self, Repo, RepoRecord, RepoViewBasic}; 15use jacquard_axum::service_auth::ExtractOptionalServiceAuth; 16use jacquard_common::types::string::{AtUri, Datetime, Did, Handle, UriValue}; 17use jacquard_common::xrpc::XrpcResp; 18use jacquard_common::{DefaultStr, IntoStatic as _}; 19use jacquard_identity::resolver::IdentityResolver; 20 21use crate::{ 22 AppState, XrpcError, XrpcQuery, fetch, json_stream, paged_tail, parse_cursor, parse_limit, 23}; 24 25const REPO_NSID: &str = "sh.tangled.repo"; 26const STAR_NSID: &str = "sh.tangled.feed.star"; 27const FOLLOW_NSID: &str = "sh.tangled.graph.follow"; 28const STAR_BY_NSID: &str = "sh.tangled.feed.star.by"; 29const FOLLOW_BY_NSID: &str = "sh.tangled.graph.follow.by"; 30 31const HYDRATE_CONCURRENCY: usize = 8; 32 33#[derive(Debug, Deserialize)] 34#[serde(rename_all = "camelCase")] 35pub(crate) struct GetTimelineQuery { 36 #[serde(default)] 37 following_only: bool, 38 limit: Option<u32>, 39 cursor: Option<String>, 40} 41 42#[axum::debug_handler] 43pub(crate) async fn get_timeline( 44 State(state): State<AppState>, 45 ExtractOptionalServiceAuth(auth): ExtractOptionalServiceAuth, 46 XrpcQuery(q): XrpcQuery<GetTimelineQuery>, 47) -> Result<Response, XrpcError> { 48 let limit = parse_limit(q.limit)?; 49 let cursor = parse_cursor(q.cursor.as_deref())?; 50 let permit = state.heavy_permit()?; 51 let viewer = auth.as_ref().map(|auth| auth.did().into_static()); 52 53 // Prepare timeline skeleton 54 let page = if q.following_only { 55 // Following feed: the viewer's followed set, then a time-ordered k-way merge. 56 let viewer = viewer.as_ref().ok_or_else(|| { 57 XrpcError::AuthRequired("followingOnly requires an authenticated viewer".into()) 58 })?; 59 let followed = followed_dids(&state, &viewer).await; 60 following_skeleton(&state, &followed, cursor, limit) 61 } else { 62 global_skeleton(&state, cursor, limit) 63 }; 64 65 // Hydrate skeleton 66 let hydrated = stream::iter( 67 page.items 68 .iter() 69 .map(|it| hydrate_item(&state, viewer.as_ref(), it)) 70 .collect::<Vec<_>>(), 71 ) 72 .buffered(HYDRATE_CONCURRENCY) 73 .collect::<Vec<_>>() 74 .await; 75 76 let mut feed: Vec<TimelineItem<DefaultStr>> = Vec::with_capacity(hydrated.len()); 77 for (item, r) in page.items.iter().zip(hydrated) { 78 match r { 79 Ok(item) => feed.push(item), 80 Err(e) => { 81 tracing::warn!(uri = %item.uri, error = %e, "skipping timeline item, hydration failed") 82 } 83 } 84 } 85 86 let items = stream::iter(feed.into_iter().map(Ok::<_, XrpcError>)); 87 Ok(json_stream::<TimelineItem<DefaultStr>, _>( 88 "feed", 89 items, 90 paged_tail(page.next.map(PageToken::encode_token)), 91 permit, 92 )) 93} 94 95async fn followed_dids(state: &AppState, viewer: &Did<DefaultStr>) -> Vec<Did<DefaultStr>> { 96 let key = EdgeKey::new(nsid_static(FOLLOW_BY_NSID), SubjectRef::Did(viewer.clone())); 97 let uris = state.edges.sources_for(&key); 98 // TODO: store follow actor information in edge 99 stream::iter(uris.into_iter().map(|uri| async move { 100 match fetch::<FollowRecord, Follow<DefaultStr>>(state, &uri).await { 101 Ok((_, follow)) => Some(follow.subject), 102 Err(e) => { 103 tracing::warn!(uri = %uri, error = %e, "skipping followed did, follow record unresolved"); 104 None 105 } 106 } 107 })) 108 .buffered(HYDRATE_CONCURRENCY) 109 .collect::<Vec<_>>() 110 .await 111 .into_iter() 112 .flatten() 113 .collect() 114} 115 116fn following_skeleton( 117 state: &AppState, 118 followed: &[Did<DefaultStr>], 119 cursor: PageCursor, 120 limit: PageLimit, 121) -> EdgePage { 122 let keys: Vec<EdgeKey> = followed 123 .iter() 124 .flat_map(|u| { 125 [REPO_NSID, STAR_BY_NSID, FOLLOW_BY_NSID] 126 .map(|kind| EdgeKey::new(nsid_static(kind), SubjectRef::Did(u.clone()))) 127 }) 128 .collect(); 129 state.edges.list_multi(&keys, cursor, limit, SortDir::Desc) 130} 131 132fn global_skeleton(state: &AppState, cursor: PageCursor, limit: PageLimit) -> EdgePage { 133 let keys = [REPO_NSID, STAR_NSID, FOLLOW_NSID] 134 .map(|kind| EdgeKey::new(nsid_static(kind), SubjectRef::Global)); 135 state.edges.list_multi(&keys, cursor, limit, SortDir::Desc) 136} 137 138async fn hydrate_item( 139 state: &AppState, 140 viewer: Option<&Did<DefaultStr>>, 141 item: &EdgeItem, 142) -> Result<TimelineItem<DefaultStr>, XrpcError> { 143 let collection = item 144 .uri 145 .collection() 146 .ok_or_else(|| XrpcError::InvalidParams("timeline uri missing collection".into()))?; 147 let (event, event_at) = match collection.as_ref() { 148 REPO_NSID => hydrate_repo_event(state, viewer, &item.uri).await?, 149 STAR_NSID => hydrate_star_event(state, viewer, &item.uri).await?, 150 FOLLOW_NSID => hydrate_follow_event(state, viewer, &item.uri).await?, 151 other => { 152 return Err(XrpcError::InvalidRecord(format!( 153 "unexpected timeline collection: {other}" 154 ))); 155 } 156 }; 157 Ok(TimelineItem::new().event(event).event_at(event_at).build()) 158} 159 160async fn hydrate_repo_event( 161 state: &AppState, 162 viewer: Option<&Did<DefaultStr>>, 163 uri: &AtUri<DefaultStr>, 164) -> Result<(TimelineItemEvent<DefaultStr>, Datetime), XrpcError> { 165 let (body, repo) = fetch::<RepoRecord, Repo<DefaultStr>>(state, uri).await?; 166 let repo_did = repo.repo_did.clone().ok_or_else(|| { 167 XrpcError::InvalidRecord("repo has no repo_did; cannot build repoViewBasic".into()) 168 })?; 169 let view = build_repo_view_basic(state, viewer, uri, &repo, repo_did).await?; 170 let source = match &repo.source { 171 Some(UriValue::Did(did)) => build_repo_view_from_did(state, viewer, did).await.ok(), 172 Some(UriValue::At(_aturi)) => None, // we drop the legacy support 173 _ => None, 174 }; 175 let event_at = repo.created_at.clone(); 176 let ev = RepoEvent::new() 177 .uri(body.uri.clone()) 178 .cid(body.cid.clone()) 179 .repo(view) 180 .source(source) 181 .build(); 182 Ok((TimelineItemEvent::RepoEvent(Box::new(ev)), event_at)) 183} 184 185async fn hydrate_star_event( 186 state: &AppState, 187 viewer: Option<&Did<DefaultStr>>, 188 uri: &AtUri<DefaultStr>, 189) -> Result<(TimelineItemEvent<DefaultStr>, Datetime), XrpcError> { 190 let (body, star) = fetch::<StarRecord, Star<DefaultStr>>(state, uri).await?; 191 let starrer = owner_did_from_aturi(uri) 192 .ok_or_else(|| XrpcError::InvalidParams("star uri authority must be a did".into()))?; 193 let actor = build_profile_basic(state, viewer, &starrer).await; 194 let repo_did = match &star.subject { 195 StarSubject::Repo(r) => r.did.clone(), 196 StarSubject::String(_) => { 197 return Err(XrpcError::InvalidRecord( 198 "timeline star subject is not a repo".into(), 199 )); 200 } 201 }; 202 let repo = build_repo_view_from_did(state, viewer, &repo_did).await?; 203 let event_at = star.created_at.clone(); 204 let ev = StarEvent::new() 205 .uri(body.uri.clone()) 206 .cid(body.cid.clone()) 207 .actor(actor) 208 .repo(repo) 209 .build(); 210 Ok((TimelineItemEvent::StarEvent(Box::new(ev)), event_at)) 211} 212 213async fn hydrate_follow_event( 214 state: &AppState, 215 viewer: Option<&Did<DefaultStr>>, 216 uri: &AtUri<DefaultStr>, 217) -> Result<(TimelineItemEvent<DefaultStr>, Datetime), XrpcError> { 218 let (body, follow) = fetch::<FollowRecord, Follow<DefaultStr>>(state, uri).await?; 219 let follower = owner_did_from_aturi(uri) 220 .ok_or_else(|| XrpcError::InvalidParams("follow uri authority must be a did".into()))?; 221 let actor = build_profile_basic(state, viewer, &follower).await; 222 let subject = build_profile_detailed(state, viewer, &follow.subject).await; 223 let event_at = follow.created_at.clone(); 224 let ev = FollowEvent::new() 225 .uri(body.uri.clone()) 226 .cid(body.cid.clone()) 227 .actor(actor) 228 .subject(subject) 229 .build(); 230 Ok((TimelineItemEvent::FollowEvent(Box::new(ev)), event_at)) 231} 232 233async fn build_repo_view_from_did( 234 state: &AppState, 235 viewer: Option<&Did<DefaultStr>>, 236 repo_did: &Did<DefaultStr>, 237) -> Result<RepoViewBasic<DefaultStr>, XrpcError> { 238 let ident = state 239 .resolver 240 .lookup_by_repo_did(repo_did) 241 .await 242 .ok_or(XrpcError::NotFound)?; 243 let uri = AtUri::<DefaultStr>::from_parts_owned( 244 ident.owner.as_str(), 245 RepoRecord::NSID, 246 ident.rkey.as_str(), 247 ) 248 .map_err(|e| XrpcError::InvalidRecord(format!("repo uri assembly: {e}")))?; 249 let (_, repo) = fetch::<RepoRecord, Repo<DefaultStr>>(state, &uri).await?; 250 build_repo_view_basic(state, viewer, &uri, &repo, repo_did.clone()).await 251} 252 253async fn build_repo_view_basic( 254 state: &AppState, 255 viewer: Option<&Did<DefaultStr>>, 256 repo_uri: &AtUri<DefaultStr>, 257 repo_record: &Repo<DefaultStr>, 258 repo_did: Did<DefaultStr>, 259) -> Result<RepoViewBasic<DefaultStr>, XrpcError> { 260 let owner_did = owner_did_from_aturi(repo_uri) 261 .ok_or_else(|| XrpcError::InvalidParams("repo uri authority must be a did".into()))?; 262 let owner = build_profile_basic(state, viewer, &owner_did).await; 263 let slug = repo_record.name.clone().unwrap_or_else(|| { 264 repo_uri 265 .rkey() 266 .map(|r| DefaultStr::from(r.as_ref())) 267 .unwrap_or_default() 268 }); 269 let star_key = EdgeKey::new(nsid_static(STAR_NSID), SubjectRef::Did(repo_did.clone())); 270 let star_count = state.edges.count(&star_key) as i64; 271 let viewer = viewer.map(|v| { 272 let mut viewer = repo::ViewerState::default(); 273 viewer.star = state.edges.viewer_source(&star_key, v.as_str()); 274 viewer 275 }); 276 Ok(RepoViewBasic::new() 277 .did(repo_did) 278 .owner(owner) 279 .slug(slug) 280 .created_at(repo_record.created_at.clone()) 281 .description(repo_record.description.clone()) 282 .star_count(star_count) 283 .viewer(viewer) 284 .build()) 285} 286 287async fn build_profile_basic( 288 state: &AppState, 289 viewer: Option<&Did<DefaultStr>>, 290 did: &Did<DefaultStr>, 291) -> ProfileViewBasic<DefaultStr> { 292 let handle = resolve_handle(state, did).await; 293 let avatar = None; // state.avatar.as_ref().and_then(|s| s.url(did)); 294 let viewer_state = build_actor_viewer_state(&state.edges, viewer, did); 295 ProfileViewBasic::new() 296 .did(did.clone()) 297 .handle(handle) 298 .maybe_avatar(avatar) 299 .maybe_viewer(viewer_state) 300 .build() 301} 302 303async fn build_profile_detailed( 304 state: &AppState, 305 viewer: Option<&Did<DefaultStr>>, 306 did: &Did<DefaultStr>, 307) -> ProfileViewDetailed<DefaultStr> { 308 let handle = resolve_handle(state, did).await; 309 let avatar = None; // state.avatar.as_ref().and_then(|s| s.url(did)); 310 let viewer_state = build_actor_viewer_state(&state.edges, viewer, did); 311 let followers = state.edges.count(&EdgeKey::new( 312 nsid_static(FOLLOW_NSID), 313 SubjectRef::Did(did.clone()), 314 )) as i64; 315 let follows = state.edges.count(&EdgeKey::new( 316 nsid_static(FOLLOW_BY_NSID), 317 SubjectRef::Did(did.clone()), 318 )) as i64; 319 ProfileViewDetailed::new() 320 .did(did.clone()) 321 .handle(handle) 322 .followers_count(followers) 323 .follows_count(follows) 324 .maybe_avatar(avatar) 325 .maybe_viewer(viewer_state) 326 .build() 327} 328 329fn build_actor_viewer_state( 330 edges: &EdgeStore, 331 viewer: Option<&Did<DefaultStr>>, 332 subject: &Did<DefaultStr>, 333) -> Option<ViewerState<DefaultStr>> { 334 let viewer = viewer?; 335 if viewer == subject { 336 return None; 337 } 338 let following = edges.viewer_source( 339 &EdgeKey::new(nsid_static(FOLLOW_NSID), SubjectRef::Did(subject.clone())), 340 viewer.as_str(), 341 ); 342 let followed_by = edges.viewer_source( 343 &EdgeKey::new(nsid_static(FOLLOW_NSID), SubjectRef::Did(viewer.clone())), 344 subject.as_str(), 345 ); 346 (following.is_some() || followed_by.is_some()).then(|| ViewerState { 347 following, 348 followed_by, 349 ..Default::default() 350 }) 351} 352 353async fn resolve_handle(state: &AppState, did: &Did<DefaultStr>) -> Handle<DefaultStr> { 354 state 355 .directory 356 .resolve_did_doc_owned(did) 357 .await 358 .ok() 359 .and_then(|doc| { 360 doc.handles() 361 .into_iter() 362 .next() 363 .map(|h| h.as_str().to_owned()) 364 }) 365 .and_then(|s| Handle::new_owned(s).ok()) 366 .unwrap_or_else(|| { 367 Handle::new_static("handle.invalid").expect("handle.invalid is a valid handle") 368 }) 369}