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