Signed-off-by: Seongmin Lee git@boltless.me
+375
-6
Diff
Round #0
+367
bobbin/crates/xrpc/src/feed.rs
+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
+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
boltless.me
submitted
#1
1 commit
Expand
Collapse
bobbin-xrpc:
sh.tangled.feed.getTimeline
Signed-off-by: Seongmin Lee <git@boltless.me>
Expand 0 comments
boltless.me
submitted
#0
1 commit
Expand
Collapse
bobbin-xrpc:
sh.tangled.feed.getTimeline
Signed-off-by: Seongmin Lee <git@boltless.me>