diff --git a/crates/volt-runner/src/js_runtime/tests/native_modules/plugins.rs b/crates/volt-runner/src/js_runtime/tests/native_modules/plugins.rs index 1588fca..9ac52f1 100644 --- a/crates/volt-runner/src/js_runtime/tests/native_modules/plugins.rs +++ b/crates/volt-runner/src/js_runtime/tests/native_modules/plugins.rs @@ -112,3 +112,92 @@ fn plugins_module_prefetch_for_is_available() { assert_eq!(result, "ok"); } + +#[test] +fn plugins_module_exposes_state_and_error_queries() { + let (manager, _) = build_plugin_manager(); + manager.fail_plugin( + "acme.search", + "PLUGIN_BROKEN", + "boom".to_string(), + None, + None, + ); + let runtime = runtime_with_plugin_manager( + unique_temp_dir("plugins-observability"), + &["fs"], + Some(manager), + ); + + let result = runtime + .client() + .eval_promise_string( + "(async () => { + const plugins = globalThis.__volt.plugins; + const state = await plugins.getPluginState('acme.search'); + const errors = await plugins.getPluginErrors('acme.search'); + return `${state.currentState}:${errors.length}:${errors[0].code}`; + })()", + ) + .expect("plugin state"); + + assert_eq!(result, "failed:1:PLUGIN_BROKEN"); +} + +#[test] +fn plugins_module_receives_lifecycle_events_via_native_bridge() { + let (manager, _) = build_plugin_manager(); + let runtime = runtime_with_plugin_manager( + unique_temp_dir("plugins-lifecycle-events"), + &["fs"], + Some(manager.clone()), + ); + let runtime_client = runtime.client(); + let _subscription = manager.on_lifecycle(Box::new(move |event| { + let payload = serde_json::to_value(event).expect("serialize event"); + runtime_client + .dispatch_native_event("plugin:lifecycle", payload) + .expect("dispatch lifecycle event"); + })); + + runtime + .client() + .eval_unit( + "(async () => { + const plugins = globalThis.__volt.plugins; + globalThis.__pluginLifecycleEvents = []; + const handler = (event) => globalThis.__pluginLifecycleEvents.push(event.newState); + globalThis.__pluginLifecycleHandler = handler; + plugins.on('plugin:lifecycle', handler); + })()", + ) + .expect("bind lifecycle handler"); + + manager.fail_plugin( + "acme.search", + "PLUGIN_BROKEN", + "boom".to_string(), + None, + None, + ); + + let seen = runtime + .client() + .eval_string("globalThis.__pluginLifecycleEvents.join(',')") + .expect("captured events"); + assert_eq!(seen, "failed"); + + runtime + .client() + .eval_unit( + "globalThis.__volt.plugins.off('plugin:lifecycle', globalThis.__pluginLifecycleHandler)", + ) + .expect("unbind lifecycle handler"); + let _ = manager.retry_plugin("acme.search"); + + let after_off = runtime + .client() + .eval_string("globalThis.__pluginLifecycleEvents.join(',')") + .expect("events after off"); + assert_eq!(after_off, "failed"); +} diff --git a/crates/volt-runner/src/main.rs b/crates/volt-runner/src/main.rs index af5f7e3..40466c9 100644 --- a/crates/volt-runner/src/main.rs +++ b/crates/volt-runner/src/main.rs @@ -66,6 +66,18 @@ fn run() -> Result<(), RunnerError> { .load_backend_bundle(&backend_bundle_source) .map_err(|err| RunnerError::App(format!("failed to load backend bundle: {err}")))?; let runtime_client = js_runtime.client(); + let lifecycle_runtime = runtime_client.clone(); + let _lifecycle_subscription = plugin_manager.on_lifecycle(Box::new(move |event| { + dispatch_plugin_lifecycle_event(&lifecycle_runtime, "plugin:lifecycle", event); + })); + let failed_runtime = runtime_client.clone(); + let _failed_subscription = plugin_manager.on_plugin_failed(Box::new(move |event| { + dispatch_plugin_lifecycle_event(&failed_runtime, "plugin:failed", event); + })); + let activated_runtime = runtime_client.clone(); + let _activated_subscription = plugin_manager.on_plugin_activated(Box::new(move |event| { + dispatch_plugin_lifecycle_event(&activated_runtime, "plugin:activated", event); + })); let ipc_bridge = ipc_bridge::IpcBridge::new_with_plugin_manager( runtime_client.clone(), Some(plugin_manager.clone()), @@ -127,3 +139,20 @@ fn run() -> Result<(), RunnerError> { Ok(()) } + +fn dispatch_plugin_lifecycle_event( + runtime_client: &js_runtime_pool::JsRuntimePoolClient, + event_type: &str, + event: &plugin_manager::PluginLifecycleEvent, +) { + let payload = match serde_json::to_value(event) { + Ok(payload) => payload, + Err(error) => { + tracing::error!(error = %error, event_type = %event_type, "failed to serialize plugin lifecycle event"); + return; + } + }; + if let Err(error) = runtime_client.dispatch_native_event(event_type, payload) { + tracing::error!(error = %error, event_type = %event_type, "failed to dispatch plugin lifecycle event"); + } +} diff --git a/crates/volt-runner/src/modules/volt_plugins.rs b/crates/volt-runner/src/modules/volt_plugins.rs index 80a2cdb..3c44add 100644 --- a/crates/volt-runner/src/modules/volt_plugins.rs +++ b/crates/volt-runner/src/modules/volt_plugins.rs @@ -1,6 +1,14 @@ -use boa_engine::{Context, IntoJsFunctionCopied, JsValue, Module}; +use boa_engine::object::builtins::JsFunction; +use boa_engine::{Context, IntoJsFunctionCopied, JsResult, JsValue, Module}; +use serde_json::Value; -use super::{native_function_module, plugin_manager, promise_from_result}; +use super::{ + bind_native_event_handler, js_error, native_function_module, plugin_manager, + promise_from_json_result, promise_from_result, +}; + +const NATIVE_EVENT_ON_GLOBAL: &str = "__volt_native_event_on__"; +const NATIVE_EVENT_OFF_GLOBAL: &str = "__volt_native_event_off__"; fn delegate_grant(plugin_id: String, grant_id: String, context: &mut Context) -> JsValue { promise_from_result(context, delegate_grant_result(plugin_id, grant_id)).into() @@ -14,6 +22,42 @@ fn prefetch_for(surface: String, context: &mut Context) -> JsValue { promise_from_result(context, prefetch_for_result(surface)).into() } +fn get_states(context: &mut Context) -> JsValue { + promise_from_json_result(context, get_states_result()).into() +} + +fn get_plugin_state(plugin_id: String, context: &mut Context) -> JsValue { + promise_from_json_result(context, get_plugin_state_result(plugin_id)).into() +} + +fn get_errors(context: &mut Context) -> JsValue { + promise_from_json_result(context, get_errors_result()).into() +} + +fn get_plugin_errors(plugin_id: String, context: &mut Context) -> JsValue { + promise_from_json_result(context, get_plugin_errors_result(plugin_id)).into() +} + +fn get_discovery_issues(context: &mut Context) -> JsValue { + promise_from_json_result(context, get_discovery_issues_result()).into() +} + +fn retry_plugin(plugin_id: String, context: &mut Context) -> JsValue { + promise_from_result(context, retry_plugin_result(plugin_id)).into() +} + +fn enable_plugin(plugin_id: String, context: &mut Context) -> JsValue { + promise_from_result(context, enable_plugin_result(plugin_id)).into() +} + +fn on(event_name: String, handler: JsFunction, context: &mut Context) -> JsResult<()> { + bind_native_plugin_event(context, "on", NATIVE_EVENT_ON_GLOBAL, event_name, handler) +} + +fn off(event_name: String, handler: JsFunction, context: &mut Context) -> JsResult<()> { + bind_native_plugin_event(context, "off", NATIVE_EVENT_OFF_GLOBAL, event_name, handler) +} + fn delegate_grant_result(plugin_id: String, grant_id: String) -> Result<(), String> { let plugin_id = required_name(plugin_id, "plugin id")?; let grant_id = required_name(grant_id, "grant id")?; @@ -35,6 +79,72 @@ fn prefetch_for_result(surface: String) -> Result<(), String> { Ok(()) } +fn get_states_result() -> Result { + serde_json::to_value(plugin_manager()?.get_states()).map_err(|error| error.to_string()) +} + +fn get_plugin_state_result(plugin_id: String) -> Result { + let plugin_id = required_name(plugin_id, "plugin id")?; + serde_json::to_value(plugin_manager()?.get_plugin_state(&plugin_id)) + .map_err(|error| error.to_string()) +} + +fn get_errors_result() -> Result { + serde_json::to_value(plugin_manager()?.get_errors()).map_err(|error| error.to_string()) +} + +fn get_plugin_errors_result(plugin_id: String) -> Result { + let plugin_id = required_name(plugin_id, "plugin id")?; + serde_json::to_value(plugin_manager()?.get_plugin_errors(&plugin_id)) + .map_err(|error| error.to_string()) +} + +fn get_discovery_issues_result() -> Result { + serde_json::to_value(plugin_manager()?.discovery_issues()).map_err(|error| error.to_string()) +} + +fn retry_plugin_result(plugin_id: String) -> Result<(), String> { + let plugin_id = required_name(plugin_id, "plugin id")?; + plugin_manager()? + .retry_plugin(&plugin_id) + .map_err(|error| error.to_string()) +} + +fn enable_plugin_result(plugin_id: String) -> Result<(), String> { + let plugin_id = required_name(plugin_id, "plugin id")?; + plugin_manager()? + .enable_plugin(&plugin_id) + .map_err(|error| error.to_string()) +} + +fn bind_native_plugin_event( + context: &mut Context, + api_function: &'static str, + global_name: &'static str, + event_name: String, + handler: JsFunction, +) -> JsResult<()> { + bind_native_event_handler( + context, + "volt:plugins", + api_function, + global_name, + normalize_event_name(event_name) + .map_err(|error| js_error("volt:plugins", api_function, error))?, + handler, + ) +} + +fn normalize_event_name(event_name: String) -> Result<&'static str, String> { + match event_name.trim() { + "plugin:lifecycle" => Ok("plugin:lifecycle"), + "plugin:failed" => Ok("plugin:failed"), + "plugin:activated" => Ok("plugin:activated"), + "" => Err("plugin event name must not be empty".to_string()), + other => Err(format!("unsupported plugin event '{other}'")), + } +} + fn required_name(value: String, label: &str) -> Result { let trimmed = value.trim(); if trimmed.is_empty() { @@ -47,6 +157,15 @@ pub fn build_module(context: &mut Context) -> Module { let delegate_grant = delegate_grant.into_js_function_copied(context); let revoke_grant = revoke_grant.into_js_function_copied(context); let prefetch_for = prefetch_for.into_js_function_copied(context); + let get_states = get_states.into_js_function_copied(context); + let get_plugin_state = get_plugin_state.into_js_function_copied(context); + let get_errors = get_errors.into_js_function_copied(context); + let get_plugin_errors = get_plugin_errors.into_js_function_copied(context); + let get_discovery_issues = get_discovery_issues.into_js_function_copied(context); + let retry_plugin = retry_plugin.into_js_function_copied(context); + let enable_plugin = enable_plugin.into_js_function_copied(context); + let on = on.into_js_function_copied(context); + let off = off.into_js_function_copied(context); native_function_module( context, @@ -54,6 +173,15 @@ pub fn build_module(context: &mut Context) -> Module { ("delegateGrant", delegate_grant), ("revokeGrant", revoke_grant), ("prefetchFor", prefetch_for), + ("getStates", get_states), + ("getPluginState", get_plugin_state), + ("getErrors", get_errors), + ("getPluginErrors", get_plugin_errors), + ("getDiscoveryIssues", get_discovery_issues), + ("retryPlugin", retry_plugin), + ("enablePlugin", enable_plugin), + ("on", on), + ("off", off), ], ) } @@ -81,4 +209,20 @@ mod tests { assert!(error.contains("plugin manager is unavailable")); } + + #[test] + fn normalize_event_name_accepts_supported_events() { + assert_eq!( + normalize_event_name("plugin:lifecycle".to_string()), + Ok("plugin:lifecycle") + ); + assert_eq!( + normalize_event_name("plugin:failed".to_string()), + Ok("plugin:failed") + ); + assert_eq!( + normalize_event_name("plugin:activated".to_string()), + Ok("plugin:activated") + ); + } } diff --git a/crates/volt-runner/src/plugin_manager/discovery.rs b/crates/volt-runner/src/plugin_manager/discovery.rs index 8895130..b201396 100644 --- a/crates/volt-runner/src/plugin_manager/discovery.rs +++ b/crates/volt-runner/src/plugin_manager/discovery.rs @@ -50,6 +50,24 @@ impl PluginManager { config: RunnerPluginConfig, factory: Arc, access_picker: Arc, + ) -> Result { + Self::with_dependencies_and_error_history_limit( + app_name, + permissions, + config, + factory, + access_picker, + super::DEFAULT_PLUGIN_ERROR_HISTORY_LIMIT, + ) + } + + pub(super) fn with_dependencies_and_error_history_limit( + app_name: String, + permissions: &[String], + config: RunnerPluginConfig, + factory: Arc, + access_picker: Arc, + error_history_limit: usize, ) -> Result { let app_permissions = permissions .iter() @@ -63,6 +81,8 @@ impl PluginManager { app_data_root, factory, access_picker, + error_history_limit, + lifecycle_bus: super::LifecycleBus::new(), registry: Mutex::new(PluginRegistry::new()), }), }; @@ -80,6 +100,7 @@ impl PluginManager { .collect::>(); let mut manifest_paths = Vec::new(); let mut registry = PluginRegistry::new(); + let mut lifecycle_events = Vec::new(); for directory in &self.inner.config.plugin_dirs { let resolved = resolve_plugin_directory(directory); @@ -124,6 +145,7 @@ impl PluginManager { continue; } } + lifecycle_events.extend(record.lifecycle.recorded_events(&record.manifest.id)); registry.plugins.insert(record.manifest.id.clone(), record); } Err(issue) => registry.discovery_issues.push(issue), @@ -144,6 +166,9 @@ impl PluginManager { if let Ok(mut guard) = self.inner.registry.lock() { *guard = registry; } + for event in lifecycle_events { + self.inner.lifecycle_bus.emit(event); + } } fn discover_plugin_record( @@ -201,16 +226,22 @@ impl PluginManager { .difference(&effective_capabilities) .cloned() .collect::>(); - lifecycle.fail( - &manifest.id, - PLUGIN_NOT_AVAILABLE_CODE, - format!( - "requested capabilities are unsatisfiable: {}", - missing.join(", ") - ), - Some(json!({ "missingCapabilities": missing })), - None, - ); + lifecycle + .fail( + &manifest.id, + PLUGIN_NOT_AVAILABLE_CODE, + format!( + "requested capabilities are unsatisfiable: {}", + missing.join(", ") + ), + Some(json!({ "missingCapabilities": missing })), + None, + self.inner.error_history_limit, + ) + .map_err(|message| PluginDiscoveryIssue { + path: Some(manifest_path.to_path_buf()), + message, + })?; } else { lifecycle .transition(PluginState::Validated) diff --git a/crates/volt-runner/src/plugin_manager/lifecycle.rs b/crates/volt-runner/src/plugin_manager/lifecycle.rs index e0c4770..0c514e1 100644 --- a/crates/volt-runner/src/plugin_manager/lifecycle.rs +++ b/crates/volt-runner/src/plugin_manager/lifecycle.rs @@ -2,30 +2,40 @@ use std::time::{Duration, SystemTime, UNIX_EPOCH}; use serde_json::Value; -use super::{PluginError, PluginRecord, PluginSnapshot, PluginState, PluginStateTransition}; +use super::{ + PluginError, PluginLifecycleEvent, PluginRecord, PluginRegistrationSnapshot, PluginState, + PluginStateSnapshot, PluginStateTransition, +}; #[derive(Debug, Clone)] pub(super) struct PluginLifecycle { state: Option, pub(super) transitions: Vec, pub(super) errors: Vec, + consecutive_failures: u32, } impl PluginRecord { - #[allow(dead_code)] - pub(super) fn snapshot(&self) -> PluginSnapshot { - PluginSnapshot { + pub(super) fn snapshot(&self) -> PluginStateSnapshot { + PluginStateSnapshot { plugin_id: self.manifest.id.clone(), - state: self.lifecycle.current_state(), + current_state: self.lifecycle.current_state(), enabled: self.enabled, manifest_path: self.manifest_path.clone(), data_root: self.data_root.clone(), requested_capabilities: self.requested_capabilities.iter().cloned().collect(), effective_capabilities: self.effective_capabilities.iter().cloned().collect(), - transitions: self.lifecycle.transitions.clone(), + transition_history: self.lifecycle.transitions.clone(), errors: self.lifecycle.errors.clone(), metrics: self.metrics.clone(), process_running: self.process.is_some(), + active_registrations: PluginRegistrationSnapshot { + command_count: self.registrations.commands.len(), + event_subscription_count: self.registrations.event_subscriptions.len(), + ipc_handler_count: self.registrations.ipc_handlers.len(), + }, + delegated_grant_count: self.delegated_grants.len(), + consecutive_failures: self.lifecycle.consecutive_failures(), } } } @@ -36,10 +46,14 @@ impl PluginLifecycle { state: None, transitions: Vec::new(), errors: Vec::new(), + consecutive_failures: 0, } } - pub(super) fn transition(&mut self, next_state: PluginState) -> Result<(), String> { + pub(super) fn transition( + &mut self, + next_state: PluginState, + ) -> Result { if let Some(current_state) = self.state && !is_valid_transition(current_state, next_state) { @@ -55,14 +69,17 @@ impl PluginLifecycle { )); } - let previous_state = self.state; - self.state = Some(next_state); - self.transitions.push(PluginStateTransition { - previous_state, + let transition = PluginStateTransition { + previous_state: self.state, new_state: next_state, timestamp_ms: now_ms(), - }); - Ok(()) + }; + self.state = Some(next_state); + self.transitions.push(transition.clone()); + if next_state == PluginState::Running { + self.consecutive_failures = 0; + } + Ok(transition) } pub(super) fn fail( @@ -72,32 +89,78 @@ impl PluginLifecycle { message: String, details: Option, stderr: Option, - ) { - if self.state.is_none() { - self.state = Some(PluginState::Failed); - self.transitions.push(PluginStateTransition { + max_errors: usize, + ) -> Result<(PluginStateTransition, PluginError, u32), String> { + let failure_state = self.state.unwrap_or(PluginState::Failed); + let transition = if self.state.is_none() { + let transition = PluginStateTransition { previous_state: None, new_state: PluginState::Failed, timestamp_ms: now_ms(), - }); + }; + self.state = Some(PluginState::Failed); + self.transitions.push(transition.clone()); + transition } else { - let _ = self.transition(PluginState::Failed); - } + self.transition(PluginState::Failed)? + }; + let error = self.push_error( + PluginError { + plugin_id: plugin_id.to_string(), + state: failure_state, + code: code.to_string(), + message, + details, + stderr, + timestamp_ms: transition.timestamp_ms, + }, + max_errors, + ); + self.consecutive_failures = self.consecutive_failures.saturating_add(1); + Ok((transition, error, self.consecutive_failures)) + } - self.errors.push(PluginError { - plugin_id: plugin_id.to_string(), - state: PluginState::Failed, - code: code.to_string(), - message, - details, - stderr, - timestamp_ms: now_ms(), - }); + pub(super) fn push_error(&mut self, error: PluginError, max_errors: usize) -> PluginError { + self.errors.push(error.clone()); + trim_error_history(&mut self.errors, max_errors); + error } pub(super) fn current_state(&self) -> PluginState { self.state.expect("plugin lifecycle must be initialized") } + + pub(super) fn consecutive_failures(&self) -> u32 { + self.consecutive_failures + } + + pub(super) fn reset_failures(&mut self) { + self.consecutive_failures = 0; + } + + pub(super) fn recorded_events(&self, plugin_id: &str) -> Vec { + self.transitions + .iter() + .map(|transition| PluginLifecycleEvent { + plugin_id: plugin_id.to_string(), + previous_state: transition.previous_state, + new_state: transition.new_state, + timestamp: transition.timestamp_ms, + error: self + .errors + .iter() + .find(|error| error.timestamp_ms == transition.timestamp_ms) + .cloned(), + }) + .collect() + } +} + +fn trim_error_history(errors: &mut Vec, max_errors: usize) { + if errors.len() > max_errors { + let drop_count = errors.len() - max_errors; + errors.drain(0..drop_count); + } } fn is_valid_transition(current: PluginState, next: PluginState) -> bool { @@ -110,6 +173,8 @@ fn is_valid_transition(current: PluginState, next: PluginState) -> bool { (PluginState::Discovered, PluginState::Validated) | (PluginState::Validated, PluginState::Spawning) | (PluginState::Terminated, PluginState::Spawning) + | (PluginState::Failed, PluginState::Spawning) + | (PluginState::Disabled, PluginState::Validated) | (PluginState::Spawning, PluginState::Loaded) | (PluginState::Loaded, PluginState::Terminated) | (PluginState::Loaded, PluginState::Active) diff --git a/crates/volt-runner/src/plugin_manager/lifecycle_bus.rs b/crates/volt-runner/src/plugin_manager/lifecycle_bus.rs new file mode 100644 index 0000000..89f1b0b --- /dev/null +++ b/crates/volt-runner/src/plugin_manager/lifecycle_bus.rs @@ -0,0 +1,115 @@ +use std::collections::HashMap; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::{Arc, Mutex}; + +use serde::Serialize; + +use super::{PluginError, PluginState}; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +pub(crate) struct SubscriptionId(pub(crate) u64); + +#[derive(Debug, Clone, Serialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct PluginLifecycleEvent { + pub(crate) plugin_id: String, + pub(crate) previous_state: Option, + pub(crate) new_state: PluginState, + pub(crate) timestamp: u64, + pub(crate) error: Option, +} + +#[derive(Clone, Copy)] +enum LifecycleTopic { + All, + Failed, + Activated, +} + +type LifecycleHandler = Arc; + +#[derive(Clone)] +pub(crate) struct LifecycleBus { + next_id: Arc, + subscribers: Arc>>, +} + +impl LifecycleBus { + pub(crate) fn new() -> Self { + Self { + next_id: Arc::new(AtomicU64::new(1)), + subscribers: Arc::new(Mutex::new(HashMap::new())), + } + } + + pub(crate) fn on_lifecycle( + &self, + handler: Box, + ) -> SubscriptionId { + self.subscribe(LifecycleTopic::All, handler) + } + + pub(crate) fn on_plugin_failed( + &self, + handler: Box, + ) -> SubscriptionId { + self.subscribe(LifecycleTopic::Failed, handler) + } + + pub(crate) fn on_plugin_activated( + &self, + handler: Box, + ) -> SubscriptionId { + self.subscribe(LifecycleTopic::Activated, handler) + } + + #[allow(dead_code)] + pub(crate) fn off(&self, subscription_id: SubscriptionId) { + if let Ok(mut subscribers) = self.subscribers.lock() { + subscribers.remove(&subscription_id); + } + } + + pub(crate) fn emit(&self, event: PluginLifecycleEvent) { + let handlers = self + .subscribers + .lock() + .map(|subscribers| { + subscribers + .values() + .filter(|(topic, _)| topic_matches(*topic, &event)) + .map(|(_, handler)| handler.clone()) + .collect::>() + }) + .unwrap_or_default(); + + for handler in handlers { + handler(&event); + } + } + + fn subscribe( + &self, + topic: LifecycleTopic, + handler: Box, + ) -> SubscriptionId { + let id = SubscriptionId(self.next_id.fetch_add(1, Ordering::Relaxed)); + if let Ok(mut subscribers) = self.subscribers.lock() { + subscribers.insert(id, (topic, Arc::from(handler))); + } + id + } +} + +fn topic_matches(topic: LifecycleTopic, event: &PluginLifecycleEvent) -> bool { + match topic { + LifecycleTopic::All => true, + LifecycleTopic::Failed => event.new_state == PluginState::Failed, + LifecycleTopic::Activated => { + matches!( + (event.previous_state, event.new_state), + (Some(PluginState::Loaded), PluginState::Active) + ) + } + } +} diff --git a/crates/volt-runner/src/plugin_manager/mod.rs b/crates/volt-runner/src/plugin_manager/mod.rs index be4d397..1a0d1d2 100644 --- a/crates/volt-runner/src/plugin_manager/mod.rs +++ b/crates/volt-runner/src/plugin_manager/mod.rs @@ -2,6 +2,7 @@ use std::collections::{BTreeSet, HashMap, HashSet}; use std::path::PathBuf; use std::sync::{Arc, Mutex}; +use serde::Serialize; use serde_json::Value; #[cfg(test)] use volt_core::ipc::IPC_HANDLER_TIMEOUT_CODE; @@ -19,6 +20,7 @@ mod host_api_helpers; mod host_api_storage; mod host_api_support; mod lifecycle; +mod lifecycle_bus; mod manifest; mod paths; mod process; @@ -27,6 +29,7 @@ mod watchdog; use self::access::{NativePluginAccessPicker, PluginAccessPicker}; use self::lifecycle::{PluginLifecycle, now_ms}; +use self::lifecycle_bus::LifecycleBus; use self::manifest::{compute_effective_capabilities, parse_plugin_manifest, parse_plugin_route}; use self::paths::{ collect_manifest_paths, ensure_plugin_data_root, resolve_app_data_root, @@ -52,9 +55,14 @@ const PLUGIN_NOT_AVAILABLE_CODE: &str = "PLUGIN_NOT_AVAILABLE"; const PLUGIN_ROUTE_INVALID_CODE: &str = "PLUGIN_ROUTE_INVALID"; const PLUGIN_RUNTIME_ERROR_CODE: &str = "PLUGIN_RUNTIME_ERROR"; const PLUGIN_STORAGE_ERROR_CODE: &str = "PLUGIN_STORAGE_ERROR"; +const PLUGIN_AUTO_DISABLED_CODE: &str = "PLUGIN_AUTO_DISABLED"; const DEFAULT_PRE_SPAWN_GRACE_MS: u64 = 50; +const DEFAULT_PLUGIN_ERROR_HISTORY_LIMIT: usize = 50; -#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) use self::lifecycle_bus::{PluginLifecycleEvent, SubscriptionId}; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[serde(rename_all = "camelCase")] pub(crate) enum PluginState { Discovered, Validated, @@ -68,15 +76,16 @@ pub(crate) enum PluginState { Disabled, } -#[derive(Debug, Clone, PartialEq, Eq)] +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] +#[serde(rename_all = "camelCase")] pub(crate) struct PluginStateTransition { pub(crate) previous_state: Option, pub(crate) new_state: PluginState, pub(crate) timestamp_ms: u64, } -#[derive(Debug, Clone)] -#[allow(dead_code)] +#[derive(Debug, Clone, Serialize)] +#[serde(rename_all = "camelCase")] pub(crate) struct PluginError { pub(crate) plugin_id: String, pub(crate) state: PluginState, @@ -87,7 +96,8 @@ pub(crate) struct PluginError { pub(crate) timestamp_ms: u64, } -#[derive(Debug, Clone, Default, PartialEq, Eq)] +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)] +#[serde(rename_all = "camelCase")] pub(crate) struct PluginResourceMetrics { pub(crate) pid: Option, pub(crate) started_at_ms: Option, @@ -98,27 +108,38 @@ pub(crate) struct PluginResourceMetrics { pub(crate) heartbeat_failures: u32, } -#[derive(Debug, Clone)] -#[allow(dead_code)] +#[derive(Debug, Clone, Serialize)] +#[serde(rename_all = "camelCase")] pub(crate) struct PluginDiscoveryIssue { pub(crate) path: Option, pub(crate) message: String, } -#[derive(Debug, Clone)] -#[allow(dead_code)] -pub(crate) struct PluginSnapshot { +#[derive(Debug, Clone, Serialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct PluginRegistrationSnapshot { + pub(crate) command_count: usize, + pub(crate) event_subscription_count: usize, + pub(crate) ipc_handler_count: usize, +} + +#[derive(Debug, Clone, Serialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct PluginStateSnapshot { pub(crate) plugin_id: String, - pub(crate) state: PluginState, + pub(crate) current_state: PluginState, pub(crate) enabled: bool, pub(crate) manifest_path: PathBuf, pub(crate) data_root: Option, pub(crate) requested_capabilities: Vec, pub(crate) effective_capabilities: Vec, - pub(crate) transitions: Vec, + pub(crate) transition_history: Vec, pub(crate) errors: Vec, pub(crate) metrics: PluginResourceMetrics, pub(crate) process_running: bool, + pub(crate) active_registrations: PluginRegistrationSnapshot, + pub(crate) delegated_grant_count: usize, + pub(crate) consecutive_failures: u32, } #[derive(Clone)] @@ -132,6 +153,8 @@ struct PluginManagerInner { app_data_root: PathBuf, factory: Arc, access_picker: Arc, + error_history_limit: usize, + lifecycle_bus: LifecycleBus, registry: Mutex, } diff --git a/crates/volt-runner/src/plugin_manager/runtime/lifecycle/shutdown.rs b/crates/volt-runner/src/plugin_manager/runtime/lifecycle/shutdown.rs index e0d5948..979c521 100644 --- a/crates/volt-runner/src/plugin_manager/runtime/lifecycle/shutdown.rs +++ b/crates/volt-runner/src/plugin_manager/runtime/lifecycle/shutdown.rs @@ -2,7 +2,7 @@ use crate::plugin_manager::{PluginManager, PluginState}; impl PluginManager { pub(in crate::plugin_manager) fn deactivate_plugin(&self, plugin_id: &str) { - let (process, state) = { + let (process, state, pre_events) = { let Ok(mut registry) = self.inner.registry.lock() else { return; }; @@ -29,12 +29,26 @@ impl PluginManager { } return; } - let record = registry.plugins.get_mut(plugin_id).expect("checked above"); - if matches!(state, PluginState::Active | PluginState::Running) { - let _ = record.lifecycle.transition(PluginState::Deactivating); + + let mut events = Vec::new(); + if matches!(state, PluginState::Active | PluginState::Running) + && let Ok(event) = self.transition_plugin_locked( + &mut registry, + plugin_id, + PluginState::Deactivating, + ) + { + events.push(event); } - (record.process.clone(), state) + let process = registry + .plugins + .get(plugin_id) + .and_then(|record| record.process.clone()); + (process, state, events) }; + for event in pre_events { + self.emit_lifecycle_event(event); + } let Some(process) = process else { return; @@ -44,7 +58,11 @@ impl PluginManager { } else { process.deactivate(self.deactivation_timeout()) }; - if let Ok(mut registry) = self.inner.registry.lock() { + + let post_events = { + let Ok(mut registry) = self.inner.registry.lock() else { + return; + }; crate::plugin_manager::host_api_helpers::clear_plugin_registrations_locked( &mut registry, plugin_id, @@ -52,26 +70,40 @@ impl PluginManager { if let Some(record) = registry.plugins.get_mut(plugin_id) { record.process = None; record.pending_requests = 0; - match result { - Ok(()) => { - if state == PluginState::Loaded { - record.metrics.pid = None; - let _ = record.lifecycle.transition(PluginState::Terminated); - } else { - let _ = record.lifecycle.transition(PluginState::Terminated); - } - } - Err(error) => { - record.lifecycle.fail( + } + match result { + Ok(()) => { + let already_terminated = registry + .plugins + .get(plugin_id) + .map(|record| record.lifecycle.current_state() == PluginState::Terminated) + .unwrap_or(true); + if already_terminated { + Vec::new() + } else { + self.transition_plugin_locked( + &mut registry, plugin_id, - &error.code, - error.message, - None, - process.stderr_snapshot(), - ); + PluginState::Terminated, + ) + .map(|event| vec![event]) + .unwrap_or_default() } } + Err(error) => self + .fail_plugin_locked( + &mut registry, + plugin_id, + &error.code, + error.message, + None, + process.stderr_snapshot(), + ) + .unwrap_or_default(), } + }; + for event in post_events { + self.emit_lifecycle_event(event); } } } diff --git a/crates/volt-runner/src/plugin_manager/runtime/lifecycle/spawn.rs b/crates/volt-runner/src/plugin_manager/runtime/lifecycle/spawn.rs index 86e94bf..fb2a910 100644 --- a/crates/volt-runner/src/plugin_manager/runtime/lifecycle/spawn.rs +++ b/crates/volt-runner/src/plugin_manager/runtime/lifecycle/spawn.rs @@ -3,42 +3,48 @@ use std::sync::Arc; use crate::plugin_manager::runtime::PluginStartupMode; use crate::plugin_manager::{ HostIpcSettings, PLUGIN_NOT_AVAILABLE_CODE, PLUGIN_RUNTIME_ERROR_CODE, PluginBootstrapConfig, - PluginManager, PluginProcessHandle, PluginRuntimeError, PluginState, now_ms, + PluginLifecycleEvent, PluginManager, PluginProcessHandle, PluginRuntimeError, PluginState, + now_ms, }; impl PluginManager { - #[allow(dead_code)] pub(in crate::plugin_manager) fn ensure_plugin_loaded( &self, plugin_id: &str, ) -> Result, PluginRuntimeError> { - self.ensure_plugin_started(plugin_id, PluginStartupMode::LoadOnly) + self.ensure_plugin_started(plugin_id, PluginStartupMode::LoadOnly, false) } pub(in crate::plugin_manager) fn ensure_plugin_running( &self, plugin_id: &str, ) -> Result, PluginRuntimeError> { - self.ensure_plugin_started(plugin_id, PluginStartupMode::Activate) + self.ensure_plugin_started(plugin_id, PluginStartupMode::Activate, false) + } + + pub(in crate::plugin_manager) fn retry_failed_plugin( + &self, + plugin_id: &str, + ) -> Result, PluginRuntimeError> { + self.ensure_plugin_started(plugin_id, PluginStartupMode::Activate, true) } fn ensure_plugin_started( &self, plugin_id: &str, mode: PluginStartupMode, + allow_retry_from_failed: bool, ) -> Result, PluginRuntimeError> { let spawn_lock = { - let registry = self.inner.registry.lock().map_err(|_| PluginRuntimeError { - code: PLUGIN_RUNTIME_ERROR_CODE.to_string(), - message: "plugin registry is unavailable".to_string(), - })?; + let registry = self + .inner + .registry + .lock() + .map_err(|_| registry_unavailable())?; let record = registry .plugins .get(plugin_id) - .ok_or_else(|| PluginRuntimeError { - code: PLUGIN_NOT_AVAILABLE_CODE.to_string(), - message: format!("plugin '{plugin_id}' is not registered"), - })?; + .ok_or_else(|| unavailable_plugin(plugin_id))?; record.spawn_lock.clone() }; let _guard = spawn_lock.lock().map_err(|_| PluginRuntimeError { @@ -55,7 +61,11 @@ impl PluginManager { return Ok(process); } - let bootstrap = self.prepare_spawn(plugin_id, mode)?; + let (bootstrap, spawn_event) = + self.prepare_spawn(plugin_id, mode, allow_retry_from_failed)?; + if let Some(event) = spawn_event { + self.emit_lifecycle_event(event); + } let process = self.inner.factory.spawn(&bootstrap)?; let manager = self.clone(); let plugin_id_for_exit = plugin_id.to_string(); @@ -78,10 +88,7 @@ impl PluginManager { process.stderr_snapshot(), ); let _ = process.kill(); - return Err(PluginRuntimeError { - code: PLUGIN_NOT_AVAILABLE_CODE.to_string(), - message: format!("plugin '{plugin_id}' failed before ready"), - }); + return Err(unavailable_message(plugin_id, "failed before ready")); } self.transition_plugin(plugin_id, PluginState::Loaded)?; if mode == PluginStartupMode::LoadOnly { @@ -106,16 +113,12 @@ impl PluginManager { process.stderr_snapshot(), ); let _ = process.kill(); - return Err(PluginRuntimeError { - code: PLUGIN_NOT_AVAILABLE_CODE.to_string(), - message: format!("plugin '{plugin_id}' failed to activate"), - }); + return Err(unavailable_message(plugin_id, "failed to activate")); } self.transition_plugin(plugin_id, PluginState::Active)?; self.transition_plugin(plugin_id, PluginState::Running)?; self.record_activity(plugin_id); self.start_watchdog(plugin_id.to_string(), process.clone()); - Ok(()) } @@ -140,84 +143,88 @@ impl PluginManager { &self, plugin_id: &str, mode: PluginStartupMode, - ) -> Result { - let mut registry = self.inner.registry.lock().map_err(|_| PluginRuntimeError { - code: PLUGIN_RUNTIME_ERROR_CODE.to_string(), - message: "plugin registry is unavailable".to_string(), - })?; - let record = registry + allow_retry_from_failed: bool, + ) -> Result<(PluginBootstrapConfig, Option), PluginRuntimeError> { + let mut registry = self + .inner + .registry + .lock() + .map_err(|_| registry_unavailable())?; + let current_state = registry .plugins - .get_mut(plugin_id) - .ok_or_else(|| PluginRuntimeError { - code: PLUGIN_NOT_AVAILABLE_CODE.to_string(), - message: format!("plugin '{plugin_id}' is not registered"), - })?; - - match record.lifecycle.current_state() { - PluginState::Validated | PluginState::Terminated => { - record - .lifecycle - .transition(PluginState::Spawning) - .map_err(|message| PluginRuntimeError { + .get(plugin_id) + .map(|record| record.lifecycle.current_state()) + .ok_or_else(|| unavailable_plugin(plugin_id))?; + let spawn_event = + match current_state { + PluginState::Validated | PluginState::Terminated => Some( + self.transition_plugin_locked(&mut registry, plugin_id, PluginState::Spawning)?, + ), + PluginState::Failed if allow_retry_from_failed => Some( + self.transition_plugin_locked(&mut registry, plugin_id, PluginState::Spawning)?, + ), + PluginState::Loaded if mode == PluginStartupMode::Activate => None, + PluginState::Active | PluginState::Running => None, + PluginState::Disabled => { + return Err(PluginRuntimeError { + code: PLUGIN_NOT_AVAILABLE_CODE.to_string(), + message: format!("plugin '{plugin_id}' is disabled"), + }); + } + PluginState::Failed => { + return Err(PluginRuntimeError { + code: PLUGIN_NOT_AVAILABLE_CODE.to_string(), + message: format!("plugin '{plugin_id}' is in failed state"), + }); + } + other => { + return Err(PluginRuntimeError { code: PLUGIN_RUNTIME_ERROR_CODE.to_string(), - message, - })?; - } - PluginState::Loaded if mode == PluginStartupMode::Activate => {} - PluginState::Active | PluginState::Running => {} - PluginState::Disabled => { - return Err(PluginRuntimeError { - code: PLUGIN_NOT_AVAILABLE_CODE.to_string(), - message: format!("plugin '{plugin_id}' is disabled"), - }); - } - PluginState::Failed => { - return Err(PluginRuntimeError { - code: PLUGIN_NOT_AVAILABLE_CODE.to_string(), - message: format!("plugin '{plugin_id}' is in failed state"), - }); - } - other => { - return Err(PluginRuntimeError { - code: PLUGIN_RUNTIME_ERROR_CODE.to_string(), - message: format!( - "plugin '{plugin_id}' cannot be spawned from state {:?}", - other - ), - }); - } - } + message: format!( + "plugin '{plugin_id}' cannot be spawned from state {:?}", + other + ), + }); + } + }; + let record = registry + .plugins + .get(plugin_id) + .expect("plugin exists after state transition"); let data_root = record.data_root.clone().ok_or_else(|| PluginRuntimeError { code: PLUGIN_RUNTIME_ERROR_CODE.to_string(), message: format!("plugin '{plugin_id}' is missing a data root"), })?; - Ok(PluginBootstrapConfig { - plugin_id: record.manifest.id.clone(), - backend_entry: record.manifest.backend_entry.display().to_string(), - manifest: record.manifest.raw_manifest.clone(), - capabilities: record.effective_capabilities.iter().cloned().collect(), - data_root: data_root.display().to_string(), - delegated_grants: record - .delegated_grants - .iter() - .filter_map(|grant_id| { - volt_core::grant_store::resolve_grant(grant_id) - .ok() - .map(|path| crate::plugin_manager::DelegatedGrant { - grant_id: grant_id.clone(), - path: path.display().to_string(), - }) - }) - .collect(), - host_ipc_settings: HostIpcSettings { - heartbeat_interval_ms: self.inner.config.limits.heartbeat_interval_ms, - heartbeat_timeout_ms: self.inner.config.limits.heartbeat_timeout_ms, - call_timeout_ms: self.inner.config.limits.call_timeout_ms, - max_inflight: 64, - max_queue_depth: 256, + Ok(( + PluginBootstrapConfig { + plugin_id: record.manifest.id.clone(), + backend_entry: record.manifest.backend_entry.display().to_string(), + manifest: record.manifest.raw_manifest.clone(), + capabilities: record.effective_capabilities.iter().cloned().collect(), + data_root: data_root.display().to_string(), + delegated_grants: record + .delegated_grants + .iter() + .filter_map(|grant_id| { + volt_core::grant_store::resolve_grant(grant_id) + .ok() + .map(|path| crate::plugin_manager::DelegatedGrant { + grant_id: grant_id.clone(), + path: path.display().to_string(), + }) + }) + .collect(), + host_ipc_settings: HostIpcSettings { + heartbeat_interval_ms: self.inner.config.limits.heartbeat_interval_ms, + heartbeat_timeout_ms: self.inner.config.limits.heartbeat_timeout_ms, + call_timeout_ms: self.inner.config.limits.call_timeout_ms, + max_inflight: 64, + max_queue_depth: 256, + }, }, - }) + spawn_event, + )) } fn register_process( @@ -235,3 +242,24 @@ impl PluginManager { } } } + +fn registry_unavailable() -> PluginRuntimeError { + PluginRuntimeError { + code: PLUGIN_RUNTIME_ERROR_CODE.to_string(), + message: "plugin registry is unavailable".to_string(), + } +} + +fn unavailable_plugin(plugin_id: &str) -> PluginRuntimeError { + PluginRuntimeError { + code: PLUGIN_NOT_AVAILABLE_CODE.to_string(), + message: format!("plugin '{plugin_id}' is not registered"), + } +} + +fn unavailable_message(plugin_id: &str, reason: &str) -> PluginRuntimeError { + PluginRuntimeError { + code: PLUGIN_NOT_AVAILABLE_CODE.to_string(), + message: format!("plugin '{plugin_id}' {reason}"), + } +} diff --git a/crates/volt-runner/src/plugin_manager/runtime/lifecycle/state.rs b/crates/volt-runner/src/plugin_manager/runtime/lifecycle/state.rs index 960d070..d21c107 100644 --- a/crates/volt-runner/src/plugin_manager/runtime/lifecycle/state.rs +++ b/crates/volt-runner/src/plugin_manager/runtime/lifecycle/state.rs @@ -1,22 +1,80 @@ use std::time::Duration; -use serde_json::Value; +use serde_json::{Value, json}; use crate::plugin_manager::{ - PLUGIN_NOT_AVAILABLE_CODE, PLUGIN_RUNTIME_ERROR_CODE, PluginManager, PluginRuntimeError, + PLUGIN_AUTO_DISABLED_CODE, PLUGIN_NOT_AVAILABLE_CODE, PLUGIN_RUNTIME_ERROR_CODE, + PluginLifecycleEvent, PluginManager, PluginRecord, PluginRegistry, PluginRuntimeError, PluginState, now_ms, }; +const MAX_CONSECUTIVE_FAILURES: u32 = 3; + impl PluginManager { pub(super) fn transition_plugin( &self, plugin_id: &str, next_state: PluginState, ) -> Result<(), PluginRuntimeError> { - let mut registry = self.inner.registry.lock().map_err(|_| PluginRuntimeError { - code: PLUGIN_RUNTIME_ERROR_CODE.to_string(), - message: "plugin registry is unavailable".to_string(), - })?; + let event = { + let mut registry = self.inner.registry.lock().map_err(|_| PluginRuntimeError { + code: PLUGIN_RUNTIME_ERROR_CODE.to_string(), + message: "plugin registry is unavailable".to_string(), + })?; + self.transition_plugin_locked(&mut registry, plugin_id, next_state)? + }; + self.emit_lifecycle_event(event); + Ok(()) + } + + pub(in crate::plugin_manager) fn transition_plugin_locked( + &self, + registry: &mut PluginRegistry, + plugin_id: &str, + next_state: PluginState, + ) -> Result { + let record = registry + .plugins + .get_mut(plugin_id) + .ok_or_else(|| PluginRuntimeError { + code: PLUGIN_NOT_AVAILABLE_CODE.to_string(), + message: format!("plugin '{plugin_id}' is not registered"), + })?; + transition_record(plugin_id, record, next_state) + } + + pub(crate) fn fail_plugin( + &self, + plugin_id: &str, + code: &str, + message: String, + details: Option, + stderr: Option, + ) { + let events = { + let Ok(mut registry) = self.inner.registry.lock() else { + return; + }; + self.fail_plugin_locked(&mut registry, plugin_id, code, message, details, stderr) + .unwrap_or_default() + }; + for event in events { + self.emit_lifecycle_event(event); + } + } + + pub(in crate::plugin_manager) fn fail_plugin_locked( + &self, + registry: &mut PluginRegistry, + plugin_id: &str, + code: &str, + message: String, + details: Option, + stderr: Option, + ) -> Result, PluginRuntimeError> { + crate::plugin_manager::host_api_helpers::clear_plugin_registrations_locked( + registry, plugin_id, + ); let record = registry .plugins .get_mut(plugin_id) @@ -24,13 +82,42 @@ impl PluginManager { code: PLUGIN_NOT_AVAILABLE_CODE.to_string(), message: format!("plugin '{plugin_id}' is not registered"), })?; - record + record.process = None; + record.pending_requests = 0; + + let (transition, error, failures) = record .lifecycle - .transition(next_state) - .map_err(|message| PluginRuntimeError { - code: PLUGIN_RUNTIME_ERROR_CODE.to_string(), + .fail( + plugin_id, + code, message, - }) + details, + stderr, + self.inner.error_history_limit, + ) + .map_err(runtime_error)?; + let mut events = vec![PluginLifecycleEvent { + plugin_id: plugin_id.to_string(), + previous_state: transition.previous_state, + new_state: transition.new_state, + timestamp: transition.timestamp_ms, + error: Some(error), + }]; + + if failures >= MAX_CONSECUTIVE_FAILURES { + events.push(auto_disable_record( + plugin_id, + record, + self.inner.error_history_limit, + failures, + )?); + } + + Ok(events) + } + + pub(in crate::plugin_manager) fn emit_lifecycle_event(&self, event: PluginLifecycleEvent) { + self.inner.lifecycle_bus.emit(event); } pub(in crate::plugin_manager) fn record_activity(&self, plugin_id: &str) { @@ -59,29 +146,6 @@ impl PluginManager { } } - pub(in crate::plugin_manager) fn fail_plugin( - &self, - plugin_id: &str, - code: &str, - message: String, - details: Option, - stderr: Option, - ) { - if let Ok(mut registry) = self.inner.registry.lock() { - crate::plugin_manager::host_api_helpers::clear_plugin_registrations_locked( - &mut registry, - plugin_id, - ); - if let Some(record) = registry.plugins.get_mut(plugin_id) { - record.process = None; - record.pending_requests = 0; - record - .lifecycle - .fail(plugin_id, code, message, details, stderr); - } - } - } - pub(super) fn activation_timeout(&self) -> Duration { Duration::from_millis(self.inner.config.limits.activation_timeout_ms) } @@ -90,3 +154,59 @@ impl PluginManager { Duration::from_millis(self.inner.config.limits.deactivation_timeout_ms) } } + +fn transition_record( + plugin_id: &str, + record: &mut PluginRecord, + next_state: PluginState, +) -> Result { + let transition = record + .lifecycle + .transition(next_state) + .map_err(runtime_error)?; + Ok(PluginLifecycleEvent { + plugin_id: plugin_id.to_string(), + previous_state: transition.previous_state, + new_state: transition.new_state, + timestamp: transition.timestamp_ms, + error: None, + }) +} + +fn auto_disable_record( + plugin_id: &str, + record: &mut PluginRecord, + max_errors: usize, + failures: u32, +) -> Result { + let transition = record + .lifecycle + .transition(PluginState::Disabled) + .map_err(runtime_error)?; + let error = record.lifecycle.push_error( + crate::plugin_manager::PluginError { + plugin_id: plugin_id.to_string(), + state: PluginState::Disabled, + code: PLUGIN_AUTO_DISABLED_CODE.to_string(), + message: format!("plugin auto-disabled after {failures} consecutive failures"), + details: Some(json!({ "consecutiveFailures": failures })), + stderr: None, + timestamp_ms: transition.timestamp_ms, + }, + max_errors, + ); + Ok(PluginLifecycleEvent { + plugin_id: plugin_id.to_string(), + previous_state: transition.previous_state, + new_state: transition.new_state, + timestamp: transition.timestamp_ms, + error: Some(error), + }) +} + +fn runtime_error(message: String) -> PluginRuntimeError { + PluginRuntimeError { + code: PLUGIN_RUNTIME_ERROR_CODE.to_string(), + message, + } +} diff --git a/crates/volt-runner/src/plugin_manager/runtime/mod.rs b/crates/volt-runner/src/plugin_manager/runtime/mod.rs index c6c25a9..bdf1714 100644 --- a/crates/volt-runner/src/plugin_manager/runtime/mod.rs +++ b/crates/volt-runner/src/plugin_manager/runtime/mod.rs @@ -7,11 +7,13 @@ use volt_core::ipc::{ use super::{ DEFAULT_PRE_SPAWN_GRACE_MS, PLUGIN_IPC_HANDLER_NOT_FOUND_CODE, PLUGIN_ROUTE_INVALID_CODE, - PluginDiscoveryIssue, PluginManager, PluginSnapshot, PluginState, parse_plugin_route, + PluginDiscoveryIssue, PluginManager, PluginState, parse_plugin_route, }; use crate::runner::config::RunnerPluginSpawningStrategy; mod lifecycle; +mod observability; +mod recovery; #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub(crate) enum PluginStartupMode { @@ -20,28 +22,6 @@ pub(crate) enum PluginStartupMode { } impl PluginManager { - #[allow(dead_code)] - pub(crate) fn get_plugin_state(&self, plugin_id: &str) -> Option { - let registry = self.inner.registry.lock().ok()?; - let record = registry.plugins.get(plugin_id)?; - Some(record.snapshot()) - } - - #[allow(dead_code)] - pub(crate) fn get_states(&self) -> Vec { - let Ok(registry) = self.inner.registry.lock() else { - return Vec::new(); - }; - let mut states = registry - .plugins - .values() - .map(super::PluginRecord::snapshot) - .collect::>(); - states.sort_by(|left, right| left.plugin_id.cmp(&right.plugin_id)); - states - } - - #[allow(dead_code)] pub(crate) fn discovery_issues(&self) -> Vec { let Ok(registry) = self.inner.registry.lock() else { return Vec::new(); diff --git a/crates/volt-runner/src/plugin_manager/runtime/observability.rs b/crates/volt-runner/src/plugin_manager/runtime/observability.rs new file mode 100644 index 0000000..3704cd4 --- /dev/null +++ b/crates/volt-runner/src/plugin_manager/runtime/observability.rs @@ -0,0 +1,76 @@ +use crate::plugin_manager::{ + PluginError, PluginLifecycleEvent, PluginManager, PluginRecord, PluginStateSnapshot, + SubscriptionId, +}; + +impl PluginManager { + pub(crate) fn on_lifecycle( + &self, + handler: Box, + ) -> SubscriptionId { + self.inner.lifecycle_bus.on_lifecycle(handler) + } + + pub(crate) fn on_plugin_failed( + &self, + handler: Box, + ) -> SubscriptionId { + self.inner.lifecycle_bus.on_plugin_failed(handler) + } + + pub(crate) fn on_plugin_activated( + &self, + handler: Box, + ) -> SubscriptionId { + self.inner.lifecycle_bus.on_plugin_activated(handler) + } + + #[allow(dead_code)] + pub(crate) fn off(&self, subscription_id: SubscriptionId) { + self.inner.lifecycle_bus.off(subscription_id); + } + + pub(crate) fn get_plugin_state(&self, plugin_id: &str) -> Option { + let registry = self.inner.registry.lock().ok()?; + let record = registry.plugins.get(plugin_id)?; + Some(record.snapshot()) + } + + pub(crate) fn get_states(&self) -> Vec { + let Ok(registry) = self.inner.registry.lock() else { + return Vec::new(); + }; + let mut states = registry + .plugins + .values() + .map(PluginRecord::snapshot) + .collect::>(); + states.sort_by(|left, right| left.plugin_id.cmp(&right.plugin_id)); + states + } + + pub(crate) fn get_errors(&self) -> Vec { + let Ok(registry) = self.inner.registry.lock() else { + return Vec::new(); + }; + let mut errors = registry + .plugins + .values() + .flat_map(|record| record.lifecycle.errors.clone()) + .collect::>(); + errors.sort_by(|left, right| right.timestamp_ms.cmp(&left.timestamp_ms)); + errors + } + + pub(crate) fn get_plugin_errors(&self, plugin_id: &str) -> Vec { + let Ok(registry) = self.inner.registry.lock() else { + return Vec::new(); + }; + let Some(record) = registry.plugins.get(plugin_id) else { + return Vec::new(); + }; + let mut errors = record.lifecycle.errors.clone(); + errors.sort_by(|left, right| right.timestamp_ms.cmp(&left.timestamp_ms)); + errors + } +} diff --git a/crates/volt-runner/src/plugin_manager/runtime/recovery.rs b/crates/volt-runner/src/plugin_manager/runtime/recovery.rs new file mode 100644 index 0000000..e5df54f --- /dev/null +++ b/crates/volt-runner/src/plugin_manager/runtime/recovery.rs @@ -0,0 +1,64 @@ +use crate::plugin_manager::{ + PLUGIN_NOT_AVAILABLE_CODE, PLUGIN_RUNTIME_ERROR_CODE, PluginManager, PluginRuntimeError, + PluginState, +}; + +impl PluginManager { + pub(crate) fn retry_plugin(&self, plugin_id: &str) -> Result<(), PluginRuntimeError> { + let state = { + let registry = self.inner.registry.lock().map_err(|_| PluginRuntimeError { + code: PLUGIN_RUNTIME_ERROR_CODE.to_string(), + message: "plugin registry is unavailable".to_string(), + })?; + registry + .plugins + .get(plugin_id) + .map(|record| record.lifecycle.current_state()) + .ok_or_else(|| PluginRuntimeError { + code: PLUGIN_NOT_AVAILABLE_CODE.to_string(), + message: format!("plugin '{plugin_id}' is not registered"), + })? + }; + if state != PluginState::Failed { + return Err(PluginRuntimeError { + code: PLUGIN_NOT_AVAILABLE_CODE.to_string(), + message: format!("plugin '{plugin_id}' is not in failed state"), + }); + } + self.retry_failed_plugin(plugin_id).map(|_| ()) + } + + pub(crate) fn enable_plugin(&self, plugin_id: &str) -> Result<(), PluginRuntimeError> { + let event = { + let mut registry = self.inner.registry.lock().map_err(|_| PluginRuntimeError { + code: PLUGIN_RUNTIME_ERROR_CODE.to_string(), + message: "plugin registry is unavailable".to_string(), + })?; + let record = registry + .plugins + .get_mut(plugin_id) + .ok_or_else(|| PluginRuntimeError { + code: PLUGIN_NOT_AVAILABLE_CODE.to_string(), + message: format!("plugin '{plugin_id}' is not registered"), + })?; + if !record.enabled { + return Err(PluginRuntimeError { + code: PLUGIN_NOT_AVAILABLE_CODE.to_string(), + message: format!("plugin '{plugin_id}' is disabled by configuration"), + }); + } + if record.lifecycle.current_state() != PluginState::Disabled { + return Err(PluginRuntimeError { + code: PLUGIN_NOT_AVAILABLE_CODE.to_string(), + message: format!("plugin '{plugin_id}' is not disabled"), + }); + } + record.process = None; + record.pending_requests = 0; + record.lifecycle.reset_failures(); + self.transition_plugin_locked(&mut registry, plugin_id, PluginState::Validated)? + }; + self.emit_lifecycle_event(event); + Ok(()) + } +} diff --git a/crates/volt-runner/src/plugin_manager/tests/activation.rs b/crates/volt-runner/src/plugin_manager/tests/activation.rs index 0d18e1a..0dbf735 100644 --- a/crates/volt-runner/src/plugin_manager/tests/activation.rs +++ b/crates/volt-runner/src/plugin_manager/tests/activation.rs @@ -56,7 +56,7 @@ fn lazy_spawn_happens_on_first_ipc_request() { assert_eq!(response.result, Some(serde_json::json!({ "ok": true }))); assert_eq!(factory.spawn_count.load(Ordering::Relaxed), 1); let snapshot = manager.get_plugin_state("acme.search").expect("plugin"); - assert_eq!(snapshot.state, PluginState::Running); + assert_eq!(snapshot.current_state, PluginState::Running); assert_eq!(snapshot.metrics.pid, Some(42)); assert!(snapshot.metrics.started_at_ms.is_some()); assert!(snapshot.metrics.last_activity_ms.is_some()); @@ -107,7 +107,7 @@ fn spawn_timeout_moves_plugin_to_failed() { manager .get_plugin_state("acme.search") .expect("plugin") - .state, + .current_state, PluginState::Failed ); let snapshot = manager.get_plugin_state("acme.search").expect("plugin"); @@ -155,7 +155,7 @@ fn spawn_crash_moves_plugin_to_failed() { Some(PLUGIN_NOT_AVAILABLE_CODE) ); let snapshot = manager.get_plugin_state("acme.search").expect("plugin"); - assert_eq!(snapshot.state, PluginState::Failed); + assert_eq!(snapshot.current_state, PluginState::Failed); assert!(snapshot.errors.iter().any(|error| { error .details @@ -206,7 +206,7 @@ fn activation_error_moves_plugin_to_failed() { Some(PLUGIN_NOT_AVAILABLE_CODE) ); let snapshot = manager.get_plugin_state("acme.search").expect("plugin"); - assert_eq!(snapshot.state, PluginState::Failed); + assert_eq!(snapshot.current_state, PluginState::Failed); assert!( snapshot .errors @@ -247,7 +247,7 @@ fn pre_spawn_forces_startup_activation_after_window_ready_hook() { manager .get_plugin_state("acme.search") .expect("plugin") - .state, + .current_state, PluginState::Running ); } diff --git a/crates/volt-runner/src/plugin_manager/tests/discovery.rs b/crates/volt-runner/src/plugin_manager/tests/discovery.rs index 7353c85..121e2e7 100644 --- a/crates/volt-runner/src/plugin_manager/tests/discovery.rs +++ b/crates/volt-runner/src/plugin_manager/tests/discovery.rs @@ -33,7 +33,7 @@ fn discovery_finds_manifests_and_reports_missing_directories() { manager .get_plugin_state("acme.search") .expect("plugin") - .state, + .current_state, PluginState::Validated ); assert_eq!(manager.discovery_issues().len(), 1); @@ -106,7 +106,7 @@ fn capability_intersection_rejects_unsatisfied_plugins_and_keeps_exact_matches() ); let failed = manager.get_plugin_state("acme.search").expect("search"); - assert_eq!(failed.state, PluginState::Failed); + assert_eq!(failed.current_state, PluginState::Failed); assert_eq!(failed.plugin_id, "acme.search"); assert!(failed.enabled); assert!( @@ -120,12 +120,15 @@ fn capability_intersection_rejects_unsatisfied_plugins_and_keeps_exact_matches() ); assert_eq!(failed.effective_capabilities, vec!["fs".to_string()]); assert!(failed.data_root.as_ref().expect("data root").exists()); - assert_eq!(failed.transitions.len(), 2); - assert_eq!(failed.transitions[0].new_state, PluginState::Discovered); - assert_eq!(failed.transitions[1].new_state, PluginState::Failed); + assert_eq!(failed.transition_history.len(), 2); + assert_eq!( + failed.transition_history[0].new_state, + PluginState::Discovered + ); + assert_eq!(failed.transition_history[1].new_state, PluginState::Failed); assert_eq!(failed.errors.len(), 1); assert_eq!(failed.errors[0].plugin_id, "acme.search"); - assert_eq!(failed.errors[0].state, PluginState::Failed); + assert_eq!(failed.errors[0].state, PluginState::Discovered); assert_eq!(failed.errors[0].code, PLUGIN_NOT_AVAILABLE_CODE); assert!(failed.errors[0].message.contains("unsatisfiable")); assert!(failed.errors[0].details.is_some()); @@ -136,7 +139,7 @@ fn capability_intersection_rejects_unsatisfied_plugins_and_keeps_exact_matches() assert!(!failed.process_running); let exact = manager.get_plugin_state("acme.clip").expect("clip"); - assert_eq!(exact.state, PluginState::Validated); + assert_eq!(exact.current_state, PluginState::Validated); assert_eq!(exact.effective_capabilities, vec!["fs".to_string()]); } @@ -169,7 +172,7 @@ fn boot_rule_validation_does_not_spawn_plugin_processes() { manager .get_plugin_state("acme.search") .expect("plugin") - .state, + .current_state, PluginState::Validated ); assert!( diff --git a/crates/volt-runner/src/plugin_manager/tests/lifecycle.rs b/crates/volt-runner/src/plugin_manager/tests/lifecycle.rs index e0ef161..3623ca7 100644 --- a/crates/volt-runner/src/plugin_manager/tests/lifecycle.rs +++ b/crates/volt-runner/src/plugin_manager/tests/lifecycle.rs @@ -3,12 +3,13 @@ use std::sync::Arc; use std::time::Duration; use serde_json::Value; +use volt_core::grant_store; use volt_core::ipc::IpcRequest; use super::super::*; use super::fs_support::{TempDir, write_manifest}; use super::process_support::{FakePlan, FakeProcessFactory, FakeRequestOutcome}; -use super::shared::{manager_with_factory, register_ipc_handler}; +use super::shared::{lock_grant_state, manager_with_factory, register_ipc_handler}; use crate::runner::config::RunnerPluginConfig; #[test] @@ -108,6 +109,60 @@ fn shutdown_all_deactivates_running_plugins_cleanly() { manager.shutdown_all(); let snapshot = manager.get_plugin_state("acme.search").expect("plugin"); - assert_eq!(snapshot.state, PluginState::Terminated); + assert_eq!(snapshot.current_state, PluginState::Terminated); assert!(!snapshot.process_running); } + +#[test] +fn state_snapshot_reports_active_registrations_and_grants() { + let _guard = lock_grant_state(); + let root = TempDir::new("snapshot-counts"); + write_manifest( + &root.join("plugins/acme.search/volt-plugin.json"), + "acme.search", + &["fs"], + ); + let grant_path = root.join("granted"); + std::fs::create_dir_all(&grant_path).expect("grant path"); + let grant_id = grant_store::create_grant(grant_path).expect("grant"); + let manager = manager_with_factory( + RunnerPluginConfig { + enabled: vec!["acme.search".to_string()], + grants: BTreeMap::from([("acme.search".to_string(), vec!["fs".to_string()])]), + plugin_dirs: vec![root.join("plugins").display().to_string()], + ..RunnerPluginConfig::default() + }, + Arc::new(FakeProcessFactory::new(HashMap::new())), + ); + + manager + .delegate_grant("acme.search", &grant_id) + .expect("delegate grant"); + let _ = manager.handle_plugin_message( + "acme.search", + crate::plugin_manager::process::WireMessage { + message_type: WireMessageType::Request, + id: "register-command".to_string(), + method: "plugin:register-command".to_string(), + payload: Some(serde_json::json!({ "id": "reindex" })), + error: None, + }, + ); + let _ = manager.handle_plugin_message( + "acme.search", + crate::plugin_manager::process::WireMessage { + message_type: WireMessageType::Request, + id: "subscribe-event".to_string(), + method: "plugin:subscribe-event".to_string(), + payload: Some(serde_json::json!({ "event": "app:focus" })), + error: None, + }, + ); + register_ipc_handler(&manager, "acme.search", "ping"); + + let snapshot = manager.get_plugin_state("acme.search").expect("plugin"); + assert_eq!(snapshot.active_registrations.command_count, 1); + assert_eq!(snapshot.active_registrations.event_subscription_count, 1); + assert_eq!(snapshot.active_registrations.ipc_handler_count, 1); + assert_eq!(snapshot.delegated_grant_count, 1); +} diff --git a/crates/volt-runner/src/plugin_manager/tests/lifecycle_bus.rs b/crates/volt-runner/src/plugin_manager/tests/lifecycle_bus.rs new file mode 100644 index 0000000..110a900 --- /dev/null +++ b/crates/volt-runner/src/plugin_manager/tests/lifecycle_bus.rs @@ -0,0 +1,165 @@ +use std::collections::{BTreeMap, HashMap}; +use std::sync::{Arc, Mutex}; +use std::time::Duration; + +use serde_json::Value; +use volt_core::ipc::IpcRequest; + +use super::super::*; +use super::fs_support::{TempDir, write_manifest}; +use super::process_support::{FakePlan, FakeProcessFactory, FakeRequestOutcome}; +use super::shared::{manager_with_factory, register_ipc_handler}; +use crate::runner::config::RunnerPluginConfig; + +fn build_manager_with_root(name: &str) -> (TempDir, PluginManager) { + let root = TempDir::new(name); + write_manifest( + &root.join("plugins/acme.search/volt-plugin.json"), + "acme.search", + &["fs"], + ); + let manager = manager_with_factory( + RunnerPluginConfig { + enabled: vec!["acme.search".to_string()], + grants: BTreeMap::from([("acme.search".to_string(), vec!["fs".to_string()])]), + plugin_dirs: vec![root.join("plugins").display().to_string()], + ..RunnerPluginConfig::default() + }, + Arc::new(FakeProcessFactory::new(HashMap::from([( + "acme.search".to_string(), + FakePlan { + requests: HashMap::from([( + "plugin:invoke-ipc".to_string(), + FakeRequestOutcome::Success(serde_json::json!({ "ok": true })), + )]), + ..FakePlan::default() + }, + )]))), + ); + (root, manager) +} + +#[test] +fn lifecycle_bus_replays_discovery_transitions_when_rediscovering() { + let (_root, manager) = build_manager_with_root("lifecycle-bus-redisco"); + let events = Arc::new(Mutex::new(Vec::new())); + let events_for_handler = events.clone(); + manager.on_lifecycle(Box::new(move |event| { + events_for_handler + .lock() + .expect("events") + .push(event.clone()); + })); + + manager.discover_plugins(); + + let states = events + .lock() + .expect("events") + .iter() + .map(|event| event.new_state) + .collect::>(); + assert_eq!( + states, + vec![PluginState::Discovered, PluginState::Validated] + ); +} + +#[test] +fn lifecycle_bus_emits_runtime_transitions_in_order_to_multiple_subscribers() { + let (_root, manager) = build_manager_with_root("lifecycle-bus-runtime"); + register_ipc_handler(&manager, "acme.search", "ping"); + let first = Arc::new(Mutex::new(Vec::new())); + let second = Arc::new(Mutex::new(Vec::new())); + let first_for_handler = first.clone(); + let second_for_handler = second.clone(); + manager.on_lifecycle(Box::new(move |event| { + first_for_handler + .lock() + .expect("first") + .push(event.new_state); + })); + manager.on_lifecycle(Box::new(move |event| { + second_for_handler + .lock() + .expect("second") + .push(event.new_state); + })); + + let _ = manager.handle_ipc_request( + &IpcRequest { + id: "req-1".to_string(), + method: "plugin:acme.search:ping".to_string(), + args: Value::Null, + }, + Duration::from_millis(50), + ); + manager.shutdown_all(); + + let expected = vec![ + PluginState::Spawning, + PluginState::Loaded, + PluginState::Active, + PluginState::Running, + PluginState::Deactivating, + PluginState::Terminated, + ]; + assert_eq!(*first.lock().expect("first"), expected); + assert_eq!(*second.lock().expect("second"), expected); +} + +#[test] +fn failed_and_activated_subscribers_are_filtered_and_off_removes_handlers() { + let (_root, manager) = build_manager_with_root("lifecycle-bus-filters"); + let failed = Arc::new(Mutex::new(Vec::new())); + let activated = Arc::new(Mutex::new(Vec::new())); + let removed = Arc::new(Mutex::new(Vec::new())); + let failed_for_handler = failed.clone(); + let activated_for_handler = activated.clone(); + let removed_for_handler = removed.clone(); + manager.on_plugin_failed(Box::new(move |event| { + failed_for_handler + .lock() + .expect("failed") + .push(event.clone()); + })); + manager.on_plugin_activated(Box::new(move |event| { + activated_for_handler + .lock() + .expect("activated") + .push(event.new_state); + })); + let removed_subscription = manager.on_lifecycle(Box::new(move |event| { + removed_for_handler + .lock() + .expect("removed") + .push(event.new_state); + })); + manager.off(removed_subscription); + + manager.fail_plugin( + "acme.search", + "PLUGIN_BROKEN", + "boom".to_string(), + Some(serde_json::json!({ "attempt": 1 })), + Some("stderr".to_string()), + ); + manager.retry_plugin("acme.search").expect("retry"); + + let failed = failed.lock().expect("failed"); + assert_eq!(failed.len(), 1); + assert_eq!(failed[0].new_state, PluginState::Failed); + assert_eq!( + failed[0].error.as_ref().expect("error").code, + "PLUGIN_BROKEN" + ); + assert_eq!( + failed[0].error.as_ref().expect("error").stderr.as_deref(), + Some("stderr") + ); + assert_eq!( + *activated.lock().expect("activated"), + vec![PluginState::Active] + ); + assert!(removed.lock().expect("removed").is_empty()); +} diff --git a/crates/volt-runner/src/plugin_manager/tests/mod.rs b/crates/volt-runner/src/plugin_manager/tests/mod.rs index 1895ac1..e5caa62 100644 --- a/crates/volt-runner/src/plugin_manager/tests/mod.rs +++ b/crates/volt-runner/src/plugin_manager/tests/mod.rs @@ -4,9 +4,11 @@ mod discovery; mod fs_support; mod grant_delegation; mod lifecycle; +mod lifecycle_bus; mod manifest; mod prefetch; mod process_support; +mod recovery; mod registrations; mod request_access; mod request_runtime; diff --git a/crates/volt-runner/src/plugin_manager/tests/prefetch.rs b/crates/volt-runner/src/plugin_manager/tests/prefetch.rs index b2e827e..2c3e13e 100644 --- a/crates/volt-runner/src/plugin_manager/tests/prefetch.rs +++ b/crates/volt-runner/src/plugin_manager/tests/prefetch.rs @@ -66,11 +66,14 @@ fn prefetch_spawns_matching_plugin_without_activation() { manager .get_plugin_state("acme.search") .expect("search") - .state, + .current_state, PluginState::Loaded ); assert_eq!( - manager.get_plugin_state("beta.index").expect("beta").state, + manager + .get_plugin_state("beta.index") + .expect("beta") + .current_state, PluginState::Validated ); } @@ -96,7 +99,7 @@ fn prefetch_ignores_non_matching_surfaces() { manager .get_plugin_state("acme.search") .expect("search") - .state, + .current_state, PluginState::Validated ); } diff --git a/crates/volt-runner/src/plugin_manager/tests/recovery.rs b/crates/volt-runner/src/plugin_manager/tests/recovery.rs new file mode 100644 index 0000000..eb5f724 --- /dev/null +++ b/crates/volt-runner/src/plugin_manager/tests/recovery.rs @@ -0,0 +1,97 @@ +use std::collections::{BTreeMap, HashMap}; +use std::sync::Arc; +use std::time::Duration; + +use super::super::*; +use super::fs_support::{TempDir, write_manifest}; +use super::process_support::{FakePlan, FakeProcessFactory}; +use super::shared::manager_with_error_history_limit; +use crate::runner::config::RunnerPluginConfig; + +fn build_manager(error_history_limit: usize) -> PluginManager { + let root = TempDir::new("recovery"); + write_manifest( + &root.join("plugins/acme.search/volt-plugin.json"), + "acme.search", + &["fs"], + ); + manager_with_error_history_limit( + RunnerPluginConfig { + enabled: vec!["acme.search".to_string()], + grants: BTreeMap::from([("acme.search".to_string(), vec!["fs".to_string()])]), + plugin_dirs: vec![root.join("plugins").display().to_string()], + ..RunnerPluginConfig::default() + }, + Arc::new(FakeProcessFactory::new(HashMap::from([( + "acme.search".to_string(), + FakePlan::default(), + )]))), + error_history_limit, + ) +} + +#[test] +fn retry_plugin_reactivates_failed_plugin() { + let manager = build_manager(50); + + manager.fail_plugin( + "acme.search", + "PLUGIN_BROKEN", + "boom".to_string(), + None, + None, + ); + manager.retry_plugin("acme.search").expect("retry"); + + let snapshot = manager.get_plugin_state("acme.search").expect("plugin"); + assert_eq!(snapshot.current_state, PluginState::Running); + assert_eq!(snapshot.consecutive_failures, 0); +} + +#[test] +fn error_history_respects_cap_and_returns_descending_timestamps() { + let manager = build_manager(2); + + manager.fail_plugin("acme.search", "E1", "boom-1".to_string(), None, None); + std::thread::sleep(Duration::from_millis(2)); + manager.fail_plugin("acme.search", "E2", "boom-2".to_string(), None, None); + std::thread::sleep(Duration::from_millis(2)); + manager.fail_plugin("acme.search", "E3", "boom-3".to_string(), None, None); + + let errors = manager.get_plugin_errors("acme.search"); + assert_eq!(errors.len(), 2); + assert!(errors[0].timestamp_ms >= errors[1].timestamp_ms); + assert!(errors.iter().all(|error| error.code != "E1")); + assert_eq!(manager.get_errors().len(), 2); +} + +#[test] +fn three_consecutive_failures_auto_disable_and_enable_resets_streak() { + let manager = build_manager(50); + + for attempt in 1..=3 { + manager.fail_plugin( + "acme.search", + "PLUGIN_BROKEN", + format!("boom-{attempt}"), + None, + None, + ); + } + + let disabled = manager.get_plugin_state("acme.search").expect("plugin"); + assert_eq!(disabled.current_state, PluginState::Disabled); + assert_eq!(disabled.consecutive_failures, 3); + assert!( + disabled + .errors + .iter() + .any(|error| error.code == PLUGIN_AUTO_DISABLED_CODE) + ); + + manager.enable_plugin("acme.search").expect("enable"); + + let enabled = manager.get_plugin_state("acme.search").expect("plugin"); + assert_eq!(enabled.current_state, PluginState::Validated); + assert_eq!(enabled.consecutive_failures, 0); +} diff --git a/crates/volt-runner/src/plugin_manager/tests/request_runtime.rs b/crates/volt-runner/src/plugin_manager/tests/request_runtime.rs index bb994bd..bd35302 100644 --- a/crates/volt-runner/src/plugin_manager/tests/request_runtime.rs +++ b/crates/volt-runner/src/plugin_manager/tests/request_runtime.rs @@ -60,7 +60,7 @@ fn request_timeout_maps_to_ipc_timeout_error_without_crashing_plugin() { manager .get_plugin_state("acme.search") .expect("plugin") - .state, + .current_state, PluginState::Running ); } @@ -110,7 +110,7 @@ fn request_crash_transitions_running_plugin_to_failed() { ); thread::sleep(Duration::from_millis(20)); let snapshot = manager.get_plugin_state("acme.search").expect("plugin"); - assert_eq!(snapshot.state, PluginState::Failed); + assert_eq!(snapshot.current_state, PluginState::Failed); assert!(snapshot.errors.iter().any(|error| { error .details @@ -167,7 +167,7 @@ fn watchdog_kills_after_two_missed_heartbeats() { manager .get_plugin_state("acme.search") .expect("plugin") - .state, + .current_state, PluginState::Failed ); } diff --git a/crates/volt-runner/src/plugin_manager/tests/shared.rs b/crates/volt-runner/src/plugin_manager/tests/shared.rs index 725e39a..e3363c8 100644 --- a/crates/volt-runner/src/plugin_manager/tests/shared.rs +++ b/crates/volt-runner/src/plugin_manager/tests/shared.rs @@ -33,7 +33,16 @@ pub(super) fn manager_with_picker( factory: Arc, picker: Arc, ) -> PluginManager { - PluginManager::with_dependencies( + manager_with_picker_and_error_history_limit(config, factory, picker, 50) +} + +pub(super) fn manager_with_picker_and_error_history_limit( + config: RunnerPluginConfig, + factory: Arc, + picker: Arc, + error_history_limit: usize, +) -> PluginManager { + PluginManager::with_dependencies_and_error_history_limit( "Volt Test".to_string(), &[ "fs".to_string(), @@ -43,10 +52,24 @@ pub(super) fn manager_with_picker( config, factory, picker, + error_history_limit, ) .expect("manager") } +pub(super) fn manager_with_error_history_limit( + config: RunnerPluginConfig, + factory: Arc, + error_history_limit: usize, +) -> PluginManager { + manager_with_picker_and_error_history_limit( + config, + factory, + Arc::new(FakeAccessPicker::default()), + error_history_limit, + ) +} + #[allow(dead_code)] pub(super) fn factory_from_empty() -> Arc { Arc::new(FakeProcessFactory::new(std::collections::HashMap::new())) diff --git a/crates/volt-runner/src/plugin_manager/watchdog.rs b/crates/volt-runner/src/plugin_manager/watchdog.rs index 7dd8e59..73fab84 100644 --- a/crates/volt-runner/src/plugin_manager/watchdog.rs +++ b/crates/volt-runner/src/plugin_manager/watchdog.rs @@ -9,7 +9,10 @@ use super::{ impl PluginManager { pub(super) fn handle_process_exit(&self, plugin_id: &str, exit: ProcessExitInfo) { - if let Ok(mut registry) = self.inner.registry.lock() { + let events = { + let Ok(mut registry) = self.inner.registry.lock() else { + return; + }; crate::plugin_manager::host_api_helpers::clear_plugin_registrations_locked( &mut registry, plugin_id, @@ -17,12 +20,23 @@ impl PluginManager { if let Some(record) = registry.plugins.get_mut(plugin_id) { record.process = None; record.pending_requests = 0; - match record.lifecycle.current_state() { - PluginState::Deactivating | PluginState::Terminated | PluginState::Disabled => { - let _ = record.lifecycle.transition(PluginState::Terminated); - } - PluginState::Failed => {} - _ => record.lifecycle.fail( + } + + match registry + .plugins + .get(plugin_id) + .map(|record| record.lifecycle.current_state()) + { + Some( + PluginState::Deactivating | PluginState::Terminated | PluginState::Disabled, + ) => self + .transition_plugin_locked(&mut registry, plugin_id, PluginState::Terminated) + .map(|event| vec![event]) + .unwrap_or_default(), + Some(PluginState::Failed) => Vec::new(), + Some(_) => self + .fail_plugin_locked( + &mut registry, plugin_id, PLUGIN_RUNTIME_ERROR_CODE, format!( @@ -31,9 +45,13 @@ impl PluginManager { ), Some(serde_json::json!({ "exitCode": exit.code })), None, - ), - } + ) + .unwrap_or_default(), + None => Vec::new(), } + }; + for event in events { + self.emit_lifecycle_event(event); } }