This repository has no description
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}