This repository has no description
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}