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_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}