From 2a8041aeae3de09936df3048723242507b672a1f Mon Sep 17 00:00:00 2001 From: Lewis Date: Sat, 16 May 2026 22:29:14 +0300 Subject: [PATCH] feat(edge-index): state index paginatione Lewis: May this revision serve well! --- crates/edge-index/src/lib.rs | 254 +++++++++++++-- crates/edge-index/src/state_index.rs | 450 +++++++++++++++++++++++++++ 2 files changed, 676 insertions(+), 28 deletions(-) create mode 100644 crates/edge-index/src/state_index.rs diff --git a/crates/edge-index/src/lib.rs b/crates/edge-index/src/lib.rs index f188173..7fd5a6c 100644 --- a/crates/edge-index/src/lib.rs +++ b/crates/edge-index/src/lib.rs @@ -1,19 +1,44 @@ use std::collections::{BTreeSet, HashMap}; use std::num::NonZeroU32; +use std::ops::{Bound, ControlFlow}; use std::sync::{Arc, Mutex}; +const FILTER_SCAN_MULTIPLIER: usize = 64; +const FILTER_SCAN_FLOOR: usize = 512; + +struct ScanState { + matched: Vec<(BucketKey, AtUri)>, + last_scanned: Option, + scanned: usize, +} + +impl ScanState { + fn with_capacity(cap: usize) -> Self { + Self { + matched: Vec::with_capacity(cap), + last_scanned: None, + scanned: 0, + } + } +} + use bobbin_runtime::RuntimeHasher; use bobbin_types::edges::Edge; use bobbin_types::ids::EdgeKey; use jacquard_common::DefaultStr; use jacquard_common::types::string::AtUri; -use jacquard_common::types::tid::Tid; use lasso::{Key, Spur, ThreadedRodeo}; use scc::HashMap as SccMap; use thiserror::Error; +type BucketKey = (u64, u32); + pub mod coverage; +pub mod state_index; pub use coverage::{Coverage, CoverageWatch, HydrantCursor, PromotionSignal}; +pub use state_index::{ + ApplyOutcome, IssueStateKind, PullStatusKind, StateIndex, StateKind, apply_record_state, +}; #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] pub struct SourceId(u32); @@ -33,28 +58,65 @@ impl SourceId { } #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] -pub struct PageToken(u64); +pub struct PageToken { + micros: u64, + source: u32, +} impl PageToken { - pub fn new(micros: u64) -> Self { - Self(micros) + pub fn new(micros: u64, source: u32) -> Self { + Self { micros, source } } pub fn micros(self) -> u64 { - self.0 + self.micros + } + + pub fn source(self) -> u32 { + self.source } pub fn encode_token(self) -> String { - Tid::from_time(self.0, 0).as_str().to_owned() + let mut bytes = [0u8; 12]; + bytes[..8].copy_from_slice(&self.micros.to_be_bytes()); + bytes[8..].copy_from_slice(&self.source.to_be_bytes()); + encode_hex(&bytes) } pub fn decode_token(token: &str) -> Result { - Tid::new(token) - .map(|t| Self(t.timestamp())) - .map_err(|_| CursorParseError::Malformed) + let bytes: [u8; 12] = decode_hex_array(token).ok_or(CursorParseError::Malformed)?; + let micros = u64::from_be_bytes(bytes[..8].try_into().unwrap()); + let source = u32::from_be_bytes(bytes[8..].try_into().unwrap()); + Ok(Self { micros, source }) } } +fn encode_hex(bytes: &[u8]) -> String { + bytes + .iter() + .fold(String::with_capacity(bytes.len() * 2), |mut acc, b| { + acc.push(char::from_digit((b >> 4) as u32, 16).unwrap()); + acc.push(char::from_digit((b & 0x0f) as u32, 16).unwrap()); + acc + }) +} + +fn decode_hex_array(token: &str) -> Option<[u8; N]> { + if token.len() != N * 2 { + return None; + } + let parsed: Vec = token + .as_bytes() + .chunks_exact(2) + .map(|pair| { + let hi = (pair[0] as char).to_digit(16)?; + let lo = (pair[1] as char).to_digit(16)?; + Some(((hi << 4) | lo) as u8) + }) + .collect::>>()?; + parsed.try_into().ok() +} + #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] struct AuthorId(u32); @@ -77,10 +139,10 @@ pub enum CursorParseError { } impl PageCursor { - fn skip_until(self) -> Option { + fn range_start(self) -> Bound { match self { - Self::Start => None, - Self::After(tok) => Some(tok.micros().saturating_add(1)), + Self::Start => Bound::Unbounded, + Self::After(tok) => Bound::Excluded((tok.micros, tok.source)), } } @@ -132,7 +194,7 @@ pub struct EdgePage { } struct EdgeBucket { - sources: BTreeSet<(u64, u32)>, + sources: BTreeSet, author_refs: HashMap, } @@ -271,10 +333,9 @@ impl EdgeStore { let limit_usize = limit.get() as usize; self.forward .read_sync(key, |_, bucket| { - let start_micros = cursor.skip_until().unwrap_or(0); - let entries: Vec<(u64, u32)> = bucket + let entries: Vec = bucket .sources - .range((start_micros, 0)..) + .range((cursor.range_start(), Bound::Unbounded)) .take(limit_usize + 1) .copied() .collect(); @@ -292,7 +353,73 @@ impl EdgeStore { let next = has_more .then(|| page.last().copied()) .flatten() - .map(|(m, _)| PageToken(m)); + .map(|(m, s)| PageToken::new(m, s)); + EdgePage { items, next } + }) + .unwrap_or(EdgePage { + items: Vec::new(), + next: None, + }) + } + + pub fn list_filtered( + &self, + key: &EdgeKey, + cursor: PageCursor, + limit: PageLimit, + predicate: F, + ) -> EdgePage + where + F: Fn(&AtUri) -> bool, + { + let limit_usize = limit.get() as usize; + let scan_cap = limit_usize + .saturating_mul(FILTER_SCAN_MULTIPLIER) + .max(FILTER_SCAN_FLOOR); + self.forward + .read_sync(key, |_, bucket| { + let init = ScanState::with_capacity(limit_usize + 1); + let outcome = bucket + .sources + .range((cursor.range_start(), Bound::Unbounded)) + .try_fold(init, |mut state, &(m, raw)| { + if state.scanned >= scan_cap && state.matched.len() <= limit_usize { + return ControlFlow::Break(state); + } + state.scanned += 1; + state.last_scanned = Some((m, raw)); + if let Some(spur) = SourceId(raw).to_spur() + && let Some(s) = self.source_interner.try_resolve(&spur) + && let Ok(uri) = AtUri::new_owned(s) + && predicate(&uri) + { + state.matched.push(((m, raw), uri)); + if state.matched.len() > limit_usize { + return ControlFlow::Break(state); + } + } + ControlFlow::Continue(state) + }); + let (state, bucket_exhausted) = match outcome { + ControlFlow::Continue(s) => (s, true), + ControlFlow::Break(s) => (s, false), + }; + let has_more_matches = state.matched.len() > limit_usize; + let visible_len = state.matched.len().min(limit_usize); + let next = if has_more_matches && visible_len > 0 { + let (m, s) = state.matched[visible_len - 1].0; + Some(PageToken::new(m, s)) + } else if !bucket_exhausted { + state.last_scanned.map(|(m, s)| PageToken::new(m, s)) + } else { + None + }; + let items = state + .matched + .into_iter() + .take(visible_len) + .map(|(_, uri)| uri) + .collect(); EdgePage { items, next } }) .unwrap_or(EdgePage { @@ -568,9 +695,9 @@ mod tests { #[test] fn cursor_token_round_trip() { - let original = PageToken::new(1_730_000_000_000_000); + let original = PageToken::new(1_730_000_000_000_000, 0x1234_abcd); let token = original.encode_token(); - assert_eq!(token.len(), 13, "TID encoding is 13 base32 chars"); + assert_eq!(token.len(), 24, "12-byte cursor encodes to 24 hex chars"); assert_eq!(PageToken::decode_token(&token).unwrap(), original); } @@ -581,7 +708,7 @@ mod tests { #[test] fn from_token_some_yields_after() { - let original = PageToken::new(1_730_000_000_000_000); + let original = PageToken::new(1_730_000_000_000_000, 42); let token = original.encode_token(); assert_eq!( PageCursor::from_token(Some(&token)).unwrap(), @@ -593,10 +720,10 @@ mod tests { fn cursor_decode_rejects_malformed() { let bad = [ "", - "deadbeef0", - "no-hex!!", - "1234567", - "way-too-long-not-a-tid", + "deadbeef", + "no-hex!!aaaaaaaaaaaaaaaa", + "12345", + "this-string-is-way-too-long-to-be-a-valid-cursor", ]; bad.into_iter().for_each(|s| { assert!( @@ -607,9 +734,80 @@ mod tests { } #[test] - fn cursor_token_is_tid_shape() { - let token = PageToken::new(1_730_000_000_000_000).encode_token(); - assert_eq!(token.len(), 13); - assert!(token.chars().all(|c| c.is_ascii_alphanumeric())); + fn cursor_token_is_hex_shape() { + let token = PageToken::new(1_730_000_000_000_000, 7).encode_token(); + assert_eq!(token.len(), 24); + assert!(token.chars().all(|c| c.is_ascii_hexdigit())); + } + + #[test] + fn list_does_not_drop_entries_at_same_sort_micros_across_pages() { + let store = store(); + let subject = "did:plc:limpet"; + let key = EdgeKey::new(nsid("sh.tangled.feed.star"), did_subj(subject)); + (0..4).for_each(|i| { + store.add(star_edge_at( + &format!("at://did:plc:{}/sh.tangled.feed.star/r{i}", NAMES[i]), + subject, + 42, + )); + }); + + let page1 = store.list(&key, PageCursor::Start, limit(2)); + assert_eq!(page1.items.len(), 2); + let token = page1.next.expect("cursor must continue across ties"); + + let page2 = store.list(&key, PageCursor::After(token), limit(2)); + assert_eq!( + page2.items.len(), + 2, + "remaining ties must surface on next page" + ); + assert!(page2.next.is_none()); + let combined: std::collections::HashSet = page1 + .items + .iter() + .chain(page2.items.iter()) + .map(|u| u.as_ref().to_owned()) + .collect(); + assert_eq!(combined.len(), 4, "every tied entry visible exactly once"); + } + + #[test] + fn list_filtered_narrows_by_predicate_and_paginates() { + let store = store(); + let subject = "did:plc:limpet"; + let key = EdgeKey::new(nsid("sh.tangled.repo.issue"), did_subj(subject)); + let by_nel = (0..3).map(|i| { + star_edge_at( + &format!("at://did:plc:nel/sh.tangled.repo.issue/r{i}"), + subject, + 100 + i as u64, + ) + }); + let by_olaren = (0..2).map(|i| { + star_edge_at( + &format!("at://did:plc:olaren/sh.tangled.repo.issue/o{i}"), + subject, + 500 + i as u64, + ) + }); + by_nel + .chain(by_olaren) + .map(|mut e| { + e.kind = nsid("sh.tangled.repo.issue"); + e + }) + .for_each(|e| store.add(e)); + + let only_nel = |u: &AtUri| u.as_ref().starts_with("at://did:plc:nel/"); + let page1 = + store.list_filtered(&key, PageCursor::Start, limit(2), only_nel); + assert_eq!(page1.items.len(), 2); + assert!(page1.next.is_some(), "cursor must allow more nel matches"); + + let page2 = store.list_filtered(&key, PageCursor::After(page1.next.unwrap()), limit(2), only_nel); + assert_eq!(page2.items.len(), 1, "only one nel issue left"); + assert!(page2.next.is_none(), "tail page must not promise more"); } } diff --git a/crates/edge-index/src/state_index.rs b/crates/edge-index/src/state_index.rs new file mode 100644 index 0000000..69aba7b --- /dev/null +++ b/crates/edge-index/src/state_index.rs @@ -0,0 +1,450 @@ +use std::collections::BTreeMap; +use std::sync::Mutex; + +use bobbin_runtime::RuntimeHasher; +use bobbin_types::edges::Record; +use bobbin_types::sh_tangled::repo::issue::state::StateState; +use bobbin_types::sh_tangled::repo::pull::status::StatusStatus; +use jacquard_common::DefaultStr; +use jacquard_common::types::string::AtUri; +use scc::HashMap as SccMap; +use scc::hash_map::Entry; + +pub trait StateKind: Copy + Eq + std::fmt::Debug + Send + Sync + 'static { + fn wire(self) -> &'static str; +} + +#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] +pub enum IssueStateKind { + Open, + Closed, +} + +impl StateKind for IssueStateKind { + fn wire(self) -> &'static str { + match self { + Self::Open => "open", + Self::Closed => "closed", + } + } +} + +#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] +pub enum PullStatusKind { + Open, + Closed, + Merged, +} + +impl StateKind for PullStatusKind { + fn wire(self) -> &'static str { + match self { + Self::Open => "open", + Self::Closed => "closed", + Self::Merged => "merged", + } + } +} + +#[derive(Clone, Debug)] +struct ReverseRef { + entity: AtUri, + sort_micros: u64, + kind: V, +} + +#[derive(Clone, Debug)] +struct SortableSource(AtUri); + +impl PartialEq for SortableSource { + fn eq(&self, other: &Self) -> bool { + self.0.as_ref() == other.0.as_ref() + } +} + +impl Eq for SortableSource {} + +impl PartialOrd for SortableSource { + fn partial_cmp(&self, other: &Self) -> Option { + Some(self.cmp(other)) + } +} + +impl Ord for SortableSource { + fn cmp(&self, other: &Self) -> std::cmp::Ordering { + self.0.as_ref().cmp(other.0.as_ref()) + } +} + +type ForwardKey = (u64, SortableSource); + +pub struct StateIndex { + forward: SccMap, BTreeMap, RuntimeHasher>, + reverse: SccMap, ReverseRef, RuntimeHasher>, + writer: Mutex<()>, +} + +impl StateIndex { + pub fn new(hasher: RuntimeHasher) -> Self { + Self { + forward: SccMap::with_hasher(hasher.clone()), + reverse: SccMap::with_hasher(hasher), + writer: Mutex::new(()), + } + } + + pub fn upsert( + &self, + source: AtUri, + entity: AtUri, + sort_micros: u64, + kind: V, + ) { + let _w = self + .writer + .lock() + .expect("state-index writer mutex poisoned"); + let rev = self.reverse.entry_sync(source.clone()); + let prior = match &rev { + Entry::Occupied(occ) => Some((occ.get().entity.clone(), occ.get().sort_micros)), + Entry::Vacant(_) => None, + }; + if let Some((prior_entity, prior_micros)) = prior.as_ref() + && (prior_entity != &entity || *prior_micros != sort_micros) + { + self.remove_forward(prior_entity, *prior_micros, &source); + } + self.insert_forward(entity.clone(), sort_micros, kind, source); + match rev { + Entry::Occupied(mut occ) => { + let slot = occ.get_mut(); + slot.entity = entity; + slot.sort_micros = sort_micros; + slot.kind = kind; + } + Entry::Vacant(vac) => { + vac.insert_entry(ReverseRef { + entity, + sort_micros, + kind, + }); + } + } + } + + pub fn remove_source(&self, source: &AtUri) { + let _w = self + .writer + .lock() + .expect("state-index writer mutex poisoned"); + let Entry::Occupied(occ) = self.reverse.entry_sync(source.clone()) else { + return; + }; + let entity = occ.get().entity.clone(); + let sort_micros = occ.get().sort_micros; + let _ = occ.remove(); + self.remove_forward(&entity, sort_micros, source); + } + + pub fn remove_entity(&self, entity: &AtUri) { + let _w = self + .writer + .lock() + .expect("state-index writer mutex poisoned"); + let Some((_, set)) = self.forward.remove_sync(entity) else { + return; + }; + set.into_iter().for_each(|((_, src), _)| { + let Entry::Occupied(occ) = self.reverse.entry_sync(src.0.clone()) else { + return; + }; + if occ.get().entity == *entity { + let _ = occ.remove(); + } + }); + } + + fn insert_forward( + &self, + entity: AtUri, + sort_micros: u64, + kind: V, + source: AtUri, + ) { + let mut entry = self.forward.entry_sync(entity).or_default(); + entry + .get_mut() + .insert((sort_micros, SortableSource(source)), kind); + } + + fn remove_forward( + &self, + entity: &AtUri, + sort_micros: u64, + source: &AtUri, + ) { + let Entry::Occupied(mut entry) = self.forward.entry_sync(entity.clone()) else { + return; + }; + entry + .get_mut() + .remove(&(sort_micros, SortableSource(source.clone()))); + if entry.get().is_empty() { + let _ = entry.remove(); + } + } + + pub fn latest(&self, entity: &AtUri) -> Option<(V, u64)> { + self.latest_by(entity, |_| true) + } + + pub fn latest_by(&self, entity: &AtUri, accept: F) -> Option<(V, u64)> + where + F: Fn(&AtUri) -> bool, + { + self.forward + .read_sync(entity, |_, set| { + set.iter() + .rev() + .find_map(|((m, src), k)| accept(&src.0).then_some((*k, *m))) + }) + .flatten() + } + + pub fn entity_count(&self) -> usize { + self.forward.len() + } + + pub fn source_count(&self) -> usize { + self.reverse.len() + } +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum ApplyOutcome { + Applied, + UnknownVariant, + NotStateRecord, + Removed, +} + +fn issue_kind_from(s: &StateState) -> Option { + match s { + StateState::ShTangledRepoIssueStateOpen => Some(IssueStateKind::Open), + StateState::ShTangledRepoIssueStateClosed => Some(IssueStateKind::Closed), + StateState::Other(_) => None, + } +} + +fn pull_kind_from(s: &StatusStatus) -> Option { + match s { + StatusStatus::ShTangledRepoPullStatusOpen => Some(PullStatusKind::Open), + StatusStatus::ShTangledRepoPullStatusClosed => Some(PullStatusKind::Closed), + StatusStatus::ShTangledRepoPullStatusMerged => Some(PullStatusKind::Merged), + StatusStatus::Other(_) => None, + } +} + +pub fn apply_record_state( + issue_idx: &StateIndex, + pull_idx: &StateIndex, + source: &AtUri, + record: &Record, +) -> ApplyOutcome { + let sort_micros = record.sort_micros_for(source); + match record { + Record::IssueState(r) => match issue_kind_from(&r.state) { + Some(kind) => { + issue_idx.upsert(source.clone(), r.issue.clone(), sort_micros, kind); + ApplyOutcome::Applied + } + None => ApplyOutcome::UnknownVariant, + }, + Record::PullStatus(r) => match pull_kind_from(&r.status) { + Some(kind) => { + pull_idx.upsert(source.clone(), r.pull.clone(), sort_micros, kind); + ApplyOutcome::Applied + } + None => ApplyOutcome::UnknownVariant, + }, + _ => ApplyOutcome::NotStateRecord, + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn idx() -> StateIndex { + StateIndex::new(RuntimeHasher::default()) + } + + fn at(s: &str) -> AtUri { + AtUri::new_owned(s).unwrap() + } + + #[test] + fn latest_returns_largest_sort_micros() { + let s = idx(); + let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); + s.upsert( + at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"), + issue.clone(), + 100, + IssueStateKind::Open, + ); + s.upsert( + at("at://did:plc:nel/sh.tangled.repo.issue.state/s2"), + issue.clone(), + 200, + IssueStateKind::Closed, + ); + assert_eq!(s.latest(&issue), Some((IssueStateKind::Closed, 200))); + } + + #[test] + fn remove_source_drops_entry_and_recovers_prior() { + let s = idx(); + let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); + let s1 = at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"); + let s2 = at("at://did:plc:nel/sh.tangled.repo.issue.state/s2"); + s.upsert(s1.clone(), issue.clone(), 100, IssueStateKind::Open); + s.upsert(s2.clone(), issue.clone(), 200, IssueStateKind::Closed); + s.remove_source(&s2); + assert_eq!(s.latest(&issue), Some((IssueStateKind::Open, 100))); + s.remove_source(&s1); + assert!(s.latest(&issue).is_none()); + assert_eq!(s.entity_count(), 0); + assert_eq!(s.source_count(), 0); + } + + #[test] + fn upsert_same_source_replaces() { + let s = idx(); + let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); + let src = at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"); + s.upsert(src.clone(), issue.clone(), 100, IssueStateKind::Open); + s.upsert(src.clone(), issue.clone(), 100, IssueStateKind::Closed); + assert_eq!(s.latest(&issue), Some((IssueStateKind::Closed, 100))); + assert_eq!(s.source_count(), 1); + } + + #[test] + fn distinct_sources_with_identical_micros_and_kind_survive_individual_removal() { + let s = idx(); + let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); + let s1 = at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"); + let s2 = at("at://did:plc:olaren/sh.tangled.repo.issue.state/s2"); + s.upsert(s1.clone(), issue.clone(), 100, IssueStateKind::Open); + s.upsert(s2.clone(), issue.clone(), 100, IssueStateKind::Open); + s.remove_source(&s1); + assert_eq!( + s.latest(&issue), + Some((IssueStateKind::Open, 100)), + "removing one source must not wipe the other's matching state" + ); + s.remove_source(&s2); + assert!(s.latest(&issue).is_none()); + } + + #[test] + fn unknown_variant_is_reported() { + use bobbin_types::sh_tangled::repo::issue::state::State as IssueStateRec; + use jacquard_common::deps::smol_str::SmolStr; + let issue_idx = idx(); + let pull_idx = StateIndex::::new(RuntimeHasher::default()); + let rec = Record::IssueState(IssueStateRec { + issue: at("at://did:plc:limpet/sh.tangled.repo.issue/i1"), + state: StateState::Other(SmolStr::new_static( + "sh.tangled.repo.issue.state.reopened", + )), + extra_data: None, + }); + let outcome = apply_record_state( + &issue_idx, + &pull_idx, + &at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"), + &rec, + ); + assert_eq!(outcome, ApplyOutcome::UnknownVariant); + assert_eq!(issue_idx.entity_count(), 0); + } + + #[test] + fn latest_by_filters_unauthorized_sources() { + let s = idx(); + let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); + let owner_state = at("at://did:plc:limpet/sh.tangled.repo.issue.state/legit"); + let attacker_state = at("at://did:plc:nautilus/sh.tangled.repo.issue.state/spoof"); + s.upsert(owner_state.clone(), issue.clone(), 100, IssueStateKind::Open); + s.upsert( + attacker_state.clone(), + issue.clone(), + 500, + IssueStateKind::Closed, + ); + let only_owner = |src: &AtUri| { + src.as_ref().starts_with("at://did:plc:limpet/") + }; + assert_eq!( + s.latest_by(&issue, only_owner), + Some((IssueStateKind::Open, 100)), + "spoofed attacker state must be ignored", + ); + assert_eq!( + s.latest(&issue), + Some((IssueStateKind::Closed, 500)), + "unfiltered latest still surfaces the spoof for sanity", + ); + } + + #[test] + fn remove_entity_drops_forward_and_reverse() { + let s = idx(); + let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); + let s1 = at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"); + let s2 = at("at://did:plc:olaren/sh.tangled.repo.issue.state/s2"); + s.upsert(s1.clone(), issue.clone(), 100, IssueStateKind::Open); + s.upsert(s2.clone(), issue.clone(), 200, IssueStateKind::Closed); + s.remove_entity(&issue); + assert!(s.latest(&issue).is_none()); + assert_eq!(s.entity_count(), 0); + assert_eq!( + s.source_count(), + 0, + "reverse entries pointing at the dead entity must clear too", + ); + } + + #[test] + fn remove_entity_does_not_touch_reverse_pointing_elsewhere() { + let s = idx(); + let dead = at("at://did:plc:limpet/sh.tangled.repo.issue/dead"); + let alive = at("at://did:plc:limpet/sh.tangled.repo.issue/alive"); + let src = at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"); + s.upsert(src.clone(), dead.clone(), 100, IssueStateKind::Open); + s.upsert(src.clone(), alive.clone(), 200, IssueStateKind::Closed); + s.remove_entity(&dead); + assert_eq!( + s.latest(&alive), + Some((IssueStateKind::Closed, 200)), + "removing the dead entity must not touch state for the live one", + ); + assert_eq!(s.source_count(), 1); + } + + #[test] + fn same_micros_distinct_kind_picks_by_source_not_kind() { + let s = StateIndex::::new(RuntimeHasher::default()); + let pull = at("at://did:plc:limpet/sh.tangled.repo.pull/p1"); + let early = at("at://did:plc:nel/sh.tangled.repo.pull.status/aaa"); + let later = at("at://did:plc:nel/sh.tangled.repo.pull.status/zzz"); + s.upsert(early.clone(), pull.clone(), 1000, PullStatusKind::Merged); + s.upsert(later.clone(), pull.clone(), 1000, PullStatusKind::Open); + assert_eq!( + s.latest(&pull), + Some((PullStatusKind::Open, 1000)), + "tied micros must use source ordering as the tiebreak", + ); + } +} -- 2.51.2