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