diff --git a/Cargo.lock b/Cargo.lock --- a/Cargo.lock +++ b/Cargo.lock @@ -373,6 +373,7 @@ "bobbin-runtime", "bobbin-types", "either", + "itertools 0.14.0", "jacquard-common", "lasso", "scc", diff --git a/Cargo.toml b/Cargo.toml --- a/Cargo.toml +++ b/Cargo.toml @@ -57,6 +57,7 @@ tokio-tungstenite = { version = "0.29", features = ["rustls-tls-webpki-roots"] } futures = "0.3" either = "1" +itertools = "0.14" serde = { version = "1", features = ["derive"] } serde_json = { version = "1", features = ["raw_value"] } diff --git a/bobbin/crates/edge-index/Cargo.toml b/bobbin/crates/edge-index/Cargo.toml --- a/bobbin/crates/edge-index/Cargo.toml +++ b/bobbin/crates/edge-index/Cargo.toml @@ -9,6 +9,7 @@ bobbin-runtime = { workspace = true } bobbin-types = { workspace = true } either = { workspace = true } +itertools = { workspace = true } jacquard-common = { workspace = true } lasso = { workspace = true } scc = { workspace = true } diff --git a/bobbin/crates/edge-index/src/lib.rs b/bobbin/crates/edge-index/src/lib.rs --- a/bobbin/crates/edge-index/src/lib.rs +++ b/bobbin/crates/edge-index/src/lib.rs @@ -28,6 +28,7 @@ use bobbin_types::edges::Edge; use bobbin_types::ids::EdgeKey; use either::Either; +use itertools::Itertools; use jacquard_common::DefaultStr; use jacquard_common::types::string::AtUri; use lasso::{Key, Spur, ThreadedRodeo}; @@ -733,6 +734,55 @@ }) } + /// Merge a page across many buckets in one global time order. + pub fn list_multi( + &self, + keys: &[EdgeKey], + cursor: PageCursor, + limit: PageLimit, + dir: SortDir, + ) -> EdgePage { + let n = limit.get() as usize + 1; + let is_first = |x: &BucketKey, y: &BucketKey| match dir { + SortDir::Asc => x <= y, + SortDir::Desc => x >= y, + }; + let mut top: Vec = Vec::with_capacity(n); + for key in keys { + let Some(id) = self.lookup_key(key) else { + continue; + }; + let bucket: Vec = self + .forward + .read_sync(&id, |_, sources| { + sources.directed(cursor, dir).take(n).collect() + }) + .unwrap_or_default(); + top = top.into_iter().merge_by(bucket, &is_first).take(n).collect(); + } + top.dedup_by_key(|k| k.source); + + 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 next = has_more + .then(|| page.last().copied()) + .flatten() + .map(BucketKey::token); + EdgePage { items, next } + } + pub fn list_filtered( &self, key: &EdgeKey, @@ -1061,6 +1111,70 @@ ); }); }); + } + + // 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 { + let rows = [ + ("alpha", "r0", 10u64), + ("bravo", "r1", 50), + ("charlie", "r2", 20), + ("alpha", "r3", 40), + ("bravo", "r4", 30), + ("charlie", "r5", 60), + ]; + for (subj, rkey, micros) in rows { + let source = at(&format!("at://did:plc:x/sh.tangled.feed.star/{rkey}")); + store.add(star_edge_at(source, did(&format!("did:plc:{subj}")), micros)); + } + ["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] + fn list_multi_paginates_disjoint_across_buckets() { + let store = store(); + let keys = fill_three_subjects(&store); + // limit=4 with 6 items across 3 buckets exercises the bounded top-N merge. + let p1 = store.list_multi(&keys, PageCursor::Start, limit(4), SortDir::Desc); + let m1: Vec = p1.items.iter().map(|it| it.sort_micros).collect(); + assert_eq!(m1, vec![60, 50, 40, 30]); + let tok = p1.next.expect("first page has more"); + let p2 = store.list_multi(&keys, PageCursor::After(tok), limit(4), SortDir::Desc); + let m2: Vec = p2.items.iter().map(|it| it.sort_micros).collect(); + assert_eq!(m2, vec![20, 10]); + assert!(p2.next.is_none()); + } + + #[test] + fn list_multi_dedups_source_under_two_keys() { + let store = store(); + // Same record (source) indexed under two subject keys, same createdAt. + let shared = at("at://did:plc:x/sh.tangled.feed.star/shared"); + store.add(star_edge_at(shared.clone(), did("did:plc:alpha"), 100)); + store.add(star_edge_at(shared, did("did:plc:bravo"), 100)); + let keys = vec![ + EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:alpha")), + EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:bravo")), + ]; + let page = store.list_multi(&keys, PageCursor::Start, limit(10), SortDir::Desc); + assert_eq!(page.items.len(), 1); + assert!(page.next.is_none()); } #[test]