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