use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use khive_gate::GateRef;
use khive_storage::EventStore;
use khive_types::Namespace;
use serde_json::Value;
use crate::error::RuntimeError;
use crate::runtime::NamespaceToken;
use super::{
DispatchHook, EndpointKind, HandlerDef, PackByIdResolver, PackRuntime, SPECIAL_RELATIONS,
};
#[derive(Clone)]
pub struct VerbRegistry {
pub(super) packs: std::sync::Arc<Vec<Box<dyn PackRuntime>>>,
pub(super) resolvers: std::sync::Arc<Vec<(String, Box<dyn PackByIdResolver>)>>,
pub(super) kg_read_resolver: Option<Arc<crate::kg_read::KgReadResolver>>,
pub(super) gate: GateRef,
pub(super) default_namespace: String,
pub(super) visible_namespaces: Vec<Namespace>,
pub(super) actor_id: Option<String>,
pub(super) event_store: Option<Arc<dyn EventStore>>,
pub(super) audit_store_read_only: bool,
pub(super) dispatch_hook: Option<Arc<dyn DispatchHook>>,
pub(super) available_verbs: Arc<Vec<&'static str>>,
pub(super) handler_by_name: Arc<HashMap<&'static str, &'static HandlerDef>>,
pub(super) degrade_safe_verbs: Arc<HashSet<&'static str>>,
pub(super) read_replay_safe_verbs: Arc<HashSet<&'static str>>,
pub(super) reference_ring: Arc<crate::reference_ring::ReferenceRing>,
pub(super) audit_batch: Option<Arc<crate::audit_batch::AuditBatch>>,
}
#[derive(Debug, Clone, PartialEq)]
pub struct InterceptedDispatchResult<M> {
pub result: Value,
pub metadata: M,
}
impl<M> InterceptedDispatchResult<M> {
pub fn new(result: Value, metadata: M) -> Self {
Self { result, metadata }
}
}
#[derive(Debug, Clone, Default)]
pub struct RequestIdentity {
pub namespace: String,
pub actor_id: Option<String>,
pub visible_namespaces: Vec<String>,
pub process_ref: Option<String>,
pub request_id: Option<u64>,
}
impl RequestIdentity {
pub fn from_token(token: &NamespaceToken) -> Self {
Self {
namespace: token.namespace().as_str().to_string(),
actor_id: token.actor().binding_id().map(str::to_string),
visible_namespaces: token
.visible_namespaces()
.iter()
.map(|namespace| namespace.as_str().to_string())
.collect(),
process_ref: token.process_ref().map(str::to_owned),
request_id: None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct VerifiedActor(String);
impl VerifiedActor {
pub fn new(id: impl Into<String>) -> Result<Self, RuntimeError> {
let id = id.into();
if id.trim().is_empty() {
return Err(RuntimeError::InvalidInput(
"VerifiedActor: identifier must not be empty or whitespace-only".to_string(),
));
}
Ok(Self(id))
}
pub fn as_str(&self) -> &str {
&self.0
}
pub(super) fn into_inner(self) -> String {
self.0
}
}
#[derive(Debug)]
pub struct PackSchemaCollisionError {
pub pack_a: &'static str,
pub pack_b: &'static str,
pub table: String,
}
impl std::fmt::Display for PackSchemaCollisionError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
if self.pack_a == self.pack_b {
write!(
f,
"pack schema boot failure for pack {:?}: {}",
self.pack_a, self.table
)
} else {
write!(
f,
"pack schema collision: packs {:?} and {:?} both declare table {:?} \
on the same backend — move one pack to a separate backend or rename the table",
self.pack_a, self.pack_b, self.table
)
}
}
}
impl std::error::Error for PackSchemaCollisionError {}
pub(super) fn extract_table_names(stmt: &str) -> Vec<String> {
enum SqlToken {
Bare(String),
Quoted(String),
Punctuation(char),
}
let mut tokens = Vec::new();
let mut chars = stmt.chars().peekable();
while let Some(ch) = chars.next() {
if ch.is_whitespace() {
continue;
}
if ch == '-' && chars.peek() == Some(&'-') {
chars.next();
for next in chars.by_ref() {
if next == '\n' {
break;
}
}
continue;
}
if ch == '/' && chars.peek() == Some(&'*') {
chars.next();
let mut previous = '\0';
for next in chars.by_ref() {
if previous == '*' && next == '/' {
break;
}
previous = next;
}
continue;
}
if matches!(ch, '"' | '`' | '[' | '\'') {
let closing = if ch == '[' { ']' } else { ch };
let mut token = String::new();
while let Some(next) = chars.next() {
if next == closing {
if chars.peek() == Some(&closing) {
chars.next();
token.push(closing);
} else {
break;
}
} else {
token.push(next);
}
}
tokens.push(SqlToken::Quoted(token));
continue;
}
if matches!(ch, '.' | '(' | ';') {
tokens.push(SqlToken::Punctuation(ch));
continue;
}
let mut token = ch.to_string();
while let Some(next) = chars.peek().copied() {
let begins_comment = (next == '-' && chars.clone().nth(1) == Some('-'))
|| (next == '/' && chars.clone().nth(1) == Some('*'));
if next.is_whitespace()
|| matches!(next, '.' | '(' | ';' | '"' | '`' | '[' | '\'')
|| begins_comment
{
break;
}
token.push(next);
chars.next();
}
tokens.push(SqlToken::Bare(token));
}
let keyword = |index: usize, word: &str| matches!(tokens.get(index), Some(SqlToken::Bare(token)) if token.eq_ignore_ascii_case(word));
if !keyword(0, "CREATE") {
return Vec::new();
}
let mut index = 1;
if keyword(index, "TEMP") || keyword(index, "TEMPORARY") {
index += 1;
}
if keyword(index, "VIRTUAL") {
index += 1;
}
if !keyword(index, "TABLE") {
return Vec::new();
}
index += 1;
if keyword(index, "IF") && keyword(index + 1, "NOT") && keyword(index + 2, "EXISTS") {
index += 3;
}
let main_qualifier = matches!(
tokens.get(index),
Some(SqlToken::Bare(name) | SqlToken::Quoted(name)) if name.eq_ignore_ascii_case("main")
);
if main_qualifier && matches!(tokens.get(index + 1), Some(SqlToken::Punctuation('.'))) {
index += 2;
}
match tokens.get(index) {
Some(SqlToken::Bare(name) | SqlToken::Quoted(name)) if !name.is_empty() => {
vec![name.to_ascii_lowercase()]
}
_ => Vec::new(),
}
}
fn endpoint_kind_label(kind: &EndpointKind) -> String {
match kind {
EndpointKind::EntityOfKind(k) => format!("entity:{k}"),
EndpointKind::NoteOfKind(k) => format!("note:{k}"),
EndpointKind::EntityOfType { kind, entity_type } => {
format!("entity:{kind}({entity_type})")
}
}
}
pub(crate) fn is_special_relation(relation: khive_types::EdgeRelation) -> bool {
SPECIAL_RELATIONS.contains(&relation)
}
pub(super) fn edge_endpoint_table(packs: &[Box<dyn PackRuntime>]) -> Vec<Value> {
let mut rows: Vec<Value> = crate::operations::base_entity_endpoint_rules()
.iter()
.map(|(src, rel, tgt)| {
serde_json::json!({
"relation": rel.as_str(),
"source": format!("entity:{src}"),
"target": format!("entity:{tgt}"),
})
})
.collect();
for rel in SPECIAL_RELATIONS {
rows.push(serde_json::json!({
"relation": rel.as_str(),
"source": "note:*",
"target": "note:*",
}));
}
for pack in packs.iter() {
for rule in pack.edge_rules().iter() {
if is_special_relation(rule.relation) {
continue;
}
rows.push(serde_json::json!({
"relation": rule.relation.as_str(),
"source": endpoint_kind_label(&rule.source),
"target": endpoint_kind_label(&rule.target),
}));
}
}
rows.push(serde_json::json!({
"relation": "annotates",
"source": "note:*",
"target": "any (entity, note, edge, or event)",
}));
rows
}