use std::collections::{HashMap, HashSet, VecDeque};
use std::sync::Arc;
#[cfg(doc)]
use khive_gate::GateRequest;
use khive_gate::{AllowAllGate, GateRef};
use khive_storage::EventStore;
#[cfg(doc)]
use khive_storage::EventView;
use khive_types::Namespace;
use serde_json::Value;
use crate::error::{
CircularPackDependency, MissingPackDependencies, MissingPackDependency, RuntimeError,
};
use crate::KhiveRuntime;
use super::{
DispatchHook, EdgeEndpointRule, EntityTypeDef, HandlerDef, PackByIdResolver, PackRuntime,
VerbCategory, VerbRegistry, Visibility, RESERVED_ENVELOPE_ARGS,
};
#[cfg(doc)]
use super::{PackFactory, PackRegistry};
pub struct VerbRegistryBuilder {
packs: Vec<Box<dyn PackRuntime>>,
pack_trusted: Vec<bool>,
resolvers: Vec<(String, Box<dyn PackByIdResolver>)>,
pub(super) kg_read_resolver: Option<Arc<crate::kg_read::KgReadResolver>>,
gate: GateRef,
default_namespace: String,
visible_namespaces: Vec<Namespace>,
actor_id: Option<String>,
event_store: Option<Arc<dyn EventStore>>,
runtime_event_store: Option<KhiveRuntime>,
audit_store_read_only: bool,
dispatch_hook: Option<Arc<dyn DispatchHook>>,
audit_batch_config: Option<crate::audit_batch::AuditBatchConfig>,
}
impl VerbRegistryBuilder {
pub fn new() -> Self {
Self {
packs: Vec::new(),
pack_trusted: Vec::new(),
resolvers: Vec::new(),
kg_read_resolver: None,
gate: std::sync::Arc::new(AllowAllGate),
default_namespace: Namespace::local().as_str().to_string(),
visible_namespaces: vec![],
actor_id: None,
event_store: None,
runtime_event_store: None,
audit_store_read_only: false,
dispatch_hook: None,
audit_batch_config: None,
}
}
pub fn with_visible_namespaces(&mut self, ns: Vec<Namespace>) -> &mut Self {
self.visible_namespaces = ns;
self
}
pub fn with_actor_id(&mut self, actor_id: Option<String>) -> &mut Self {
self.actor_id = actor_id;
self
}
pub fn register<P: khive_types::Pack + PackRuntime + 'static>(&mut self, pack: P) -> &mut Self {
self.packs.push(Box::new(pack));
self.pack_trusted.push(false);
self
}
pub(crate) fn register_boxed(&mut self, pack: Box<dyn PackRuntime>) -> &mut Self {
self.packs.push(pack);
self.pack_trusted.push(true);
self
}
pub fn register_mounted(
&mut self,
pack: Box<dyn PackRuntime>,
) -> Result<&mut Self, RuntimeError> {
if pack.mounted_namespace() != Some(pack.name()) || !pack.handlers().is_empty() {
return Err(RuntimeError::InvalidInput(
"invalid mounted namespace registration".into(),
));
}
self.packs.push(pack);
self.pack_trusted.push(false);
Ok(self)
}
#[cfg(any(test, feature = "test-internals"))]
pub fn register_trusted<P: khive_types::Pack + PackRuntime + 'static>(
&mut self,
pack: P,
) -> &mut Self {
self.packs.push(Box::new(pack));
self.pack_trusted.push(true);
self
}
pub fn register_resolver(
&mut self,
name: impl Into<String>,
resolver: Box<dyn PackByIdResolver>,
) -> &mut Self {
self.resolvers.push((name.into(), resolver));
self
}
pub fn with_gate(&mut self, gate: GateRef) -> &mut Self {
self.gate = gate;
self
}
pub fn with_default_namespace(&mut self, ns: impl Into<String>) -> &mut Self {
self.default_namespace = ns.into();
self
}
pub fn with_event_store(&mut self, store: Arc<dyn EventStore>) -> &mut Self {
self.event_store = Some(store);
self.runtime_event_store = None;
self.audit_store_read_only = false;
self
}
pub fn with_runtime_event_store(
&mut self,
runtime: &KhiveRuntime,
) -> Result<&mut Self, RuntimeError> {
self.event_store = None;
self.runtime_event_store = Some(runtime.clone());
self.audit_store_read_only = false;
Ok(self)
}
pub fn with_audit_batch_config(
&mut self,
config: crate::audit_batch::AuditBatchConfig,
) -> &mut Self {
self.audit_batch_config = Some(config);
self
}
pub fn with_read_only_audit_store(&mut self) -> &mut Self {
self.event_store = None;
self.runtime_event_store = None;
self.audit_store_read_only = true;
self
}
pub fn with_dispatch_hook(&mut self, hook: Arc<dyn DispatchHook>) -> &mut Self {
self.dispatch_hook = Some(hook);
self
}
pub fn build(self) -> Result<VerbRegistry, RuntimeError> {
self.build_registry(true)
}
pub fn build_metadata(mut self) -> Result<PackMetadataRegistry, RuntimeError> {
self.event_store = None;
self.runtime_event_store = None;
self.dispatch_hook = None;
self.resolvers.clear();
self.build_registry(false)
.map(|registry| PackMetadataRegistry { registry })
}
fn build_registry(self, activate: bool) -> Result<VerbRegistry, RuntimeError> {
let packs = self.packs;
let mut name_to_idx: HashMap<&str, usize> = HashMap::with_capacity(packs.len());
for (idx, pack) in packs.iter().enumerate() {
if let Some(prev_idx) = name_to_idx.insert(pack.name(), idx) {
return Err(RuntimeError::PackRedeclared {
name: pack.name().to_string(),
first_idx: prev_idx,
second_idx: idx,
});
}
}
for mounted in packs
.iter()
.filter(|pack| pack.mounted_namespace().is_some())
{
let prefix = format!("{}.", mounted.name());
if packs
.iter()
.flat_map(|pack| pack.handlers())
.any(|handler| handler.name.starts_with(&prefix))
{
return Err(RuntimeError::InvalidInput(
"mounted namespace collides with a native verb".into(),
));
}
}
for pack in &packs {
for handler in pack.handlers() {
for parameter in handler.params {
if RESERVED_ENVELOPE_ARGS.contains(¶meter.name) {
return Err(RuntimeError::ReservedEnvelopeParam {
pack: pack.name().to_string(),
verb: handler.name.to_string(),
param: parameter.name.to_string(),
});
}
}
}
}
let mut missing: Vec<MissingPackDependency> = Vec::new();
let mut indegree = vec![0usize; packs.len()];
let mut dependents: Vec<Vec<usize>> = vec![Vec::new(); packs.len()];
for (idx, pack) in packs.iter().enumerate() {
for &requires in pack.requires() {
match name_to_idx.get(requires).copied() {
Some(dep_idx) => {
dependents[dep_idx].push(idx);
indegree[idx] += 1;
}
None => missing.push(MissingPackDependency {
from: pack.name().to_string(),
requires: requires.to_string(),
}),
}
}
}
if !missing.is_empty() {
return if missing.len() == 1 {
Err(RuntimeError::MissingPackDependency(missing.remove(0)))
} else {
Err(RuntimeError::MissingPackDependencies(
MissingPackDependencies { missing },
))
};
}
let mut ready: VecDeque<usize> = indegree
.iter()
.enumerate()
.filter_map(|(idx, degree)| (*degree == 0).then_some(idx))
.collect();
let mut ordered_indices = Vec::with_capacity(packs.len());
while let Some(idx) = ready.pop_front() {
ordered_indices.push(idx);
for &dep_idx in &dependents[idx] {
indegree[dep_idx] -= 1;
if indegree[dep_idx] == 0 {
ready.push_back(dep_idx);
}
}
}
if ordered_indices.len() != packs.len() {
let cycle_nodes: HashSet<usize> = indegree
.iter()
.enumerate()
.filter_map(|(idx, degree)| (*degree > 0).then_some(idx))
.collect();
let cycle = find_pack_dependency_cycle(&packs, &name_to_idx, &cycle_nodes);
return Err(RuntimeError::CircularPackDependency(
CircularPackDependency { cycle },
));
}
let mut pack_slots: Vec<Option<Box<dyn PackRuntime>>> =
packs.into_iter().map(Some).collect();
let mut trusted_slots: Vec<Option<bool>> =
self.pack_trusted.into_iter().map(Some).collect();
let mut ordered_packs: Vec<Box<dyn PackRuntime>> = Vec::with_capacity(pack_slots.len());
let mut ordered_trusted: Vec<bool> = Vec::with_capacity(trusted_slots.len());
for idx in ordered_indices {
ordered_packs.push(
pack_slots[idx]
.take()
.expect("topological index must exist"),
);
ordered_trusted.push(
trusted_slots[idx]
.take()
.expect("topological index must exist"),
);
}
validate_unique_note_kinds(&ordered_packs)?;
validate_unique_verb_names(&ordered_packs)?;
validate_unique_entity_types(&ordered_packs)?;
validate_entity_type_note_kind_collisions(&ordered_packs)?;
validate_brain_consumer_kinds(&ordered_packs)?;
if activate {
for pack in &ordered_packs {
pack.validate_config()?;
}
}
let available_verbs: Vec<&'static str> = ordered_packs
.iter()
.flat_map(|p| p.handlers().iter())
.filter(|h| matches!(h.visibility, Visibility::Verb))
.map(|h| h.name)
.collect();
let mut handler_by_name: HashMap<&'static str, &'static HandlerDef> = HashMap::new();
for pack in &ordered_packs {
for handler in pack.handlers() {
handler_by_name.entry(handler.name).or_insert(handler);
}
}
let mut degrade_safe_verbs: HashSet<&'static str> = HashSet::new();
let mut read_replay_safe_verbs = HashSet::new();
for (pack, &trusted) in ordered_packs.iter().zip(ordered_trusted.iter()) {
if !trusted {
continue;
}
let pack_name = pack.name();
for handler in pack.handlers() {
let canonical_owner = handler
.name
.split_once('.')
.map_or("kg", |(owner, _)| owner);
if matches!(handler.visibility, Visibility::Verb)
&& pack_name == canonical_owner
&& crate::classify_operation(handler.name) == Some(crate::OperationAccess::Read)
&& !VerbRegistry::SIDE_EFFECTING_ASSERTIVE_VERBS.contains(&handler.name)
{
read_replay_safe_verbs.insert(handler.name);
}
if !matches!(handler.visibility, Visibility::Verb)
|| handler.category != VerbCategory::Assertive
{
continue;
}
let eligible = VerbRegistry::admission_degrade_safe_sorted()
.binary_search_by(|&(p, v)| p.cmp(pack_name).then_with(|| v.cmp(handler.name)))
.is_ok();
if eligible {
degrade_safe_verbs.insert(handler.name);
}
}
}
let event_store = match self.runtime_event_store {
Some(runtime) => Some(runtime.raw_events_for_namespace(&self.default_namespace)?),
None => self.event_store,
};
if let Some(store) = &event_store {
if !store.supports_idempotent_audit_batch() {
return Err(RuntimeError::IncompatibleEventStore(
"the configured EventStore does not implement ADR-133's \
preflight_event/append_events_idempotent pair \
(supports_idempotent_audit_batch() returned false); every \
audited dispatch would silently lose its audit row while \
still reporting success. Implement both methods and \
override supports_idempotent_audit_batch() to opt in, or \
do not call with_event_store() for this backend."
.to_string(),
));
}
}
let audit_batch = event_store.clone().map(|store| {
crate::audit_batch::AuditBatch::new(
store,
self.audit_batch_config.clone().unwrap_or_default(),
)
});
Ok(VerbRegistry {
packs: Arc::new(ordered_packs),
resolvers: Arc::new(self.resolvers),
kg_read_resolver: self.kg_read_resolver,
gate: self.gate,
default_namespace: self.default_namespace,
visible_namespaces: self.visible_namespaces,
actor_id: self.actor_id,
event_store,
audit_store_read_only: self.audit_store_read_only,
dispatch_hook: self.dispatch_hook,
available_verbs: Arc::new(available_verbs),
handler_by_name: Arc::new(handler_by_name),
degrade_safe_verbs: Arc::new(degrade_safe_verbs),
read_replay_safe_verbs: Arc::new(read_replay_safe_verbs),
reference_ring: Arc::new(crate::reference_ring::ReferenceRing::new()),
audit_batch,
})
}
}
fn validate_unique_note_kinds(packs: &[Box<dyn PackRuntime>]) -> Result<(), RuntimeError> {
let mut seen: HashMap<&str, &str> = HashMap::new();
for pack in packs {
for &kind in pack.note_kinds() {
if let Some(first_pack) = seen.insert(kind, pack.name()) {
return Err(RuntimeError::InvalidInput(format!(
"duplicate note kind {kind:?}: claimed by both {first_pack:?} and {:?}",
pack.name()
)));
}
}
}
Ok(())
}
fn validate_brain_consumer_kinds(packs: &[Box<dyn PackRuntime>]) -> Result<(), RuntimeError> {
for pack in packs {
for &kind in pack.brain_consumer_kinds() {
if kind == "*" || kind.trim().is_empty() || kind.trim() != kind {
return Err(RuntimeError::InvalidInput(format!(
"pack {:?} declares invalid brain consumer kind {kind:?}; declarations must be non-empty exact wire values and must not use the registry-owned \"*\" wildcard",
pack.name()
)));
}
}
}
Ok(())
}
fn validate_unique_verb_names(packs: &[Box<dyn PackRuntime>]) -> Result<(), RuntimeError> {
let mut seen: HashMap<&str, &str> = HashMap::new();
for pack in packs {
for handler in pack.handlers() {
if !matches!(handler.visibility, Visibility::Verb) {
continue;
}
if let Some(first_pack) = seen.insert(handler.name, pack.name()) {
return Err(RuntimeError::VerbCollision {
verb: handler.name.to_string(),
first_pack: first_pack.to_string(),
second_pack: pack.name().to_string(),
});
}
}
}
Ok(())
}
fn validate_unique_entity_types(packs: &[Box<dyn PackRuntime>]) -> Result<(), RuntimeError> {
let owned_defs = packs
.iter()
.flat_map(|p| p.entity_types().iter().map(move |def| (p.name(), def)));
khive_types::EntityTypeRegistry::check_extra_collisions(owned_defs)
.map_err(RuntimeError::InvalidInput)
}
fn validate_entity_type_note_kind_collisions(
packs: &[Box<dyn PackRuntime>],
) -> Result<(), RuntimeError> {
let mut note_kinds = HashMap::new();
for pack in packs {
for &kind in pack.note_kinds() {
note_kinds
.entry(khive_types::to_snake_case(kind))
.or_insert(pack.name());
}
}
let check_definition = |definition: &EntityTypeDef, owner: &str| {
for name in std::iter::once(definition.type_name).chain(definition.aliases.iter().copied())
{
let normalized = khive_types::to_snake_case(name);
if let Some(note_owner) = note_kinds.get(&normalized) {
return Err(RuntimeError::InvalidInput(format!(
"entity subtype {name:?} from {owner:?} collides with note kind {normalized:?} from pack {note_owner:?}"
)));
}
}
Ok(())
};
let builtin = khive_types::EntityTypeRegistry::builtin();
for definition in builtin.definitions() {
check_definition(definition, "builtin")?;
}
for pack in packs {
for definition in pack.entity_types() {
check_definition(definition, pack.name())?;
}
}
Ok(())
}
fn find_pack_dependency_cycle(
packs: &[Box<dyn PackRuntime>],
name_to_idx: &HashMap<&str, usize>,
cycle_nodes: &HashSet<usize>,
) -> Vec<String> {
fn visit(
idx: usize,
packs: &[Box<dyn PackRuntime>],
name_to_idx: &HashMap<&str, usize>,
cycle_nodes: &HashSet<usize>,
visiting: &mut Vec<usize>,
visited: &mut HashSet<usize>,
) -> Option<Vec<String>> {
if let Some(pos) = visiting.iter().position(|&seen| seen == idx) {
let mut cycle: Vec<String> = visiting[pos..]
.iter()
.map(|&i| packs[i].name().to_string())
.collect();
cycle.push(packs[idx].name().to_string());
return Some(cycle);
}
if !visited.insert(idx) {
return None;
}
visiting.push(idx);
for &req in packs[idx].requires() {
let Some(&dep_idx) = name_to_idx.get(req) else {
continue;
};
if cycle_nodes.contains(&dep_idx) {
if let Some(cycle) =
visit(dep_idx, packs, name_to_idx, cycle_nodes, visiting, visited)
{
return Some(cycle);
}
}
}
visiting.pop();
None
}
let mut visited = HashSet::new();
for &idx in cycle_nodes {
let mut visiting = Vec::new();
if let Some(cycle) = visit(
idx,
packs,
name_to_idx,
cycle_nodes,
&mut visiting,
&mut visited,
) {
return cycle;
}
}
cycle_nodes
.iter()
.map(|&idx| packs[idx].name().to_string())
.collect()
}
impl Default for VerbRegistryBuilder {
fn default() -> Self {
Self::new()
}
}
pub struct PackMetadataRegistry {
pub(super) registry: VerbRegistry,
}
impl PackMetadataRegistry {
pub fn has_verb(&self, verb: &str) -> bool {
self.registry.has_verb(verb)
}
pub fn describe_verb(&self, verb: &str) -> Result<Value, RuntimeError> {
self.registry.describe_verb(verb)
}
pub fn all_handlers_with_names(&self) -> Vec<(&str, &'static HandlerDef)> {
self.registry.all_handlers_with_names()
}
pub fn all_verbs(&self) -> Vec<&'static HandlerDef> {
self.registry.all_verbs()
}
pub fn pack_names(&self) -> Vec<&str> {
self.registry.pack_names()
}
pub fn pack_requires(&self, name: &str) -> Option<&'static [&'static str]> {
self.registry.pack_requires(name)
}
pub fn pack_note_kinds(&self, name: &str) -> Option<&'static [&'static str]> {
self.registry.pack_note_kinds(name)
}
pub fn pack_entity_kinds(&self, name: &str) -> Option<&'static [&'static str]> {
self.registry.pack_entity_kinds(name)
}
pub fn pack_verbs(&self, name: &str) -> Option<&'static [HandlerDef]> {
self.registry.pack_verbs(name)
}
pub fn all_entity_kinds(&self) -> Vec<&'static str> {
self.registry.all_entity_kinds()
}
pub fn all_note_kinds(&self) -> Vec<&'static str> {
self.registry.all_note_kinds()
}
pub fn all_edge_rules(&self) -> Vec<EdgeEndpointRule> {
self.registry.all_edge_rules()
}
}