use std::sync::Arc;
use std::time::Instant;
use crate::config::PluginsConfig;
use crate::engine::{FunctionEntry, PluginBinding};
use crate::storage::models::{Plugin, PluginHealth};
use crate::storage::repositories::plugins::PluginRepository;
use super::handler::PluginFunctionHandler;
use super::limits::Limits;
use super::manifest::Manifest;
use super::runtime::WasmRuntime;
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct PluginLoadIssue {
pub plugin: String,
pub version: i64,
pub digest: String,
pub stage: &'static str,
pub reason: String,
}
pub struct LoadedPlugin {
pub id: String,
pub version: i64,
pub digest: String,
pub functions: Vec<String>,
pub compile_ms: u64,
entries: Vec<FunctionEntry>,
handlers: Vec<PluginFunctionHandler>,
}
#[derive(Default)]
pub struct PluginSet {
pub plugins: Vec<LoadedPlugin>,
pub issues: Vec<PluginLoadIssue>,
pub unavailable: Vec<(String, String)>,
fingerprint: String,
}
impl PluginSet {
pub fn empty() -> Self {
Self::default()
}
pub fn fingerprint_of(rows: &[Plugin]) -> String {
let mut parts: Vec<String> = rows
.iter()
.map(|r| format!("{}@{}:{}", r.plugin_id, r.version, r.digest))
.collect();
parts.sort();
parts.join(";")
}
pub fn fingerprint(&self) -> &str {
&self.fingerprint
}
pub fn entries(&self) -> Vec<FunctionEntry> {
self.plugins
.iter()
.flat_map(|p| p.entries.iter().cloned())
.collect()
}
pub fn handlers(&self) -> Vec<(String, dataflow_rs::BoxedFunctionHandler)> {
self.plugins
.iter()
.flat_map(|p| p.handlers.iter())
.map(|h| {
(
h.name().to_string(),
Box::new(h.clone()) as dataflow_rs::BoxedFunctionHandler,
)
})
.collect()
}
pub fn health_of(&self, plugin_id: &str, version: i64, enabled: bool) -> PluginHealth {
if let Some(p) = self
.plugins
.iter()
.find(|p| p.id == plugin_id && p.version == version)
{
return PluginHealth {
state: "loaded".to_string(),
compile_ms: Some(p.compile_ms),
reason: None,
};
}
if let Some(issue) = self
.issues
.iter()
.find(|i| i.plugin == plugin_id && i.version == version)
{
return PluginHealth {
state: if issue.stage == "disabled" {
"disabled".to_string()
} else {
"failed".to_string()
},
compile_ms: None,
reason: Some(format!("{}: {}", issue.stage, issue.reason)),
};
}
PluginHealth {
state: if enabled { "inactive" } else { "disabled" }.to_string(),
compile_ms: None,
reason: None,
}
}
pub fn annotate(&self, reason: &mut String) {
for (function, why) in &self.unavailable {
if reason.contains(function.as_str()) {
reason.push_str(&format!(" — plugin function '{function}': {why}"));
}
}
}
}
pub async fn load_active(
rows: Vec<Plugin>,
repo: &dyn PluginRepository,
runtime: Option<&Arc<WasmRuntime>>,
config: &PluginsConfig,
) -> PluginSet {
let fingerprint = PluginSet::fingerprint_of(&rows);
let mut set = PluginSet {
fingerprint,
..PluginSet::default()
};
for row in rows {
match load_one(&row, repo, runtime, config).await {
Ok(loaded) => {
crate::metrics::record_plugin_load(
&row.plugin_id,
"ok",
Some(loaded.compile_ms as f64 / 1000.0),
);
tracing::info!(
plugin = %row.plugin_id,
version = row.version,
digest = %row.digest,
functions = ?loaded.functions,
compile_ms = loaded.compile_ms,
"Plugin loaded"
);
set.plugins.push(loaded);
}
Err((stage, reason)) => {
crate::metrics::record_plugin_load(&row.plugin_id, "error", None);
tracing::error!(
plugin = %row.plugin_id,
version = row.version,
digest = %row.digest,
stage,
reason = %reason,
"Plugin not loaded: the workflows naming its functions are quarantined"
);
if let Ok(manifest) = serde_json::from_str::<Manifest>(&row.manifest_json) {
for f in manifest.function_names() {
set.unavailable.push((
f.to_string(),
format!("{} v{} {stage}: {reason}", row.plugin_id, row.version),
));
}
}
set.issues.push(PluginLoadIssue {
plugin: row.plugin_id.clone(),
version: row.version,
digest: row.digest.clone(),
stage,
reason,
});
}
}
}
set
}
async fn load_one(
row: &Plugin,
repo: &dyn PluginRepository,
runtime: Option<&Arc<WasmRuntime>>,
config: &PluginsConfig,
) -> Result<LoadedPlugin, (&'static str, String)> {
let manifest: Manifest = serde_json::from_str(&row.manifest_json)
.map_err(|e| ("manifest", format!("stored manifest does not parse: {e}")))?;
let Some(runtime) = runtime else {
return Err((
"disabled",
"plugins are disabled on this node (plugins.enabled = false)".to_string(),
));
};
super::trust::verify(
&config.trust.public_keys,
&row.digest,
row.signature.as_deref(),
)
.map_err(|reason| ("signature", reason))?;
let loaded = match runtime.cached(&row.digest) {
Some(loaded) => loaded,
None => {
let bytes = repo
.get_artifact(&row.digest)
.await
.map_err(|e| ("artifact", e.to_string()))?
.ok_or_else(|| {
(
"artifact",
format!("no component is stored under {}", row.digest),
)
})?;
let started = Instant::now();
let loaded = runtime
.load(bytes)
.await
.map_err(|e| (load_stage(&e), e.to_string()))?;
let _ = started;
loaded
}
};
if loaded.digest != row.digest {
return Err((
"artifact",
format!(
"stored bytes hash to {} but the row names {}",
loaded.digest, row.digest
),
));
}
let limits = Limits::effective(config, &row.plugin_id);
let functions: Vec<String> = manifest.function_names().map(str::to_string).collect();
let names: Vec<&str> = functions.iter().map(String::as_str).collect();
runtime
.self_test(&loaded, &limits, &names)
.await
.map_err(|e| ("self_test", e.to_string()))?;
let binding = PluginBinding {
id: row.plugin_id.clone(),
version: row.version,
digest: row.digest.clone(),
abi: manifest.abi.clone(),
};
let entries = manifest.entries(&binding);
let handlers = entries
.iter()
.map(|e| {
PluginFunctionHandler::new(Arc::new(e.clone()), loaded.clone(), runtime.clone(), limits)
})
.collect();
Ok(LoadedPlugin {
id: row.plugin_id.clone(),
version: row.version,
digest: row.digest.clone(),
functions,
compile_ms: loaded.compile_time.as_millis() as u64,
entries,
handlers,
})
}
fn load_stage(e: &super::runtime::LoadError) -> &'static str {
use super::runtime::LoadError;
match e {
LoadError::TooLarge { .. } => "size",
LoadError::Compile(_) => "compile",
LoadError::Link(_) => "link",
LoadError::SelfTest { .. } => "self_test",
}
}