use std::collections::HashMap;
use async_trait::async_trait;
use uuid::Uuid;
use khive_mcp::coordinator::{
BackendSearchFailure as CoordBackendFailure,
BackendSearchFailureKind as CoordBackendFailureKind, BackendSearchResult as CoordBackendResult,
CoordError, CoordLinkResult, CoordSearchResult, CoordinatorService,
};
use khive_pack_kg::handlers::ValidatedSearchRequest;
use khive_runtime::BackendId;
use khive_runtime::Namespace;
use khive_storage::EdgeRelation;
use super::dispatch::{BackendSearchFailureKind, SubstrateCoordinator};
pub struct SubstrateCoordinatorService {
inner: SubstrateCoordinator,
}
impl SubstrateCoordinatorService {
pub fn new(coordinator: SubstrateCoordinator) -> Self {
Self { inner: coordinator }
}
pub fn primary_backend_id_inner(&self) -> Option<BackendId> {
self.inner.primary_runtime().map(|_| BackendId::main())
}
}
#[async_trait]
impl CoordinatorService for SubstrateCoordinatorService {
async fn locate(&self, id: Uuid) -> Option<BackendId> {
let ns = Namespace::local();
self.inner.locate(id, &ns).await
}
fn record_created(&self, id: Uuid, backend_id: BackendId) {
self.inner.record_created(id, backend_id);
}
fn primary_backend_id(&self) -> Option<BackendId> {
self.inner.registry().primary().map(|e| e.id.clone())
}
#[allow(clippy::too_many_arguments)]
async fn link(
&self,
namespace: &Namespace,
source_id: Uuid,
target_id: Uuid,
relation: EdgeRelation,
weight: f64,
metadata: Option<serde_json::Value>,
resurrect: bool,
) -> Result<CoordLinkResult, CoordError> {
self.inner
.link_cross_backend_observed(
namespace, source_id, target_id, relation, weight, metadata, resurrect,
)
.await
.and_then(|row| {
let cross_backend = row.edge.target_backend.is_some();
let target_backend_id = row
.edge
.target_backend
.as_deref()
.map(BackendId::parse)
.transpose()
.map_err(|error| format!("stored target backend is invalid: {error}"))?;
Ok(CoordLinkResult {
edge: row.edge,
cross_backend,
target_backend_id,
mutation: row.disposition,
})
})
.map_err(|msg| {
if msg.contains("not found on any backend") {
CoordError::Backend(msg)
} else if msg.contains("edge rule violation")
|| msg.contains("self-loop")
|| msg.contains("must be a note")
{
CoordError::EdgeRuleViolation(msg)
} else {
CoordError::Backend(msg)
}
})
}
async fn fan_out_search(
&self,
request: &ValidatedSearchRequest,
namespace: &Namespace,
extra_visible: &[Namespace],
) -> CoordSearchResult {
let (entity_hits, note_hits, per_backend) = self
.inner
.fan_out_search_with_visibility(request, namespace, extra_visible)
.await;
let partial = per_backend.iter().any(|r| r.error.is_some());
let mut entity_kinds: HashMap<Uuid, String> = HashMap::new();
let mut entity_created_at: HashMap<Uuid, i64> = HashMap::new();
let mut entity_updated_at: HashMap<Uuid, i64> = HashMap::new();
let mut entity_versions: HashMap<Uuid, i64> = HashMap::new();
for hit in &entity_hits {
if khive_storage::request_read_is_cancelled() {
break;
}
let backend_id = self.inner.locate(hit.entity_id, namespace).await;
if let Some(bid) = backend_id {
if let Some(entry) = self.inner.registry().get(&bid) {
let rt = &entry.runtime;
if let Ok(token) = rt.authorize(namespace.clone()) {
if let Ok(entity) = rt.get_entity(&token, hit.entity_id).await {
entity_created_at.insert(hit.entity_id, entity.created_at);
entity_updated_at.insert(hit.entity_id, entity.updated_at);
entity_versions.insert(hit.entity_id, entity.version);
entity_kinds.insert(hit.entity_id, entity.kind);
}
}
}
}
}
let mut note_kinds: HashMap<Uuid, String> = HashMap::new();
let mut note_created_at: HashMap<Uuid, i64> = HashMap::new();
let mut note_updated_at: HashMap<Uuid, i64> = HashMap::new();
let mut note_versions: HashMap<Uuid, i64> = HashMap::new();
let mut note_names: HashMap<Uuid, Option<String>> = HashMap::new();
for hit in ¬e_hits {
if khive_storage::request_read_is_cancelled() {
break;
}
let backend_id = self.inner.locate(hit.note_id, namespace).await;
if let Some(bid) = backend_id {
if let Some(entry) = self.inner.registry().get(&bid) {
let rt = &entry.runtime;
if let Ok(token) = rt.authorize(namespace.clone()) {
if let Ok(store) = rt.notes(&token) {
if let Ok(Some(note)) = store.get_note(hit.note_id).await {
note_versions.insert(hit.note_id, note.version);
note_created_at.insert(hit.note_id, note.created_at);
note_updated_at.insert(hit.note_id, note.updated_at);
note_names.insert(hit.note_id, note.name.clone());
note_kinds.insert(hit.note_id, note.kind);
}
}
}
}
}
}
let coord_per_backend: Vec<CoordBackendResult> = per_backend
.into_iter()
.map(|r| {
let vector_selected = self
.inner
.registry()
.get(&r.backend_id)
.is_some_and(|entry| entry.runtime.vector_arm_selected());
CoordBackendResult {
backend_id: r.backend_id,
entity_hits: r.hits,
note_hits: r.note_hits,
vector_selected,
error: r.error.map(|failure| CoordBackendFailure {
kind: match failure.kind {
BackendSearchFailureKind::BackendError => {
CoordBackendFailureKind::BackendError
}
BackendSearchFailureKind::Timeout => CoordBackendFailureKind::Timeout,
},
message: failure.message,
}),
vector_error: r.vector_error,
}
})
.collect();
CoordSearchResult {
entity_hits,
note_hits,
per_backend: coord_per_backend,
partial,
entity_kinds,
note_kinds,
entity_created_at,
entity_updated_at,
entity_versions,
note_created_at,
note_updated_at,
note_versions,
note_names,
}
}
fn is_single_backend(&self) -> bool {
self.inner.is_single_backend()
}
}