This repository has no description
0

Configure Feed

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

core / bobbin / crates / edge-index / src / lib.rs
54 kB 1642 lines
1use std::collections::{BTreeSet, HashMap}; 2use std::marker::PhantomData; 3use std::num::NonZeroU32; 4use std::ops::{Bound, ControlFlow}; 5use std::sync::atomic::{AtomicU32, Ordering}; 6use std::sync::{Arc, Mutex}; 7 8const FILTER_SCAN_MULTIPLIER: usize = 64; 9const FILTER_SCAN_FLOOR: usize = 512; 10 11struct ScanState { 12 matched: Vec<(BucketKey, AtUri<DefaultStr>)>, 13 last_scanned: Option<BucketKey>, 14 scanned: usize, 15} 16 17impl ScanState { 18 fn with_capacity(cap: usize) -> Self { 19 Self { 20 matched: Vec::with_capacity(cap), 21 last_scanned: None, 22 scanned: 0, 23 } 24 } 25} 26 27use bobbin_runtime::RuntimeHasher; 28use bobbin_types::edges::Edge; 29use bobbin_types::ids::EdgeKey; 30use either::Either; 31use jacquard_common::DefaultStr; 32use jacquard_common::types::string::AtUri; 33use lasso::{Key, Spur, ThreadedRodeo}; 34use scc::HashMap as SccMap; 35use scc::hash_map::Entry; 36use smallvec::SmallVec; 37use thiserror::Error; 38 39#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq, Ord, PartialOrd)] 40struct SortMicros(u64); 41 42#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq, Ord, PartialOrd)] 43struct BucketKey { 44 micros: SortMicros, 45 source: SourceId, 46} 47 48impl BucketKey { 49 fn new(micros: u64, source: SourceId) -> Self { 50 Self { 51 micros: SortMicros(micros), 52 source, 53 } 54 } 55 56 fn token(self) -> PageToken { 57 PageToken::new(self.micros.0, self.source.index()) 58 } 59 60 fn from_token(tok: PageToken) -> Self { 61 Self { 62 micros: SortMicros(tok.micros), 63 source: SourceId::from_raw(tok.source), 64 } 65 } 66} 67 68pub mod coverage; 69pub mod state_index; 70pub use coverage::{Coverage, CoverageWatch, HydrantCursor, PromotionSignal}; 71pub use state_index::{ 72 ApplyOutcome, IssueStateKind, PullStatusKind, StateIndex, StateKind, apply_record_state, 73}; 74 75#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 76pub struct SourceTag; 77#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 78struct AuthorTag; 79#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 80struct CollectionTag; 81 82#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 83pub struct Interned<T>(u32, PhantomData<T>); 84 85impl<T: Copy> Interned<T> { 86 fn from_spur(spur: Spur) -> Self { 87 Self(spur.into_usize() as u32, PhantomData) 88 } 89 90 fn from_raw(raw: u32) -> Self { 91 Self(raw, PhantomData) 92 } 93 94 fn index(self) -> u32 { 95 self.0 96 } 97 98 fn to_spur(self) -> Option<Spur> { 99 Spur::try_from_usize(self.0 as usize) 100 } 101} 102 103pub type SourceId = Interned<SourceTag>; 104type AuthorId = Interned<AuthorTag>; 105type CollectionId = Interned<CollectionTag>; 106 107#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 108pub struct PageToken { 109 micros: u64, 110 source: u32, 111} 112 113impl PageToken { 114 pub fn new(micros: u64, source: u32) -> Self { 115 Self { micros, source } 116 } 117 118 pub fn micros(self) -> u64 { 119 self.micros 120 } 121 122 pub fn source(self) -> u32 { 123 self.source 124 } 125 126 pub fn encode_token(self) -> String { 127 let mut bytes = [0u8; 12]; 128 bytes[..8].copy_from_slice(&self.micros.to_be_bytes()); 129 bytes[8..].copy_from_slice(&self.source.to_be_bytes()); 130 encode_hex(&bytes) 131 } 132 133 pub fn decode_token(token: &str) -> Result<Self, CursorParseError> { 134 let bytes: [u8; 12] = decode_hex_array(token).ok_or(CursorParseError::Malformed)?; 135 let micros = u64::from_be_bytes(bytes[..8].try_into().unwrap()); 136 let source = u32::from_be_bytes(bytes[8..].try_into().unwrap()); 137 Ok(Self { micros, source }) 138 } 139} 140 141fn encode_hex(bytes: &[u8]) -> String { 142 bytes 143 .iter() 144 .fold(String::with_capacity(bytes.len() * 2), |mut acc, b| { 145 acc.push(char::from_digit((b >> 4) as u32, 16).unwrap()); 146 acc.push(char::from_digit((b & 0x0f) as u32, 16).unwrap()); 147 acc 148 }) 149} 150 151fn decode_hex_array<const N: usize>(token: &str) -> Option<[u8; N]> { 152 if token.len() != N * 2 { 153 return None; 154 } 155 let parsed: Vec<u8> = token 156 .as_bytes() 157 .chunks_exact(2) 158 .map(|pair| { 159 let hi = (pair[0] as char).to_digit(16)?; 160 let lo = (pair[1] as char).to_digit(16)?; 161 Some(((hi << 4) | lo) as u8) 162 }) 163 .collect::<Option<Vec<u8>>>()?; 164 parsed.try_into().ok() 165} 166 167#[derive(Clone, Copy, Debug, Eq, PartialEq)] 168pub enum PageCursor { 169 Start, 170 After(PageToken), 171} 172 173#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] 174pub enum SortDir { 175 Asc, 176 #[default] 177 Desc, 178} 179 180#[derive(Clone, Copy, Debug, Eq, PartialEq, Error)] 181pub enum CursorParseError { 182 #[error("cursor token must be a valid TID")] 183 Malformed, 184} 185 186impl PageCursor { 187 pub fn from_token(raw: Option<&str>) -> Result<Self, CursorParseError> { 188 raw.map_or(Ok(Self::Start), |t| { 189 PageToken::decode_token(t).map(Self::After) 190 }) 191 } 192} 193 194#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 195pub struct PageLimit(u32); 196 197#[derive(Clone, Copy, Debug, Eq, PartialEq, Error)] 198pub enum PageLimitError { 199 #[error("page limit {value} below minimum {min}")] 200 TooSmall { value: u32, min: u32 }, 201 #[error("page limit {value} above maximum {max}")] 202 TooLarge { value: u32, max: u32 }, 203} 204 205impl PageLimit { 206 pub const MIN: u32 = 1; 207 pub const MAX: u32 = 1000; 208 209 pub fn new(value: u32) -> Result<Self, PageLimitError> { 210 match value { 211 v if v < Self::MIN => Err(PageLimitError::TooSmall { 212 value: v, 213 min: Self::MIN, 214 }), 215 v if v > Self::MAX => Err(PageLimitError::TooLarge { 216 value: v, 217 max: Self::MAX, 218 }), 219 v => Ok(Self(v)), 220 } 221 } 222 223 pub const fn get(self) -> u32 { 224 self.0 225 } 226} 227 228#[derive(Debug)] 229pub struct EdgePage { 230 pub items: Vec<EdgeItem>, 231 pub next: Option<PageToken>, 232} 233 234#[derive(Clone, Debug, Eq, PartialEq)] 235pub struct EdgeItem { 236 pub uri: AtUri<DefaultStr>, 237 pub sort_micros: u64, 238} 239 240impl AsRef<str> for EdgeItem { 241 fn as_ref(&self) -> &str { 242 self.uri.as_ref() 243 } 244} 245 246#[derive(Clone, Copy, Debug, Default)] 247pub struct EdgeMemReport { 248 pub key_count: u64, 249 pub source_count: u64, 250 pub source_interner_bytes: u64, 251 pub did_interner_bytes: u64, 252 pub collection_interner_bytes: u64, 253 pub key_interner_bytes: u64, 254 pub edges_total: u64, 255 pub author_refs_total: u64, 256 pub reverse_entries: u64, 257 pub reverse_cap: u64, 258 pub forward_struct_bytes: u64, 259 pub reverse_struct_bytes: u64, 260 pub bucket_struct_bytes: u64, 261 pub max_bucket: u64, 262 pub bucket_size_classes: [u64; BUCKET_CLASS_COUNT], 263} 264 265const BUCKET_CLASS_BOUNDS: [u64; 16] = [ 266 1, 267 2, 268 4, 269 8, 270 16, 271 32, 272 64, 273 128, 274 256, 275 512, 276 1024, 277 2048, 278 8192, 279 32768, 280 131072, 281 u64::MAX, 282]; 283const BUCKET_CLASS_COUNT: usize = BUCKET_CLASS_BOUNDS.len(); 284 285fn bucket_class(n: u64) -> usize { 286 BUCKET_CLASS_BOUNDS 287 .iter() 288 .position(|&bound| n <= bound) 289 .unwrap_or(BUCKET_CLASS_COUNT - 1) 290} 291 292impl EdgeMemReport { 293 pub fn bucket_histogram(&self) -> impl Iterator<Item = (u64, u64)> { 294 BUCKET_CLASS_BOUNDS 295 .into_iter() 296 .zip(self.bucket_size_classes) 297 } 298} 299 300#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 301struct EdgeKeyId(u32); 302 303fn bump_author(authors: &mut HashMap<AuthorId, NonZeroU32, RuntimeHasher>, author: AuthorId) { 304 authors 305 .entry(author) 306 .and_modify(|c| *c = c.saturating_add(1)) 307 .or_insert(NonZeroU32::MIN); 308} 309 310fn drop_author(authors: &mut HashMap<AuthorId, NonZeroU32, RuntimeHasher>, author: AuthorId) { 311 match authors.get(&author).map(|c| c.get() - 1) { 312 Some(0) | None => { 313 authors.remove(&author); 314 } 315 Some(next) => { 316 authors.insert(author, NonZeroU32::new(next).unwrap()); 317 } 318 } 319} 320 321struct LargeBucket { 322 keys: BTreeSet<BucketKey>, 323 authors: HashMap<AuthorId, NonZeroU32, RuntimeHasher>, 324} 325 326const BUCKET_PROMOTE_AT: usize = 256; 327const SOURCE_BYTES: u64 = 16; 328 329enum Sources { 330 Small(SmallVec<[BucketKey; 2]>), 331 Large(Box<LargeBucket>), 332} 333 334impl Default for Sources { 335 fn default() -> Self { 336 Self::Small(SmallVec::new()) 337 } 338} 339 340impl Sources { 341 fn len(&self) -> usize { 342 match self { 343 Self::Small(v) => v.len(), 344 Self::Large(big) => big.keys.len(), 345 } 346 } 347 348 fn is_empty(&self) -> bool { 349 self.len() == 0 350 } 351 352 fn insert(&mut self, key: BucketKey, author: Option<AuthorId>) -> bool { 353 match self { 354 Self::Small(v) => match v.binary_search(&key) { 355 Ok(_) => false, 356 Err(pos) => { 357 v.insert(pos, key); 358 true 359 } 360 }, 361 Self::Large(big) => { 362 let inserted = big.keys.insert(key); 363 if inserted && let Some(a) = author { 364 bump_author(&mut big.authors, a); 365 } 366 inserted 367 } 368 } 369 } 370 371 fn remove(&mut self, key: &BucketKey, author: Option<AuthorId>) { 372 match self { 373 Self::Small(v) => { 374 if let Ok(pos) = v.binary_search(key) { 375 v.remove(pos); 376 } 377 } 378 Self::Large(big) => { 379 if big.keys.remove(key) 380 && let Some(a) = author 381 { 382 drop_author(&mut big.authors, a); 383 } 384 } 385 } 386 } 387 388 fn directed( 389 &self, 390 cursor: PageCursor, 391 dir: SortDir, 392 ) -> Box<dyn Iterator<Item = BucketKey> + '_> { 393 match self { 394 Self::Small(v) => Box::new(directed_slice(v, cursor, dir)), 395 Self::Large(big) => Box::new(directed_tree(&big.keys, cursor, dir)), 396 } 397 } 398 399 fn heap_bytes(&self) -> u64 { 400 const BTREE_BYTES_PER_KEY: u64 = 32; 401 const HASHMAP_FIXED: u64 = 48; 402 const HASHMAP_PER_CAP: u64 = 9; 403 match self { 404 Self::Small(v) if v.spilled() => v.capacity() as u64 * SOURCE_BYTES, 405 Self::Small(_) => 0, 406 Self::Large(big) => { 407 std::mem::size_of::<LargeBucket>() as u64 408 + big.keys.len() as u64 * BTREE_BYTES_PER_KEY 409 + HASHMAP_FIXED 410 + big.authors.capacity() as u64 * HASHMAP_PER_CAP 411 } 412 } 413 } 414} 415 416#[derive(Clone, Copy, Debug, Eq, PartialEq)] 417struct ReverseEntry { 418 key_id: EdgeKeyId, 419 sort_micros: u64, 420} 421 422pub struct EdgeStore { 423 source_interner: Arc<ThreadedRodeo<Spur, RuntimeHasher>>, 424 did_interner: Arc<ThreadedRodeo<Spur, RuntimeHasher>>, 425 collection_interner: Arc<ThreadedRodeo<Spur, RuntimeHasher>>, 426 key_ids: SccMap<EdgeKey, EdgeKeyId, RuntimeHasher>, 427 next_key_id: AtomicU32, 428 forward: SccMap<EdgeKeyId, Sources, RuntimeHasher>, 429 reverse: SccMap<SourceId, SmallVec<[ReverseEntry; 1]>, RuntimeHasher>, 430 hasher: RuntimeHasher, 431 writer: Mutex<()>, 432} 433 434impl EdgeStore { 435 pub fn new(hasher: RuntimeHasher) -> Self { 436 Self { 437 source_interner: Arc::new(ThreadedRodeo::with_hasher(hasher.clone())), 438 did_interner: Arc::new(ThreadedRodeo::with_hasher(hasher.clone())), 439 collection_interner: Arc::new(ThreadedRodeo::with_hasher(hasher.clone())), 440 key_ids: SccMap::with_hasher(hasher.clone()), 441 next_key_id: AtomicU32::new(0), 442 forward: SccMap::with_hasher(hasher.clone()), 443 reverse: SccMap::with_hasher(hasher.clone()), 444 hasher, 445 writer: Mutex::new(()), 446 } 447 } 448 449 fn intern_key(&self, key: EdgeKey) -> EdgeKeyId { 450 match self.key_ids.entry_sync(key) { 451 Entry::Occupied(e) => *e.get(), 452 Entry::Vacant(e) => { 453 let id = EdgeKeyId(self.next_key_id.fetch_add(1, Ordering::Relaxed)); 454 e.insert_entry(id); 455 id 456 } 457 } 458 } 459 460 fn lookup_key(&self, key: &EdgeKey) -> Option<EdgeKeyId> { 461 self.key_ids.read_sync(key, |_, id| *id) 462 } 463 464 pub fn add(&self, edge: Edge) { 465 let _w = self 466 .writer 467 .lock() 468 .expect("edge-store writer mutex poisoned"); 469 self.add_locked(edge); 470 } 471 472 pub fn intern_source(&self, source: &AtUri<DefaultStr>) -> SourceId { 473 let author = self.intern_author(source); 474 SourceId::from_spur( 475 self.source_interner 476 .get_or_intern(self.source_key(source, author)), 477 ) 478 } 479 480 pub fn upsert_source(&self, source: &AtUri<DefaultStr>, edges: Vec<Edge>) { 481 let _w = self 482 .writer 483 .lock() 484 .expect("edge-store writer mutex poisoned"); 485 self.clear_source_locked(source); 486 edges.into_iter().for_each(|e| self.add_locked(e)); 487 } 488 489 pub fn remove_source(&self, source: &AtUri<DefaultStr>) { 490 let _w = self 491 .writer 492 .lock() 493 .expect("edge-store writer mutex poisoned"); 494 self.clear_source_locked(source); 495 } 496 497 fn add_locked(&self, edge: Edge) { 498 let author = self.intern_author(&edge.source); 499 let source_key = self.source_key(&edge.source, author); 500 let id = SourceId::from_spur(self.source_interner.get_or_intern(source_key)); 501 let sort_micros = edge.sort_micros; 502 let key_id = self.intern_key(EdgeKey::new(edge.kind, edge.subject)); 503 504 let key = BucketKey::new(sort_micros, id); 505 let mut entry = self.forward.entry_sync(key_id).or_default(); 506 let inserted = entry.get_mut().insert(key, author); 507 let promote_keys = match entry.get() { 508 Sources::Small(v) if v.len() > BUCKET_PROMOTE_AT => { 509 Some(v.iter().copied().collect::<Vec<_>>()) 510 } 511 _ => None, 512 }; 513 drop(entry); 514 515 if let Some(keys) = promote_keys { 516 let authors = self.build_author_map(&keys); 517 let large = LargeBucket { 518 keys: keys.into_iter().collect(), 519 authors, 520 }; 521 if let Entry::Occupied(mut e) = self.forward.entry_sync(key_id) { 522 *e.get_mut() = Sources::Large(Box::new(large)); 523 } 524 } 525 526 if inserted { 527 let mut rev = self.reverse.entry_sync(id).or_default(); 528 rev.get_mut().push(ReverseEntry { 529 key_id, 530 sort_micros, 531 }); 532 } 533 } 534 535 fn clear_source_locked(&self, source: &AtUri<DefaultStr>) { 536 let author = source_authority_did(source) 537 .and_then(|s| self.did_interner.get(s)) 538 .map(AuthorId::from_spur); 539 let Some(source_spur) = self.source_interner.get(self.source_key(source, author)) else { 540 return; 541 }; 542 let id = SourceId::from_spur(source_spur); 543 let Some((_, entries)) = self.reverse.remove_sync(&id) else { 544 return; 545 }; 546 entries.into_iter().for_each( 547 |ReverseEntry { 548 key_id, 549 sort_micros, 550 }| { 551 self.forward.update_sync(&key_id, |_, sources| { 552 sources.remove(&BucketKey::new(sort_micros, id), author); 553 }); 554 self.forward 555 .remove_if_sync(&key_id, |sources| sources.is_empty()); 556 }, 557 ); 558 } 559 560 fn intern_author(&self, source: &AtUri<DefaultStr>) -> Option<AuthorId> { 561 let did = source_authority_did(source)?; 562 Some(AuthorId::from_spur(self.did_interner.get_or_intern(did))) 563 } 564 565 fn source_key(&self, source: &AtUri<DefaultStr>, author: Option<AuthorId>) -> String { 566 match (split_record_uri(source.as_ref()), author) { 567 (Some((_, collection, rkey)), Some(author)) => { 568 let collection = 569 CollectionId::from_spur(self.collection_interner.get_or_intern(collection)); 570 format!("{}/{}/{}", author.index(), collection.index(), rkey) 571 } 572 _ => source.as_ref().to_owned(), 573 } 574 } 575 576 fn decode_source(&self, stored: &str) -> Option<String> { 577 if stored.starts_with("at://") { 578 return Some(stored.to_owned()); 579 } 580 let mut parts = stored.splitn(3, '/'); 581 let author = AuthorId::from_raw(parts.next()?.parse().ok()?); 582 let collection = CollectionId::from_raw(parts.next()?.parse().ok()?); 583 let rkey = parts.next()?; 584 let did = self.did_interner.try_resolve(&author.to_spur()?)?; 585 let collection = self 586 .collection_interner 587 .try_resolve(&collection.to_spur()?)?; 588 Some(format!("at://{did}/{collection}/{rkey}")) 589 } 590 591 fn author_of_stored(&self, stored: &str) -> Option<AuthorId> { 592 match stored.strip_prefix("at://") { 593 Some(rest) => { 594 let authority = rest.split('/').next().unwrap_or(rest); 595 authority 596 .starts_with("did:") 597 .then(|| self.did_interner.get(authority).map(AuthorId::from_spur)) 598 .flatten() 599 } 600 None => stored 601 .split('/') 602 .next()? 603 .parse::<u32>() 604 .ok() 605 .map(AuthorId::from_raw), 606 } 607 } 608 609 fn author_of(&self, source: SourceId) -> Option<AuthorId> { 610 let spur = source.to_spur()?; 611 let stored = self.source_interner.try_resolve(&spur)?; 612 self.author_of_stored(stored) 613 } 614 615 fn distinct_authors_small(&self, keys: &[BucketKey]) -> u64 { 616 keys.iter() 617 .filter_map(|key| self.author_of(key.source)) 618 .collect::<std::collections::HashSet<AuthorId>>() 619 .len() as u64 620 } 621 622 fn build_author_map(&self, keys: &[BucketKey]) -> HashMap<AuthorId, NonZeroU32, RuntimeHasher> { 623 keys.iter().fold( 624 HashMap::with_hasher(self.hasher.clone()), 625 |mut authors, key| { 626 if let Some(a) = self.author_of(key.source) { 627 bump_author(&mut authors, a); 628 } 629 authors 630 }, 631 ) 632 } 633 634 pub fn count(&self, key: &EdgeKey) -> u64 { 635 self.lookup_key(key) 636 .and_then(|id| { 637 self.forward 638 .read_sync(&id, |_, sources| sources.len() as u64) 639 }) 640 .unwrap_or(0) 641 } 642 643 pub fn count_distinct_authors(&self, key: &EdgeKey) -> u64 { 644 self.lookup_key(key) 645 .and_then(|id| { 646 self.forward.read_sync(&id, |_, sources| match sources { 647 Sources::Large(big) => big.authors.len() as u64, 648 Sources::Small(v) => self.distinct_authors_small(v), 649 }) 650 }) 651 .unwrap_or(0) 652 } 653 /// answers "did the viewer star/follow/etc. this subject, and with what rkey" 654 pub fn viewer_source(&self, key: &EdgeKey, viewer: &str) -> Option<AtUri<DefaultStr>> { 655 let author_spur = self.did_interner.get(viewer)?; 656 let author_id = AuthorId::from_spur(author_spur); 657 let key_id = self.lookup_key(key)?; 658 659 self.forward 660 .read_sync(&key_id, |_, sources| { 661 // large buckets track authors, so a missing author fast-fails the scan 662 if let Sources::Large(big) = sources 663 && !big.authors.contains_key(&author_id) 664 { 665 return None; 666 } 667 sources 668 .directed(PageCursor::Start, SortDir::Desc) 669 .find_map(|bucket| { 670 let spur = bucket.source.to_spur()?; 671 let stored = self.source_interner.try_resolve(&spur)?; 672 let uri = AtUri::new_owned(self.decode_source(stored)?).ok()?; 673 (source_authority_did(&uri) == Some(viewer)).then_some(uri) 674 }) 675 }) 676 .flatten() 677 } 678 679 pub fn sources_for(&self, key: &EdgeKey) -> Vec<AtUri<DefaultStr>> { 680 self.lookup_key(key) 681 .and_then(|id| { 682 self.forward.read_sync(&id, |_, sources| { 683 sources 684 .directed(PageCursor::Start, SortDir::Desc) 685 .filter_map(|bucket| { 686 let spur = bucket.source.to_spur()?; 687 let stored = self.source_interner.try_resolve(&spur)?; 688 AtUri::new_owned(self.decode_source(stored)?).ok() 689 }) 690 .collect::<Vec<_>>() 691 }) 692 }) 693 .unwrap_or_default() 694 } 695 696 pub fn list( 697 &self, 698 key: &EdgeKey, 699 cursor: PageCursor, 700 limit: PageLimit, 701 dir: SortDir, 702 ) -> EdgePage { 703 let limit_usize = limit.get() as usize; 704 self.lookup_key(key) 705 .and_then(|id| { 706 self.forward.read_sync(&id, |_, sources| { 707 let iter = sources.directed(cursor, dir); 708 let entries: Vec<BucketKey> = iter.take(limit_usize + 1).collect(); 709 let has_more = entries.len() > limit_usize; 710 let page = &entries[..entries.len().min(limit_usize)]; 711 let items = page 712 .iter() 713 .filter_map(|&key| { 714 let spur = key.source.to_spur()?; 715 let stored = self.source_interner.try_resolve(&spur)?; 716 let uri = AtUri::new_owned(self.decode_source(stored)?).ok()?; 717 Some(EdgeItem { 718 uri, 719 sort_micros: key.micros.0, 720 }) 721 }) 722 .collect(); 723 let next = has_more 724 .then(|| page.last().copied()) 725 .flatten() 726 .map(BucketKey::token); 727 EdgePage { items, next } 728 }) 729 }) 730 .unwrap_or(EdgePage { 731 items: Vec::new(), 732 next: None, 733 }) 734 } 735 736 pub fn list_filtered<F>( 737 &self, 738 key: &EdgeKey, 739 cursor: PageCursor, 740 limit: PageLimit, 741 dir: SortDir, 742 predicate: F, 743 ) -> EdgePage 744 where 745 F: Fn(&AtUri<DefaultStr>) -> bool, 746 { 747 let limit_usize = limit.get() as usize; 748 let scan_cap = limit_usize 749 .saturating_mul(FILTER_SCAN_MULTIPLIER) 750 .max(FILTER_SCAN_FLOOR); 751 self.lookup_key(key) 752 .and_then(|id| { 753 self.forward.read_sync(&id, |_, sources| { 754 let init = ScanState::with_capacity(limit_usize + 1); 755 let outcome = sources 756 .directed(cursor, dir) 757 .try_fold(init, |mut state, key| { 758 if state.scanned >= scan_cap && state.matched.len() <= limit_usize { 759 return ControlFlow::Break(state); 760 } 761 state.scanned += 1; 762 state.last_scanned = Some(key); 763 if let Some(spur) = key.source.to_spur() 764 && let Some(stored) = self.source_interner.try_resolve(&spur) 765 && let Some(decoded) = self.decode_source(stored) 766 && let Ok(uri) = AtUri::new_owned(decoded) 767 && predicate(&uri) 768 { 769 state.matched.push((key, uri)); 770 if state.matched.len() > limit_usize { 771 return ControlFlow::Break(state); 772 } 773 } 774 ControlFlow::Continue(state) 775 }); 776 let (state, bucket_exhausted) = match outcome { 777 ControlFlow::Continue(s) => (s, true), 778 ControlFlow::Break(s) => (s, false), 779 }; 780 let has_more_matches = state.matched.len() > limit_usize; 781 let visible_len = state.matched.len().min(limit_usize); 782 let next = if has_more_matches && visible_len > 0 { 783 Some(state.matched[visible_len - 1].0.token()) 784 } else if !bucket_exhausted { 785 state.last_scanned.map(BucketKey::token) 786 } else { 787 None 788 }; 789 let items = state 790 .matched 791 .into_iter() 792 .take(visible_len) 793 .map(|(key, uri)| EdgeItem { 794 uri, 795 sort_micros: key.micros.0, 796 }) 797 .collect(); 798 EdgePage { items, next } 799 }) 800 }) 801 .unwrap_or(EdgePage { 802 items: Vec::new(), 803 next: None, 804 }) 805 } 806 807 pub fn key_count(&self) -> usize { 808 self.forward.len() 809 } 810 811 pub fn source_count(&self) -> usize { 812 self.reverse.len() 813 } 814 815 pub fn mem_report(&self) -> EdgeMemReport { 816 const SCC_SLOT: u64 = 32; 817 let edge_key = std::mem::size_of::<EdgeKey>() as u64; 818 let rev_entry = std::mem::size_of::<ReverseEntry>() as u64; 819 let rev_smallvec = std::mem::size_of::<SmallVec<[ReverseEntry; 1]>>() as u64; 820 let bucket_struct = std::mem::size_of::<Sources>() as u64; 821 822 let mut edges_total = 0u64; 823 let mut author_refs_total = 0u64; 824 let mut forward_struct_bytes = 0u64; 825 let mut max_bucket = 0u64; 826 let mut bucket_size_classes = [0u64; BUCKET_CLASS_COUNT]; 827 self.forward.iter_sync(|_, sources| { 828 let bucket_len = sources.len() as u64; 829 edges_total += bucket_len; 830 max_bucket = max_bucket.max(bucket_len); 831 bucket_size_classes[bucket_class(bucket_len)] += 1; 832 if let Sources::Large(big) = sources { 833 author_refs_total += big.authors.len() as u64; 834 } 835 forward_struct_bytes += SCC_SLOT + bucket_struct + sources.heap_bytes(); 836 true 837 }); 838 let key_interner_bytes = self.key_ids.len() as u64 839 * (edge_key + std::mem::size_of::<EdgeKeyId>() as u64 + SCC_SLOT); 840 841 let mut reverse_entries = 0u64; 842 let mut reverse_cap = 0u64; 843 let mut reverse_struct_bytes = 0u64; 844 self.reverse.iter_sync(|_, refs| { 845 let cap = refs.capacity() as u64; 846 reverse_entries += refs.len() as u64; 847 reverse_cap += cap; 848 let heap = if refs.spilled() { cap * rev_entry } else { 0 }; 849 reverse_struct_bytes += SCC_SLOT + rev_smallvec + heap; 850 true 851 }); 852 853 EdgeMemReport { 854 key_count: self.forward.len() as u64, 855 source_count: self.reverse.len() as u64, 856 source_interner_bytes: self.source_interner.current_memory_usage() as u64, 857 did_interner_bytes: self.did_interner.current_memory_usage() as u64, 858 collection_interner_bytes: self.collection_interner.current_memory_usage() as u64, 859 key_interner_bytes, 860 edges_total, 861 author_refs_total, 862 reverse_entries, 863 reverse_cap, 864 forward_struct_bytes, 865 reverse_struct_bytes, 866 bucket_struct_bytes: bucket_struct, 867 max_bucket, 868 bucket_size_classes, 869 } 870 } 871} 872 873fn directed_slice( 874 sources: &[BucketKey], 875 cursor: PageCursor, 876 dir: SortDir, 877) -> impl Iterator<Item = BucketKey> + '_ { 878 match dir { 879 SortDir::Asc => { 880 let start = match cursor { 881 PageCursor::Start => 0, 882 PageCursor::After(tok) => { 883 sources.partition_point(|k| *k <= BucketKey::from_token(tok)) 884 } 885 }; 886 Either::Left(sources[start..].iter().copied()) 887 } 888 SortDir::Desc => { 889 let end = match cursor { 890 PageCursor::Start => sources.len(), 891 PageCursor::After(tok) => { 892 sources.partition_point(|k| *k < BucketKey::from_token(tok)) 893 } 894 }; 895 Either::Right(sources[..end].iter().rev().copied()) 896 } 897 } 898} 899 900fn directed_tree( 901 keys: &BTreeSet<BucketKey>, 902 cursor: PageCursor, 903 dir: SortDir, 904) -> impl Iterator<Item = BucketKey> + '_ { 905 match dir { 906 SortDir::Asc => { 907 let lower = match cursor { 908 PageCursor::Start => Bound::Unbounded, 909 PageCursor::After(tok) => Bound::Excluded(BucketKey::from_token(tok)), 910 }; 911 Either::Left(keys.range((lower, Bound::Unbounded)).copied()) 912 } 913 SortDir::Desc => { 914 let upper = match cursor { 915 PageCursor::Start => Bound::Unbounded, 916 PageCursor::After(tok) => Bound::Excluded(BucketKey::from_token(tok)), 917 }; 918 Either::Right(keys.range((Bound::Unbounded, upper)).rev().copied()) 919 } 920 } 921} 922 923fn split_record_uri(source: &str) -> Option<(&str, &str, &str)> { 924 let rest = source.strip_prefix("at://")?; 925 let mut parts = rest.split('/'); 926 let authority = parts.next()?; 927 let collection = parts.next()?; 928 let rkey = parts.next()?; 929 if parts.next().is_some() { 930 return None; 931 } 932 (authority.starts_with("did:") && !collection.is_empty() && !rkey.is_empty()) 933 .then_some((authority, collection, rkey)) 934} 935 936fn source_authority_did(source: &AtUri<DefaultStr>) -> Option<&str> { 937 let rest = source.as_ref().strip_prefix("at://")?; 938 let end = rest.find('/').unwrap_or(rest.len()); 939 let candidate = &rest[..end]; 940 candidate.starts_with("did:").then_some(candidate) 941} 942 943#[cfg(test)] 944mod tests { 945 use super::*; 946 use bobbin_types::ids::SubjectRef; 947 use jacquard_common::types::did::Did; 948 use jacquard_common::types::nsid::Nsid; 949 use jacquard_common::types::string::AtUri; 950 951 fn store() -> EdgeStore { 952 EdgeStore::new(RuntimeHasher::default()) 953 } 954 955 fn nsid(s: &'static str) -> Nsid<DefaultStr> { 956 Nsid::new_static(s).unwrap() 957 } 958 959 fn at(s: &str) -> AtUri<DefaultStr> { 960 AtUri::new_owned(s).unwrap() 961 } 962 963 fn did(s: &str) -> Did<DefaultStr> { 964 Did::new_owned(s).unwrap() 965 } 966 967 fn did_subj(s: &str) -> SubjectRef { 968 SubjectRef::Did(did(s)) 969 } 970 971 fn limit(n: u32) -> PageLimit { 972 PageLimit::new(n).unwrap() 973 } 974 975 const NAMES: [&str; 5] = ["nel", "olaren", "teq", "lyna", "bailey"]; 976 977 fn star_edge(source: AtUri<DefaultStr>, subject: Did<DefaultStr>) -> Edge { 978 star_edge_at(source, subject, 0) 979 } 980 981 fn star_edge_at(source: AtUri<DefaultStr>, subject: Did<DefaultStr>, sort_micros: u64) -> Edge { 982 Edge { 983 kind: nsid("sh.tangled.feed.star"), 984 subject: SubjectRef::Did(subject), 985 source, 986 sort_micros, 987 } 988 } 989 990 fn shuffled_micros(i: usize) -> u64 { 991 (i as u64).wrapping_mul(2_654_435_761) % 64 992 } 993 994 fn source_uri(i: usize) -> String { 995 format!( 996 "at://did:plc:{}/sh.tangled.feed.star/r{i}", 997 NAMES[i % NAMES.len()] 998 ) 999 } 1000 1001 fn fill_subject(store: &EdgeStore, n: usize) -> EdgeKey { 1002 let subject = did("did:plc:squid"); 1003 (0..n).for_each(|i| { 1004 store.add(star_edge_at( 1005 at(&source_uri(i)), 1006 subject.clone(), 1007 shuffled_micros(i), 1008 )); 1009 }); 1010 EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:squid")) 1011 } 1012 1013 fn reference(kept: impl Iterator<Item = usize>) -> Vec<String> { 1014 let mut rows: Vec<(u64, usize)> = kept.map(|i| (shuffled_micros(i), i)).collect(); 1015 rows.sort_unstable(); 1016 rows.into_iter().map(|(_, i)| source_uri(i)).collect() 1017 } 1018 1019 fn paginate_all(store: &EdgeStore, key: &EdgeKey, dir: SortDir, page: u32) -> Vec<String> { 1020 std::iter::successors( 1021 Some(store.list(key, PageCursor::Start, limit(page), dir)), 1022 |prev| { 1023 prev.next 1024 .map(|tok| store.list(key, PageCursor::After(tok), limit(page), dir)) 1025 }, 1026 ) 1027 .flat_map(|p| { 1028 p.items 1029 .into_iter() 1030 .map(|u| u.as_ref().to_owned()) 1031 .collect::<Vec<_>>() 1032 }) 1033 .collect() 1034 } 1035 1036 #[test] 1037 fn pagination_matches_reference_across_small_and_large() { 1038 [64usize, 1000].into_iter().for_each(|n| { 1039 let store = store(); 1040 let key = fill_subject(&store, n); 1041 assert_eq!(store.count(&key), n as u64, "count mismatch at n={n}"); 1042 assert_eq!( 1043 store.count_distinct_authors(&key), 1044 NAMES.len() as u64, 1045 "distinct authors at n={n}" 1046 ); 1047 1048 let asc = reference(0..n); 1049 let desc: Vec<String> = asc.iter().rev().cloned().collect(); 1050 1051 [3u32, 7, 50].into_iter().for_each(|page| { 1052 assert_eq!( 1053 paginate_all(&store, &key, SortDir::Asc, page), 1054 asc, 1055 "asc mismatch n={n} page={page}" 1056 ); 1057 assert_eq!( 1058 paginate_all(&store, &key, SortDir::Desc, page), 1059 desc, 1060 "desc mismatch n={n} page={page}" 1061 ); 1062 }); 1063 }); 1064 } 1065 1066 #[test] 1067 fn remove_from_large_bucket_keeps_pagination_exact() { 1068 let store = store(); 1069 let n = 1000usize; 1070 let key = fill_subject(&store, n); 1071 (0..n).step_by(3).for_each(|i| { 1072 store.remove_source(&at(&source_uri(i))); 1073 }); 1074 let expected = reference((0..n).filter(|i| i % 3 != 0)); 1075 assert_eq!(store.count(&key), expected.len() as u64); 1076 assert_eq!( 1077 store.count_distinct_authors(&key), 1078 NAMES.len() as u64, 1079 "every author keeps sources after partial removal" 1080 ); 1081 assert_eq!(paginate_all(&store, &key, SortDir::Asc, 7), expected); 1082 assert_eq!( 1083 paginate_all(&store, &key, SortDir::Desc, 11), 1084 expected.iter().rev().cloned().collect::<Vec<_>>() 1085 ); 1086 } 1087 1088 #[test] 1089 fn add_then_count() { 1090 let store = store(); 1091 let key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:abalone")); 1092 1093 store.add(star_edge( 1094 at("at://did:plc:nel/sh.tangled.feed.star/r1"), 1095 did("did:plc:abalone"), 1096 )); 1097 store.add(star_edge( 1098 at("at://did:plc:olaren/sh.tangled.feed.star/r2"), 1099 did("did:plc:abalone"), 1100 )); 1101 store.add(star_edge( 1102 at("at://did:plc:nel/sh.tangled.feed.star/r3"), 1103 did("did:plc:abalone"), 1104 )); 1105 1106 assert_eq!(store.count(&key), 3); 1107 assert_eq!(store.count_distinct_authors(&key), 2); 1108 } 1109 1110 #[test] 1111 fn duplicate_add_is_idempotent() { 1112 let store = store(); 1113 let key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:abalone")); 1114 let edge = star_edge( 1115 at("at://did:plc:nel/sh.tangled.feed.star/r1"), 1116 did("did:plc:abalone"), 1117 ); 1118 store.add(edge.clone()); 1119 store.add(edge); 1120 assert_eq!(store.count(&key), 1); 1121 assert_eq!(store.count_distinct_authors(&key), 1); 1122 } 1123 1124 #[test] 1125 fn remove_source_clears_all_keys_for_that_source() { 1126 let store = store(); 1127 let star_key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:abalone")); 1128 let follow_key = EdgeKey::new(nsid("sh.tangled.graph.follow"), did_subj("did:plc:lyna")); 1129 let source = "at://did:plc:nel/sh.tangled.feed.star/r1"; 1130 1131 store.add(Edge { 1132 kind: nsid("sh.tangled.feed.star"), 1133 subject: did_subj("did:plc:abalone"), 1134 source: at(source), 1135 sort_micros: 0, 1136 }); 1137 store.add(Edge { 1138 kind: nsid("sh.tangled.graph.follow"), 1139 subject: did_subj("did:plc:lyna"), 1140 source: at(source), 1141 sort_micros: 0, 1142 }); 1143 1144 assert_eq!(store.count(&star_key), 1); 1145 assert_eq!(store.count(&follow_key), 1); 1146 1147 store.remove_source(&at(source)); 1148 assert_eq!(store.count(&star_key), 0); 1149 assert_eq!(store.count(&follow_key), 0); 1150 } 1151 1152 #[test] 1153 fn upsert_source_replaces_old_edges() { 1154 let store = store(); 1155 let source = at("at://did:plc:teq/sh.tangled.feed.star/r1"); 1156 let old_subject = did_subj("did:plc:abalone"); 1157 let new_subject = did_subj("did:plc:uni"); 1158 let kind = nsid("sh.tangled.feed.star"); 1159 1160 store.upsert_source( 1161 &source, 1162 vec![Edge { 1163 kind: kind.clone(), 1164 subject: old_subject.clone(), 1165 source: source.clone(), 1166 sort_micros: 0, 1167 }], 1168 ); 1169 assert_eq!( 1170 store.count(&EdgeKey::new(kind.clone(), old_subject.clone())), 1171 1 1172 ); 1173 1174 store.upsert_source( 1175 &source, 1176 vec![Edge { 1177 kind: kind.clone(), 1178 subject: new_subject.clone(), 1179 source: source.clone(), 1180 sort_micros: 0, 1181 }], 1182 ); 1183 assert_eq!(store.count(&EdgeKey::new(kind.clone(), old_subject)), 0); 1184 assert_eq!(store.count(&EdgeKey::new(kind, new_subject)), 1); 1185 } 1186 1187 #[test] 1188 fn list_pages_in_sort_order() { 1189 let store = store(); 1190 let key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:abalone")); 1191 (0..5).for_each(|i| { 1192 store.add(star_edge_at( 1193 at(&format!( 1194 "at://did:plc:{}/sh.tangled.feed.star/r{i}", 1195 NAMES[i] 1196 )), 1197 did("did:plc:abalone"), 1198 1_000_000 + i as u64 * 1_000_000, 1199 )); 1200 }); 1201 1202 let page1 = store.list(&key, PageCursor::Start, limit(2), SortDir::Asc); 1203 assert_eq!(page1.items.len(), 2); 1204 let cursor = PageCursor::After(page1.next.expect("has next")); 1205 1206 let page2 = store.list(&key, cursor, limit(2), SortDir::Asc); 1207 assert_eq!(page2.items.len(), 2); 1208 1209 let cursor2 = PageCursor::After(page2.next.expect("has next")); 1210 let page3 = store.list(&key, cursor2, limit(2), SortDir::Asc); 1211 assert_eq!(page3.items.len(), 1); 1212 assert!( 1213 page3.next.is_none(), 1214 "final partial page should signal exhaustion" 1215 ); 1216 } 1217 1218 #[test] 1219 fn list_exact_fill_signals_exhaustion() { 1220 let store = store(); 1221 let key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:abalone")); 1222 (0..2).for_each(|i| { 1223 store.add(star_edge( 1224 at(&format!( 1225 "at://did:plc:{}/sh.tangled.feed.star/r{i}", 1226 NAMES[i] 1227 )), 1228 did("did:plc:abalone"), 1229 )); 1230 }); 1231 1232 let page = store.list(&key, PageCursor::Start, limit(2), SortDir::Asc); 1233 assert_eq!(page.items.len(), 2); 1234 assert!(page.next.is_none(), "exact-fill page must not promise more"); 1235 } 1236 1237 #[test] 1238 fn list_on_unknown_key_is_empty() { 1239 let store = store(); 1240 let page = store.list( 1241 &EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:periwinkle")), 1242 PageCursor::Start, 1243 limit(10), 1244 SortDir::Asc, 1245 ); 1246 assert!(page.items.is_empty()); 1247 assert!(page.next.is_none()); 1248 } 1249 1250 #[test] 1251 fn distinct_authors_decreases_when_last_source_from_author_removed() { 1252 let store = store(); 1253 let key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:abalone")); 1254 let s1 = at("at://did:plc:nel/sh.tangled.feed.star/r1"); 1255 let s2 = at("at://did:plc:nel/sh.tangled.feed.star/r2"); 1256 let s3 = at("at://did:plc:olaren/sh.tangled.feed.star/r3"); 1257 1258 store.add(star_edge(s1.clone(), did("did:plc:abalone"))); 1259 store.add(star_edge(s2.clone(), did("did:plc:abalone"))); 1260 store.add(star_edge(s3, did("did:plc:abalone"))); 1261 assert_eq!(store.count_distinct_authors(&key), 2); 1262 1263 store.remove_source(&s1); 1264 assert_eq!(store.count_distinct_authors(&key), 2, "user1 still has s2"); 1265 1266 store.remove_source(&s2); 1267 assert_eq!(store.count_distinct_authors(&key), 1, "user1 fully gone"); 1268 } 1269 1270 #[test] 1271 fn page_limit_rejects_zero_and_oversize() { 1272 assert!(matches!( 1273 PageLimit::new(0), 1274 Err(PageLimitError::TooSmall { value: 0, min: 1 }) 1275 )); 1276 assert!(matches!( 1277 PageLimit::new(PageLimit::MAX + 1), 1278 Err(PageLimitError::TooLarge { .. }) 1279 )); 1280 assert_eq!(PageLimit::new(50).unwrap().get(), 50); 1281 } 1282 1283 #[test] 1284 fn cursor_token_round_trip() { 1285 let original = PageToken::new(1_730_000_000_000_000, 0x1234_abcd); 1286 let token = original.encode_token(); 1287 assert_eq!(token.len(), 24, "12-byte cursor encodes to 24 hex chars"); 1288 assert_eq!(PageToken::decode_token(&token).unwrap(), original); 1289 } 1290 1291 #[test] 1292 fn from_token_none_yields_start() { 1293 assert_eq!(PageCursor::from_token(None).unwrap(), PageCursor::Start); 1294 } 1295 1296 #[test] 1297 fn from_token_some_yields_after() { 1298 let original = PageToken::new(1_730_000_000_000_000, 42); 1299 let token = original.encode_token(); 1300 assert_eq!( 1301 PageCursor::from_token(Some(&token)).unwrap(), 1302 PageCursor::After(original), 1303 ); 1304 } 1305 1306 #[test] 1307 fn cursor_decode_rejects_malformed() { 1308 let bad = [ 1309 "", 1310 "deadbeef", 1311 "no-hex!!aaaaaaaaaaaaaaaa", 1312 "12345", 1313 "this-string-is-way-too-long-to-be-a-valid-cursor", 1314 ]; 1315 bad.into_iter().for_each(|s| { 1316 assert!( 1317 matches!(PageToken::decode_token(s), Err(CursorParseError::Malformed)), 1318 "expected malformed for {s:?}", 1319 ); 1320 }); 1321 } 1322 1323 #[test] 1324 fn cursor_token_is_hex_shape() { 1325 let token = PageToken::new(1_730_000_000_000_000, 7).encode_token(); 1326 assert_eq!(token.len(), 24); 1327 assert!(token.chars().all(|c| c.is_ascii_hexdigit())); 1328 } 1329 1330 #[test] 1331 fn list_does_not_drop_entries_at_same_sort_micros_across_pages() { 1332 let store = store(); 1333 let subject_did = did("did:plc:limpet"); 1334 let key = EdgeKey::new( 1335 nsid("sh.tangled.feed.star"), 1336 SubjectRef::Did(subject_did.clone()), 1337 ); 1338 (0..4).for_each(|i| { 1339 store.add(star_edge_at( 1340 at(&format!( 1341 "at://did:plc:{}/sh.tangled.feed.star/r{i}", 1342 NAMES[i] 1343 )), 1344 subject_did.clone(), 1345 42, 1346 )); 1347 }); 1348 1349 let page1 = store.list(&key, PageCursor::Start, limit(2), SortDir::Asc); 1350 assert_eq!(page1.items.len(), 2); 1351 let token = page1.next.expect("cursor must continue across ties"); 1352 1353 let page2 = store.list(&key, PageCursor::After(token), limit(2), SortDir::Asc); 1354 assert_eq!( 1355 page2.items.len(), 1356 2, 1357 "remaining ties must surface on next page" 1358 ); 1359 assert!(page2.next.is_none()); 1360 let combined: std::collections::HashSet<String> = page1 1361 .items 1362 .iter() 1363 .chain(page2.items.iter()) 1364 .map(|u| u.as_ref().to_owned()) 1365 .collect(); 1366 assert_eq!(combined.len(), 4, "every tied entry visible exactly once"); 1367 } 1368 1369 #[test] 1370 fn list_filtered_narrows_by_predicate_and_paginates() { 1371 let store = store(); 1372 let subject_did = did("did:plc:limpet"); 1373 let key = EdgeKey::new( 1374 nsid("sh.tangled.repo.issue"), 1375 SubjectRef::Did(subject_did.clone()), 1376 ); 1377 let by_nel = (0..3).map(|i| { 1378 star_edge_at( 1379 at(&format!("at://did:plc:nel/sh.tangled.repo.issue/r{i}")), 1380 subject_did.clone(), 1381 100 + i as u64, 1382 ) 1383 }); 1384 let by_olaren = (0..2).map(|i| { 1385 star_edge_at( 1386 at(&format!("at://did:plc:olaren/sh.tangled.repo.issue/o{i}")), 1387 subject_did.clone(), 1388 500 + i as u64, 1389 ) 1390 }); 1391 by_nel 1392 .chain(by_olaren) 1393 .map(|mut e| { 1394 e.kind = nsid("sh.tangled.repo.issue"); 1395 e 1396 }) 1397 .for_each(|e| store.add(e)); 1398 1399 let only_nel = |u: &AtUri<DefaultStr>| u.as_ref().starts_with("at://did:plc:nel/"); 1400 let page1 = store.list_filtered(&key, PageCursor::Start, limit(2), SortDir::Asc, only_nel); 1401 assert_eq!(page1.items.len(), 2); 1402 assert!(page1.next.is_some(), "cursor must allow more nel matches"); 1403 1404 let page2 = store.list_filtered( 1405 &key, 1406 PageCursor::After(page1.next.unwrap()), 1407 limit(2), 1408 SortDir::Asc, 1409 only_nel, 1410 ); 1411 assert_eq!(page2.items.len(), 1, "only one nel issue left"); 1412 assert!(page2.next.is_none(), "tail page must not promise more"); 1413 } 1414 1415 #[test] 1416 fn list_descending_returns_newest_first() { 1417 let store = store(); 1418 let subject_did = did("did:plc:limpet"); 1419 let key = EdgeKey::new( 1420 nsid("sh.tangled.feed.star"), 1421 SubjectRef::Did(subject_did.clone()), 1422 ); 1423 (0..5).for_each(|i| { 1424 store.add(star_edge_at( 1425 at(&format!( 1426 "at://did:plc:{}/sh.tangled.feed.star/r{i}", 1427 NAMES[i] 1428 )), 1429 subject_did.clone(), 1430 1_000_000 + i as u64 * 1_000_000, 1431 )); 1432 }); 1433 1434 let asc = store.list(&key, PageCursor::Start, limit(5), SortDir::Asc); 1435 let desc = store.list(&key, PageCursor::Start, limit(5), SortDir::Desc); 1436 assert_eq!(asc.items.len(), 5); 1437 assert_eq!(desc.items.len(), 5); 1438 let asc_uris: Vec<_> = asc.items.iter().map(|u| u.as_ref().to_owned()).collect(); 1439 let mut reversed = asc_uris.clone(); 1440 reversed.reverse(); 1441 let desc_uris: Vec<_> = desc.items.iter().map(|u| u.as_ref().to_owned()).collect(); 1442 assert_eq!(desc_uris, reversed, "desc must be exact reverse of asc"); 1443 } 1444 1445 #[test] 1446 fn list_descending_paginates_with_cursor() { 1447 let store = store(); 1448 let subject_did = did("did:plc:whelk"); 1449 let key = EdgeKey::new( 1450 nsid("sh.tangled.feed.star"), 1451 SubjectRef::Did(subject_did.clone()), 1452 ); 1453 (0..5).for_each(|i| { 1454 store.add(star_edge_at( 1455 at(&format!( 1456 "at://did:plc:{}/sh.tangled.feed.star/r{i}", 1457 NAMES[i] 1458 )), 1459 subject_did.clone(), 1460 1_000_000 + i as u64 * 1_000_000, 1461 )); 1462 }); 1463 1464 let page1 = store.list(&key, PageCursor::Start, limit(2), SortDir::Desc); 1465 assert_eq!(page1.items.len(), 2); 1466 let cursor = PageCursor::After(page1.next.expect("desc page1 must continue")); 1467 let page2 = store.list(&key, cursor, limit(2), SortDir::Desc); 1468 assert_eq!(page2.items.len(), 2); 1469 let cursor2 = PageCursor::After(page2.next.expect("desc page2 must continue")); 1470 let page3 = store.list(&key, cursor2, limit(2), SortDir::Desc); 1471 assert_eq!(page3.items.len(), 1); 1472 assert!(page3.next.is_none()); 1473 1474 let combined: std::collections::HashSet<String> = page1 1475 .items 1476 .iter() 1477 .chain(page2.items.iter()) 1478 .chain(page3.items.iter()) 1479 .map(|u| u.as_ref().to_owned()) 1480 .collect(); 1481 assert_eq!( 1482 combined.len(), 1483 5, 1484 "desc pagination visits each item exactly once" 1485 ); 1486 } 1487 1488 #[test] 1489 fn list_filtered_descending_respects_predicate() { 1490 let store = store(); 1491 let subject_did = did("did:plc:scallop"); 1492 let key = EdgeKey::new( 1493 nsid("sh.tangled.repo.issue"), 1494 SubjectRef::Did(subject_did.clone()), 1495 ); 1496 let by_nel = (0..3).map(|i| { 1497 star_edge_at( 1498 at(&format!("at://did:plc:nel/sh.tangled.repo.issue/r{i}")), 1499 subject_did.clone(), 1500 100 + i as u64, 1501 ) 1502 }); 1503 let by_olaren = (0..2).map(|i| { 1504 star_edge_at( 1505 at(&format!("at://did:plc:olaren/sh.tangled.repo.issue/o{i}")), 1506 subject_did.clone(), 1507 500 + i as u64, 1508 ) 1509 }); 1510 by_nel 1511 .chain(by_olaren) 1512 .map(|mut e| { 1513 e.kind = nsid("sh.tangled.repo.issue"); 1514 e 1515 }) 1516 .for_each(|e| store.add(e)); 1517 1518 let only_nel = |u: &AtUri<DefaultStr>| u.as_ref().starts_with("at://did:plc:nel/"); 1519 let page = store.list_filtered(&key, PageCursor::Start, limit(5), SortDir::Desc, only_nel); 1520 assert_eq!(page.items.len(), 3, "all three nel issues visible"); 1521 let last = page.items.last().unwrap().as_ref(); 1522 let first = page.items.first().unwrap().as_ref(); 1523 assert!( 1524 first > last, 1525 "desc order: first item rkey must be greater than last (got first={first}, last={last})", 1526 ); 1527 } 1528 1529 #[test] 1530 fn non_did_source_round_trips_via_raw_fallback() { 1531 let store = store(); 1532 let subject = did("did:plc:limpet"); 1533 let key = EdgeKey::new( 1534 nsid("sh.tangled.feed.star"), 1535 SubjectRef::Did(subject.clone()), 1536 ); 1537 let did_source = at("at://did:plc:nel/sh.tangled.feed.star/r1"); 1538 let handle_source = at("at://witchcraft.systems/sh.tangled.feed.star/r2"); 1539 store.add(star_edge_at(did_source.clone(), subject.clone(), 1)); 1540 store.add(star_edge_at(handle_source.clone(), subject.clone(), 2)); 1541 1542 let page = store.list(&key, PageCursor::Start, limit(10), SortDir::Asc); 1543 let got: std::collections::HashSet<String> = 1544 page.items.iter().map(|u| u.as_ref().to_owned()).collect(); 1545 assert!( 1546 got.contains(did_source.as_ref()), 1547 "did source must decode exactly" 1548 ); 1549 assert!( 1550 got.contains(handle_source.as_ref()), 1551 "non-did authority must round-trip through the raw fallback" 1552 ); 1553 1554 store.remove_source(&handle_source); 1555 assert_eq!( 1556 store.count(&key), 1557 1, 1558 "raw-keyed source removable by its uri" 1559 ); 1560 } 1561 1562 #[test] 1563 fn distinct_collections_decode_with_their_own_collection() { 1564 let store = store(); 1565 let subject = did("did:plc:limpet"); 1566 let key = EdgeKey::new( 1567 nsid("sh.tangled.feed.star"), 1568 SubjectRef::Did(subject.clone()), 1569 ); 1570 let star_src = at("at://did:plc:nel/sh.tangled.feed.star/aaa"); 1571 let issue_src = at("at://did:plc:nel/sh.tangled.repo.issue/bbb"); 1572 store.add(Edge { 1573 kind: nsid("sh.tangled.feed.star"), 1574 subject: SubjectRef::Did(subject.clone()), 1575 source: star_src.clone(), 1576 sort_micros: 1, 1577 }); 1578 store.add(Edge { 1579 kind: nsid("sh.tangled.feed.star"), 1580 subject: SubjectRef::Did(subject.clone()), 1581 source: issue_src.clone(), 1582 sort_micros: 2, 1583 }); 1584 1585 let page = store.list(&key, PageCursor::Start, limit(10), SortDir::Asc); 1586 let got: std::collections::HashSet<String> = 1587 page.items.iter().map(|u| u.as_ref().to_owned()).collect(); 1588 assert!( 1589 got.contains(star_src.as_ref()), 1590 "star-collection source decodes exactly" 1591 ); 1592 assert!( 1593 got.contains(issue_src.as_ref()), 1594 "issue-collection source must keep its own collection, not borrow the star one" 1595 ); 1596 } 1597 1598 #[test] 1599 fn author_refs_spill_preserves_distinct_count() { 1600 let store = store(); 1601 let subject = did("did:plc:scallop"); 1602 let key = EdgeKey::new( 1603 nsid("sh.tangled.feed.star"), 1604 SubjectRef::Did(subject.clone()), 1605 ); 1606 (0..5).for_each(|i| { 1607 store.add(star_edge_at( 1608 at(&format!( 1609 "at://did:plc:{}/sh.tangled.feed.star/r{i}", 1610 NAMES[i] 1611 )), 1612 subject.clone(), 1613 i as u64, 1614 )); 1615 }); 1616 assert_eq!( 1617 store.count_distinct_authors(&key), 1618 5, 1619 "five authors exceed the inline cap and spill to the map" 1620 ); 1621 1622 store.add(star_edge_at( 1623 at("at://did:plc:nel/sh.tangled.feed.star/r99"), 1624 subject.clone(), 1625 99, 1626 )); 1627 assert_eq!( 1628 store.count_distinct_authors(&key), 1629 5, 1630 "second nel source adds no author" 1631 ); 1632 1633 store.remove_source(&at("at://did:plc:nel/sh.tangled.feed.star/r0")); 1634 assert_eq!( 1635 store.count_distinct_authors(&key), 1636 5, 1637 "nel still present via r99" 1638 ); 1639 store.remove_source(&at("at://did:plc:nel/sh.tangled.feed.star/r99")); 1640 assert_eq!(store.count_distinct_authors(&key), 4, "nel fully removed"); 1641 } 1642}