Skip to main content

khive_runtime/
kg_read.rs

1//! Shared KG handle resolution and entity-delete routing for pack-scoped backends.
2
3use std::collections::{BTreeSet, HashMap, HashSet};
4
5use khive_storage::types::{Direction, NeighborCursor, NeighborHit, NeighborQuery};
6use uuid::Uuid;
7
8use crate::{KhiveRuntime, Namespace, NamespaceToken, Resolved, RuntimeError, VerbRegistry};
9
10/// Query options for a KG neighbor read whose origin may live on a pack backend.
11pub struct KgNeighborRead {
12    pub query: NeighborQuery,
13    pub after: Option<NeighborCursor>,
14    pub neighbor_kinds: Option<Vec<String>>,
15    pub enrich: bool,
16    /// Narrow graph selection to an already visible namespace while retaining
17    /// the caller token's identity and originating request metadata.
18    pub namespace: Option<Namespace>,
19}
20
21pub(crate) fn neighbor_read_namespaces<'a>(
22    token: &'a NamespaceToken,
23    namespace: Option<&'a Namespace>,
24) -> Result<&'a [Namespace], RuntimeError> {
25    match namespace {
26        Some(namespace) if token.visible_namespaces().contains(namespace) => {
27            Ok(std::slice::from_ref(namespace))
28        }
29        Some(_) => Err(RuntimeError::InvalidInput(
30            "KG neighbor namespace must already be visible to the caller".into(),
31        )),
32        None => Ok(token.visible_namespaces()),
33    }
34}
35
36impl VerbRegistry {
37    /// Resolve a live KG origin across configured backends before expanding
38    /// adjacency on the graph runtime with the original caller token.
39    pub async fn neighbors_for_kg_read(
40        &self,
41        runtime: &KhiveRuntime,
42        token: &NamespaceToken,
43        node_id: Uuid,
44        options: KgNeighborRead,
45    ) -> Result<Vec<NeighborHit>, RuntimeError> {
46        self.neighbors_for_kg_read_inner(runtime, token, node_id, options, false)
47            .await
48            .map(|(hits, _)| hits)
49    }
50
51    /// Resolve a live origin and share entity-kind hints from its graph's
52    /// deletion screen with mailbox filtering, retaining the original token.
53    pub async fn neighbors_for_kg_read_with_entity_kinds(
54        &self,
55        runtime: &KhiveRuntime,
56        token: &NamespaceToken,
57        node_id: Uuid,
58        options: KgNeighborRead,
59    ) -> Result<(Vec<NeighborHit>, HashMap<Uuid, String>), RuntimeError> {
60        self.neighbors_for_kg_read_inner(runtime, token, node_id, options, true)
61            .await
62    }
63
64    async fn neighbors_for_kg_read_inner(
65        &self,
66        runtime: &KhiveRuntime,
67        token: &NamespaceToken,
68        node_id: Uuid,
69        options: KgNeighborRead,
70        with_entity_kinds: bool,
71    ) -> Result<(Vec<NeighborHit>, HashMap<Uuid, String>), RuntimeError> {
72        neighbor_read_namespaces(token, options.namespace.as_ref())?;
73        if self
74            .resolve_kg_read_by_id(runtime, token, node_id, false)
75            .await?
76            .is_none()
77            && !runtime.substrate_exists_by_id(token, node_id).await?
78        {
79            return Err(RuntimeError::NotFound(format!(
80                "neighbor anchor {node_id} not found"
81            )));
82        }
83        if with_entity_kinds {
84            runtime
85                .neighbors_for_resolved_kg_read_with_entity_kinds(token, node_id, options)
86                .await
87        } else {
88            runtime
89                .neighbors_for_resolved_kg_read(token, node_id, options)
90                .await
91                .map(|hits| (hits, HashMap::new()))
92        }
93    }
94
95    /// The directed form of [`Self::neighbors_for_kg_read`], retaining stored
96    /// edge direction and the existing graph namespace selection.
97    pub async fn directed_neighbors_for_kg_read(
98        &self,
99        runtime: &KhiveRuntime,
100        token: &NamespaceToken,
101        node_id: Uuid,
102        query: NeighborQuery,
103        namespace: Option<&Namespace>,
104    ) -> Result<Vec<(NeighborHit, Direction)>, RuntimeError> {
105        neighbor_read_namespaces(token, namespace)?;
106        if self
107            .resolve_kg_read_by_id(runtime, token, node_id, false)
108            .await?
109            .is_none()
110            && !runtime.substrate_exists_by_id(token, node_id).await?
111        {
112            return Err(RuntimeError::NotFound(format!(
113                "neighbor anchor {node_id} not found"
114            )));
115        }
116        runtime
117            .directed_neighbors_for_resolved_kg_read(token, node_id, query, namespace)
118            .await
119    }
120}
121
122pub(crate) struct KgReadResolver {
123    runtimes: Vec<KhiveRuntime>,
124    primary: KhiveRuntime,
125    by_pack: HashMap<String, KhiveRuntime>,
126}
127
128impl KgReadResolver {
129    pub(crate) fn new(primary: &KhiveRuntime, runtimes: &HashMap<String, KhiveRuntime>) -> Self {
130        // Keep primary precedence deterministic and query each assigned backend
131        // once, even when several packs share it. These handles own no registry,
132        // so the registry-held topology creates no reference cycle.
133        let mut seen = HashSet::new();
134        let mut unique = Vec::new();
135        let mut packs: Vec<_> = runtimes.iter().collect();
136        packs.sort_by_key(|(name, _)| *name);
137        for runtime in std::iter::once(primary).chain(packs.into_iter().map(|(_, rt)| rt)) {
138            if seen.insert(runtime.backend_id().clone()) {
139                unique.push(runtime.clone());
140            }
141        }
142        Self {
143            runtimes: unique,
144            primary: primary.clone(),
145            by_pack: runtimes.clone(),
146        }
147    }
148
149    pub(crate) fn runtime_for_pack(&self, pack: &str) -> &KhiveRuntime {
150        self.by_pack.get(pack).unwrap_or(&self.primary)
151    }
152
153    pub(crate) async fn by_id(
154        &self,
155        token: &NamespaceToken,
156        id: Uuid,
157        include_deleted: bool,
158    ) -> Result<Option<Resolved>, RuntimeError> {
159        let mut found = None;
160        for runtime in &self.runtimes {
161            let candidate = if include_deleted {
162                runtime.resolve_by_id_including_deleted(token, id).await?
163            } else {
164                runtime.resolve_by_id(token, id).await?
165            };
166            // Preserve entity-before-note substrate precedence. A success must
167            // not hide a failure from a later configured backend.
168            if found.is_none()
169                || (matches!(&found, Some(Resolved::Note(_)))
170                    && matches!(&candidate, Some(Resolved::Entity(_))))
171            {
172                found = candidate;
173            }
174        }
175        Ok(found)
176    }
177
178    pub(crate) async fn entity_runtime(
179        &self,
180        token: &NamespaceToken,
181        id: Uuid,
182    ) -> Result<Option<KhiveRuntime>, RuntimeError> {
183        let mut owner = None;
184        for runtime in &self.runtimes {
185            let store = runtime.entities(token)?;
186            let entity = store.get_entity_including_deleted(id).await?;
187            if entity.is_some() {
188                if owner.is_some() {
189                    return Err(RuntimeError::InvalidInput(format!(
190                        "entity {id} exists on multiple backends; deletion is ambiguous"
191                    )));
192                }
193                owner = Some(runtime.clone());
194            }
195        }
196        Ok(owner)
197    }
198
199    pub(crate) async fn prefix(
200        &self,
201        prefix: &str,
202        include_deleted: bool,
203    ) -> Result<Option<Uuid>, RuntimeError> {
204        let mut matches = BTreeSet::new();
205        for runtime in &self.runtimes {
206            let candidate = runtime
207                .resolve_prefix_for_kg_read(prefix, include_deleted)
208                .await;
209            match candidate {
210                Ok(Some(id)) => {
211                    matches.insert(id);
212                }
213                Ok(None) => {}
214                Err(RuntimeError::AmbiguousPrefix { matches: ids, .. }) => matches.extend(ids),
215                Err(error) => return Err(error),
216            }
217        }
218        match matches.len() {
219            0 => Ok(None),
220            1 => Ok(matches.into_iter().next()),
221            _ => Err(RuntimeError::AmbiguousPrefix {
222                prefix: prefix.into(),
223                matches: matches.into_iter().collect(),
224            }),
225        }
226    }
227}
228
229#[cfg(test)]
230mod tests {
231    use super::*;
232    use crate::{BackendId, Namespace, RuntimeConfig};
233    use chrono::Utc;
234    use khive_storage::types::{Edge, LinkId};
235    use khive_storage::{Entity, Note};
236
237    #[tokio::test]
238    async fn kg_neighbor_namespace_selection_is_narrowing_and_keeps_directed_self_loops() {
239        let project = Namespace::parse("project").unwrap();
240        let main = KhiveRuntime::new(RuntimeConfig {
241            db_path: None,
242            backend_id: BackendId::main(),
243            actor_id: Some("reader".into()),
244            visible_namespaces: vec![project.clone()],
245            ..RuntimeConfig::no_embeddings()
246        })
247        .unwrap();
248        let token = main
249            .authorize_with_visibility(Namespace::local(), vec![project.clone()])
250            .unwrap();
251        let anchor = Entity::new("local", "concept", "namespace selector anchor");
252        let neighbor = Entity::new("project", "concept", "namespace selector neighbor");
253        for entity in [&anchor, &neighbor] {
254            main.entities(&token)
255                .unwrap()
256                .upsert_entity(entity.clone())
257                .await
258                .unwrap();
259        }
260        let self_loop = Uuid::new_v4();
261        for (namespace, id, source) in [
262            (Namespace::local(), self_loop, anchor.id),
263            (project.clone(), Uuid::new_v4(), neighbor.id),
264        ] {
265            let edge_token = main.authorize(namespace.clone()).unwrap();
266            let now = Utc::now();
267            main.graph(&edge_token)
268                .unwrap()
269                .upsert_edge(Edge {
270                    id: LinkId(id),
271                    namespace: namespace.to_string(),
272                    source_id: source,
273                    target_id: anchor.id,
274                    relation: khive_storage::EdgeRelation::Extends,
275                    weight: 1.0,
276                    created_at: now,
277                    updated_at: now,
278                    deleted_at: None,
279                    metadata: None,
280                    target_backend: None,
281                })
282                .await
283                .unwrap();
284        }
285        let registry = crate::VerbRegistryBuilder::new().build().unwrap();
286        let query = NeighborQuery {
287            direction: Direction::Both,
288            relations: None,
289            limit: Some(10),
290            min_weight: None,
291        };
292        let baseline = main
293            .neighbors_with_query_directed(&token, anchor.id, query.clone())
294            .await
295            .unwrap();
296        let full = registry
297            .directed_neighbors_for_kg_read(&main, &token, anchor.id, query.clone(), None)
298            .await
299            .unwrap();
300        let keys = |hits: &[(NeighborHit, Direction)]| {
301            hits.iter()
302                .map(|(hit, direction)| (hit.node_id, hit.edge_id, direction.clone()))
303                .collect::<Vec<_>>()
304        };
305        assert_eq!(keys(&full), keys(&baseline));
306        let local = registry
307            .directed_neighbors_for_kg_read(
308                &main,
309                &token,
310                anchor.id,
311                query.clone(),
312                Some(&Namespace::local()),
313            )
314            .await
315            .unwrap();
316        assert_eq!(local.len(), 2);
317        assert_eq!(local[0].0.edge_id, self_loop);
318        assert_eq!(local[0].1, Direction::Out);
319        assert_eq!(local[1].1, Direction::In);
320        let project_hits = registry
321            .neighbors_for_kg_read(
322                &main,
323                &token,
324                anchor.id,
325                KgNeighborRead {
326                    query: query.clone(),
327                    after: None,
328                    neighbor_kinds: None,
329                    enrich: true,
330                    namespace: Some(project),
331                },
332            )
333            .await
334            .unwrap();
335        assert_eq!(project_hits.len(), 1);
336        assert_eq!(project_hits[0].node_id, neighbor.id);
337        let outside = Namespace::parse("outside").unwrap();
338        assert!(matches!(
339            registry
340                .directed_neighbors_for_kg_read(
341                    &main,
342                    &token,
343                    anchor.id,
344                    query.clone(),
345                    Some(&outside)
346                )
347                .await,
348            Err(RuntimeError::InvalidInput(_))
349        ));
350        assert!(matches!(
351            registry
352                .neighbors_for_kg_read(
353                    &main,
354                    &token,
355                    anchor.id,
356                    KgNeighborRead {
357                        query,
358                        after: None,
359                        neighbor_kinds: None,
360                        enrich: true,
361                        namespace: Some(outside),
362                    },
363                )
364                .await,
365            Err(RuntimeError::InvalidInput(_))
366        ));
367    }
368
369    fn runtime(name: &str) -> KhiveRuntime {
370        KhiveRuntime::new(RuntimeConfig {
371            db_path: None,
372            backend_id: BackendId::parse(name).unwrap(),
373            ..RuntimeConfig::no_embeddings()
374        })
375        .unwrap()
376    }
377
378    async fn note(runtime: &KhiveRuntime, id: Uuid, namespace: &str) -> Note {
379        let token = runtime
380            .authorize(Namespace::parse(namespace).unwrap())
381            .unwrap();
382        let mut note = Note::new(namespace, "observation", "read routing fixture");
383        note.id = id;
384        runtime
385            .notes(&token)
386            .unwrap()
387            .upsert_note(note.clone())
388            .await
389            .unwrap();
390        note
391    }
392
393    #[tokio::test]
394    async fn issue2992_shared_reads_deduplicate_backends_and_preserve_global_prefix_ambiguity() {
395        let main = runtime("main");
396        let secondary = runtime("comm");
397        let resolver = KgReadResolver::new(
398            &main,
399            &HashMap::from([
400                ("kg".into(), main.clone()),
401                ("comm".into(), secondary.clone()),
402                ("another-pack".into(), secondary.clone()),
403            ]),
404        );
405        assert_eq!(resolver.runtimes.len(), 2);
406        let first = Uuid::parse_str("cafe1234-0000-4000-8000-000000000001").unwrap();
407        let second = Uuid::parse_str("cafe1234-0000-4000-8000-000000000002").unwrap();
408        let remote = note(&secondary, first, "elsewhere").await;
409        let token = main.authorize(Namespace::local()).unwrap();
410        assert!(
411            matches!(resolver.by_id(&token, first, false).await.unwrap(), Some(Resolved::Note(found)) if found == remote)
412        );
413        assert_eq!(
414            resolver.prefix("cafe1234", false).await.unwrap(),
415            Some(first)
416        );
417        // The same global UUID in another backend is not a distinct prefix candidate.
418        note(&main, first, "local").await;
419        assert_eq!(
420            resolver.prefix("cafe1234", false).await.unwrap(),
421            Some(first)
422        );
423        note(&main, second, "local").await;
424        let error = resolver.prefix("cafe1234", false).await.unwrap_err();
425        assert!(
426            matches!(error, RuntimeError::AmbiguousPrefix { matches, .. } if matches == vec![first, second])
427        );
428        assert!(resolver.prefix("deadbeef", false).await.unwrap().is_none());
429        assert!(resolver
430            .by_id(&token, Uuid::new_v4(), false)
431            .await
432            .unwrap()
433            .is_none());
434    }
435
436    #[tokio::test]
437    async fn issue2992_shared_reads_surface_backend_read_failures_instead_of_partial_results() {
438        let main = runtime("main");
439        let secondary = runtime("comm");
440        let resolver =
441            KgReadResolver::new(&main, &HashMap::from([("comm".into(), secondary.clone())]));
442        let id = Uuid::parse_str("bad01234-0000-4000-8000-000000000001").unwrap();
443        note(&main, id, "local").await;
444        let token = main.authorize(Namespace::local()).unwrap();
445        assert!(resolver.by_id(&token, id, false).await.unwrap().is_some());
446        assert_eq!(resolver.prefix("bad01234", false).await.unwrap(), Some(id));
447        // A backend read that cannot be admitted is a storage failure of the
448        // whole lookup, never a quiet miss: a pooled reader checkout refused by
449        // an exhausted request read deadline surfaces through both entry points.
450        // (Schema faults injected through the SQL writer are invisible to the
451        // pooled readers' snapshots, so admission is the fault a test can inject.)
452        let expired = std::time::Duration::ZERO;
453        assert!(matches!(
454            khive_storage::scope_request_read_deadline(expired, resolver.by_id(&token, id, false))
455                .await,
456            Err(RuntimeError::Storage(_) | RuntimeError::Sqlite(_))
457        ));
458        assert!(matches!(
459            khive_storage::scope_request_read_deadline(expired, resolver.prefix("bad01234", false))
460                .await,
461            Err(RuntimeError::Storage(_) | RuntimeError::Sqlite(_))
462        ));
463        assert!(matches!(
464            khive_storage::scope_request_read_deadline(expired, resolver.prefix("deadbeef", false))
465                .await,
466            Err(RuntimeError::Storage(_) | RuntimeError::Sqlite(_))
467        ));
468    }
469
470    #[tokio::test]
471    async fn issue2992_shared_reads_include_live_entities_and_explicit_deleted_notes() {
472        let main = runtime("main");
473        let secondary = runtime("archive");
474        let resolver = KgReadResolver::new(
475            &main,
476            &HashMap::from([("archive".into(), secondary.clone())]),
477        );
478        let token = main.authorize(Namespace::local()).unwrap();
479        let entity = Entity::new("elsewhere", "concept", "remote entity");
480        let entity_id = entity.id;
481        secondary
482            .entities(&token)
483            .unwrap()
484            .upsert_entity(entity)
485            .await
486            .unwrap();
487        assert!(
488            matches!(resolver.by_id(&token, entity_id, false).await.unwrap(), Some(Resolved::Entity(found)) if found.id == entity_id)
489        );
490        let mut deleted = Note::new("local", "observation", "deleted remote note");
491        deleted.deleted_at = Some(deleted.updated_at);
492        secondary
493            .notes(&token)
494            .unwrap()
495            .upsert_note(deleted.clone())
496            .await
497            .unwrap();
498        assert!(resolver
499            .by_id(&token, deleted.id, false)
500            .await
501            .unwrap()
502            .is_none());
503        assert!(
504            matches!(resolver.by_id(&token, deleted.id, true).await.unwrap(), Some(Resolved::Note(found)) if found == deleted)
505        );
506        let prefix = deleted.id.simple().to_string();
507        assert_eq!(
508            resolver.prefix(&prefix, true).await.unwrap(),
509            Some(deleted.id)
510        );
511    }
512}