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