1use 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
10pub struct KgNeighborRead {
12 pub query: NeighborQuery,
13 pub after: Option<NeighborCursor>,
14 pub neighbor_kinds: Option<Vec<String>>,
15 pub enrich: bool,
16 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 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 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 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 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 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 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 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}