#![allow(unexpected_cfgs)]
use std::path::Path;
use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use chrono::Utc;
use serde::Serialize;
use cloacina_constructor_contract::{
AccumulatorInvocation, AccumulatorOutcome, ConstructorManifest, PollOutcome, PrimitiveKind,
ProviderManifest, ProviderRuntime, ReactorInvocation, ReactorOutcome, TaskInvocation,
TaskOutcome, TriggerInvocation, ACCUMULATOR_CONSTRUCTOR_INTERFACE_VERSION, METHOD_EVALUATE,
METHOD_EXECUTE, METHOD_INGEST, METHOD_POLL, METHOD_SOURCE,
REACTOR_CONSTRUCTOR_INTERFACE_VERSION, STREAM_ACCUMULATOR_CONSTRUCTOR_INTERFACE_VERSION,
TASK_CONSTRUCTOR_INTERFACE_VERSION, TRIGGER_CONSTRUCTOR_INTERFACE_VERSION,
};
use fidius_host::PluginHost;
use super::grants::{lint_unmet_intents, translate, GrantSpec, ResolvedGrants};
use crate::computation_graph::accumulator::{Accumulator, AccumulatorError, EventSource};
use crate::computation_graph::reactor::ReactorFireDecider;
use crate::computation_graph::types::InputCache;
use crate::context::Context;
use crate::error::TaskError;
use crate::registry::error::LoaderError;
use crate::runtime::Runtime;
use crate::task::{Task, TaskNamespace};
use crate::trigger::{Trigger, TriggerResult};
use cloacina_workflow::TriggerError;
#[fidius_macro::plugin_interface(version = 1, buffer = PluginAllocated, crate = "fidius_core")]
pub trait TaskConstructor: Send + Sync {
fn execute(&self, invocation_json: String) -> String;
}
#[fidius_macro::plugin_interface(version = 1, buffer = PluginAllocated, crate = "fidius_core")]
pub trait TriggerConstructor: Send + Sync {
fn poll(&self, invocation_json: String) -> String;
}
pub const CONSTRUCTOR_MANIFEST_FILE: &str = "constructor.json";
pub const PROVIDER_MANIFEST_FILE: &str = "provider.json";
#[derive(serde::Serialize)]
struct ProviderConfigure {
name: String,
config: Vec<u8>,
}
pub fn read_provider_manifest(package_dir: &Path) -> Result<ProviderManifest, LoaderError> {
let path = package_dir.join(PROVIDER_MANIFEST_FILE);
let raw = std::fs::read_to_string(&path).map_err(|e| LoaderError::FileSystem {
path: path.display().to_string(),
error: e.to_string(),
})?;
ProviderManifest::from_json(&raw).map_err(|e| LoaderError::ManifestParse {
reason: format!("{PROVIDER_MANIFEST_FILE}: {e}"),
})
}
fn read_member_manifest(
package_dir: &Path,
constructor_name: &str,
) -> Result<ConstructorManifest, LoaderError> {
let provider = read_provider_manifest(package_dir)?;
provider
.constructor(constructor_name)
.cloned()
.ok_or_else(|| LoaderError::Validation {
reason: format!(
"provider '{}' has no constructor '{}'; members: [{}]",
provider.name,
constructor_name,
provider
.constructors
.iter()
.map(|c| c.name.as_str())
.collect::<Vec<_>>()
.join(", ")
),
})
}
fn wrap_member_config<C: Serialize>(
constructor_name: &str,
config: &C,
) -> Result<ProviderConfigure, LoaderError> {
let inner = fidius_core::wire::serialize(config).map_err(|e| LoaderError::Validation {
reason: format!("serialize config for constructor '{constructor_name}': {e}"),
})?;
Ok(ProviderConfigure {
name: constructor_name.to_string(),
config: inner,
})
}
fn resolve_native_provider(
search_path: &Path,
package_name: &str,
) -> Result<Option<(std::path::PathBuf, ProviderManifest)>, LoaderError> {
if !search_path.is_dir() {
return Ok(None);
}
let entries = std::fs::read_dir(search_path).map_err(|e| LoaderError::FileSystem {
path: search_path.display().to_string(),
error: e.to_string(),
})?;
for entry in entries.flatten() {
let path = entry.path();
if !path.is_dir() {
continue;
}
if let Ok(provider) = read_provider_manifest(&path) {
if provider.name == package_name && provider.runtime == ProviderRuntime::Native {
return Ok(Some((path, provider)));
}
}
}
Ok(None)
}
fn load_native_member<C: Serialize>(
search_path: &Path,
package_name: &str,
constructor_name: &str,
config: &C,
expected_kind: PrimitiveKind,
expected_iface_version: u32,
holder_plugin: &str,
) -> Result<Option<(fidius_host::PluginHandle, ConstructorManifest)>, LoaderError> {
let Some((dir, provider)) = resolve_native_provider(search_path, package_name)? else {
return Ok(None);
};
let member = provider
.constructor(constructor_name)
.cloned()
.ok_or_else(|| LoaderError::Validation {
reason: format!(
"native provider '{}' has no constructor '{}'; members: [{}]",
provider.name,
constructor_name,
provider
.constructors
.iter()
.map(|c| c.name.as_str())
.collect::<Vec<_>>()
.join(", ")
),
})?;
if member.primitive_kind != expected_kind {
return Err(LoaderError::Validation {
reason: format!(
"constructor '{}' is {:?}, not {:?}",
member.name, member.primitive_kind, expected_kind
),
});
}
if member.interface_version != expected_iface_version {
return Err(LoaderError::Validation {
reason: format!(
"native constructor '{}' declares interface v{}, loader supports v{}",
member.name, member.interface_version, expected_iface_version
),
});
}
let configure = wrap_member_config(constructor_name, config)?;
let lib_path = dir.join(&provider.component);
let loaded =
fidius_host::loader::load_library(&lib_path).map_err(|e| LoaderError::LibraryLoad {
path: lib_path.display().to_string(),
error: format!(
"load_library (native provider, host triple '{}'): {e} — if this is a \
wrong-arch cdylib, the per-target compiler has not filled this arch yet",
crate::fleet::protocol::host_target_triple()
),
})?;
let plugin = loaded
.plugins
.into_iter()
.find(|p| p.info.name == holder_plugin)
.ok_or_else(|| LoaderError::Validation {
reason: format!(
"native provider '{}' cdylib '{}' has no '{}' plugin (kind {:?})",
provider.name, provider.component, holder_plugin, expected_kind
),
})?;
let handle =
fidius_host::PluginHandle::configure_from_loaded(plugin, &configure).map_err(|e| {
LoaderError::LibraryLoad {
path: lib_path.display().to_string(),
error: format!("configure_from_loaded (native provider): {e}"),
}
})?;
tracing::info!(
provider = %provider.name,
constructor = %constructor_name,
"loaded NATIVE constructor provider (trusted, unsandboxed in-process; grants advisory)"
);
Ok(Some((handle, member)))
}
pub struct WasmTaskConstructor {
id: String,
handle: Arc<fidius_host::PluginHandle>,
dependencies: Vec<TaskNamespace>,
}
impl WasmTaskConstructor {
pub fn name(&self) -> &str {
&self.id
}
}
impl std::fmt::Debug for WasmTaskConstructor {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("WasmTaskConstructor")
.field("id", &self.id)
.field("dependencies", &self.dependencies)
.finish_non_exhaustive()
}
}
#[async_trait]
impl Task for WasmTaskConstructor {
async fn execute(
&self,
context: Context<serde_json::Value>,
) -> Result<Context<serde_json::Value>, TaskError> {
let id = self.id.clone();
let context_json = context
.to_json()
.map_err(|e| exec_err(&id, format!("serialize context: {e}")))?;
let inv_json = serde_json::to_string(&TaskInvocation { context_json })
.map_err(|e| exec_err(&id, format!("serialize invocation: {e}")))?;
let handle = self.handle.clone();
let call_id = id.clone();
let out_json: String = tokio::task::spawn_blocking(move || {
handle.call_method::<_, String>(METHOD_EXECUTE, &(inv_json,))
})
.await
.map_err(|e| exec_err(&call_id, format!("constructor task join: {e}")))?
.map_err(|e| exec_err(&call_id, format!("constructor FFI call: {e}")))?;
let outcome: TaskOutcome = serde_json::from_str(&out_json)
.map_err(|e| exec_err(&id, format!("parse outcome: {e}")))?;
if outcome.success {
let updated = outcome
.context_json
.ok_or_else(|| exec_err(&id, "successful outcome missing context_json"))?;
Context::from_json(updated).map_err(|e| exec_err(&id, format!("rebuild context: {e}")))
} else {
Err(exec_err(
&id,
outcome
.error
.unwrap_or_else(|| "constructor reported failure with no message".to_string()),
))
}
}
fn id(&self) -> &str {
&self.id
}
fn dependencies(&self) -> &[TaskNamespace] {
&self.dependencies
}
}
fn exec_err(task_id: &str, message: impl Into<String>) -> TaskError {
TaskError::ExecutionFailed {
message: message.into(),
task_id: task_id.to_string(),
timestamp: Utc::now(),
}
}
pub fn read_constructor_manifest(package_dir: &Path) -> Result<ConstructorManifest, LoaderError> {
let path = package_dir.join(CONSTRUCTOR_MANIFEST_FILE);
let raw = std::fs::read_to_string(&path).map_err(|e| LoaderError::FileSystem {
path: path.display().to_string(),
error: e.to_string(),
})?;
ConstructorManifest::from_json(&raw).map_err(|e| LoaderError::ManifestParse {
reason: format!("{CONSTRUCTOR_MANIFEST_FILE}: {e}"),
})
}
pub fn load_task_constructor<C: Serialize>(
search_path: impl AsRef<Path>,
package_name: &str,
constructor_name: &str,
config: &C,
grants: &ResolvedGrants,
) -> Result<Arc<dyn Task>, LoaderError> {
let search_path = search_path.as_ref();
if let Some((handle, member)) = load_native_member(
search_path,
package_name,
constructor_name,
config,
PrimitiveKind::Task,
TASK_CONSTRUCTOR_INTERFACE_VERSION,
"__ProviderTask",
)? {
return Ok(Arc::new(WasmTaskConstructor {
id: member.name,
handle: Arc::new(handle),
dependencies: Vec::new(),
}));
}
let host = PluginHost::builder()
.search_path(search_path)
.build()
.map_err(|e| LoaderError::LibraryLoad {
path: search_path.display().to_string(),
error: format!("build plugin host: {e}"),
})?;
let dir = host
.find_wasm_package(package_name)
.map_err(|e| LoaderError::Validation {
reason: format!("locate wasm constructor package '{package_name}': {e}"),
})?;
let manifest = read_member_manifest(&dir, constructor_name)?;
if manifest.primitive_kind != PrimitiveKind::Task {
return Err(LoaderError::Validation {
reason: format!(
"constructor '{}' is {:?}, not Task; trigger/accumulator/reactor loading is CLOACI-T-0824",
manifest.name, manifest.primitive_kind
),
});
}
if manifest.interface_version != TASK_CONSTRUCTOR_INTERFACE_VERSION {
return Err(LoaderError::Validation {
reason: format!(
"constructor '{}' declares task-constructor interface v{}, loader supports v{}",
manifest.name, manifest.interface_version, TASK_CONSTRUCTOR_INTERFACE_VERSION
),
});
}
let configure = wrap_member_config(constructor_name, config)?;
let handle = host
.load_wasm_configured_with_grants(
package_name,
&__fidius_TaskConstructor::TaskConstructor_WASM_DESCRIPTOR,
&configure,
grants.capabilities.clone(),
grants.egress.clone(),
)
.map_err(|e| LoaderError::LibraryLoad {
path: dir.display().to_string(),
error: format!("load_wasm_configured_with_grants: {e}"),
})?;
Ok(Arc::new(WasmTaskConstructor {
id: manifest.name,
handle: Arc::new(handle),
dependencies: Vec::new(),
}))
}
pub fn unpack_provider_archive(
archive: impl AsRef<Path>,
dest: impl AsRef<Path>,
verify_keys: &[ed25519_dalek::VerifyingKey],
) -> Result<std::path::PathBuf, LoaderError> {
let archive = archive.as_ref();
let dest = dest.as_ref();
let pkg_dir = fidius_core::package::unpack_package(archive, dest).map_err(|e| {
LoaderError::Validation {
reason: format!("unpack provider package '{}': {e}", archive.display()),
}
})?;
if !verify_keys.is_empty() {
fidius_host::package::verify_package(&pkg_dir, verify_keys).map_err(|e| {
LoaderError::Validation {
reason: format!("verify provider signature for '{}': {e}", pkg_dir.display()),
}
})?;
}
Ok(pkg_dir)
}
pub fn load_task_constructor_from_package<C: Serialize>(
archive: impl AsRef<Path>,
dest: impl AsRef<Path>,
package_name: &str,
constructor_name: &str,
config: &C,
verify_keys: &[ed25519_dalek::VerifyingKey],
grants: &ResolvedGrants,
) -> Result<Arc<dyn Task>, LoaderError> {
let dest = dest.as_ref();
let _pkg_dir = unpack_provider_archive(archive, dest, verify_keys)?;
load_task_constructor(dest, package_name, constructor_name, config, grants)
}
#[derive(Debug, Clone)]
pub struct TriggerBinding {
pub poll_interval: Duration,
pub allow_concurrent: bool,
pub workflow_name: String,
pub cron_expression: Option<String>,
}
impl Default for TriggerBinding {
fn default() -> Self {
Self {
poll_interval: Duration::from_secs(60),
allow_concurrent: false,
workflow_name: String::new(),
cron_expression: None,
}
}
}
pub struct WasmTriggerConstructor {
name: String,
handle: Arc<fidius_host::PluginHandle>,
binding: TriggerBinding,
}
impl std::fmt::Debug for WasmTriggerConstructor {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("WasmTriggerConstructor")
.field("name", &self.name)
.field("binding", &self.binding)
.finish_non_exhaustive()
}
}
#[async_trait]
impl Trigger for WasmTriggerConstructor {
fn name(&self) -> &str {
&self.name
}
fn poll_interval(&self) -> Duration {
self.binding.poll_interval
}
fn allow_concurrent(&self) -> bool {
self.binding.allow_concurrent
}
fn cron_expression(&self) -> Option<String> {
self.binding.cron_expression.clone()
}
fn workflow_name(&self) -> &str {
&self.binding.workflow_name
}
async fn poll(&self) -> Result<TriggerResult, TriggerError> {
let name = self.name.clone();
let inv_json = serde_json::to_string(&TriggerInvocation::default()).map_err(|e| {
TriggerError::PollError {
message: format!("trigger '{name}': serialize invocation: {e}"),
}
})?;
let handle = self.handle.clone();
let call_name = name.clone();
let out_json: String = tokio::task::spawn_blocking(move || {
handle.call_method::<_, String>(METHOD_POLL, &(inv_json,))
})
.await
.map_err(|e| TriggerError::PollError {
message: format!("trigger '{call_name}': poll join: {e}"),
})?
.map_err(|e| TriggerError::PollError {
message: format!("trigger '{call_name}': poll FFI call: {e}"),
})?;
let outcome: PollOutcome =
serde_json::from_str(&out_json).map_err(|e| TriggerError::PollError {
message: format!("trigger '{name}': parse poll outcome: {e}"),
})?;
if let Some(err) = outcome.error {
return Err(TriggerError::PollError {
message: format!("trigger '{name}': {err}"),
});
}
if !outcome.fire {
return Ok(TriggerResult::Skip);
}
match outcome.context_json {
None => Ok(TriggerResult::Fire(None)),
Some(ctx_json) => {
let ctx = Context::from_json(ctx_json).map_err(|e| TriggerError::PollError {
message: format!("trigger '{name}': rebuild fire context: {e}"),
})?;
Ok(TriggerResult::Fire(Some(ctx)))
}
}
}
}
pub fn load_trigger_constructor<C: Serialize>(
search_path: impl AsRef<Path>,
package_name: &str,
constructor_name: &str,
config: &C,
binding: TriggerBinding,
grants: &ResolvedGrants,
) -> Result<Arc<dyn Trigger>, LoaderError> {
let search_path = search_path.as_ref();
if let Some((handle, member)) = load_native_member(
search_path,
package_name,
constructor_name,
config,
PrimitiveKind::Trigger,
TRIGGER_CONSTRUCTOR_INTERFACE_VERSION,
"__ProviderTrigger",
)? {
return Ok(Arc::new(WasmTriggerConstructor {
name: member.name,
handle: Arc::new(handle),
binding,
}));
}
let host = PluginHost::builder()
.search_path(search_path)
.build()
.map_err(|e| LoaderError::LibraryLoad {
path: search_path.display().to_string(),
error: format!("build plugin host: {e}"),
})?;
let dir = host
.find_wasm_package(package_name)
.map_err(|e| LoaderError::Validation {
reason: format!("locate wasm constructor package '{package_name}': {e}"),
})?;
let manifest = read_member_manifest(&dir, constructor_name)?;
if manifest.primitive_kind != PrimitiveKind::Trigger {
return Err(LoaderError::Validation {
reason: format!(
"constructor '{}' is {:?}, not Trigger",
manifest.name, manifest.primitive_kind
),
});
}
if manifest.interface_version != TRIGGER_CONSTRUCTOR_INTERFACE_VERSION {
return Err(LoaderError::Validation {
reason: format!(
"constructor '{}' declares trigger-constructor interface v{}, loader supports v{}",
manifest.name, manifest.interface_version, TRIGGER_CONSTRUCTOR_INTERFACE_VERSION
),
});
}
let configure = wrap_member_config(constructor_name, config)?;
let handle = host
.load_wasm_configured_with_grants(
package_name,
&__fidius_TriggerConstructor::TriggerConstructor_WASM_DESCRIPTOR,
&configure,
grants.capabilities.clone(),
grants.egress.clone(),
)
.map_err(|e| LoaderError::LibraryLoad {
path: dir.display().to_string(),
error: format!("load_wasm_configured_with_grants: {e}"),
})?;
Ok(Arc::new(WasmTriggerConstructor {
name: manifest.name,
handle: Arc::new(handle),
binding,
}))
}
pub enum ConstructorBinding {
Task {
namespace: TaskNamespace,
},
Trigger(TriggerBinding),
}
pub fn load_constructor<C: Serialize>(
runtime: &Runtime,
search_path: impl AsRef<Path>,
package_name: &str,
constructor_name: &str,
config: &C,
binding: ConstructorBinding,
grants: &ResolvedGrants,
) -> Result<(), LoaderError> {
let search_path = search_path.as_ref();
let host = PluginHost::builder()
.search_path(search_path)
.build()
.map_err(|e| LoaderError::LibraryLoad {
path: search_path.display().to_string(),
error: format!("build plugin host: {e}"),
})?;
let dir = host
.find_wasm_package(package_name)
.map_err(|e| LoaderError::Validation {
reason: format!("locate wasm constructor package '{package_name}': {e}"),
})?;
let manifest = read_member_manifest(&dir, constructor_name)?;
match (&binding, manifest.primitive_kind) {
(ConstructorBinding::Task { namespace }, PrimitiveKind::Task) => {
let task =
load_task_constructor(search_path, package_name, constructor_name, config, grants)?;
let namespace = namespace.clone();
runtime.register_task(namespace, move || task.clone());
Ok(())
}
(ConstructorBinding::Trigger(tb), PrimitiveKind::Trigger) => {
let name = manifest.name.clone();
let trigger = load_trigger_constructor(
search_path,
package_name,
constructor_name,
config,
tb.clone(),
grants,
)?;
runtime.register_trigger(name, move || trigger.clone());
Ok(())
}
(_, kind) => Err(LoaderError::Validation {
reason: format!(
"constructor '{}' is {:?}, which does not match the supplied binding \
(accumulator/reactor registration is a CLOACI-T-0824 continuation)",
manifest.name, kind
),
}),
}
}
#[fidius_macro::plugin_interface(version = 1, buffer = PluginAllocated, crate = "fidius_core")]
pub trait AccumulatorConstructor: Send + Sync {
fn ingest(&self, invocation_json: String) -> String;
}
pub struct WasmAccumulatorConstructor {
name: String,
handle: Arc<fidius_host::PluginHandle>,
}
impl WasmAccumulatorConstructor {
pub fn name(&self) -> &str {
&self.name
}
}
impl std::fmt::Debug for WasmAccumulatorConstructor {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("WasmAccumulatorConstructor")
.field("name", &self.name)
.finish_non_exhaustive()
}
}
#[async_trait::async_trait]
impl Accumulator for WasmAccumulatorConstructor {
type Output = Vec<u8>;
fn process(&mut self, event: Vec<u8>) -> Option<Vec<u8>> {
let event_json = match String::from_utf8(event) {
Ok(s) => s,
Err(e) => {
tracing::error!(name = %self.name, "accumulator constructor: event is not UTF-8 JSON: {e}");
return None;
}
};
let inv_json = match serde_json::to_string(&AccumulatorInvocation { event_json }) {
Ok(s) => s,
Err(e) => {
tracing::error!(name = %self.name, "accumulator constructor: serialize invocation: {e}");
return None;
}
};
let out_json: String = match self
.handle
.call_method::<_, String>(METHOD_INGEST, &(inv_json,))
{
Ok(s) => s,
Err(e) => {
tracing::error!(name = %self.name, "accumulator constructor: ingest FFI call: {e}");
return None;
}
};
let outcome: AccumulatorOutcome = match serde_json::from_str(&out_json) {
Ok(o) => o,
Err(e) => {
tracing::error!(name = %self.name, "accumulator constructor: parse ingest outcome: {e}");
return None;
}
};
if let Some(err) = outcome.error {
tracing::error!(name = %self.name, "accumulator constructor ingest error: {err}");
return None;
}
match outcome.boundary_json {
None => None,
Some(boundary) => {
if let Err(e) = serde_json::from_str::<serde_json::Value>(&boundary) {
tracing::error!(name = %self.name, "accumulator constructor: boundary is not valid JSON: {e}");
return None;
}
Some(boundary.into_bytes())
}
}
}
}
pub fn load_accumulator_constructor<C: Serialize>(
search_path: impl AsRef<Path>,
package_name: &str,
constructor_name: &str,
config: &C,
grants: &ResolvedGrants,
) -> Result<WasmAccumulatorConstructor, LoaderError> {
let search_path = search_path.as_ref();
if let Some((handle, member)) = load_native_member(
search_path,
package_name,
constructor_name,
config,
PrimitiveKind::Accumulator,
ACCUMULATOR_CONSTRUCTOR_INTERFACE_VERSION,
"__ProviderAccumulator",
)? {
return Ok(WasmAccumulatorConstructor {
name: member.name,
handle: Arc::new(handle),
});
}
let host = PluginHost::builder()
.search_path(search_path)
.build()
.map_err(|e| LoaderError::LibraryLoad {
path: search_path.display().to_string(),
error: format!("build plugin host: {e}"),
})?;
let dir = host
.find_wasm_package(package_name)
.map_err(|e| LoaderError::Validation {
reason: format!("locate wasm constructor package '{package_name}': {e}"),
})?;
let manifest = read_member_manifest(&dir, constructor_name)?;
if manifest.primitive_kind != PrimitiveKind::Accumulator {
return Err(LoaderError::Validation {
reason: format!(
"constructor '{}' is {:?}, not Accumulator",
manifest.name, manifest.primitive_kind
),
});
}
if manifest.interface_version != ACCUMULATOR_CONSTRUCTOR_INTERFACE_VERSION {
return Err(LoaderError::Validation {
reason: format!(
"constructor '{}' declares accumulator-constructor interface v{}, loader supports v{}",
manifest.name, manifest.interface_version, ACCUMULATOR_CONSTRUCTOR_INTERFACE_VERSION
),
});
}
let configure = wrap_member_config(constructor_name, config)?;
let handle = host
.load_wasm_configured_with_grants(
package_name,
&__fidius_AccumulatorConstructor::AccumulatorConstructor_WASM_DESCRIPTOR,
&configure,
grants.capabilities.clone(),
grants.egress.clone(),
)
.map_err(|e| LoaderError::LibraryLoad {
path: dir.display().to_string(),
error: format!("load_wasm_configured_with_grants: {e}"),
})?;
Ok(WasmAccumulatorConstructor {
name: manifest.name,
handle: Arc::new(handle),
})
}
pub struct ProviderStreamSource {
name: String,
_handle: Arc<fidius_host::PluginHandle>,
stream: fidius_host::ChunkStream,
}
#[async_trait]
impl EventSource for ProviderStreamSource {
async fn run(
mut self,
events: tokio::sync::mpsc::Sender<Vec<u8>>,
mut shutdown: tokio::sync::watch::Receiver<bool>,
) -> Result<(), AccumulatorError> {
use futures::StreamExt as _;
loop {
tokio::select! {
item = self.stream.next() => {
match item {
Some(Ok(value)) => {
let boundary: String = fidius_core::from_value(value).map_err(|e| {
AccumulatorError::Init(format!("decode provider stream item: {e}"))
})?;
if boundary.is_empty() {
continue;
}
if events.send(boundary.into_bytes()).await.is_err() {
break; }
}
Some(Err(e)) => {
tracing::warn!(accumulator = %self.name, "provider stream error: {e}");
}
None => break, }
}
_ = shutdown.changed() => {
tracing::debug!(accumulator = %self.name, "provider stream source shutting down");
break; }
}
}
Ok(())
}
}
pub async fn load_stream_accumulator_source<C: Serialize>(
search_path: impl AsRef<Path>,
package_name: &str,
constructor_name: &str,
config: &C,
) -> Result<ProviderStreamSource, LoaderError> {
let search_path = search_path.as_ref();
let (handle, member) = load_native_member(
search_path,
package_name,
constructor_name,
config,
PrimitiveKind::Accumulator,
STREAM_ACCUMULATOR_CONSTRUCTOR_INTERFACE_VERSION,
"__ProviderStreamAccumulator",
)?
.ok_or_else(|| LoaderError::Validation {
reason: format!(
"provider '{package_name}' member '{constructor_name}' is not a native stream \
accumulator (needs runtime = \"native\" + a stream-accumulator member; wasm \
streaming parity is CLOACI-T-0906/0907)"
),
})?;
let handle = Arc::new(handle);
let stream = handle
.call_streaming::<(u64,), String>(METHOD_SOURCE, &(0u64,))
.await
.map_err(|e| LoaderError::LibraryLoad {
path: package_name.to_string(),
error: format!("call_streaming(source) on native stream accumulator: {e}"),
})?;
Ok(ProviderStreamSource {
name: member.name,
_handle: handle,
stream,
})
}
pub async fn load_stream_accumulator_source_from_config(
provider_name: &str,
constructor_name: &str,
config: &std::collections::HashMap<String, String>,
) -> Result<ProviderStreamSource, LoaderError> {
load_stream_accumulator_source_from_config_in(
&provider_search_path(),
provider_name,
constructor_name,
config,
)
.await
}
pub async fn load_stream_accumulator_source_from_config_in(
search_path: &Path,
provider_name: &str,
constructor_name: &str,
config: &std::collections::HashMap<String, String>,
) -> Result<ProviderStreamSource, LoaderError> {
let search_path = search_path.to_path_buf();
let (_dir, provider) =
resolve_native_provider(&search_path, provider_name)?.ok_or_else(|| {
LoaderError::Validation {
reason: format!(
"stream accumulator declares provider '{provider_name}' but no NATIVE provider \
by that name exists under '{}' (bundle it with the package, or check the \
provider/runtime — wasm stream providers are CLOACI-T-0906/0907 follow-up)",
search_path.display()
),
}
})?;
let member = provider
.constructor(constructor_name)
.ok_or_else(|| LoaderError::Validation {
reason: format!(
"provider '{provider_name}' has no constructor '{constructor_name}'; members: [{}]",
provider
.constructors
.iter()
.map(|c| c.name.as_str())
.collect::<Vec<_>>()
.join(", ")
),
})?
.clone();
let author: Vec<(String, serde_json::Value)> = config
.iter()
.map(|(k, v)| {
let field_ty = member
.config_fields
.iter()
.find(|f| f.name == *k)
.map(|f| f.ty.as_str());
let value = match field_ty {
Some("String") | Some("str") | None => serde_json::Value::String(v.clone()),
Some("bool") => v
.parse::<bool>()
.map(serde_json::Value::Bool)
.unwrap_or_else(|_| serde_json::Value::String(v.clone())),
Some(_) => serde_json::from_str::<serde_json::Value>(v)
.ok()
.filter(serde_json::Value::is_number)
.unwrap_or_else(|| serde_json::Value::String(v.clone())),
};
(k.clone(), value)
})
.collect();
let ordered = bind_config_by_name(
&format!("accumulator '{constructor_name}'"),
provider_name,
constructor_name,
&member,
author,
)?;
load_stream_accumulator_source(&search_path, provider_name, constructor_name, &ordered).await
}
#[fidius_macro::plugin_interface(version = 1, buffer = PluginAllocated, crate = "fidius_core")]
pub trait ReactorConstructor: Send + Sync {
fn evaluate(&self, invocation_json: String) -> String;
}
pub struct WasmReactorConstructor {
name: String,
handle: Arc<fidius_host::PluginHandle>,
}
impl WasmReactorConstructor {
pub fn name(&self) -> &str {
&self.name
}
pub async fn evaluate(&self, boundaries_json: String) -> Result<ReactorOutcome, LoaderError> {
let name = self.name.clone();
let inv_json =
serde_json::to_string(&ReactorInvocation { boundaries_json }).map_err(|e| {
LoaderError::Validation {
reason: format!("reactor constructor '{name}': serialize invocation: {e}"),
}
})?;
let handle = self.handle.clone();
let call_name = name.clone();
let out_json: String = tokio::task::spawn_blocking(move || {
handle.call_method::<_, String>(METHOD_EVALUATE, &(inv_json,))
})
.await
.map_err(|e| LoaderError::Validation {
reason: format!("reactor constructor '{call_name}': evaluate join: {e}"),
})?
.map_err(|e| LoaderError::Validation {
reason: format!("reactor constructor '{call_name}': evaluate FFI call: {e}"),
})?;
serde_json::from_str::<ReactorOutcome>(&out_json).map_err(|e| LoaderError::Validation {
reason: format!("reactor constructor '{name}': parse evaluate outcome: {e}"),
})
}
}
impl std::fmt::Debug for WasmReactorConstructor {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("WasmReactorConstructor")
.field("name", &self.name)
.finish_non_exhaustive()
}
}
#[async_trait]
impl ReactorFireDecider for WasmReactorConstructor {
async fn should_fire(&self, snapshot: &InputCache) -> bool {
let boundaries: std::collections::HashMap<String, serde_json::Value> =
crate::computation_graph::reactor::capture_fire_inputs(snapshot);
let boundaries_json = match serde_json::to_string(&boundaries) {
Ok(s) => s,
Err(e) => {
tracing::error!(name = %self.name, "reactor constructor: serialize boundaries: {e}");
return false;
}
};
match self.evaluate(boundaries_json).await {
Ok(outcome) => {
if let Some(err) = outcome.error {
tracing::error!(name = %self.name, "reactor constructor evaluate error: {err}");
return false;
}
outcome.fire
}
Err(e) => {
tracing::error!(name = %self.name, "reactor constructor evaluate failed: {e}");
false
}
}
}
}
pub fn load_reactor_constructor<C: Serialize>(
search_path: impl AsRef<Path>,
package_name: &str,
constructor_name: &str,
config: &C,
grants: &ResolvedGrants,
) -> Result<WasmReactorConstructor, LoaderError> {
let search_path = search_path.as_ref();
if let Some((handle, member)) = load_native_member(
search_path,
package_name,
constructor_name,
config,
PrimitiveKind::Reactor,
REACTOR_CONSTRUCTOR_INTERFACE_VERSION,
"__ProviderReactor",
)? {
return Ok(WasmReactorConstructor {
name: member.name,
handle: Arc::new(handle),
});
}
let host = PluginHost::builder()
.search_path(search_path)
.build()
.map_err(|e| LoaderError::LibraryLoad {
path: search_path.display().to_string(),
error: format!("build plugin host: {e}"),
})?;
let dir = host
.find_wasm_package(package_name)
.map_err(|e| LoaderError::Validation {
reason: format!("locate wasm constructor package '{package_name}': {e}"),
})?;
let manifest = read_member_manifest(&dir, constructor_name)?;
if manifest.primitive_kind != PrimitiveKind::Reactor {
return Err(LoaderError::Validation {
reason: format!(
"constructor '{}' is {:?}, not Reactor",
manifest.name, manifest.primitive_kind
),
});
}
if manifest.interface_version != REACTOR_CONSTRUCTOR_INTERFACE_VERSION {
return Err(LoaderError::Validation {
reason: format!(
"constructor '{}' declares reactor-constructor interface v{}, loader supports v{}",
manifest.name, manifest.interface_version, REACTOR_CONSTRUCTOR_INTERFACE_VERSION
),
});
}
let configure = wrap_member_config(constructor_name, config)?;
let handle = host
.load_wasm_configured_with_grants(
package_name,
&__fidius_ReactorConstructor::ReactorConstructor_WASM_DESCRIPTOR,
&configure,
grants.capabilities.clone(),
grants.egress.clone(),
)
.map_err(|e| LoaderError::LibraryLoad {
path: dir.display().to_string(),
error: format!("load_wasm_configured_with_grants: {e}"),
})?;
Ok(WasmReactorConstructor {
name: manifest.name,
handle: Arc::new(handle),
})
}
static PROVIDER_SEARCH_PATH: std::sync::RwLock<Option<std::path::PathBuf>> =
std::sync::RwLock::new(None);
pub const PROVIDER_PATH_ENV: &str = "CLOACINA_PROVIDER_PATH";
pub const DEFAULT_PROVIDER_DIR: &str = "providers";
pub fn set_provider_search_path(path: impl Into<std::path::PathBuf>) {
*PROVIDER_SEARCH_PATH.write().unwrap() = Some(path.into());
}
pub fn clear_provider_search_path() {
*PROVIDER_SEARCH_PATH.write().unwrap() = None;
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ProviderScope {
Staged(std::path::PathBuf),
Unbundled,
}
impl ProviderScope {
pub fn from_staged_root(root: Option<std::path::PathBuf>) -> Self {
match root {
Some(root) => ProviderScope::Staged(root),
None => ProviderScope::Unbundled,
}
}
pub fn search_path(&self) -> std::path::PathBuf {
match self {
ProviderScope::Staged(root) => root.clone(),
ProviderScope::Unbundled => provider_search_path_from_env(),
}
}
}
thread_local! {
static CURRENT_PROVIDER_SCOPE: std::cell::RefCell<Option<ProviderScope>> =
const { std::cell::RefCell::new(None) };
}
pub fn current_provider_scope() -> Option<ProviderScope> {
CURRENT_PROVIDER_SCOPE.with(|slot| slot.borrow().clone())
}
#[must_use = "the scope is uninstalled when the guard drops"]
pub struct ScopedProviderSearch {
previous: Option<ProviderScope>,
}
impl ScopedProviderSearch {
pub fn enter(scope: ProviderScope) -> Self {
Self::enter_opt(Some(scope))
}
pub fn enter_opt(scope: Option<ProviderScope>) -> Self {
let previous = CURRENT_PROVIDER_SCOPE.with(|slot| slot.replace(scope));
Self { previous }
}
pub fn for_staged_root(root: Option<std::path::PathBuf>) -> Self {
Self::enter(ProviderScope::from_staged_root(root))
}
}
impl Drop for ScopedProviderSearch {
fn drop(&mut self) {
let previous = self.previous.take();
CURRENT_PROVIDER_SCOPE.with(|slot| *slot.borrow_mut() = previous);
}
}
fn provider_search_path_from_env() -> std::path::PathBuf {
if let Ok(p) = std::env::var(PROVIDER_PATH_ENV) {
if !p.is_empty() {
return std::path::PathBuf::from(p);
}
}
std::path::PathBuf::from(DEFAULT_PROVIDER_DIR)
}
pub fn provider_search_path() -> std::path::PathBuf {
if let Some(scope) = current_provider_scope() {
return scope.search_path();
}
if let Some(p) = PROVIDER_SEARCH_PATH.read().unwrap().clone() {
return p;
}
provider_search_path_from_env()
}
fn provider_package_name(from: &str) -> &str {
match from.split_once('@') {
Some((name, _ver)) => name,
None => from,
}
}
fn provider_version_pin(from: &str) -> Option<&str> {
from.split_once('@').map(|(_name, ver)| ver)
}
fn version_satisfies_pin(actual: &str, pin: &str) -> bool {
actual == pin || actual.starts_with(&format!("{pin}."))
}
fn enforce_version_pin(context: &str, from: &str, dir: &Path) -> Result<(), LoaderError> {
let Some(pin) = provider_version_pin(from) else {
return Ok(());
};
let provider = read_provider_manifest(dir)?;
if version_satisfies_pin(&provider.version, pin) {
return Ok(());
}
Err(LoaderError::Validation {
reason: format!(
"{context}: provider '{}' resolved at version {} but the reference \
pins @{pin} — restage the pinned provider version or update the pin",
provider.name, provider.version
),
})
}
fn runtime_literal(runtime: ProviderRuntime) -> &'static str {
match runtime {
ProviderRuntime::Wasm => "wasm",
ProviderRuntime::Native => "native",
}
}
pub fn parse_runtime_pin(context: &str, literal: &str) -> Result<ProviderRuntime, LoaderError> {
match literal {
"wasm" => Ok(ProviderRuntime::Wasm),
"native" => Ok(ProviderRuntime::Native),
other => Err(LoaderError::Validation {
reason: format!(
"{context}: unknown runtime pin '{other}'; valid values are \"wasm\" and \"native\""
),
}),
}
}
fn enforce_runtime_pin(
context: &str,
from: &str,
search_path: &Path,
package_name: &str,
pin: Option<ProviderRuntime>,
grants: &GrantSpec,
) -> Result<Option<std::path::PathBuf>, LoaderError> {
let native_dir = resolve_native_provider(search_path, package_name)?.map(|(dir, _)| dir);
let resolved = match native_dir {
Some(_) => ProviderRuntime::Native,
None => ProviderRuntime::Wasm,
};
if let Some(pinned) = pin {
if pinned != resolved {
return Err(LoaderError::Validation {
reason: format!(
"{context} (from = '{from}'): pinned runtime = \"{}\" but provider \
'{package_name}' resolved as runtime = \"{}\". A provider changing runtime \
between versions is a BREAKING change (grants are enforced on wasm, advisory \
on native) — restage the pinned runtime, or update the pin deliberately",
runtime_literal(pinned),
runtime_literal(resolved),
),
});
}
}
if resolved == ProviderRuntime::Native && !grants.is_empty() {
if pin != Some(ProviderRuntime::Native) {
return Err(LoaderError::Validation {
reason: format!(
"{context} (from = '{from}'): grants are not enforced on native providers; \
either pin runtime = \"wasm\", acknowledge with runtime = \"native\", or \
remove grants"
),
});
}
tracing::warn!(
context = %context,
from = %from,
provider = %package_name,
"grants declared against a NATIVE provider are ADVISORY ONLY (acknowledged via \
runtime = \"native\"): the constructor runs unsandboxed in-process with full host \
trust"
);
}
Ok(native_dir)
}
pub struct ConstructorNode {
id: String,
inner: Arc<dyn Task>,
dependencies: Vec<TaskNamespace>,
}
impl std::fmt::Debug for ConstructorNode {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ConstructorNode")
.field("id", &self.id)
.field("dependencies", &self.dependencies)
.finish_non_exhaustive()
}
}
#[async_trait]
impl Task for ConstructorNode {
async fn execute(
&self,
context: Context<serde_json::Value>,
) -> Result<Context<serde_json::Value>, TaskError> {
self.inner.execute(context).await
}
fn id(&self) -> &str {
&self.id
}
fn dependencies(&self) -> &[TaskNamespace] {
&self.dependencies
}
fn retry_policy(&self) -> crate::retry::RetryPolicy {
self.inner.retry_policy()
}
fn trigger_rules(&self) -> serde_json::Value {
self.inner.trigger_rules()
}
fn code_fingerprint(&self) -> Option<String> {
self.inner.code_fingerprint()
}
fn requires_handle(&self) -> bool {
self.inner.requires_handle()
}
}
enum TypedConfigValue {
Str(String),
Bool(bool),
I8(i8),
I16(i16),
I32(i32),
I64(i64),
U8(u8),
U16(u16),
U32(u32),
U64(u64),
F32(f32),
F64(f64),
}
impl Serialize for TypedConfigValue {
fn serialize<S: serde::Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
match self {
TypedConfigValue::Str(v) => s.serialize_str(v),
TypedConfigValue::Bool(v) => s.serialize_bool(*v),
TypedConfigValue::I8(v) => s.serialize_i8(*v),
TypedConfigValue::I16(v) => s.serialize_i16(*v),
TypedConfigValue::I32(v) => s.serialize_i32(*v),
TypedConfigValue::I64(v) => s.serialize_i64(*v),
TypedConfigValue::U8(v) => s.serialize_u8(*v),
TypedConfigValue::U16(v) => s.serialize_u16(*v),
TypedConfigValue::U32(v) => s.serialize_u32(*v),
TypedConfigValue::U64(v) => s.serialize_u64(*v),
TypedConfigValue::F32(v) => s.serialize_f32(*v),
TypedConfigValue::F64(v) => s.serialize_f64(*v),
}
}
}
struct OrderedConfig(Vec<TypedConfigValue>);
impl Serialize for OrderedConfig {
fn serialize<S: serde::Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
use serde::ser::SerializeTuple;
let mut tup = s.serialize_tuple(self.0.len())?;
for v in &self.0 {
tup.serialize_element(v)?;
}
tup.end()
}
}
fn coerce_config_value(
key: &str,
ty: &str,
value: &serde_json::Value,
) -> Result<TypedConfigValue, String> {
let wrong =
|want: &str| format!("config field '{key}' expects {want} (declared `{ty}`), got {value}");
match ty {
"String" | "str" => value
.as_str()
.map(|s| TypedConfigValue::Str(s.to_string()))
.ok_or_else(|| wrong("a string")),
"bool" => value
.as_bool()
.map(TypedConfigValue::Bool)
.ok_or_else(|| wrong("a boolean")),
"i8" => int_in_range::<i8>(value)
.map(TypedConfigValue::I8)
.ok_or_else(|| wrong("an i8")),
"i16" => int_in_range::<i16>(value)
.map(TypedConfigValue::I16)
.ok_or_else(|| wrong("an i16")),
"i32" => int_in_range::<i32>(value)
.map(TypedConfigValue::I32)
.ok_or_else(|| wrong("an i32")),
"i64" | "isize" => value
.as_i64()
.map(TypedConfigValue::I64)
.ok_or_else(|| wrong("an i64")),
"u8" => uint_in_range::<u8>(value)
.map(TypedConfigValue::U8)
.ok_or_else(|| wrong("a u8")),
"u16" => uint_in_range::<u16>(value)
.map(TypedConfigValue::U16)
.ok_or_else(|| wrong("a u16")),
"u32" => uint_in_range::<u32>(value)
.map(TypedConfigValue::U32)
.ok_or_else(|| wrong("a u32")),
"u64" | "usize" => value
.as_u64()
.map(TypedConfigValue::U64)
.ok_or_else(|| wrong("a u64")),
"f32" => value
.as_f64()
.map(|f| TypedConfigValue::F32(f as f32))
.ok_or_else(|| wrong("a number")),
"f64" => value
.as_f64()
.map(TypedConfigValue::F64)
.ok_or_else(|| wrong("a number")),
other => Err(format!(
"config field '{key}' has unsupported declared type `{other}`; \
constructor config supports string/bool/integer/float literals"
)),
}
}
fn int_in_range<T: TryFrom<i64>>(value: &serde_json::Value) -> Option<T> {
value.as_i64().and_then(|v| T::try_from(v).ok())
}
fn uint_in_range<T: TryFrom<u64>>(value: &serde_json::Value) -> Option<T> {
value.as_u64().and_then(|v| T::try_from(v).ok())
}
fn bind_config_by_name(
node_id: &str,
from: &str,
constructor_name: &str,
manifest: &ConstructorManifest,
mut author: Vec<(String, serde_json::Value)>,
) -> Result<OrderedConfig, LoaderError> {
let ctx = |reason: String| {
LoaderError::Validation {
reason: format!(
"constructor node '{node_id}' (from = '{from}', constructor = '{constructor_name}'): {reason}"
),
}
};
let declared: Vec<&str> = manifest
.config_fields
.iter()
.map(|f| f.name.as_str())
.collect();
for (k, _) in &author {
if !declared.contains(&k.as_str()) {
return Err(ctx(format!(
"config key '{k}' is not a #[config] field of constructor '{}'. \
Declared config fields: [{}]",
manifest.name,
declared.join(", ")
)));
}
}
let mut ordered = Vec::with_capacity(manifest.config_fields.len());
for field in &manifest.config_fields {
let pos = author.iter().position(|(k, _)| k == &field.name);
let value = match pos {
Some(i) => author.remove(i).1,
None => {
return Err(ctx(format!(
"missing required config field '{}' for constructor '{}'",
field.name, manifest.name
)))
}
};
let typed = coerce_config_value(&field.name, &field.ty, &value).map_err(ctx)?;
ordered.push(typed);
}
if let Some((dup, _)) = author.into_iter().next() {
return Err(ctx(format!("duplicate config key '{dup}'")));
}
Ok(OrderedConfig(ordered))
}
fn lint_constructor_grants(node_id: &str, from: &str, dir: &Path, grants: &GrantSpec) {
let Ok(pkg) = fidius_core::package::load_manifest_untyped(dir) else {
return;
};
let Some(wasm) = pkg.wasm.as_ref() else {
return;
};
for warning in lint_unmet_intents(&wasm.capabilities, grants) {
tracing::warn!(node = %node_id, from = %from, "{warning}");
}
}
pub fn load_constructor_node(
node_id: &str,
from: &str,
constructor_name: &str,
config: Vec<(String, serde_json::Value)>,
dependencies: Vec<TaskNamespace>,
grants: GrantSpec,
) -> Result<Arc<dyn Task>, LoaderError> {
load_constructor_node_pinned(
node_id,
from,
constructor_name,
config,
dependencies,
grants,
None,
)
}
pub fn load_constructor_node_pinned(
node_id: &str,
from: &str,
constructor_name: &str,
config: Vec<(String, serde_json::Value)>,
dependencies: Vec<TaskNamespace>,
grants: GrantSpec,
runtime_pin: Option<ProviderRuntime>,
) -> Result<Arc<dyn Task>, LoaderError> {
load_constructor_node_pinned_in(
&provider_search_path(),
node_id,
from,
constructor_name,
config,
dependencies,
grants,
runtime_pin,
)
}
pub fn load_constructor_node_in(
search_path: &Path,
node_id: &str,
from: &str,
constructor_name: &str,
config: Vec<(String, serde_json::Value)>,
dependencies: Vec<TaskNamespace>,
grants: GrantSpec,
) -> Result<Arc<dyn Task>, LoaderError> {
load_constructor_node_pinned_in(
search_path,
node_id,
from,
constructor_name,
config,
dependencies,
grants,
None,
)
}
#[allow(clippy::too_many_arguments)]
pub fn load_constructor_node_pinned_in(
search_path: &Path,
node_id: &str,
from: &str,
constructor_name: &str,
config: Vec<(String, serde_json::Value)>,
dependencies: Vec<TaskNamespace>,
grants: GrantSpec,
runtime_pin: Option<ProviderRuntime>,
) -> Result<Arc<dyn Task>, LoaderError> {
let search_path = search_path.to_path_buf();
let package_name = provider_package_name(from);
let native_dir = enforce_runtime_pin(
&format!("constructor node '{node_id}'"),
from,
&search_path,
package_name,
runtime_pin,
&grants,
)?;
let resolved = translate(&grants).map_err(|e| LoaderError::Validation {
reason: format!("constructor node '{node_id}' (from = '{from}'): {e}"),
})?;
let dir = match native_dir {
Some(dir) => dir,
None => {
let host = PluginHost::builder()
.search_path(&search_path)
.build()
.map_err(|e| LoaderError::LibraryLoad {
path: search_path.display().to_string(),
error: format!("build plugin host: {e}"),
})?;
host.find_wasm_package(package_name)
.map_err(|e| LoaderError::Validation {
reason: format!(
"resolve constructor node '{node_id}' (from = '{from}', constructor = \
'{constructor_name}'): locate provider package '{package_name}' in \
provider search path '{}': {e}",
search_path.display()
),
})?
}
};
enforce_version_pin(&format!("constructor node '{node_id}'"), from, &dir)?;
let manifest = read_member_manifest(&dir, constructor_name)?;
lint_constructor_grants(node_id, from, &dir, &grants);
let ordered_config = bind_config_by_name(node_id, from, constructor_name, &manifest, config)?;
let task = load_task_constructor(
&search_path,
package_name,
constructor_name,
&ordered_config,
&resolved,
)
.map_err(|e| LoaderError::Validation {
reason: format!(
"resolve constructor node '{node_id}' (from = '{from}', constructor = \
'{constructor_name}') in provider search path '{}': {e}",
search_path.display()
),
})?;
Ok(Arc::new(ConstructorNode {
id: node_id.to_string(),
inner: task,
dependencies,
}))
}
pub fn load_reactor_constructor_node(
from: &str,
constructor_name: &str,
config: Vec<(String, serde_json::Value)>,
grants: GrantSpec,
) -> Result<Arc<dyn ReactorFireDecider>, LoaderError> {
load_reactor_constructor_node_pinned(from, constructor_name, config, grants, None)
}
pub fn load_reactor_constructor_node_pinned(
from: &str,
constructor_name: &str,
config: Vec<(String, serde_json::Value)>,
grants: GrantSpec,
runtime_pin: Option<ProviderRuntime>,
) -> Result<Arc<dyn ReactorFireDecider>, LoaderError> {
load_reactor_constructor_node_pinned_in(
&provider_search_path(),
from,
constructor_name,
config,
grants,
runtime_pin,
)
}
pub fn load_reactor_constructor_node_pinned_in(
search_path: &Path,
from: &str,
constructor_name: &str,
config: Vec<(String, serde_json::Value)>,
grants: GrantSpec,
runtime_pin: Option<ProviderRuntime>,
) -> Result<Arc<dyn ReactorFireDecider>, LoaderError> {
let search_path = search_path.to_path_buf();
let package_name = provider_package_name(from);
let native_dir = enforce_runtime_pin(
&format!("reactor constructor '{constructor_name}'"),
from,
&search_path,
package_name,
runtime_pin,
&grants,
)?;
let resolved = translate(&grants).map_err(|e| LoaderError::Validation {
reason: format!("reactor constructor '{constructor_name}' (from = '{from}'): {e}"),
})?;
let dir = match native_dir {
Some(dir) => dir,
None => {
let host = PluginHost::builder()
.search_path(&search_path)
.build()
.map_err(|e| LoaderError::LibraryLoad {
path: search_path.display().to_string(),
error: format!("build plugin host: {e}"),
})?;
host.find_wasm_package(package_name)
.map_err(|e| LoaderError::Validation {
reason: format!(
"resolve reactor constructor (from = '{from}', constructor = \
'{constructor_name}'): locate provider package '{package_name}' in \
provider search path '{}': {e}",
search_path.display()
),
})?
}
};
enforce_version_pin(
&format!("reactor constructor '{constructor_name}'"),
from,
&dir,
)?;
let manifest = read_member_manifest(&dir, constructor_name)?;
lint_constructor_grants(constructor_name, from, &dir, &grants);
let ordered_config =
bind_config_by_name(constructor_name, from, constructor_name, &manifest, config)?;
let reactor_constructor = load_reactor_constructor(
&search_path,
package_name,
constructor_name,
&ordered_config,
&resolved,
)
.map_err(|e| LoaderError::Validation {
reason: format!(
"resolve reactor constructor (from = '{from}', constructor = \
'{constructor_name}') in provider search path '{}': {e}",
search_path.display()
),
})?;
Ok(Arc::new(reactor_constructor) as Arc<dyn ReactorFireDecider>)
}
#[cfg(test)]
mod provider_scope_tests {
use super::*;
use std::path::PathBuf;
#[test]
#[serial_test::serial(provider_search_path)]
fn scope_wins_over_the_process_override() {
let host = PathBuf::from("/tmp/t0925-host-configured");
let tenant_a = PathBuf::from("/tmp/t0925-tenant-a");
set_provider_search_path(&host);
assert_eq!(provider_search_path(), host, "unscoped = the process knob");
{
let _scope = ScopedProviderSearch::enter(ProviderScope::Staged(tenant_a.clone()));
assert_eq!(
provider_search_path(),
tenant_a,
"a staged scope must win over the process-wide override"
);
}
assert_eq!(
provider_search_path(),
host,
"the override is visible again once the scope drops"
);
clear_provider_search_path();
}
#[test]
#[serial_test::serial(provider_search_path)]
fn unbundled_scope_skips_the_process_override() {
let host = PathBuf::from("/tmp/t0925-host-configured-2");
set_provider_search_path(&host);
{
let _scope = ScopedProviderSearch::enter(ProviderScope::Unbundled);
assert_eq!(
provider_search_path(),
provider_search_path_from_env(),
"a package that bundles nothing must not inherit the host/other-tenant \
override — it falls through to env/default"
);
}
clear_provider_search_path();
}
#[test]
fn scopes_nest_and_restore() {
let outer = PathBuf::from("/tmp/t0925-outer");
let inner = PathBuf::from("/tmp/t0925-inner");
assert_eq!(current_provider_scope(), None);
let _o = ScopedProviderSearch::enter(ProviderScope::Staged(outer.clone()));
{
let _i = ScopedProviderSearch::enter(ProviderScope::Staged(inner.clone()));
assert_eq!(provider_search_path(), inner);
}
assert_eq!(
provider_search_path(),
outer,
"an inner scope must not strand the outer one"
);
}
#[test]
fn capture_and_reinstall_carries_a_scope_across_threads() {
let staged = PathBuf::from("/tmp/t0925-staged-for-import");
let _scope = ScopedProviderSearch::enter(ProviderScope::Staged(staged.clone()));
let captured = current_provider_scope();
let seen = std::thread::spawn(move || {
let unscoped = current_provider_scope();
let _s = ScopedProviderSearch::enter_opt(captured);
(unscoped, provider_search_path())
})
.join()
.unwrap();
assert_eq!(seen.0, None, "a fresh thread starts unscoped");
assert_eq!(seen.1, staged, "the captured scope resolves on that thread");
}
}
#[cfg(test)]
mod version_pin_tests {
use super::*;
#[test]
fn pin_parsing_splits_name_and_version() {
assert_eq!(provider_package_name("p@0.1.0"), "p");
assert_eq!(provider_package_name("p"), "p");
assert_eq!(provider_version_pin("p@0.1.0"), Some("0.1.0"));
assert_eq!(provider_version_pin("p"), None);
}
#[test]
fn pin_matching_requires_segment_boundary() {
assert!(version_satisfies_pin("0.1.0", "0.1.0"));
assert!(version_satisfies_pin("0.1.5", "0.1"));
assert!(version_satisfies_pin("1.2.0", "1"));
assert!(!version_satisfies_pin("0.10.0", "0.1"));
assert!(!version_satisfies_pin("10.0.0", "1"));
assert!(!version_satisfies_pin("0.2.0", "0.1.0"));
}
}