1use std::collections::{HashMap, VecDeque};
14use std::sync::Mutex;
15use std::time::{Duration, Instant};
16
17use serde_json::Value;
18use uuid::Uuid;
19
20pub const DEFAULT_RING_CAPACITY: usize = 64;
22
23pub const DEFAULT_RING_TTL: Duration = Duration::from_secs(30 * 60);
25
26pub const DEFAULT_MAX_OUTER_KEYS: usize = 4096;
33
34#[derive(Clone, Debug)]
38pub struct RingEntry {
39 pub id: Uuid,
40 pub name: Option<String>,
41 pub touched_at: Instant,
42}
43
44type RingKey = (String, String);
45
46#[derive(Default)]
47struct RingState {
48 rings: HashMap<RingKey, VecDeque<RingEntry>>,
49}
50
51pub struct ReferenceRing {
57 state: Mutex<RingState>,
58 capacity: usize,
59 ttl: Duration,
60 max_outer_keys: usize,
61}
62
63impl Default for ReferenceRing {
64 fn default() -> Self {
65 Self::new()
66 }
67}
68
69impl ReferenceRing {
70 pub fn new() -> Self {
71 Self::with_bounds(DEFAULT_RING_CAPACITY, DEFAULT_RING_TTL)
72 }
73
74 pub fn with_bounds(capacity: usize, ttl: Duration) -> Self {
75 Self::with_bounds_and_outer_limit(capacity, ttl, DEFAULT_MAX_OUTER_KEYS)
76 }
77
78 pub fn with_bounds_and_outer_limit(
79 capacity: usize,
80 ttl: Duration,
81 max_outer_keys: usize,
82 ) -> Self {
83 Self {
84 state: Mutex::new(RingState::default()),
85 capacity,
86 ttl,
87 max_outer_keys,
88 }
89 }
90
91 fn key(namespace: &str, actor: &str) -> RingKey {
92 (namespace.to_owned(), actor.to_owned())
93 }
94
95 fn lock_state(&self) -> std::sync::MutexGuard<'_, RingState> {
99 match self.state.lock() {
100 Ok(guard) => guard,
101 Err(poisoned) => {
102 tracing::warn!(
103 "reference ring mutex poisoned by a prior panic; recovering inner state \
104 (ring admission/lookup is best-effort and must never fail a dispatch)"
105 );
106 poisoned.into_inner()
107 }
108 }
109 }
110
111 fn evict_stale(ring: &mut VecDeque<RingEntry>, ttl: Duration, now: Instant) {
114 while let Some(front) = ring.front() {
115 if now.duration_since(front.touched_at) > ttl {
116 ring.pop_front();
117 } else {
118 break;
119 }
120 }
121 }
122
123 fn prune_outer_map(
128 rings: &mut HashMap<RingKey, VecDeque<RingEntry>>,
129 ttl: Duration,
130 max_outer_keys: usize,
131 now: Instant,
132 exempt: &RingKey,
133 ) {
134 rings.retain(|_, ring| {
135 Self::evict_stale(ring, ttl, now);
136 !ring.is_empty()
137 });
138 if rings.len() <= max_outer_keys {
139 return;
140 }
141 let mut by_recency: Vec<(RingKey, Instant)> = rings
142 .iter()
143 .filter(|(k, _)| *k != exempt)
144 .filter_map(|(k, ring)| ring.back().map(|e| (k.clone(), e.touched_at)))
145 .collect();
146 by_recency.sort_by_key(|(_, touched_at)| *touched_at);
147 let overflow = rings.len().saturating_sub(max_outer_keys);
148 for (key, _) in by_recency.into_iter().take(overflow) {
149 rings.remove(&key);
150 }
151 }
152
153 pub fn admit(&self, namespace: &str, actor: &str, id: Uuid, name: Option<String>) {
160 let now = Instant::now();
161 let mut state = self.lock_state();
162 let key = Self::key(namespace, actor);
163 {
164 let ring = state.rings.entry(key.clone()).or_default();
165 Self::evict_stale(ring, self.ttl, now);
166 ring.retain(|e| e.id != id);
167 ring.push_back(RingEntry {
168 id,
169 name,
170 touched_at: now,
171 });
172 while ring.len() > self.capacity {
173 ring.pop_front();
174 }
175 }
176 Self::prune_outer_map(&mut state.rings, self.ttl, self.max_outer_keys, now, &key);
177 }
178
179 pub fn snapshot(&self, namespace: &str, actor: &str) -> Vec<RingEntry> {
188 let now = Instant::now();
189 let mut state = self.lock_state();
190 let key = Self::key(namespace, actor);
191 let snap: Vec<RingEntry> = match state.rings.get_mut(&key) {
192 Some(ring) => {
193 Self::evict_stale(ring, self.ttl, now);
194 ring.iter().rev().cloned().collect()
195 }
196 None => Vec::new(),
197 };
198 Self::prune_outer_map(&mut state.rings, self.ttl, self.max_outer_keys, now, &key);
199 snap
200 }
201}
202
203fn display_name(result: &Value) -> Option<String> {
207 let name = result.get("name").and_then(Value::as_str)?;
208 let trimmed = name.trim();
209 (!trimmed.is_empty()).then(|| trimmed.to_string())
210}
211
212const ENTITY_KINDS: [&str; 9] = [
215 "concept", "document", "dataset", "project", "person", "org", "artifact", "service", "resource",
216];
217
218fn is_entity_kind_value(v: &str) -> bool {
219 v == "entity" || ENTITY_KINDS.contains(&v)
220}
221
222fn substrate_admits_as_entity(obj: &serde_json::Map<String, Value>) -> bool {
241 if matches!(
242 obj.get("kind").and_then(Value::as_str),
243 Some("edge") | Some("event")
244 ) {
245 return false;
246 }
247 if obj.contains_key("content") {
248 return false;
249 }
250 if obj.contains_key("entity_type") {
251 return true;
252 }
253 obj.get("kind")
254 .and_then(Value::as_str)
255 .is_some_and(is_entity_kind_value)
256}
257
258pub(crate) fn ring_admissions_for(verb: &str, result: &Value) -> Vec<(Uuid, Option<String>)> {
269 let Some(obj) = result.as_object() else {
270 return Vec::new();
271 };
272 if obj.contains_key("attempted") {
273 return Vec::new();
274 }
275 let parse_id = |key: &str| -> Option<Uuid> {
276 obj.get(key)
277 .and_then(Value::as_str)
278 .and_then(|s| Uuid::parse_str(s).ok())
279 };
280 match verb {
281 "create" | "get" | "update" | "delete" => {
282 if !substrate_admits_as_entity(obj) {
283 return Vec::new();
284 }
285 match parse_id("id") {
286 Some(id) => vec![(id, display_name(result))],
287 None => Vec::new(),
288 }
289 }
290 "merge" => match parse_id("kept_id") {
293 Some(id) => vec![(id, None)],
294 None => Vec::new(),
295 },
296 "link" => {
297 let mut out = Vec::new();
298 if let Some(id) = parse_id("source_id") {
299 out.push((id, None));
300 }
301 if let Some(id) = parse_id("target_id") {
302 out.push((id, None));
303 }
304 out
305 }
306 _ => Vec::new(),
307 }
308}
309
310#[cfg(test)]
311mod tests {
312 use super::*;
313 use serde_json::json;
314
315 #[test]
316 fn admits_by_id_ops_and_extracts_name() {
317 let ring = ReferenceRing::new();
318 ring.admit("local", "actor:a", Uuid::nil(), Some("Alpha".into()));
319 let snap = ring.snapshot("local", "actor:a");
320 assert_eq!(snap.len(), 1);
321 assert_eq!(snap[0].id, Uuid::nil());
322 assert_eq!(snap[0].name.as_deref(), Some("Alpha"));
323 }
324
325 #[test]
326 fn ring_admissions_for_search_and_list_is_empty() {
327 let result = json!([{"id": Uuid::nil().to_string(), "name": "hit"}]);
328 assert!(ring_admissions_for("search", &result).is_empty());
329 assert!(ring_admissions_for("list", &result).is_empty());
330 }
331
332 #[test]
333 fn ring_admissions_for_get_extracts_id_and_name() {
334 let id = Uuid::new_v4();
335 let result = json!({"id": id.to_string(), "name": "Concept", "entity_type": null});
339 let admissions = ring_admissions_for("get", &result);
340 assert_eq!(admissions, vec![(id, Some("Concept".to_string()))]);
341 }
342
343 #[test]
344 fn ring_admissions_for_note_result_is_empty() {
345 let id = Uuid::new_v4();
346 let result = json!({"id": id.to_string(), "name": "a note", "content": "body text"});
349 assert!(ring_admissions_for("create", &result).is_empty());
350 assert!(ring_admissions_for("get", &result).is_empty());
351 assert!(ring_admissions_for("update", &result).is_empty());
352 assert!(ring_admissions_for("delete", &result).is_empty());
353 }
354
355 #[test]
356 fn ring_admissions_for_edge_and_event_kind_is_empty() {
357 let id = Uuid::new_v4();
358 let edge_result = json!({"id": id.to_string(), "kind": "edge"});
359 assert!(ring_admissions_for("get", &edge_result).is_empty());
360 let event_result = json!({"id": id.to_string(), "kind": "event"});
361 assert!(ring_admissions_for("get", &event_result).is_empty());
362 }
363
364 #[test]
365 fn ring_admissions_for_delete_uses_resolved_kind() {
366 let id = Uuid::new_v4();
367 let entity_delete = json!({"deleted": true, "id": id.to_string(), "kind": "concept"});
371 assert_eq!(
372 ring_admissions_for("delete", &entity_delete),
373 vec![(id, None)]
374 );
375 let generic_entity_delete =
376 json!({"deleted": true, "id": id.to_string(), "kind": "entity"});
377 assert_eq!(
378 ring_admissions_for("delete", &generic_entity_delete),
379 vec![(id, None)]
380 );
381 let note_delete = json!({"deleted": true, "id": id.to_string(), "kind": "observation"});
382 assert!(ring_admissions_for("delete", ¬e_delete).is_empty());
383 let unspecified_delete = json!({"deleted": true, "id": id.to_string(), "kind": null});
386 assert!(ring_admissions_for("delete", &unspecified_delete).is_empty());
387 }
388
389 #[test]
390 fn display_name_never_falls_back_to_content() {
391 let result = json!({"id": Uuid::new_v4().to_string(), "content": "some note body"});
392 assert_eq!(display_name(&result), None);
393 }
394
395 #[test]
396 fn ring_admissions_for_link_extracts_both_endpoints() {
397 let source = Uuid::new_v4();
398 let target = Uuid::new_v4();
399 let result = json!({
400 "id": Uuid::new_v4().to_string(),
401 "source_id": source.to_string(),
402 "target_id": target.to_string(),
403 });
404 let admissions = ring_admissions_for("link", &result);
405 assert_eq!(admissions, vec![(source, None), (target, None)]);
406 }
407
408 #[test]
409 fn ring_admissions_for_merge_uses_kept_id() {
410 let kept = Uuid::new_v4();
411 let removed = Uuid::new_v4();
412 let result = json!({"kept_id": kept.to_string(), "removed_id": removed.to_string()});
413 let admissions = ring_admissions_for("merge", &result);
414 assert_eq!(admissions, vec![(kept, None)]);
415 }
416
417 #[test]
418 fn ring_admissions_for_bulk_shapes_is_empty() {
419 let bulk_create = json!({"attempted": 3, "created": 3});
420 assert!(ring_admissions_for("create", &bulk_create).is_empty());
421 let bulk_link = json!({"attempted": 2, "created": 2, "skipped": 0, "failed": 0});
422 assert!(ring_admissions_for("link", &bulk_link).is_empty());
423 }
424
425 #[test]
426 fn admission_bounds_by_size() {
427 let ring = ReferenceRing::with_bounds(3, DEFAULT_RING_TTL);
428 let ids: Vec<Uuid> = (0..5).map(|_| Uuid::new_v4()).collect();
429 for id in &ids {
430 ring.admit("local", "actor:a", *id, None);
431 }
432 let snap = ring.snapshot("local", "actor:a");
433 assert_eq!(snap.len(), 3);
434 assert_eq!(snap[0].id, ids[4]);
436 assert_eq!(snap[1].id, ids[3]);
437 assert_eq!(snap[2].id, ids[2]);
438 }
439
440 #[test]
441 fn admission_bounds_by_age() {
442 let ring = ReferenceRing::with_bounds(64, Duration::from_millis(20));
443 let old = Uuid::new_v4();
444 ring.admit("local", "actor:a", old, None);
445 std::thread::sleep(Duration::from_millis(40));
446 let fresh = Uuid::new_v4();
447 ring.admit("local", "actor:a", fresh, None);
448 let snap = ring.snapshot("local", "actor:a");
449 assert_eq!(snap.len(), 1);
450 assert_eq!(snap[0].id, fresh);
451 }
452
453 #[test]
454 fn actor_isolation_never_crosses_boundary() {
455 let ring = ReferenceRing::new();
456 let id_a = Uuid::new_v4();
457 ring.admit("local", "actor:a", id_a, Some("A-only".into()));
458 let snap_b = ring.snapshot("local", "actor:b");
459 assert!(snap_b.is_empty(), "actor b must never see actor a's ring");
460 let snap_a = ring.snapshot("local", "actor:a");
461 assert_eq!(snap_a.len(), 1);
462 }
463
464 #[test]
465 fn namespace_isolation_is_independent_of_actor_isolation() {
466 let ring = ReferenceRing::new();
467 let id = Uuid::new_v4();
468 ring.admit("tenant-a", "actor:a", id, None);
469 assert!(ring.snapshot("tenant-b", "actor:a").is_empty());
470 assert_eq!(ring.snapshot("tenant-a", "actor:a").len(), 1);
471 }
472
473 #[test]
474 fn re_admitting_an_id_moves_it_to_most_recent_without_duplicating() {
475 let ring = ReferenceRing::new();
476 let a = Uuid::new_v4();
477 let b = Uuid::new_v4();
478 ring.admit("local", "actor:a", a, Some("A".into()));
479 ring.admit("local", "actor:a", b, Some("B".into()));
480 ring.admit("local", "actor:a", a, Some("A-renamed".into()));
481 let snap = ring.snapshot("local", "actor:a");
482 assert_eq!(snap.len(), 2, "re-admission must not duplicate the entry");
483 assert_eq!(snap[0].id, a, "re-admitted id must be most-recent");
484 assert_eq!(snap[0].name.as_deref(), Some("A-renamed"));
485 }
486
487 #[test]
488 fn snapshot_prunes_a_key_that_ages_out_entirely() {
489 let ring = ReferenceRing::with_bounds(64, Duration::from_millis(20));
490 ring.admit("local", "actor:a", Uuid::new_v4(), None);
491 std::thread::sleep(Duration::from_millis(40));
492 assert!(ring.snapshot("local", "actor:a").is_empty());
494 let state = ring.state.lock().unwrap();
495 assert!(
496 !state
497 .rings
498 .contains_key(&("local".to_string(), "actor:a".to_string())),
499 "a key whose ring emptied via TTL eviction must not linger in the outer map"
500 );
501 }
502
503 #[test]
509 fn snapshot_prunes_other_stale_keys_it_did_not_query() {
510 let ring = ReferenceRing::with_bounds(64, Duration::from_millis(20));
511 ring.admit("local", "actor:queried", Uuid::new_v4(), None);
512 ring.admit("local", "actor:other", Uuid::new_v4(), None);
513 std::thread::sleep(Duration::from_millis(40));
514
515 assert!(ring.snapshot("local", "actor:queried").is_empty());
517
518 let state = ring.state.lock().unwrap();
519 assert!(
520 !state
521 .rings
522 .contains_key(&("local".to_string(), "actor:other".to_string())),
523 "snapshotting one actor must also prune OTHER actors' fully-stale keys, \
524 not just the one queried"
525 );
526 }
527
528 #[test]
529 fn outer_map_evicts_least_recently_touched_keys_over_budget() {
530 let ring = ReferenceRing::with_bounds_and_outer_limit(64, DEFAULT_RING_TTL, 3);
531 for i in 0..4 {
535 ring.admit(
536 "local",
537 &format!("actor:{i}"),
538 Uuid::new_v4(),
539 Some(format!("actor {i}")),
540 );
541 }
542 assert!(
543 ring.snapshot("local", "actor:0").is_empty(),
544 "the least-recently-touched key must be evicted once the outer-key budget is exceeded"
545 );
546 for i in 1..4 {
547 assert!(
548 !ring.snapshot("local", &format!("actor:{i}")).is_empty(),
549 "actor:{i} must survive the budget eviction"
550 );
551 }
552 }
553
554 #[test]
555 fn outer_map_budget_eviction_never_evicts_the_key_just_admitted() {
556 let ring = ReferenceRing::with_bounds_and_outer_limit(64, DEFAULT_RING_TTL, 1);
557 ring.admit("local", "actor:a", Uuid::new_v4(), None);
558 ring.admit("local", "actor:b", Uuid::new_v4(), None);
562 assert!(!ring.snapshot("local", "actor:b").is_empty());
563 }
564
565 #[test]
570 fn admit_and_snapshot_recover_from_poisoned_mutex() {
571 let ring = std::sync::Arc::new(ReferenceRing::new());
572 let poison_ring = ring.clone();
573 let _ = std::thread::spawn(move || {
574 let _guard = poison_ring.state.lock().unwrap();
575 panic!("deliberately poisoning the reference ring mutex for the recovery test");
576 })
577 .join();
578
579 ring.admit(
581 "local",
582 "actor:a",
583 Uuid::new_v4(),
584 Some("post-poison".into()),
585 );
586 let snap = ring.snapshot("local", "actor:a");
587 assert_eq!(snap.len(), 1);
588 assert_eq!(snap[0].name.as_deref(), Some("post-poison"));
589 }
590}