Skip to main content

khive_runtime/
reference_ring.rs

1//! Recently-referenced ring: a bounded, per-`(namespace, actor)` cache of ids
2//! this actor recently touched by name, held in daemon-warm memory only.
3//!
4//! No schema, no persistence, no migration — a daemon restart empties it, and
5//! that is fine: its miss path is the hybrid-search fallback in
6//! `reference_resolution::resolve_reference`. Admission is gated by the
7//! dispatch boundary (`pack.rs::dispatch_with_identity`) under a strict rule:
8//! only by-id touches (create/get/update/delete/merge/link) admit an id.
9//! `search`/`list` result-sets never enter the ring — the anaphora signal is
10//! the sparsity, so admitting every search hit would drown "the old record"
11//! in noise.
12
13use std::collections::{HashMap, VecDeque};
14use std::sync::Mutex;
15use std::time::{Duration, Instant};
16
17use serde_json::Value;
18use uuid::Uuid;
19
20/// Default ring size per `(namespace, actor)` key.
21pub const DEFAULT_RING_CAPACITY: usize = 64;
22
23/// Default eviction age.
24pub const DEFAULT_RING_TTL: Duration = Duration::from_secs(30 * 60);
25
26/// Maximum distinct `(namespace, actor)` keys held at once. The per-key ring
27/// is capacity/TTL-bounded (above), but the outer map itself has no such
28/// bound by default — a daemon serving many transient actor ids needs an
29/// eviction rule for the map too, or it grows without limit. LRU by
30/// most-recent touch, sized generously above the 64-entry inner bound since
31/// this is the actor-fanout axis, not the per-actor recency axis.
32pub const DEFAULT_MAX_OUTER_KEYS: usize = 4096;
33
34/// One admitted id: the id itself, a best-effort display name (drawn from the
35/// dispatch result JSON — never a fresh fetch, to keep admission cheap on the
36/// Tier-1 latency budget), and the instant it was touched.
37#[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
51/// Daemon-warm, actor-scoped recently-referenced ring.
52///
53/// Privacy is structural: rings are keyed by `(namespace, actor)` and a
54/// snapshot for one actor never observes another actor's entries, even
55/// within the same namespace.
56pub 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    /// Lock the shared state, recovering from mutex poisoning instead of
96    /// panicking: ring admission is best-effort cache maintenance and must
97    /// never turn an unrelated panic into a failure for every future caller.
98    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    /// Evict entries older than `ttl` from the front of `ring` (front = oldest,
112    /// since admission always pushes to the back with the current instant).
113    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    /// Drop empty/aged-out outer-map keys, then if still over `max_outer_keys`
124    /// evict the least-recently-touched surviving keys until back within
125    /// budget. `exempt` (the key that triggered this admission) is never
126    /// evicted: admitting an id must never immediately evict its own ring.
127    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    /// Admit `id` into the `(namespace, actor)` ring, touching it to "now".
154    ///
155    /// Re-admitting an id already present moves it to the back (most-recent)
156    /// instead of duplicating the entry. Eviction is size-or-age, then the
157    /// outer map itself is pruned (see `prune_outer_map`) so it never grows
158    /// without bound.
159    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    /// Snapshot the live (non-stale) entries for `(namespace, actor)`,
180    /// most-recently-touched first. Returns an empty vec for an unknown or
181    /// empty key: never an error, since a ring miss is always legitimate
182    /// (fresh session, restarted daemon, cross-actor query).
183    ///
184    /// Also runs `prune_outer_map` over the whole outer map, not just the
185    /// queried key: otherwise a daemon that only ever reads (never admits)
186    /// would never prune any other actor's stale key.
187    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
203/// Display name extracted from a dispatch result's `name` field. No fallback
204/// to a note's `content`: the ring only ever admits entities, and free text
205/// must never stand in for an entity's actual name.
206fn 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
212/// The nine closed entity kinds. Duplicated here rather than depending on
213/// khive-pack-kg's vocab, which would invert the crate dependency direction.
214const 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
222/// Whether a `create`/`get`/`update`/`delete` result JSON denotes an entity
223/// (the ring only ever admits entity ids, plus `link` endpoints). No extra
224/// storage read: the discriminator is read off the shape khive-pack-kg's
225/// handlers already return.
226///
227/// - `edge`/`event` results carry an explicit top-level `"kind": "edge"` /
228///   `"kind": "event"`: never entities.
229/// - `create`/`get`/`update` return the raw storage record: entities always
230///   serialize an `entity_type` key (even as `null`), notes always serialize
231///   a `content` key: the two shapes never overlap, so key presence alone
232///   is a reliable substrate discriminator.
233/// - `delete` returns a synthetic `{deleted, id, kind}` summary carrying
234///   neither field; its `kind` is the one the handler RESOLVED off the row
235///   before removing it, so a caller who deleted by a bare id or a hex prefix
236///   still names a substrate here. A pack resolver's own delete path can still
237///   answer without one. An absent or unrecognised `kind` never admits: a
238///   delete admission is a nice-to-have, so skipping is always the safe choice
239///   over guessing.
240fn 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
258/// Compute the `(id, name)` pairs a successful `verb` dispatch admits to the
259/// ring, from its already-serialized JSON `result` alone — no extra storage
260/// reads, so admission stays cheap on the Tier-1 latency budget.
261///
262/// Only singleton by-id touches admit: `create`, `get`, `update`, `delete`,
263/// `merge`, and `link` (both endpoints). Bulk shapes (`items=[...]`,
264/// `links=[...]`) are identifiable by an `attempted` count in their response
265/// and are excluded: they name multiple ids, not the one-caller-named-id
266/// semantic the ring exists to serve. `search`/`list` never reach this
267/// function — they are not in the verb match below.
268pub(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 is entity-only by construction: no substrate check needed;
291        // `MergeSummary` carries no `kind` field to check even if one were wanted.
292        "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        // Real entity get/create/update responses always carry `entity_type`
336        // (even as `null`) — that key's presence is the entity-substrate
337        // discriminator `substrate_admits_as_entity` checks for.
338        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        // Real note get/create/update responses carry `content`, never
347        // `entity_type`: the ring only ever admits entity ids.
348        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        // delete's synthetic summary has neither `content` nor `entity_type`;
368        // it falls back to `kind`, which the handler resolves off the row it
369        // removed rather than echoing back what the caller typed.
370        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", &note_delete).is_empty());
383        // A pack resolver's own delete path answers without a kind. That still
384        // has to skip rather than guess, so the arm stays.
385        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        // Most-recently-touched first; the two oldest (ids[0], ids[1]) were evicted.
435        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        // The read itself must observe (and clean up) the now-fully-stale key.
493        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    /// Regression: `snapshot` must sweep the whole outer map, not just the
504    /// queried key: otherwise a read-only daemon (one that only ever calls
505    /// `snapshot`, never `admit`) never prunes any other actor's stale key.
506    /// Two actors age out; only `actor:queried` is snapshotted, but
507    /// `actor:other`'s stale key must be removed too.
508    #[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        // Only `actor:queried` is ever snapshotted directly.
516        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        // Admit four distinct actors in order; the outer-key budget is 3, so
532        // the least-recently-touched one (actor:0) must be evicted once the
533        // fourth actor is admitted.
534        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        // Admitting actor:b pushes the outer map to 2 keys against a budget
559        // of 1; the key that triggered this admission (actor:b) must survive
560        // even though it is, by construction, also the most-recently-touched.
561        ring.admit("local", "actor:b", Uuid::new_v4(), None);
562        assert!(!ring.snapshot("local", "actor:b").is_empty());
563    }
564
565    /// A prior panic while holding the ring's mutex must never turn a
566    /// later, unrelated `admit`/`snapshot` call into a panic too — admission
567    /// is best-effort cache maintenance riding on an already-successful
568    /// dispatch.
569    #[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        // Both calls must complete normally (not panic) despite the poison.
580        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}