Monorepo for Tangled
0

Configure Feed

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

tangled-core / bobbin / crates / edge-index / src / state_index.rs
15 kB 453 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::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>) { 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; 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 } 150 151 pub fn remove_entity(&self, entity: &AtUri<DefaultStr>) { 152 let _w = self 153 .writer 154 .lock() 155 .expect("state-index writer mutex poisoned"); 156 let Some((_, set)) = self.forward.remove_sync(entity) else { 157 return; 158 }; 159 set.into_iter().for_each(|((_, src), _)| { 160 let Entry::Occupied(occ) = self.reverse.entry_sync(src.0.clone()) else { 161 return; 162 }; 163 if occ.get().entity == *entity { 164 let _ = occ.remove(); 165 } 166 }); 167 } 168 169 fn insert_forward( 170 &self, 171 entity: AtUri<DefaultStr>, 172 sort_micros: u64, 173 kind: V, 174 source: AtUri<DefaultStr>, 175 ) { 176 let mut entry = self.forward.entry_sync(entity).or_default(); 177 entry 178 .get_mut() 179 .insert((sort_micros, SortableSource(source)), kind); 180 } 181 182 fn remove_forward( 183 &self, 184 entity: &AtUri<DefaultStr>, 185 sort_micros: u64, 186 source: &AtUri<DefaultStr>, 187 ) { 188 let Entry::Occupied(mut entry) = self.forward.entry_sync(entity.clone()) else { 189 return; 190 }; 191 entry 192 .get_mut() 193 .remove(&(sort_micros, SortableSource(source.clone()))); 194 if entry.get().is_empty() { 195 let _ = entry.remove(); 196 } 197 } 198 199 pub fn latest(&self, entity: &AtUri<DefaultStr>) -> Option<(V, u64)> { 200 self.latest_by(entity, |_| true) 201 } 202 203 pub fn latest_by<F>(&self, entity: &AtUri<DefaultStr>, accept: F) -> Option<(V, u64)> 204 where 205 F: Fn(&AtUri<DefaultStr>) -> bool, 206 { 207 self.forward 208 .read_sync(entity, |_, set| { 209 set.iter() 210 .rev() 211 .find_map(|((m, src), k)| accept(&src.0).then_some((*k, *m))) 212 }) 213 .flatten() 214 } 215 216 pub fn entity_count(&self) -> usize { 217 self.forward.len() 218 } 219 220 pub fn source_count(&self) -> usize { 221 self.reverse.len() 222 } 223} 224 225#[derive(Clone, Copy, Debug, Eq, PartialEq)] 226pub enum ApplyOutcome { 227 Applied, 228 UnknownVariant, 229 NotStateRecord, 230 Removed, 231} 232 233fn issue_kind_from(s: &StateState<DefaultStr>) -> Option<IssueStateKind> { 234 match s { 235 StateState::ShTangledRepoIssueStateOpen => Some(IssueStateKind::Open), 236 StateState::ShTangledRepoIssueStateClosed => Some(IssueStateKind::Closed), 237 StateState::Other(_) => None, 238 } 239} 240 241fn pull_kind_from(s: &StatusStatus<DefaultStr>) -> Option<PullStatusKind> { 242 match s { 243 StatusStatus::ShTangledRepoPullStatusOpen => Some(PullStatusKind::Open), 244 StatusStatus::ShTangledRepoPullStatusClosed => Some(PullStatusKind::Closed), 245 StatusStatus::ShTangledRepoPullStatusMerged => Some(PullStatusKind::Merged), 246 StatusStatus::Other(_) => None, 247 } 248} 249 250pub fn apply_record_state( 251 issue_idx: &StateIndex<IssueStateKind>, 252 pull_idx: &StateIndex<PullStatusKind>, 253 source: &AtUri<DefaultStr>, 254 record: &Record, 255) -> ApplyOutcome { 256 let sort_micros = record.sort_micros_for(source); 257 match record { 258 Record::IssueState(r) => match issue_kind_from(&r.state) { 259 Some(kind) => { 260 issue_idx.upsert(source.clone(), r.issue.clone(), sort_micros, kind); 261 ApplyOutcome::Applied 262 } 263 None => ApplyOutcome::UnknownVariant, 264 }, 265 Record::PullStatus(r) => match pull_kind_from(&r.status) { 266 Some(kind) => { 267 pull_idx.upsert(source.clone(), r.pull.clone(), sort_micros, kind); 268 ApplyOutcome::Applied 269 } 270 None => ApplyOutcome::UnknownVariant, 271 }, 272 _ => ApplyOutcome::NotStateRecord, 273 } 274} 275 276#[cfg(test)] 277mod tests { 278 use super::*; 279 280 fn idx() -> StateIndex<IssueStateKind> { 281 StateIndex::new(RuntimeHasher::default()) 282 } 283 284 fn at(s: &str) -> AtUri<DefaultStr> { 285 AtUri::new_owned(s).unwrap() 286 } 287 288 #[test] 289 fn latest_returns_largest_sort_micros() { 290 let s = idx(); 291 let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); 292 s.upsert( 293 at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"), 294 issue.clone(), 295 100, 296 IssueStateKind::Open, 297 ); 298 s.upsert( 299 at("at://did:plc:nel/sh.tangled.repo.issue.state/s2"), 300 issue.clone(), 301 200, 302 IssueStateKind::Closed, 303 ); 304 assert_eq!(s.latest(&issue), Some((IssueStateKind::Closed, 200))); 305 } 306 307 #[test] 308 fn remove_source_drops_entry_and_recovers_prior() { 309 let s = idx(); 310 let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); 311 let s1 = at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"); 312 let s2 = at("at://did:plc:nel/sh.tangled.repo.issue.state/s2"); 313 s.upsert(s1.clone(), issue.clone(), 100, IssueStateKind::Open); 314 s.upsert(s2.clone(), issue.clone(), 200, IssueStateKind::Closed); 315 s.remove_source(&s2); 316 assert_eq!(s.latest(&issue), Some((IssueStateKind::Open, 100))); 317 s.remove_source(&s1); 318 assert!(s.latest(&issue).is_none()); 319 assert_eq!(s.entity_count(), 0); 320 assert_eq!(s.source_count(), 0); 321 } 322 323 #[test] 324 fn upsert_same_source_replaces() { 325 let s = idx(); 326 let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); 327 let src = at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"); 328 s.upsert(src.clone(), issue.clone(), 100, IssueStateKind::Open); 329 s.upsert(src.clone(), issue.clone(), 100, IssueStateKind::Closed); 330 assert_eq!(s.latest(&issue), Some((IssueStateKind::Closed, 100))); 331 assert_eq!(s.source_count(), 1); 332 } 333 334 #[test] 335 fn distinct_sources_with_identical_micros_and_kind_survive_individual_removal() { 336 let s = idx(); 337 let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); 338 let s1 = at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"); 339 let s2 = at("at://did:plc:olaren/sh.tangled.repo.issue.state/s2"); 340 s.upsert(s1.clone(), issue.clone(), 100, IssueStateKind::Open); 341 s.upsert(s2.clone(), issue.clone(), 100, IssueStateKind::Open); 342 s.remove_source(&s1); 343 assert_eq!( 344 s.latest(&issue), 345 Some((IssueStateKind::Open, 100)), 346 "removing one source must not wipe the other's matching state" 347 ); 348 s.remove_source(&s2); 349 assert!(s.latest(&issue).is_none()); 350 } 351 352 #[test] 353 fn unknown_variant_is_reported() { 354 use bobbin_types::sh_tangled::repo::issue::state::State as IssueStateRec; 355 use jacquard_common::deps::smol_str::SmolStr; 356 let issue_idx = idx(); 357 let pull_idx = StateIndex::<PullStatusKind>::new(RuntimeHasher::default()); 358 let rec = Record::IssueState(IssueStateRec { 359 issue: at("at://did:plc:limpet/sh.tangled.repo.issue/i1"), 360 state: StateState::Other(SmolStr::new_static("sh.tangled.repo.issue.state.reopened")), 361 extra_data: None, 362 }); 363 let outcome = apply_record_state( 364 &issue_idx, 365 &pull_idx, 366 &at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"), 367 &rec, 368 ); 369 assert_eq!(outcome, ApplyOutcome::UnknownVariant); 370 assert_eq!(issue_idx.entity_count(), 0); 371 } 372 373 #[test] 374 fn latest_by_filters_unauthorized_sources() { 375 let s = idx(); 376 let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); 377 let owner_state = at("at://did:plc:limpet/sh.tangled.repo.issue.state/legit"); 378 let attacker_state = at("at://did:plc:nautilus/sh.tangled.repo.issue.state/spoof"); 379 s.upsert( 380 owner_state.clone(), 381 issue.clone(), 382 100, 383 IssueStateKind::Open, 384 ); 385 s.upsert( 386 attacker_state.clone(), 387 issue.clone(), 388 500, 389 IssueStateKind::Closed, 390 ); 391 let only_owner = |src: &AtUri<DefaultStr>| src.as_ref().starts_with("at://did:plc:limpet/"); 392 assert_eq!( 393 s.latest_by(&issue, only_owner), 394 Some((IssueStateKind::Open, 100)), 395 "spoofed attacker state must be ignored", 396 ); 397 assert_eq!( 398 s.latest(&issue), 399 Some((IssueStateKind::Closed, 500)), 400 "unfiltered latest still surfaces the spoof for sanity", 401 ); 402 } 403 404 #[test] 405 fn remove_entity_drops_forward_and_reverse() { 406 let s = idx(); 407 let issue = at("at://did:plc:limpet/sh.tangled.repo.issue/i1"); 408 let s1 = at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"); 409 let s2 = at("at://did:plc:olaren/sh.tangled.repo.issue.state/s2"); 410 s.upsert(s1.clone(), issue.clone(), 100, IssueStateKind::Open); 411 s.upsert(s2.clone(), issue.clone(), 200, IssueStateKind::Closed); 412 s.remove_entity(&issue); 413 assert!(s.latest(&issue).is_none()); 414 assert_eq!(s.entity_count(), 0); 415 assert_eq!( 416 s.source_count(), 417 0, 418 "reverse entries pointing at the dead entity must clear too", 419 ); 420 } 421 422 #[test] 423 fn remove_entity_does_not_touch_reverse_pointing_elsewhere() { 424 let s = idx(); 425 let dead = at("at://did:plc:limpet/sh.tangled.repo.issue/dead"); 426 let alive = at("at://did:plc:limpet/sh.tangled.repo.issue/alive"); 427 let src = at("at://did:plc:nel/sh.tangled.repo.issue.state/s1"); 428 s.upsert(src.clone(), dead.clone(), 100, IssueStateKind::Open); 429 s.upsert(src.clone(), alive.clone(), 200, IssueStateKind::Closed); 430 s.remove_entity(&dead); 431 assert_eq!( 432 s.latest(&alive), 433 Some((IssueStateKind::Closed, 200)), 434 "removing the dead entity must not touch state for the live one", 435 ); 436 assert_eq!(s.source_count(), 1); 437 } 438 439 #[test] 440 fn same_micros_distinct_kind_picks_by_source_not_kind() { 441 let s = StateIndex::<PullStatusKind>::new(RuntimeHasher::default()); 442 let pull = at("at://did:plc:limpet/sh.tangled.repo.pull/p1"); 443 let early = at("at://did:plc:nel/sh.tangled.repo.pull.status/aaa"); 444 let later = at("at://did:plc:nel/sh.tangled.repo.pull.status/zzz"); 445 s.upsert(early.clone(), pull.clone(), 1000, PullStatusKind::Merged); 446 s.upsert(later.clone(), pull.clone(), 1000, PullStatusKind::Open); 447 assert_eq!( 448 s.latest(&pull), 449 Some((PullStatusKind::Open, 1000)), 450 "tied micros must use source ordering as the tiebreak", 451 ); 452 } 453}