use super::{
Arc, DispatchHook, HashMap, KhiveRuntime, PackByIdResolver, PackRuntime, RuntimeError,
VerbRegistry, VerbRegistryBuilder, Visibility, CHANNEL_INGEST_CAPABLE_PACKS,
};
pub struct PackInstall {
pub runtime: Box<dyn PackRuntime>,
pub resolver: Option<Box<dyn PackByIdResolver>>,
pub dispatch_hook: Option<Arc<dyn DispatchHook>>,
}
pub struct ChannelIngestCapability {
pub(crate) _sealed: (),
}
impl ChannelIngestCapability {
pub fn grant_for_direct_composition() -> Self {
Self { _sealed: () }
}
}
pub trait PackFactory: Send + Sync + 'static {
fn name(&self) -> &'static str;
fn requires(&self) -> &'static [&'static str] {
&[]
}
fn intentionally_verbless(&self) -> bool {
false
}
fn create(&self, runtime: KhiveRuntime) -> Box<dyn PackRuntime>;
fn create_install(&self, runtime: KhiveRuntime) -> PackInstall {
let resolver = self.create_resolver(runtime.clone());
PackInstall {
runtime: self.create(runtime),
resolver,
dispatch_hook: None,
}
}
fn create_resolver(&self, _runtime: KhiveRuntime) -> Option<Box<dyn PackByIdResolver>> {
None
}
}
pub struct PackRegistration(pub &'static dyn PackFactory);
inventory::collect!(PackRegistration);
#[derive(Debug)]
pub enum PackLoadError {
UnknownPack(String),
DuplicatePack(String),
MissingDependency {
pack: String,
dep: String,
},
NoPublicVerbs {
pack: String,
},
}
impl std::fmt::Display for PackLoadError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
PackLoadError::UnknownPack(name) => write!(f, "unknown pack {name:?}"),
PackLoadError::DuplicatePack(name) => write!(f, "duplicate pack {name:?}"),
PackLoadError::MissingDependency { pack, dep } => write!(
f,
"pack {pack:?} requires {dep:?}, which is not in the requested pack list; \
add --pack {dep} before --pack {pack}"
),
PackLoadError::NoPublicVerbs { pack } => write!(
f,
"declared pack {pack:?} registers no public verbs; if this pack is \
intentionally vocabulary- or ontology-only, its factory must declare \
intentionally_verbless() = true"
),
}
}
}
impl std::error::Error for PackLoadError {}
fn check_pack_has_public_verbs(
factory: &dyn PackFactory,
install: &PackInstall,
name: &str,
) -> Result<(), PackLoadError> {
if !factory.intentionally_verbless()
&& !install
.runtime
.handlers()
.iter()
.any(|handler| matches!(handler.visibility, Visibility::Verb))
{
return Err(PackLoadError::NoPublicVerbs {
pack: name.to_string(),
});
}
Ok(())
}
pub struct PackRegistry;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum IngestAuditStore {
Attach,
Detach,
}
impl PackRegistry {
pub fn discovered_names() -> Vec<&'static str> {
inventory::iter::<PackRegistration>
.into_iter()
.map(|r| r.0.name())
.collect()
}
pub fn validate_pack_selection(names: &[String]) -> Result<(), PackLoadError> {
let all: Vec<&'static dyn PackFactory> = inventory::iter::<PackRegistration>
.into_iter()
.map(|r| r.0)
.collect();
Self::validate_pack_selection_from(&all, names)
}
fn validate_pack_selection_from(
factories: &[&'static dyn PackFactory],
names: &[String],
) -> Result<(), PackLoadError> {
let factory_for = |name: &str| factories.iter().copied().find(|f| f.name() == name);
let mut requested = std::collections::HashSet::new();
for name in names {
factory_for(name).ok_or_else(|| PackLoadError::UnknownPack(name.clone()))?;
if !requested.insert(name.as_str()) {
return Err(PackLoadError::DuplicatePack(name.clone()));
}
}
for name in names {
let factory = factory_for(name).unwrap(); for &dep in factory.requires() {
if !requested.contains(dep) {
return Err(PackLoadError::MissingDependency {
pack: name.clone(),
dep: dep.to_string(),
});
}
}
}
Ok(())
}
pub fn register_packs(
names: &[String],
runtime: KhiveRuntime,
builder: &mut VerbRegistryBuilder,
) -> Result<(), PackLoadError> {
let all: Vec<&'static dyn PackFactory> = inventory::iter::<PackRegistration>
.into_iter()
.map(|r| r.0)
.collect();
let factory_for = |name: &str| -> Option<&'static dyn PackFactory> {
all.iter().copied().find(|f| f.name() == name)
};
Self::validate_pack_selection_from(&all, names)?;
for name in names {
let factory = factory_for(name.as_str()).unwrap(); let install = factory.create_install(runtime.clone());
check_pack_has_public_verbs(factory, &install, name)?;
if CHANNEL_INGEST_CAPABLE_PACKS.contains(&name.as_str()) {
install
.runtime
.accept_channel_ingest_capability(ChannelIngestCapability { _sealed: () });
}
builder.register_boxed(install.runtime);
if let Some(resolver) = install.resolver {
builder.register_resolver(name.clone(), resolver);
}
if let Some(hook) = install.dispatch_hook {
builder.with_dispatch_hook(hook);
}
}
Ok(())
}
pub fn build_ingest_registry(
runtime: &KhiveRuntime,
audit_store: IngestAuditStore,
) -> Result<VerbRegistry, RuntimeError> {
let mut builder = VerbRegistryBuilder::new();
builder.with_gate(runtime.config().gate.clone());
builder.with_default_namespace(runtime.config().default_namespace.as_str());
builder.with_visible_namespaces(runtime.config().visible_namespaces.clone());
builder.with_actor_id(runtime.config().actor_id.clone());
if audit_store == IngestAuditStore::Attach {
if runtime.is_read_only() {
builder.with_read_only_audit_store();
} else {
builder.with_runtime_event_store(runtime)?;
}
}
Self::register_packs(
&runtime.config().packs.clone(),
runtime.clone(),
&mut builder,
)
.map_err(|e| RuntimeError::Internal(format!("pack registration failed: {e:?}")))?;
let registry = builder.build()?;
runtime.install_edge_rules(registry.all_edge_rules());
Ok(registry)
}
pub fn register_packs_with_runtimes(
names: &[String],
runtimes: &HashMap<String, KhiveRuntime>,
default_runtime: &KhiveRuntime,
builder: &mut VerbRegistryBuilder,
) -> Result<(), PackLoadError> {
let all: Vec<&'static dyn PackFactory> = inventory::iter::<PackRegistration>
.into_iter()
.map(|r| r.0)
.collect();
Self::register_packs_with_runtimes_from(&all, names, runtimes, default_runtime, builder)
}
pub fn register_packs_with_runtimes_with_extra_factories(
extra_factories: &[&'static dyn PackFactory],
names: &[String],
runtimes: &HashMap<String, KhiveRuntime>,
default_runtime: &KhiveRuntime,
builder: &mut VerbRegistryBuilder,
) -> Result<(), PackLoadError> {
let mut all: Vec<&'static dyn PackFactory> = inventory::iter::<PackRegistration>
.into_iter()
.map(|r| r.0)
.collect();
all.extend(extra_factories.iter().copied());
Self::register_packs_with_runtimes_from(&all, names, runtimes, default_runtime, builder)
}
fn register_packs_with_runtimes_from(
factories: &[&'static dyn PackFactory],
names: &[String],
runtimes: &HashMap<String, KhiveRuntime>,
default_runtime: &KhiveRuntime,
builder: &mut VerbRegistryBuilder,
) -> Result<(), PackLoadError> {
let factory_for = |name: &str| -> Option<&'static dyn PackFactory> {
factories.iter().copied().find(|f| f.name() == name)
};
Self::validate_pack_selection_from(factories, names)?;
builder.kg_read_resolver = Some(Arc::new(crate::kg_read::KgReadResolver::new(
default_runtime,
runtimes,
)));
for name in names {
let factory = factory_for(name.as_str()).unwrap();
let runtime = runtimes
.get(name.as_str())
.cloned()
.unwrap_or_else(|| default_runtime.clone());
let install = factory.create_install(runtime);
check_pack_has_public_verbs(factory, &install, name)?;
if CHANNEL_INGEST_CAPABLE_PACKS.contains(&name.as_str()) {
install
.runtime
.accept_channel_ingest_capability(ChannelIngestCapability { _sealed: () });
}
builder.register_boxed(install.runtime);
if let Some(resolver) = install.resolver {
builder.register_resolver(name.clone(), resolver);
}
if let Some(hook) = install.dispatch_hook {
builder.with_dispatch_hook(hook);
}
}
Ok(())
}
}