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