diff --git a/Cargo.lock b/Cargo.lock index a20c546fb..1911f4e8d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -796,6 +796,7 @@ version = "0.0.1" dependencies = [ "bobbin-runtime", "bobbin-types", + "divan", "either", "ftree", "itertools 0.14.0", diff --git a/api/tangled/actorgetProfileActivity.go b/api/tangled/actorgetProfileActivity.go new file mode 100644 index 000000000..5a723cb18 --- /dev/null +++ b/api/tangled/actorgetProfileActivity.go @@ -0,0 +1,82 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package tangled + +// schema: sh.tangled.actor.getProfileActivity + +import ( + "context" + + "github.com/bluesky-social/indigo/lex/util" +) + +const ( + ActorGetProfileActivityNSID = "sh.tangled.actor.getProfileActivity" +) + +// ActorGetProfileActivity_ActivityMonth is a "activityMonth" in the sh.tangled.actor.getProfileActivity schema. +type ActorGetProfileActivity_ActivityMonth struct { + Commits int64 `json:"commits" cborgen:"commits"` + Issues []*ActorGetProfileActivity_IssueActivity `json:"issues" cborgen:"issues"` + Pulls []*ActorGetProfileActivity_PullActivity `json:"pulls" cborgen:"pulls"` + Repos []*ActorGetProfileActivity_RepoActivity `json:"repos" cborgen:"repos"` + // start: Start of this UTC calendar month. + Start string `json:"start" cborgen:"start"` +} + +// ActorGetProfileActivity_IssueActivity is a "issueActivity" in the sh.tangled.actor.getProfileActivity schema. +type ActorGetProfileActivity_IssueActivity struct { + Repo *ActorGetProfileActivity_RepoRef `json:"repo" cborgen:"repo"` + State string `json:"state" cborgen:"state"` + Title string `json:"title" cborgen:"title"` + Uri string `json:"uri" cborgen:"uri"` +} + +// ActorGetProfileActivity_Output is the output of a sh.tangled.actor.getProfileActivity call. +type ActorGetProfileActivity_Output struct { + Months []*ActorGetProfileActivity_ActivityMonth `json:"months" cborgen:"months"` + // truncated: True when the result reached the server activity-item safety limit. + Truncated bool `json:"truncated" cborgen:"truncated"` +} + +// ActorGetProfileActivity_PullActivity is a "pullActivity" in the sh.tangled.actor.getProfileActivity schema. +type ActorGetProfileActivity_PullActivity struct { + Repo *ActorGetProfileActivity_RepoRef `json:"repo" cborgen:"repo"` + State string `json:"state" cborgen:"state"` + Title string `json:"title" cborgen:"title"` + Uri string `json:"uri" cborgen:"uri"` +} + +// ActorGetProfileActivity_RepoActivity is a "repoActivity" in the sh.tangled.actor.getProfileActivity schema. +type ActorGetProfileActivity_RepoActivity struct { + Kind string `json:"kind" cborgen:"kind"` + Repo *ActorGetProfileActivity_RepoRef `json:"repo" cborgen:"repo"` +} + +// ActorGetProfileActivity_RepoRef is a "repoRef" in the sh.tangled.actor.getProfileActivity schema. +type ActorGetProfileActivity_RepoRef struct { + Did *string `json:"did,omitempty" cborgen:"did,omitempty"` + Name string `json:"name" cborgen:"name"` + OwnerDid string `json:"ownerDid" cborgen:"ownerDid"` + OwnerHandle string `json:"ownerHandle" cborgen:"ownerHandle"` + Uri string `json:"uri" cborgen:"uri"` +} + +// ActorGetProfileActivity calls the XRPC method "sh.tangled.actor.getProfileActivity". +// +// actor: Actor whose activity to return. +// months: Number of UTC calendar months to include, including the current month. +func ActorGetProfileActivity(ctx context.Context, c util.LexClient, actor string, months int64) (*ActorGetProfileActivity_Output, error) { + var out ActorGetProfileActivity_Output + + params := map[string]interface{}{} + params["actor"] = actor + if months != 0 { + params["months"] = months + } + if err := c.LexDo(ctx, util.Query, "", "sh.tangled.actor.getProfileActivity", params, nil, &out); err != nil { + return nil, err + } + + return &out, nil +} diff --git a/bobbin/crates/edge-index/Cargo.toml b/bobbin/crates/edge-index/Cargo.toml index 212e6f458..1d0bb76e5 100644 --- a/bobbin/crates/edge-index/Cargo.toml +++ b/bobbin/crates/edge-index/Cargo.toml @@ -21,4 +21,9 @@ thiserror = { workspace = true } tokio = { workspace = true } [dev-dependencies] +divan = { workspace = true } tokio = { workspace = true, features = ["macros", "rt", "test-util"] } + +[[bench]] +name = "profile_activity" +harness = false diff --git a/bobbin/crates/edge-index/benches/profile_activity.rs b/bobbin/crates/edge-index/benches/profile_activity.rs new file mode 100644 index 000000000..49949c81b --- /dev/null +++ b/bobbin/crates/edge-index/benches/profile_activity.rs @@ -0,0 +1,133 @@ +use bobbin_edge_index::{EdgeStore, PageCursor, PageLimit, SortDir}; +use bobbin_runtime::RuntimeHasher; +use bobbin_types::edges::Edge; +use bobbin_types::ids::{EdgeKey, SubjectRef, nsid_static}; +use divan::Bencher; +use divan::counter::ItemsCount; +use jacquard_common::{DefaultStr, types::did::Did, types::string::AtUri}; + +#[global_allocator] +static ALLOC: divan::AllocProfiler = divan::AllocProfiler::system(); + +const EDGE_COUNT: usize = 100_000; +const GLOBAL_ACTOR_COUNT: usize = 250_000; +const PAGE_SIZE: u32 = 1_000; +const RECENT_COUNT: usize = 1_000; +const KINDS: [(&str, &str); 3] = [ + ("sh.tangled.repo", "sh.tangled.repo"), + ("sh.tangled.repo.issue.by", "sh.tangled.repo.issue"), + ("sh.tangled.repo.pull.by", "sh.tangled.repo.pull"), +]; + +fn did(value: &str) -> Did { + Did::new_owned(value).expect("valid did") +} + +fn populated_store() -> (EdgeStore, Vec) { + let store = EdgeStore::new(RuntimeHasher::default()); + let actor = did("did:plc:profileactor"); + let keys = KINDS + .iter() + .map(|(kind, _)| EdgeKey::new(nsid_static(kind), SubjectRef::Did(actor.clone()))) + .collect(); + for index in 0..EDGE_COUNT { + let (kind, collection) = KINDS[index % KINDS.len()]; + let source = AtUri::new_owned(format!( + "at://did:plc:owner{}/{collection}/item{index}", + index % 1_000, + )) + .expect("valid at-uri"); + store.add(Edge { + kind: nsid_static(kind), + subject: SubjectRef::Did(actor.clone()), + source, + sort_micros: index as u64, + }); + } + (store, keys) +} + +fn globally_populated_store() -> (EdgeStore, Vec) { + let store = EdgeStore::new(RuntimeHasher::default()); + for index in 0..GLOBAL_ACTOR_COUNT { + let actor = did(&format!("did:plc:actor{index}")); + store.add(Edge { + kind: nsid_static("sh.tangled.repo"), + subject: SubjectRef::Did(actor.clone()), + source: AtUri::new_owned(format!("at://{}/sh.tangled.repo/repo", actor.as_ref(),)) + .expect("valid at-uri"), + sort_micros: index as u64, + }); + } + let target = did("did:plc:targetactor"); + let keys = KINDS + .iter() + .map(|(kind, collection)| { + store.add(Edge { + kind: nsid_static(kind), + subject: SubjectRef::Did(target.clone()), + source: AtUri::new_owned(format!("at://did:plc:targetactor/{collection}/item",)) + .expect("valid at-uri"), + sort_micros: GLOBAL_ACTOR_COUNT as u64, + }); + EdgeKey::new(nsid_static(kind), SubjectRef::Did(target.clone())) + }) + .collect(); + (store, keys) +} + +fn collect_recent(store: &EdgeStore, keys: &[EdgeKey], cutoff: u64) -> usize { + let limit = PageLimit::new(PAGE_SIZE).expect("valid page size"); + let mut cursor = PageCursor::Start; + let mut count = 0; + loop { + let page = store.list_multi(keys, cursor, limit, SortDir::Desc); + let mut reached_cutoff = false; + for item in page.items { + if item.sort_micros < cutoff { + reached_cutoff = true; + break; + } + count += 1; + } + if reached_cutoff { + break; + } + let Some(next) = page.next else { + break; + }; + cursor = PageCursor::After(next); + } + count +} + +fn main() { + divan::main(); +} + +#[divan::bench] +fn profile_activity_skeleton_page(bencher: Bencher) { + let (store, keys) = populated_store(); + let limit = PageLimit::new(PAGE_SIZE).expect("valid page size"); + bencher + .counter(ItemsCount::new(PAGE_SIZE)) + .bench_local(|| store.list_multi(&keys, PageCursor::Start, limit, SortDir::Desc)); +} + +#[divan::bench] +fn profile_activity_skeleton_window(bencher: Bencher) { + let (store, keys) = populated_store(); + let cutoff = (EDGE_COUNT - RECENT_COUNT) as u64; + bencher + .counter(ItemsCount::new(RECENT_COUNT)) + .bench_local(|| collect_recent(&store, &keys, cutoff)); +} + +#[divan::bench] +fn profile_activity_among_users(bencher: Bencher) { + let (store, keys) = globally_populated_store(); + let limit = PageLimit::new(PAGE_SIZE).expect("valid page size"); + bencher + .counter(ItemsCount::new(keys.len())) + .bench_local(|| store.list_multi(&keys, PageCursor::Start, limit, SortDir::Desc)); +} diff --git a/bobbin/crates/xrpc/src/commit_stats.rs b/bobbin/crates/xrpc/src/commit_stats.rs new file mode 100644 index 000000000..604e6a1de --- /dev/null +++ b/bobbin/crates/xrpc/src/commit_stats.rs @@ -0,0 +1,667 @@ +use std::collections::{BTreeMap, HashMap}; +use std::sync::{Arc, Mutex}; +use std::time::{Duration, Instant}; + +use bobbin_knot_proxy::MirrorProxy; +use chrono::{DateTime, Datelike, Utc}; +use futures::StreamExt; +use http::HeaderMap; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use jacquard_common::types::nsid::Nsid; +use serde::Deserialize; +use tokio::sync::{Notify, Semaphore}; + +const COMMIT_STATS_NSID: &str = "sh.tangled.git.temp2.getCommitStats"; +const FETCH_MONTHS: u32 = 24; +const MAX_RESPONSE_BYTES: usize = 64 * 1024; +const DEFAULT_CAPACITY: usize = 4_096; +const DEFAULT_FRESHNESS: Duration = Duration::from_secs(30); +const FAILURE_RETRY: Duration = Duration::from_secs(5); +const MAX_CONCURRENT_REFRESHES: usize = 16; + +#[derive(Clone, Debug, Deserialize)] +struct CommitMonth { + start: String, + commits: i64, +} + +#[derive(Debug, Deserialize)] +struct CommitStats { + months: Vec, +} + +#[derive(Clone)] +struct Entry { + counts: BTreeMap, + refresh_after: Option, + refreshing: bool, + loaded: bool, + notify: Arc, + access: u64, +} + +struct CacheState { + entries: HashMap, + access: u64, +} + +pub(crate) struct CommitStatsCache { + state: Mutex, + capacity: usize, + freshness: Duration, + refreshes: Arc, + capacity_notify: Arc, +} + +impl Default for CommitStatsCache { + fn default() -> Self { + Self::new(DEFAULT_CAPACITY, DEFAULT_FRESHNESS) + } +} + +impl CommitStatsCache { + fn new(capacity: usize, freshness: Duration) -> Self { + Self { + state: Mutex::new(CacheState { + entries: HashMap::new(), + access: 0, + }), + capacity: capacity.max(1), + freshness, + refreshes: Arc::new(Semaphore::new(MAX_CONCURRENT_REFRESHES)), + capacity_notify: Arc::new(Notify::new()), + } + } + + pub(crate) async fn get( + self: &Arc, + mirror: Option>, + actor: &Did, + ) -> BTreeMap { + let Some(mirror) = mirror else { + return BTreeMap::new(); + }; + let key = actor.to_string(); + + loop { + let (counts, loaded, should_refresh, waiter) = { + let mut state = self + .state + .lock() + .expect("commit stats cache mutex poisoned"); + state.access = state.access.wrapping_add(1); + let access = state.access; + let capacity_waiter = if !state.entries.contains_key(&key) { + if evict_one(&mut state.entries, self.capacity) { + state.entries.insert( + key.clone(), + Entry { + counts: BTreeMap::new(), + refresh_after: None, + refreshing: false, + loaded: false, + notify: Arc::new(Notify::new()), + access, + }, + ); + None + } else { + Some(Arc::clone(&self.capacity_notify).notified_owned()) + } + } else { + None + }; + if let Some(waiter) = capacity_waiter { + (BTreeMap::new(), false, false, Some(waiter)) + } else { + let entry = state.entries.get_mut(&key).expect("entry inserted above"); + entry.access = access; + if !entry.loaded && entry.refreshing { + ( + BTreeMap::new(), + false, + false, + Some(Arc::clone(&entry.notify).notified_owned()), + ) + } else { + let fresh = entry + .refresh_after + .is_some_and(|refresh_after| Instant::now() < refresh_after); + let should_refresh = !fresh && !entry.refreshing; + if should_refresh { + entry.refreshing = true; + } + (entry.counts.clone(), entry.loaded, should_refresh, None) + } + } + }; + + if loaded { + if should_refresh { + match Arc::clone(&self.refreshes).try_acquire_owned() { + Ok(permit) => { + let cache = Arc::clone(self); + let mirror = Arc::clone(&mirror); + let key = key.clone(); + let mut guard = RefreshGuard::new(Arc::clone(&cache), key.clone()); + tokio::spawn(async move { + let result = fetch_commit_stats(&mirror, &key).await; + drop(permit); + match result { + Ok(counts) => cache.finish(key, Some(counts)), + Err(error) => { + tracing::warn!(actor = %key, %error, "refreshing commit stats failed"); + cache.finish(key, None); + } + } + guard.disarm(); + }); + } + Err(_) => self.finish(key.clone(), None), + } + } + return counts; + } + + if let Some(waiter) = waiter { + waiter.await; + continue; + } + + let mut guard = RefreshGuard::new(Arc::clone(self), key.clone()); + let Ok(permit) = Arc::clone(&self.refreshes).acquire_owned().await else { + self.finish(key, None); + guard.disarm(); + return BTreeMap::new(); + }; + let result = fetch_commit_stats(&mirror, &key).await; + drop(permit); + return match result { + Ok(counts) => { + self.finish(key, Some(counts.clone())); + guard.disarm(); + counts + } + Err(error) => { + tracing::warn!(actor = %key, %error, "loading commit stats failed"); + self.finish(key, None); + guard.disarm(); + BTreeMap::new() + } + }; + } + } + + fn finish(&self, key: String, counts: Option>) { + let mut state = self + .state + .lock() + .expect("commit stats cache mutex poisoned"); + state.access = state.access.wrapping_add(1); + let access = state.access; + let Some(entry) = state.entries.get_mut(&key) else { + return; + }; + entry.refreshing = false; + entry.loaded = true; + entry.access = access; + match counts { + Some(counts) => { + entry.counts = counts; + entry.refresh_after = Some(Instant::now() + self.freshness); + } + None => entry.refresh_after = Some(Instant::now() + FAILURE_RETRY), + } + let notify = Arc::clone(&entry.notify); + drop(state); + notify.notify_waiters(); + self.capacity_notify.notify_waiters(); + } + + fn cancel_refresh(&self, key: &str) { + let mut state = self + .state + .lock() + .expect("commit stats cache mutex poisoned"); + let Some(entry) = state.entries.get_mut(key) else { + return; + }; + entry.refreshing = false; + let notify = Arc::clone(&entry.notify); + drop(state); + notify.notify_waiters(); + self.capacity_notify.notify_waiters(); + } + + #[cfg(test)] + fn seed(&self, actor: &str, counts: BTreeMap, refresh_after: Option) { + self.state + .lock() + .expect("commit stats cache mutex poisoned") + .entries + .insert( + actor.to_owned(), + Entry { + counts, + refresh_after, + refreshing: false, + loaded: true, + notify: Arc::new(Notify::new()), + access: 0, + }, + ); + } +} + +struct RefreshGuard { + cache: Arc, + key: String, + active: bool, +} + +impl RefreshGuard { + fn new(cache: Arc, key: String) -> Self { + Self { + cache, + key, + active: true, + } + } + + fn disarm(&mut self) { + self.active = false; + } +} + +impl Drop for RefreshGuard { + fn drop(&mut self) { + if self.active { + self.cache.cancel_refresh(&self.key); + } + } +} + +fn evict_one(entries: &mut HashMap, capacity: usize) -> bool { + if entries.len() < capacity { + return true; + } + let Some(key) = entries + .iter() + .filter(|(_, entry)| !entry.refreshing) + .min_by_key(|(_, entry)| entry.access) + .map(|(key, _)| key.clone()) + else { + return false; + }; + entries.remove(&key); + true +} + +async fn fetch_commit_stats( + mirror: &MirrorProxy, + actor: &str, +) -> Result, String> { + let nsid = Nsid::new_static(COMMIT_STATS_NSID).expect("commit stats NSID is valid"); + let months = FETCH_MONTHS.to_string(); + let response = mirror + .forward_raw( + &nsid, + &[("actor", actor), ("months", months.as_str())], + 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 commit stats response is too large".to_owned()); + } + bytes.extend_from_slice(&chunk); + } + let output: CommitStats = serde_json::from_slice(&bytes).map_err(|error| error.to_string())?; + if output.months.len() > FETCH_MONTHS as usize { + return Err("mirror returned too many commit stat months".to_owned()); + } + output + .months + .into_iter() + .try_fold(BTreeMap::new(), |mut counts, month| -> Result<_, String> { + let start = DateTime::parse_from_rfc3339(&month.start) + .map_err(|error| format!("invalid month start: {error}"))? + .with_timezone(&Utc); + if month.commits < 0 { + return Err("mirror returned a negative commit count".to_owned()); + } + counts.insert(month_number(start), month.commits); + Ok(counts) + }) +} + +pub(crate) fn month_number(value: DateTime) -> i32 { + value.year() * 12 + value.month0() as i32 +} + +#[cfg(test)] +mod tests { + use super::*; + use bobbin_runtime::{RuntimeHasher, SystemClock}; + use serde_json::json; + use wiremock::matchers::{method, path, query_param}; + use wiremock::{Mock, MockServer, ResponseTemplate}; + + fn actor() -> Did { + Did::new_static("did:plc:pushcacheactor").unwrap() + } + + fn mirror(server: &MockServer) -> Arc { + Arc::new( + MirrorProxy::new( + &url::Url::parse(&server.uri()).unwrap(), + Arc::new(SystemClock::new()), + RuntimeHasher::default(), + ) + .unwrap(), + ) + } + + #[tokio::test] + async fn a_cold_read_waits_for_the_first_value() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/sh.tangled.git.temp2.getCommitStats")) + .and(query_param("actor", actor().as_ref())) + .and(query_param("months", "24")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + "months": [{"start": "2026-09-01T00:00:00Z", "commits": 7}] + }))) + .expect(1) + .mount(&server) + .await; + + let cache = Arc::new(CommitStatsCache::new(4, Duration::from_secs(60))); + let counts = cache.get(Some(mirror(&server)), &actor()).await; + assert_eq!(counts.get(&(2026 * 12 + 8)), Some(&7)); + assert_eq!( + cache + .get(Some(mirror(&server)), &actor()) + .await + .get(&(2026 * 12 + 8)), + Some(&7), + ); + server.verify().await; + } + + #[tokio::test] + async fn rejects_more_months_than_requested() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/sh.tangled.git.temp2.getCommitStats")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + "months": (0..=FETCH_MONTHS) + .map(|month| json!({ + "start": format!("2024-{:02}-01T00:00:00Z", month % 12 + 1), + "commits": 1, + })) + .collect::>() + }))) + .mount(&server) + .await; + + let error = fetch_commit_stats(&mirror(&server), actor().as_ref()) + .await + .unwrap_err(); + assert!(error.contains("too many commit stat months")); + } + + #[tokio::test] + async fn concurrent_cold_reads_share_one_request() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/sh.tangled.git.temp2.getCommitStats")) + .respond_with( + ResponseTemplate::new(200) + .set_delay(Duration::from_millis(25)) + .set_body_json(json!({ + "months": [{"start": "2026-09-01T00:00:00Z", "commits": 7}] + })), + ) + .expect(1) + .mount(&server) + .await; + let cache = Arc::new(CommitStatsCache::new(4, Duration::from_secs(60))); + let actor = actor(); + let first = cache.get(Some(mirror(&server)), &actor); + let second = cache.get(Some(mirror(&server)), &actor); + + let (first, second) = tokio::join!(first, second); + assert_eq!(first.get(&(2026 * 12 + 8)), Some(&7)); + assert_eq!(second.get(&(2026 * 12 + 8)), Some(&7)); + server.verify().await; + } + + #[tokio::test] + async fn a_cold_read_waits_for_a_bounded_refresh_permit() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/sh.tangled.git.temp2.getCommitStats")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({"months": []}))) + .expect(1) + .mount(&server) + .await; + let cache = Arc::new(CommitStatsCache::new(4, Duration::ZERO)); + let permits = Arc::clone(&cache.refreshes) + .acquire_many_owned(MAX_CONCURRENT_REFRESHES as u32) + .await + .unwrap(); + let read = tokio::spawn({ + let cache = Arc::clone(&cache); + let mirror = mirror(&server); + async move { cache.get(Some(mirror), &actor()).await } + }); + + for _ in 0..50 { + if cache + .state + .lock() + .unwrap() + .entries + .get(actor().as_ref()) + .is_some_and(|entry| entry.refreshing) + { + break; + } + tokio::task::yield_now().await; + } + assert!( + cache + .state + .lock() + .unwrap() + .entries + .get(actor().as_ref()) + .is_some_and(|entry| entry.refreshing && !entry.loaded) + ); + assert!(server.received_requests().await.unwrap().is_empty()); + drop(permits); + assert!(read.await.unwrap().is_empty()); + server.verify().await; + } + + #[tokio::test] + async fn cancelling_a_cold_read_releases_waiters() { + let server = MockServer::start().await; + let cache = Arc::new(CommitStatsCache::new(4, Duration::ZERO)); + let permits = Arc::clone(&cache.refreshes) + .acquire_many_owned(MAX_CONCURRENT_REFRESHES as u32) + .await + .unwrap(); + let read = tokio::spawn({ + let cache = Arc::clone(&cache); + let mirror = mirror(&server); + async move { cache.get(Some(mirror), &actor()).await } + }); + for _ in 0..50 { + if cache + .state + .lock() + .unwrap() + .entries + .get(actor().as_ref()) + .is_some_and(|entry| entry.refreshing) + { + break; + } + tokio::task::yield_now().await; + } + assert!( + cache + .state + .lock() + .unwrap() + .entries + .get(actor().as_ref()) + .is_some_and(|entry| entry.refreshing && !entry.loaded) + ); + read.abort(); + assert!(read.await.unwrap_err().is_cancelled()); + assert!( + cache + .state + .lock() + .unwrap() + .entries + .get(actor().as_ref()) + .is_some_and(|entry| !entry.refreshing && !entry.loaded) + ); + drop(permits); + } + + #[tokio::test] + async fn cold_entries_do_not_exceed_capacity_while_refreshes_wait() { + let server = MockServer::start().await; + let cache = Arc::new(CommitStatsCache::new(1, Duration::ZERO)); + let permits = Arc::clone(&cache.refreshes) + .acquire_many_owned(MAX_CONCURRENT_REFRESHES as u32) + .await + .unwrap(); + let first_actor = Did::new_static("did:plc:firstcoldactor").unwrap(); + let second_actor = Did::new_static("did:plc:secondcoldactor").unwrap(); + let first = tokio::spawn({ + let cache = Arc::clone(&cache); + let mirror = mirror(&server); + async move { cache.get(Some(mirror), &first_actor).await } + }); + for _ in 0..50 { + if cache.state.lock().unwrap().entries.len() == 1 { + break; + } + tokio::task::yield_now().await; + } + let second = tokio::spawn({ + let cache = Arc::clone(&cache); + let mirror = mirror(&server); + async move { cache.get(Some(mirror), &second_actor).await } + }); + for _ in 0..50 { + tokio::task::yield_now().await; + } + assert_eq!(cache.state.lock().unwrap().entries.len(), 1); + + first.abort(); + assert!(first.await.unwrap_err().is_cancelled()); + for _ in 0..50 { + if cache + .state + .lock() + .unwrap() + .entries + .contains_key("did:plc:secondcoldactor") + { + break; + } + tokio::task::yield_now().await; + } + let state = cache.state.lock().unwrap(); + assert_eq!(state.entries.len(), 1); + assert!(state.entries.contains_key("did:plc:secondcoldactor")); + drop(state); + second.abort(); + assert!(second.await.unwrap_err().is_cancelled()); + drop(permits); + } + + #[test] + fn dropping_a_refresh_guard_clears_stale_refresh_state() { + let cache = Arc::new(CommitStatsCache::new(1, Duration::ZERO)); + cache.seed(actor().as_ref(), BTreeMap::new(), Some(Instant::now())); + cache + .state + .lock() + .unwrap() + .entries + .get_mut(actor().as_ref()) + .unwrap() + .refreshing = true; + drop(RefreshGuard::new(Arc::clone(&cache), actor().to_string())); + assert!( + !cache + .state + .lock() + .unwrap() + .entries + .get(actor().as_ref()) + .unwrap() + .refreshing + ); + } + + #[tokio::test] + async fn a_failed_refresh_keeps_stale_counts() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/xrpc/sh.tangled.git.temp2.getCommitStats")) + .respond_with(ResponseTemplate::new(503)) + .expect(1) + .mount(&server) + .await; + + let cache = Arc::new(CommitStatsCache::new(4, Duration::ZERO)); + cache.seed( + actor().as_ref(), + BTreeMap::from([(2026 * 12 + 8, 5)]), + Some(Instant::now()), + ); + let counts = cache.get(Some(mirror(&server)), &actor()).await; + assert_eq!(counts.get(&(2026 * 12 + 8)), Some(&5)); + for _ in 0..50 { + let refreshing = cache + .state + .lock() + .unwrap() + .entries + .get(actor().as_ref()) + .is_some_and(|entry| entry.refreshing); + if !refreshing { + break; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + assert_eq!( + cache + .get(Some(mirror(&server)), &actor()) + .await + .get(&(2026 * 12 + 8)), + Some(&5), + ); + server.verify().await; + } +} diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index a1082e671..f40be3b50 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -114,9 +114,11 @@ mod actor; mod backpressure; mod client_address; mod codesearch; +mod commit_stats; mod enrich; mod feed; mod filter; +mod profile_activity; mod recordpath; mod repo; mod view; @@ -172,6 +174,7 @@ pub struct AppState { service_auth_config: service_auth::ServiceAuthConfig>, enrich_router: Arc>, trending_cache: Arc>>, + commit_stats_cache: Arc, } impl AppState { @@ -220,6 +223,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()), } } @@ -327,6 +331,10 @@ pub fn router(state: AppState) -> Router { .route("/xrpc/sh.tangled.repo.getRepoByName", get(get_repo_by_name)) .route("/xrpc/sh.tangled.actor.getProfile", get(get_profile)) .route("/xrpc/sh.tangled.actor.getProfiles", get(get_profiles)) + .route( + "/xrpc/sh.tangled.actor.getProfileActivity", + get(profile_activity::get_profile_activity), + ) .route( "/xrpc/sh.tangled.actor.getTrending", get(actor::get_trending), @@ -1238,7 +1246,6 @@ edge_kinds! { "org.tangled.feed.subscription" => SubscriptionRecord, SubjectShape::BareDidOrAnyAtUri; } -// the fork edge is not a collection so it is not in the table above const FORK_SUBJECT_SHAPE: SubjectShape = SubjectShape::BareDid; fn parse_subject(raw: &SubjectQuery, shape: SubjectShape) -> Result { diff --git a/bobbin/crates/xrpc/src/profile_activity.rs b/bobbin/crates/xrpc/src/profile_activity.rs new file mode 100644 index 000000000..9f0d27ea1 --- /dev/null +++ b/bobbin/crates/xrpc/src/profile_activity.rs @@ -0,0 +1,499 @@ +use std::collections::{BTreeMap, HashMap}; + +use axum::{extract::State, response::Json}; +use bobbin_edge_index::{EdgeItem, EdgeStore, PageCursor, PageLimit, SortDir}; +use bobbin_types::ids::{EdgeKey, SubjectRef, nsid_static, owner_did_from_aturi}; +use bobbin_types::sh_tangled::actor::get_profile_activity::{ + ActivityMonth, GetProfileActivityOutput, IssueActivity, IssueActivityState, PullActivity, + PullActivityState, RepoActivity, RepoActivityKind, RepoRef, +}; +use bobbin_types::sh_tangled::repo::issue::{Issue, IssueRecord}; +use bobbin_types::sh_tangled::repo::pull::{Pull, PullRecord}; +use bobbin_types::sh_tangled::repo::{Repo, RepoRecord}; +use chrono::{DateTime, TimeZone, Utc}; +use futures::{StreamExt, stream}; +use jacquard_common::DefaultStr; +use jacquard_common::types::string::{AtUri, Datetime, Did}; +use jacquard_common::xrpc::XrpcResp; +use serde::Deserialize; + +use crate::{ + AppState, XrpcError, XrpcQuery, accept_state_source, commit_stats::month_number, fetch, + source_authority_did, view::build_profile_basic, +}; + +const DEFAULT_MONTHS: u32 = 7; +const MAX_MONTHS: u32 = 12; +const PAGE_SIZE: u32 = 1_000; +const MAX_ACTIVITY_ITEMS: usize = 1_000; +const HYDRATE_CONCURRENCY: usize = 8; + +const ISSUE_BY_NSID: &str = "sh.tangled.repo.issue.by"; +const PULL_BY_NSID: &str = "sh.tangled.repo.pull.by"; + +#[derive(Debug, Deserialize)] +pub(crate) struct GetProfileActivityQuery { + actor: Did, + months: Option, +} + +#[derive(Clone, Debug)] +struct SkeletonItem { + edge: EdgeItem, + months_ago: usize, +} + +enum RawActivity { + Repo { + months_ago: usize, + repo: RepoRef, + kind: RepoActivityKind, + }, + Issue { + months_ago: usize, + uri: AtUri, + title: DefaultStr, + state: IssueActivityState, + repo_did: Did, + }, + Pull { + months_ago: usize, + uri: AtUri, + title: DefaultStr, + state: PullActivityState, + repo_did: Did, + }, +} + +pub(crate) async fn get_profile_activity( + State(state): State, + XrpcQuery(query): XrpcQuery, +) -> Result>, XrpcError> { + let months = query.months.unwrap_or(DEFAULT_MONTHS); + if !(1..=MAX_MONTHS).contains(&months) { + return Err(XrpcError::InvalidParams(format!( + "months must be between 1 and {MAX_MONTHS}", + ))); + } + let _permit = state.heavy_permit()?; + let now = Utc::now(); + let commit_counts = state + .commit_stats_cache + .get(state.mirror_v2.clone(), &query.actor); + let (skeleton, truncated) = collect_skeleton(&state.edges, &query.actor, now, months); + let hydration = stream::iter( + skeleton + .into_iter() + .map(|item| hydrate_activity(&state, item)), + ) + .buffered(HYDRATE_CONCURRENCY) + .collect::>(); + let (commit_counts, hydrated) = tokio::join!(commit_counts, hydration); + + let mut activity = Vec::with_capacity(hydrated.len()); + for item in hydrated { + match item { + Ok(item) => activity.push(item), + Err((uri, error)) => { + tracing::warn!(%uri, %error, "skipping profile activity item, hydration failed") + } + } + } + + let months = assemble_months(&state, activity, &commit_counts, now, months).await; + Ok(Json(GetProfileActivityOutput { + months, + truncated, + extra_data: None, + })) +} + +fn activity_keys(subject: &Did) -> [EdgeKey; 3] { + [ + EdgeKey::new( + nsid_static(RepoRecord::NSID), + SubjectRef::Did(subject.clone()), + ), + EdgeKey::new(nsid_static(ISSUE_BY_NSID), SubjectRef::Did(subject.clone())), + EdgeKey::new(nsid_static(PULL_BY_NSID), SubjectRef::Did(subject.clone())), + ] +} + +fn collect_skeleton( + edges: &EdgeStore, + subject: &Did, + now: DateTime, + months: u32, +) -> (Vec, bool) { + let keys = activity_keys(subject); + let limit = PageLimit::new(PAGE_SIZE).expect("profile activity page size is valid"); + let current_month = month_number(now); + let earliest_month = current_month - (months as i32 - 1); + let now_micros = now.timestamp_micros(); + let mut cursor = PageCursor::Start; + let mut items = Vec::new(); + let mut truncated = false; + + 'pages: loop { + let page = edges.list_multi(&keys, cursor, limit, SortDir::Desc); + for edge in page.items { + let Ok(micros) = i64::try_from(edge.sort_micros) else { + continue; + }; + if micros > now_micros { + continue; + } + let Some(at) = DateTime::::from_timestamp_micros(micros) else { + continue; + }; + let event_month = month_number(at); + if event_month < earliest_month { + break 'pages; + } + if event_month > current_month { + continue; + } + if items.len() == MAX_ACTIVITY_ITEMS { + truncated = true; + break 'pages; + } + items.push(SkeletonItem { + edge, + months_ago: (current_month - event_month) as usize, + }); + } + let Some(next) = page.next else { + break; + }; + cursor = PageCursor::After(next); + } + (items, truncated) +} + +async fn hydrate_activity( + state: &AppState, + item: SkeletonItem, +) -> Result, XrpcError)> { + let uri = item.edge.uri; + let collection = uri.collection(); + let result: Result = async { + match collection.as_ref().map(|value| value.as_ref()) { + Some(RepoRecord::NSID) => { + let (_, record) = fetch::>(state, &uri).await?; + let kind = if record.source.is_some() { + RepoActivityKind::Forked + } else { + RepoActivityKind::Created + }; + let repo = repo_ref_from_record(state, &uri, &record).await?; + Ok(RawActivity::Repo { + months_ago: item.months_ago, + repo, + kind, + }) + } + Some(IssueRecord::NSID) => { + let (_, record) = fetch::>(state, &uri).await?; + let state_value = issue_state(state, &uri, &record); + Ok(RawActivity::Issue { + months_ago: item.months_ago, + uri: uri.clone(), + title: record.title, + state: state_value, + repo_did: record.repo, + }) + } + Some(PullRecord::NSID) => { + let (_, record) = fetch::>(state, &uri).await?; + let state_value = pull_state(state, &uri, &record); + Ok(RawActivity::Pull { + months_ago: item.months_ago, + uri: uri.clone(), + title: record.title, + state: state_value, + repo_did: record.target.repo, + }) + } + Some(other) => Err(XrpcError::InvalidRecord(format!( + "unexpected profile activity collection: {other}", + ))), + None => Err(XrpcError::InvalidRecord( + "profile activity uri has no collection".into(), + )), + } + } + .await; + result.map_err(|error| (uri, error)) +} + +async fn assemble_months( + state: &AppState, + activity: Vec, + commit_counts: &BTreeMap, + now: DateTime, + count: u32, +) -> Vec> { + let mut refs = HashMap::>::new(); + let mut targets = BTreeMap::>::new(); + for item in &activity { + if let RawActivity::Repo { repo, .. } = item + && let Some(did) = &repo.did + { + refs.insert(did.to_string(), repo.clone()); + } + } + for item in &activity { + if let RawActivity::Issue { repo_did, .. } | RawActivity::Pull { repo_did, .. } = item + && !refs.contains_key(repo_did.as_ref()) + { + targets.insert(repo_did.to_string(), repo_did.clone()); + } + } + + let hydrated = stream::iter(targets.into_values().map(|did| async move { + let result = repo_ref_by_did(state, &did).await; + (did, result) + })) + .buffered(HYDRATE_CONCURRENCY) + .collect::>() + .await; + for (did, result) in hydrated { + match result { + Ok(repo) => { + refs.insert(did.to_string(), repo); + } + Err(error) => { + tracing::warn!(repo = %did, %error, "skipping unresolved profile activity repo") + } + } + } + + let current_month = month_number(now); + let mut months = (0..count) + .map(|months_ago| { + let month = current_month - months_ago as i32; + ActivityMonth { + start: Datetime::from(month_start(month).fixed_offset()), + commits: commit_counts.get(&month).copied().unwrap_or(0), + repos: Vec::new(), + issues: Vec::new(), + pulls: Vec::new(), + extra_data: None, + } + }) + .collect::>(); + + for item in activity { + match item { + RawActivity::Repo { + months_ago, + repo, + kind, + } => months[months_ago].repos.push(RepoActivity { + repo, + kind, + extra_data: None, + }), + RawActivity::Issue { + months_ago, + uri, + title, + state, + repo_did, + } => { + let Some(repo) = refs.get(repo_did.as_ref()) else { + continue; + }; + months[months_ago].issues.push(IssueActivity { + uri, + title, + state, + repo: repo.clone(), + extra_data: None, + }); + } + RawActivity::Pull { + months_ago, + uri, + title, + state, + repo_did, + } => { + let Some(repo) = refs.get(repo_did.as_ref()) else { + continue; + }; + months[months_ago].pulls.push(PullActivity { + uri, + title, + state, + repo: repo.clone(), + extra_data: None, + }); + } + } + } + + months.retain(|month| { + month.commits > 0 + || !month.repos.is_empty() + || !month.issues.is_empty() + || !month.pulls.is_empty() + }); + months +} + +async fn repo_ref_by_did( + state: &AppState, + did: &Did, +) -> Result, XrpcError> { + let ident = state + .resolver + .lookup_by_repo_did(did) + .await + .ok_or(XrpcError::NotFound)?; + let uri = AtUri::from_parts_owned(ident.owner.as_str(), RepoRecord::NSID, ident.rkey.as_str()) + .map_err(|error| XrpcError::InvalidRecord(format!("repo uri assembly: {error}")))?; + let (_, record) = fetch::>(state, &uri).await?; + repo_ref_from_record(state, &uri, &record).await +} + +async fn repo_ref_from_record( + state: &AppState, + uri: &AtUri, + record: &Repo, +) -> Result, XrpcError> { + let owner_did = owner_did_from_aturi(uri) + .ok_or_else(|| XrpcError::InvalidRecord("repo uri authority must be a did".into()))?; + let owner = build_profile_basic(state, None, &owner_did).await; + let name = record.name.clone().unwrap_or_else(|| { + uri.rkey() + .map(|rkey| DefaultStr::from(rkey.as_ref())) + .unwrap_or_default() + }); + Ok(RepoRef { + uri: uri.clone(), + did: record.repo_did.clone(), + owner_did: owner.did, + owner_handle: owner.handle, + name, + extra_data: None, + }) +} + +fn issue_state( + state: &AppState, + uri: &AtUri, + issue: &Issue, +) -> IssueActivityState { + let author = source_authority_did(uri); + state + .issue_states + .latest_by(uri, |source| { + accept_state_source(source, author.as_ref(), &issue.repo) + }) + .map_or(IssueActivityState::Open, |(kind, _)| match kind { + bobbin_edge_index::IssueStateKind::Open => IssueActivityState::Open, + bobbin_edge_index::IssueStateKind::Closed => IssueActivityState::Closed, + }) +} + +fn pull_state( + state: &AppState, + uri: &AtUri, + pull: &Pull, +) -> PullActivityState { + let author = source_authority_did(uri); + state + .pull_statuses + .latest_by(uri, |source| { + accept_state_source(source, author.as_ref(), &pull.target.repo) + }) + .map_or(PullActivityState::Open, |(kind, _)| match kind { + bobbin_edge_index::PullStatusKind::Open => PullActivityState::Open, + bobbin_edge_index::PullStatusKind::Closed => PullActivityState::Closed, + bobbin_edge_index::PullStatusKind::Merged => PullActivityState::Merged, + }) +} + +fn month_start(value: i32) -> DateTime { + let year = value.div_euclid(12); + let month = value.rem_euclid(12) as u32 + 1; + Utc.with_ymd_and_hms(year, month, 1, 0, 0, 0) + .single() + .expect("calendar month start is valid") +} + +#[cfg(test)] +mod tests { + use super::*; + use bobbin_runtime::RuntimeHasher; + use bobbin_types::edges::Edge; + + fn did(value: &str) -> Did { + Did::new_owned(value).expect("valid did") + } + + #[test] + fn skeleton_merges_kinds_and_stops_at_month_cutoff() { + let store = EdgeStore::new(RuntimeHasher::default()); + let actor = did("did:plc:profileactor"); + let now = Utc.with_ymd_and_hms(2026, 9, 15, 0, 0, 0).single().unwrap(); + let samples = [ + ("sh.tangled.repo", "sh.tangled.repo", "2026-09-01T00:00:00Z"), + ( + ISSUE_BY_NSID, + "sh.tangled.repo.issue", + "2026-08-01T00:00:00Z", + ), + (PULL_BY_NSID, "sh.tangled.repo.pull", "2026-03-01T00:00:00Z"), + ]; + for (index, (kind, collection, at)) in samples.into_iter().enumerate() { + let timestamp = DateTime::parse_from_rfc3339(at).unwrap().timestamp_micros() as u64; + store.add(Edge { + kind: nsid_static(kind), + subject: SubjectRef::Did(actor.clone()), + source: AtUri::new_owned(format!("at://did:plc:owner/{collection}/item{index}",)) + .unwrap(), + sort_micros: timestamp, + }); + } + + let (items, truncated) = collect_skeleton(&store, &actor, now, 7); + assert!(!truncated); + assert_eq!( + items.iter().map(|item| item.months_ago).collect::>(), + vec![0, 1, 6], + ); + } + #[test] + fn skeleton_reports_only_omitted_activity_as_truncated() { + let actor = did("did:plc:profileactor"); + let now = Utc.with_ymd_and_hms(2026, 9, 15, 0, 0, 0).single().unwrap(); + let base = Utc + .with_ymd_and_hms(2026, 9, 1, 0, 0, 0) + .single() + .unwrap() + .timestamp_micros() as u64; + + for (item_count, expected_truncated) in + [(MAX_ACTIVITY_ITEMS, false), (MAX_ACTIVITY_ITEMS + 1, true)] + { + let store = EdgeStore::new(RuntimeHasher::default()); + for index in 0..item_count { + store.add(Edge { + kind: nsid_static("sh.tangled.repo"), + subject: SubjectRef::Did(actor.clone()), + source: AtUri::new_owned(format!( + "at://did:plc:owner/sh.tangled.repo/item{index}", + )) + .unwrap(), + sort_micros: base + index as u64, + }); + } + + let (items, truncated) = collect_skeleton(&store, &actor, now, 7); + assert_eq!(truncated, expected_truncated, "item count {item_count}"); + assert_eq!(items.len(), MAX_ACTIVITY_ITEMS); + } + } +} diff --git a/bobbin/crates/xrpc/tests/aggregation.rs b/bobbin/crates/xrpc/tests/aggregation.rs index 98bdb9115..f0696f33b 100644 --- a/bobbin/crates/xrpc/tests/aggregation.rs +++ b/bobbin/crates/xrpc/tests/aggregation.rs @@ -5,7 +5,7 @@ use bobbin_edge_index::{ Coverage, CoverageWatch, EdgeStore, HydrantCursor, InviteIndex, InviteScope, IssueStateKind, Offer, PageToken, PullStatusKind, StateIndex, }; -use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig}; +use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig, MirrorProxy}; use bobbin_record_lru::{CacheCapacity, LruRecordStore}; use bobbin_resolver::RepoIdResolver; use bobbin_runtime::{RuntimeHasher, SystemClock, UnixMicros}; @@ -14,6 +14,7 @@ use bobbin_slingshot_client::SlingshotClient; use bobbin_types::edges::Edge; use bobbin_types::ids::SubjectRef; use bobbin_xrpc::{AppState, router}; +use chrono::Utc; use futures::stream::{self, StreamExt}; use http::{Request, StatusCode}; use jacquard_common::DefaultStr; @@ -2097,6 +2098,169 @@ async fn extractor_to_xrpc_round_trip_for_star() { ); } +#[tokio::test] +async fn profile_activity_returns_one_hydrated_aggregate() { + let mut h = Harness::new().await; + let actor = did("did:plc:lyna"); + let created_repo_did = did("did:plc:freshrepo"); + let target_repo_did = did("did:plc:limpet"); + let target_rkey = rkey("target"); + let created_repo_uri = at(&format!("at://{}/sh.tangled.repo/fresh", actor.as_ref(),)); + let issue_uri = at(&format!("at://{}/sh.tangled.repo.issue/i1", actor.as_ref(),)); + let pull_uri = at(&format!("at://{}/sh.tangled.repo.pull/p1", actor.as_ref(),)); + let event_micros = Utc::now().timestamp_micros() as u64; + Mock::given(method("GET")) + .and(path("/xrpc/sh.tangled.git.temp2.getCommitStats")) + .and(query_param("actor", actor.as_ref())) + .and(query_param("months", "24")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + "months": [{ + "start": Utc::now().format("%Y-%m-01T00:00:00Z").to_string(), + "commits": 3, + }] + }))) + .expect(1) + .mount(&h.server) + .await; + h.state = h.state.clone().with_mirror_v2(Some(Arc::new( + MirrorProxy::new( + &Url::parse(&h.server.uri()).unwrap(), + Arc::new(SystemClock::new()), + RuntimeHasher::default(), + ) + .unwrap(), + ))); + + for (kind, source) in [ + ("sh.tangled.repo", &created_repo_uri), + ("sh.tangled.repo.issue.by", &issue_uri), + ("sh.tangled.repo.pull.by", &pull_uri), + ] { + h.edges.add(Edge { + kind: nsid(kind), + subject: SubjectRef::Did(actor.clone()), + source: source.clone(), + sort_micros: event_micros, + }); + } + + h.mount( + &actor, + &nsid("sh.tangled.repo"), + &rkey("fresh"), + json!({ + "$type": "sh.tangled.repo", + "name": "fresh", + "knot": "oyster.cafe", + "repoDid": created_repo_did.as_ref(), + "createdAt": Utc::now().to_rfc3339(), + }), + ) + .await; + h.mount( + &actor, + &nsid("sh.tangled.repo"), + &target_rkey, + json!({ + "$type": "sh.tangled.repo", + "name": "core", + "knot": "oyster.cafe", + "repoDid": target_repo_did.as_ref(), + "createdAt": "2025-01-01T00:00:00Z", + }), + ) + .await; + h.state + .resolver + .observe( + actor.clone(), + target_rkey, + Some(target_repo_did.clone()), + Some(DefaultStr::from("core")), + ) + .await; + h.mount( + &actor, + &nsid("sh.tangled.repo.issue"), + &rkey("i1"), + issue_body(&target_repo_did, "fix the thing"), + ) + .await; + h.mount( + &actor, + &nsid("sh.tangled.repo.pull"), + &rkey("p1"), + pull_body(&target_repo_did, "ship the thing"), + ) + .await; + h.upsert_issue_state( + at(&format!( + "at://{}/sh.tangled.repo.issue.state/s1", + actor.as_ref(), + )), + issue_uri, + event_micros + 1, + IssueStateKind::Closed, + ); + h.upsert_pull_status( + at(&format!( + "at://{}/sh.tangled.repo.pull.status/s1", + actor.as_ref(), + )), + pull_uri, + event_micros + 1, + PullStatusKind::Merged, + ); + + let (status, body) = json_response( + router(h.state.clone()) + .oneshot( + Request::builder() + .uri(format!( + "/xrpc/sh.tangled.actor.getProfileActivity?actor={}", + actor.as_ref(), + )) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(), + ) + .await; + + assert_eq!(status, StatusCode::OK, "body was {body}"); + assert_eq!(body["truncated"], false); + let months = body["months"].as_array().expect("months array"); + assert_eq!(months.len(), 1); + assert_eq!(months[0]["commits"], 3); + assert_eq!(months[0]["repos"][0]["kind"], "created"); + assert_eq!(months[0]["repos"][0]["repo"]["name"], "fresh"); + assert_eq!(months[0]["issues"][0]["state"], "closed"); + assert_eq!(months[0]["issues"][0]["repo"]["name"], "core"); + assert_eq!(months[0]["pulls"][0]["state"], "merged"); + assert_eq!(months[0]["pulls"][0]["repo"]["name"], "core"); +} + +#[tokio::test] +async fn profile_activity_validates_the_month_window() { + let h = Harness::new().await; + let (status, body) = json_response( + router(h.state.clone()) + .oneshot( + Request::builder() + .uri("/xrpc/sh.tangled.actor.getProfileActivity?actor=did:plc:lyna&months=13") + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(), + ) + .await; + + assert_eq!(status, StatusCode::BAD_REQUEST); + assert_eq!(body["error"], "InvalidRequest"); +} + #[tokio::test] async fn list_issues_includes_state_comment_count_and_state_updated_at() { let h = Harness::new().await; diff --git a/lexicons/actor/getProfileActivity.json b/lexicons/actor/getProfileActivity.json new file mode 100644 index 000000000..89d34ae7f --- /dev/null +++ b/lexicons/actor/getProfileActivity.json @@ -0,0 +1,204 @@ +{ + "lexicon": 1, + "id": "sh.tangled.actor.getProfileActivity", + "defs": { + "main": { + "type": "query", + "description": "Get recent Tangled activity for an actor, grouped by calendar month.", + "parameters": { + "type": "params", + "required": [ + "actor" + ], + "properties": { + "actor": { + "type": "string", + "format": "did", + "description": "Actor whose activity to return." + }, + "months": { + "type": "integer", + "minimum": 1, + "maximum": 12, + "default": 7, + "description": "Number of UTC calendar months to include, including the current month." + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": [ + "months", + "truncated" + ], + "properties": { + "months": { + "type": "array", + "items": { + "type": "ref", + "ref": "#activityMonth" + } + }, + "truncated": { + "type": "boolean", + "description": "True when the result reached the server activity-item safety limit." + } + } + } + } + }, + "activityMonth": { + "type": "object", + "required": [ + "start", + "commits", + "repos", + "issues", + "pulls" + ], + "properties": { + "start": { + "type": "string", + "format": "datetime", + "description": "Start of this UTC calendar month." + }, + "commits": { + "type": "integer", + "minimum": 0 + }, + "repos": { + "type": "array", + "items": { + "type": "ref", + "ref": "#repoActivity" + } + }, + "issues": { + "type": "array", + "items": { + "type": "ref", + "ref": "#issueActivity" + } + }, + "pulls": { + "type": "array", + "items": { + "type": "ref", + "ref": "#pullActivity" + } + } + } + }, + "repoRef": { + "type": "object", + "required": [ + "uri", + "ownerDid", + "ownerHandle", + "name" + ], + "properties": { + "uri": { + "type": "string", + "format": "at-uri" + }, + "did": { + "type": "string", + "format": "did" + }, + "ownerDid": { + "type": "string", + "format": "did" + }, + "ownerHandle": { + "type": "string", + "format": "handle" + }, + "name": { + "type": "string" + } + } + }, + "repoActivity": { + "type": "object", + "required": [ + "repo", + "kind" + ], + "properties": { + "repo": { + "type": "ref", + "ref": "#repoRef" + }, + "kind": { + "type": "string", + "knownValues": [ + "created", + "forked" + ] + } + } + }, + "issueActivity": { + "type": "object", + "required": [ + "uri", + "title", + "state", + "repo" + ], + "properties": { + "uri": { + "type": "string", + "format": "at-uri" + }, + "title": { + "type": "string" + }, + "state": { + "type": "string", + "knownValues": [ + "open", + "closed" + ] + }, + "repo": { + "type": "ref", + "ref": "#repoRef" + } + } + }, + "pullActivity": { + "type": "object", + "required": [ + "uri", + "title", + "state", + "repo" + ], + "properties": { + "uri": { + "type": "string", + "format": "at-uri" + }, + "title": { + "type": "string" + }, + "state": { + "type": "string", + "knownValues": [ + "open", + "closed", + "merged" + ] + }, + "repo": { + "type": "ref", + "ref": "#repoRef" + } + } + } + } +}