This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

bobbin-edge-index: `list_multi`

Signed-off-by: Seongmin Lee <git@boltless.me>

author
Seongmin Lee
committer
dawn
date (Jul 31, 2026, 10:57 PM +0300) commit fc1600ca parent 133a92f0 change-id kowwlutr
+129
+1
Cargo.lock
··· 711 711 "bobbin-runtime", 712 712 "bobbin-types", 713 713 "either", 714 + "itertools 0.14.0", 714 715 "jacquard-common", 715 716 "lasso", 716 717 "scc",
+1
Cargo.toml
··· 97 97 futures = "0.3" 98 98 either = "1" 99 99 async-trait = "0.1" 100 + itertools = "0.14" 100 101 101 102 serde = { version = "1", features = ["derive"] } 102 103 serde_json = { version = "1", features = ["raw_value"] }
+1
bobbin/crates/edge-index/Cargo.toml
··· 9 9 bobbin-runtime = { workspace = true } 10 10 bobbin-types = { workspace = true } 11 11 either = { workspace = true } 12 + itertools = { workspace = true } 12 13 jacquard-common = { workspace = true } 13 14 lasso = { workspace = true } 14 15 scc = { workspace = true }
+126
bobbin/crates/edge-index/src/lib.rs
··· 28 28 use bobbin_types::edges::{Edge, Record}; 29 29 use bobbin_types::ids::EdgeKey; 30 30 use either::Either; 31 + use itertools::Itertools; 31 32 use jacquard_common::DefaultStr; 32 33 use jacquard_common::types::did::Did; 33 34 use jacquard_common::types::nsid::Nsid; ··· 1036 1037 }) 1037 1038 } 1038 1039 1040 + /// Merge a page across many buckets in one global time order. 1041 + pub fn list_multi( 1042 + &self, 1043 + keys: &[EdgeKey], 1044 + cursor: PageCursor, 1045 + limit: PageLimit, 1046 + dir: SortDir, 1047 + ) -> EdgePage { 1048 + let n = limit.get() as usize + 1; 1049 + let is_first = |x: &BucketKey, y: &BucketKey| match dir { 1050 + SortDir::Asc => x <= y, 1051 + SortDir::Desc => x >= y, 1052 + }; 1053 + let mut top: Vec<BucketKey> = Vec::with_capacity(n); 1054 + for key in keys { 1055 + let Some(id) = self.lookup_key(key) else { 1056 + continue; 1057 + }; 1058 + let bucket: Vec<BucketKey> = self 1059 + .forward 1060 + .read_sync(&id, |_, sources| { 1061 + sources.directed(cursor, dir).take(n).collect() 1062 + }) 1063 + .unwrap_or_default(); 1064 + top = top 1065 + .into_iter() 1066 + .merge_by(bucket, &is_first) 1067 + .take(n) 1068 + .collect(); 1069 + } 1070 + top.dedup_by_key(|k| k.source); 1071 + 1072 + let has_more = top.len() > limit.get() as usize; 1073 + let page = &top[..top.len().min(limit.get() as usize)]; 1074 + let items = page 1075 + .iter() 1076 + .filter_map(|&k| { 1077 + let spur = k.source.to_spur()?; 1078 + let stored = self.source_interner.try_resolve(&spur)?; 1079 + let uri = AtUri::new_owned(self.decode_source(stored)?).ok()?; 1080 + Some(EdgeItem { 1081 + uri, 1082 + sort_micros: k.micros.0, 1083 + }) 1084 + }) 1085 + .collect(); 1086 + let next = has_more 1087 + .then(|| page.last().copied()) 1088 + .flatten() 1089 + .map(BucketKey::token); 1090 + EdgePage { items, next } 1091 + } 1092 + 1039 1093 pub fn list_filtered<F>( 1040 1094 &self, 1041 1095 key: &EdgeKey, ··· 1438 1492 ); 1439 1493 }); 1440 1494 }); 1495 + } 1496 + 1497 + // Build three subject buckets with interleaved micros so the global time 1498 + // order crosses buckets, and return the three keys. 1499 + fn fill_three_subjects(store: &EdgeStore) -> Vec<EdgeKey> { 1500 + [ 1501 + ("alpha", "r0", 10u64), 1502 + ("bravo", "r1", 50), 1503 + ("charlie", "r2", 20), 1504 + ("alpha", "r3", 40), 1505 + ("bravo", "r4", 30), 1506 + ("charlie", "r5", 60), 1507 + ] 1508 + .into_iter() 1509 + .map(|(subj, rkey, micros)| { 1510 + star_edge_at( 1511 + at(&format!("at://did:plc:x/sh.tangled.feed.star/{rkey}")), 1512 + did(&format!("did:plc:{subj}")), 1513 + micros, 1514 + ) 1515 + }) 1516 + .for_each(|edge| store.add(edge)); 1517 + ["alpha", "bravo", "charlie"] 1518 + .iter() 1519 + .map(|s| { 1520 + EdgeKey::new( 1521 + nsid("sh.tangled.feed.star"), 1522 + did_subj(&format!("did:plc:{s}")), 1523 + ) 1524 + }) 1525 + .collect() 1526 + } 1527 + 1528 + #[test] 1529 + fn list_multi_merges_subjects_newest_first() { 1530 + let store = store(); 1531 + let keys = fill_three_subjects(&store); 1532 + let page = store.list_multi(&keys, PageCursor::Start, limit(10), SortDir::Desc); 1533 + let micros: Vec<u64> = page.items.iter().map(|it| it.sort_micros).collect(); 1534 + assert_eq!(micros, vec![60, 50, 40, 30, 20, 10]); 1535 + assert!(page.next.is_none()); 1536 + } 1537 + 1538 + #[test] 1539 + fn list_multi_paginates_disjoint_across_buckets() { 1540 + let store = store(); 1541 + let keys = fill_three_subjects(&store); 1542 + // limit=4 with 6 items across 3 buckets exercises the bounded top-N merge. 1543 + let p1 = store.list_multi(&keys, PageCursor::Start, limit(4), SortDir::Desc); 1544 + let m1: Vec<u64> = p1.items.iter().map(|it| it.sort_micros).collect(); 1545 + assert_eq!(m1, vec![60, 50, 40, 30]); 1546 + let tok = p1.next.expect("first page has more"); 1547 + let p2 = store.list_multi(&keys, PageCursor::After(tok), limit(4), SortDir::Desc); 1548 + let m2: Vec<u64> = p2.items.iter().map(|it| it.sort_micros).collect(); 1549 + assert_eq!(m2, vec![20, 10]); 1550 + assert!(p2.next.is_none()); 1551 + } 1552 + 1553 + #[test] 1554 + fn list_multi_dedups_source_under_two_keys() { 1555 + let store = store(); 1556 + // Same record (source) indexed under two subject keys, same createdAt. 1557 + let shared = at("at://did:plc:x/sh.tangled.feed.star/shared"); 1558 + store.add(star_edge_at(shared.clone(), did("did:plc:alpha"), 100)); 1559 + store.add(star_edge_at(shared, did("did:plc:bravo"), 100)); 1560 + let keys = vec![ 1561 + EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:alpha")), 1562 + EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:bravo")), 1563 + ]; 1564 + let page = store.list_multi(&keys, PageCursor::Start, limit(10), SortDir::Desc); 1565 + assert_eq!(page.items.len(), 1); 1566 + assert!(page.next.is_none()); 1441 1567 } 1442 1568 1443 1569 #[test]