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