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