diff --git a/Cargo.lock b/Cargo.lock index 8e85f0592..f2b3bdf44 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -712,10 +712,12 @@ dependencies = [ "bobbin-runtime", "bobbin-types", "either", + "ftree", "itertools 0.14.0", "jacquard-common", "lasso", "scc", + "serde", "smallvec", "thiserror 2.0.18", "tokio", @@ -2294,6 +2296,12 @@ version = "1.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "42703706b716c37f96a77aea830392ad231f44c9e9a67872fa5548707e11b11c" +[[package]] +name = "ftree" +version = "1.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "eeef4a9366aaf0ed0bb5292b7c489d80600a7431fca0d96271e10e23ac0bc2b0" + [[package]] name = "futures" version = "0.3.32" diff --git a/bobbin/crates/edge-index/Cargo.toml b/bobbin/crates/edge-index/Cargo.toml index 8c3e68a6a..bfdff6a6c 100644 --- a/bobbin/crates/edge-index/Cargo.toml +++ b/bobbin/crates/edge-index/Cargo.toml @@ -9,10 +9,12 @@ rust-version.workspace = true bobbin-runtime = { workspace = true } bobbin-types = { workspace = true } either = { workspace = true } +ftree = "1.3.0" itertools = { workspace = true } jacquard-common = { workspace = true } lasso = { workspace = true } scc = { workspace = true } smallvec = "1" +serde = { workspace = true } thiserror = { workspace = true } tokio = { workspace = true } diff --git a/bobbin/crates/edge-index/src/lib.rs b/bobbin/crates/edge-index/src/lib.rs index 3370e9d4b..91035cec1 100644 --- a/bobbin/crates/edge-index/src/lib.rs +++ b/bobbin/crates/edge-index/src/lib.rs @@ -1,28 +1,8 @@ -use std::collections::{BTreeSet, HashMap}; +use std::collections::HashMap; use std::marker::PhantomData; use std::num::NonZeroU32; -use std::ops::{Bound, ControlFlow}; use std::sync::atomic::{AtomicU32, Ordering}; -use std::sync::{Arc, Mutex, RwLock}; - -const FILTER_SCAN_MULTIPLIER: usize = 64; -const FILTER_SCAN_FLOOR: usize = 512; - -struct ScanState { - matched: Vec<(BucketKey, AtUri)>, - last_scanned: Option, - scanned: usize, -} - -impl ScanState { - fn with_capacity(cap: usize) -> Self { - Self { - matched: Vec::with_capacity(cap), - last_scanned: None, - scanned: 0, - } - } -} +use std::sync::{Arc, Mutex}; use bobbin_runtime::RuntimeHasher; use bobbin_types::edges::{Edge, Record}; @@ -41,16 +21,16 @@ use smallvec::SmallVec; use thiserror::Error; #[derive(Clone, Copy, Debug, Eq, Hash, PartialEq, Ord, PartialOrd)] -struct SortMicros(u64); +pub(crate) struct SortMicros(u64); #[derive(Clone, Copy, Debug, Eq, Hash, PartialEq, Ord, PartialOrd)] -struct BucketKey { +pub(crate) struct BucketKey { micros: SortMicros, source: SourceId, } impl BucketKey { - fn new(micros: u64, source: SourceId) -> Self { + pub(crate) fn new(micros: u64, source: SourceId) -> Self { Self { micros: SortMicros(micros), source, @@ -70,11 +50,15 @@ impl BucketKey { } pub mod coverage; +mod sorted_blocks; pub mod state_index; +mod state_view; pub use coverage::{Coverage, CoverageWatch, HydrantCursor, PromotionSignal}; +use sorted_blocks::{PageKey, SortedBlocks}; pub use state_index::{ ApplyOutcome, IssueStateKind, PullStatusKind, StateIndex, StateKind, apply_record_state, }; +use state_view::{ProjectedState, StateViewIndex}; #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] pub struct SourceTag; @@ -105,7 +89,7 @@ impl Interned { } pub type SourceId = Interned; -type AuthorId = Interned; +pub(crate) type AuthorId = Interned; type CollectionId = Interned; #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] @@ -174,6 +158,115 @@ pub enum PageCursor { After(PageToken), } +#[derive( + Clone, Copy, Debug, Default, Eq, Hash, Ord, PartialEq, PartialOrd, serde::Deserialize, serde::Serialize, +)] +#[serde(transparent)] +pub struct PageOffset(u64); + +impl PageOffset { + pub const ZERO: Self = Self(0); + + pub const fn new(value: u64) -> Self { + Self(value) + } + + pub const fn get(self) -> u64 { + self.0 + } + + pub const fn as_usize(self) -> usize { + self.0 as usize + } +} + +impl From for PageOffset { + fn from(v: u64) -> Self { + Self(v) + } +} + +impl From for PageOffset { + fn from(v: usize) -> Self { + Self(v as u64) + } +} + +#[derive(Clone, Copy, Debug, Default, Eq, Hash, Ord, PartialEq, PartialOrd)] +pub struct PageRank(usize); + +impl PageRank { + pub const ZERO: Self = Self(0); + + pub const fn new(value: usize) -> Self { + Self(value) + } + + pub const fn get(self) -> usize { + self.0 + } +} + +impl From for PageRank { + fn from(v: usize) -> Self { + Self(v) + } +} + +#[derive( + Clone, Copy, Debug, Default, Eq, Hash, Ord, PartialEq, PartialOrd, serde::Deserialize, serde::Serialize, +)] +#[serde(transparent)] +pub struct TotalCount(u64); + +impl TotalCount { + pub const ZERO: Self = Self(0); + + pub const fn new(value: u64) -> Self { + Self(value) + } + + pub const fn get(self) -> u64 { + self.0 + } +} + +impl From for TotalCount { + fn from(v: u64) -> Self { + Self(v) + } +} + +impl From for TotalCount { + fn from(v: usize) -> Self { + Self(v as u64) + } +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum PageStart { + Cursor(PageCursor), + Offset(PageOffset), +} + +impl From for PageStart { + fn from(c: PageCursor) -> Self { + Self::Cursor(c) + } +} + +impl From for PageStart { + fn from(o: PageOffset) -> Self { + Self::Offset(o) + } +} + +impl From for PageStart { + fn from(v: u64) -> Self { + Self::Offset(PageOffset::new(v)) + } +} + #[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] pub enum SortDir { Asc, @@ -233,6 +326,7 @@ impl PageLimit { pub struct EdgePage { pub items: Vec, pub next: Option, + pub total: Option, } #[derive(Clone, Debug, Eq, PartialEq)] @@ -343,16 +437,22 @@ impl EdgeMemReport { } #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] -struct EdgeKeyId(u32); +pub(crate) struct EdgeKeyId(u32); -fn bump_author(authors: &mut HashMap, author: AuthorId) { +pub(crate) fn bump_author( + authors: &mut HashMap, + author: AuthorId, +) { authors .entry(author) .and_modify(|c| *c = c.saturating_add(1)) .or_insert(NonZeroU32::MIN); } -fn drop_author(authors: &mut HashMap, author: AuthorId) { +pub(crate) fn drop_author( + authors: &mut HashMap, + author: AuthorId, +) { match authors.get(&author).map(|c| c.get() - 1) { Some(0) | None => { authors.remove(&author); @@ -364,13 +464,25 @@ fn drop_author(authors: &mut HashMap, autho } struct LargeBucket { - keys: BTreeSet, + keys: BlockStore, authors: HashMap, } const BUCKET_PROMOTE_AT: usize = 256; const SOURCE_BYTES: u64 = 16; +type BlockStore = SortedBlocks; + +impl PageKey for BucketKey { + fn token(self) -> PageToken { + self.token() + } + + fn from_token(tok: PageToken) -> Self { + Self::from_token(tok) + } +} + enum Sources { Small(SmallVec<[BucketKey; 2]>), Large(Box), @@ -437,12 +549,38 @@ impl Sources { ) -> Box + '_> { match self { Self::Small(v) => Box::new(directed_slice(v, cursor, dir)), - Self::Large(big) => Box::new(directed_tree(&big.keys, cursor, dir)), + Self::Large(big) => big.keys.directed(cursor, dir), + } + } + + fn select(&self, rank: PageRank) -> Option { + match self { + Self::Small(v) => v.get(rank.get()).copied(), + Self::Large(big) => big.keys.select(rank), + } + } + + fn resolve_start(&self, start: PageStart, dir: SortDir) -> Option { + match start { + PageStart::Cursor(c) => Some(c), + PageStart::Offset(o) if o.get() == 0 => Some(PageCursor::Start), + PageStart::Offset(o) => { + let len = self.len() as u64; + let o = o.get(); + if o >= len { + return None; + } + let pred_rank = match dir { + SortDir::Asc => o - 1, + SortDir::Desc => len - o, + }; + self.select(PageRank::new(pred_rank as usize)) + .map(|k| PageCursor::After(k.token())) + } } } fn heap_bytes(&self) -> u64 { - const BTREE_BYTES_PER_KEY: u64 = 32; const HASHMAP_FIXED: u64 = 48; const HASHMAP_PER_CAP: u64 = 9; match self { @@ -450,7 +588,7 @@ impl Sources { Self::Small(_) => 0, Self::Large(big) => { std::mem::size_of::() as u64 - + big.keys.len() as u64 * BTREE_BYTES_PER_KEY + + big.keys.heap_bytes() + HASHMAP_FIXED + big.authors.capacity() as u64 * HASHMAP_PER_CAP } @@ -464,68 +602,6 @@ struct ReverseEntry { sort_micros: u64, } -#[derive(Clone, Copy)] -struct ProjectedState { - key: EdgeKeyId, - kind: K, - author: Option, -} - -struct CountBucket { - count: u64, - authors: HashMap, -} - -struct StateCountInner { - projected: HashMap; 1]>, RuntimeHasher>, - buckets: HashMap<(EdgeKeyId, K), CountBucket, RuntimeHasher>, -} - -struct StateCountIndex { - inner: RwLock>, - hasher: RuntimeHasher, -} - -impl StateCountIndex { - fn new(hasher: RuntimeHasher) -> Self { - Self { - inner: RwLock::new(StateCountInner { - projected: HashMap::with_hasher(hasher.clone()), - buckets: HashMap::with_hasher(hasher.clone()), - }), - hasher, - } - } - - fn heap_bytes(&self) -> u64 { - let inner = self - .inner - .read() - .expect("state-count index rwlock poisoned"); - let projected = inner.projected.capacity() - * (std::mem::size_of::() - + std::mem::size_of::; 1]>>() - + 1) - + inner - .projected - .values() - .filter(|states| states.spilled()) - .map(|states| states.capacity() * std::mem::size_of::>()) - .sum::(); - let buckets = inner.buckets.capacity() - * (std::mem::size_of::<(EdgeKeyId, K)>() + std::mem::size_of::() + 1); - let author_slots = inner - .buckets - .values() - .map(|bucket| { - bucket.authors.capacity() - * (std::mem::size_of::() + std::mem::size_of::() + 1) - }) - .sum::(); - (projected + buckets + author_slots) as u64 - } -} - pub struct EdgeStore { source_interner: Arc>, did_interner: Arc>, @@ -536,8 +612,8 @@ pub struct EdgeStore { forward: SccMap, reverse: SccMap, RuntimeHasher>, source_collections: SccMap, RuntimeHasher>, - issue_counts: StateCountIndex, - pull_counts: StateCountIndex, + issue_counts: StateViewIndex, + pull_counts: StateViewIndex, hasher: RuntimeHasher, writer: Mutex<()>, } @@ -554,8 +630,8 @@ impl EdgeStore { forward: SccMap::with_hasher(hasher.clone()), reverse: SccMap::with_hasher(hasher.clone()), source_collections: SccMap::with_hasher(hasher.clone()), - issue_counts: StateCountIndex::new(hasher.clone()), - pull_counts: StateCountIndex::new(hasher.clone()), + issue_counts: StateViewIndex::new(hasher.clone()), + pull_counts: StateViewIndex::new(hasher.clone()), hasher, writer: Mutex::new(()), } @@ -665,7 +741,7 @@ impl EdgeStore { if let Some(keys) = promote_keys { let authors = self.build_author_map(&keys); let large = LargeBucket { - keys: keys.into_iter().collect(), + keys: BlockStore::from_sorted(keys), authors, }; if let Entry::Occupied(mut e) = self.forward.entry_sync(key_id) { @@ -830,7 +906,7 @@ impl EdgeStore { &self, states: &StateIndex, entity: &AtUri, - edge_kind: &str, + edge_kinds: &[&str], ) -> (SourceId, SmallVec<[ProjectedState; 1]>) where K: StateKind + Default, @@ -846,7 +922,8 @@ impl EdgeStore { .into_iter() .filter_map(|entry| { let key = self.keys.read_sync(&entry.key_id, |_, key| key.clone())?; - if key.kind.as_ref() != edge_kind { + // repo buckets accept owner or issue author while by-buckets accept only issue author + if !edge_kinds.contains(&key.kind.as_ref()) { return None; } let repo_owner = key.subject.as_did()?; @@ -861,6 +938,8 @@ impl EdgeStore { key: entry.key_id, kind, author, + micros: entry.sort_micros, + by_bucket: edge_kinds[1..].contains(&key.kind.as_ref()), }) }) .collect(); @@ -869,54 +948,17 @@ impl EdgeStore { fn refresh_state_counts( &self, - index: &StateCountIndex, + index: &StateViewIndex, states: &StateIndex, entity: &AtUri, - edge_kind: &str, + edge_kinds: &[&str], ) where K: StateKind + Default, { // compute under the write lock so concurrent updates dont write a stale projection - let mut inner = index - .inner - .write() - .expect("state-count index rwlock poisoned"); - let (source, current) = self.current_state_projection(states, entity, edge_kind); - - if let Some(previous) = inner.projected.remove(&source) { - for projected in previous { - let bucket_key = (projected.key, projected.kind); - let remove = if let Some(bucket) = inner.buckets.get_mut(&bucket_key) { - bucket.count -= 1; - if let Some(author) = projected.author { - drop_author(&mut bucket.authors, author); - } - bucket.count == 0 - } else { - false - }; - if remove { - inner.buckets.remove(&bucket_key); - } - } - } - - for projected in current.iter().copied() { - let bucket = inner - .buckets - .entry((projected.key, projected.kind)) - .or_insert_with(|| CountBucket { - count: 0, - authors: HashMap::with_hasher(index.hasher.clone()), - }); - bucket.count += 1; - if let Some(author) = projected.author { - bump_author(&mut bucket.authors, author); - } - } - if !current.is_empty() { - inner.projected.insert(source, current); - } + let mut inner = index.write_inner(); + let (source, current) = self.current_state_projection(states, entity, edge_kinds); + StateViewIndex::apply_projection(&mut inner, source, ¤t); } pub fn refresh_issue_counts( @@ -924,7 +966,12 @@ impl EdgeStore { states: &StateIndex, entity: &AtUri, ) { - self.refresh_state_counts(&self.issue_counts, states, entity, "sh.tangled.repo.issue"); + self.refresh_state_counts( + &self.issue_counts, + states, + entity, + &["sh.tangled.repo.issue", "sh.tangled.repo.issue.by"], + ); } pub fn refresh_pull_counts( @@ -932,12 +979,17 @@ impl EdgeStore { states: &StateIndex, entity: &AtUri, ) { - self.refresh_state_counts(&self.pull_counts, states, entity, "sh.tangled.repo.pull"); + self.refresh_state_counts( + &self.pull_counts, + states, + entity, + &["sh.tangled.repo.pull", "sh.tangled.repo.pull.by"], + ); } fn count_state( &self, - index: &StateCountIndex, + index: &StateViewIndex, key: &EdgeKey, kind: K, author: Option<&Did>, @@ -955,25 +1007,85 @@ impl EdgeStore { }, None => None, }; - let inner = index - .inner - .read() - .expect("state-count index rwlock poisoned"); - let Some(bucket) = inner.buckets.get(&(key, kind)) else { - return FilteredCount::default(); + let (count, distinct) = index.counts(key, kind, author); + FilteredCount::new(count, distinct) + } + + fn page_state( + &self, + index: &StateViewIndex, + key: &EdgeKey, + kind: Option, + start: impl Into, + limit: PageLimit, + dir: SortDir, + author: Option<&Did>, + ) -> EdgePage + where + K: StateKind, + { + let empty = |total: u64| EdgePage { + items: Vec::new(), + next: None, + total: Some(TotalCount::from(total)), }; - match author { - Some(author) => { - let count = bucket - .authors - .get(&author) - .map_or(0, |count| count.get() as u64); - FilteredCount::new(count, u64::from(count != 0)) - } - None => FilteredCount::new(bucket.count, bucket.authors.len() as u64), + let Some(key_id) = self.lookup_key(key) else { + return empty(0); + }; + let author_id = match author { + Some(a) => match self.did_interner.get(a.as_ref()).map(AuthorId::from_spur) { + Some(id) => Some(id), + None => return empty(0), + }, + None => None, + }; + let Some(page) = index.page(key_id, kind, start.into(), limit, dir, author_id) else { + return empty(0); + }; + 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, + }) + }) + .collect(); + EdgePage { + items, + next: page.next, + total: Some(page.total), } } + pub fn page_issue_state( + &self, + key: &EdgeKey, + kind: Option, + start: impl Into, + limit: PageLimit, + dir: SortDir, + author: Option<&Did>, + ) -> EdgePage { + self.page_state(&self.issue_counts, key, kind, start, limit, dir, author) + } + + pub fn page_pull_status( + &self, + key: &EdgeKey, + kind: Option, + start: impl Into, + limit: PageLimit, + dir: SortDir, + author: Option<&Did>, + ) -> EdgePage { + self.page_state(&self.pull_counts, key, kind, start, limit, dir, author) + } + pub fn count_issue_state( &self, key: &EdgeKey, @@ -1038,14 +1150,23 @@ impl EdgeStore { pub fn list( &self, key: &EdgeKey, - cursor: PageCursor, + start: impl Into, limit: PageLimit, 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 Some(cursor) = sources.resolve_start(start, dir) else { + return EdgePage { + items: Vec::new(), + next: None, + total, + }; + }; let iter = sources.directed(cursor, dir); let entries: Vec = iter.take(limit_usize + 1).collect(); let has_more = entries.len() > limit_usize; @@ -1066,12 +1187,13 @@ impl EdgeStore { .then(|| page.last().copied()) .flatten() .map(BucketKey::token); - EdgePage { items, next } + EdgePage { items, next, total } }) }) .unwrap_or(EdgePage { items: Vec::new(), next: None, + total: None, }) } @@ -1125,78 +1247,11 @@ impl EdgeStore { .then(|| page.last().copied()) .flatten() .map(BucketKey::token); - EdgePage { items, next } - } - - pub fn list_filtered( - &self, - key: &EdgeKey, - cursor: PageCursor, - limit: PageLimit, - dir: SortDir, - predicate: F, - ) -> EdgePage - where - F: Fn(&AtUri) -> bool, - { - let limit_usize = limit.get() as usize; - let scan_cap = limit_usize - .saturating_mul(FILTER_SCAN_MULTIPLIER) - .max(FILTER_SCAN_FLOOR); - self.lookup_key(key) - .and_then(|id| { - self.forward.read_sync(&id, |_, sources| { - let init = ScanState::with_capacity(limit_usize + 1); - let outcome = sources - .directed(cursor, dir) - .try_fold(init, |mut state, key| { - if state.scanned >= scan_cap && state.matched.len() <= limit_usize { - return ControlFlow::Break(state); - } - state.scanned += 1; - state.last_scanned = Some(key); - if let Some(spur) = key.source.to_spur() - && let Some(stored) = self.source_interner.try_resolve(&spur) - && let Some(decoded) = self.decode_source(stored) - && let Ok(uri) = AtUri::new_owned(decoded) - && predicate(&uri) - { - state.matched.push((key, uri)); - if state.matched.len() > limit_usize { - return ControlFlow::Break(state); - } - } - ControlFlow::Continue(state) - }); - let (state, bucket_exhausted) = match outcome { - ControlFlow::Continue(s) => (s, true), - ControlFlow::Break(s) => (s, false), - }; - let has_more_matches = state.matched.len() > limit_usize; - let visible_len = state.matched.len().min(limit_usize); - let next = if has_more_matches && visible_len > 0 { - Some(state.matched[visible_len - 1].0.token()) - } else if !bucket_exhausted { - state.last_scanned.map(BucketKey::token) - } else { - None - }; - let items = state - .matched - .into_iter() - .take(visible_len) - .map(|(key, uri)| EdgeItem { - uri, - sort_micros: key.micros.0, - }) - .collect(); - EdgePage { items, next } - }) - }) - .unwrap_or(EdgePage { - items: Vec::new(), - next: None, - }) + EdgePage { + items, + next, + total: None, + } } pub fn key_count(&self) -> usize { @@ -1366,29 +1421,6 @@ fn directed_slice( } } -fn directed_tree( - keys: &BTreeSet, - cursor: PageCursor, - dir: SortDir, -) -> impl Iterator + '_ { - match dir { - SortDir::Asc => { - let lower = match cursor { - PageCursor::Start => Bound::Unbounded, - PageCursor::After(tok) => Bound::Excluded(BucketKey::from_token(tok)), - }; - Either::Left(keys.range((lower, Bound::Unbounded)).copied()) - } - SortDir::Desc => { - let upper = match cursor { - PageCursor::Start => Bound::Unbounded, - PageCursor::After(tok) => Bound::Excluded(BucketKey::from_token(tok)), - }; - Either::Right(keys.range((Bound::Unbounded, upper)).rev().copied()) - } - } -} - fn split_record_uri(source: &str) -> Option<(&str, &str, &str)> { let rest = source.strip_prefix("at://")?; let mut parts = rest.split('/'); @@ -1485,20 +1517,15 @@ mod tests { rows.into_iter().map(|(_, i)| source_uri(i)).collect() } - fn paginate_all(store: &EdgeStore, key: &EdgeKey, dir: SortDir, page: u32) -> Vec { + fn paginate_all(store: &EdgeStore, key: &EdgeKey, dir: SortDir, page: PageLimit) -> Vec { std::iter::successors( - Some(store.list(key, PageCursor::Start, limit(page), dir)), + Some(store.list(key, PageCursor::Start, page, dir)), |prev| { prev.next - .map(|tok| store.list(key, PageCursor::After(tok), limit(page), dir)) + .map(|tok| store.list(key, PageCursor::After(tok), page, dir)) }, ) - .flat_map(|p| { - p.items - .into_iter() - .map(|u| u.as_ref().to_owned()) - .collect::>() - }) + .flat_map(|p| p.items.into_iter().map(|u| u.as_ref().to_owned())) .collect() } @@ -1519,12 +1546,12 @@ mod tests { [3u32, 7, 50].into_iter().for_each(|page| { assert_eq!( - paginate_all(&store, &key, SortDir::Asc, page), + paginate_all(&store, &key, SortDir::Asc, limit(page)), asc, "asc mismatch n={n} page={page}" ); assert_eq!( - paginate_all(&store, &key, SortDir::Desc, page), + paginate_all(&store, &key, SortDir::Desc, limit(page)), desc, "desc mismatch n={n} page={page}" ); @@ -1532,45 +1559,370 @@ mod tests { }); } - // Build three subject buckets with interleaved micros so the global time - // order crosses buckets, and return the three keys. - fn fill_three_subjects(store: &EdgeStore) -> Vec { - [ - ("alpha", "r0", 10u64), - ("bravo", "r1", 50), - ("charlie", "r2", 20), - ("alpha", "r3", 40), - ("bravo", "r4", 30), - ("charlie", "r5", 60), - ] - .into_iter() - .map(|(subj, rkey, micros)| { - star_edge_at( - at(&format!("at://did:plc:x/sh.tangled.feed.star/{rkey}")), - did(&format!("did:plc:{subj}")), - micros, - ) - }) - .for_each(|edge| store.add(edge)); - ["alpha", "bravo", "charlie"] - .iter() - .map(|s| { - EdgeKey::new( - nsid("sh.tangled.feed.star"), - did_subj(&format!("did:plc:{s}")), - ) - }) - .collect() - } - - #[test] - fn list_multi_merges_subjects_newest_first() { - let store = store(); - let keys = fill_three_subjects(&store); - let page = store.list_multi(&keys, PageCursor::Start, limit(10), SortDir::Desc); - let micros: Vec = page.items.iter().map(|it| it.sort_micros).collect(); - assert_eq!(micros, vec![60, 50, 40, 30, 20, 10]); - assert!(page.next.is_none()); + fn paginate_by_offsets( + store: &EdgeStore, + key: &EdgeKey, + dir: SortDir, + page: PageLimit, + ) -> Vec { + std::iter::successors( + Some((PageOffset::ZERO, store.list(key, PageOffset::ZERO, page, dir))), + |(offset, prev)| { + prev.next.map(|_| { + let next = PageOffset::from(offset.get() + page.get() as u64); + (next, store.list(key, next, page, dir)) + }) + }, + ) + .flat_map(|(_, p)| p.items.into_iter().map(|u| u.as_ref().to_owned())) + .collect() + } + + #[test] + fn offset_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 = reference(0..n); + let desc: Vec = asc.iter().rev().cloned().collect(); + + [3u32, 7, 50].into_iter().for_each(|page| { + assert_eq!( + paginate_by_offsets(&store, &key, SortDir::Asc, limit(page)), + asc, + "asc offset walk mismatch n={n} page={page}" + ); + assert_eq!( + paginate_by_offsets(&store, &key, SortDir::Desc, limit(page)), + desc, + "desc offset walk mismatch n={n} page={page}" + ); + }); + + for o in [1u64, 5, 63, 255, 256, 257, 999] { + let o = o.min(n as u64 - 1); + let page = store.list(&key, PageStart::Offset(PageOffset::from(o)), limit(10), SortDir::Asc); + let want: Vec = asc.iter().skip(o as usize).take(10).cloned().collect(); + let got: Vec = page.items.iter().map(|u| u.as_ref().to_owned()).collect(); + assert_eq!(got, want, "asc offset {o} mismatch n={n}"); + assert_eq!(page.total, Some(TotalCount::from(n as u64)), "total n={n}"); + } + + let page = store.list(&key, PageStart::Offset(PageOffset::from(n as u64)), limit(10), SortDir::Asc); + assert!(page.items.is_empty(), "offset==len must be empty n={n}"); + assert!(page.next.is_none()); + assert_eq!(page.total, Some(TotalCount::from(n as u64))); + }); + } + + #[test] + fn offset_page_cursor_continues_sequentially() { + let n = 1000usize; + let store = store(); + let key = fill_subject(&store, n); + let asc = reference(0..n); + + let first = store.list(&key, PageStart::Offset(PageOffset::new(600)), limit(7), SortDir::Asc); + let got: Vec = std::iter::successors(Some(first), |prev| { + prev.next.map(|tok| store.list(&key, PageCursor::After(tok), limit(7), SortDir::Asc)) + }) + .flat_map(|p| p.items.into_iter().map(|u| u.as_ref().to_owned())) + .collect(); + } + + #[test] + fn offset_desc_matches_cursor_desc_small_bucket() { + let store = store(); + let key = fill_subject(&store, 5); + let desc: Vec = reference(0..5).into_iter().rev().collect(); + + let first = store.list(&key, PageStart::Offset(PageOffset::new(2)), limit(3), SortDir::Desc); + let got: Vec = std::iter::successors(Some(first), |prev| { + prev.next.map(|tok| store.list(&key, PageCursor::After(tok), limit(3), SortDir::Desc)) + }) + .flat_map(|p| p.items.into_iter().map(|u| u.as_ref().to_owned())) + .collect(); + } + + #[test] + fn offset_desc_matches_cursor_desc_large_bucket() { + let n = 1000usize; + let store = store(); + let key = fill_subject(&store, n); + let desc: Vec = reference(0..n).into_iter().rev().collect(); + + let first = store.list(&key, PageStart::Offset(PageOffset::new(600)), limit(7), SortDir::Desc); + let got: Vec = std::iter::successors(Some(first), |prev| { + prev.next.map(|tok| store.list(&key, PageCursor::After(tok), limit(7), SortDir::Desc)) + }) + .flat_map(|p| p.items.into_iter().map(|u| u.as_ref().to_owned())) + .collect(); + } + + #[test] + fn offset_boundary_matrix_asc_and_desc() { + // offsets {0, 1, len-1, len} across asc and desc pin resolve_start pred_rank math + for &n in &[5usize, 300] { + let store = store(); + let key = fill_subject(&store, n); + let asc = reference(0..n); + let desc: Vec = asc.iter().rev().cloned().collect(); + + for &(dir, ref_order, label) in + &[(SortDir::Asc, &asc, "asc"), (SortDir::Desc, &desc, "desc")] + { + let p = store.list(&key, PageStart::Offset(PageOffset::ZERO), limit(n as u32), dir); + let got: Vec = p.items.iter().map(|u| u.as_ref().to_owned()).collect(); + assert_eq!(got, ref_order[..n].to_vec(), "n={n} {label} offset=0"); + + let p = store.list(&key, PageStart::Offset(PageOffset::new(1)), limit(n as u32), dir); + let got: Vec = p.items.iter().map(|u| u.as_ref().to_owned()).collect(); + assert_eq!(got, ref_order[1..].to_vec(), "n={n} {label} offset=1"); + + let p = store.list(&key, PageStart::Offset(PageOffset::from(n as u64 - 1)), limit(n as u32), dir); + let got: Vec = p.items.iter().map(|u| u.as_ref().to_owned()).collect(); + assert_eq!( + got, + vec![ref_order[n - 1].clone()], + "n={n} {label} offset=len-1" + ); + assert!(p.next.is_none(), "last item has no continuation"); + + let p = store.list(&key, PageStart::Offset(PageOffset::from(n as u64)), limit(n as u32), dir); + assert!(p.items.is_empty(), "n={n} {label} offset=len is empty"); + assert_eq!(p.total, Some(TotalCount::from(n as u64)), "past-end page still reports total"); + } + } + } + #[test] + fn descending_cursor_block_boundary_edges() { + let store = store(); + let n = 600usize; + let key = fill_subject(&store, n); + let asc = reference(0..n); + let desc: Vec = asc.iter().rev().cloned().collect(); + + for target_idx in [255usize, 256, 257, 511, 512] { + let p_single = store.list( + &key, + PageStart::Offset(PageOffset::from(target_idx as u64)), + limit(1), + SortDir::Asc, + ); + let tok = p_single.next.unwrap(); + let page = store.list( + &key, + PageStart::Cursor(PageCursor::After(tok)), + limit(600), + SortDir::Desc, + ); + let got: Vec = page.items.iter().map(|u| u.as_ref().to_owned()).collect(); + let expected = &desc[n - target_idx..]; + assert_eq!(got, expected, "descending cursor after index {target_idx}"); + } + + let page = store.list(&key, PageStart::Offset(PageOffset::new(256)), limit(10), SortDir::Desc); + let got: Vec = page.items.iter().map(|u| u.as_ref().to_owned()).collect(); + assert_eq!(got, desc[256..266].to_vec(), "descending offset 256"); + } + #[test] + fn empty_store_descending_after_cursor_does_not_underflow() { + let store = store(); + let key = EdgeKey::new( + nsid("sh.tangled.feed.star"), + SubjectRef::Did(did("did:plc:empty_test")), + ); + let tok = PageToken::new(12345, 1); + let page = store.list( + &key, + PageStart::Cursor(PageCursor::After(tok)), + limit(10), + SortDir::Desc, + ); + assert!(page.items.is_empty()); + } + + #[test] + fn property_random_offsets_match_reference() { + // seeded random inserts with duplicate attempts swept across offsets + fn xorshift(mut x: u64) -> u64 { + x ^= x << 13; + x ^= x >> 7; + x ^= x << 17; + x + } + + let store = store(); + let key = EdgeKey::new( + nsid("sh.tangled.repo.issue"), + SubjectRef::Did(did("did:plc:prop")), + ); + + let mut rng = 0xdead_beef_cafe_babe_u64; + let mut order: Vec = (0..500).collect(); + for i in (1..order.len()).rev() { + rng = xorshift(rng); + let j = rng as usize % (i + 1); + order.swap(i, j); + } + for _ in 0..100 { + rng = xorshift(rng); + order.push(rng as usize % 500); + } + + for &i in &order { + store.add(Edge { + kind: nsid("sh.tangled.repo.issue"), + subject: SubjectRef::Did(did("did:plc:prop")), + source: at(&format!("at://did:plc:prop/sh.tangled.repo.issue/i{i}")), + sort_micros: (i as u64 + 1) * 1000, + }); + } + + let asc: Vec = (0..500) + .map(|i| format!("at://did:plc:prop/sh.tangled.repo.issue/i{i}")) + .collect(); + let desc: Vec = asc.iter().rev().cloned().collect(); + let n = 500usize; + + let mut offsets: Vec = (0..n as u64).step_by(50).collect(); + offsets.extend([n as u64 - 1, n as u64]); + offsets.sort(); + offsets.dedup(); + + for o in offsets { + for &(dir, ref_order, label) in + &[(SortDir::Asc, &asc, "asc"), (SortDir::Desc, &desc, "desc")] + { + let first = store.list(&key, PageStart::Offset(PageOffset::from(o)), limit(25), dir); + let got: Vec = std::iter::successors(Some(first), |prev| { + prev.next.map(|tok| store.list(&key, PageCursor::After(tok), limit(25), dir)) + }) + .flat_map(|p| p.items.into_iter().map(|u| u.as_ref().to_owned())) + .collect(); + let expected = if (o as usize) < n { + ref_order[o as usize..].to_vec() + } else { + vec![] + }; + assert_eq!(got, expected, "random n=500 {label} offset={o}"); + } + } + } + + #[test] + fn block_store_matches_sorted_vec_oracle() { + let mut blocks = BlockStore::from_sorted(Vec::new()); + let mut oracle: Vec = Vec::new(); + let key_at = |i: usize| BucketKey::new(shuffled_micros(i), SourceId::from_raw(i as u32)); + + let mut inserted: Vec = (0..1500).collect(); + inserted.sort_by_key(|i| shuffled_micros(*i ^ 0x5a5a)); + for i in inserted { + let key = key_at(i); + assert!(blocks.insert(key), "insert {i}"); + let pos = oracle.partition_point(|k| *k < key); + oracle.insert(pos, key); + } + assert_eq!(blocks.len(), oracle.len()); + + // remove every third key in a different shuffled order + let mut removed: Vec = oracle.iter().step_by(3).copied().collect(); + removed.sort_by_key(|k| shuffled_micros(k.source.index() as usize ^ 0x33)); + for key in &removed { + assert!(blocks.remove(key), "remove present key"); + let pos = oracle.partition_point(|k| *k < *key); + oracle.remove(pos); + } + assert_eq!(blocks.len(), oracle.len()); + assert!(!blocks.remove(&removed[0]), "double remove must be false"); + + for rank in [ + 0usize, + 1, + 255, + 256, + 300, + 511, + oracle.len() / 2, + oracle.len() - 1, + ] { + assert_eq!( + blocks.select(PageRank::new(rank)), + Some(oracle[rank]), + "select rank {rank}" + ); + } + assert_eq!(blocks.select(PageRank::new(oracle.len())), None, "select past end"); + + let asc: Vec = blocks.directed(PageCursor::Start, SortDir::Asc).collect(); + assert_eq!(asc, oracle, "asc full iteration"); + let desc: Vec = blocks.directed(PageCursor::Start, SortDir::Desc).collect(); + assert_eq!( + desc, + oracle.iter().rev().copied().collect::>(), + "desc full iteration" + ); + + // cursor-resume from a mid-block token in both directions + let pivot = oracle[oracle.len() / 2]; + let after_asc: Vec = blocks + .directed(PageCursor::After(pivot.token()), SortDir::Asc) + .collect(); + let want_asc: Vec = oracle.iter().copied().filter(|k| *k > pivot).collect(); + assert_eq!(after_asc, want_asc, "asc after pivot"); + let after_desc: Vec = blocks + .directed(PageCursor::After(pivot.token()), SortDir::Desc) + .collect(); + let want_desc: Vec = oracle + .iter() + .rev() + .copied() + .filter(|k| *k < pivot) + .collect(); + assert_eq!(after_desc, want_desc, "desc after pivot"); + } + + // Build three subject buckets with interleaved micros so the global time + // order crosses buckets, and return the three keys. + fn fill_three_subjects(store: &EdgeStore) -> Vec { + [ + ("alpha", "r0", 10u64), + ("bravo", "r1", 50), + ("charlie", "r2", 20), + ("alpha", "r3", 40), + ("bravo", "r4", 30), + ("charlie", "r5", 60), + ] + .into_iter() + .map(|(subj, rkey, micros)| { + star_edge_at( + at(&format!("at://did:plc:x/sh.tangled.feed.star/{rkey}")), + did(&format!("did:plc:{subj}")), + micros, + ) + }) + .for_each(|edge| store.add(edge)); + ["alpha", "bravo", "charlie"] + .iter() + .map(|s| { + EdgeKey::new( + nsid("sh.tangled.feed.star"), + did_subj(&format!("did:plc:{s}")), + ) + }) + .collect() + } + + #[test] + fn list_multi_merges_subjects_newest_first() { + let store = store(); + let keys = fill_three_subjects(&store); + let page = store.list_multi(&keys, PageCursor::Start, limit(10), SortDir::Desc); + let micros: Vec = page.items.iter().map(|it| it.sort_micros).collect(); + assert_eq!(micros, vec![60, 50, 40, 30, 20, 10]); + assert!(page.next.is_none()); } #[test] @@ -1619,9 +1971,9 @@ mod tests { NAMES.len() as u64, "every author keeps sources after partial removal" ); - assert_eq!(paginate_all(&store, &key, SortDir::Asc, 7), expected); + assert_eq!(paginate_all(&store, &key, SortDir::Asc, limit(7)), expected); assert_eq!( - paginate_all(&store, &key, SortDir::Desc, 11), + paginate_all(&store, &key, SortDir::Desc, limit(11)), expected.iter().rev().cloned().collect::>() ); } @@ -1778,6 +2130,314 @@ mod tests { ); } + fn seed_issues( + store: &EdgeStore, + states: &StateIndex, + repo: &str, + specs: &[(u64, &str)], + ) -> Vec> { + let kind = nsid("sh.tangled.repo.issue"); + specs + .iter() + .map(|(micros, name)| { + let issue = at(&format!( + "at://did:plc:{name}/sh.tangled.repo.issue/{name}{micros}" + )); + store.upsert_source( + &issue, + vec![Edge { + kind: kind.clone(), + subject: SubjectRef::Did(did(repo)), + source: issue.clone(), + sort_micros: *micros, + }], + ); + store.refresh_issue_counts(states, &issue); + issue + }) + .collect() + } + + #[test] + fn state_view_pages_default_open_by_rank_with_total() { + let store = store(); + let states = StateIndex::new(RuntimeHasher::default()); + let key = EdgeKey::new( + nsid("sh.tangled.repo.issue"), + SubjectRef::Did(did("did:plc:limpet")), + ); + let issues = seed_issues( + &store, + &states, + "did:plc:limpet", + &[ + (1, "nel"), + (2, "wren"), + (3, "otis"), + (4, "moss"), + (5, "fern"), + ], + ); + + let page = store.page_issue_state( + &key, + Some(IssueStateKind::Open), + PageStart::Offset(PageOffset::new(2)), + limit(2), + SortDir::Asc, + None, + ); + assert_eq!( + page.total, + Some(TotalCount::new(5)), + ); + assert_eq!( + page.items + .iter() + .map(|i| i.uri.as_ref()) + .collect::>(), + vec![issues[2].as_ref(), issues[3].as_ref()], + ); + let follow = store.page_issue_state( + &key, + Some(IssueStateKind::Open), + PageStart::Cursor(PageCursor::After(page.next.unwrap())), + limit(2), + SortDir::Asc, + None, + ); + assert_eq!(follow.items.len(), 1); + assert_eq!(follow.items[0].uri, issues[4]); + } + + #[test] + fn state_view_transitions_repage_across_kinds() { + let store = store(); + let states = StateIndex::new(RuntimeHasher::default()); + let key = EdgeKey::new( + nsid("sh.tangled.repo.issue"), + SubjectRef::Did(did("did:plc:limpet")), + ); + let issues = seed_issues( + &store, + &states, + "did:plc:limpet", + &[(1, "nel"), (2, "wren")], + ); + + states.upsert( + at("at://did:plc:limpet/sh.tangled.repo.issue.state/s1"), + issues[0].clone(), + 100, + IssueStateKind::Closed, + ); + store.refresh_issue_counts(&states, &issues[0]); + + let open = store.page_issue_state( + &key, + Some(IssueStateKind::Open), + PageStart::Offset(PageOffset::ZERO), + limit(10), + SortDir::Asc, + None, + ); + assert_eq!(open.total, Some(TotalCount::new(1))); + assert_eq!(open.items[0].uri, issues[1]); + let closed = store.page_issue_state( + &key, + Some(IssueStateKind::Closed), + PageStart::Offset(PageOffset::ZERO), + limit(10), + SortDir::Asc, + None, + ); + assert_eq!(closed.total, Some(TotalCount::new(1))); + assert_eq!(closed.items[0].uri, issues[0]); + + // unauthorized state records cannot move the view + states.upsert( + at("at://did:plc:rando/sh.tangled.repo.issue.state/s2"), + issues[1].clone(), + 200, + IssueStateKind::Closed, + ); + store.refresh_issue_counts(&states, &issues[1]); + let open = store.page_issue_state( + &key, + Some(IssueStateKind::Open), + PageStart::Offset(PageOffset::ZERO), + limit(10), + SortDir::Asc, + None, + ); + assert_eq!( + open.total, + Some(TotalCount::new(1)), + "untrusted state record leaves the view alone" + ); + } + + #[test] + fn state_view_offset_desc_matches_cursor_desc() { + // views resolve offsets through SortedBlocks::resolve_start directly: + // issues i1..i5 ascending page as i5 i4 i3 i2 i1 descending, offset 2 + // starts at i3 + let store = store(); + let states = StateIndex::new(RuntimeHasher::default()); + let key = EdgeKey::new( + nsid("sh.tangled.repo.issue"), + SubjectRef::Did(did("did:plc:limpet")), + ); + let issues = seed_issues( + &store, + &states, + "did:plc:limpet", + &[(1, "nel"), (2, "nel"), (3, "nel"), (4, "nel"), (5, "nel")], + ); + + let page = store.page_issue_state( + &key, + Some(IssueStateKind::Open), + PageStart::Offset(PageOffset::new(2)), + limit(2), + SortDir::Desc, + None, + ); + assert_eq!(page.total, Some(TotalCount::new(5))); + assert_eq!( + page.items + .iter() + .map(|i| i.uri.as_ref()) + .collect::>(), + vec![issues[2].as_ref(), issues[1].as_ref()], + "desc offset starts at the third-newest", + ); + + let tail = store.page_issue_state( + &key, + Some(IssueStateKind::Open), + PageCursor::After(page.next.unwrap()), + limit(5), + SortDir::Desc, + None, + ); + assert_eq!( + tail.items + .iter() + .map(|i| i.uri.as_ref()) + .collect::>(), + vec![issues[0].as_ref()], + "cursor continuation finishes the desc walk", + ); + assert!(tail.next.is_none()); + } + + #[test] + fn state_view_author_filter_rank_selects_exactly() { + let store = store(); + let states = StateIndex::new(RuntimeHasher::default()); + let key = EdgeKey::new( + nsid("sh.tangled.repo.issue"), + SubjectRef::Did(did("did:plc:limpet")), + ); + let issues = seed_issues( + &store, + &states, + "did:plc:limpet", + &[(1, "nel"), (2, "wren"), (3, "nel"), (4, "wren"), (5, "nel")], + ); + let nel = did("did:plc:nel"); + + let page = store.page_issue_state( + &key, + Some(IssueStateKind::Open), + PageStart::Offset(PageOffset::new(1)), + limit(2), + SortDir::Asc, + Some(&nel), + ); + assert_eq!( + page.total, + Some(TotalCount::new(3)), + "author totals come from the author view" + ); + assert_eq!( + page.items + .iter() + .map(|i| i.uri.as_ref()) + .collect::>(), + vec![issues[2].as_ref(), issues[4].as_ref()], + "offset rank-selects inside the author view", + ); + + // a cursor from an offset page continues keyset-style within the + // author view + let page = store.page_issue_state( + &key, + Some(IssueStateKind::Open), + PageStart::Offset(PageOffset::new(2)), + limit(1), + SortDir::Asc, + Some(&nel), + ); + assert!(page.next.is_none(), "third nel issue is the last"); + assert_eq!(page.items[0].uri, issues[4]); + } + + #[test] + fn state_view_any_kind_author_view_ignores_state() { + let store = store(); + let states = StateIndex::new(RuntimeHasher::default()); + let key = EdgeKey::new( + nsid("sh.tangled.repo.issue"), + SubjectRef::Did(did("did:plc:limpet")), + ); + let issues = seed_issues( + &store, + &states, + "did:plc:limpet", + &[(1, "nel"), (2, "wren"), (3, "nel"), (4, "wren"), (5, "nel")], + ); + let nel = did("did:plc:nel"); + // close one of nel's issues: the any-kind view still contains it + states.upsert( + at("at://did:plc:limpet/sh.tangled.repo.issue.state/s1"), + issues[2].clone(), + 42, + IssueStateKind::Closed, + ); + store.refresh_issue_counts(&states, &issues[2]); + + let page = store.page_issue_state( + &key, + None, + PageStart::Offset(PageOffset::ZERO), + limit(10), + SortDir::Asc, + Some(&nel), + ); + assert_eq!(page.total, Some(TotalCount::new(3))); + assert_eq!( + page.items + .iter() + .map(|i| i.uri.as_ref()) + .collect::>(), + vec![issues[0].as_ref(), issues[2].as_ref(), issues[4].as_ref()], + "any-kind author view spans open and closed", + ); + + // the kind-scoped author view does not + let page = store.page_issue_state( + &key, + Some(IssueStateKind::Open), + PageStart::Offset(PageOffset::ZERO), + limit(10), + SortDir::Asc, + Some(&nel), + ); + assert_eq!(page.total, Some(TotalCount::new(2))); + } + #[test] fn record_index_mutations_move_materialized_state_counts() { use bobbin_types::sh_tangled::repo::issue::state::{State as IssueStateRecord, StateState}; @@ -2020,52 +2680,6 @@ mod tests { assert_eq!(combined.len(), 4, "every tied entry visible exactly once"); } - #[test] - fn list_filtered_narrows_by_predicate_and_paginates() { - let store = store(); - let subject_did = did("did:plc:limpet"); - let key = EdgeKey::new( - nsid("sh.tangled.repo.issue"), - SubjectRef::Did(subject_did.clone()), - ); - let by_nel = (0..3).map(|i| { - star_edge_at( - at(&format!("at://did:plc:nel/sh.tangled.repo.issue/r{i}")), - subject_did.clone(), - 100 + i as u64, - ) - }); - let by_olaren = (0..2).map(|i| { - star_edge_at( - at(&format!("at://did:plc:olaren/sh.tangled.repo.issue/o{i}")), - subject_did.clone(), - 500 + i as u64, - ) - }); - by_nel - .chain(by_olaren) - .map(|mut e| { - e.kind = nsid("sh.tangled.repo.issue"); - e - }) - .for_each(|e| store.add(e)); - - let only_nel = |u: &AtUri| u.as_ref().starts_with("at://did:plc:nel/"); - let page1 = store.list_filtered(&key, PageCursor::Start, limit(2), SortDir::Asc, only_nel); - assert_eq!(page1.items.len(), 2); - assert!(page1.next.is_some(), "cursor must allow more nel matches"); - - let page2 = store.list_filtered( - &key, - PageCursor::After(page1.next.unwrap()), - limit(2), - SortDir::Asc, - only_nel, - ); - assert_eq!(page2.items.len(), 1, "only one nel issue left"); - assert!(page2.next.is_none(), "tail page must not promise more"); - } - #[test] fn list_descending_returns_newest_first() { let store = store(); @@ -2139,47 +2753,6 @@ mod tests { ); } - #[test] - fn list_filtered_descending_respects_predicate() { - let store = store(); - let subject_did = did("did:plc:scallop"); - let key = EdgeKey::new( - nsid("sh.tangled.repo.issue"), - SubjectRef::Did(subject_did.clone()), - ); - let by_nel = (0..3).map(|i| { - star_edge_at( - at(&format!("at://did:plc:nel/sh.tangled.repo.issue/r{i}")), - subject_did.clone(), - 100 + i as u64, - ) - }); - let by_olaren = (0..2).map(|i| { - star_edge_at( - at(&format!("at://did:plc:olaren/sh.tangled.repo.issue/o{i}")), - subject_did.clone(), - 500 + i as u64, - ) - }); - by_nel - .chain(by_olaren) - .map(|mut e| { - e.kind = nsid("sh.tangled.repo.issue"); - e - }) - .for_each(|e| store.add(e)); - - let only_nel = |u: &AtUri| u.as_ref().starts_with("at://did:plc:nel/"); - let page = store.list_filtered(&key, PageCursor::Start, limit(5), SortDir::Desc, only_nel); - assert_eq!(page.items.len(), 3, "all three nel issues visible"); - let last = page.items.last().unwrap().as_ref(); - let first = page.items.first().unwrap().as_ref(); - assert!( - first > last, - "desc order: first item rkey must be greater than last (got first={first}, last={last})", - ); - } - #[test] fn non_did_source_round_trips_via_raw_fallback() { let store = store(); diff --git a/bobbin/crates/edge-index/src/sorted_blocks.rs b/bobbin/crates/edge-index/src/sorted_blocks.rs new file mode 100644 index 000000000..bd41dc732 --- /dev/null +++ b/bobbin/crates/edge-index/src/sorted_blocks.rs @@ -0,0 +1,279 @@ +use ftree::FenwickTree; + +use crate::{PageCursor, PageRank, PageStart, PageToken, SortDir}; + +const BLOCK_TARGET: usize = 256; +const BLOCK_SPLIT_AT: usize = 2 * BLOCK_TARGET; +const BLOCK_MERGE_AT: usize = 64; + +pub(crate) trait PageKey: Copy + Ord { + fn token(self) -> PageToken; + fn from_token(tok: PageToken) -> Self; +} + +struct Block { + keys: Vec, + min: K, + max: K, +} + +impl Block { + fn new(mut 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 } + } + + fn update_bounds(&mut self) { + debug_assert!(!self.keys.is_empty()); + self.min = self.keys[0]; + self.max = *self.keys.last().unwrap(); + } + + fn len(&self) -> usize { + self.keys.len() + } +} + +pub(crate) struct SortedBlocks { + blocks: Vec>, + fenwick: FenwickTree, + len: usize, +} + +impl SortedBlocks { + pub fn from_sorted(keys: Vec) -> Self { + debug_assert!(keys.windows(2).all(|w| w[0] <= w[1])); + let len = keys.len(); + let blocks = keys + .chunks(BLOCK_TARGET) + .map(|c| Block::new(c.to_vec())) + .collect(); + let mut store = Self { + blocks, + fenwick: FenwickTree::new(), + len, + }; + store.rebuild_fenwick(); + store + } + + pub fn len(&self) -> usize { + self.len + } + + 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]) + } + + // converts offset to exclusive cursor so directed() handles iteration + pub fn resolve_start(&self, start: PageStart, dir: SortDir) -> Option { + match start { + PageStart::Cursor(c) => Some(c), + PageStart::Offset(o) if o.get() == 0 => Some(PageCursor::Start), + PageStart::Offset(o) => { + let len = self.len as u64; + let o = o.get(); + if o >= len { + return None; + } + let pred_rank = match dir { + SortDir::Asc => o - 1, + SortDir::Desc => len - o, + }; + self.select(PageRank::new(pred_rank as usize)) + .map(|k| PageCursor::After(k.token())) + } + } + } + + pub fn insert(&mut self, key: K) -> bool { + let Some(last_block) = self.blocks.last_mut() else { + self.blocks.push(Block::new(vec![key])); + self.len = 1; + self.rebuild_fenwick(); + return true; + }; + // fast path for append-heavy ingestion avoids block binary search + if key > last_block.max { + last_block.keys.push(key); + last_block.max = key; + let bi = self.blocks.len() - 1; + self.len += 1; + self.after_insert(bi); + return true; + } + let bi = self.blocks.partition_point(|b| b.max < key); + match self.blocks[bi].keys.binary_search(&key) { + Ok(_) => false, + Err(pos) => { + self.blocks[bi].keys.insert(pos, key); + self.blocks[bi].update_bounds(); + self.len += 1; + self.after_insert(bi); + true + } + } + } + + // 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); + return; + } + let tail_keys = self.blocks[bi].keys.split_off(BLOCK_TARGET); + self.blocks[bi].update_bounds(); + self.blocks.insert(bi + 1, Block::new(tail_keys)); + self.rebuild_fenwick(); + } + + pub fn remove(&mut self, key: &K) -> bool { + let bi = self.blocks.partition_point(|b| b.max < *key); + if bi == self.blocks.len() { + return false; + } + match self.blocks[bi].keys.binary_search(key) { + Ok(pos) => { + let old_len = self.blocks[bi].keys.len(); + self.blocks[bi].keys.remove(pos); + self.len -= 1; + if self.blocks[bi].keys.is_empty() { + self.blocks.remove(bi); + self.rebuild_fenwick(); + } else { + if pos == 0 || pos + 1 == old_len { + self.blocks[bi].update_bounds(); + } + self.after_remove(bi); + } + true + } + Err(_) => false, + } + } + + fn merge_blocks(&mut self, target: usize, source: usize) { + let removed = self.blocks.remove(source); + self.blocks[target].keys.extend(removed.keys); + self.blocks[target].update_bounds(); + self.rebuild_fenwick(); + } + + // merge underfilled block into adjacent neighbor when combined size fits split limit + 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 { + 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.merge_blocks(bi, bi + 1); + return; + } + } + self.fenwick.sub_at(bi, 1); + } + + pub fn directed(&self, cursor: PageCursor, dir: SortDir) -> Box + '_> { + match dir { + SortDir::Asc => { + let (bi, skip) = match cursor { + PageCursor::Start => (0, 0), + PageCursor::After(tok) => { + let key = K::from_token(tok); + let bi = self.blocks.partition_point(|b| b.max <= key); + if bi == self.blocks.len() { + return Box::new(std::iter::empty()); + } + (bi, self.blocks[bi].keys.partition_point(|k| *k <= key)) + } + }; + Box::new( + self.blocks[bi..] + .iter() + .enumerate() + .flat_map(move |(i, b)| { + let from = if i == 0 { skip } else { 0 }; + b.keys[from..].iter().copied() + }), + ) + } + SortDir::Desc => { + if self.blocks.is_empty() { + return Box::new(std::iter::empty()); + } + let (bi, take) = match cursor { + PageCursor::Start => (self.blocks.len() - 1, usize::MAX), + PageCursor::After(tok) => { + let key = K::from_token(tok); + // locate block for key, stepping to predecessor block if key is at block boundary + let bi = self.blocks.partition_point(|b| b.max < key); + if bi == self.blocks.len() { + (self.blocks.len() - 1, usize::MAX) + } else { + let take = self.blocks[bi].keys.partition_point(|k| *k < key); + if take > 0 { + (bi, take) + } else if bi > 0 { + (bi - 1, usize::MAX) + } else { + return Box::new(std::iter::empty()); + } + } + } + }; + Box::new( + self.blocks[..=bi] + .iter() + .enumerate() + .rev() + .flat_map(move |(i, b)| { + let upto = if i == bi { take.min(b.len()) } else { b.len() }; + b.keys[..upto].iter().rev().copied() + }), + ) + } + } + } + + pub fn heap_bytes(&self) -> u64 { + let block_overhead = std::mem::size_of::>() as u64; + self.blocks.capacity() as u64 * block_overhead + + self + .blocks + .iter() + .map(|b| (b.keys.capacity() * std::mem::size_of::()) as u64) + .sum::() + + (self.fenwick.len() * std::mem::size_of::()) as u64 + } +} diff --git a/bobbin/crates/edge-index/src/state_view.rs b/bobbin/crates/edge-index/src/state_view.rs new file mode 100644 index 000000000..579b52a94 --- /dev/null +++ b/bobbin/crates/edge-index/src/state_view.rs @@ -0,0 +1,251 @@ + +use std::collections::HashMap; +use std::num::NonZeroU32; +use std::sync::{RwLock, RwLockReadGuard, RwLockWriteGuard}; + +use bobbin_runtime::RuntimeHasher; +use smallvec::SmallVec; + +use crate::sorted_blocks::SortedBlocks; +use crate::{ + AuthorId, BucketKey, EdgeKeyId, PageLimit, PageStart, PageToken, SortDir, SourceId, StateKind, + TotalCount, bump_author, drop_author, +}; + +#[derive(Clone, Copy)] +pub(crate) struct ProjectedState { + pub key: EdgeKeyId, + pub kind: K, + pub author: Option, + pub micros: u64, + // by-buckets never get author sub-views since the bucket key is already the author + pub by_bucket: bool, +} + +pub(crate) struct ViewBucket { + blocks: SortedBlocks, + authors: HashMap, +} + +pub(crate) struct StateViewInner { + pub projected: HashMap; 1]>, RuntimeHasher>, + buckets: HashMap<(EdgeKeyId, K), ViewBucket, RuntimeHasher>, + // eager author sub-views so author queries rank select instead of scanning + author_views: HashMap<(EdgeKeyId, Option, AuthorId), SortedBlocks, RuntimeHasher>, +} + +pub(crate) struct ViewPage { + pub keys: Vec, + pub next: Option, + pub total: TotalCount, +} +// maps start offset/cursor onto a rank-selected page and total count +fn rank_page( + blocks: &SortedBlocks, + start: PageStart, + limit: PageLimit, + dir: SortDir, +) -> ViewPage { + let total = TotalCount::from(blocks.len()); + let limit_usize = limit.get() as usize; + let Some(cursor) = blocks.resolve_start(start, dir) else { + return ViewPage { + keys: Vec::new(), + next: None, + total, + }; + }; + let mut keys: Vec = blocks.directed(cursor, dir).take(limit_usize + 1).collect(); + let next = if keys.len() > limit_usize { + keys.truncate(limit_usize); + keys.last().map(|k| k.token()) + } else { + None + }; + ViewPage { keys, next, total } +} + +pub(crate) struct StateViewIndex { + inner: RwLock>, +} + +impl StateViewIndex { + pub fn new(hasher: RuntimeHasher) -> Self { + Self { + inner: RwLock::new(StateViewInner { + projected: HashMap::with_hasher(hasher.clone()), + buckets: HashMap::with_hasher(hasher.clone()), + author_views: HashMap::with_hasher(hasher.clone()), + }), + } + } + + // refresh runs under write lock to avoid committing stale projections + pub fn write_inner(&self) -> RwLockWriteGuard<'_, StateViewInner> { + self.inner + .write() + .expect("state-view index rwlock poisoned") + } + + fn read_inner(&self) -> RwLockReadGuard<'_, StateViewInner> { + self.inner.read().expect("state-view index rwlock poisoned") + } + + // updates state view buckets and per-author sub-views on state transition + pub fn apply_projection( + inner: &mut StateViewInner, + source: SourceId, + current: &SmallVec<[ProjectedState; 1]>, + ) { + if let Some(previous) = inner.projected.remove(&source) { + for projected in previous { + let remove = + if let Some(view) = inner.buckets.get_mut(&(projected.key, projected.kind)) { + view.blocks + .remove(&BucketKey::new(projected.micros, source)); + if let Some(author) = projected.author { + drop_author(&mut view.authors, author); + } + view.blocks.len() == 0 + } else { + false + }; + if remove { + inner.buckets.remove(&(projected.key, projected.kind)); + } + if let Some(author) = projected.author + && !projected.by_bucket + { + for kind in [Some(projected.kind), None] { + let remove = if let Some(blocks) = + inner.author_views.get_mut(&(projected.key, kind, author)) + { + blocks.remove(&BucketKey::new(projected.micros, source)); + blocks.len() == 0 + } else { + false + }; + if remove { + inner.author_views.remove(&(projected.key, kind, author)); + } + } + } + } + } + + for projected in current.iter().copied() { + let view = inner + .buckets + .entry((projected.key, projected.kind)) + .or_insert_with(|| ViewBucket { + blocks: SortedBlocks::from_sorted(Vec::new()), + authors: HashMap::with_hasher(RuntimeHasher::default()), + }); + view.blocks.insert(BucketKey::new(projected.micros, source)); + if let Some(author) = projected.author { + bump_author(&mut view.authors, author); + if !projected.by_bucket { + for kind in [Some(projected.kind), None] { + inner + .author_views + .entry((projected.key, kind, author)) + .or_insert_with(|| SortedBlocks::from_sorted(Vec::new())) + .insert(BucketKey::new(projected.micros, source)); + } + } + } + } + if !current.is_empty() { + inner.projected.insert(source, current.clone()); + } + } + + pub fn counts(&self, key: EdgeKeyId, kind: K, author: Option) -> (u64, u64) { + let inner = self.read_inner(); + match author { + Some(author) => { + let count = inner + .author_views + .get(&(key, Some(kind), author)) + .map_or(0, |blocks| blocks.len() as u64); + (count, u64::from(count != 0)) + } + None => { + let Some(view) = inner.buckets.get(&(key, kind)) else { + return (0, 0); + }; + (view.blocks.len() as u64, view.authors.len() as u64) + } + } + } + + // kind is None for author-only queries which hit the any-kind view + pub fn page( + &self, + key: EdgeKeyId, + kind: Option, + start: PageStart, + limit: PageLimit, + dir: SortDir, + author: Option, + ) -> Option { + let inner = self.read_inner(); + match author { + None => { + let view = inner.buckets.get(&(key, kind?))?; + Some(rank_page(&view.blocks, start, limit, dir)) + } + Some(author) => { + let page = inner + .author_views + .get(&(key, kind, author)) + .map(|blocks| rank_page(blocks, start, limit, dir)) + .unwrap_or(ViewPage { + keys: Vec::new(), + next: None, + total: TotalCount::ZERO, + }); + Some(page) + } + } + } + + pub fn heap_bytes(&self) -> u64 { + let inner = self.read_inner(); + let projected = inner.projected.capacity() + * (std::mem::size_of::() + + std::mem::size_of::; 1]>>() + + 1) + + inner + .projected + .values() + .filter(|states| states.spilled()) + .map(|states| states.capacity() * std::mem::size_of::>()) + .sum::(); + let buckets = inner.buckets.capacity() + * (std::mem::size_of::<(EdgeKeyId, K)>() + std::mem::size_of::() + 1); + let author_slots = inner + .buckets + .values() + .map(|bucket| { + bucket.authors.capacity() + * (std::mem::size_of::() + std::mem::size_of::() + 1) + }) + .sum::(); + let author_views = inner.author_views.capacity() + * (std::mem::size_of::<(EdgeKeyId, Option, AuthorId)>() + + std::mem::size_of::>() + + 1); + let blocks: u64 = inner + .buckets + .values() + .map(|b| b.blocks.heap_bytes()) + .sum::() + + inner + .author_views + .values() + .map(|b| b.heap_bytes()) + .sum::(); + (projected + buckets + author_slots + author_views) as u64 + blocks + } +}