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
75 kB 2258 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, RwLock}; 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, Record}; 29use bobbin_types::ids::EdgeKey; 30use either::Either; 31use itertools::Itertools; 32use jacquard_common::DefaultStr; 33use jacquard_common::types::did::Did; 34use jacquard_common::types::nsid::Nsid; 35use jacquard_common::types::string::AtUri; 36use lasso::{Key, Spur, ThreadedRodeo}; 37use scc::HashMap as SccMap; 38use scc::hash_map::Entry; 39use smallvec::SmallVec; 40use thiserror::Error; 41 42#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq, Ord, PartialOrd)] 43struct SortMicros(u64); 44 45#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq, Ord, PartialOrd)] 46struct BucketKey { 47 micros: SortMicros, 48 source: SourceId, 49} 50 51impl BucketKey { 52 fn new(micros: u64, source: SourceId) -> Self { 53 Self { 54 micros: SortMicros(micros), 55 source, 56 } 57 } 58 59 fn token(self) -> PageToken { 60 PageToken::new(self.micros.0, self.source.index()) 61 } 62 63 fn from_token(tok: PageToken) -> Self { 64 Self { 65 micros: SortMicros(tok.micros), 66 source: SourceId::from_raw(tok.source), 67 } 68 } 69} 70 71pub mod coverage; 72pub mod state_index; 73pub use coverage::{Coverage, CoverageWatch, HydrantCursor, PromotionSignal}; 74pub use state_index::{ 75 ApplyOutcome, IssueStateKind, PullStatusKind, StateIndex, StateKind, apply_record_state, 76}; 77 78#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 79pub struct SourceTag; 80#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 81struct AuthorTag; 82#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 83struct CollectionTag; 84 85#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 86pub struct Interned<T>(u32, PhantomData<T>); 87 88impl<T: Copy> Interned<T> { 89 fn from_spur(spur: Spur) -> Self { 90 Self(spur.into_usize() as u32, PhantomData) 91 } 92 93 fn from_raw(raw: u32) -> Self { 94 Self(raw, PhantomData) 95 } 96 97 fn index(self) -> u32 { 98 self.0 99 } 100 101 fn to_spur(self) -> Option<Spur> { 102 Spur::try_from_usize(self.0 as usize) 103 } 104} 105 106pub type SourceId = Interned<SourceTag>; 107type AuthorId = Interned<AuthorTag>; 108type CollectionId = Interned<CollectionTag>; 109 110#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 111pub struct PageToken { 112 micros: u64, 113 source: u32, 114} 115 116impl PageToken { 117 pub fn new(micros: u64, source: u32) -> Self { 118 Self { micros, source } 119 } 120 121 pub fn micros(self) -> u64 { 122 self.micros 123 } 124 125 pub fn source(self) -> u32 { 126 self.source 127 } 128 129 pub fn encode_token(self) -> String { 130 let mut bytes = [0u8; 12]; 131 bytes[..8].copy_from_slice(&self.micros.to_be_bytes()); 132 bytes[8..].copy_from_slice(&self.source.to_be_bytes()); 133 encode_hex(&bytes) 134 } 135 136 pub fn decode_token(token: &str) -> Result<Self, CursorParseError> { 137 let bytes: [u8; 12] = decode_hex_array(token).ok_or(CursorParseError::Malformed)?; 138 let micros = u64::from_be_bytes(bytes[..8].try_into().unwrap()); 139 let source = u32::from_be_bytes(bytes[8..].try_into().unwrap()); 140 Ok(Self { micros, source }) 141 } 142} 143 144fn encode_hex(bytes: &[u8]) -> String { 145 bytes 146 .iter() 147 .fold(String::with_capacity(bytes.len() * 2), |mut acc, b| { 148 acc.push(char::from_digit((b >> 4) as u32, 16).unwrap()); 149 acc.push(char::from_digit((b & 0x0f) as u32, 16).unwrap()); 150 acc 151 }) 152} 153 154fn decode_hex_array<const N: usize>(token: &str) -> Option<[u8; N]> { 155 if token.len() != N * 2 { 156 return None; 157 } 158 let parsed: Vec<u8> = token 159 .as_bytes() 160 .chunks_exact(2) 161 .map(|pair| { 162 let hi = (pair[0] as char).to_digit(16)?; 163 let lo = (pair[1] as char).to_digit(16)?; 164 Some(((hi << 4) | lo) as u8) 165 }) 166 .collect::<Option<Vec<u8>>>()?; 167 parsed.try_into().ok() 168} 169 170#[derive(Clone, Copy, Debug, Eq, PartialEq)] 171pub enum PageCursor { 172 Start, 173 After(PageToken), 174} 175 176#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] 177pub enum SortDir { 178 Asc, 179 #[default] 180 Desc, 181} 182 183#[derive(Clone, Copy, Debug, Eq, PartialEq, Error)] 184pub enum CursorParseError { 185 #[error("cursor token must be a valid TID")] 186 Malformed, 187} 188 189impl PageCursor { 190 pub fn from_token(raw: Option<&str>) -> Result<Self, CursorParseError> { 191 raw.map_or(Ok(Self::Start), |t| { 192 PageToken::decode_token(t).map(Self::After) 193 }) 194 } 195} 196 197#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 198pub struct PageLimit(u32); 199 200#[derive(Clone, Copy, Debug, Eq, PartialEq, Error)] 201pub enum PageLimitError { 202 #[error("page limit {value} below minimum {min}")] 203 TooSmall { value: u32, min: u32 }, 204 #[error("page limit {value} above maximum {max}")] 205 TooLarge { value: u32, max: u32 }, 206} 207 208impl PageLimit { 209 pub const MIN: u32 = 1; 210 pub const MAX: u32 = 1000; 211 212 pub fn new(value: u32) -> Result<Self, PageLimitError> { 213 match value { 214 v if v < Self::MIN => Err(PageLimitError::TooSmall { 215 value: v, 216 min: Self::MIN, 217 }), 218 v if v > Self::MAX => Err(PageLimitError::TooLarge { 219 value: v, 220 max: Self::MAX, 221 }), 222 v => Ok(Self(v)), 223 } 224 } 225 226 pub const fn get(self) -> u32 { 227 self.0 228 } 229} 230 231#[derive(Debug)] 232pub struct EdgePage { 233 pub items: Vec<EdgeItem>, 234 pub next: Option<PageToken>, 235} 236 237#[derive(Clone, Debug, Eq, PartialEq)] 238pub struct EdgeItem { 239 pub uri: AtUri<DefaultStr>, 240 pub sort_micros: u64, 241} 242 243impl AsRef<str> for EdgeItem { 244 fn as_ref(&self) -> &str { 245 self.uri.as_ref() 246 } 247} 248 249#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] 250pub struct FilteredCount { 251 pub count: Count, 252 pub distinct_authors: DistinctAuthorCount, 253} 254 255#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] 256pub struct Count(u64); 257 258impl Count { 259 pub const fn new(value: u64) -> Self { 260 Self(value) 261 } 262 263 pub const fn get(self) -> u64 { 264 self.0 265 } 266} 267 268#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] 269pub struct DistinctAuthorCount(u64); 270 271impl DistinctAuthorCount { 272 pub const fn new(value: u64) -> Self { 273 Self(value) 274 } 275 276 pub const fn get(self) -> u64 { 277 self.0 278 } 279} 280 281impl FilteredCount { 282 pub const fn new(count: u64, distinct_authors: u64) -> Self { 283 Self { 284 count: Count::new(count), 285 distinct_authors: DistinctAuthorCount::new(distinct_authors), 286 } 287 } 288} 289 290#[derive(Clone, Copy, Debug, Default)] 291pub struct EdgeMemReport { 292 pub key_count: u64, 293 pub source_count: u64, 294 pub source_interner_bytes: u64, 295 pub did_interner_bytes: u64, 296 pub collection_interner_bytes: u64, 297 pub key_interner_bytes: u64, 298 pub edges_total: u64, 299 pub author_refs_total: u64, 300 pub reverse_entries: u64, 301 pub reverse_cap: u64, 302 pub forward_struct_bytes: u64, 303 pub reverse_struct_bytes: u64, 304 pub bucket_struct_bytes: u64, 305 pub max_bucket: u64, 306 pub bucket_size_classes: [u64; BUCKET_CLASS_COUNT], 307} 308 309const BUCKET_CLASS_BOUNDS: [u64; 16] = [ 310 1, 311 2, 312 4, 313 8, 314 16, 315 32, 316 64, 317 128, 318 256, 319 512, 320 1024, 321 2048, 322 8192, 323 32768, 324 131072, 325 u64::MAX, 326]; 327const BUCKET_CLASS_COUNT: usize = BUCKET_CLASS_BOUNDS.len(); 328 329fn bucket_class(n: u64) -> usize { 330 BUCKET_CLASS_BOUNDS 331 .iter() 332 .position(|&bound| n <= bound) 333 .unwrap_or(BUCKET_CLASS_COUNT - 1) 334} 335 336impl EdgeMemReport { 337 pub fn bucket_histogram(&self) -> impl Iterator<Item = (u64, u64)> { 338 BUCKET_CLASS_BOUNDS 339 .into_iter() 340 .zip(self.bucket_size_classes) 341 } 342} 343 344#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 345struct EdgeKeyId(u32); 346 347fn bump_author(authors: &mut HashMap<AuthorId, NonZeroU32, RuntimeHasher>, author: AuthorId) { 348 authors 349 .entry(author) 350 .and_modify(|c| *c = c.saturating_add(1)) 351 .or_insert(NonZeroU32::MIN); 352} 353 354fn drop_author(authors: &mut HashMap<AuthorId, NonZeroU32, RuntimeHasher>, author: AuthorId) { 355 match authors.get(&author).map(|c| c.get() - 1) { 356 Some(0) | None => { 357 authors.remove(&author); 358 } 359 Some(next) => { 360 authors.insert(author, NonZeroU32::new(next).unwrap()); 361 } 362 } 363} 364 365struct LargeBucket { 366 keys: BTreeSet<BucketKey>, 367 authors: HashMap<AuthorId, NonZeroU32, RuntimeHasher>, 368} 369 370const BUCKET_PROMOTE_AT: usize = 256; 371const SOURCE_BYTES: u64 = 16; 372 373enum Sources { 374 Small(SmallVec<[BucketKey; 2]>), 375 Large(Box<LargeBucket>), 376} 377 378impl Default for Sources { 379 fn default() -> Self { 380 Self::Small(SmallVec::new()) 381 } 382} 383 384impl Sources { 385 fn len(&self) -> usize { 386 match self { 387 Self::Small(v) => v.len(), 388 Self::Large(big) => big.keys.len(), 389 } 390 } 391 392 fn is_empty(&self) -> bool { 393 self.len() == 0 394 } 395 396 fn insert(&mut self, key: BucketKey, author: Option<AuthorId>) -> bool { 397 match self { 398 Self::Small(v) => match v.binary_search(&key) { 399 Ok(_) => false, 400 Err(pos) => { 401 v.insert(pos, key); 402 true 403 } 404 }, 405 Self::Large(big) => { 406 let inserted = big.keys.insert(key); 407 if inserted && let Some(a) = author { 408 bump_author(&mut big.authors, a); 409 } 410 inserted 411 } 412 } 413 } 414 415 fn remove(&mut self, key: &BucketKey, author: Option<AuthorId>) { 416 match self { 417 Self::Small(v) => { 418 if let Ok(pos) = v.binary_search(key) { 419 v.remove(pos); 420 } 421 } 422 Self::Large(big) => { 423 if big.keys.remove(key) 424 && let Some(a) = author 425 { 426 drop_author(&mut big.authors, a); 427 } 428 } 429 } 430 } 431 432 fn directed( 433 &self, 434 cursor: PageCursor, 435 dir: SortDir, 436 ) -> Box<dyn Iterator<Item = BucketKey> + '_> { 437 match self { 438 Self::Small(v) => Box::new(directed_slice(v, cursor, dir)), 439 Self::Large(big) => Box::new(directed_tree(&big.keys, cursor, dir)), 440 } 441 } 442 443 fn heap_bytes(&self) -> u64 { 444 const BTREE_BYTES_PER_KEY: u64 = 32; 445 const HASHMAP_FIXED: u64 = 48; 446 const HASHMAP_PER_CAP: u64 = 9; 447 match self { 448 Self::Small(v) if v.spilled() => v.capacity() as u64 * SOURCE_BYTES, 449 Self::Small(_) => 0, 450 Self::Large(big) => { 451 std::mem::size_of::<LargeBucket>() as u64 452 + big.keys.len() as u64 * BTREE_BYTES_PER_KEY 453 + HASHMAP_FIXED 454 + big.authors.capacity() as u64 * HASHMAP_PER_CAP 455 } 456 } 457 } 458} 459 460#[derive(Clone, Copy, Debug, Eq, PartialEq)] 461struct ReverseEntry { 462 key_id: EdgeKeyId, 463 sort_micros: u64, 464} 465 466#[derive(Clone, Copy)] 467struct ProjectedState<K> { 468 key: EdgeKeyId, 469 kind: K, 470 author: Option<AuthorId>, 471} 472 473struct CountBucket { 474 count: u64, 475 authors: HashMap<AuthorId, NonZeroU32, RuntimeHasher>, 476} 477 478struct StateCountInner<K> { 479 projected: HashMap<SourceId, SmallVec<[ProjectedState<K>; 1]>, RuntimeHasher>, 480 buckets: HashMap<(EdgeKeyId, K), CountBucket, RuntimeHasher>, 481} 482 483struct StateCountIndex<K> { 484 inner: RwLock<StateCountInner<K>>, 485 hasher: RuntimeHasher, 486} 487 488impl<K> StateCountIndex<K> { 489 fn new(hasher: RuntimeHasher) -> Self { 490 Self { 491 inner: RwLock::new(StateCountInner { 492 projected: HashMap::with_hasher(hasher.clone()), 493 buckets: HashMap::with_hasher(hasher.clone()), 494 }), 495 hasher, 496 } 497 } 498 499 fn heap_bytes(&self) -> u64 { 500 let inner = self 501 .inner 502 .read() 503 .expect("state-count index rwlock poisoned"); 504 let projected = inner.projected.capacity() 505 * (std::mem::size_of::<SourceId>() 506 + std::mem::size_of::<SmallVec<[ProjectedState<K>; 1]>>() 507 + 1) 508 + inner 509 .projected 510 .values() 511 .filter(|states| states.spilled()) 512 .map(|states| states.capacity() * std::mem::size_of::<ProjectedState<K>>()) 513 .sum::<usize>(); 514 let buckets = inner.buckets.capacity() 515 * (std::mem::size_of::<(EdgeKeyId, K)>() + std::mem::size_of::<CountBucket>() + 1); 516 let author_slots = inner 517 .buckets 518 .values() 519 .map(|bucket| { 520 bucket.authors.capacity() 521 * (std::mem::size_of::<AuthorId>() + std::mem::size_of::<NonZeroU32>() + 1) 522 }) 523 .sum::<usize>(); 524 (projected + buckets + author_slots) as u64 525 } 526} 527 528pub struct EdgeStore { 529 source_interner: Arc<ThreadedRodeo<Spur, RuntimeHasher>>, 530 did_interner: Arc<ThreadedRodeo<Spur, RuntimeHasher>>, 531 collection_interner: Arc<ThreadedRodeo<Spur, RuntimeHasher>>, 532 key_ids: SccMap<EdgeKey, EdgeKeyId, RuntimeHasher>, 533 keys: SccMap<EdgeKeyId, EdgeKey, RuntimeHasher>, 534 next_key_id: AtomicU32, 535 forward: SccMap<EdgeKeyId, Sources, RuntimeHasher>, 536 reverse: SccMap<SourceId, SmallVec<[ReverseEntry; 1]>, RuntimeHasher>, 537 issue_counts: StateCountIndex<IssueStateKind>, 538 pull_counts: StateCountIndex<PullStatusKind>, 539 hasher: RuntimeHasher, 540 writer: Mutex<()>, 541} 542 543impl EdgeStore { 544 pub fn new(hasher: RuntimeHasher) -> Self { 545 Self { 546 source_interner: Arc::new(ThreadedRodeo::with_hasher(hasher.clone())), 547 did_interner: Arc::new(ThreadedRodeo::with_hasher(hasher.clone())), 548 collection_interner: Arc::new(ThreadedRodeo::with_hasher(hasher.clone())), 549 key_ids: SccMap::with_hasher(hasher.clone()), 550 keys: SccMap::with_hasher(hasher.clone()), 551 next_key_id: AtomicU32::new(0), 552 forward: SccMap::with_hasher(hasher.clone()), 553 reverse: SccMap::with_hasher(hasher.clone()), 554 issue_counts: StateCountIndex::new(hasher.clone()), 555 pull_counts: StateCountIndex::new(hasher.clone()), 556 hasher, 557 writer: Mutex::new(()), 558 } 559 } 560 561 fn intern_key(&self, key: EdgeKey) -> EdgeKeyId { 562 match self.key_ids.entry_sync(key.clone()) { 563 Entry::Occupied(e) => *e.get(), 564 Entry::Vacant(e) => { 565 let id = EdgeKeyId(self.next_key_id.fetch_add(1, Ordering::Relaxed)); 566 e.insert_entry(id); 567 let _ = self.keys.insert_sync(id, key); 568 id 569 } 570 } 571 } 572 573 fn lookup_key(&self, key: &EdgeKey) -> Option<EdgeKeyId> { 574 self.key_ids.read_sync(key, |_, id| *id) 575 } 576 577 pub fn add(&self, edge: Edge) { 578 let _w = self 579 .writer 580 .lock() 581 .expect("edge-store writer mutex poisoned"); 582 self.add_locked(edge); 583 } 584 585 pub fn intern_source(&self, source: &AtUri<DefaultStr>) -> SourceId { 586 let author = self.intern_author(source); 587 SourceId::from_spur( 588 self.source_interner 589 .get_or_intern(self.source_key(source, author)), 590 ) 591 } 592 593 pub fn upsert_source(&self, source: &AtUri<DefaultStr>, edges: Vec<Edge>) { 594 let _w = self 595 .writer 596 .lock() 597 .expect("edge-store writer mutex poisoned"); 598 self.clear_source_locked(source); 599 edges.into_iter().for_each(|e| self.add_locked(e)); 600 } 601 602 pub fn remove_source(&self, source: &AtUri<DefaultStr>) { 603 let _w = self 604 .writer 605 .lock() 606 .expect("edge-store writer mutex poisoned"); 607 self.clear_source_locked(source); 608 } 609 610 fn add_locked(&self, edge: Edge) { 611 let author = self.intern_author(&edge.source); 612 let source_key = self.source_key(&edge.source, author); 613 let id = SourceId::from_spur(self.source_interner.get_or_intern(source_key)); 614 let sort_micros = edge.sort_micros; 615 let key_id = self.intern_key(EdgeKey::new(edge.kind, edge.subject)); 616 617 let key = BucketKey::new(sort_micros, id); 618 let mut entry = self.forward.entry_sync(key_id).or_default(); 619 let inserted = entry.get_mut().insert(key, author); 620 let promote_keys = match entry.get() { 621 Sources::Small(v) if v.len() > BUCKET_PROMOTE_AT => { 622 Some(v.iter().copied().collect::<Vec<_>>()) 623 } 624 _ => None, 625 }; 626 drop(entry); 627 628 if let Some(keys) = promote_keys { 629 let authors = self.build_author_map(&keys); 630 let large = LargeBucket { 631 keys: keys.into_iter().collect(), 632 authors, 633 }; 634 if let Entry::Occupied(mut e) = self.forward.entry_sync(key_id) { 635 *e.get_mut() = Sources::Large(Box::new(large)); 636 } 637 } 638 639 if inserted { 640 let mut rev = self.reverse.entry_sync(id).or_default(); 641 rev.get_mut().push(ReverseEntry { 642 key_id, 643 sort_micros, 644 }); 645 } 646 } 647 648 fn clear_source_locked(&self, source: &AtUri<DefaultStr>) { 649 let author = source_authority_did(source) 650 .and_then(|s| self.did_interner.get(s)) 651 .map(AuthorId::from_spur); 652 let Some(source_spur) = self.source_interner.get(self.source_key(source, author)) else { 653 return; 654 }; 655 let id = SourceId::from_spur(source_spur); 656 let Some((_, entries)) = self.reverse.remove_sync(&id) else { 657 return; 658 }; 659 entries.into_iter().for_each( 660 |ReverseEntry { 661 key_id, 662 sort_micros, 663 }| { 664 self.forward.update_sync(&key_id, |_, sources| { 665 sources.remove(&BucketKey::new(sort_micros, id), author); 666 }); 667 self.forward 668 .remove_if_sync(&key_id, |sources| sources.is_empty()); 669 }, 670 ); 671 } 672 673 fn intern_author(&self, source: &AtUri<DefaultStr>) -> Option<AuthorId> { 674 let did = source_authority_did(source)?; 675 Some(AuthorId::from_spur(self.did_interner.get_or_intern(did))) 676 } 677 678 fn source_key(&self, source: &AtUri<DefaultStr>, author: Option<AuthorId>) -> String { 679 match (split_record_uri(source.as_ref()), author) { 680 (Some((_, collection, rkey)), Some(author)) => { 681 let collection = 682 CollectionId::from_spur(self.collection_interner.get_or_intern(collection)); 683 format!("{}/{}/{}", author.index(), collection.index(), rkey) 684 } 685 _ => source.as_ref().to_owned(), 686 } 687 } 688 689 fn decode_source(&self, stored: &str) -> Option<String> { 690 if stored.starts_with("at://") { 691 return Some(stored.to_owned()); 692 } 693 let mut parts = stored.splitn(3, '/'); 694 let author = AuthorId::from_raw(parts.next()?.parse().ok()?); 695 let collection = CollectionId::from_raw(parts.next()?.parse().ok()?); 696 let rkey = parts.next()?; 697 let did = self.did_interner.try_resolve(&author.to_spur()?)?; 698 let collection = self 699 .collection_interner 700 .try_resolve(&collection.to_spur()?)?; 701 Some(format!("at://{did}/{collection}/{rkey}")) 702 } 703 704 fn author_of_stored(&self, stored: &str) -> Option<AuthorId> { 705 match stored.strip_prefix("at://") { 706 Some(rest) => { 707 let authority = rest.split('/').next().unwrap_or(rest); 708 authority 709 .starts_with("did:") 710 .then(|| self.did_interner.get(authority).map(AuthorId::from_spur)) 711 .flatten() 712 } 713 None => stored 714 .split('/') 715 .next()? 716 .parse::<u32>() 717 .ok() 718 .map(AuthorId::from_raw), 719 } 720 } 721 722 fn author_of(&self, source: SourceId) -> Option<AuthorId> { 723 let spur = source.to_spur()?; 724 let stored = self.source_interner.try_resolve(&spur)?; 725 self.author_of_stored(stored) 726 } 727 728 fn distinct_authors_small(&self, keys: &[BucketKey]) -> u64 { 729 keys.iter() 730 .filter_map(|key| self.author_of(key.source)) 731 .collect::<std::collections::HashSet<AuthorId>>() 732 .len() as u64 733 } 734 735 fn build_author_map(&self, keys: &[BucketKey]) -> HashMap<AuthorId, NonZeroU32, RuntimeHasher> { 736 keys.iter().fold( 737 HashMap::with_hasher(self.hasher.clone()), 738 |mut authors, key| { 739 if let Some(a) = self.author_of(key.source) { 740 bump_author(&mut authors, a); 741 } 742 authors 743 }, 744 ) 745 } 746 747 pub fn count(&self, key: &EdgeKey) -> u64 { 748 self.lookup_key(key) 749 .and_then(|id| { 750 self.forward 751 .read_sync(&id, |_, sources| sources.len() as u64) 752 }) 753 .unwrap_or(0) 754 } 755 756 pub fn count_distinct_authors(&self, key: &EdgeKey) -> u64 { 757 self.lookup_key(key) 758 .and_then(|id| { 759 self.forward.read_sync(&id, |_, sources| match sources { 760 Sources::Large(big) => big.authors.len() as u64, 761 Sources::Small(v) => self.distinct_authors_small(v), 762 }) 763 }) 764 .unwrap_or(0) 765 } 766 767 pub fn count_by_author(&self, key: &EdgeKey, author: &Did<DefaultStr>) -> u64 { 768 let Some(author) = self 769 .did_interner 770 .get(author.as_ref()) 771 .map(AuthorId::from_spur) 772 else { 773 return 0; 774 }; 775 self.lookup_key(key) 776 .and_then(|id| { 777 self.forward.read_sync(&id, |_, sources| match sources { 778 Sources::Large(big) => big 779 .authors 780 .get(&author) 781 .map_or(0, |count| count.get() as u64), 782 Sources::Small(keys) => keys 783 .iter() 784 .filter(|key| self.author_of(key.source) == Some(author)) 785 .count() as u64, 786 }) 787 }) 788 .unwrap_or(0) 789 } 790 791 fn current_state_projection<K>( 792 &self, 793 states: &StateIndex<K>, 794 entity: &AtUri<DefaultStr>, 795 edge_kind: &str, 796 ) -> (SourceId, SmallVec<[ProjectedState<K>; 1]>) 797 where 798 K: StateKind + Default, 799 { 800 let source = self.intern_source(entity); 801 let author = self.author_of(source); 802 let entity_author = source_authority_did(entity); 803 let reverse = self 804 .reverse 805 .read_sync(&source, |_, entries| entries.clone()) 806 .unwrap_or_default(); 807 let projected = reverse 808 .into_iter() 809 .filter_map(|entry| { 810 let key = self.keys.read_sync(&entry.key_id, |_, key| key.clone())?; 811 if key.kind.as_ref() != edge_kind { 812 return None; 813 } 814 let repo_owner = key.subject.as_did()?; 815 let kind = states 816 .latest_by(entity, |state_source| { 817 let state_author = source_authority_did(state_source); 818 state_author == entity_author || state_author == Some(repo_owner.as_ref()) 819 }) 820 .map(|(kind, _)| kind) 821 .unwrap_or_default(); 822 Some(ProjectedState { 823 key: entry.key_id, 824 kind, 825 author, 826 }) 827 }) 828 .collect(); 829 (source, projected) 830 } 831 832 fn refresh_state_counts<K>( 833 &self, 834 index: &StateCountIndex<K>, 835 states: &StateIndex<K>, 836 entity: &AtUri<DefaultStr>, 837 edge_kind: &str, 838 ) where 839 K: StateKind + Default, 840 { 841 // compute under the write lock so concurrent updates dont write a stale projection 842 let mut inner = index 843 .inner 844 .write() 845 .expect("state-count index rwlock poisoned"); 846 let (source, current) = self.current_state_projection(states, entity, edge_kind); 847 848 if let Some(previous) = inner.projected.remove(&source) { 849 for projected in previous { 850 let bucket_key = (projected.key, projected.kind); 851 let remove = if let Some(bucket) = inner.buckets.get_mut(&bucket_key) { 852 bucket.count -= 1; 853 if let Some(author) = projected.author { 854 drop_author(&mut bucket.authors, author); 855 } 856 bucket.count == 0 857 } else { 858 false 859 }; 860 if remove { 861 inner.buckets.remove(&bucket_key); 862 } 863 } 864 } 865 866 for projected in current.iter().copied() { 867 let bucket = inner 868 .buckets 869 .entry((projected.key, projected.kind)) 870 .or_insert_with(|| CountBucket { 871 count: 0, 872 authors: HashMap::with_hasher(index.hasher.clone()), 873 }); 874 bucket.count += 1; 875 if let Some(author) = projected.author { 876 bump_author(&mut bucket.authors, author); 877 } 878 } 879 if !current.is_empty() { 880 inner.projected.insert(source, current); 881 } 882 } 883 884 pub fn refresh_issue_counts( 885 &self, 886 states: &StateIndex<IssueStateKind>, 887 entity: &AtUri<DefaultStr>, 888 ) { 889 self.refresh_state_counts(&self.issue_counts, states, entity, "sh.tangled.repo.issue"); 890 } 891 892 pub fn refresh_pull_counts( 893 &self, 894 states: &StateIndex<PullStatusKind>, 895 entity: &AtUri<DefaultStr>, 896 ) { 897 self.refresh_state_counts(&self.pull_counts, states, entity, "sh.tangled.repo.pull"); 898 } 899 900 fn count_state<K>( 901 &self, 902 index: &StateCountIndex<K>, 903 key: &EdgeKey, 904 kind: K, 905 author: Option<&Did<DefaultStr>>, 906 ) -> FilteredCount 907 where 908 K: StateKind, 909 { 910 let Some(key) = self.lookup_key(key) else { 911 return FilteredCount::default(); 912 }; 913 let author = match author { 914 Some(author) => match self.did_interner.get(author).map(AuthorId::from_spur) { 915 Some(author) => Some(author), 916 None => return FilteredCount::default(), 917 }, 918 None => None, 919 }; 920 let inner = index 921 .inner 922 .read() 923 .expect("state-count index rwlock poisoned"); 924 let Some(bucket) = inner.buckets.get(&(key, kind)) else { 925 return FilteredCount::default(); 926 }; 927 match author { 928 Some(author) => { 929 let count = bucket 930 .authors 931 .get(&author) 932 .map_or(0, |count| count.get() as u64); 933 FilteredCount::new(count, u64::from(count != 0)) 934 } 935 None => FilteredCount::new(bucket.count, bucket.authors.len() as u64), 936 } 937 } 938 939 pub fn count_issue_state( 940 &self, 941 key: &EdgeKey, 942 kind: IssueStateKind, 943 author: Option<&Did<DefaultStr>>, 944 ) -> FilteredCount { 945 self.count_state(&self.issue_counts, key, kind, author) 946 } 947 948 pub fn count_pull_status( 949 &self, 950 key: &EdgeKey, 951 kind: PullStatusKind, 952 author: Option<&Did<DefaultStr>>, 953 ) -> FilteredCount { 954 self.count_state(&self.pull_counts, key, kind, author) 955 } 956 957 /// answers "did the viewer star/follow/etc. this subject, and with what rkey" 958 pub fn viewer_source(&self, key: &EdgeKey, viewer: &str) -> Option<AtUri<DefaultStr>> { 959 let author_spur = self.did_interner.get(viewer)?; 960 let author_id = AuthorId::from_spur(author_spur); 961 let key_id = self.lookup_key(key)?; 962 963 self.forward 964 .read_sync(&key_id, |_, sources| { 965 // large buckets track authors, so a missing author fast-fails the scan 966 if let Sources::Large(big) = sources 967 && !big.authors.contains_key(&author_id) 968 { 969 return None; 970 } 971 sources 972 .directed(PageCursor::Start, SortDir::Desc) 973 .find_map(|bucket| { 974 let spur = bucket.source.to_spur()?; 975 let stored = self.source_interner.try_resolve(&spur)?; 976 let uri = AtUri::new_owned(self.decode_source(stored)?).ok()?; 977 (source_authority_did(&uri) == Some(viewer)).then_some(uri) 978 }) 979 }) 980 .flatten() 981 } 982 983 pub fn sources_for(&self, key: &EdgeKey) -> Vec<AtUri<DefaultStr>> { 984 self.lookup_key(key) 985 .and_then(|id| { 986 self.forward.read_sync(&id, |_, sources| { 987 sources 988 .directed(PageCursor::Start, SortDir::Desc) 989 .filter_map(|bucket| { 990 let spur = bucket.source.to_spur()?; 991 let stored = self.source_interner.try_resolve(&spur)?; 992 AtUri::new_owned(self.decode_source(stored)?).ok() 993 }) 994 .collect::<Vec<_>>() 995 }) 996 }) 997 .unwrap_or_default() 998 } 999 1000 pub fn list( 1001 &self, 1002 key: &EdgeKey, 1003 cursor: PageCursor, 1004 limit: PageLimit, 1005 dir: SortDir, 1006 ) -> EdgePage { 1007 let limit_usize = limit.get() as usize; 1008 self.lookup_key(key) 1009 .and_then(|id| { 1010 self.forward.read_sync(&id, |_, sources| { 1011 let iter = sources.directed(cursor, dir); 1012 let entries: Vec<BucketKey> = iter.take(limit_usize + 1).collect(); 1013 let has_more = entries.len() > limit_usize; 1014 let page = &entries[..entries.len().min(limit_usize)]; 1015 let items = page 1016 .iter() 1017 .filter_map(|&key| { 1018 let spur = key.source.to_spur()?; 1019 let stored = self.source_interner.try_resolve(&spur)?; 1020 let uri = AtUri::new_owned(self.decode_source(stored)?).ok()?; 1021 Some(EdgeItem { 1022 uri, 1023 sort_micros: key.micros.0, 1024 }) 1025 }) 1026 .collect(); 1027 let next = has_more 1028 .then(|| page.last().copied()) 1029 .flatten() 1030 .map(BucketKey::token); 1031 EdgePage { items, next } 1032 }) 1033 }) 1034 .unwrap_or(EdgePage { 1035 items: Vec::new(), 1036 next: None, 1037 }) 1038 } 1039 1040 /// Merge a page across many buckets in one global time order. 1041 pub fn list_multi( 1042 &self, 1043 keys: &[EdgeKey], 1044 cursor: PageCursor, 1045 limit: PageLimit, 1046 dir: SortDir, 1047 ) -> EdgePage { 1048 let n = limit.get() as usize + 1; 1049 let is_first = |x: &BucketKey, y: &BucketKey| match dir { 1050 SortDir::Asc => x <= y, 1051 SortDir::Desc => x >= y, 1052 }; 1053 let mut top: Vec<BucketKey> = Vec::with_capacity(n); 1054 for key in keys { 1055 let Some(id) = self.lookup_key(key) else { 1056 continue; 1057 }; 1058 let bucket: Vec<BucketKey> = self 1059 .forward 1060 .read_sync(&id, |_, sources| { 1061 sources.directed(cursor, dir).take(n).collect() 1062 }) 1063 .unwrap_or_default(); 1064 top = top 1065 .into_iter() 1066 .merge_by(bucket, &is_first) 1067 .take(n) 1068 .collect(); 1069 } 1070 top.dedup_by_key(|k| k.source); 1071 1072 let has_more = top.len() > limit.get() as usize; 1073 let page = &top[..top.len().min(limit.get() as usize)]; 1074 let items = page 1075 .iter() 1076 .filter_map(|&k| { 1077 let spur = k.source.to_spur()?; 1078 let stored = self.source_interner.try_resolve(&spur)?; 1079 let uri = AtUri::new_owned(self.decode_source(stored)?).ok()?; 1080 Some(EdgeItem { 1081 uri, 1082 sort_micros: k.micros.0, 1083 }) 1084 }) 1085 .collect(); 1086 let next = has_more 1087 .then(|| page.last().copied()) 1088 .flatten() 1089 .map(BucketKey::token); 1090 EdgePage { items, next } 1091 } 1092 1093 pub fn list_filtered<F>( 1094 &self, 1095 key: &EdgeKey, 1096 cursor: PageCursor, 1097 limit: PageLimit, 1098 dir: SortDir, 1099 predicate: F, 1100 ) -> EdgePage 1101 where 1102 F: Fn(&AtUri<DefaultStr>) -> bool, 1103 { 1104 let limit_usize = limit.get() as usize; 1105 let scan_cap = limit_usize 1106 .saturating_mul(FILTER_SCAN_MULTIPLIER) 1107 .max(FILTER_SCAN_FLOOR); 1108 self.lookup_key(key) 1109 .and_then(|id| { 1110 self.forward.read_sync(&id, |_, sources| { 1111 let init = ScanState::with_capacity(limit_usize + 1); 1112 let outcome = sources 1113 .directed(cursor, dir) 1114 .try_fold(init, |mut state, key| { 1115 if state.scanned >= scan_cap && state.matched.len() <= limit_usize { 1116 return ControlFlow::Break(state); 1117 } 1118 state.scanned += 1; 1119 state.last_scanned = Some(key); 1120 if let Some(spur) = key.source.to_spur() 1121 && let Some(stored) = self.source_interner.try_resolve(&spur) 1122 && let Some(decoded) = self.decode_source(stored) 1123 && let Ok(uri) = AtUri::new_owned(decoded) 1124 && predicate(&uri) 1125 { 1126 state.matched.push((key, uri)); 1127 if state.matched.len() > limit_usize { 1128 return ControlFlow::Break(state); 1129 } 1130 } 1131 ControlFlow::Continue(state) 1132 }); 1133 let (state, bucket_exhausted) = match outcome { 1134 ControlFlow::Continue(s) => (s, true), 1135 ControlFlow::Break(s) => (s, false), 1136 }; 1137 let has_more_matches = state.matched.len() > limit_usize; 1138 let visible_len = state.matched.len().min(limit_usize); 1139 let next = if has_more_matches && visible_len > 0 { 1140 Some(state.matched[visible_len - 1].0.token()) 1141 } else if !bucket_exhausted { 1142 state.last_scanned.map(BucketKey::token) 1143 } else { 1144 None 1145 }; 1146 let items = state 1147 .matched 1148 .into_iter() 1149 .take(visible_len) 1150 .map(|(key, uri)| EdgeItem { 1151 uri, 1152 sort_micros: key.micros.0, 1153 }) 1154 .collect(); 1155 EdgePage { items, next } 1156 }) 1157 }) 1158 .unwrap_or(EdgePage { 1159 items: Vec::new(), 1160 next: None, 1161 }) 1162 } 1163 1164 pub fn key_count(&self) -> usize { 1165 self.forward.len() 1166 } 1167 1168 pub fn source_count(&self) -> usize { 1169 self.reverse.len() 1170 } 1171 1172 pub fn mem_report(&self) -> EdgeMemReport { 1173 const SCC_SLOT: u64 = 32; 1174 let edge_key = std::mem::size_of::<EdgeKey>() as u64; 1175 let rev_entry = std::mem::size_of::<ReverseEntry>() as u64; 1176 let rev_smallvec = std::mem::size_of::<SmallVec<[ReverseEntry; 1]>>() as u64; 1177 let bucket_struct = std::mem::size_of::<Sources>() as u64; 1178 1179 let mut edges_total = 0u64; 1180 let mut author_refs_total = 0u64; 1181 let mut forward_struct_bytes = 1182 self.issue_counts.heap_bytes() + self.pull_counts.heap_bytes(); 1183 let mut max_bucket = 0u64; 1184 let mut bucket_size_classes = [0u64; BUCKET_CLASS_COUNT]; 1185 self.forward.iter_sync(|_, sources| { 1186 let bucket_len = sources.len() as u64; 1187 edges_total += bucket_len; 1188 max_bucket = max_bucket.max(bucket_len); 1189 bucket_size_classes[bucket_class(bucket_len)] += 1; 1190 if let Sources::Large(big) = sources { 1191 author_refs_total += big.authors.len() as u64; 1192 } 1193 forward_struct_bytes += SCC_SLOT + bucket_struct + sources.heap_bytes(); 1194 true 1195 }); 1196 let key_interner_bytes = self.key_ids.len() as u64 1197 * (edge_key + std::mem::size_of::<EdgeKeyId>() as u64 + SCC_SLOT) 1198 + self.keys.len() as u64 1199 * (edge_key + std::mem::size_of::<EdgeKeyId>() as u64 + SCC_SLOT); 1200 1201 let mut reverse_entries = 0u64; 1202 let mut reverse_cap = 0u64; 1203 let mut reverse_struct_bytes = 0u64; 1204 self.reverse.iter_sync(|_, refs| { 1205 let cap = refs.capacity() as u64; 1206 reverse_entries += refs.len() as u64; 1207 reverse_cap += cap; 1208 let heap = if refs.spilled() { cap * rev_entry } else { 0 }; 1209 reverse_struct_bytes += SCC_SLOT + rev_smallvec + heap; 1210 true 1211 }); 1212 1213 EdgeMemReport { 1214 key_count: self.forward.len() as u64, 1215 source_count: self.reverse.len() as u64, 1216 source_interner_bytes: self.source_interner.current_memory_usage() as u64, 1217 did_interner_bytes: self.did_interner.current_memory_usage() as u64, 1218 collection_interner_bytes: self.collection_interner.current_memory_usage() as u64, 1219 key_interner_bytes, 1220 edges_total, 1221 author_refs_total, 1222 reverse_entries, 1223 reverse_cap, 1224 forward_struct_bytes, 1225 reverse_struct_bytes, 1226 bucket_struct_bytes: bucket_struct, 1227 max_bucket, 1228 bucket_size_classes, 1229 } 1230 } 1231} 1232 1233pub fn upsert_record_indexes( 1234 edges: &EdgeStore, 1235 issue_states: &StateIndex<IssueStateKind>, 1236 pull_statuses: &StateIndex<PullStatusKind>, 1237 source: &AtUri<DefaultStr>, 1238 record_edges: Vec<Edge>, 1239 record: &Record, 1240) -> ApplyOutcome { 1241 let previous_state_entity = match record { 1242 Record::IssueState(_) => issue_states.entity_for_source(source), 1243 Record::PullStatus(_) => pull_statuses.entity_for_source(source), 1244 _ => None, 1245 }; 1246 edges.upsert_source(source, record_edges); 1247 let outcome = apply_record_state(issue_states, pull_statuses, source, record); 1248 1249 match record { 1250 Record::Issue(_) => edges.refresh_issue_counts(issue_states, source), 1251 Record::Pull(_) => edges.refresh_pull_counts(pull_statuses, source), 1252 Record::IssueState(state) => { 1253 if let Some(previous) = previous_state_entity.as_ref() 1254 && previous != &state.issue 1255 { 1256 edges.refresh_issue_counts(issue_states, previous); 1257 } 1258 edges.refresh_issue_counts(issue_states, &state.issue); 1259 } 1260 Record::PullStatus(status) => { 1261 if let Some(previous) = previous_state_entity.as_ref() 1262 && previous != &status.pull 1263 { 1264 edges.refresh_pull_counts(pull_statuses, previous); 1265 } 1266 edges.refresh_pull_counts(pull_statuses, &status.pull); 1267 } 1268 _ => {} 1269 } 1270 outcome 1271} 1272 1273pub fn delete_record_indexes( 1274 edges: &EdgeStore, 1275 issue_states: &StateIndex<IssueStateKind>, 1276 pull_statuses: &StateIndex<PullStatusKind>, 1277 source: &AtUri<DefaultStr>, 1278 nsid: &Nsid<DefaultStr>, 1279) { 1280 edges.remove_source(source); 1281 match nsid.as_ref() { 1282 "sh.tangled.repo.issue" => { 1283 issue_states.remove_entity(source); 1284 edges.refresh_issue_counts(issue_states, source); 1285 } 1286 "sh.tangled.repo.pull" => { 1287 pull_statuses.remove_entity(source); 1288 edges.refresh_pull_counts(pull_statuses, source); 1289 } 1290 "sh.tangled.repo.issue.state" => { 1291 if let Some(entity) = issue_states.remove_source(source) { 1292 edges.refresh_issue_counts(issue_states, &entity); 1293 } 1294 } 1295 "sh.tangled.repo.pull.status" => { 1296 if let Some(entity) = pull_statuses.remove_source(source) { 1297 edges.refresh_pull_counts(pull_statuses, &entity); 1298 } 1299 } 1300 _ => {} 1301 } 1302} 1303 1304fn directed_slice( 1305 sources: &[BucketKey], 1306 cursor: PageCursor, 1307 dir: SortDir, 1308) -> impl Iterator<Item = BucketKey> + '_ { 1309 match dir { 1310 SortDir::Asc => { 1311 let start = match cursor { 1312 PageCursor::Start => 0, 1313 PageCursor::After(tok) => { 1314 sources.partition_point(|k| *k <= BucketKey::from_token(tok)) 1315 } 1316 }; 1317 Either::Left(sources[start..].iter().copied()) 1318 } 1319 SortDir::Desc => { 1320 let end = match cursor { 1321 PageCursor::Start => sources.len(), 1322 PageCursor::After(tok) => { 1323 sources.partition_point(|k| *k < BucketKey::from_token(tok)) 1324 } 1325 }; 1326 Either::Right(sources[..end].iter().rev().copied()) 1327 } 1328 } 1329} 1330 1331fn directed_tree( 1332 keys: &BTreeSet<BucketKey>, 1333 cursor: PageCursor, 1334 dir: SortDir, 1335) -> impl Iterator<Item = BucketKey> + '_ { 1336 match dir { 1337 SortDir::Asc => { 1338 let lower = match cursor { 1339 PageCursor::Start => Bound::Unbounded, 1340 PageCursor::After(tok) => Bound::Excluded(BucketKey::from_token(tok)), 1341 }; 1342 Either::Left(keys.range((lower, Bound::Unbounded)).copied()) 1343 } 1344 SortDir::Desc => { 1345 let upper = match cursor { 1346 PageCursor::Start => Bound::Unbounded, 1347 PageCursor::After(tok) => Bound::Excluded(BucketKey::from_token(tok)), 1348 }; 1349 Either::Right(keys.range((Bound::Unbounded, upper)).rev().copied()) 1350 } 1351 } 1352} 1353 1354fn split_record_uri(source: &str) -> Option<(&str, &str, &str)> { 1355 let rest = source.strip_prefix("at://")?; 1356 let mut parts = rest.split('/'); 1357 let authority = parts.next()?; 1358 let collection = parts.next()?; 1359 let rkey = parts.next()?; 1360 if parts.next().is_some() { 1361 return None; 1362 } 1363 (authority.starts_with("did:") && !collection.is_empty() && !rkey.is_empty()) 1364 .then_some((authority, collection, rkey)) 1365} 1366 1367fn source_authority_did(source: &AtUri<DefaultStr>) -> Option<&str> { 1368 let rest = source.as_ref().strip_prefix("at://")?; 1369 let end = rest.find('/').unwrap_or(rest.len()); 1370 let candidate = &rest[..end]; 1371 candidate.starts_with("did:").then_some(candidate) 1372} 1373 1374#[cfg(test)] 1375mod tests { 1376 use super::*; 1377 use bobbin_types::ids::SubjectRef; 1378 use jacquard_common::types::did::Did; 1379 use jacquard_common::types::nsid::Nsid; 1380 use jacquard_common::types::string::AtUri; 1381 1382 fn store() -> EdgeStore { 1383 EdgeStore::new(RuntimeHasher::default()) 1384 } 1385 1386 fn nsid(s: &'static str) -> Nsid<DefaultStr> { 1387 Nsid::new_static(s).unwrap() 1388 } 1389 1390 fn at(s: &str) -> AtUri<DefaultStr> { 1391 AtUri::new_owned(s).unwrap() 1392 } 1393 1394 fn did(s: &str) -> Did<DefaultStr> { 1395 Did::new_owned(s).unwrap() 1396 } 1397 1398 fn did_subj(s: &str) -> SubjectRef { 1399 SubjectRef::Did(did(s)) 1400 } 1401 1402 fn limit(n: u32) -> PageLimit { 1403 PageLimit::new(n).unwrap() 1404 } 1405 1406 const NAMES: [&str; 5] = ["nel", "olaren", "teq", "lyna", "bailey"]; 1407 1408 fn star_edge(source: AtUri<DefaultStr>, subject: Did<DefaultStr>) -> Edge { 1409 star_edge_at(source, subject, 0) 1410 } 1411 1412 fn star_edge_at(source: AtUri<DefaultStr>, subject: Did<DefaultStr>, sort_micros: u64) -> Edge { 1413 Edge { 1414 kind: nsid("sh.tangled.feed.star"), 1415 subject: SubjectRef::Did(subject), 1416 source, 1417 sort_micros, 1418 } 1419 } 1420 1421 fn shuffled_micros(i: usize) -> u64 { 1422 (i as u64).wrapping_mul(2_654_435_761) % 64 1423 } 1424 1425 fn source_uri(i: usize) -> String { 1426 format!( 1427 "at://did:plc:{}/sh.tangled.feed.star/r{i}", 1428 NAMES[i % NAMES.len()] 1429 ) 1430 } 1431 1432 fn fill_subject(store: &EdgeStore, n: usize) -> EdgeKey { 1433 let subject = did("did:plc:squid"); 1434 (0..n).for_each(|i| { 1435 store.add(star_edge_at( 1436 at(&source_uri(i)), 1437 subject.clone(), 1438 shuffled_micros(i), 1439 )); 1440 }); 1441 EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:squid")) 1442 } 1443 1444 fn reference(kept: impl Iterator<Item = usize>) -> Vec<String> { 1445 let mut rows: Vec<(u64, usize)> = kept.map(|i| (shuffled_micros(i), i)).collect(); 1446 rows.sort_unstable(); 1447 rows.into_iter().map(|(_, i)| source_uri(i)).collect() 1448 } 1449 1450 fn paginate_all(store: &EdgeStore, key: &EdgeKey, dir: SortDir, page: u32) -> Vec<String> { 1451 std::iter::successors( 1452 Some(store.list(key, PageCursor::Start, limit(page), dir)), 1453 |prev| { 1454 prev.next 1455 .map(|tok| store.list(key, PageCursor::After(tok), limit(page), dir)) 1456 }, 1457 ) 1458 .flat_map(|p| { 1459 p.items 1460 .into_iter() 1461 .map(|u| u.as_ref().to_owned()) 1462 .collect::<Vec<_>>() 1463 }) 1464 .collect() 1465 } 1466 1467 #[test] 1468 fn pagination_matches_reference_across_small_and_large() { 1469 [64usize, 1000].into_iter().for_each(|n| { 1470 let store = store(); 1471 let key = fill_subject(&store, n); 1472 assert_eq!(store.count(&key), n as u64, "count mismatch at n={n}"); 1473 assert_eq!( 1474 store.count_distinct_authors(&key), 1475 NAMES.len() as u64, 1476 "distinct authors at n={n}" 1477 ); 1478 1479 let asc = reference(0..n); 1480 let desc: Vec<String> = asc.iter().rev().cloned().collect(); 1481 1482 [3u32, 7, 50].into_iter().for_each(|page| { 1483 assert_eq!( 1484 paginate_all(&store, &key, SortDir::Asc, page), 1485 asc, 1486 "asc mismatch n={n} page={page}" 1487 ); 1488 assert_eq!( 1489 paginate_all(&store, &key, SortDir::Desc, page), 1490 desc, 1491 "desc mismatch n={n} page={page}" 1492 ); 1493 }); 1494 }); 1495 } 1496 1497 // Build three subject buckets with interleaved micros so the global time 1498 // order crosses buckets, and return the three keys. 1499 fn fill_three_subjects(store: &EdgeStore) -> Vec<EdgeKey> { 1500 [ 1501 ("alpha", "r0", 10u64), 1502 ("bravo", "r1", 50), 1503 ("charlie", "r2", 20), 1504 ("alpha", "r3", 40), 1505 ("bravo", "r4", 30), 1506 ("charlie", "r5", 60), 1507 ] 1508 .into_iter() 1509 .map(|(subj, rkey, micros)| { 1510 star_edge_at( 1511 at(&format!("at://did:plc:x/sh.tangled.feed.star/{rkey}")), 1512 did(&format!("did:plc:{subj}")), 1513 micros, 1514 ) 1515 }) 1516 .for_each(|edge| store.add(edge)); 1517 ["alpha", "bravo", "charlie"] 1518 .iter() 1519 .map(|s| { 1520 EdgeKey::new( 1521 nsid("sh.tangled.feed.star"), 1522 did_subj(&format!("did:plc:{s}")), 1523 ) 1524 }) 1525 .collect() 1526 } 1527 1528 #[test] 1529 fn list_multi_merges_subjects_newest_first() { 1530 let store = store(); 1531 let keys = fill_three_subjects(&store); 1532 let page = store.list_multi(&keys, PageCursor::Start, limit(10), SortDir::Desc); 1533 let micros: Vec<u64> = page.items.iter().map(|it| it.sort_micros).collect(); 1534 assert_eq!(micros, vec![60, 50, 40, 30, 20, 10]); 1535 assert!(page.next.is_none()); 1536 } 1537 1538 #[test] 1539 fn list_multi_paginates_disjoint_across_buckets() { 1540 let store = store(); 1541 let keys = fill_three_subjects(&store); 1542 // limit=4 with 6 items across 3 buckets exercises the bounded top-N merge. 1543 let p1 = store.list_multi(&keys, PageCursor::Start, limit(4), SortDir::Desc); 1544 let m1: Vec<u64> = p1.items.iter().map(|it| it.sort_micros).collect(); 1545 assert_eq!(m1, vec![60, 50, 40, 30]); 1546 let tok = p1.next.expect("first page has more"); 1547 let p2 = store.list_multi(&keys, PageCursor::After(tok), limit(4), SortDir::Desc); 1548 let m2: Vec<u64> = p2.items.iter().map(|it| it.sort_micros).collect(); 1549 assert_eq!(m2, vec![20, 10]); 1550 assert!(p2.next.is_none()); 1551 } 1552 1553 #[test] 1554 fn list_multi_dedups_source_under_two_keys() { 1555 let store = store(); 1556 // Same record (source) indexed under two subject keys, same createdAt. 1557 let shared = at("at://did:plc:x/sh.tangled.feed.star/shared"); 1558 store.add(star_edge_at(shared.clone(), did("did:plc:alpha"), 100)); 1559 store.add(star_edge_at(shared, did("did:plc:bravo"), 100)); 1560 let keys = vec![ 1561 EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:alpha")), 1562 EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:bravo")), 1563 ]; 1564 let page = store.list_multi(&keys, PageCursor::Start, limit(10), SortDir::Desc); 1565 assert_eq!(page.items.len(), 1); 1566 assert!(page.next.is_none()); 1567 } 1568 1569 #[test] 1570 fn remove_from_large_bucket_keeps_pagination_exact() { 1571 let store = store(); 1572 let n = 1000usize; 1573 let key = fill_subject(&store, n); 1574 (0..n).step_by(3).for_each(|i| { 1575 store.remove_source(&at(&source_uri(i))); 1576 }); 1577 let expected = reference((0..n).filter(|i| i % 3 != 0)); 1578 assert_eq!(store.count(&key), expected.len() as u64); 1579 assert_eq!( 1580 store.count_distinct_authors(&key), 1581 NAMES.len() as u64, 1582 "every author keeps sources after partial removal" 1583 ); 1584 assert_eq!(paginate_all(&store, &key, SortDir::Asc, 7), expected); 1585 assert_eq!( 1586 paginate_all(&store, &key, SortDir::Desc, 11), 1587 expected.iter().rev().cloned().collect::<Vec<_>>() 1588 ); 1589 } 1590 1591 #[test] 1592 fn add_then_count() { 1593 let store = store(); 1594 let key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:abalone")); 1595 1596 store.add(star_edge( 1597 at("at://did:plc:nel/sh.tangled.feed.star/r1"), 1598 did("did:plc:abalone"), 1599 )); 1600 store.add(star_edge( 1601 at("at://did:plc:olaren/sh.tangled.feed.star/r2"), 1602 did("did:plc:abalone"), 1603 )); 1604 store.add(star_edge( 1605 at("at://did:plc:nel/sh.tangled.feed.star/r3"), 1606 did("did:plc:abalone"), 1607 )); 1608 1609 assert_eq!(store.count(&key), 3); 1610 assert_eq!(store.count_distinct_authors(&key), 2); 1611 } 1612 1613 #[test] 1614 fn duplicate_add_is_idempotent() { 1615 let store = store(); 1616 let key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:abalone")); 1617 let edge = star_edge( 1618 at("at://did:plc:nel/sh.tangled.feed.star/r1"), 1619 did("did:plc:abalone"), 1620 ); 1621 store.add(edge.clone()); 1622 store.add(edge); 1623 assert_eq!(store.count(&key), 1); 1624 assert_eq!(store.count_distinct_authors(&key), 1); 1625 } 1626 1627 #[test] 1628 fn remove_source_clears_all_keys_for_that_source() { 1629 let store = store(); 1630 let star_key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:abalone")); 1631 let follow_key = EdgeKey::new(nsid("sh.tangled.graph.follow"), did_subj("did:plc:lyna")); 1632 let source = "at://did:plc:nel/sh.tangled.feed.star/r1"; 1633 1634 store.add(Edge { 1635 kind: nsid("sh.tangled.feed.star"), 1636 subject: did_subj("did:plc:abalone"), 1637 source: at(source), 1638 sort_micros: 0, 1639 }); 1640 store.add(Edge { 1641 kind: nsid("sh.tangled.graph.follow"), 1642 subject: did_subj("did:plc:lyna"), 1643 source: at(source), 1644 sort_micros: 0, 1645 }); 1646 1647 assert_eq!(store.count(&star_key), 1); 1648 assert_eq!(store.count(&follow_key), 1); 1649 1650 store.remove_source(&at(source)); 1651 assert_eq!(store.count(&star_key), 0); 1652 assert_eq!(store.count(&follow_key), 0); 1653 } 1654 1655 #[test] 1656 fn upsert_source_replaces_old_edges() { 1657 let store = store(); 1658 let source = at("at://did:plc:teq/sh.tangled.feed.star/r1"); 1659 let old_subject = did_subj("did:plc:abalone"); 1660 let new_subject = did_subj("did:plc:uni"); 1661 let kind = nsid("sh.tangled.feed.star"); 1662 1663 store.upsert_source( 1664 &source, 1665 vec![Edge { 1666 kind: kind.clone(), 1667 subject: old_subject.clone(), 1668 source: source.clone(), 1669 sort_micros: 0, 1670 }], 1671 ); 1672 assert_eq!( 1673 store.count(&EdgeKey::new(kind.clone(), old_subject.clone())), 1674 1 1675 ); 1676 1677 store.upsert_source( 1678 &source, 1679 vec![Edge { 1680 kind: kind.clone(), 1681 subject: new_subject.clone(), 1682 source: source.clone(), 1683 sort_micros: 0, 1684 }], 1685 ); 1686 assert_eq!(store.count(&EdgeKey::new(kind.clone(), old_subject)), 0); 1687 assert_eq!(store.count(&EdgeKey::new(kind, new_subject)), 1); 1688 } 1689 1690 #[test] 1691 fn state_counts_reproject_after_state_first_ingest_and_rekey() { 1692 let store = store(); 1693 let states = StateIndex::new(RuntimeHasher::default()); 1694 let issue = at("at://did:plc:nel/sh.tangled.repo.issue/i1"); 1695 let old_repo = did("did:plc:limpet"); 1696 let new_repo = did("did:plc:scallop"); 1697 let kind = nsid("sh.tangled.repo.issue"); 1698 let old_key = EdgeKey::new(kind.clone(), SubjectRef::Did(old_repo.clone())); 1699 let new_key = EdgeKey::new(kind.clone(), SubjectRef::Did(new_repo.clone())); 1700 1701 states.upsert( 1702 at("at://did:plc:limpet/sh.tangled.repo.issue.state/s1"), 1703 issue.clone(), 1704 100, 1705 IssueStateKind::Closed, 1706 ); 1707 store.upsert_source( 1708 &issue, 1709 vec![Edge { 1710 kind: kind.clone(), 1711 subject: SubjectRef::Did(old_repo), 1712 source: issue.clone(), 1713 sort_micros: 1, 1714 }], 1715 ); 1716 store.refresh_issue_counts(&states, &issue); 1717 assert_eq!( 1718 store.count_issue_state(&old_key, IssueStateKind::Closed, None), 1719 FilteredCount::new(1, 1) 1720 ); 1721 1722 store.upsert_source( 1723 &issue, 1724 vec![Edge { 1725 kind, 1726 subject: SubjectRef::Did(new_repo), 1727 source: issue.clone(), 1728 sort_micros: 2, 1729 }], 1730 ); 1731 store.refresh_issue_counts(&states, &issue); 1732 assert_eq!( 1733 store.count_issue_state(&old_key, IssueStateKind::Closed, None), 1734 FilteredCount::default() 1735 ); 1736 assert_eq!( 1737 store.count_issue_state(&new_key, IssueStateKind::Open, Some(&did("did:plc:nel")),), 1738 FilteredCount::new(1, 1), 1739 "the old repo owner's state stops being accepted after the rekey", 1740 ); 1741 } 1742 1743 #[test] 1744 fn record_index_mutations_move_materialized_state_counts() { 1745 use bobbin_types::sh_tangled::repo::issue::state::{State as IssueStateRecord, StateState}; 1746 1747 let store = store(); 1748 let issue_states = StateIndex::new(RuntimeHasher::default()); 1749 let pull_statuses = StateIndex::new(RuntimeHasher::default()); 1750 let issue = at("at://did:plc:nel/sh.tangled.repo.issue/i1"); 1751 let repo = did("did:plc:limpet"); 1752 let key = EdgeKey::new(nsid("sh.tangled.repo.issue"), SubjectRef::Did(repo.clone())); 1753 store.upsert_source( 1754 &issue, 1755 vec![Edge { 1756 kind: nsid("sh.tangled.repo.issue"), 1757 subject: SubjectRef::Did(repo), 1758 source: issue.clone(), 1759 sort_micros: 1, 1760 }], 1761 ); 1762 store.refresh_issue_counts(&issue_states, &issue); 1763 1764 let state_source = at("at://did:plc:limpet/sh.tangled.repo.issue.state/s1"); 1765 let state_record = Record::IssueState(IssueStateRecord { 1766 issue: issue.clone(), 1767 created_at: jacquard_common::types::string::Datetime::raw_str("2026-05-01T00:00:00Z"), 1768 state: StateState::ShTangledRepoIssueStateClosed, 1769 extra_data: None, 1770 }); 1771 upsert_record_indexes( 1772 &store, 1773 &issue_states, 1774 &pull_statuses, 1775 &state_source, 1776 Vec::new(), 1777 &state_record, 1778 ); 1779 assert_eq!( 1780 store 1781 .count_issue_state(&key, IssueStateKind::Closed, None) 1782 .count 1783 .get(), 1784 1, 1785 ); 1786 1787 delete_record_indexes( 1788 &store, 1789 &issue_states, 1790 &pull_statuses, 1791 &state_source, 1792 &nsid("sh.tangled.repo.issue.state"), 1793 ); 1794 assert_eq!( 1795 store 1796 .count_issue_state(&key, IssueStateKind::Open, None) 1797 .count 1798 .get(), 1799 1, 1800 ); 1801 } 1802 1803 #[test] 1804 fn list_pages_in_sort_order() { 1805 let store = store(); 1806 let key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:abalone")); 1807 (0..5).for_each(|i| { 1808 store.add(star_edge_at( 1809 at(&format!( 1810 "at://did:plc:{}/sh.tangled.feed.star/r{i}", 1811 NAMES[i] 1812 )), 1813 did("did:plc:abalone"), 1814 1_000_000 + i as u64 * 1_000_000, 1815 )); 1816 }); 1817 1818 let page1 = store.list(&key, PageCursor::Start, limit(2), SortDir::Asc); 1819 assert_eq!(page1.items.len(), 2); 1820 let cursor = PageCursor::After(page1.next.expect("has next")); 1821 1822 let page2 = store.list(&key, cursor, limit(2), SortDir::Asc); 1823 assert_eq!(page2.items.len(), 2); 1824 1825 let cursor2 = PageCursor::After(page2.next.expect("has next")); 1826 let page3 = store.list(&key, cursor2, limit(2), SortDir::Asc); 1827 assert_eq!(page3.items.len(), 1); 1828 assert!( 1829 page3.next.is_none(), 1830 "final partial page should signal exhaustion" 1831 ); 1832 } 1833 1834 #[test] 1835 fn list_exact_fill_signals_exhaustion() { 1836 let store = store(); 1837 let key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:abalone")); 1838 (0..2).for_each(|i| { 1839 store.add(star_edge( 1840 at(&format!( 1841 "at://did:plc:{}/sh.tangled.feed.star/r{i}", 1842 NAMES[i] 1843 )), 1844 did("did:plc:abalone"), 1845 )); 1846 }); 1847 1848 let page = store.list(&key, PageCursor::Start, limit(2), SortDir::Asc); 1849 assert_eq!(page.items.len(), 2); 1850 assert!(page.next.is_none(), "exact-fill page must not promise more"); 1851 } 1852 1853 #[test] 1854 fn list_on_unknown_key_is_empty() { 1855 let store = store(); 1856 let page = store.list( 1857 &EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:periwinkle")), 1858 PageCursor::Start, 1859 limit(10), 1860 SortDir::Asc, 1861 ); 1862 assert!(page.items.is_empty()); 1863 assert!(page.next.is_none()); 1864 } 1865 1866 #[test] 1867 fn distinct_authors_decreases_when_last_source_from_author_removed() { 1868 let store = store(); 1869 let key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj("did:plc:abalone")); 1870 let s1 = at("at://did:plc:nel/sh.tangled.feed.star/r1"); 1871 let s2 = at("at://did:plc:nel/sh.tangled.feed.star/r2"); 1872 let s3 = at("at://did:plc:olaren/sh.tangled.feed.star/r3"); 1873 1874 store.add(star_edge(s1.clone(), did("did:plc:abalone"))); 1875 store.add(star_edge(s2.clone(), did("did:plc:abalone"))); 1876 store.add(star_edge(s3, did("did:plc:abalone"))); 1877 assert_eq!(store.count_distinct_authors(&key), 2); 1878 1879 store.remove_source(&s1); 1880 assert_eq!(store.count_distinct_authors(&key), 2, "user1 still has s2"); 1881 1882 store.remove_source(&s2); 1883 assert_eq!(store.count_distinct_authors(&key), 1, "user1 fully gone"); 1884 } 1885 1886 #[test] 1887 fn page_limit_rejects_zero_and_oversize() { 1888 assert!(matches!( 1889 PageLimit::new(0), 1890 Err(PageLimitError::TooSmall { value: 0, min: 1 }) 1891 )); 1892 assert!(matches!( 1893 PageLimit::new(PageLimit::MAX + 1), 1894 Err(PageLimitError::TooLarge { .. }) 1895 )); 1896 assert_eq!(PageLimit::new(50).unwrap().get(), 50); 1897 } 1898 1899 #[test] 1900 fn cursor_token_round_trip() { 1901 let original = PageToken::new(1_730_000_000_000_000, 0x1234_abcd); 1902 let token = original.encode_token(); 1903 assert_eq!(token.len(), 24, "12-byte cursor encodes to 24 hex chars"); 1904 assert_eq!(PageToken::decode_token(&token).unwrap(), original); 1905 } 1906 1907 #[test] 1908 fn from_token_none_yields_start() { 1909 assert_eq!(PageCursor::from_token(None).unwrap(), PageCursor::Start); 1910 } 1911 1912 #[test] 1913 fn from_token_some_yields_after() { 1914 let original = PageToken::new(1_730_000_000_000_000, 42); 1915 let token = original.encode_token(); 1916 assert_eq!( 1917 PageCursor::from_token(Some(&token)).unwrap(), 1918 PageCursor::After(original), 1919 ); 1920 } 1921 1922 #[test] 1923 fn cursor_decode_rejects_malformed() { 1924 let bad = [ 1925 "", 1926 "deadbeef", 1927 "no-hex!!aaaaaaaaaaaaaaaa", 1928 "12345", 1929 "this-string-is-way-too-long-to-be-a-valid-cursor", 1930 ]; 1931 bad.into_iter().for_each(|s| { 1932 assert!( 1933 matches!(PageToken::decode_token(s), Err(CursorParseError::Malformed)), 1934 "expected malformed for {s:?}", 1935 ); 1936 }); 1937 } 1938 1939 #[test] 1940 fn cursor_token_is_hex_shape() { 1941 let token = PageToken::new(1_730_000_000_000_000, 7).encode_token(); 1942 assert_eq!(token.len(), 24); 1943 assert!(token.chars().all(|c| c.is_ascii_hexdigit())); 1944 } 1945 1946 #[test] 1947 fn list_does_not_drop_entries_at_same_sort_micros_across_pages() { 1948 let store = store(); 1949 let subject_did = did("did:plc:limpet"); 1950 let key = EdgeKey::new( 1951 nsid("sh.tangled.feed.star"), 1952 SubjectRef::Did(subject_did.clone()), 1953 ); 1954 (0..4).for_each(|i| { 1955 store.add(star_edge_at( 1956 at(&format!( 1957 "at://did:plc:{}/sh.tangled.feed.star/r{i}", 1958 NAMES[i] 1959 )), 1960 subject_did.clone(), 1961 42, 1962 )); 1963 }); 1964 1965 let page1 = store.list(&key, PageCursor::Start, limit(2), SortDir::Asc); 1966 assert_eq!(page1.items.len(), 2); 1967 let token = page1.next.expect("cursor must continue across ties"); 1968 1969 let page2 = store.list(&key, PageCursor::After(token), limit(2), SortDir::Asc); 1970 assert_eq!( 1971 page2.items.len(), 1972 2, 1973 "remaining ties must surface on next page" 1974 ); 1975 assert!(page2.next.is_none()); 1976 let combined: std::collections::HashSet<String> = page1 1977 .items 1978 .iter() 1979 .chain(page2.items.iter()) 1980 .map(|u| u.as_ref().to_owned()) 1981 .collect(); 1982 assert_eq!(combined.len(), 4, "every tied entry visible exactly once"); 1983 } 1984 1985 #[test] 1986 fn list_filtered_narrows_by_predicate_and_paginates() { 1987 let store = store(); 1988 let subject_did = did("did:plc:limpet"); 1989 let key = EdgeKey::new( 1990 nsid("sh.tangled.repo.issue"), 1991 SubjectRef::Did(subject_did.clone()), 1992 ); 1993 let by_nel = (0..3).map(|i| { 1994 star_edge_at( 1995 at(&format!("at://did:plc:nel/sh.tangled.repo.issue/r{i}")), 1996 subject_did.clone(), 1997 100 + i as u64, 1998 ) 1999 }); 2000 let by_olaren = (0..2).map(|i| { 2001 star_edge_at( 2002 at(&format!("at://did:plc:olaren/sh.tangled.repo.issue/o{i}")), 2003 subject_did.clone(), 2004 500 + i as u64, 2005 ) 2006 }); 2007 by_nel 2008 .chain(by_olaren) 2009 .map(|mut e| { 2010 e.kind = nsid("sh.tangled.repo.issue"); 2011 e 2012 }) 2013 .for_each(|e| store.add(e)); 2014 2015 let only_nel = |u: &AtUri<DefaultStr>| u.as_ref().starts_with("at://did:plc:nel/"); 2016 let page1 = store.list_filtered(&key, PageCursor::Start, limit(2), SortDir::Asc, only_nel); 2017 assert_eq!(page1.items.len(), 2); 2018 assert!(page1.next.is_some(), "cursor must allow more nel matches"); 2019 2020 let page2 = store.list_filtered( 2021 &key, 2022 PageCursor::After(page1.next.unwrap()), 2023 limit(2), 2024 SortDir::Asc, 2025 only_nel, 2026 ); 2027 assert_eq!(page2.items.len(), 1, "only one nel issue left"); 2028 assert!(page2.next.is_none(), "tail page must not promise more"); 2029 } 2030 2031 #[test] 2032 fn list_descending_returns_newest_first() { 2033 let store = store(); 2034 let subject_did = did("did:plc:limpet"); 2035 let key = EdgeKey::new( 2036 nsid("sh.tangled.feed.star"), 2037 SubjectRef::Did(subject_did.clone()), 2038 ); 2039 (0..5).for_each(|i| { 2040 store.add(star_edge_at( 2041 at(&format!( 2042 "at://did:plc:{}/sh.tangled.feed.star/r{i}", 2043 NAMES[i] 2044 )), 2045 subject_did.clone(), 2046 1_000_000 + i as u64 * 1_000_000, 2047 )); 2048 }); 2049 2050 let asc = store.list(&key, PageCursor::Start, limit(5), SortDir::Asc); 2051 let desc = store.list(&key, PageCursor::Start, limit(5), SortDir::Desc); 2052 assert_eq!(asc.items.len(), 5); 2053 assert_eq!(desc.items.len(), 5); 2054 let asc_uris: Vec<_> = asc.items.iter().map(|u| u.as_ref().to_owned()).collect(); 2055 let mut reversed = asc_uris.clone(); 2056 reversed.reverse(); 2057 let desc_uris: Vec<_> = desc.items.iter().map(|u| u.as_ref().to_owned()).collect(); 2058 assert_eq!(desc_uris, reversed, "desc must be exact reverse of asc"); 2059 } 2060 2061 #[test] 2062 fn list_descending_paginates_with_cursor() { 2063 let store = store(); 2064 let subject_did = did("did:plc:whelk"); 2065 let key = EdgeKey::new( 2066 nsid("sh.tangled.feed.star"), 2067 SubjectRef::Did(subject_did.clone()), 2068 ); 2069 (0..5).for_each(|i| { 2070 store.add(star_edge_at( 2071 at(&format!( 2072 "at://did:plc:{}/sh.tangled.feed.star/r{i}", 2073 NAMES[i] 2074 )), 2075 subject_did.clone(), 2076 1_000_000 + i as u64 * 1_000_000, 2077 )); 2078 }); 2079 2080 let page1 = store.list(&key, PageCursor::Start, limit(2), SortDir::Desc); 2081 assert_eq!(page1.items.len(), 2); 2082 let cursor = PageCursor::After(page1.next.expect("desc page1 must continue")); 2083 let page2 = store.list(&key, cursor, limit(2), SortDir::Desc); 2084 assert_eq!(page2.items.len(), 2); 2085 let cursor2 = PageCursor::After(page2.next.expect("desc page2 must continue")); 2086 let page3 = store.list(&key, cursor2, limit(2), SortDir::Desc); 2087 assert_eq!(page3.items.len(), 1); 2088 assert!(page3.next.is_none()); 2089 2090 let combined: std::collections::HashSet<String> = page1 2091 .items 2092 .iter() 2093 .chain(page2.items.iter()) 2094 .chain(page3.items.iter()) 2095 .map(|u| u.as_ref().to_owned()) 2096 .collect(); 2097 assert_eq!( 2098 combined.len(), 2099 5, 2100 "desc pagination visits each item exactly once" 2101 ); 2102 } 2103 2104 #[test] 2105 fn list_filtered_descending_respects_predicate() { 2106 let store = store(); 2107 let subject_did = did("did:plc:scallop"); 2108 let key = EdgeKey::new( 2109 nsid("sh.tangled.repo.issue"), 2110 SubjectRef::Did(subject_did.clone()), 2111 ); 2112 let by_nel = (0..3).map(|i| { 2113 star_edge_at( 2114 at(&format!("at://did:plc:nel/sh.tangled.repo.issue/r{i}")), 2115 subject_did.clone(), 2116 100 + i as u64, 2117 ) 2118 }); 2119 let by_olaren = (0..2).map(|i| { 2120 star_edge_at( 2121 at(&format!("at://did:plc:olaren/sh.tangled.repo.issue/o{i}")), 2122 subject_did.clone(), 2123 500 + i as u64, 2124 ) 2125 }); 2126 by_nel 2127 .chain(by_olaren) 2128 .map(|mut e| { 2129 e.kind = nsid("sh.tangled.repo.issue"); 2130 e 2131 }) 2132 .for_each(|e| store.add(e)); 2133 2134 let only_nel = |u: &AtUri<DefaultStr>| u.as_ref().starts_with("at://did:plc:nel/"); 2135 let page = store.list_filtered(&key, PageCursor::Start, limit(5), SortDir::Desc, only_nel); 2136 assert_eq!(page.items.len(), 3, "all three nel issues visible"); 2137 let last = page.items.last().unwrap().as_ref(); 2138 let first = page.items.first().unwrap().as_ref(); 2139 assert!( 2140 first > last, 2141 "desc order: first item rkey must be greater than last (got first={first}, last={last})", 2142 ); 2143 } 2144 2145 #[test] 2146 fn non_did_source_round_trips_via_raw_fallback() { 2147 let store = store(); 2148 let subject = did("did:plc:limpet"); 2149 let key = EdgeKey::new( 2150 nsid("sh.tangled.feed.star"), 2151 SubjectRef::Did(subject.clone()), 2152 ); 2153 let did_source = at("at://did:plc:nel/sh.tangled.feed.star/r1"); 2154 let handle_source = at("at://witchcraft.systems/sh.tangled.feed.star/r2"); 2155 store.add(star_edge_at(did_source.clone(), subject.clone(), 1)); 2156 store.add(star_edge_at(handle_source.clone(), subject.clone(), 2)); 2157 2158 let page = store.list(&key, PageCursor::Start, limit(10), SortDir::Asc); 2159 let got: std::collections::HashSet<String> = 2160 page.items.iter().map(|u| u.as_ref().to_owned()).collect(); 2161 assert!( 2162 got.contains(did_source.as_ref()), 2163 "did source must decode exactly" 2164 ); 2165 assert!( 2166 got.contains(handle_source.as_ref()), 2167 "non-did authority must round-trip through the raw fallback" 2168 ); 2169 2170 store.remove_source(&handle_source); 2171 assert_eq!( 2172 store.count(&key), 2173 1, 2174 "raw-keyed source removable by its uri" 2175 ); 2176 } 2177 2178 #[test] 2179 fn distinct_collections_decode_with_their_own_collection() { 2180 let store = store(); 2181 let subject = did("did:plc:limpet"); 2182 let key = EdgeKey::new( 2183 nsid("sh.tangled.feed.star"), 2184 SubjectRef::Did(subject.clone()), 2185 ); 2186 let star_src = at("at://did:plc:nel/sh.tangled.feed.star/aaa"); 2187 let issue_src = at("at://did:plc:nel/sh.tangled.repo.issue/bbb"); 2188 store.add(Edge { 2189 kind: nsid("sh.tangled.feed.star"), 2190 subject: SubjectRef::Did(subject.clone()), 2191 source: star_src.clone(), 2192 sort_micros: 1, 2193 }); 2194 store.add(Edge { 2195 kind: nsid("sh.tangled.feed.star"), 2196 subject: SubjectRef::Did(subject.clone()), 2197 source: issue_src.clone(), 2198 sort_micros: 2, 2199 }); 2200 2201 let page = store.list(&key, PageCursor::Start, limit(10), SortDir::Asc); 2202 let got: std::collections::HashSet<String> = 2203 page.items.iter().map(|u| u.as_ref().to_owned()).collect(); 2204 assert!( 2205 got.contains(star_src.as_ref()), 2206 "star-collection source decodes exactly" 2207 ); 2208 assert!( 2209 got.contains(issue_src.as_ref()), 2210 "issue-collection source must keep its own collection, not borrow the star one" 2211 ); 2212 } 2213 2214 #[test] 2215 fn author_refs_spill_preserves_distinct_count() { 2216 let store = store(); 2217 let subject = did("did:plc:scallop"); 2218 let key = EdgeKey::new( 2219 nsid("sh.tangled.feed.star"), 2220 SubjectRef::Did(subject.clone()), 2221 ); 2222 (0..5).for_each(|i| { 2223 store.add(star_edge_at( 2224 at(&format!( 2225 "at://did:plc:{}/sh.tangled.feed.star/r{i}", 2226 NAMES[i] 2227 )), 2228 subject.clone(), 2229 i as u64, 2230 )); 2231 }); 2232 assert_eq!( 2233 store.count_distinct_authors(&key), 2234 5, 2235 "five authors exceed the inline cap and spill to the map" 2236 ); 2237 2238 store.add(star_edge_at( 2239 at("at://did:plc:nel/sh.tangled.feed.star/r99"), 2240 subject.clone(), 2241 99, 2242 )); 2243 assert_eq!( 2244 store.count_distinct_authors(&key), 2245 5, 2246 "second nel source adds no author" 2247 ); 2248 2249 store.remove_source(&at("at://did:plc:nel/sh.tangled.feed.star/r0")); 2250 assert_eq!( 2251 store.count_distinct_authors(&key), 2252 5, 2253 "nel still present via r99" 2254 ); 2255 store.remove_source(&at("at://did:plc:nel/sh.tangled.feed.star/r99")); 2256 assert_eq!(store.count_distinct_authors(&key), 4, "nel fully removed"); 2257 } 2258}