Monorepo for Tangled tangled.org
1

Configure Feed

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

Labels

None yet.

Participants 1
AT URI
at://did:plc:xasnlahkri4ewmbuzly2rlc5/sh.tangled.repo.pull/3mrmrvzsutq22
+375 -6
Diff #1
+367
bobbin/crates/xrpc/src/feed.rs
··· 1 + use axum::extract::State; 2 + use axum::response::Response; 3 + use futures::stream::{self, StreamExt}; 4 + use serde::Deserialize; 5 + 6 + use bobbin_edge_index::{EdgeItem, EdgePage, EdgeStore, PageCursor, PageLimit, PageToken, SortDir}; 7 + use bobbin_types::ids::{nsid_static, owner_did_from_aturi, EdgeKey, SubjectRef}; 8 + use bobbin_types::sh_tangled::actor::{ProfileViewBasic, ProfileViewDetailed, ViewerState}; 9 + use bobbin_types::sh_tangled::feed::get_timeline::{ 10 + FollowEvent, RepoEvent, StarEvent, TimelineItem, TimelineItemEvent, 11 + }; 12 + use bobbin_types::sh_tangled::feed::star::{Star, StarRecord, StarSubject}; 13 + use bobbin_types::sh_tangled::graph::follow::{Follow, FollowRecord}; 14 + use bobbin_types::sh_tangled::repo::{self, Repo, RepoRecord, RepoViewBasic}; 15 + use jacquard_common::types::string::{AtUri, Datetime, Did, Handle, UriValue}; 16 + use jacquard_common::xrpc::XrpcResp; 17 + use jacquard_common::DefaultStr; 18 + use jacquard_identity::resolver::IdentityResolver; 19 + 20 + use crate::{ 21 + fetch, json_stream, paged_tail, parse_cursor, parse_limit, AppState, XrpcError, XrpcQuery, 22 + }; 23 + 24 + const REPO_NSID: &str = "sh.tangled.repo"; 25 + const STAR_NSID: &str = "sh.tangled.feed.star"; 26 + const FOLLOW_NSID: &str = "sh.tangled.graph.follow"; 27 + const STAR_BY_NSID: &str = "sh.tangled.feed.star.by"; 28 + const FOLLOW_BY_NSID: &str = "sh.tangled.graph.follow.by"; 29 + 30 + const HYDRATE_CONCURRENCY: usize = 8; 31 + 32 + #[derive(Debug, Deserialize)] 33 + #[serde(rename_all = "camelCase")] 34 + pub(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 + 42 + pub(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 + 93 + async 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 + 114 + fn 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 + 130 + fn 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 + 136 + async 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 + 158 + async 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 + 183 + async 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 + 211 + async 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 + 231 + async 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 + 251 + async 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 + 285 + async 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 + 301 + async 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 + 327 + fn 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 + 351 + async 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 + }
+8 -6
bobbin/crates/xrpc/src/lib.rs
··· 99 99 100 100 mod backpressure; 101 101 mod enrich; 102 + mod feed; 102 103 mod filter; 103 104 mod recordpath; 104 105 ··· 168 169 .clone() 169 170 } 170 171 171 - fn heavy_permit(&self) -> Result<Option<HeavyPermit>, XrpcError> { 172 + pub(crate) fn heavy_permit(&self) -> Result<Option<HeavyPermit>, XrpcError> { 172 173 self.limiter.as_ref().map(|l| l.try_enter()).transpose() 173 174 } 174 175 } ··· 260 261 ) 261 262 .route("/xrpc/sh.tangled.graph.listVouches", get(list_vouches)) 262 263 .route("/xrpc/sh.tangled.graph.countVouches", get(count_vouches)) 264 + .route("/xrpc/sh.tangled.feed.getTimeline", get(feed::get_timeline)) 263 265 .route("/xrpc/sh.tangled.feed.listStarsBy", get(list_stars_by)) 264 266 .route("/xrpc/sh.tangled.feed.countStarsBy", get(count_stars_by)) 265 267 .route( ··· 1086 1088 }) 1087 1089 } 1088 1090 1089 - fn parse_cursor(raw: Option<&str>) -> Result<PageCursor, XrpcError> { 1091 + pub(crate) fn parse_cursor(raw: Option<&str>) -> Result<PageCursor, XrpcError> { 1090 1092 PageCursor::from_token(raw) 1091 1093 .map_err(|e: CursorParseError| XrpcError::InvalidParams(format!("cursor: {e}"))) 1092 1094 } 1093 1095 1094 - fn parse_limit(raw: Option<u32>) -> Result<PageLimit, XrpcError> { 1096 + pub(crate) fn parse_limit(raw: Option<u32>) -> Result<PageLimit, XrpcError> { 1095 1097 PageLimit::new(raw.unwrap_or(DEFAULT_LIMIT)) 1096 1098 .map_err(|e| XrpcError::InvalidParams(format!("limit: {e}"))) 1097 1099 } ··· 1249 1251 Ok((body, value)) 1250 1252 } 1251 1253 1252 - async fn fetch<R, V>( 1254 + pub(crate) async fn fetch<R, V>( 1253 1255 state: &AppState, 1254 1256 uri: &AtUri<DefaultStr>, 1255 1257 ) -> Result<(Arc<RecordBody>, V), XrpcError> ··· 1633 1635 permit: Option<HeavyPermit>, 1634 1636 } 1635 1637 1636 - fn paged_tail(cursor: Option<String>) -> Vec<u8> { 1638 + pub(crate) fn paged_tail(cursor: Option<String>) -> Vec<u8> { 1637 1639 let encoded = serde_json::to_string(&cursor).unwrap_or_else(|_| "null".to_owned()); 1638 1640 format!("],\"cursor\":{encoded}}}").into_bytes() 1639 1641 } ··· 1642 1644 b"]}".to_vec() 1643 1645 } 1644 1646 1645 - fn json_stream<V, S>( 1647 + pub(crate) fn json_stream<V, S>( 1646 1648 array_key: &'static str, 1647 1649 items: S, 1648 1650 tail: Vec<u8>,

History

2 rounds 0 comments
Sign up or Login to add to the discussion
1 commit
Expand
bobbin-xrpc: sh.tangled.feed.getTimeline
Checking mergeability…
Expand 0 comments
1 commit
Expand
bobbin-xrpc: sh.tangled.feed.getTimeline
Expand 0 comments