From dfc5b4e3803947925d7f842ebcc9bf01974f3632 Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Wed, 16 Sep 2026 00:23:13 +0900 Subject: [PATCH] bobbin-xrpc: hydrated `feed.listComments` Signed-off-by: Seongmin Lee --- Cargo.lock | 1 + bobbin/crates/xrpc/Cargo.toml | 1 + bobbin/crates/xrpc/src/change_ids.rs | 85 +++++++++++++++++++++ bobbin/crates/xrpc/src/feed.rs | 98 ++++++++++++++++++++++++- bobbin/crates/xrpc/src/lib.rs | 6 +- bobbin/crates/xrpc/tests/aggregation.rs | 16 ++-- 6 files changed, 194 insertions(+), 13 deletions(-) create mode 100644 bobbin/crates/xrpc/src/change_ids.rs diff --git a/Cargo.lock b/Cargo.lock index 877864ca6..c8b5c38c4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1037,6 +1037,7 @@ dependencies = [ "jacquard-common", "jacquard-identity", "lexicons", + "quick_cache", "reqwest 0.13.1", "scc", "serde", diff --git a/bobbin/crates/xrpc/Cargo.toml b/bobbin/crates/xrpc/Cargo.toml index fafc5d695..061dd61a4 100644 --- a/bobbin/crates/xrpc/Cargo.toml +++ b/bobbin/crates/xrpc/Cargo.toml @@ -24,6 +24,7 @@ axum = { workspace = true } chrono = { workspace = true } futures = { workspace = true } http = { workspace = true } +quick_cache = { workspace = true } scc = { workspace = true } serde = { workspace = true } serde_json = { workspace = true } diff --git a/bobbin/crates/xrpc/src/change_ids.rs b/bobbin/crates/xrpc/src/change_ids.rs new file mode 100644 index 000000000..b2ead30fd --- /dev/null +++ b/bobbin/crates/xrpc/src/change_ids.rs @@ -0,0 +1,85 @@ +use bobbin_knot_proxy::MirrorProxy; +use futures::StreamExt; +use http::HeaderMap; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use jacquard_common::types::nsid::Nsid; +use lexicons::sh_tangled::git::temp2::list_commits::ListCommitsOutput; +use quick_cache::sync::Cache; + +const LIST_COMMITS_NSID: &str = "sh.tangled.git.temp2.listCommits"; +const CHANGE_ID_HEADER: &str = "change-id"; +const MAX_RESPONSE_BYTES: usize = 64 * 1024; +const DEFAULT_CAPACITY: usize = 8_192; + +pub(crate) struct ChangeIdCache(Cache>); + +impl Default for ChangeIdCache { + fn default() -> Self { + Self(Cache::new(DEFAULT_CAPACITY)) + } +} + +impl ChangeIdCache { + pub(crate) async fn get( + &self, + mirror: Option<&MirrorProxy>, + repo: &Did, + commit: &str, + ) -> Option { + if let Some(cached) = self.0.get(commit) { + return cached; + } + let change = match fetch_change_id(mirror?, repo.as_ref(), commit).await { + Ok(change) => change, + Err(error) => { + tracing::warn!(%repo, commit, %error, "looking up a commit change-id failed"); + return None; + } + }; + self.0.insert(DefaultStr::from(commit), change.clone()); + change + } +} + +async fn fetch_change_id( + mirror: &MirrorProxy, + repo: &str, + commit: &str, +) -> Result, String> { + let nsid = Nsid::new_static(LIST_COMMITS_NSID).expect("list commits NSID is valid"); + let response = mirror + .forward_raw( + &nsid, + &[("repo", repo), ("ranges", commit), ("limit", "1")], + HeaderMap::new(), + ) + .await + .map_err(|error| error.to_string())?; + if !response.status().is_success() { + let status = response.status(); + response.discard().await; + return Err(format!("mirror returned {status}")); + } + + let mut body = response.into_body_stream(); + let mut bytes = Vec::new(); + while let Some(chunk) = body.next().await { + let chunk = chunk.map_err(|error| error.to_string())?; + if bytes.len().saturating_add(chunk.len()) > MAX_RESPONSE_BYTES { + return Err("mirror list commits response is too large".to_owned()); + } + bytes.extend_from_slice(&chunk); + } + let output: ListCommitsOutput = + serde_json::from_slice(&bytes).map_err(|error| error.to_string())?; + let found = output + .commits + .first() + .ok_or_else(|| "mirror returned no commit".to_owned())?; + Ok(found + .extra_headers + .iter() + .find(|header| header.key == CHANGE_ID_HEADER) + .map(|header| header.value.clone())) +} diff --git a/bobbin/crates/xrpc/src/feed.rs b/bobbin/crates/xrpc/src/feed.rs index 29ed942ad..76e998dda 100644 --- a/bobbin/crates/xrpc/src/feed.rs +++ b/bobbin/crates/xrpc/src/feed.rs @@ -1,12 +1,18 @@ use axum::extract::State; use axum::response::Response; -use futures::stream::{self, StreamExt}; +use futures::future::OptionFuture; +use futures::stream::{self, StreamExt, TryStreamExt}; use lexicons::sh_tangled::feed::list_comments::ListComments; use serde::Deserialize; -use bobbin_edge_index::{EdgeItem, EdgePage, EdgeStore, PageCursor, PageLimit, PageToken, SortDir}; +use bobbin_edge_index::{ + EdgeItem, EdgePage, EdgeStore, PageCursor, PageLimit, PageOffset, PageToken, SortDir, +}; use bobbin_types::ids::{EdgeKey, SubjectRef, nsid_static, owner_did_from_aturi}; use bobbin_types::sh_tangled::actor::{ProfileViewBasic, ProfileViewDetailed, ViewerState}; +use bobbin_types::sh_tangled::embed::commit::View as CommitEmbedView; +use bobbin_types::sh_tangled::feed::CommentView; +use bobbin_types::sh_tangled::feed::comment::{Comment, CommentRecord}; use bobbin_types::sh_tangled::feed::get_timeline::{ FollowEvent, RepoEvent, StarEvent, TimelineItem, TimelineItemEvent, }; @@ -19,7 +25,9 @@ use jacquard_common::xrpc::XrpcResp; use jacquard_common::{DefaultStr, IntoStatic as _}; use crate::{ - AppState, XrpcError, XrpcQuery, fetch, json_stream, paged_tail, parse_cursor, parse_limit, + AppState, HasSubject, HitProvenance, RecordView, SubjectQuery, XrpcError, XrpcQuery, fetch, + hydrate_record_stream, json_stream, paged_tail, parse_cursor, parse_limit, parse_start, + parse_subject, }; const REPO_NSID: &str = "sh.tangled.repo"; @@ -34,7 +42,89 @@ pub(crate) async fn list_comments( State(state): State, XrpcQuery(q): XrpcQuery, ) -> Result { - todo!() + let subject = parse_subject(&SubjectQuery::Uri(q.subject), CommentRecord::SHAPE)?; + let offset = q + .offset + .map(|offset| { + u64::try_from(offset) + .map(PageOffset::new) + .map_err(|_| XrpcError::InvalidParams("offset must not be negative".into())) + }) + .transpose()?; + let start = parse_start(q.cursor.as_deref(), offset)?; + let limit = parse_limit( + q.limit + .map(|limit| { + u32::try_from(limit).map_err(|_| { + XrpcError::InvalidParams("limit must be between 1 and 1000".into()) + }) + }) + .transpose()?, + )?; + let dir = match q.order.as_deref() { + Some("asc") => SortDir::Asc, + _ => SortDir::Desc, + }; + + let nsid = nsid_static(CommentRecord::NSID); + let page = state + .edges + .list(&EdgeKey::new(nsid.clone(), subject), start, limit, dir); + let permit = state.heavy_permit()?; + let owned = state.clone(); + let items = hydrate_record_stream::>( + &state, + nsid, + page.items, + HitProvenance::Indexed, + ) + .map_ok(move |view| { + let state = owned.clone(); + async move { Ok(comment_view(&state, view).await) } + }) + .try_buffered(HYDRATE_CONCURRENCY) + .try_filter_map(|view| async move { Ok(view) }); + + Ok(json_stream::, _>( + "items", + items, + paged_tail(page.next.map(PageToken::encode_token), page.total), + permit, + )) +} + +async fn comment_view( + state: &AppState, + view: RecordView>, +) -> Option> { + let RecordView { uri, cid, value } = view; + let embed = OptionFuture::from(value.embed.map(|commit| async { + let change = state + .change_ids + .get( + state.mirror_v2.as_deref(), + &commit.repo, + commit.commit.oid.as_ref(), + ) + .await; + CommitEmbedView { + change, + commit: commit.commit.oid, + repo: commit.repo, + extra_data: None, + } + })) + .await; + Some(CommentView { + body: value.body, + cid: cid?, + created_at: value.created_at, + embed, + reply_to: value.reply_to, + subject: value.subject, + uri, + extra_data: None, + }) } #[derive(Debug, Deserialize)] diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index 37672bb2a..6e576b7a0 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -112,6 +112,7 @@ use tracing::{Level, Span}; mod actor; mod backpressure; +mod change_ids; mod client_address; mod codesearch; mod commit_stats; @@ -178,6 +179,7 @@ pub struct AppState { enrich_router: Arc>, trending_cache: Arc>>, commit_stats_cache: Arc, + change_ids: Arc, } impl AppState { @@ -228,6 +230,7 @@ impl AppState { enrich_router: Arc::new(std::sync::OnceLock::new()), trending_cache: Arc::new(tokio::sync::Mutex::new(None)), commit_stats_cache: Arc::new(commit_stats::CommitStatsCache::default()), + change_ids: Arc::new(change_ids::ChangeIdCache::default()), } } @@ -364,7 +367,7 @@ pub fn router(state: AppState) -> Router { .route("/xrpc/sh.tangled.repo.countPulls", get(count_pulls)) .route( "/xrpc/sh.tangled.feed.listComments", - get(list_feed_comments), + get(feed::list_comments), ) .route( "/xrpc/sh.tangled.feed.countComments", @@ -2437,6 +2440,7 @@ async fn count_pulls( count_typed_for::(&state, q).map(Json) } +#[allow(dead_code)] async fn list_feed_comments( State(state): State, XrpcQuery(q): XrpcQuery>, diff --git a/bobbin/crates/xrpc/tests/aggregation.rs b/bobbin/crates/xrpc/tests/aggregation.rs index 06ec7ae69..d46b72323 100644 --- a/bobbin/crates/xrpc/tests/aggregation.rs +++ b/bobbin/crates/xrpc/tests/aggregation.rs @@ -1609,11 +1609,10 @@ async fn list_feed_comments_hydrates_end_to_end() { assert_eq!(status, StatusCode::OK); let items = body["items"].as_array().unwrap(); assert_eq!(items.len(), 1); - assert_eq!(items[0]["value"]["body"]["text"], json!("thoughts")); - assert_eq!( - items[0]["value"]["subject"]["uri"], - json!(issue_uri.as_ref()) - ); + assert_eq!(items[0]["body"]["text"], json!("thoughts")); + assert_eq!(items[0]["subject"]["uri"], json!(issue_uri.as_ref())); + assert_eq!(items[0]["createdAt"], json!("2026-05-01T00:00:00Z")); + assert!(items[0]["cid"].is_string()); } #[tokio::test] @@ -1987,10 +1986,11 @@ async fn feed_comment_endpoints_reject_bare_did_or_wrong_collection() { let (status, body) = json_response(resp).await; assert_eq!(status, StatusCode::BAD_REQUEST, "{endpoint} input={input}"); let msg = body["message"].as_str().unwrap_or_default(); + let names_collections = msg.contains("sh.tangled.repo.issue") + && msg.contains("sh.tangled.repo.pull") + && msg.contains("sh.tangled.string"); assert!( - msg.contains("sh.tangled.repo.issue") - && msg.contains("sh.tangled.repo.pull") - && msg.contains("sh.tangled.string"), + names_collections || msg.contains("at:////"), "{endpoint} input={input}: {msg}", ); } -- 2.51.2