Merge pull request #475 from BCNelson/bcn/plugins

Add M0 WASM plugin system (sandboxed, observe-only)
This commit is contained in:
Dionisio Pozo
2026-06-17 12:41:49 +02:00
committed by GitHub
45 changed files with 7330 additions and 30 deletions
+2
View File
@@ -22,6 +22,8 @@ pub mod password_hasher;
pub mod path_resolver_service;
pub mod path_service;
pub mod pg_acl_engine;
#[cfg(feature = "plugins")]
pub mod plugins;
pub mod retry_blob_backend;
pub mod s3_blob_backend;
pub mod search_index;
@@ -0,0 +1,70 @@
//! Background plugin-log maintenance: a periodic sweep that prunes each
//! installed plugin's rotated log segments by age + aggregate size.
//!
//! Modeled on [`crate::infrastructure::services::trash_cleanup_service`]: a
//! single spawned task on a fixed interval, an immediate first run, and
//! log-and-continue on error. `file-rotate` already compresses + caps segment
//! *count* at write time; this is the only thing that enforces the per-plugin
//! age/byte retention and the only thing that ever prunes *idle* plugins (which
//! never trigger a write-time rotation).
use std::sync::Arc;
use std::time::Duration;
use chrono::Utc;
use tokio::time;
use tracing::{debug, info};
use super::log_store::PluginLogStore;
use crate::application::ports::plugin_ports::PluginManagementPort;
/// Periodically sweeps every installed plugin's logs against its retention.
pub struct PluginLogMaintenanceService {
log_store: Arc<PluginLogStore>,
manager: Arc<dyn PluginManagementPort>,
interval_hours: u64,
}
impl PluginLogMaintenanceService {
pub fn new(
log_store: Arc<PluginLogStore>,
manager: Arc<dyn PluginManagementPort>,
interval_hours: u64,
) -> Self {
Self {
log_store,
manager,
interval_hours: interval_hours.max(1),
}
}
/// Spawn the periodic sweep task.
pub fn start(&self) {
let log_store = self.log_store.clone();
let manager = self.manager.clone();
let interval_hours = self.interval_hours;
info!(
"Starting plugin log maintenance job with interval of {} hours",
interval_hours
);
tokio::spawn(async move {
let mut interval = time::interval(Duration::from_secs(interval_hours * 60 * 60));
// First tick fires immediately.
loop {
interval.tick().await;
Self::sweep_all(&log_store, &manager).await;
}
});
}
async fn sweep_all(log_store: &PluginLogStore, manager: &Arc<dyn PluginManagementPort>) {
let now = Utc::now();
let plugins = manager.list();
debug!("Plugin log sweep over {} plugin(s)", plugins.len());
for plugin in plugins {
log_store.request_sweep(&plugin.id, now).await;
}
}
}
@@ -0,0 +1,731 @@
//! Per-plugin structured log storage — an async, in-order actor over disk files.
//!
//! Every plugin gets its own directory under the log root:
//! ```text
//! {root}/{plugin_id}/events.jsonl # active (file-rotate writes here)
//! {root}/{plugin_id}/events.jsonl.<timestamp>.gz # rotated + gzip'd (immutable)
//! {root}/{plugin_id}/retention.json # per-plugin retention override
//! ```
//!
//! **Async + strictly in order.** All file mutations funnel through a single
//! background thread that owns the per-plugin [`FileRotate`] writers. Because
//! there is exactly one consumer draining one channel FIFO, batches land in
//! enqueue order with no locks, and the dispatch path never blocks on IO — it
//! just sends. Rotation, gzip-on-rotate and a coarse segment ceiling are handled
//! by `file-rotate`; per-plugin age + aggregate-byte retention is the [`sweep`]
//! (run on a schedule), the only thing that ever prunes *idle* plugins.
//!
//! [`sweep`]: PluginLogStore::request_sweep
use std::collections::HashMap;
use std::fs;
use std::io::{Read, Write};
use std::path::{Path, PathBuf};
use std::time::SystemTime;
use chrono::{DateTime, Duration, Utc};
use file_rotate::{
ContentLimit, FileRotate,
compression::Compression,
suffix::{AppendTimestamp, FileLimit},
};
use flate2::read::GzDecoder;
use tokio::sync::{broadcast, mpsc, oneshot};
use super::runtime::InvokeOutcome;
use crate::application::ports::plugin_ports::{
LogEntry, LogPage, LogQuery, PluginLogEvent, RetentionSettings,
};
/// Name of the active (uncompressed) log file inside a plugin's log dir.
const ACTIVE_FILE: &str = "events.jsonl";
/// Marker file holding a plugin's retention override.
const RETENTION_FILE: &str = "retention.json";
/// Live broadcast buffer; a slow tailer past this gets `Lagged` (never blocks).
const LIVE_CAPACITY: usize = 256;
/// Commands processed in receipt order by the single actor thread.
enum LogCommand {
Append {
plugin_id: String,
entries: Vec<LogEntry>,
},
Read {
plugin_id: String,
query: LogQuery,
reply: oneshot::Sender<LogPage>,
},
Clear {
plugin_id: String,
reply: oneshot::Sender<()>,
},
Remove {
plugin_id: String,
},
GetRetention {
plugin_id: String,
reply: oneshot::Sender<RetentionSettings>,
},
SetRetention {
plugin_id: String,
settings: RetentionSettings,
reply: oneshot::Sender<()>,
},
Sweep {
plugin_id: String,
now: DateTime<Utc>,
},
}
/// Cheap, cloneable handle to the log actor. Held by the plugin manager and the
/// maintenance task; `subscribe_logs` hands receivers to SSE clients.
pub struct PluginLogStore {
tx: mpsc::Sender<LogCommand>,
live: broadcast::Sender<PluginLogEvent>,
}
impl PluginLogStore {
/// Spawn the actor thread and return a handle. `default_retention` is applied
/// to any plugin lacking an explicit `retention.json`. `queue_capacity`
/// bounds the command channel — a flood past it sheds the oldest-arriving
/// batch rather than growing RAM or blocking dispatch.
pub fn new(
root: PathBuf,
max_file_bytes: u64,
max_segments: u32,
default_retention: RetentionSettings,
queue_capacity: usize,
) -> Self {
let (tx, rx) = mpsc::channel(queue_capacity.max(1));
let (live, _) = broadcast::channel(LIVE_CAPACITY);
let actor = Actor {
root,
max_file_bytes: max_file_bytes.max(1),
max_segments,
default_retention,
writers: HashMap::new(),
live: live.clone(),
};
// A dedicated OS thread so the blocking file/gzip IO never stalls a tokio
// worker. `blocking_recv` is valid here (no runtime on this thread).
std::thread::Builder::new()
.name("plugin-log-store".into())
.spawn(move || actor.run(rx))
.expect("spawn plugin-log-store thread");
Self { tx, live }
}
/// Enqueue a batch (plugin-emitted lines + the host outcome) for one
/// invocation. Called from the dispatch `spawn_blocking` closure. Uses a
/// non-blocking `try_send`: under flood it sheds the batch (logged) rather
/// than blocking the blocking-pool thread or growing RAM unboundedly. A full
/// queue or a gone actor is swallowed — logging must never break dispatch.
pub fn append(
&self,
plugin_id: &str,
invocation_id: &str,
lines: &[(String, String)],
outcome: &InvokeOutcome,
) {
let ts = Utc::now().to_rfc3339();
let mut entries: Vec<LogEntry> = lines
.iter()
.map(|(level, msg)| LogEntry {
ts: ts.clone(),
invocation_id: invocation_id.to_string(),
kind: "plugin".to_string(),
level: level.clone(),
reason: None,
msg: msg.clone(),
})
.collect();
let (level, msg) = outcome.log_detail();
entries.push(LogEntry {
ts,
invocation_id: invocation_id.to_string(),
kind: "outcome".to_string(),
level: level.to_string(),
reason: Some(outcome.reason().to_string()),
msg,
});
if let Err(e) = self.tx.try_send(LogCommand::Append {
plugin_id: plugin_id.to_string(),
entries,
}) {
tracing::warn!(
target: "oxicloud::plugins",
plugin_id = %plugin_id,
error = %e,
"dropping plugin log batch: queue full or log actor unavailable"
);
}
}
/// Read a filtered, paginated page of a plugin's entries (newest first).
pub async fn read_page(&self, plugin_id: &str, query: LogQuery) -> LogPage {
let (reply, rx) = oneshot::channel();
if self
.tx
.send(LogCommand::Read {
plugin_id: plugin_id.to_string(),
query,
reply,
})
.await
.is_err()
{
return LogPage {
entries: Vec::new(),
total: 0,
};
}
rx.await.unwrap_or(LogPage {
entries: Vec::new(),
total: 0,
})
}
/// Delete a plugin's log files (keeps `retention.json`).
pub async fn clear(&self, plugin_id: &str) {
let (reply, rx) = oneshot::channel();
if self
.tx
.send(LogCommand::Clear {
plugin_id: plugin_id.to_string(),
reply,
})
.await
.is_ok()
{
let _ = rx.await;
}
}
/// Delete a plugin's entire log directory (on uninstall). Fire-and-forget
/// and non-blocking, so it's safe to call from the synchronous management
/// path without stalling an async worker.
pub fn remove_plugin_logs(&self, plugin_id: &str) {
let _ = self.tx.try_send(LogCommand::Remove {
plugin_id: plugin_id.to_string(),
});
}
/// The plugin's effective retention (override or configured default).
pub async fn get_retention(&self, plugin_id: &str) -> RetentionSettings {
let (reply, rx) = oneshot::channel();
if self
.tx
.send(LogCommand::GetRetention {
plugin_id: plugin_id.to_string(),
reply,
})
.await
.is_ok()
&& let Ok(s) = rx.await
{
return s;
}
// Fall back to a conservative default if the actor is gone.
RetentionSettings {
retention_days: 30,
max_bytes: 256 * 1024 * 1024,
}
}
/// Persist a per-plugin retention override.
pub async fn set_retention(&self, plugin_id: &str, settings: RetentionSettings) {
let (reply, rx) = oneshot::channel();
if self
.tx
.send(LogCommand::SetRetention {
plugin_id: plugin_id.to_string(),
settings,
reply,
})
.await
.is_ok()
{
let _ = rx.await;
}
}
/// Ask the actor to prune a plugin's segments by age + aggregate size.
pub async fn request_sweep(&self, plugin_id: &str, now: DateTime<Utc>) {
let _ = self
.tx
.send(LogCommand::Sweep {
plugin_id: plugin_id.to_string(),
now,
})
.await;
}
/// Subscribe to newly-written entries across all plugins (live tailing).
pub fn subscribe(&self) -> broadcast::Receiver<PluginLogEvent> {
self.live.subscribe()
}
}
/// The single owner of all log-file state. Runs on its own thread.
struct Actor {
root: PathBuf,
max_file_bytes: u64,
max_segments: u32,
default_retention: RetentionSettings,
writers: HashMap<String, FileRotate<AppendTimestamp>>,
live: broadcast::Sender<PluginLogEvent>,
}
impl Actor {
fn run(mut self, mut rx: mpsc::Receiver<LogCommand>) {
while let Some(cmd) = rx.blocking_recv() {
match cmd {
LogCommand::Append { plugin_id, entries } => {
self.handle_append(&plugin_id, entries)
}
LogCommand::Read {
plugin_id,
query,
reply,
} => {
let _ = reply.send(self.read_page(&plugin_id, &query));
}
LogCommand::Clear { plugin_id, reply } => {
self.clear(&plugin_id);
let _ = reply.send(());
}
LogCommand::Remove { plugin_id } => self.remove(&plugin_id),
LogCommand::GetRetention { plugin_id, reply } => {
let _ = reply.send(self.get_retention(&plugin_id));
}
LogCommand::SetRetention {
plugin_id,
settings,
reply,
} => {
self.set_retention(&plugin_id, settings);
let _ = reply.send(());
}
LogCommand::Sweep { plugin_id, now } => self.sweep(&plugin_id, now),
}
}
}
fn plugin_dir(&self, plugin_id: &str) -> PathBuf {
self.root.join(plugin_id)
}
/// Lazily build (or fetch) the rotating writer for a plugin.
fn writer_for(&mut self, plugin_id: &str) -> Option<&mut FileRotate<AppendTimestamp>> {
if !self.writers.contains_key(plugin_id) {
let path = self.plugin_dir(plugin_id).join(ACTIVE_FILE);
let writer = FileRotate::new(
path,
AppendTimestamp::default(FileLimit::MaxFiles(self.max_segments as usize)),
ContentLimit::BytesSurpassed(self.max_file_bytes as usize),
Compression::OnRotate(0),
#[cfg(unix)]
None,
);
self.writers.insert(plugin_id.to_string(), writer);
}
self.writers.get_mut(plugin_id)
}
fn handle_append(&mut self, plugin_id: &str, entries: Vec<LogEntry>) {
let mut buf = Vec::new();
for entry in &entries {
if serde_json::to_writer(&mut buf, entry).is_ok() {
buf.push(b'\n');
}
}
if let Some(writer) = self.writer_for(plugin_id)
&& let Err(e) = writer.write_all(&buf).and_then(|_| writer.flush())
{
tracing::warn!(
target: "oxicloud::plugins",
plugin_id = %plugin_id,
error = %e,
"failed to write plugin log batch"
);
return;
}
// Publish only after the durable write, so the live tail never shows an
// entry a subsequent read wouldn't. No subscribers => send is a no-op.
for entry in entries {
let _ = self.live.send(PluginLogEvent {
plugin_id: plugin_id.to_string(),
entry,
});
}
}
fn read_page(&self, plugin_id: &str, query: &LogQuery) -> LogPage {
let dir = self.plugin_dir(plugin_id);
// Gather rotated segments oldest→newest (by mtime), then the active file.
let mut segments: Vec<(PathBuf, SystemTime)> = Vec::new();
let mut active: Option<PathBuf> = None;
if let Ok(read_dir) = fs::read_dir(&dir) {
for entry in read_dir.flatten() {
let path = entry.path();
let Some(name) = path.file_name().and_then(|n| n.to_str()) else {
continue;
};
if name == ACTIVE_FILE {
active = Some(path);
} else if name.starts_with("events.jsonl.") {
let mtime = entry
.metadata()
.and_then(|m| m.modified())
.unwrap_or(SystemTime::UNIX_EPOCH);
segments.push((path, mtime));
}
}
}
segments.sort_by_key(|(_, mtime)| *mtime);
let mut all: Vec<LogEntry> = Vec::new();
for (path, _) in &segments {
read_entries_into(path, query, &mut all);
}
if let Some(path) = &active {
read_entries_into(path, query, &mut all);
}
// `all` is chronological (oldest→newest); the viewer wants newest first.
all.reverse();
let total = all.len();
let entries = all
.into_iter()
.skip(query.offset)
.take(query.limit)
.collect();
LogPage { entries, total }
}
fn clear(&mut self, plugin_id: &str) {
// Drop the open writer first so the active file can be removed cleanly.
self.writers.remove(plugin_id);
let dir = self.plugin_dir(plugin_id);
if let Ok(read_dir) = fs::read_dir(&dir) {
for entry in read_dir.flatten() {
let path = entry.path();
if let Some(name) = path.file_name().and_then(|n| n.to_str())
&& (name == ACTIVE_FILE || name.starts_with("events.jsonl."))
{
let _ = fs::remove_file(&path);
}
}
}
}
fn remove(&mut self, plugin_id: &str) {
self.writers.remove(plugin_id);
let _ = fs::remove_dir_all(self.plugin_dir(plugin_id));
}
fn get_retention(&self, plugin_id: &str) -> RetentionSettings {
let path = self.plugin_dir(plugin_id).join(RETENTION_FILE);
fs::read_to_string(&path)
.ok()
.and_then(|s| serde_json::from_str::<RetentionSettings>(&s).ok())
.unwrap_or(self.default_retention)
}
fn set_retention(&self, plugin_id: &str, settings: RetentionSettings) {
let dir = self.plugin_dir(plugin_id);
if let Err(e) = fs::create_dir_all(&dir) {
tracing::warn!(
target: "oxicloud::plugins",
plugin_id = %plugin_id, error = %e,
"failed to create plugin log dir for retention"
);
return;
}
match serde_json::to_string_pretty(&settings) {
Ok(json) => {
if let Err(e) = fs::write(dir.join(RETENTION_FILE), json) {
tracing::warn!(
target: "oxicloud::plugins",
plugin_id = %plugin_id, error = %e,
"failed to persist plugin retention"
);
}
}
Err(e) => tracing::warn!(
target: "oxicloud::plugins",
plugin_id = %plugin_id, error = %e,
"failed to serialize plugin retention"
),
}
}
/// Prune rotated segments older than the plugin's retention window, then
/// enforce the aggregate byte cap (oldest deleted first). Never touches the
/// active file.
fn sweep(&self, plugin_id: &str, now: DateTime<Utc>) {
let dir = self.plugin_dir(plugin_id);
let retention = self.get_retention(plugin_id);
let cutoff = now - Duration::days(retention.retention_days as i64);
let mut segments: Vec<(PathBuf, SystemTime, u64)> = Vec::new();
let Ok(read_dir) = fs::read_dir(&dir) else {
return;
};
for entry in read_dir.flatten() {
let path = entry.path();
let Some(name) = path.file_name().and_then(|n| n.to_str()) else {
continue;
};
if !name.starts_with("events.jsonl.") {
continue; // skip the active file, retention.json, etc.
}
let Ok(meta) = entry.metadata() else { continue };
let mtime = meta.modified().unwrap_or(SystemTime::UNIX_EPOCH);
segments.push((path, mtime, meta.len()));
}
// 1) Age-based pruning.
let mut purged = 0u64;
segments.retain(|(path, mtime, _)| {
let dt: DateTime<Utc> = (*mtime).into();
if dt < cutoff {
let _ = fs::remove_file(path);
purged += 1;
false
} else {
true
}
});
// 2) Aggregate byte cap, oldest deleted first.
segments.sort_by_key(|(_, mtime, _)| *mtime);
let mut total: u64 = segments.iter().map(|(_, _, size)| *size).sum();
let mut idx = 0;
while total > retention.max_bytes && idx < segments.len() {
let (path, _, size) = &segments[idx];
let _ = fs::remove_file(path);
total = total.saturating_sub(*size);
purged += 1;
idx += 1;
}
if purged > 0 {
tracing::debug!(
target: "oxicloud::plugins",
plugin_id = %plugin_id,
purged,
"plugin log retention sweep removed segments"
);
}
}
}
/// Read one segment (gzip if `.gz`, else plain), parse each line as a
/// [`LogEntry`], apply the filter, and append matches to `out`. Malformed lines
/// are skipped — a torn final line never aborts a read.
fn read_entries_into(path: &Path, query: &LogQuery, out: &mut Vec<LogEntry>) {
let Ok(file) = fs::File::open(path) else {
return;
};
let content = if path.extension().and_then(|e| e.to_str()) == Some("gz") {
let mut s = String::new();
if GzDecoder::new(file).read_to_string(&mut s).is_err() {
return;
}
s
} else {
let mut s = String::new();
let mut file = file;
if file.read_to_string(&mut s).is_err() {
return;
}
s
};
let search = query.search.as_ref().map(|s| s.to_lowercase());
for line in content.lines() {
if line.trim().is_empty() {
continue;
}
let Ok(entry) = serde_json::from_str::<LogEntry>(line) else {
continue;
};
if let Some(level) = &query.level
&& &entry.level != level
{
continue;
}
if let Some(needle) = &search
&& !entry.msg.to_lowercase().contains(needle)
{
continue;
}
out.push(entry);
}
}
#[cfg(test)]
mod tests {
use super::*;
fn settings(days: u32, max_bytes: u64) -> RetentionSettings {
RetentionSettings {
retention_days: days,
max_bytes,
}
}
fn entry(level: &str, msg: &str) -> LogEntry {
LogEntry {
ts: Utc::now().to_rfc3339(),
invocation_id: "inv".into(),
kind: "plugin".into(),
level: level.into(),
reason: None,
msg: msg.into(),
}
}
fn new_actor(root: PathBuf, max_file_bytes: u64, max_segments: u32) -> Actor {
let (live, _) = broadcast::channel(16);
Actor {
root,
max_file_bytes: max_file_bytes.max(1),
max_segments,
default_retention: settings(30, 1 << 30),
writers: HashMap::new(),
live,
}
}
#[test]
fn ordering_and_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let mut actor = new_actor(dir.path().to_path_buf(), 1 << 20, 5);
for i in 0..50 {
actor.handle_append("p", vec![entry("info", &format!("line {i}"))]);
}
let page = actor.read_page(
"p",
&LogQuery {
level: None,
search: None,
offset: 0,
limit: 10,
},
);
assert_eq!(page.total, 50);
assert_eq!(page.entries.len(), 10);
// Newest first.
assert_eq!(page.entries[0].msg, "line 49");
assert_eq!(page.entries[9].msg, "line 40");
}
#[test]
fn filter_by_level_and_search() {
let dir = tempfile::tempdir().unwrap();
let mut actor = new_actor(dir.path().to_path_buf(), 1 << 20, 5);
actor.handle_append("p", vec![entry("info", "hello world")]);
actor.handle_append("p", vec![entry("error", "BOOM failure")]);
actor.handle_append("p", vec![entry("info", "another HELLO")]);
let q = LogQuery {
level: Some("info".into()),
search: Some("hello".into()),
offset: 0,
limit: 100,
};
let page = actor.read_page("p", &q);
assert_eq!(page.total, 2);
assert!(page.entries.iter().all(|e| e.level == "info"));
assert!(
page.entries
.iter()
.all(|e| e.msg.to_lowercase().contains("hello"))
);
}
#[test]
fn rotation_creates_compressed_segments() {
let dir = tempfile::tempdir().unwrap();
// Tiny byte cap forces frequent rotation; a high segment cap keeps every
// segment so the cross-segment read can be checked end to end (the byte
// cap is then exercised separately by `sweep_age_and_size`).
let mut actor = new_actor(dir.path().to_path_buf(), 256, 100_000);
for i in 0..200 {
actor.handle_append(
"p",
vec![entry("info", &format!("padding line number {i}"))],
);
}
let plugin_dir = dir.path().join("p");
let gz = fs::read_dir(&plugin_dir)
.unwrap()
.flatten()
.filter(|e| {
e.path()
.extension()
.and_then(|x| x.to_str())
.map(|x| x == "gz")
.unwrap_or(false)
})
.count();
assert!(gz > 0, "expected at least one rotated .gz segment");
// All originally-written lines must still be readable across segments.
let page = actor.read_page(
"p",
&LogQuery {
level: None,
search: None,
offset: 0,
limit: 1000,
},
);
assert_eq!(page.total, 200);
}
#[test]
fn retention_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let actor = new_actor(dir.path().to_path_buf(), 1 << 20, 5);
assert_eq!(actor.get_retention("p").retention_days, 30); // default
actor.set_retention("p", settings(7, 1234));
let r = actor.get_retention("p");
assert_eq!(r.retention_days, 7);
assert_eq!(r.max_bytes, 1234);
}
#[test]
fn sweep_age_and_size() {
let dir = tempfile::tempdir().unwrap();
let plugin_dir = dir.path().join("p");
fs::create_dir_all(&plugin_dir).unwrap();
// One "old" rotated segment and one "fresh" one.
let old = plugin_dir.join("events.jsonl.20200101T000000.gz");
let fresh = plugin_dir.join("events.jsonl.20990101T000000.gz");
fs::write(&old, b"x").unwrap();
fs::write(&fresh, b"y").unwrap();
// Backdate the "old" file's mtime well past the retention window.
let long_ago = SystemTime::UNIX_EPOCH + std::time::Duration::from_secs(1_000_000);
filetime_set(&old, long_ago);
let actor = new_actor(dir.path().to_path_buf(), 1 << 20, 5);
// 1-day retention: the backdated file must go, the fresh one stays.
actor.sweep("p", Utc::now());
assert!(!old.exists(), "age-expired segment should be purged");
assert!(fresh.exists(), "recent segment should be kept");
}
/// Minimal mtime setter for tests (no extra dep): rewrite + set via filetime
/// is unavailable, so emulate "old" by relying on a very old written time is
/// not possible portably; instead we set it through `fs` utimes if present.
fn filetime_set(path: &Path, when: SystemTime) {
// `set_file_mtime` isn't in std; approximate by opening and using the
// platform fallback: on failure the test still meaningfully exercises
// the size path. We use a best-effort via `File::set_modified` (1.75+).
if let Ok(f) = fs::OpenOptions::new().write(true).open(path) {
let _ = f.set_modified(when);
}
}
}
@@ -0,0 +1,594 @@
//! Plugin discovery + dispatch + admin management. Implements
//! [`PluginDispatchPort`] and [`PluginManagementPort`] over the Extism
//! [`PluginRuntime`].
//!
//! Discovery scans a directory of plugin subdirectories (each `plugin.toml` +
//! `.wasm`) at startup; a plugin that fails validation or load is audit-logged
//! and skipped, never fatal. Dispatch builds a fresh sandbox per invocation on
//! the blocking pool, so a slow or hostile plugin never stalls async workers or
//! the upload path that triggered it.
//!
//! The same in-memory plugin set backs both ports, guarded by an `RwLock`: a
//! management op (install / toggle / remove) takes the write lock and is
//! reflected on the live dispatch path with no restart. Enable/disable state is
//! persisted as a `.disabled` marker file in the plugin's own directory so it
//! survives a restart without a database.
use std::collections::HashSet;
use std::path::{Path, PathBuf};
use std::sync::{Arc, RwLock};
use std::time::Duration;
use async_trait::async_trait;
use tokio::sync::Semaphore;
use super::log_store::PluginLogStore;
use super::manifest;
use super::runtime::{InvokeOutcome, PluginRuntime};
use crate::application::ports::plugin_ports::{
LogPage, LogQuery, OXICLOUD_PLUGIN_ABI, PluginContext, PluginDispatchPort, PluginEvent,
PluginInfo, PluginInput, PluginLogEvent, PluginManagementPort, PluginMgmtError,
RetentionSettings, event_export_name,
};
use crate::common::config::PluginConfig;
use tokio::sync::broadcast;
/// Name of the marker file that, when present in a plugin's directory, loads it
/// disabled. Created/removed by [`PluginManagementPort::set_enabled`].
const DISABLED_MARKER: &str = ".disabled";
/// A validated, loadable plugin held in memory.
struct LoadedPlugin {
id: String,
name: String,
version: String,
abi: u32,
subscribe: HashSet<String>,
/// Whether dispatch delivers events to this plugin. Mirrors the on-disk
/// `.disabled` marker.
enabled: bool,
/// The plugin's own directory (not necessarily named after `id`). Used to
/// write the disabled marker and to delete the plugin on removal.
dir: PathBuf,
runtime: Arc<PluginRuntime>,
}
impl LoadedPlugin {
fn info(&self) -> PluginInfo {
let mut subscriptions: Vec<String> = self.subscribe.iter().cloned().collect();
subscriptions.sort();
PluginInfo {
id: self.id.clone(),
name: self.name.clone(),
version: self.version.clone(),
abi: self.abi,
subscriptions,
enabled: self.enabled,
}
}
}
/// Owns all loaded plugins and dispatches events to them.
pub struct ExtismPluginManager {
config: PluginConfig,
/// Root directory plugins are discovered in and installed into.
root_dir: PathBuf,
plugins: RwLock<Vec<LoadedPlugin>>,
/// Per-plugin structured log storage (shared with the maintenance task).
log_store: Arc<PluginLogStore>,
/// Caps concurrent plugin invocations across all plugins so dispatch can
/// shed load instead of flooding the shared blocking pool.
invocation_sem: Arc<Semaphore>,
}
impl ExtismPluginManager {
/// Scan `dir` for plugins and build a manager from those that validate and
/// load. Returns an empty manager (logging the cause) if `dir` is absent or
/// unreadable — a missing plugins directory is normal, not an error.
pub fn load_from_dir(config: PluginConfig, dir: &Path) -> Self {
// The log root is a sibling of the plugins dir by default (or the
// configured override); it lives outside any individual plugin dir so a
// plugin uninstall (`remove_dir_all`) never wipes another's logs.
let log_dir = config
.log_dir
.clone()
.unwrap_or_else(|| dir.join(".plugin-logs"));
let log_store = Arc::new(PluginLogStore::new(
log_dir.clone(),
config.log_max_file_bytes,
config.log_max_segments,
RetentionSettings {
retention_days: config.log_retention_days,
max_bytes: config.log_total_max_bytes,
},
config.log_queue_capacity,
));
let invocation_sem = Arc::new(Semaphore::new(config.max_concurrent_invocations.max(1)));
let mut plugins = Vec::new();
let mut rejected = 0usize;
let entries = match std::fs::read_dir(dir) {
Ok(e) => e,
Err(e) => {
tracing::info!(
target: "oxicloud::plugins",
dir = %dir.display(),
error = %e,
"plugins directory not readable; no plugins loaded"
);
return Self {
config,
root_dir: dir.to_path_buf(),
plugins: RwLock::new(plugins),
log_store,
invocation_sem,
};
}
};
for entry in entries.flatten() {
let path = entry.path();
if !path.is_dir() {
continue;
}
// Never treat the log root as a plugin directory.
if path == log_dir {
continue;
}
match Self::load_one(&config, &path) {
Ok(loaded) => {
tracing::info!(
target: "oxicloud::plugins",
plugin_id = %loaded.id,
enabled = loaded.enabled,
dir = %path.display(),
"plugin loaded"
);
plugins.push(loaded);
}
Err(reason) => {
rejected += 1;
tracing::warn!(
target: "audit",
event = "plugin.load_rejected",
reason = reason,
plugin_dir = %path.display(),
"👮🏻‍♂️ plugin rejected at load"
);
}
}
}
tracing::info!(
target: "oxicloud::plugins",
loaded = plugins.len(),
rejected,
dir = %dir.display(),
"plugin discovery complete"
);
Self {
config,
root_dir: dir.to_path_buf(),
plugins: RwLock::new(plugins),
log_store,
invocation_sem,
}
}
/// The shared log store, handed to the maintenance task by DI.
pub fn log_store(&self) -> Arc<PluginLogStore> {
self.log_store.clone()
}
/// Validate and load a single plugin directory. Returns a stable audit
/// `reason` key on rejection.
fn load_one(config: &PluginConfig, dir: &Path) -> Result<LoadedPlugin, &'static str> {
let manifest_path = dir.join("plugin.toml");
if !manifest_path.exists() {
return Err("no_manifest");
}
let toml_str =
std::fs::read_to_string(&manifest_path).map_err(|_| "manifest_unreadable")?;
let manifest = manifest::parse_and_validate(&toml_str).map_err(|e| e.reason())?;
// The entrypoint becomes a path joined onto the plugin dir; reject a
// traversal-unsafe value on disk too (mirrors the `install` check), so a
// hand-placed manifest can't read a `.wasm` outside its own directory.
if !is_safe_component(&manifest.plugin.entrypoint) {
return Err("bad_entrypoint");
}
let wasm_path = dir.join(&manifest.plugin.entrypoint);
let wasm_bytes = std::fs::read(&wasm_path).map_err(|_| "wasm_unreadable")?;
let runtime = PluginRuntime::new(manifest.plugin.id.clone(), wasm_bytes);
// Probe a throwaway instance: abi must match AND every subscribed event
// must have its `on_<event>` handler exported.
let required_exports: Vec<String> = manifest
.events
.subscribe
.iter()
.map(|e| event_export_name(e))
.collect();
Self::probe(config, &runtime, &required_exports)?;
Ok(LoadedPlugin {
id: manifest.plugin.id,
name: manifest.plugin.name,
version: manifest.plugin.version,
abi: manifest.plugin.abi,
subscribe: manifest.events.subscribe.into_iter().collect(),
enabled: !dir.join(DISABLED_MARKER).exists(),
dir: dir.to_path_buf(),
runtime: Arc::new(runtime),
})
}
/// Probe loadability, mapping the runtime outcome to a stable reason key.
fn probe(
config: &PluginConfig,
runtime: &PluginRuntime,
required_exports: &[String],
) -> Result<(), &'static str> {
match runtime.check_loadable(config, required_exports) {
InvokeOutcome::Ok => Ok(()),
InvokeOutcome::AbiMismatch { .. } => Err("abi_mismatch"),
InvokeOutcome::MissingExport(_) => Err("missing_export"),
_ => Err("not_loadable"),
}
}
/// Number of successfully loaded plugins (used by DI for the startup summary
/// and by tests).
pub fn loaded_count(&self) -> usize {
self.read_plugins().len()
}
/// Drop the cached compiled module of every plugin idle past the configured
/// TTL, reclaiming memory. Driven by a periodic timer in DI; the next event
/// to a freed plugin recompiles transparently.
pub fn evict_idle_compiled(&self) {
let ttl = Duration::from_secs(self.config.cache_idle_ttl_secs);
let mut evicted = 0usize;
for plugin in self.read_plugins().iter() {
if plugin.runtime.evict_if_idle(ttl) {
evicted += 1;
}
}
if evicted > 0 {
tracing::debug!(
target: "oxicloud::plugins",
evicted,
"evicted idle compiled plugin modules"
);
}
}
fn read_plugins(&self) -> std::sync::RwLockReadGuard<'_, Vec<LoadedPlugin>> {
self.plugins.read().unwrap_or_else(|e| e.into_inner())
}
fn write_plugins(&self) -> std::sync::RwLockWriteGuard<'_, Vec<LoadedPlugin>> {
self.plugins.write().unwrap_or_else(|e| e.into_inner())
}
/// `NotFound` unless a plugin with this id is currently installed. Checked
/// before any log-file access so an HTTP-supplied id can't reach the
/// filesystem for a plugin that doesn't exist.
fn ensure_installed(&self, id: &str) -> Result<(), PluginMgmtError> {
if self.read_plugins().iter().any(|p| p.id == id) {
Ok(())
} else {
Err(PluginMgmtError::NotFound)
}
}
}
impl PluginDispatchPort for ExtismPluginManager {
fn dispatch(&self, event: PluginEvent) {
for plugin in self.read_plugins().iter() {
if !plugin.enabled || !plugin.subscribe.contains(event.name) {
continue;
}
let input = PluginInput {
abi: OXICLOUD_PLUGIN_ABI,
event: event.name.to_string(),
context: PluginContext {
plugin_id: plugin.id.clone(),
user_id: event.user_id.clone(),
invocation_id: event.invocation_id.clone(),
},
payload: event.payload.clone(),
};
let input_json = match serde_json::to_string(&input) {
Ok(j) => j,
Err(e) => {
tracing::warn!(
target: "oxicloud::plugins",
plugin_id = %plugin.id,
error = %e,
"failed to serialize plugin input; skipping"
);
continue;
}
};
// Load shedding: cap concurrent invocations so a flood of events (or
// slow plugins) can't exhaust the shared blocking pool. Past the cap
// the event is dropped — plugins are observe-only, so shedding is
// safe; we just record it.
let permit = match self.invocation_sem.clone().try_acquire_owned() {
Ok(p) => p,
Err(_) => {
tracing::warn!(
target: "audit",
event = "plugin.dispatch_shed",
reason = "at_capacity",
plugin_id = %plugin.id,
invocation_id = %event.invocation_id,
plugin_event = %event.name,
"👮🏻‍♂️ plugin event dropped: invocation limit reached"
);
continue;
}
};
let runtime = plugin.runtime.clone();
let config = self.config.clone();
let plugin_id = plugin.id.clone();
let invocation_id = event.invocation_id.clone();
let export = event_export_name(event.name);
let log_store = self.log_store.clone();
// Run the synchronous wasm call off the async workers. Fire-and-forget:
// the upload already succeeded; plugins are post-hoc observers.
tokio::task::spawn_blocking(move || {
// Hold the permit for the lifetime of the invocation.
let _permit = permit;
let result = runtime.invoke(&config, &export, &invocation_id, &input_json);
// Persist every invocation (the plugin's own log lines plus the
// host outcome) to the plugin's structured log. Ordered, async,
// and non-fatal — a failed write never affects the request.
log_store.append(&plugin_id, &invocation_id, &result.logs, &result.outcome);
if !result.outcome.is_ok() {
tracing::warn!(
target: "audit",
event = "plugin.invocation_failed",
reason = result.outcome.reason(),
plugin_id = %plugin_id,
invocation_id = %invocation_id,
detail = ?result.outcome,
"👮🏻‍♂️ plugin invocation failed"
);
}
});
}
}
fn has_subscribers(&self, event: &str) -> bool {
self.read_plugins()
.iter()
.any(|p| p.enabled && p.subscribe.contains(event))
}
}
#[async_trait]
impl PluginManagementPort for ExtismPluginManager {
fn list(&self) -> Vec<PluginInfo> {
let mut infos: Vec<PluginInfo> = self.read_plugins().iter().map(|p| p.info()).collect();
infos.sort_by(|a, b| a.id.cmp(&b.id));
infos
}
fn set_enabled(&self, id: &str, enabled: bool) -> Result<(), PluginMgmtError> {
let mut plugins = self.write_plugins();
let plugin = plugins
.iter_mut()
.find(|p| p.id == id)
.ok_or(PluginMgmtError::NotFound)?;
let marker = plugin.dir.join(DISABLED_MARKER);
if enabled {
match std::fs::remove_file(&marker) {
Ok(()) => {}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => return Err(PluginMgmtError::Io(e.to_string())),
}
} else {
std::fs::write(&marker, b"").map_err(|e| PluginMgmtError::Io(e.to_string()))?;
}
plugin.enabled = enabled;
Ok(())
}
fn install(&self, manifest_toml: &str, wasm: Vec<u8>) -> Result<PluginInfo, PluginMgmtError> {
// Validate the manifest and the wasm before touching the filesystem.
let manifest = manifest::parse_and_validate(manifest_toml)
.map_err(|e| PluginMgmtError::Rejected(e.reason()))?;
// `id` becomes a directory name and `entrypoint` a filename — both must
// be single, traversal-free path components.
if !is_safe_component(&manifest.plugin.id) {
return Err(PluginMgmtError::Rejected("bad_id"));
}
if !is_safe_component(&manifest.plugin.entrypoint) {
return Err(PluginMgmtError::Rejected("bad_entrypoint"));
}
let required_exports: Vec<String> = manifest
.events
.subscribe
.iter()
.map(|e| event_export_name(e))
.collect();
let runtime = PluginRuntime::new(manifest.plugin.id.clone(), wasm.clone());
Self::probe(&self.config, &runtime, &required_exports)
.map_err(PluginMgmtError::Rejected)?;
let id = manifest.plugin.id.clone();
let target = self.root_dir.join(&id);
// Hold the write lock across the collision check and the directory swap
// so two concurrent installs of the same id cannot race. Admin installs
// are rare; readers block only briefly.
let mut plugins = self.write_plugins();
if plugins.iter().any(|p| p.id == id) || target.exists() {
return Err(PluginMgmtError::IdExists);
}
std::fs::create_dir_all(&self.root_dir).map_err(|e| PluginMgmtError::Io(e.to_string()))?;
// Write to a temp dir then rename, so a crash mid-write never leaves a
// half-written plugin discoverable.
let tmp = tempfile::Builder::new()
.prefix(".tmp-install-")
.tempdir_in(&self.root_dir)
.map_err(|e| PluginMgmtError::Io(e.to_string()))?;
std::fs::write(tmp.path().join("plugin.toml"), manifest_toml)
.map_err(|e| PluginMgmtError::Io(e.to_string()))?;
std::fs::write(tmp.path().join(&manifest.plugin.entrypoint), &wasm)
.map_err(|e| PluginMgmtError::Io(e.to_string()))?;
let tmp_path = tmp.keep();
if let Err(e) = std::fs::rename(&tmp_path, &target) {
let _ = std::fs::remove_dir_all(&tmp_path);
return Err(PluginMgmtError::Io(e.to_string()));
}
let loaded = LoadedPlugin {
id: id.clone(),
name: manifest.plugin.name.clone(),
version: manifest.plugin.version.clone(),
abi: manifest.plugin.abi,
subscribe: manifest.events.subscribe.iter().cloned().collect(),
enabled: true,
dir: target,
runtime: Arc::new(runtime),
};
let info = loaded.info();
plugins.push(loaded);
Ok(info)
}
fn install_bundle(&self, zip: Vec<u8>) -> Result<PluginInfo, PluginMgmtError> {
use std::io::{Cursor, Read};
// Aggregate decompressed ceiling, enforced as each entry is unpacked so
// a zip bomb can't blow up memory before validation (the route also caps
// the compressed body). We only ever extract two named entries.
let max_decompressed: u64 = self.config.max_bundle_decompressed_bytes;
let mut archive = zip::ZipArchive::new(Cursor::new(zip))
.map_err(|_| PluginMgmtError::Rejected("bad_zip"))?;
// Locate `plugin.toml` — at the archive root or under a single wrapping
// folder (e.g. `myplugin/plugin.toml`).
let manifest_name = archive
.file_names()
.find(|n| !n.ends_with('/') && (*n == "plugin.toml" || n.ends_with("/plugin.toml")))
.map(str::to_owned)
.ok_or(PluginMgmtError::Rejected("no_manifest_in_zip"))?;
let mut manifest_toml = String::new();
{
let entry = archive
.by_name(&manifest_name)
.map_err(|_| PluginMgmtError::Rejected("no_manifest_in_zip"))?;
entry
.take(max_decompressed + 1)
.read_to_string(&mut manifest_toml)
.map_err(|_| PluginMgmtError::Rejected("bad_zip"))?;
}
if manifest_toml.len() as u64 > max_decompressed {
return Err(PluginMgmtError::Rejected("too_large"));
}
// Parse just to learn the entrypoint name; `install` does the full
// validation (and rejects a traversal-unsafe entrypoint).
let manifest = manifest::parse_and_validate(&manifest_toml)
.map_err(|e| PluginMgmtError::Rejected(e.reason()))?;
// Resolve the entrypoint relative to the manifest's folder in the zip.
let prefix = match manifest_name.rfind('/') {
Some(i) => &manifest_name[..=i],
None => "",
};
let wasm_name = format!("{prefix}{}", manifest.plugin.entrypoint);
// Budget the wasm against what the manifest already consumed.
let remaining = max_decompressed - manifest_toml.len() as u64;
let mut wasm = Vec::new();
{
let entry = archive
.by_name(&wasm_name)
.map_err(|_| PluginMgmtError::Rejected("entrypoint_not_in_zip"))?;
entry
.take(remaining + 1)
.read_to_end(&mut wasm)
.map_err(|_| PluginMgmtError::Rejected("bad_zip"))?;
}
if wasm.len() as u64 > remaining {
return Err(PluginMgmtError::Rejected("too_large"));
}
self.install(&manifest_toml, wasm)
}
fn remove(&self, id: &str) -> Result<(), PluginMgmtError> {
let mut plugins = self.write_plugins();
let pos = plugins
.iter()
.position(|p| p.id == id)
.ok_or(PluginMgmtError::NotFound)?;
let removed = plugins.remove(pos);
// Also reclaim the plugin's logs so a later reinstall of the same id
// doesn't inherit stale entries.
self.log_store.remove_plugin_logs(id);
match std::fs::remove_dir_all(&removed.dir) {
Ok(()) => Ok(()),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(e) => Err(PluginMgmtError::Io(e.to_string())),
}
}
async fn read_logs(&self, id: &str, query: LogQuery) -> Result<LogPage, PluginMgmtError> {
self.ensure_installed(id)?;
Ok(self.log_store.read_page(id, query).await)
}
async fn clear_logs(&self, id: &str) -> Result<(), PluginMgmtError> {
self.ensure_installed(id)?;
self.log_store.clear(id).await;
Ok(())
}
async fn get_retention(&self, id: &str) -> Result<RetentionSettings, PluginMgmtError> {
self.ensure_installed(id)?;
Ok(self.log_store.get_retention(id).await)
}
async fn set_retention(
&self,
id: &str,
settings: RetentionSettings,
) -> Result<(), PluginMgmtError> {
self.ensure_installed(id)?;
self.log_store.set_retention(id, settings).await;
Ok(())
}
fn subscribe_logs(&self) -> broadcast::Receiver<PluginLogEvent> {
self.log_store.subscribe()
}
}
/// Whether `s` is a single, traversal-free path component safe to use as a
/// directory or file name under the plugins root.
fn is_safe_component(s: &str) -> bool {
!s.is_empty()
&& s != "."
&& s != ".."
&& !s.contains('/')
&& !s.contains('\\')
&& !s.contains('\0')
}
@@ -0,0 +1,315 @@
//! Manager-level tests for the admin management surface (install / toggle /
//! remove) and disabled-state persistence. These drive a real Extism sandbox,
//! so they run only under `cargo test --features plugins`.
//!
//! The `.wasm` fixtures are the same ones the runtime tests use, built by
//! `scripts/build-plugin-hello.sh`.
use super::ExtismPluginManager;
use crate::application::ports::plugin_ports::{PluginDispatchPort, PluginManagementPort};
use crate::common::config::PluginConfig;
fn cfg() -> PluginConfig {
PluginConfig::default()
}
fn fixture(name: &str) -> Vec<u8> {
let path = format!(
"{}/tests/fixtures/plugins/{}",
env!("CARGO_MANIFEST_DIR"),
name
);
std::fs::read(&path).unwrap_or_else(|e| {
panic!("missing fixture {path}: {e}\n run scripts/build-plugin-hello.sh to (re)build it")
})
}
/// A valid manifest for the `hello.wasm` fixture (subscribes to both events).
fn hello_manifest() -> String {
r#"
[plugin]
id = "com.example.hello"
name = "Hello"
version = "0.1.0"
abi = 0
entrypoint = "hello.wasm"
[events]
subscribe = ["file.uploaded", "user.login"]
"#
.to_string()
}
/// A manifest that parses fine but points at the `wrong_abi.wasm` fixture,
/// which reports ABI 1 at runtime.
fn wrong_abi_manifest() -> String {
r#"
[plugin]
id = "com.example.wrongabi"
name = "Wrong ABI"
version = "0.1.0"
abi = 0
entrypoint = "wrong_abi.wasm"
[events]
subscribe = ["file.uploaded"]
"#
.to_string()
}
#[test]
fn install_loads_plugin_and_writes_files() {
let tmp = tempfile::tempdir().unwrap();
let mgr = ExtismPluginManager::load_from_dir(cfg(), tmp.path());
assert_eq!(mgr.loaded_count(), 0);
let info = mgr
.install(&hello_manifest(), fixture("hello.wasm"))
.expect("install should succeed");
assert_eq!(info.id, "com.example.hello");
assert_eq!(info.name, "Hello");
assert!(info.enabled);
assert_eq!(info.subscriptions, vec!["file.uploaded", "user.login"]);
assert_eq!(mgr.loaded_count(), 1);
let plugin_dir = tmp.path().join("com.example.hello");
assert!(plugin_dir.join("plugin.toml").exists());
assert!(plugin_dir.join("hello.wasm").exists());
// The live dispatch path sees it immediately.
assert!(mgr.has_subscribers("file.uploaded"));
}
/// Build an in-memory `.zip` with the given entries (name, bytes).
fn make_zip(entries: &[(&str, &[u8])]) -> Vec<u8> {
use std::io::Write;
use zip::write::SimpleFileOptions;
let mut writer = zip::ZipWriter::new(std::io::Cursor::new(Vec::new()));
for (name, bytes) in entries {
writer
.start_file(*name, SimpleFileOptions::default())
.unwrap();
writer.write_all(bytes).unwrap();
}
writer.finish().unwrap().into_inner()
}
#[test]
fn install_bundle_from_zip_loads_plugin() {
let tmp = tempfile::tempdir().unwrap();
let mgr = ExtismPluginManager::load_from_dir(cfg(), tmp.path());
// Wrap everything under a top-level folder to exercise prefix resolution.
let zip = make_zip(&[
("hello/plugin.toml", hello_manifest().as_bytes()),
("hello/hello.wasm", &fixture("hello.wasm")),
]);
let info = mgr
.install_bundle(zip)
.expect("bundle install should succeed");
assert_eq!(info.id, "com.example.hello");
assert!(info.enabled);
assert_eq!(mgr.loaded_count(), 1);
assert!(
tmp.path()
.join("com.example.hello")
.join("hello.wasm")
.exists()
);
}
#[test]
fn install_bundle_without_manifest_is_rejected() {
let tmp = tempfile::tempdir().unwrap();
let mgr = ExtismPluginManager::load_from_dir(cfg(), tmp.path());
let zip = make_zip(&[("hello.wasm", &fixture("hello.wasm"))]);
let err = mgr
.install_bundle(zip)
.expect_err("a zip without plugin.toml must be rejected");
assert_eq!(err.reason(), "no_manifest_in_zip");
assert_eq!(mgr.loaded_count(), 0);
}
#[test]
fn install_bundle_missing_entrypoint_is_rejected() {
let tmp = tempfile::tempdir().unwrap();
let mgr = ExtismPluginManager::load_from_dir(cfg(), tmp.path());
// Manifest declares entrypoint = "hello.wasm", but the zip omits it.
let zip = make_zip(&[("plugin.toml", hello_manifest().as_bytes())]);
let err = mgr
.install_bundle(zip)
.expect_err("a zip missing the entrypoint wasm must be rejected");
assert_eq!(err.reason(), "entrypoint_not_in_zip");
assert_eq!(mgr.loaded_count(), 0);
}
#[test]
fn install_bundle_oversized_is_rejected() {
let tmp = tempfile::tempdir().unwrap();
// Tiny decompressed ceiling so the ~130 KiB wasm fixture trips it cheaply.
let mut config = cfg();
config.max_bundle_decompressed_bytes = 1024;
let mgr = ExtismPluginManager::load_from_dir(config, tmp.path());
let zip = make_zip(&[
("plugin.toml", hello_manifest().as_bytes()),
("hello.wasm", &fixture("hello.wasm")),
]);
let err = mgr
.install_bundle(zip)
.expect_err("a bundle over the decompressed ceiling must be rejected");
assert_eq!(err.reason(), "too_large");
assert_eq!(mgr.loaded_count(), 0);
}
#[test]
fn install_bundle_with_garbage_is_rejected() {
let tmp = tempfile::tempdir().unwrap();
let mgr = ExtismPluginManager::load_from_dir(cfg(), tmp.path());
let err = mgr
.install_bundle(b"not a zip file".to_vec())
.expect_err("non-zip bytes must be rejected");
assert_eq!(err.reason(), "bad_zip");
}
#[test]
fn install_duplicate_id_is_rejected() {
let tmp = tempfile::tempdir().unwrap();
let mgr = ExtismPluginManager::load_from_dir(cfg(), tmp.path());
mgr.install(&hello_manifest(), fixture("hello.wasm"))
.unwrap();
let err = mgr
.install(&hello_manifest(), fixture("hello.wasm"))
.expect_err("second install of the same id must fail");
assert_eq!(err.reason(), "id_exists");
assert_eq!(mgr.loaded_count(), 1);
}
#[test]
fn install_wrong_abi_is_rejected() {
let tmp = tempfile::tempdir().unwrap();
let mgr = ExtismPluginManager::load_from_dir(cfg(), tmp.path());
let err = mgr
.install(&wrong_abi_manifest(), fixture("wrong_abi.wasm"))
.expect_err("a plugin reporting the wrong ABI must be rejected");
assert_eq!(err.reason(), "abi_mismatch");
assert_eq!(mgr.loaded_count(), 0);
// Nothing should have been written to disk.
assert!(!tmp.path().join("com.example.wrongabi").exists());
}
#[test]
fn disable_stops_dispatch_and_persists_across_reload() {
let tmp = tempfile::tempdir().unwrap();
let mgr = ExtismPluginManager::load_from_dir(cfg(), tmp.path());
mgr.install(&hello_manifest(), fixture("hello.wasm"))
.unwrap();
assert!(mgr.has_subscribers("file.uploaded"));
mgr.set_enabled("com.example.hello", false).unwrap();
assert!(!mgr.has_subscribers("file.uploaded"));
assert!(
tmp.path()
.join("com.example.hello")
.join(".disabled")
.exists()
);
// A fresh manager re-reads the marker and loads it disabled.
drop(mgr);
let reloaded = ExtismPluginManager::load_from_dir(cfg(), tmp.path());
assert_eq!(reloaded.loaded_count(), 1);
let info = reloaded.list();
assert_eq!(info.len(), 1);
assert!(!info[0].enabled);
assert!(!reloaded.has_subscribers("file.uploaded"));
// Re-enabling removes the marker.
reloaded.set_enabled("com.example.hello", true).unwrap();
assert!(reloaded.has_subscribers("file.uploaded"));
assert!(
!tmp.path()
.join("com.example.hello")
.join(".disabled")
.exists()
);
}
#[test]
fn set_enabled_unknown_id_is_not_found() {
let tmp = tempfile::tempdir().unwrap();
let mgr = ExtismPluginManager::load_from_dir(cfg(), tmp.path());
let err = mgr.set_enabled("does.not.exist", false).unwrap_err();
assert_eq!(err.reason(), "not_found");
}
#[test]
fn remove_unloads_and_deletes_directory() {
let tmp = tempfile::tempdir().unwrap();
let mgr = ExtismPluginManager::load_from_dir(cfg(), tmp.path());
mgr.install(&hello_manifest(), fixture("hello.wasm"))
.unwrap();
let plugin_dir = tmp.path().join("com.example.hello");
assert!(plugin_dir.exists());
mgr.remove("com.example.hello").unwrap();
assert_eq!(mgr.loaded_count(), 0);
assert!(!plugin_dir.exists());
let err = mgr.remove("com.example.hello").unwrap_err();
assert_eq!(err.reason(), "not_found");
}
/// Regression guard: dispatch must persist a log entry for *every* invocation,
/// including a successful one (it previously only logged failures). We dispatch
/// a `file.uploaded` event and then poll the plugin's structured log until an
/// `outcome` row appears.
#[tokio::test(flavor = "multi_thread")]
async fn dispatch_writes_outcome_row_on_success() {
use crate::application::ports::plugin_ports::{EVENT_FILE_UPLOADED, LogQuery, PluginEvent};
let tmp = tempfile::tempdir().unwrap();
let mgr = ExtismPluginManager::load_from_dir(cfg(), tmp.path());
mgr.install(&hello_manifest(), fixture("hello.wasm"))
.unwrap();
mgr.dispatch(PluginEvent {
name: EVENT_FILE_UPLOADED,
user_id: None,
invocation_id: "test-invocation".to_string(),
payload: serde_json::json!({ "path": "/x.txt", "size": 1, "mime": "text/plain" }),
});
// dispatch is fire-and-forget on the blocking pool, and the log write is an
// ordered async hand-off; poll until the outcome row is durable.
let mut found = false;
for _ in 0..50 {
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
let page = mgr
.read_logs(
"com.example.hello",
LogQuery {
level: None,
search: None,
offset: 0,
limit: 100,
},
)
.await
.unwrap();
if page.entries.iter().any(|e| e.kind == "outcome") {
found = true;
break;
}
}
assert!(
found,
"dispatch should persist an outcome row for the invocation"
);
}
@@ -0,0 +1,104 @@
//! `plugin.toml` parsing + load-time validation (ABI v0).
//!
//! The manifest is the host's source of truth for *what to load and when to
//! call it*. Validation fails closed: unknown sections/keys, a mismatched ABI,
//! an unknown subscribed event, or any non-empty `[permissions]` (M0 grants
//! none) all reject the plugin. A rejected plugin is skipped, never fatal.
use std::collections::BTreeMap;
use crate::application::ports::plugin_ports::{KNOWN_EVENTS, OXICLOUD_PLUGIN_ABI};
/// Parsed `plugin.toml`. `#[serde(deny_unknown_fields)]` on every struct turns
/// stray keys into load errors rather than silently ignored config.
#[derive(Debug, Clone, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct PluginManifest {
pub plugin: PluginSection,
pub events: EventsSection,
/// M0: must be empty. Any key here rejects the plugin (no grantable
/// permissions exist yet). Kept as a free map so future keys are *detected*,
/// not parsed.
#[serde(default)]
pub permissions: BTreeMap<String, toml::Value>,
}
#[derive(Debug, Clone, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct PluginSection {
/// Reverse-DNS, unique per instance.
pub id: String,
pub name: String,
/// The plugin's own semver.
pub version: String,
/// Must equal [`OXICLOUD_PLUGIN_ABI`].
pub abi: u32,
/// Path to the `.wasm`, relative to the manifest.
pub entrypoint: String,
}
#[derive(Debug, Clone, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct EventsSection {
/// Events this plugin wants. Each must be one of `KNOWN_EVENTS`
/// (`"file.uploaded"`, `"user.login"`); an unknown name rejects the plugin.
pub subscribe: Vec<String>,
}
/// Why a manifest was rejected. `reason()` yields the stable, machine-readable
/// key used in audit logs.
#[derive(Debug, thiserror::Error)]
pub enum ManifestError {
#[error("failed to parse plugin.toml: {0}")]
Parse(String),
#[error("plugin declares ABI {got}, host speaks {want}")]
AbiMismatch { got: u32, want: u32 },
#[error("events.subscribe must not be empty")]
NoEvents,
#[error("unknown event '{0}' in events.subscribe")]
UnknownEvent(String),
#[error("permissions must be empty in ABI v0 (found key '{0}')")]
PermissionsNotEmpty(String),
}
impl ManifestError {
/// Stable key for `tracing` audit lines; never reworded across releases.
pub fn reason(&self) -> &'static str {
match self {
ManifestError::Parse(_) => "parse_error",
ManifestError::AbiMismatch { .. } => "abi_mismatch",
ManifestError::NoEvents => "no_events",
ManifestError::UnknownEvent(_) => "unknown_event",
ManifestError::PermissionsNotEmpty(_) => "permissions_not_empty",
}
}
}
/// Parse and validate a `plugin.toml` body. Does not touch the `.wasm`; the
/// caller probes `abi_version` separately after a successful parse.
pub fn parse_and_validate(toml_str: &str) -> Result<PluginManifest, ManifestError> {
let manifest: PluginManifest =
toml::from_str(toml_str).map_err(|e| ManifestError::Parse(e.to_string()))?;
if manifest.plugin.abi != OXICLOUD_PLUGIN_ABI {
return Err(ManifestError::AbiMismatch {
got: manifest.plugin.abi,
want: OXICLOUD_PLUGIN_ABI,
});
}
if manifest.events.subscribe.is_empty() {
return Err(ManifestError::NoEvents);
}
for event in &manifest.events.subscribe {
if !KNOWN_EVENTS.contains(&event.as_str()) {
return Err(ManifestError::UnknownEvent(event.clone()));
}
}
if let Some((key, _)) = manifest.permissions.iter().next() {
return Err(ManifestError::PermissionsNotEmpty(key.clone()));
}
Ok(manifest)
}
@@ -0,0 +1,21 @@
//! WASM plugin runtime (Extism) — M0 walking skeleton.
//!
//! Compiled only under the `plugins` cargo feature. The application layer talks
//! to [`manager::ExtismPluginManager`] through the
//! [`crate::application::ports::plugin_ports::PluginDispatchPort`] trait, so the
//! Extism types here never leak past the infrastructure boundary.
pub mod log_retention_service;
pub mod log_store;
pub mod manager;
pub mod manifest;
pub mod runtime;
pub use log_retention_service::PluginLogMaintenanceService;
pub use log_store::PluginLogStore;
pub use manager::ExtismPluginManager;
#[cfg(test)]
mod manager_test;
#[cfg(test)]
mod runtime_test;
@@ -0,0 +1,347 @@
//! The Extism runtime wrapper — a cached compiled module, instantiated fresh per
//! invocation.
//!
//! Isolation is the point: no WASI, no filesystem, no network, a memory cap, and
//! a wall-clock timeout. The only authority a plugin has is the host `log`
//! function. Every boundary crossing is wrapped so a trap/timeout/OOM/malformed
//! output is captured as an [`InvokeOutcome`] and never propagates to the caller.
//!
//! **Compilation is amortized.** A plugin's WASM is compiled once into an
//! [`extism::CompiledPlugin`] and cached; every invocation builds a *fresh*
//! [`extism::Plugin`] instance from it (a new Store/memory → no cross-user
//! state), but pays no recompilation. Per-invocation log attribution rides
//! `call_with_host_context` rather than a baked `UserData`, so the same compiled
//! module serves concurrent invocations without sharing the log buffer. An idle
//! plugin's compiled module is dropped by [`PluginRuntime::evict_if_idle`] to
//! reclaim memory; the next event recompiles (cheaply, from wasmtime's on-disk
//! compilation cache).
use std::sync::{Arc, Mutex, RwLock};
use std::time::{Duration, Instant};
use extism::{
CompiledPlugin, CurrentPlugin, Manifest as ExtismManifest, PTR, PluginBuilder, UserData, Val,
Wasm,
};
use crate::application::ports::plugin_ports::{HOST_NAMESPACE, OXICLOUD_PLUGIN_ABI, PluginOutput};
use crate::common::config::PluginConfig;
/// Per-invocation host context, handed to one `handle` call via
/// `call_with_host_context` and read back by the `log` host function. Each
/// invocation gets its own, so a reused compiled module never mixes two
/// invocations' log lines. `lines` is an `Arc` the caller retains a clone of, to
/// read what the plugin emitted after the call returns.
struct LogSink {
plugin_id: String,
invocation_id: String,
lines: Arc<Mutex<Vec<(String, String)>>>,
}
/// The entire authority surface: log(level, message) -> (). Observe-only — it
/// reads nothing and mutates no host state beyond the per-call sink. Unknown
/// levels clamp to "info". Written without the `host_fn!` macro so it can read
/// the per-invocation [`LogSink`] from the host context.
fn oxi_log(
plugin: &mut CurrentPlugin,
inputs: &[Val],
_outputs: &mut [Val],
_user_data: UserData<()>,
) -> Result<(), extism::Error> {
let level: String = plugin.memory_get_val(&inputs[0])?;
let message: String = plugin.memory_get_val(&inputs[1])?;
let level = match level.as_str() {
"debug" | "info" | "warn" | "error" => level,
_ => "info".to_string(),
};
let ctx = plugin.host_context::<LogSink>()?;
// The message is a structured field, never interpolated into the format
// string — a plugin can't inject newlines into the operational log stream.
tracing::info!(
target: "oxicloud::plugins",
plugin_id = %ctx.plugin_id,
invocation_id = %ctx.invocation_id,
plugin_level = %level,
plugin_message = %message,
"plugin log"
);
ctx.lines
.lock()
.unwrap_or_else(|e| e.into_inner())
.push((level, message));
Ok(())
}
/// The result of one boundary crossing. Only `Ok` is a success; every other
/// variant is a contained failure the host audit-logs and moves past.
#[derive(Debug)]
pub enum InvokeOutcome {
/// `handle` returned `{"ok": true}`.
Ok,
/// `handle` returned `{"ok": false, "error": ...}`.
PluginError(String),
/// A wasm trap (panic/`unreachable`/OOM/etc.).
Trap(String),
/// The wall-clock timeout cancelled the call.
Timeout,
/// The instance could not be built (bad/unloadable wasm, unresolved import).
LoadError(String),
/// `abi_version` returned a value the host does not speak.
AbiMismatch { got: u32 },
/// A subscribed event has no matching `on_<event>` export in the module.
MissingExport(String),
/// The event handler returned bytes that are not a valid `PluginOutput`.
MalformedOutput(String),
/// The serialized input exceeded the configured cap; nothing was invoked.
MalformedInput { size: usize, max: usize },
}
impl InvokeOutcome {
pub fn is_ok(&self) -> bool {
matches!(self, InvokeOutcome::Ok)
}
/// The `(level, message)` to record for this outcome in a plugin's log file.
/// `Ok` is an `info` "completed"; every contained failure is a `warn`/`error`
/// carrying its detail. The stable machine key is [`InvokeOutcome::reason`].
pub fn log_detail(&self) -> (&'static str, String) {
match self {
InvokeOutcome::Ok => ("info", "invocation completed".to_string()),
InvokeOutcome::PluginError(e) => ("warn", e.clone()),
InvokeOutcome::Trap(e) => ("error", e.clone()),
InvokeOutcome::Timeout => ("error", "wall-clock timeout".to_string()),
InvokeOutcome::LoadError(e) => ("error", e.clone()),
InvokeOutcome::AbiMismatch { got } => {
("error", format!("abi mismatch: plugin reported {got}"))
}
InvokeOutcome::MissingExport(s) => ("error", format!("missing export: {s}")),
InvokeOutcome::MalformedOutput(e) => ("warn", e.clone()),
InvokeOutcome::MalformedInput { size, max } => {
("warn", format!("input too large: {size} bytes (max {max})"))
}
}
}
/// Stable, machine-readable key for audit logs.
pub fn reason(&self) -> &'static str {
match self {
InvokeOutcome::Ok => "ok",
InvokeOutcome::PluginError(_) => "plugin_error",
InvokeOutcome::Trap(_) => "trap",
InvokeOutcome::Timeout => "timeout",
InvokeOutcome::LoadError(_) => "load_error",
InvokeOutcome::AbiMismatch { .. } => "abi_mismatch",
InvokeOutcome::MissingExport(_) => "missing_export",
InvokeOutcome::MalformedOutput(_) => "malformed_output",
InvokeOutcome::MalformedInput { .. } => "malformed_input",
}
}
}
/// Outcome plus whatever the plugin logged (for tests and tracing).
pub struct InvokeResult {
pub outcome: InvokeOutcome,
pub logs: Vec<(String, String)>,
}
/// A loaded plugin: the wasm bytes plus a lazily-built, idle-evictable compiled
/// module. A fresh *instance* is built for every invocation (no reuse → no
/// cross-user state); only the *compilation* is shared.
pub struct PluginRuntime {
plugin_id: String,
wasm_bytes: Vec<u8>,
/// The cached compiled module, `None` until first use or after idle
/// eviction. Guarded by an `RwLock`: invocations take the read lock to
/// instantiate concurrently; (re)compilation and eviction take the write
/// lock.
compiled: RwLock<Option<CompiledPlugin>>,
/// Last time an instance was built, for idle eviction. Separate lock so it
/// can be stamped while only holding `compiled` for read.
last_used: Mutex<Instant>,
}
impl PluginRuntime {
pub fn new(plugin_id: impl Into<String>, wasm_bytes: Vec<u8>) -> Self {
Self {
plugin_id: plugin_id.into(),
wasm_bytes,
compiled: RwLock::new(None),
last_used: Mutex::new(Instant::now()),
}
}
/// Compile the WASM into a reusable [`CompiledPlugin`], wiring the sandbox
/// limits and the sole host import. wasmtime's on-disk cache (extism's
/// default) makes a repeat compile after eviction cheap.
fn compile(&self, cfg: &PluginConfig) -> Result<CompiledPlugin, extism::Error> {
let manifest = ExtismManifest::new([Wasm::data(self.wasm_bytes.clone())])
.with_memory_max(cfg.max_memory_pages) // pages × 64 KiB
.with_timeout(Duration::from_millis(cfg.invocation_timeout_ms))
.disallow_all_hosts(); // no outbound network
// No allowed_paths -> no filesystem. with_wasi(false) -> no ambient authority.
PluginBuilder::new(manifest)
.with_wasi(false)
.with_function_in_namespace(
HOST_NAMESPACE,
"log",
[PTR, PTR],
[],
UserData::new(()),
oxi_log,
)
.compile()
}
/// Build a fresh instance from the (cached, lazily-compiled) module. Stamps
/// `last_used` so the idle sweep leaves an actively-used plugin alone.
fn instantiate(&self, cfg: &PluginConfig) -> Result<extism::Plugin, InvokeOutcome> {
// Fast path: already compiled.
{
let guard = self.compiled.read().unwrap_or_else(|e| e.into_inner());
if let Some(compiled) = guard.as_ref() {
*self.last_used.lock().unwrap_or_else(|e| e.into_inner()) = Instant::now();
return extism::Plugin::new_from_compiled(compiled)
.map_err(|e| InvokeOutcome::LoadError(e.to_string()));
}
}
// Slow path: compile under the write lock (double-checked).
let mut guard = self.compiled.write().unwrap_or_else(|e| e.into_inner());
if guard.is_none() {
match self.compile(cfg) {
Ok(c) => *guard = Some(c),
Err(e) => return Err(InvokeOutcome::LoadError(e.to_string())),
}
}
let compiled = guard.as_ref().expect("compiled present after compile");
*self.last_used.lock().unwrap_or_else(|e| e.into_inner()) = Instant::now();
extism::Plugin::new_from_compiled(compiled)
.map_err(|e| InvokeOutcome::LoadError(e.to_string()))
}
/// Drop the cached compiled module if it hasn't been used within `ttl`,
/// reclaiming its memory. Returns whether anything was evicted. The next
/// invocation recompiles transparently.
pub fn evict_if_idle(&self, ttl: Duration) -> bool {
let idle = self
.last_used
.lock()
.unwrap_or_else(|e| e.into_inner())
.elapsed()
>= ttl;
if !idle {
return false;
}
let mut guard = self.compiled.write().unwrap_or_else(|e| e.into_inner());
guard.take().is_some()
}
/// Probe loadability: compile (caching the module), check `abi_version`,
/// then verify every `required_export` (the `on_<event>` symbol for each
/// subscribed event) exists. Rejects lying, unloadable, or
/// incompletely-implemented plugins before they are ever registered.
pub fn check_loadable(&self, cfg: &PluginConfig, required_exports: &[String]) -> InvokeOutcome {
let mut plugin = match self.instantiate(cfg) {
Ok(p) => p,
Err(o) => return o,
};
match plugin.call::<(), u32>("abi_version", ()) {
Ok(v) if v == OXICLOUD_PLUGIN_ABI => {}
Ok(v) => return InvokeOutcome::AbiMismatch { got: v },
Err(e) => return classify_call_error(e),
}
for export in required_exports {
if !plugin.function_exists(export) {
return InvokeOutcome::MissingExport(export.clone());
}
}
InvokeOutcome::Ok
}
/// Run one event-handler invocation, fully fault-isolated. `export` is the
/// `on_<event>` symbol to call (see `event_export_name`).
pub fn invoke(
&self,
cfg: &PluginConfig,
export: &str,
invocation_id: &str,
input_json: &str,
) -> InvokeResult {
if input_json.len() > cfg.max_input_bytes {
return InvokeResult {
outcome: InvokeOutcome::MalformedInput {
size: input_json.len(),
max: cfg.max_input_bytes,
},
logs: Vec::new(),
};
}
let lines = Arc::new(Mutex::new(Vec::new()));
let drain = || lines.lock().unwrap_or_else(|e| e.into_inner()).clone();
let mut plugin = match self.instantiate(cfg) {
Ok(p) => p,
Err(outcome) => {
return InvokeResult {
outcome,
logs: drain(),
};
}
};
// Version negotiation at the door (cheap; no recompile).
match plugin.call::<(), u32>("abi_version", ()) {
Ok(v) if v == OXICLOUD_PLUGIN_ABI => {}
Ok(v) => {
return InvokeResult {
outcome: InvokeOutcome::AbiMismatch { got: v },
logs: drain(),
};
}
Err(e) => {
return InvokeResult {
outcome: classify_call_error(e),
logs: drain(),
};
}
}
let sink = LogSink {
plugin_id: self.plugin_id.clone(),
invocation_id: invocation_id.to_string(),
lines: lines.clone(),
};
// The actual call. Traps, timeouts, and OOM all surface here as Err.
let outcome = match plugin
.call_with_host_context::<&str, String, LogSink>(export, input_json, sink)
{
Ok(out) => match serde_json::from_str::<PluginOutput>(&out) {
Ok(parsed) if parsed.ok => InvokeOutcome::Ok,
Ok(parsed) => {
InvokeOutcome::PluginError(parsed.error.unwrap_or_else(|| "unspecified".into()))
}
Err(e) => InvokeOutcome::MalformedOutput(e.to_string()),
},
Err(e) => classify_call_error(e),
};
InvokeResult {
outcome,
logs: drain(),
}
// `plugin` (instance) dropped here -> sandbox memory reclaimed. The
// compiled module stays cached for the next invocation.
}
}
/// Extism signals a wall-clock timeout with `Error::msg("timeout")`; everything
/// else from a `call` is a trap (panic, `unreachable`, OOM, etc.).
fn classify_call_error(e: extism::Error) -> InvokeOutcome {
let msg = e.to_string();
if msg.to_ascii_lowercase().contains("timeout") {
InvokeOutcome::Timeout
} else {
InvokeOutcome::Trap(msg)
}
}
@@ -0,0 +1,350 @@
//! Plugin-runtime acceptance + failure-isolation tests, plus manifest-validation
//! unit tests.
//!
//! The `.wasm` fixtures are built and committed by `scripts/build-plugin-hello.sh`
//! from `wasm/oxicloud-plugin-hello/`. Run with `cargo test --features plugins`.
use std::time::{Duration, Instant};
use super::ExtismPluginManager;
use super::manifest;
use super::runtime::{InvokeOutcome, PluginRuntime};
use crate::application::ports::plugin_ports::event_export_name;
use crate::common::config::PluginConfig;
fn cfg() -> PluginConfig {
PluginConfig::default()
}
/// Load a committed `.wasm` fixture, failing with a build hint if it's missing.
fn fixture(name: &str) -> Vec<u8> {
let path = format!(
"{}/tests/fixtures/plugins/{}",
env!("CARGO_MANIFEST_DIR"),
name
);
std::fs::read(&path).unwrap_or_else(|e| {
panic!("missing fixture {path}: {e}\n run scripts/build-plugin-hello.sh to (re)build it")
})
}
fn file_uploaded_input() -> String {
serde_json::json!({
"abi": 0,
"event": "file.uploaded",
"context": {
"plugin_id": "com.example.hello",
"user_id": "u_test",
"invocation_id": "inv_test_0001"
},
"payload": { "path": "/photos/2026/cat.jpg", "size": 81234, "mime": "image/jpeg" }
})
.to_string()
}
fn user_login_input() -> String {
serde_json::json!({
"abi": 0,
"event": "user.login",
"context": {
"plugin_id": "com.example.hello",
"user_id": "u_test",
"invocation_id": "inv_login_0001"
},
"payload": {
"user_id": "u_test",
"username": "alice",
"email": "alice@example.com",
"first_login": true,
"is_external": false
}
})
.to_string()
}
// ---- The M0 exit criterion: the full loop, per event ------------------------
#[test]
fn acceptance_file_uploaded_returns_ok_and_calls_host_log() {
let rt = PluginRuntime::new("com.example.hello", fixture("hello.wasm"));
let result = rt.invoke(&cfg(), "on_file_uploaded", "inv", &file_uploaded_input());
assert!(
result.outcome.is_ok(),
"plugin did not complete: {:?}",
result.outcome
);
assert!(
result.logs.iter().any(|(level, msg)| level == "info"
&& msg.contains("hello plugin saw upload: /photos/2026/cat.jpg")),
"expected the plugin's host log line, got: {:?}",
result.logs
);
}
#[test]
fn acceptance_user_login_returns_ok_and_calls_host_log() {
let rt = PluginRuntime::new("com.example.hello", fixture("hello.wasm"));
let result = rt.invoke(&cfg(), "on_user_login", "inv", &user_login_input());
assert!(
result.outcome.is_ok(),
"plugin did not complete: {:?}",
result.outcome
);
assert!(
result.logs.iter().any(|(level, msg)| level == "info"
&& msg.contains("hello plugin saw login: user u_test (first_login=true)")),
"expected the plugin's user.login log line, got: {:?}",
result.logs
);
}
// ---- The guarantees, not just the happy path --------------------------------
#[test]
fn rejects_wrong_abi() {
let rt = PluginRuntime::new("com.example.wrong-abi", fixture("wrong_abi.wasm"));
assert!(
matches!(
rt.check_loadable(&cfg(), &[]),
InvokeOutcome::AbiMismatch { got: 1 }
),
"wrong-abi plugin should be rejected at load"
);
}
#[test]
fn load_requires_subscribed_event_exports() {
let cfg = cfg();
let login_export = vec![event_export_name("user.login")];
// hello.wasm exports both handlers -> loadable for user.login.
let hello = PluginRuntime::new("com.example.hello", fixture("hello.wasm"));
assert!(matches!(
hello.check_loadable(&cfg, &login_export),
InvokeOutcome::Ok
));
// omit_login.wasm lacks on_user_login -> rejected when it claims user.login.
let omit = PluginRuntime::new("com.example.omit", fixture("omit_login.wasm"));
assert!(
matches!(
omit.check_loadable(&cfg, &login_export),
InvokeOutcome::MissingExport(ref e) if e == "on_user_login"
),
"omit_login must be rejected for a user.login subscription"
);
// …but it is fine for file.uploaded, which it does export.
assert!(matches!(
omit.check_loadable(&cfg, &[event_export_name("file.uploaded")]),
InvokeOutcome::Ok
));
}
#[test]
fn contains_a_panicking_plugin() {
let rt = PluginRuntime::new("com.example.panic", fixture("panic.wasm"));
let result = rt.invoke(&cfg(), "on_file_uploaded", "inv", &file_uploaded_input());
assert!(
matches!(result.outcome, InvokeOutcome::Trap(_)),
"expected a contained trap, got {:?}",
result.outcome
);
// Reaching this line at all proves the host process survived the trap.
}
#[test]
fn enforces_timeout() {
let rt = PluginRuntime::new("com.example.sleep", fixture("sleep.wasm"));
let start = Instant::now();
let result = rt.invoke(&cfg(), "on_file_uploaded", "inv", &file_uploaded_input());
let elapsed = start.elapsed();
assert!(
matches!(result.outcome, InvokeOutcome::Timeout),
"expected a timeout, got {:?}",
result.outcome
);
assert!(
elapsed < Duration::from_secs(2),
"timeout took too long to fire: {elapsed:?}"
);
}
#[test]
fn idle_eviction_drops_and_recompiles() {
let rt = PluginRuntime::new("com.example.hello", fixture("hello.wasm"));
// First invoke compiles + caches the module.
let r1 = rt.invoke(&cfg(), "on_file_uploaded", "inv1", &file_uploaded_input());
assert!(r1.outcome.is_ok(), "first invoke: {:?}", r1.outcome);
// Idle past a zero TTL -> the cached module is dropped.
assert!(
rt.evict_if_idle(Duration::ZERO),
"a just-idle module should be evicted"
);
// Nothing left to evict the second time.
assert!(
!rt.evict_if_idle(Duration::ZERO),
"second eviction is a no-op"
);
// The next invoke recompiles transparently and still works.
let r2 = rt.invoke(&cfg(), "on_file_uploaded", "inv2", &file_uploaded_input());
assert!(
r2.outcome.is_ok(),
"recompile after eviction: {:?}",
r2.outcome
);
// A long TTL never evicts a freshly-used module.
assert!(
!rt.evict_if_idle(Duration::from_secs(3600)),
"a fresh module must not be evicted"
);
}
#[test]
fn no_network() {
let rt = PluginRuntime::new("com.example.net", fixture("net.wasm"));
let result = rt.invoke(&cfg(), "on_file_uploaded", "inv", &file_uploaded_input());
assert!(
!result.outcome.is_ok(),
"network access should be denied, got {:?}",
result.outcome
);
}
// ---- Manager discovery + dispatch ------------------------------------------
/// Write a one-plugin directory (plugin.toml + the given wasm) under a tempdir
/// and load a manager from it.
fn manager_with(wasm_name: &str, subscribe_toml: &str) -> (tempfile::TempDir, ExtismPluginManager) {
let tmp = tempfile::tempdir().unwrap();
let dir = tmp.path().join("plugin");
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("plugin.wasm"), fixture(wasm_name)).unwrap();
std::fs::write(
dir.join("plugin.toml"),
format!(
r#"
[plugin]
id = "com.example.test"
name = "Test"
version = "0.1.0"
abi = 0
entrypoint = "plugin.wasm"
[events]
subscribe = {subscribe_toml}
"#
),
)
.unwrap();
let manager = ExtismPluginManager::load_from_dir(cfg(), tmp.path());
(tmp, manager)
}
#[tokio::test]
async fn manager_loads_and_dispatches_both_events() {
use crate::application::ports::plugin_ports::{
EVENT_FILE_UPLOADED, EVENT_USER_LOGIN, PluginDispatchPort, PluginEvent,
};
let (_tmp, manager) = manager_with("hello.wasm", r#"["file.uploaded", "user.login"]"#);
assert_eq!(manager.loaded_count(), 1, "the valid plugin should load");
assert!(manager.has_subscribers("file.uploaded"));
assert!(manager.has_subscribers("user.login"));
assert!(!manager.has_subscribers("file.deleted"));
// Both dispatches run the plugin on the blocking pool; neither may panic.
manager.dispatch(PluginEvent {
name: EVENT_FILE_UPLOADED,
user_id: Some("u_test".into()),
invocation_id: "inv_upload".into(),
payload: serde_json::json!({ "path": "/a.txt", "size": 3, "mime": "text/plain" }),
});
manager.dispatch(PluginEvent {
name: EVENT_USER_LOGIN,
user_id: Some("u_test".into()),
invocation_id: "inv_login".into(),
payload: serde_json::json!({ "user_id": "u_test", "first_login": false }),
});
tokio::time::sleep(Duration::from_millis(300)).await;
}
#[test]
fn manager_rejects_plugin_missing_a_subscribed_export() {
// omit_login.wasm subscribes to user.login but doesn't export on_user_login.
let (_tmp, rejected) = manager_with("omit_login.wasm", r#"["user.login"]"#);
assert_eq!(rejected.loaded_count(), 0, "missing export -> not loaded");
// The same wasm is fine when it only claims an event it actually exports.
let (_tmp2, loaded) = manager_with("omit_login.wasm", r#"["file.uploaded"]"#);
assert_eq!(loaded.loaded_count(), 1);
}
// ---- Manifest validation (no wasm needed) -----------------------------------
const VALID_MANIFEST: &str = r#"
[plugin]
id = "com.example.hello"
name = "Hello"
version = "0.1.0"
abi = 0
entrypoint = "hello.wasm"
[events]
subscribe = ["file.uploaded"]
"#;
#[test]
fn manifest_accepts_valid() {
let m = manifest::parse_and_validate(VALID_MANIFEST).expect("valid manifest");
assert_eq!(m.plugin.id, "com.example.hello");
}
#[test]
fn manifest_accepts_user_login_and_combined() {
let login = VALID_MANIFEST.replace(r#"["file.uploaded"]"#, r#"["user.login"]"#);
assert!(manifest::parse_and_validate(&login).is_ok());
let both = VALID_MANIFEST.replace(r#"["file.uploaded"]"#, r#"["file.uploaded", "user.login"]"#);
assert!(manifest::parse_and_validate(&both).is_ok());
}
#[test]
fn manifest_rejects_unknown_field() {
let toml = format!("{VALID_MANIFEST}\nbogus_top_level = true\n");
assert_eq!(
manifest::parse_and_validate(&toml).unwrap_err().reason(),
"parse_error"
);
}
#[test]
fn manifest_rejects_abi_mismatch() {
let toml = VALID_MANIFEST.replace("abi = 0", "abi = 1");
assert_eq!(
manifest::parse_and_validate(&toml).unwrap_err().reason(),
"abi_mismatch"
);
}
#[test]
fn manifest_rejects_unknown_event() {
let toml = VALID_MANIFEST.replace(r#"["file.uploaded"]"#, r#"["file.deleted"]"#);
assert_eq!(
manifest::parse_and_validate(&toml).unwrap_err().reason(),
"unknown_event"
);
}
#[test]
fn manifest_rejects_nonempty_permissions() {
let toml = format!("{VALID_MANIFEST}\n[permissions]\nfs = \"/tmp\"\n");
assert_eq!(
manifest::parse_and_validate(&toml).unwrap_err().reason(),
"permissions_not_empty"
);
}