This repository has no description
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}