This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / bobbin / crates / edge-index / src / state_index.rs
15 kB 461 lines
1use std::collections::BTreeMap; 2use std::sync::Mutex; 3 4use bobbin_runtime::RuntimeHasher; 5use bobbin_types::edges::Record; 6use bobbin_types::sh_tangled::repo::issue::state::StateState; 7use bobbin_types::sh_tangled::repo::pull::status::StatusStatus; 8use jacquard_common::DefaultStr; 9use jacquard_common::types::string::AtUri; 10use scc::HashMap as SccMap; 11use scc::hash_map::Entry; 12 13pub trait StateKind: Copy + Eq + std::hash::Hash + std::fmt::Debug + Send + Sync + 'static { 14 fn wire(self) -> &'static str; 15} 16 17#[derive(Clone, Copy, Debug, Default, Eq, Hash, PartialEq)] 18pub enum IssueStateKind { 19 #[default] 20 Open, 21 Closed, 22} 23 24impl StateKind for IssueStateKind { 25 fn wire(self) -> &'static str { 26 match self { 27 Self::Open => "open", 28 Self::Closed => "closed", 29 } 30 } 31} 32 33#[derive(Clone, Copy, Debug, Default, Eq, Hash, PartialEq)] 34pub enum PullStatusKind { 35 #[default] 36 Open, 37 Closed, 38 Merged, 39} 40 41impl StateKind for PullStatusKind { 42 fn wire(self) -> &'static str { 43 match self { 44 Self::Open => "open", 45 Self::Closed => "closed", 46 Self::Merged => "merged", 47 } 48 } 49} 50 51#[derive(Clone, Debug)] 52struct ReverseRef<V: StateKind> { 53 entity: AtUri<DefaultStr>, 54 sort_micros: u64, 55 kind: V, 56} 57 58#[derive(Clone, Debug)] 59struct SortableSource(AtUri<DefaultStr>); 60 61impl PartialEq for SortableSource { 62 fn eq(&self, other: &Self) -> bool { 63 self.0.as_ref() == other.0.as_ref() 64 } 65} 66 67impl Eq for SortableSource {} 68 69impl PartialOrd for SortableSource { 70 fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> { 71 Some(self.cmp(other)) 72 } 73} 74 75impl Ord for SortableSource { 76 fn cmp(&self, other: &Self) -> std::cmp::Ordering { 77 self.0.as_ref().cmp(other.0.as_ref()) 78 } 79} 80 81type ForwardKey = (u64, SortableSource); 82 83pub struct StateIndex<V: StateKind> { 84 forward: SccMap<AtUri<DefaultStr>, BTreeMap<ForwardKey, V>, RuntimeHasher>, 85 reverse: SccMap<AtUri<DefaultStr>, ReverseRef<V>, RuntimeHasher>, 86 writer: Mutex<()>, 87} 88 89impl<V: StateKind> StateIndex<V> { 90 pub fn new(hasher: RuntimeHasher) -> Self { 91 Self { 92 forward: SccMap::with_hasher(hasher.clone()), 93 reverse: SccMap::with_hasher(hasher), 94 writer: Mutex::new(()), 95 } 96 } 97 98 pub fn upsert( 99 &self, 100 source: AtUri<DefaultStr>, 101 entity: AtUri<DefaultStr>, 102 sort_micros: u64, 103 kind: V, 104 ) { 105 let _w = self 106 .writer 107 .lock() 108 .expect("state-index writer mutex poisoned"); 109 let rev = self.reverse.entry_sync(source.clone()); 110 let prior = match &rev { 111 Entry::Occupied(occ) => Some((occ.get().entity.clone(), occ.get().sort_micros)), 112 Entry::Vacant(_) => None, 113 }; 114 if let Some((prior_entity, prior_micros)) = prior.as_ref() 115 && (prior_entity != &entity || *prior_micros != sort_micros) 116 { 117 self.remove_forward(prior_entity, *prior_micros, &source); 118 } 119 self.insert_forward(entity.clone(), sort_micros, kind, source); 120 match rev { 121 Entry::Occupied(mut occ) => { 122 let slot = occ.get_mut(); 123 slot.entity = entity; 124 slot.sort_micros = sort_micros; 125 slot.kind = kind; 126 } 127 Entry::Vacant(vac) => { 128 vac.insert_entry(ReverseRef { 129 entity, 130 sort_micros, 131 kind, 132 }); 133 } 134 } 135 } 136 137 pub fn remove_source(&self, source: &AtUri<DefaultStr>) -> Option<AtUri<DefaultStr>> { 138 let _w = self 139 .writer 140 .lock() 141 .expect("state-index writer mutex poisoned"); 142 let Entry::Occupied(occ) = self.reverse.entry_sync(source.clone()) else { 143 return None; 144 }; 145 let entity = occ.get().entity.clone(); 146 let sort_micros = occ.get().sort_micros; 147 let _ = occ.remove(); 148 self.remove_forward(&entity, sort_micros, source); 149 Some(entity) 150 } 151 152 pub fn remove_entity(&self, entity: &AtUri<DefaultStr>) { 153 let _w = self 154 .writer 155 .lock() 156 .expect("state-index writer mutex poisoned"); 157 let Some((_, set)) = self.forward.remove_sync(entity) else { 158 return; 159 }; 160 set.into_iter().for_each(|((_, src), _)| { 161 let Entry::Occupied(occ) = self.reverse.entry_sync(src.0.clone()) else { 162 return; 163 }; 164 if occ.get().entity == *entity { 165 let _ = occ.remove(); 166 } 167 }); 168 } 169 170 fn insert_forward( 171 &self, 172 entity: AtUri<DefaultStr>, 173 sort_micros: u64, 174 kind: V, 175 source: AtUri<DefaultStr>, 176 ) { 177 let mut entry = self.forward.entry_sync(entity).or_default(); 178 entry 179 .get_mut() 180 .insert((sort_micros, SortableSource(source)), kind); 181 } 182 183 fn remove_forward( 184 &self, 185 entity: &AtUri<DefaultStr>, 186 sort_micros: u64, 187 source: &AtUri<DefaultStr>, 188 ) { 189 let Entry::Occupied(mut entry) = self.forward.entry_sync(entity.clone()) else { 190 return; 191 }; 192 entry 193 .get_mut() 194 .remove(&(sort_micros, SortableSource(source.clone()))); 195 if entry.get().is_empty() { 196 let _ = entry.remove(); 197 } 198 } 199 200 pub fn latest(&self, entity: &AtUri<DefaultStr>) -> Option<(V, u64)> { 201 self.latest_by(entity, |_| true) 202 } 203 204 pub fn latest_by<F>(&self, entity: &AtUri<DefaultStr>, accept: F) -> Option<(V, u64)> 205 where 206 F: Fn(&AtUri<DefaultStr>) -> bool, 207 { 208 self.forward 209 .read_sync(entity, |_, set| { 210 set.iter() 211 .rev() 212 .find_map(|((m, src), k)| accept(&src.0).then_some((*k, *m))) 213 }) 214 .flatten() 215 } 216 217 pub fn entity_count(&self) -> usize { 218 self.forward.len() 219 } 220 221 pub fn source_count(&self) -> usize { 222 self.reverse.len() 223 } 224 225 pub fn entity_for_source(&self, source: &AtUri<DefaultStr>) -> Option<AtUri<DefaultStr>> { 226 self.reverse 227 .read_sync(source, |_, entry| entry.entity.clone()) 228 } 229} 230 231#[derive(Clone, Copy, Debug, Eq, PartialEq)] 232pub enum ApplyOutcome { 233 Applied, 234 UnknownVariant, 235 NotStateRecord, 236 Removed, 237} 238 239fn issue_kind_from(s: &StateState<DefaultStr>) -> Option<IssueStateKind> { 240 match s { 241 StateState::ShTangledRepoIssueStateOpen => Some(IssueStateKind::Open), 242 StateState::ShTangledRepoIssueStateClosed => Some(IssueStateKind::Closed), 243 StateState::Other(_) => None, 244 } 245} 246 247fn pull_kind_from(s: &StatusStatus<DefaultStr>) -> Option<PullStatusKind> { 248 match s { 249 StatusStatus::ShTangledRepoPullStatusOpen => Some(PullStatusKind::Open), 250 StatusStatus::ShTangledRepoPullStatusClosed => Some(PullStatusKind::Closed), 251 StatusStatus::ShTangledRepoPullStatusMerged => Some(PullStatusKind::Merged), 252 StatusStatus::Other(_) => None, 253 } 254} 255 256pub fn apply_record_state( 257 issue_idx: &StateIndex<IssueStateKind>, 258 pull_idx: &StateIndex<PullStatusKind>, 259 source: &AtUri<DefaultStr>, 260 record: &Record, 261) -> ApplyOutcome { 262 let sort_micros = record.sort_micros_for(source); 263 match record { 264 Record::IssueState(r) => match issue_kind_from(&r.state) { 265 Some(kind) => { 266 issue_idx.upsert(source.clone(), r.issue.clone(), sort_micros, kind); 267 ApplyOutcome::Applied 268 } 269 None => ApplyOutcome::UnknownVariant, 270 }, 271 Record::PullStatus(r) => match pull_kind_from(&r.status) { 272 Some(kind) => { 273 pull_idx.upsert(source.clone(), r.pull.clone(), sort_micros, kind); 274 ApplyOutcome::Applied 275 } 276 None => ApplyOutcome::UnknownVariant, 277 }, 278 _ => ApplyOutcome::NotStateRecord, 279 } 280} 281 282#[cfg(test)] 283mod tests { 284 use super::*; 285 286 fn idx() -> StateIndex<IssueStateKind> { 287 StateIndex::new(RuntimeHasher::default()) 288 } 289 290 fn at(s: &str) -> AtUri<DefaultStr> { 291 AtUri::new_owned(s).unwrap() 292 } 293 294 #[test] 295 fn latest_returns_largest_sort_micros() { 296 let s = idx(); 297 let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); 298 s.upsert( 299 at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"), 300 issue.clone(), 301 100, 302 IssueStateKind::Open, 303 ); 304 s.upsert( 305 at("at://did:plc:nel/sh.tangled.repo.issue.state/s2"), 306 issue.clone(), 307 200, 308 IssueStateKind::Closed, 309 ); 310 assert_eq!(s.latest(&issue), Some((IssueStateKind::Closed, 200))); 311 } 312 313 #[test] 314 fn remove_source_drops_entry_and_recovers_prior() { 315 let s = idx(); 316 let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); 317 let s1 = at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"); 318 let s2 = at("at://did:plc:nel/sh.tangled.repo.issue.state/s2"); 319 s.upsert(s1.clone(), issue.clone(), 100, IssueStateKind::Open); 320 s.upsert(s2.clone(), issue.clone(), 200, IssueStateKind::Closed); 321 s.remove_source(&s2); 322 assert_eq!(s.latest(&issue), Some((IssueStateKind::Open, 100))); 323 s.remove_source(&s1); 324 assert!(s.latest(&issue).is_none()); 325 assert_eq!(s.entity_count(), 0); 326 assert_eq!(s.source_count(), 0); 327 } 328 329 #[test] 330 fn upsert_same_source_replaces() { 331 let s = idx(); 332 let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); 333 let src = at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"); 334 s.upsert(src.clone(), issue.clone(), 100, IssueStateKind::Open); 335 s.upsert(src.clone(), issue.clone(), 100, IssueStateKind::Closed); 336 assert_eq!(s.latest(&issue), Some((IssueStateKind::Closed, 100))); 337 assert_eq!(s.source_count(), 1); 338 } 339 340 #[test] 341 fn distinct_sources_with_identical_micros_and_kind_survive_individual_removal() { 342 let s = idx(); 343 let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); 344 let s1 = at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"); 345 let s2 = at("at://did:plc:olaren/sh.tangled.repo.issue.state/s2"); 346 s.upsert(s1.clone(), issue.clone(), 100, IssueStateKind::Open); 347 s.upsert(s2.clone(), issue.clone(), 100, IssueStateKind::Open); 348 s.remove_source(&s1); 349 assert_eq!( 350 s.latest(&issue), 351 Some((IssueStateKind::Open, 100)), 352 "removing one source must not wipe the other's matching state" 353 ); 354 s.remove_source(&s2); 355 assert!(s.latest(&issue).is_none()); 356 } 357 358 #[test] 359 fn unknown_variant_is_reported() { 360 use bobbin_types::sh_tangled::repo::issue::state::State as IssueStateRec; 361 use jacquard_common::deps::smol_str::SmolStr; 362 use jacquard_common::types::string::Datetime; 363 let issue_idx = idx(); 364 let pull_idx = StateIndex::<PullStatusKind>::new(RuntimeHasher::default()); 365 let rec = Record::IssueState(IssueStateRec { 366 created_at: Datetime::raw_str("2026-06-11T00:00:00Z"), 367 issue: at("at://did:plc:limpet/sh.tangled.repo.issue/i1"), 368 state: StateState::Other(SmolStr::new_static("sh.tangled.repo.issue.state.reopened")), 369 extra_data: None, 370 }); 371 let outcome = apply_record_state( 372 &issue_idx, 373 &pull_idx, 374 &at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"), 375 &rec, 376 ); 377 assert_eq!(outcome, ApplyOutcome::UnknownVariant); 378 assert_eq!(issue_idx.entity_count(), 0); 379 } 380 381 #[test] 382 fn latest_by_filters_unauthorized_sources() { 383 let s = idx(); 384 let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); 385 let owner_state = at("at://did:plc:limpet/sh.tangled.repo.issue.state/legit"); 386 let attacker_state = at("at://did:plc:nautilus/sh.tangled.repo.issue.state/spoof"); 387 s.upsert( 388 owner_state.clone(), 389 issue.clone(), 390 100, 391 IssueStateKind::Open, 392 ); 393 s.upsert( 394 attacker_state.clone(), 395 issue.clone(), 396 500, 397 IssueStateKind::Closed, 398 ); 399 let only_owner = |src: &AtUri<DefaultStr>| src.as_ref().starts_with("at://did:plc:limpet/"); 400 assert_eq!( 401 s.latest_by(&issue, only_owner), 402 Some((IssueStateKind::Open, 100)), 403 "spoofed attacker state must be ignored", 404 ); 405 assert_eq!( 406 s.latest(&issue), 407 Some((IssueStateKind::Closed, 500)), 408 "unfiltered latest still surfaces the spoof for sanity", 409 ); 410 } 411 412 #[test] 413 fn remove_entity_drops_forward_and_reverse() { 414 let s = idx(); 415 let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); 416 let s1 = at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"); 417 let s2 = at("at://did:plc:olaren/sh.tangled.repo.issue.state/s2"); 418 s.upsert(s1.clone(), issue.clone(), 100, IssueStateKind::Open); 419 s.upsert(s2.clone(), issue.clone(), 200, IssueStateKind::Closed); 420 s.remove_entity(&issue); 421 assert!(s.latest(&issue).is_none()); 422 assert_eq!(s.entity_count(), 0); 423 assert_eq!( 424 s.source_count(), 425 0, 426 "reverse entries pointing at the dead entity must clear too", 427 ); 428 } 429 430 #[test] 431 fn remove_entity_does_not_touch_reverse_pointing_elsewhere() { 432 let s = idx(); 433 let dead = at("at://did:plc:limpet/sh.tangled.repo.issue/dead"); 434 let alive = at("at://did:plc:limpet/sh.tangled.repo.issue/alive"); 435 let src = at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"); 436 s.upsert(src.clone(), dead.clone(), 100, IssueStateKind::Open); 437 s.upsert(src.clone(), alive.clone(), 200, IssueStateKind::Closed); 438 s.remove_entity(&dead); 439 assert_eq!( 440 s.latest(&alive), 441 Some((IssueStateKind::Closed, 200)), 442 "removing the dead entity must not touch state for the live one", 443 ); 444 assert_eq!(s.source_count(), 1); 445 } 446 447 #[test] 448 fn same_micros_distinct_kind_picks_by_source_not_kind() { 449 let s = StateIndex::<PullStatusKind>::new(RuntimeHasher::default()); 450 let pull = at("at://did:plc:limpet/sh.tangled.repo.pull/p1"); 451 let early = at("at://did:plc:nel/sh.tangled.repo.pull.status/aaa"); 452 let later = at("at://did:plc:nel/sh.tangled.repo.pull.status/zzz"); 453 s.upsert(early.clone(), pull.clone(), 1000, PullStatusKind::Merged); 454 s.upsert(later.clone(), pull.clone(), 1000, PullStatusKind::Open); 455 assert_eq!( 456 s.latest(&pull), 457 Some((PullStatusKind::Open, 1000)), 458 "tied micros must use source ordering as the tiebreak", 459 ); 460 } 461}