From b27da201f662c5c902f3d7ab2d2aa2b5c9fe170d Mon Sep 17 00:00:00 2001 From: dawn Date: Wed, 9 Sep 2026 04:46:46 +0300 Subject: [PATCH] bobbin/{edge-index,xrpc}: index authors distinctly for follow, star, etc. --- bobbin/crates/edge-index/Cargo.toml | 4 + .../edge-index/benches/distinct_authors.rs | 154 ++++ bobbin/crates/edge-index/src/lib.rs | 683 +++++++++++++++--- bobbin/crates/edge-index/src/sorted_blocks.rs | 258 ++++++- bobbin/crates/xrpc/src/filter.rs | 4 + bobbin/crates/xrpc/src/lib.rs | 379 ++++++++-- bobbin/crates/xrpc/tests/aggregation.rs | 119 +++ 7 files changed, 1416 insertions(+), 185 deletions(-) create mode 100644 bobbin/crates/edge-index/benches/distinct_authors.rs diff --git a/bobbin/crates/edge-index/Cargo.toml b/bobbin/crates/edge-index/Cargo.toml index 1d0bb76e5..777cc5566 100644 --- a/bobbin/crates/edge-index/Cargo.toml +++ b/bobbin/crates/edge-index/Cargo.toml @@ -27,3 +27,7 @@ tokio = { workspace = true, features = ["macros", "rt", "test-util"] } [[bench]] name = "profile_activity" harness = false + +[[bench]] +name = "distinct_authors" +harness = false diff --git a/bobbin/crates/edge-index/benches/distinct_authors.rs b/bobbin/crates/edge-index/benches/distinct_authors.rs new file mode 100644 index 000000000..4e81eac85 --- /dev/null +++ b/bobbin/crates/edge-index/benches/distinct_authors.rs @@ -0,0 +1,154 @@ +use bobbin_edge_index::{EdgeStore, PageCursor, PageLimit, PageOffset, 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 PAGE_SIZE: u32 = 100; + +fn edge_counts() -> Vec { + std::env::var("BOBBIN_DISTINCT_BENCH_EDGES") + .map(|raw| { + raw.split(',') + .map(str::trim) + .filter(|entry| !entry.is_empty()) + .map(|entry| entry.parse().expect("edge count")) + .collect() + }) + .unwrap_or_else(|_| vec![1_000, 10_000]) +} + +fn did(value: &str) -> Did { + Did::new_owned(value).expect("valid did") +} + +fn populated_store(edges: usize, author_for: impl Fn(usize) -> usize) -> (EdgeStore, EdgeKey) { + let store = EdgeStore::new(RuntimeHasher::default()); + let subject = did("did:plc:subject"); + let kind = nsid_static("sh.tangled.feed.star"); + let key = EdgeKey::new(kind.clone(), SubjectRef::Did(subject.clone())); + + (0..edges).for_each(|index| { + let author = author_for(index); + let source = AtUri::new_owned(format!( + "at://did:plc:author{author}/sh.tangled.feed.star/r{index}", + )) + .expect("valid at-uri"); + store.add(Edge { + kind: kind.clone(), + subject: SubjectRef::Did(subject.clone()), + source, + sort_micros: index as u64, + }); + }); + + (store, key) +} + +fn unique_store(edges: usize) -> (EdgeStore, EdgeKey) { + populated_store(edges, |index| index) +} + +fn paired_author_store(edges: usize) -> (EdgeStore, EdgeKey) { + populated_store(edges, |index| index / 2) +} + +fn hot_author_store(edges: usize) -> (EdgeStore, EdgeKey) { + let authors = (edges / 10).max(PAGE_SIZE as usize + 1); + populated_store(edges, |index| index.min(authors - 1)) +} + +fn page_limit() -> PageLimit { + PageLimit::new(PAGE_SIZE).expect("valid page size") +} + +fn walk(store: &EdgeStore, key: &EdgeKey) -> usize { + std::iter::successors( + Some(store.list_distinct_authors(key, PageCursor::Start, page_limit(), SortDir::Desc)), + |page| { + page.next.map(|token| { + store.list_distinct_authors( + key, + PageCursor::After(token), + page_limit(), + SortDir::Desc, + ) + }) + }, + ) + .map(|page| page.items.len()) + .sum() +} + +fn main() { + divan::main(); +} + +#[divan::bench(args = edge_counts())] +fn build_unique(bencher: Bencher, edges: usize) { + bencher + .counter(ItemsCount::new(edges)) + .bench_local(|| unique_store(edges)); +} + +#[divan::bench(args = edge_counts())] +fn build_paired_authors(bencher: Bencher, edges: usize) { + bencher + .counter(ItemsCount::new(edges)) + .bench_local(|| paired_author_store(edges)); +} + +#[divan::bench(args = edge_counts())] +fn build_hot_author(bencher: Bencher, edges: usize) { + bencher + .counter(ItemsCount::new(edges)) + .bench_local(|| hot_author_store(edges)); +} + +#[divan::bench(args = edge_counts())] +fn first_page_unique(bencher: Bencher, edges: usize) { + let (store, key) = unique_store(edges); + bencher.counter(ItemsCount::new(PAGE_SIZE)).bench_local(|| { + store.list_distinct_authors(&key, PageCursor::Start, page_limit(), SortDir::Desc) + }); +} + +#[divan::bench(args = edge_counts())] +fn first_page_hot_author(bencher: Bencher, edges: usize) { + let (store, key) = hot_author_store(edges); + bencher.counter(ItemsCount::new(PAGE_SIZE)).bench_local(|| { + store.list_distinct_authors(&key, PageCursor::Start, page_limit(), SortDir::Desc) + }); +} + +#[divan::bench(args = edge_counts())] +fn middle_cursor_unique(bencher: Bencher, edges: usize) { + let (store, key) = unique_store(edges); + let offset = edges / 2 - PAGE_SIZE as usize; + let cursor = store + .list_distinct_authors( + &key, + PageOffset::new(offset as u64), + page_limit(), + SortDir::Desc, + ) + .next + .expect("middle cursor"); + + bencher.counter(ItemsCount::new(PAGE_SIZE)).bench_local(|| { + store.list_distinct_authors(&key, PageCursor::After(cursor), page_limit(), SortDir::Desc) + }); +} + +#[divan::bench(args = edge_counts())] +fn full_walk_unique(bencher: Bencher, edges: usize) { + let (store, key) = unique_store(edges); + bencher + .counter(ItemsCount::new(edges)) + .bench_local(|| walk(&store, &key)); +} diff --git a/bobbin/crates/edge-index/src/lib.rs b/bobbin/crates/edge-index/src/lib.rs index fb1c45a0f..07ac10d41 100644 --- a/bobbin/crates/edge-index/src/lib.rs +++ b/bobbin/crates/edge-index/src/lib.rs @@ -1,4 +1,4 @@ -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::marker::PhantomData; use std::num::NonZeroU32; use std::sync::atomic::{AtomicU32, Ordering}; @@ -20,9 +20,27 @@ use scc::hash_map::Entry; use smallvec::SmallVec; use thiserror::Error; +#[repr(C)] #[derive(Clone, Copy, Debug, Eq, Hash, PartialEq, Ord, PartialOrd)] -pub(crate) struct SortMicros(u64); +pub(crate) struct SortMicros { + high: u32, + low: u32, +} + +impl SortMicros { + const fn new(value: u64) -> Self { + Self { + high: (value >> 32) as u32, + low: value as u32, + } + } + const fn get(self) -> u64 { + ((self.high as u64) << 32) | self.low as u64 + } +} + +#[repr(C)] #[derive(Clone, Copy, Debug, Eq, Hash, PartialEq, Ord, PartialOrd)] pub(crate) struct BucketKey { micros: SortMicros, @@ -30,22 +48,19 @@ pub(crate) struct BucketKey { } impl BucketKey { - pub(crate) fn new(micros: u64, source: SourceId) -> Self { + pub(crate) const fn new(micros: u64, source: SourceId) -> Self { Self { - micros: SortMicros(micros), + micros: SortMicros::new(micros), source, } } fn token(self) -> PageToken { - PageToken::new(self.micros.0, self.source.index()) + PageToken::new(self.micros.get(), self.source.index()) } fn from_token(tok: PageToken) -> Self { - Self { - micros: SortMicros(tok.micros), - source: SourceId::from_raw(tok.source), - } + Self::new(tok.micros, SourceId::from_raw(tok.source)) } } @@ -357,6 +372,13 @@ pub struct EdgePage { pub struct EdgeItem { pub uri: AtUri, pub sort_micros: u64, + pub sort_source: u32, +} + +impl EdgeItem { + pub fn sort_token(&self) -> PageToken { + PageToken::new(self.sort_micros, self.sort_source) + } } impl AsRef for EdgeItem { @@ -487,13 +509,215 @@ pub(crate) fn drop_author( } } +enum AuthorHistory { + Small(SmallVec<[BucketKey; 4]>), + Large(BlockStore), +} + +impl AuthorHistory { + fn pair(left: BucketKey, right: BucketKey) -> Self { + let mut keys = SmallVec::from_slice(&[left, right]); + keys.sort_unstable(); + Self::Small(keys) + } + + fn len(&self) -> usize { + match self { + Self::Small(keys) => keys.len(), + Self::Large(keys) => keys.len(), + } + } + + fn newest(&self) -> BucketKey { + match self { + Self::Small(keys) => *keys.last().unwrap(), + Self::Large(keys) => keys.last().expect("nonempty author history"), + } + } + + fn insert(&mut self, key: BucketKey) -> bool { + match self { + Self::Small(keys) => match keys.binary_search(&key) { + Ok(_) => false, + Err(position) if keys.len() < keys.inline_size() => { + keys.insert(position, key); + true + } + Err(position) => { + let mut promoted = keys.to_vec(); + promoted.insert(position, key); + *self = Self::Large(BlockStore::from_sorted(promoted)); + true + } + }, + Self::Large(keys) => keys.insert(key), + } + } + + fn remove(&mut self, key: &BucketKey) -> bool { + match self { + Self::Small(keys) => keys + .binary_search(key) + .map(|position| { + keys.remove(position); + }) + .is_ok(), + Self::Large(keys) => { + if !keys.remove(key) { + return false; + } + if keys.len() <= 4 { + *self = Self::Small(keys.directed(PageCursor::Start, SortDir::Asc).collect()); + } + true + } + } + } + + fn heap_bytes(&self) -> u64 { + match self { + Self::Small(keys) if keys.spilled() => keys.capacity() as u64 * BUCKET_KEY_BYTES, + Self::Small(_) => 0, + Self::Large(keys) => keys.heap_bytes(), + } + } +} + +fn estimated_hashmap_heap_bytes(map: &HashMap) -> u64 { + let slot = std::mem::size_of::<(K, V)>() as u64 + 1; + map.capacity() as u64 * slot +} + +struct AuthorIndex { + newest: HashMap, + histories: HashMap, RuntimeHasher>, + representatives: BlockStore, +} + +impl AuthorIndex { + fn new(hasher: RuntimeHasher) -> Self { + Self { + newest: HashMap::with_hasher(hasher.clone()), + histories: HashMap::with_hasher(hasher), + representatives: BlockStore::from_sorted(Vec::new()), + } + } + + fn len(&self) -> usize { + self.newest.len() + } + + fn contains(&self, author: &AuthorId) -> bool { + self.newest.contains_key(author) + } + + fn count(&self, author: &AuthorId) -> usize { + self.newest + .contains_key(author) + .then(|| { + self.histories + .get(author) + .map_or(1, |history| history.len()) + }) + .unwrap_or(0) + } + + fn replace_representative(&mut self, previous: BucketKey, next: BucketKey) { + if self.representatives.replace_last(previous, next) { + return; + } + let removed = self.representatives.remove(&previous); + let inserted = self.representatives.insert(next); + debug_assert!(removed && inserted); + } + + fn insert(&mut self, author: AuthorId, key: BucketKey) { + let Some(previous) = self.newest.get(&author).copied() else { + let inserted = self.representatives.insert(key); + debug_assert!(inserted); + self.newest.insert(author, key); + return; + }; + let inserted = match self.histories.get_mut(&author) { + Some(history) => history.insert(key), + None => { + self.histories + .insert(author, Box::new(AuthorHistory::pair(previous, key))); + true + } + }; + debug_assert!(inserted); + if key > previous { + self.replace_representative(previous, key); + self.newest.insert(author, key); + } + } + + fn remove(&mut self, author: AuthorId, key: &BucketKey) { + let Some(newest) = self.newest.get(&author).copied() else { + debug_assert!(false, "large bucket missing author"); + return; + }; + let Some(history) = self.histories.get_mut(&author) else { + debug_assert_eq!(*key, newest, "single author key is not newest"); + self.newest.remove(&author); + let removed = self.representatives.remove(&newest); + debug_assert!(removed); + return; + }; + let removed = history.remove(key); + debug_assert!(removed); + let next = history.newest(); + if history.len() == 1 { + self.histories.remove(&author); + } + if *key == newest { + self.replace_representative(newest, next); + self.newest.insert(author, next); + } + } + + fn heap_bytes(&self) -> u64 { + estimated_hashmap_heap_bytes(&self.newest) + + estimated_hashmap_heap_bytes(&self.histories) + + self.representatives.heap_bytes() + + self + .histories + .values() + .map(|history| std::mem::size_of::() as u64 + history.heap_bytes()) + .sum::() + } +} + struct LargeBucket { keys: BlockStore, - authors: HashMap, + authors: AuthorIndex, +} + +impl LargeBucket { + fn insert(&mut self, key: BucketKey, author: Option) -> bool { + if !self.keys.insert(key) { + return false; + } + if let Some(author) = author { + self.authors.insert(author, key); + } + true + } + + fn remove(&mut self, key: &BucketKey, author: Option) -> bool { + if !self.keys.remove(key) { + return false; + } + if let Some(author) = author { + self.authors.remove(author, key); + } + true + } } const BUCKET_PROMOTE_AT: usize = 256; -const SOURCE_BYTES: u64 = 16; +const BUCKET_KEY_BYTES: u64 = std::mem::size_of::() as u64; type BlockStore = SortedBlocks; @@ -539,13 +763,7 @@ impl Sources { true } }, - Self::Large(big) => { - let inserted = big.keys.insert(key); - if inserted && let Some(a) = author { - bump_author(&mut big.authors, a); - } - inserted - } + Self::Large(big) => big.insert(key, author), } } @@ -557,11 +775,7 @@ impl Sources { } } Self::Large(big) => { - if big.keys.remove(key) - && let Some(a) = author - { - drop_author(&mut big.authors, a); - } + big.remove(key, author); } } } @@ -605,16 +819,13 @@ impl Sources { } fn heap_bytes(&self) -> u64 { - const HASHMAP_FIXED: u64 = 48; - const HASHMAP_PER_CAP: u64 = 9; match self { - Self::Small(v) if v.spilled() => v.capacity() as u64 * SOURCE_BYTES, + Self::Small(v) if v.spilled() => v.capacity() as u64 * BUCKET_KEY_BYTES, Self::Small(_) => 0, Self::Large(big) => { std::mem::size_of::() as u64 + big.keys.heap_bytes() - + HASHMAP_FIXED - + big.authors.capacity() as u64 * HASHMAP_PER_CAP + + big.authors.heap_bytes() } } } @@ -623,9 +834,20 @@ impl Sources { #[derive(Clone, Copy, Debug, Eq, PartialEq)] struct ReverseEntry { key_id: EdgeKeyId, - sort_micros: u64, + sort_micros: SortMicros, } +#[cfg(target_pointer_width = "64")] +const _: () = { + assert!(std::mem::size_of::() == 8); + assert!(std::mem::align_of::() == 4); + assert!(std::mem::size_of::() == 12); + assert!(std::mem::align_of::() == 4); + assert!(std::mem::size_of::() == 12); + assert!(std::mem::size_of::<(AuthorId, BucketKey)>() == 16); + assert!(std::mem::size_of::<(AuthorId, Box)>() == 16); +}; + pub struct EdgeStore { source_interner: Arc>, did_interner: Arc>, @@ -762,7 +984,7 @@ impl EdgeStore { drop(entry); if let Some(keys) = promote_keys { - let authors = self.build_author_map(&keys); + let authors = self.build_author_index(&keys); let large = LargeBucket { keys: BlockStore::from_sorted(keys), authors, @@ -776,7 +998,7 @@ impl EdgeStore { let mut rev = self.reverse.entry_sync(id).or_default(); rev.get_mut().push(ReverseEntry { key_id, - sort_micros, + sort_micros: SortMicros::new(sort_micros), }); } } @@ -799,7 +1021,7 @@ impl EdgeStore { sort_micros, }| { self.forward.update_sync(&key_id, |_, sources| { - sources.remove(&BucketKey::new(sort_micros, id), author); + sources.remove(&BucketKey::new(sort_micros.get(), id), author); }); self.forward .remove_if_sync(&key_id, |sources| sources.is_empty()); @@ -838,6 +1060,39 @@ impl EdgeStore { Some(format!("at://{did}/{collection}/{rkey}")) } + fn edge_item(&self, key: BucketKey) -> Option { + let spur = key.source.to_spur()?; + let stored = self.source_interner.try_resolve(&spur)?; + let uri = AtUri::new_owned(self.decode_source(stored)?).ok()?; + Some(EdgeItem { + uri, + sort_micros: key.micros.get(), + sort_source: key.token().source(), + }) + } + + fn edge_page( + &self, + total: usize, + keys: impl Iterator, + limit: PageLimit, + ) -> EdgePage { + let limit = limit.get() as usize; + let entries = keys.take(limit + 1).collect::>(); + let has_more = entries.len() > limit; + let page = &entries[..entries.len().min(limit)]; + let items = page.iter().filter_map(|&key| self.edge_item(key)).collect(); + let next = has_more + .then(|| page.last().copied()) + .flatten() + .map(BucketKey::token); + EdgePage { + items, + next, + total: Some(TotalCount::from(total)), + } + } + fn author_of_stored(&self, stored: &str) -> Option { match stored.strip_prefix("at://") { Some(rest) => { @@ -869,16 +1124,15 @@ impl EdgeStore { .len() as u64 } - fn build_author_map(&self, keys: &[BucketKey]) -> HashMap { - keys.iter().fold( - HashMap::with_hasher(self.hasher.clone()), - |mut authors, key| { - if let Some(a) = self.author_of(key.source) { - bump_author(&mut authors, a); - } - authors - }, - ) + fn build_author_index(&self, keys: &[BucketKey]) -> AuthorIndex { + debug_assert!(keys.windows(2).all(|pair| pair[0] < pair[1])); + let mut authors = AuthorIndex::new(self.hasher.clone()); + keys.iter().for_each(|&key| { + if let Some(author) = self.author_of(key.source) { + authors.insert(author, key); + } + }); + authors } fn distinct_authors(&self, sources: &Sources) -> u64 { match sources { @@ -890,7 +1144,7 @@ impl EdgeStore { fn distinct_authors_since(&self, sources: &Sources, since_micros: u64) -> u64 { sources .directed(PageCursor::Start, SortDir::Asc) - .filter(|key| key.micros.0 >= since_micros) + .filter(|key| key.micros.get() >= since_micros) .filter_map(|key| self.author_of(key.source)) .collect::>() .len() as u64 @@ -962,10 +1216,7 @@ impl EdgeStore { self.lookup_key(key) .and_then(|id| { self.forward.read_sync(&id, |_, sources| match sources { - Sources::Large(big) => big - .authors - .get(&author) - .map_or(0, |count| count.get() as u64), + Sources::Large(big) => big.authors.count(&author) as u64, Sources::Small(keys) => keys .iter() .filter(|key| self.author_of(key.source) == Some(author)) @@ -1011,7 +1262,7 @@ impl EdgeStore { key: entry.key_id, kind, author, - micros: entry.sort_micros, + micros: entry.sort_micros.get(), by_bucket: edge_kinds[1..].contains(&key.kind.as_ref()), }) }) @@ -1118,15 +1369,7 @@ impl EdgeStore { let items = page .keys .iter() - .filter_map(|&k| { - let spur = k.source.to_spur()?; - let stored = self.source_interner.try_resolve(&spur)?; - let uri = AtUri::new_owned(self.decode_source(stored)?).ok()?; - Some(EdgeItem { - uri, - sort_micros: k.micros.0, - }) - }) + .filter_map(|&key| self.edge_item(key)) .collect(); EdgePage { items, @@ -1187,7 +1430,7 @@ impl EdgeStore { .read_sync(&key_id, |_, sources| { // large buckets track authors, so a missing author fast-fails the scan if let Sources::Large(big) = sources - && !big.authors.contains_key(&author_id) + && !big.authors.contains(&author_id) { return None; } @@ -1228,39 +1471,78 @@ impl EdgeStore { dir: SortDir, ) -> EdgePage { let start = start.into(); - let limit_usize = limit.get() as usize; self.lookup_key(key) .and_then(|id| { self.forward.read_sync(&id, |_, sources| { - let total = Some(TotalCount::from(sources.len())); + let total = sources.len(); let Some(cursor) = sources.resolve_start(start, dir) else { return EdgePage { items: Vec::new(), next: None, - total, + total: Some(TotalCount::from(total)), }; }; - let iter = sources.directed(cursor, dir); - let entries: Vec = iter.take(limit_usize + 1).collect(); - let has_more = entries.len() > limit_usize; - let page = &entries[..entries.len().min(limit_usize)]; - let items = page - .iter() - .filter_map(|&key| { - let spur = key.source.to_spur()?; - let stored = self.source_interner.try_resolve(&spur)?; - let uri = AtUri::new_owned(self.decode_source(stored)?).ok()?; - Some(EdgeItem { - uri, - sort_micros: key.micros.0, + self.edge_page(total, sources.directed(cursor, dir), limit) + }) + }) + .unwrap_or(EdgePage { + items: Vec::new(), + next: None, + total: None, + }) + } + + pub fn list_distinct_authors( + &self, + key: &EdgeKey, + start: impl Into, + limit: PageLimit, + dir: SortDir, + ) -> EdgePage { + let start = start.into(); + self.lookup_key(key) + .and_then(|id| { + self.forward.read_sync(&id, |_, sources| match sources { + Sources::Small(keys) => { + let mut seen = HashSet::with_hasher(self.hasher.clone()); + let mut representatives = keys + .iter() + .rev() + .filter(|key| { + self.author_of(key.source) + .is_some_and(|author| seen.insert(author)) }) - }) - .collect(); - let next = has_more - .then(|| page.last().copied()) - .flatten() - .map(BucketKey::token); - EdgePage { items, next, total } + .copied() + .collect::>(); + representatives.reverse(); + let representatives = Sources::Small(representatives); + let total = representatives.len(); + let Some(cursor) = representatives.resolve_start(start, dir) else { + return EdgePage { + items: Vec::new(), + next: None, + total: Some(TotalCount::from(total)), + }; + }; + self.edge_page(total, representatives.directed(cursor, dir), limit) + } + Sources::Large(big) => { + let total = big.authors.representatives.len(); + debug_assert_eq!(total, big.authors.len()); + let Some(cursor) = big.authors.representatives.resolve_start(start, dir) + else { + return EdgePage { + items: Vec::new(), + next: None, + total: Some(TotalCount::from(total)), + }; + }; + self.edge_page( + total, + big.authors.representatives.directed(cursor, dir), + limit, + ) + } }) }) .unwrap_or(EdgePage { @@ -1304,18 +1586,7 @@ impl EdgeStore { let has_more = top.len() > limit.get() as usize; let page = &top[..top.len().min(limit.get() as usize)]; - let items = page - .iter() - .filter_map(|&k| { - let spur = k.source.to_spur()?; - let stored = self.source_interner.try_resolve(&spur)?; - let uri = AtUri::new_owned(self.decode_source(stored)?).ok()?; - Some(EdgeItem { - uri, - sort_micros: k.micros.0, - }) - }) - .collect(); + let items = page.iter().filter_map(|&key| self.edge_item(key)).collect(); let next = has_more .then(|| page.last().copied()) .flatten() @@ -1620,6 +1891,40 @@ mod tests { .collect() } + fn paginate_distinct_authors( + store: &EdgeStore, + key: &EdgeKey, + dir: SortDir, + page: PageLimit, + ) -> Vec { + std::iter::successors( + Some(store.list_distinct_authors(key, PageCursor::Start, page, dir)), + |prev| { + prev.next.map(|token| { + store.list_distinct_authors(key, PageCursor::After(token), page, dir) + }) + }, + ) + .flat_map(|page| page.items.into_iter().map(|item| item.as_ref().to_owned())) + .collect() + } + + fn newest_author_reference(rows: impl IntoIterator) -> Vec { + let mut authors = HashSet::new(); + let mut newest = rows + .into_iter() + .collect::>() + .into_iter() + .rev() + .filter(|uri| { + let author = split_record_uri(uri).unwrap().0; + authors.insert(author.to_owned()) + }) + .collect::>(); + newest.reverse(); + newest + } + #[test] fn pagination_matches_reference_across_small_and_large() { [64usize, 1000].into_iter().for_each(|n| { @@ -1672,6 +1977,147 @@ mod tests { .collect() } + #[test] + fn distinct_author_pagination_matches_reference_across_small_and_large() { + [64usize, 1000].into_iter().for_each(|n| { + let store = store(); + let key = fill_subject(&store, n); + let asc = newest_author_reference(reference(0..n)); + let desc = asc.iter().rev().cloned().collect::>(); + + [1u32, 2, 5].into_iter().for_each(|page| { + assert_eq!( + paginate_distinct_authors(&store, &key, SortDir::Asc, limit(page)), + asc, + "asc distinct mismatch n={n} page={page}", + ); + assert_eq!( + paginate_distinct_authors(&store, &key, SortDir::Desc, limit(page)), + desc, + "desc distinct mismatch n={n} page={page}", + ); + }); + + let offset = + store.list_distinct_authors(&key, PageOffset::new(2), limit(2), SortDir::Asc); + assert_eq!( + offset + .items + .iter() + .map(|item| item.as_ref().to_owned()) + .collect::>(), + asc.iter().skip(2).take(2).cloned().collect::>(), + ); + assert_eq!( + offset.total, + Some(TotalCount::from(NAMES.len())), + "distinct total n={n}", + ); + }); + } + + #[test] + fn large_distinct_index_tracks_newest_representative() { + let store = store(); + let subject = did("did:plc:large-subject"); + let key = EdgeKey::new( + nsid("sh.tangled.feed.star"), + SubjectRef::Did(subject.clone()), + ); + (0..255).for_each(|index| { + store.add(star_edge_at( + at(&format!( + "at://did:plc:author{index}/sh.tangled.feed.star/r{index}" + )), + subject.clone(), + index as u64, + )); + }); + let older = at("at://did:plc:repeat/sh.tangled.feed.star/older"); + let newest = at("at://did:plc:repeat/sh.tangled.feed.star/newest"); + store.add(star_edge_at(older.clone(), subject.clone(), 1_000)); + store.add(star_edge_at(newest.clone(), subject.clone(), 2_000)); + + let listed = || { + store + .list_distinct_authors( + &key, + PageCursor::Start, + limit(PageLimit::MAX), + SortDir::Desc, + ) + .items + .into_iter() + .map(|item| item.uri) + .collect::>() + }; + assert_eq!(store.count(&key), 257); + assert_eq!(store.count_distinct_authors(&key), 256); + assert!(listed().contains(&newest)); + assert!(!listed().contains(&older)); + + store.remove_source(&newest); + assert!(listed().contains(&older)); + + store.add(star_edge_at(newest.clone(), subject, 3_000)); + store.remove_source(&older); + assert!(listed().contains(&newest)); + assert_eq!(store.count_distinct_authors(&key), 256); + } + + #[test] + fn author_history_promotes_and_shrinks_without_losing_the_newest() { + let store = store(); + let subject = did("did:plc:large-subject"); + let key = EdgeKey::new( + nsid("sh.tangled.feed.star"), + SubjectRef::Did(subject.clone()), + ); + (0..255).for_each(|index| { + store.add(star_edge_at( + at(&format!( + "at://did:plc:author{index}/sh.tangled.feed.star/r{index}" + )), + subject.clone(), + index as u64, + )); + }); + let repeated = (0..6) + .map(|index| { + at(&format!( + "at://did:plc:repeat/sh.tangled.feed.star/r{index}" + )) + }) + .collect::>(); + repeated.iter().enumerate().for_each(|(index, source)| { + store.add(star_edge_at( + source.clone(), + subject.clone(), + 1_000 + index as u64, + )); + }); + let newest = || { + store + .list_distinct_authors(&key, PageCursor::Start, limit(1), SortDir::Desc) + .items[0] + .uri + .clone() + }; + + repeated + .iter() + .enumerate() + .rev() + .for_each(|(index, source)| { + assert_eq!(newest(), *source); + store.remove_source(source); + assert_eq!( + store.count_distinct_authors(&key), + 255 + u64::from(index != 0), + ); + }); + } + #[test] fn offset_pagination_matches_reference_across_small_and_large() { [64usize, 1000].into_iter().for_each(|n| { @@ -1739,6 +2185,7 @@ mod tests { }) .flat_map(|p| p.items.into_iter().map(|u| u.as_ref().to_owned())) .collect(); + assert_eq!(got, asc.into_iter().skip(600).collect::>()); } #[test] @@ -1759,6 +2206,7 @@ mod tests { }) .flat_map(|p| p.items.into_iter().map(|u| u.as_ref().to_owned())) .collect(); + assert_eq!(got, desc.into_iter().skip(2).collect::>()); } #[test] @@ -1780,6 +2228,7 @@ mod tests { }) .flat_map(|p| p.items.into_iter().map(|u| u.as_ref().to_owned())) .collect(); + assert_eq!(got, desc.into_iter().skip(600).collect::>()); } #[test] @@ -2735,6 +3184,54 @@ mod tests { assert_eq!(store.count_distinct_authors(&key), 1, "user1 fully gone"); } + #[test] + fn packed_sort_micros_preserves_u64_ordering_and_cursor_values() { + let values = [ + 0, + u32::MAX as u64, + u32::MAX as u64 + 1, + 1_730_000_000_000_000, + u64::MAX, + ]; + let keys = values.map(|micros| BucketKey::new(micros, SourceId::from_raw(7))); + assert!(keys.windows(2).all(|pair| pair[0] < pair[1])); + keys.into_iter().zip(values).for_each(|(key, micros)| { + assert_eq!(key.micros.get(), micros); + assert_eq!(key.token(), PageToken::new(micros, 7)); + }); + } + + #[test] + fn author_index_separates_and_rejoins_history() { + let author = AuthorId::from_raw(1); + let older = BucketKey::new(10, SourceId::from_raw(10)); + let middle = BucketKey::new(20, SourceId::from_raw(20)); + let newest = BucketKey::new(30, SourceId::from_raw(30)); + let mut index = AuthorIndex::new(RuntimeHasher::default()); + + index.insert(author, newest); + assert_eq!(index.newest[&author], newest); + assert!(index.histories.is_empty()); + + index.insert(author, older); + index.insert(author, middle); + assert_eq!(index.newest[&author], newest); + assert_eq!(index.histories[&author].len(), 3); + + index.remove(author, &older); + assert_eq!(index.histories[&author].len(), 2); + assert_eq!(index.newest[&author], newest); + + index.remove(author, &newest); + assert!(index.histories.is_empty()); + assert_eq!(index.newest[&author], middle); + assert_eq!(index.representatives.last(), Some(middle)); + + index.remove(author, &middle); + assert!(index.newest.is_empty()); + assert!(index.representatives.last().is_none()); + } + #[test] fn page_limit_rejects_zero_and_oversize() { assert!(matches!( diff --git a/bobbin/crates/edge-index/src/sorted_blocks.rs b/bobbin/crates/edge-index/src/sorted_blocks.rs index bd41dc732..552e35390 100644 --- a/bobbin/crates/edge-index/src/sorted_blocks.rs +++ b/bobbin/crates/edge-index/src/sorted_blocks.rs @@ -18,12 +18,9 @@ struct Block { } impl Block { - fn new(mut keys: Vec) -> Self { + fn new(keys: Vec) -> Self { debug_assert!(!keys.is_empty()); debug_assert!(keys.windows(2).all(|w| w[0] <= w[1])); - if keys.capacity() < BLOCK_SPLIT_AT { - keys.reserve(BLOCK_SPLIT_AT - keys.len()); - } let min = keys[0]; let max = *keys.last().unwrap(); Self { keys, min, max } @@ -67,29 +64,23 @@ impl SortedBlocks { self.len } + pub fn last(&self) -> Option { + self.blocks.last().map(|block| block.max) + } + fn rebuild_fenwick(&mut self) { self.fenwick = FenwickTree::from_iter(self.blocks.iter().map(|b| b.len() as u32)); } - // binary search fenwick prefix sums to map 0-indexed rank onto (block, offset) pub fn select(&self, rank: PageRank) -> Option { let rank = rank.get(); if rank >= self.len { return None; } - let mut low = 0; - let mut high = self.blocks.len(); - while low < high { - let mid = low + (high - low) / 2; - if (self.fenwick.prefix_sum(mid + 1, 0) as usize) <= rank { - low = mid + 1; - } else { - high = mid; - } - } - let bi = low; - let prev_sum = self.fenwick.prefix_sum(bi, 0) as usize; - Some(self.blocks[bi].keys[rank - prev_sum]) + // ftree is 1-indexed, so rank shifts up and the remainder shifts back + let position = u32::try_from(rank.checked_add(1)?).ok()?; + let (bi, offset) = self.fenwick.index_of_with_remainder(position); + Some(self.blocks[bi].keys[offset as usize - 1]) } // converts offset to exclusive cursor so directed() handles iteration @@ -122,6 +113,14 @@ impl SortedBlocks { }; // fast path for append-heavy ingestion avoids block binary search if key > last_block.max { + if last_block.len() >= BLOCK_TARGET { + let mut keys = Vec::with_capacity(BLOCK_TARGET); + keys.push(key); + self.blocks.push(Block::new(keys)); + self.fenwick.push(1); + self.len += 1; + return true; + } last_block.keys.push(key); last_block.max = key; let bi = self.blocks.len() - 1; @@ -142,16 +141,22 @@ impl SortedBlocks { } } - // split block in half when split limit reached to bound insert copy costs fn after_insert(&mut self, bi: usize) { - if self.blocks[bi].len() <= BLOCK_SPLIT_AT { - self.fenwick.add_at(bi, 1); + self.fenwick.add_at(bi, 1); + if self.blocks[bi].len() < BLOCK_SPLIT_AT { return; } let tail_keys = self.blocks[bi].keys.split_off(BLOCK_TARGET); + let tail_len = tail_keys.len() as u32; self.blocks[bi].update_bounds(); - self.blocks.insert(bi + 1, Block::new(tail_keys)); - self.rebuild_fenwick(); + if bi + 1 == self.blocks.len() { + self.blocks.push(Block::new(tail_keys)); + self.fenwick.sub_at(bi, tail_len); + self.fenwick.push(tail_len); + } else { + self.blocks.insert(bi + 1, Block::new(tail_keys)); + self.rebuild_fenwick(); + } } pub fn remove(&mut self, key: &K) -> bool { @@ -165,8 +170,14 @@ impl SortedBlocks { self.blocks[bi].keys.remove(pos); self.len -= 1; if self.blocks[bi].keys.is_empty() { - self.blocks.remove(bi); - self.rebuild_fenwick(); + if bi + 1 == self.blocks.len() { + let _ = self.blocks.pop(); + let popped = self.fenwick.pop(); + debug_assert!(popped); + } else { + self.blocks.remove(bi); + self.rebuild_fenwick(); + } } else { if pos == 0 || pos + 1 == old_len { self.blocks[bi].update_bounds(); @@ -190,12 +201,12 @@ impl SortedBlocks { fn after_remove(&mut self, bi: usize) { debug_assert!(!self.blocks[bi].keys.is_empty()); if self.blocks[bi].len() < BLOCK_MERGE_AT { - if bi > 0 && self.blocks[bi - 1].len() + self.blocks[bi].len() <= BLOCK_SPLIT_AT { + if bi > 0 && self.blocks[bi - 1].len() + self.blocks[bi].len() < BLOCK_SPLIT_AT { self.merge_blocks(bi - 1, bi); return; } if bi + 1 < self.blocks.len() - && self.blocks[bi].len() + self.blocks[bi + 1].len() <= BLOCK_SPLIT_AT + && self.blocks[bi].len() + self.blocks[bi + 1].len() < BLOCK_SPLIT_AT { self.merge_blocks(bi, bi + 1); return; @@ -266,6 +277,18 @@ impl SortedBlocks { } } + pub fn replace_last(&mut self, old: K, new: K) -> bool { + let Some(last) = self.blocks.last_mut() else { + return false; + }; + if last.max != old || new <= old { + return false; + } + *last.keys.last_mut().expect("nonempty block") = new; + last.update_bounds(); + true + } + pub fn heap_bytes(&self) -> u64 { let block_overhead = std::mem::size_of::>() as u64; self.blocks.capacity() as u64 * block_overhead @@ -277,3 +300,184 @@ impl SortedBlocks { + (self.fenwick.len() * std::mem::size_of::()) as u64 } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::PageOffset; + use std::collections::BTreeSet; + + #[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)] + struct TestKey { + micros: u64, + source: u32, + } + + impl PageKey for TestKey { + fn token(self) -> PageToken { + PageToken::new(self.micros, self.source) + } + fn from_token(token: PageToken) -> Self { + Self { + micros: token.micros(), + source: token.source(), + } + } + } + + fn key(value: u32) -> TestKey { + TestKey { + micros: u64::from(value / 3), + source: value, + } + } + + fn check(store: &SortedBlocks, expected: &BTreeSet) { + let asc = expected.iter().copied().collect::>(); + assert_eq!(store.len(), asc.len()); + assert_eq!( + store + .directed(PageCursor::Start, SortDir::Asc) + .collect::>(), + asc, + ); + assert_eq!( + store + .directed(PageCursor::Start, SortDir::Desc) + .collect::>(), + asc.iter().rev().copied().collect::>(), + ); + for (rank, expected) in asc.iter().enumerate() { + assert_eq!(store.select(PageRank::new(rank)), Some(*expected)); + } + assert_eq!(store.select(PageRank::new(asc.len())), None); + for dir in [SortDir::Asc, SortDir::Desc] { + let mut expected = asc.clone(); + if dir == SortDir::Desc { + expected.reverse(); + } + for offset in [0, 1, 63, 255, 256, 257, 511, 512, asc.len()] { + let got = store + .resolve_start(PageOffset::new(offset as u64).into(), dir) + .map(|cursor| store.directed(cursor, dir).take(7).collect::>()) + .unwrap_or_default(); + assert_eq!( + got, + expected + .iter() + .skip(offset) + .take(7) + .copied() + .collect::>() + ); + } + } + let total = store.fenwick.prefix_sum(store.blocks.len(), 0) as usize; + assert_eq!(total, store.len()); + for block in &store.blocks { + assert!(!block.keys.is_empty()); + assert!(block.len() < BLOCK_SPLIT_AT); + assert_eq!(block.min, block.keys[0]); + assert_eq!(block.max, *block.keys.last().unwrap()); + } + } + + #[test] + fn singleton_does_not_reserve_a_full_block() { + let mut store = SortedBlocks::from_sorted(Vec::new()); + assert!(store.insert(key(0))); + assert!(store.blocks[0].keys.capacity() < BLOCK_TARGET); + } + + #[test] + fn monotone_append_keeps_completed_blocks_dense() { + let mut store = SortedBlocks::from_sorted(Vec::new()); + let mut expected = BTreeSet::new(); + for value in 0..(BLOCK_TARGET as u32 * 4 + 7) { + assert!(store.insert(key(value))); + expected.insert(key(value)); + } + assert_eq!(store.blocks.len(), 5); + for block in &store.blocks[..4] { + assert_eq!(block.len(), BLOCK_TARGET); + assert_eq!(block.keys.capacity(), BLOCK_TARGET); + } + check(&store, &expected); + } + + #[test] + fn mixed_mutations_preserve_rank_cursor_and_merge_invariants() { + let mut store = SortedBlocks::from_sorted(Vec::new()); + let mut expected = BTreeSet::new(); + for value in 0..1500 { + assert!(store.insert(key(value))); + expected.insert(key(value)); + } + let mut rng = 62_446_u64; + for step in 0..12_000 { + rng = rng + .wrapping_mul(6364136223846793005) + .wrapping_add(1442695040888963407); + let candidate = key(((rng >> 32) % 2400) as u32); + if rng & 2 == 0 { + assert_eq!(store.insert(candidate), expected.insert(candidate)); + } else { + assert_eq!(store.remove(&candidate), expected.remove(&candidate)); + } + if step % 251 == 0 { + check(&store, &expected); + } + } + check(&store, &expected); + let remaining = expected.iter().copied().collect::>(); + for candidate in remaining.into_iter().rev() { + assert!(store.remove(&candidate)); + expected.remove(&candidate); + if expected.len() % 127 == 0 { + check(&store, &expected); + } + } + check(&store, &expected); + } + + #[test] + fn last_and_singleton_tail_removal_preserve_ranks() { + let mut store = SortedBlocks::from_sorted(Vec::new()); + assert_eq!(store.last(), None); + assert!(store.insert(key(0))); + assert_eq!(store.last(), Some(key(0))); + assert!(store.remove(&key(0))); + assert_eq!(store.last(), None); + assert_eq!(store.fenwick.len(), 0); + + for blocks in [1u32, 2, 4, 8, 16] { + let count = BLOCK_TARGET as u32 * blocks + 1; + let mut expected = (0..count).map(key).collect::>(); + let mut store = SortedBlocks::from_sorted(expected.iter().copied().collect()); + for _ in 0..8 { + assert!(store.remove(&key(count - 1))); + expected.remove(&key(count - 1)); + assert_eq!(store.last(), Some(key(count - 2))); + assert_eq!(store.fenwick.len(), blocks as usize); + check(&store, &expected); + assert!(store.insert(key(count - 1))); + expected.insert(key(count - 1)); + assert_eq!(store.last(), Some(key(count - 1))); + check(&store, &expected); + } + } + } + + #[test] + fn replacing_greatest_key_preserves_rank_and_count() { + let mut store = SortedBlocks::from_sorted((0..1000).map(key).collect()); + assert!(store.replace_last(key(999), key(2000))); + assert_eq!(store.len(), 1000); + assert_eq!(store.select(PageRank::new(999)), Some(key(2000))); + assert!(!store.replace_last(key(20), key(3000))); + assert!(!store.replace_last(key(2000), key(0))); + let mut expected = (0..999).map(key).collect::>(); + expected.insert(key(2000)); + check(&store, &expected); + } +} diff --git a/bobbin/crates/xrpc/src/filter.rs b/bobbin/crates/xrpc/src/filter.rs index 2d64c9dd1..16cd3f102 100644 --- a/bobbin/crates/xrpc/src/filter.rs +++ b/bobbin/crates/xrpc/src/filter.rs @@ -169,6 +169,10 @@ impl VouchNetworkFilter { self.network.contains(did) } + pub fn is_empty(&self) -> bool { + self.network.is_empty() + } + pub fn matches(&self, item: &EdgeItem) -> bool { crate::source_authority_did(&item.uri) .is_some_and(|did| self.network.contains(did.as_ref())) diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index f40be3b50..e4781dd66 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -1878,7 +1878,9 @@ where let owned = owned.clone(); let nsid = nsid.clone(); async move { - let EdgeItem { uri, sort_micros } = item; + let EdgeItem { + uri, sort_micros, .. + } = item; let result = hydrate_record_view::(&owned, &nsid, uri.clone(), sort_micros).await; if matches!(provenance, HitProvenance::Indexed) && let Err(err) = &result @@ -2033,6 +2035,7 @@ where .map(|uri| EdgeItem { uri, sort_micros: 0, + sort_source: 0, }) .collect(); let views = hydrate_record_stream::(state, nsid, items, HitProvenance::ClientSupplied); @@ -2083,16 +2086,14 @@ where Ok((page, nsid)) } -async fn list_records( +fn record_list_response( state: &AppState, - q: TypedListQuery, + page: EdgePage, + nsid: Nsid, ) -> Result where - R: XrpcResp + HasSubject, V: serde::de::DeserializeOwned + Serialize + NormalizeRepoRefs + Send + 'static, - F: ListFilter, { - let (page, nsid) = record_edge_page::(state, &q)?; let permit = state.heavy_permit()?; let views = hydrate_record_stream::(state, nsid, page.items, HitProvenance::Indexed); Ok(json_stream::, _>( @@ -2103,6 +2104,38 @@ where )) } +async fn list_records( + state: &AppState, + q: TypedListQuery, +) -> Result +where + R: XrpcResp + HasSubject, + V: serde::de::DeserializeOwned + Serialize + NormalizeRepoRefs + Send + 'static, + F: ListFilter, +{ + let (page, nsid) = record_edge_page::(state, &q)?; + record_list_response::(state, page, nsid) +} + +async fn list_records_by_distinct_author( + state: &AppState, + q: TypedListQuery, +) -> Result +where + R: XrpcResp + HasSubject, + V: serde::de::DeserializeOwned + Serialize + NormalizeRepoRefs + Send + 'static, +{ + let subject = parse_subject(&q.subject, R::SHAPE)?; + let start = parse_start(q.cursor.as_deref(), q.offset)?; + let limit = parse_limit(q.limit)?; + let nsid = nsid_static(R::NSID); + let key = EdgeKey::new(nsid.clone(), subject); + let page = state + .edges + .list_distinct_authors(&key, start, limit, q.dir()); + record_list_response::(state, page, nsid) +} + fn count_for( state: &AppState, q: CountQuery, @@ -2216,7 +2249,7 @@ async fn list_stars( State(state): State, XrpcQuery(q): XrpcQuery>, ) -> Result { - list_records::, _>(&state, q).await + list_records_by_distinct_author::>(&state, q).await } async fn count_stars( @@ -2243,7 +2276,7 @@ async fn list_follows( State(state): State, XrpcQuery(q): XrpcQuery>, ) -> Result { - list_records::, _>(&state, q).await + list_records_by_distinct_author::>(&state, q).await } async fn count_follows( @@ -2465,66 +2498,267 @@ fn vouch_network(state: &AppState, actor: &SubjectRef) -> std::collections::Hash .collect() } -fn vouches_received( - state: &AppState, - actor: &SubjectRef, - network: &VouchNetworkFilter, -) -> Vec { - let vouch = nsid_static("sh.tangled.graph.vouch"); - state - .edges - .list( - &EdgeKey::new(vouch, actor.clone()), - PageCursor::Start, - PageLimit::new(PageLimit::MAX).unwrap(), - SortDir::Desc, - ) - .items - .into_iter() - .filter(|item| network.matches(item)) - .collect() +#[derive(Clone, Copy)] +enum VouchList { + Received, + Sent, } -fn vouches_sent(state: &AppState, actor: &SubjectRef) -> Vec { - let vouch_by = nsid_static("sh.tangled.graph.vouch.by"); - state - .edges - .list( - &EdgeKey::new(vouch_by, actor.clone()), - PageCursor::Start, - PageLimit::new(PageLimit::MAX).unwrap(), - SortDir::Desc, - ) - .items +fn vouch_page(edges: &EdgeStore, key: &EdgeKey, list: VouchList, cursor: PageCursor) -> EdgePage { + let limit = PageLimit::new(PageLimit::MAX).unwrap(); + match list { + VouchList::Received => edges.list_distinct_authors(key, cursor, limit, SortDir::Desc), + VouchList::Sent => edges.list(key, cursor, limit, SortDir::Desc), + } +} + +fn vouch_items( + edges: &EdgeStore, + key: EdgeKey, + list: VouchList, + start: PageCursor, +) -> impl Iterator + '_ { + vouch_items_from(start, move |cursor| vouch_page(edges, &key, list, cursor)) +} + +fn vouch_items_from( + start: PageCursor, + mut fetch: impl FnMut(PageCursor) -> EdgePage, +) -> impl Iterator { + let mut next = Some(start); + let mut items = Vec::::new().into_iter(); + std::iter::from_fn(move || { + loop { + if let Some(item) = items.next() { + return Some(item); + } + let cursor = next.take()?; + let page = fetch(cursor); + next = page.next.map(PageCursor::After); + items = page.items.into_iter(); + } + }) +} + +fn merge_vouches( + received: impl Iterator, + sent: impl Iterator, +) -> impl Iterator { + let mut received = received.peekable(); + let mut sent = sent.peekable(); + std::iter::from_fn(move || match (received.peek(), sent.peek()) { + (Some(left), Some(right)) => match left.sort_token().cmp(&right.sort_token()) { + std::cmp::Ordering::Greater => received.next(), + std::cmp::Ordering::Less => sent.next(), + std::cmp::Ordering::Equal => { + sent.next(); + received.next() + } + }, + (Some(_), None) => received.next(), + (None, Some(_)) => sent.next(), + (None, None) => None, + }) +} + +fn parse_vouch_cursor(raw: Option<&str>) -> Result { + let Some(raw) = raw else { + return Ok(PageCursor::Start); + }; + PageToken::decode_token(raw) + .or_else(|_| { + raw.parse::() + .map(|micros| PageToken::new(micros, 0)) + .map_err(|_| CursorParseError::Malformed) + }) + .map(PageCursor::After) + .map_err(|error| XrpcError::InvalidParams(format!("cursor: {error}"))) } fn paginate_vouches( candidates: impl Iterator, - cursor: Option<&str>, limit: PageLimit, ) -> (Vec, Option) { - let mut seen = std::collections::HashSet::new(); - let mut combined: Vec = candidates - .filter(|item| seen.insert(item.uri.as_ref().to_string())) - .collect(); - combined.sort_by(|a, b| b.sort_micros.cmp(&a.sort_micros)); - - let start = match cursor.and_then(|s| s.parse::().ok()) { - None => 0, - Some(micros) => combined - .iter() - .position(|item| item.sort_micros < micros) - .unwrap_or(combined.len()), - }; let limit = limit.get() as usize; - let has_more = combined.len() > start + limit; - let items: Vec = combined.into_iter().skip(start).take(limit).collect(); - let next = has_more - .then(|| items.last().map(|item| item.sort_micros.to_string())) - .flatten(); + let mut seen = std::collections::HashSet::new(); + let mut items = candidates + .filter(|item| seen.insert(item.uri.clone())) + .take(limit + 1) + .collect::>(); + let has_more = items.len() > limit; + items.truncate(limit); + let next = has_more.then(|| items.last().unwrap().sort_token().encode_token()); (items, next) } +#[cfg(test)] +mod vouch_pagination_tests { + use std::collections::HashSet; + + use super::*; + use bobbin_runtime::RuntimeHasher; + use bobbin_types::edges::Edge; + + fn did(value: &str) -> Did { + Did::new_owned(value).unwrap() + } + + fn uri(value: String) -> AtUri { + AtUri::new_owned(value).unwrap() + } + + fn edge(kind: &'static str, subject: &Did, author: usize, micros: u64) -> Edge { + Edge { + kind: nsid_static(kind), + subject: SubjectRef::Did(subject.clone()), + source: uri(format!( + "at://did:plc:author{author}/sh.tangled.graph.vouch/r{author}" + )), + sort_micros: micros, + } + } + + #[test] + fn received_filter_reaches_beyond_the_first_candidate_page() { + let store = EdgeStore::new(RuntimeHasher::from_seeds(1, 2, 3, 4)); + let subject = did("did:plc:subject"); + (0..=PageLimit::MAX as usize).for_each(|author| { + store.add(edge( + "sh.tangled.graph.vouch", + &subject, + author, + author as u64, + )); + }); + let key = EdgeKey::new( + nsid_static("sh.tangled.graph.vouch"), + SubjectRef::Did(subject), + ); + let candidates = vouch_items(&store, key, VouchList::Received, PageCursor::Start) + .filter(|item| item.uri.as_ref().starts_with("at://did:plc:author0/")); + let (items, cursor) = paginate_vouches(candidates, PageLimit::new(1).unwrap()); + + assert_eq!(items.len(), 1); + assert!(items[0].uri.as_ref().starts_with("at://did:plc:author0/")); + assert!(cursor.is_none()); + } + + #[test] + fn cursor_walk_keeps_every_timestamp_tie() { + let store = EdgeStore::new(RuntimeHasher::from_seeds(1, 2, 3, 4)); + let subject = did("did:plc:subject"); + store.add(edge("sh.tangled.graph.vouch", &subject, 1, 42)); + store.add(edge("sh.tangled.graph.vouch.by", &subject, 2, 42)); + store.add(edge("sh.tangled.graph.vouch", &subject, 3, 42)); + let received = EdgeKey::new( + nsid_static("sh.tangled.graph.vouch"), + SubjectRef::Did(subject.clone()), + ); + let sent = EdgeKey::new( + nsid_static("sh.tangled.graph.vouch.by"), + SubjectRef::Did(subject), + ); + let mut start = PageCursor::Start; + let mut uris = Vec::new(); + + loop { + let candidates = merge_vouches( + vouch_items(&store, received.clone(), VouchList::Received, start), + vouch_items(&store, sent.clone(), VouchList::Sent, start), + ); + let (items, cursor) = paginate_vouches(candidates, PageLimit::new(1).unwrap()); + uris.extend(items.into_iter().map(|item| item.uri)); + let Some(cursor) = cursor else { + break; + }; + assert_eq!(cursor.len(), 24); + start = parse_vouch_cursor(Some(&cursor)).unwrap(); + } + + assert_eq!(uris.len(), 3); + assert_eq!(uris.iter().collect::>().len(), 3); + } + + #[test] + fn page_fetch_is_deferred_until_items_are_needed() { + let calls = std::cell::Cell::new(0usize); + let item = |source: u32| EdgeItem { + uri: uri(format!( + "at://did:plc:author/sh.tangled.graph.vouch/r{source}" + )), + sort_micros: 42, + sort_source: source, + }; + let mut items = vouch_items_from(PageCursor::Start, |cursor| { + calls.set(calls.get() + 1); + match cursor { + PageCursor::Start => EdgePage { + items: vec![item(3), item(2)], + next: Some(PageToken::new(42, 2)), + total: None, + }, + PageCursor::After(token) => { + assert_eq!(token, PageToken::new(42, 2)); + EdgePage { + items: vec![item(1)], + next: None, + total: None, + } + } + } + }); + assert_eq!(calls.get(), 0); + assert_eq!(items.next().unwrap().sort_source, 3); + assert_eq!(calls.get(), 1); + assert_eq!(items.next().unwrap().sort_source, 2); + assert_eq!(calls.get(), 1); + assert_eq!(items.next().unwrap().sort_source, 1); + assert_eq!(calls.get(), 2); + assert!(items.next().is_none()); + assert!(items.next().is_none()); + assert_eq!(calls.get(), 2); + } + + #[test] + fn deferred_fetch_continues_through_an_empty_page_with_a_cursor() { + let calls = std::cell::Cell::new(0usize); + let mut items = vouch_items_from(PageCursor::Start, |cursor| { + calls.set(calls.get() + 1); + match cursor { + PageCursor::Start => EdgePage { + items: Vec::new(), + next: Some(PageToken::new(42, 2)), + total: None, + }, + PageCursor::After(token) => { + assert_eq!(token, PageToken::new(42, 2)); + EdgePage { + items: vec![EdgeItem { + uri: uri("at://did:plc:author/sh.tangled.graph.vouch/r1".to_owned()), + sort_micros: 42, + sort_source: 1, + }], + next: None, + total: None, + } + } + } + }); + assert_eq!(items.next().unwrap().sort_source, 1); + assert_eq!(calls.get(), 2); + assert!(items.next().is_none()); + assert_eq!(calls.get(), 2); + } + + #[test] + fn decimal_vouch_cursor_keeps_legacy_resume_semantics() { + assert_eq!( + parse_vouch_cursor(Some("42")).unwrap(), + PageCursor::After(PageToken::new(42, 0)), + ); + } +} + async fn list_network_vouches( State(state): State, ExtractOptionalServiceAuth(auth): ExtractOptionalServiceAuth, @@ -2532,6 +2766,8 @@ async fn list_network_vouches( ) -> Result { let actor = SubjectRef::Did(q.actor.clone()); let limit = parse_limit(q.limit.and_then(|n| u32::try_from(n).ok()))?; + let start = parse_vouch_cursor(q.cursor.as_deref())?; + let permit = state.heavy_permit()?; let vouch_nsid = nsid_static("sh.tangled.graph.vouch"); let viewer = auth @@ -2543,18 +2779,31 @@ async fn list_network_vouches( .map(|viewer| vouch_network(&state, viewer)) .unwrap_or_default(), ); - let received = vouches_received(&state, &actor, &network); - - let sent = if network.contains(q.actor.as_ref()) { - vouches_sent(&state, &actor) + let (page_items, next_cursor) = if network.is_empty() { + (Vec::new(), None) } else { - Vec::new() + let received = vouch_items( + &state.edges, + EdgeKey::new(vouch_nsid.clone(), actor.clone()), + VouchList::Received, + start, + ) + .filter(|item| network.matches(item)); + let sent = network + .contains(q.actor.as_ref()) + .then(|| { + vouch_items( + &state.edges, + EdgeKey::new(nsid_static("sh.tangled.graph.vouch.by"), actor), + VouchList::Sent, + start, + ) + }) + .into_iter() + .flatten(); + paginate_vouches(merge_vouches(received, sent), limit) }; - let (page_items, next_cursor) = - paginate_vouches(received.into_iter().chain(sent), q.cursor.as_deref(), limit); - - let permit = state.heavy_permit()?; let views = hydrate_record_stream::>( &state, vouch_nsid, diff --git a/bobbin/crates/xrpc/tests/aggregation.rs b/bobbin/crates/xrpc/tests/aggregation.rs index 9035da0f6..582a958d7 100644 --- a/bobbin/crates/xrpc/tests/aggregation.rs +++ b/bobbin/crates/xrpc/tests/aggregation.rs @@ -394,6 +394,125 @@ async fn count_distinct_authors_dedupes_per_author() { assert_eq!(body["distinctAuthors"], json!(2)); } +#[tokio::test] +async fn list_stars_pages_one_newest_record_per_author() { + let h = Harness::new().await; + let subject_did = did("did:plc:abalone"); + let subject = at(&format!("at://{}", subject_did.as_ref())); + let stars = [ + ("did:plc:nel", "s1"), + ("did:plc:olaren", "s2"), + ("did:plc:nel", "s3"), + ]; + for (author, record_key) in stars { + let author = did(author); + let record_key = rkey(record_key); + h.add_edge( + &nsid("sh.tangled.feed.star"), + &subject, + &at(&format!( + "at://{}/sh.tangled.feed.star/{}", + author.as_ref(), + record_key.as_ref(), + )), + ); + h.mount( + &author, + &nsid("sh.tangled.feed.star"), + &record_key, + star_body(&subject_did), + ) + .await; + } + + let app = router(h.state.clone()); + let (_, first) = json_response( + app.clone() + .oneshot(list_request( + "sh.tangled.feed.listStars", + subject.as_ref(), + &[("limit", "1")], + )) + .await + .unwrap(), + ) + .await; + assert_eq!(first["total"], json!(2)); + assert_eq!(first["items"].as_array().unwrap().len(), 1); + assert!(first["items"][0]["uri"].as_str().unwrap().ends_with("/s3")); + + let cursor = first["cursor"] + .as_str() + .expect("one distinct author remains"); + let (_, second) = json_response( + app.oneshot(list_request( + "sh.tangled.feed.listStars", + subject.as_ref(), + &[("limit", "1"), ("cursor", cursor)], + )) + .await + .unwrap(), + ) + .await; + assert_eq!(second["items"].as_array().unwrap().len(), 1); + assert!(second["items"][0]["uri"].as_str().unwrap().ends_with("/s2")); + assert!(second.get("cursor").is_none()); +} + +#[tokio::test] +async fn list_follows_returns_one_newest_record_per_follower() { + let h = Harness::new().await; + let subject_did = did("did:plc:abalone"); + let subject = at(&format!("at://{}", subject_did.as_ref())); + let follows = [ + ("did:plc:nel", "f1"), + ("did:plc:olaren", "f2"), + ("did:plc:nel", "f3"), + ]; + for (author, record_key) in follows { + let author = did(author); + let record_key = rkey(record_key); + h.add_edge( + &nsid("sh.tangled.graph.follow"), + &subject, + &at(&format!( + "at://{}/sh.tangled.graph.follow/{}", + author.as_ref(), + record_key.as_ref(), + )), + ); + h.mount( + &author, + &nsid("sh.tangled.graph.follow"), + &record_key, + follow_body(&subject_did), + ) + .await; + } + + let (_, body) = json_response( + router(h.state.clone()) + .oneshot(list_request( + "sh.tangled.graph.listFollows", + subject.as_ref(), + &[], + )) + .await + .unwrap(), + ) + .await; + let uris = body["items"] + .as_array() + .unwrap() + .iter() + .map(|item| item["uri"].as_str().unwrap()) + .collect::>(); + assert_eq!(body["total"], json!(2)); + assert_eq!(uris.len(), 2); + assert!(uris.iter().any(|uri| uri.ends_with("/f3"))); + assert!(!uris.iter().any(|uri| uri.ends_with("/f1"))); +} + #[tokio::test] async fn count_forks_counts_repos_pointing_at_the_source() { let h = Harness::new().await; -- 2.51.2