//! Wasmtime sandbox and resident Actor lifecycle. //! //! Each started service revision owns one thread, Store, Component Instance, //! and bounded synchronous mailbox. The Actor processes commands serially, so //! component memory remains resident without allowing concurrent entry into the //! same Store. use std::{ collections::HashMap, fs, sync::{ Arc, Mutex, atomic::{AtomicBool, Ordering}, mpsc::{self, Receiver, SyncSender, TrySendError}, }, thread, time::{Duration, Instant}, }; use wasmtime::{ Config, Engine, Store, StoreLimits, StoreLimitsBuilder, component::{Component, HasData, Linker, ResourceTable}, }; use wasmtime_wasi::{ WasiCtx, WasiCtxBuilder, WasiCtxView, WasiView, cli::{WasiCli, WasiCliView as _}, clocks::{WasiClocks, WasiClocksView as _}, filesystem::{WasiFilesystem, WasiFilesystemView as _}, p2::bindings::sync::{cli, clocks, filesystem, io}, }; use crate::{ CapabilityDescriptor, KvBackend, RuntimeError, ServiceKey, ServiceManifest, bindings::{ServiceComponent, ServiceError}, capability::{Capability, CapabilityRegistry, HostCapabilities}, manifest::ResourceLimits, }; const DEFAULT_EPOCH_TICK: Duration = Duration::from_millis(5); const MAX_EPOCH_TICK: Duration = Duration::from_millis(10); const DEFAULT_MAX_WASM_STACK: usize = 512 * 1024; const ACTOR_RESOURCE_COUNT_LIMIT: usize = 16; const ACTOR_TABLE_ELEMENT_LIMIT: usize = 16 * 1024; const ALLOWED_WASI_IMPORTS: &[&str] = &[ "wasi:cli/environment@", "wasi:cli/exit@", "wasi:cli/stderr@", "wasi:cli/stdin@", "wasi:cli/stdout@", "wasi:clocks/wall-clock@", "wasi:filesystem/preopens@", "wasi:filesystem/types@", "wasi:io/error@", "wasi:io/streams@", ]; /// Process-wide Wasmtime Engine settings. #[derive(Clone, Debug)] pub struct RuntimeConfig { /// Frequency used to advance Wasmtime epoch interruption. pub epoch_tick: Duration, /// Maximum native stack reservation for WebAssembly execution. pub max_wasm_stack: usize, /// Service-scoped storage used by Components importing the KV capability. pub kv_backend: Option>, } impl Default for RuntimeConfig { fn default() -> Self { Self { epoch_tick: DEFAULT_EPOCH_TICK, max_wasm_stack: DEFAULT_MAX_WASM_STACK, kv_backend: None, } } } impl RuntimeConfig { fn validate(&self) -> Result<(), RuntimeError> { if self.epoch_tick.is_zero() { return Err(RuntimeError::InvalidManifest( "runtime epoch_tick must be greater than zero".to_owned(), )); } if self.epoch_tick > MAX_EPOCH_TICK { return Err(RuntimeError::InvalidManifest(format!( "runtime epoch_tick must not exceed {MAX_EPOCH_TICK:?}" ))); } if self.max_wasm_stack == 0 { return Err(RuntimeError::InvalidManifest( "runtime max_wasm_stack must be greater than zero".to_owned(), )); } Ok(()) } } /// Registry and lifecycle manager for resident Component Actors. pub struct Runtime { inner: Arc, } /// Point-in-time Runtime service and Actor counts. #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub struct RuntimeStats { /// Number of compiled, registered service revisions. pub registered_services: usize, /// Number of Actors currently accepting calls. pub running_actors: usize, } struct RuntimeInner { engine: Arc, epoch_tick: Duration, services: Mutex>, actors: Mutex>, lifecycle: Mutex<()>, ticker: Arc, kv_backend: Option>, } struct EpochTicker { shutdown: mpsc::Sender<()>, worker: Mutex>>, } impl EpochTicker { fn start(engine: Arc, tick: Duration) -> Result, RuntimeError> { let (shutdown, receiver) = mpsc::channel(); let worker = thread::Builder::new() .name("wasmeld-runtime-epoch".to_owned()) .spawn(move || { loop { match receiver.recv_timeout(tick) { Ok(()) | Err(mpsc::RecvTimeoutError::Disconnected) => return, Err(mpsc::RecvTimeoutError::Timeout) => engine.increment_epoch(), } } }) .map_err(RuntimeError::ActorThread)?; Ok(Arc::new(Self { shutdown, worker: Mutex::new(Some(worker)), })) } } impl Drop for EpochTicker { fn drop(&mut self) { let _ = self.shutdown.send(()); if let Ok(mut worker) = self.worker.lock() && let Some(handle) = worker.take() { let _ = handle.join(); } } } #[derive(Clone)] struct RegisteredService { manifest: ServiceManifest, component: Arc, capabilities: Vec, } impl Runtime { /// Creates the shared Wasmtime Engine and epoch interruption worker. pub fn new(config: RuntimeConfig) -> Result { config.validate()?; let mut wasmtime_config = Config::new(); wasmtime_config.wasm_component_model(true); wasmtime_config.async_support(false); wasmtime_config.consume_fuel(true); wasmtime_config.epoch_interruption(true); wasmtime_config.max_wasm_stack(config.max_wasm_stack); let engine = Arc::new(Engine::new(&wasmtime_config).map_err(RuntimeError::RuntimeCreation)?); let ticker = EpochTicker::start(Arc::clone(&engine), config.epoch_tick)?; Ok(Self { inner: Arc::new(RuntimeInner { engine, epoch_tick: config.epoch_tick, services: Mutex::new(HashMap::new()), actors: Mutex::new(HashMap::new()), lifecycle: Mutex::new(()), ticker, kv_backend: config.kv_backend, }), }) } /// Compiles and registers a Component without starting an Actor. /// /// Imports are checked against the restricted WASI and explicit Wasmeld /// capability allowlist before the service becomes visible. pub fn register( &self, manifest: ServiceManifest, component_bytes: impl AsRef<[u8]>, ) -> Result { manifest.validate()?; let key = manifest.key()?; let component = Component::new(&self.inner.engine, component_bytes.as_ref()) .map_err(RuntimeError::ComponentCompilation)?; let capabilities = validate_component_imports(&component, &self.inner.engine)?; if capabilities.contains(&Capability::KvStore) && self.inner.kv_backend.is_none() { return Err(RuntimeError::CapabilityUnavailable( Capability::KvStore.descriptor().interface().to_owned(), )); } let mut services = self .inner .services .lock() .map_err(|_| RuntimeError::LockPoisoned)?; if services.contains_key(&key) { return Err(RuntimeError::ServiceAlreadyRegistered(key)); } services.insert( key.clone(), RegisteredService { manifest, component: Arc::new(component), capabilities, }, ); Ok(key) } /// Reads a Component from the manifest path and delegates to [`Self::register`]. pub fn register_from_file( &self, manifest: ServiceManifest, ) -> Result { manifest.validate()?; let component_path = manifest.component.clone(); let bytes = fs::read(&component_path).map_err(|source| RuntimeError::ArtifactRead { path: component_path, source, })?; self.register(manifest, bytes) } /// Starts or returns the resident Actor for one registered service revision. pub fn start( &self, key: &ServiceKey, init_config: Vec, ) -> Result { let _lifecycle = self .inner .lifecycle .lock() .map_err(|_| RuntimeError::LockPoisoned)?; { let mut actors = self .inner .actors .lock() .map_err(|_| RuntimeError::LockPoisoned)?; if let Some(actor) = actors.get(key) { if actor.is_available() { return Ok(actor.clone()); } if actor.is_alive() { return Err(RuntimeError::ActorUnavailable(key.clone())); } } actors.remove(key); } let service = self .inner .services .lock() .map_err(|_| RuntimeError::LockPoisoned)? .get(key) .cloned() .ok_or_else(|| RuntimeError::ServiceNotRegistered(key.clone()))?; if service.manifest.limits.deadline() < self.inner.epoch_tick { return Err(RuntimeError::InvalidManifest(format!( "service deadline for {key} must be at least the runtime epoch_tick" ))); } let actor = spawn_actor( Arc::clone(&self.inner.engine), Arc::clone(&self.inner.ticker), self.inner.epoch_tick, key.clone(), service, init_config, self.inner.kv_backend.clone(), )?; self.inner .actors .lock() .map_err(|_| RuntimeError::LockPoisoned)? .insert(key.clone(), actor.clone()); Ok(actor) } /// Enqueues one call on the service Actor's bounded serial mailbox. pub fn invoke(&self, key: &ServiceKey, input: Vec) -> Result, RuntimeError> { let actor = self .inner .actors .lock() .map_err(|_| RuntimeError::LockPoisoned)? .get(key) .cloned() .ok_or_else(|| RuntimeError::ActorUnavailable(key.clone()))?; actor.invoke(input) } /// Returns the exact versioned host capabilities imported by a service. pub fn capabilities( &self, key: &ServiceKey, ) -> Result, RuntimeError> { let services = self .inner .services .lock() .map_err(|_| RuntimeError::LockPoisoned)?; let service = services .get(key) .ok_or_else(|| RuntimeError::ServiceNotRegistered(key.clone()))?; Ok(service .capabilities .iter() .map(|capability| capability.descriptor()) .collect()) } /// Stops one Actor and releases its Store and Component Instance. pub fn stop(&self, key: &ServiceKey) -> Result<(), RuntimeError> { let _lifecycle = self .inner .lifecycle .lock() .map_err(|_| RuntimeError::LockPoisoned)?; let actor = self .inner .actors .lock() .map_err(|_| RuntimeError::LockPoisoned)? .get(key) .cloned() .ok_or_else(|| RuntimeError::ActorUnavailable(key.clone()))?; actor.stop()?; self.inner .actors .lock() .map_err(|_| RuntimeError::LockPoisoned)? .remove(key); Ok(()) } /// Registers a new revision and immediately starts its Actor. pub fn reload( &self, manifest: ServiceManifest, component_bytes: impl AsRef<[u8]>, init_config: Vec, ) -> Result { let key = self.register(manifest, component_bytes)?; self.start(&key, init_config) } /// Stops and removes one registered service revision. /// /// Callers must switch any external deployment reference away from this /// revision before unregistering it. pub fn unregister(&self, key: &ServiceKey) -> Result<(), RuntimeError> { let _lifecycle = self .inner .lifecycle .lock() .map_err(|_| RuntimeError::LockPoisoned)?; if !self .inner .services .lock() .map_err(|_| RuntimeError::LockPoisoned)? .contains_key(key) { return Err(RuntimeError::ServiceNotRegistered(key.clone())); } let actor = self .inner .actors .lock() .map_err(|_| RuntimeError::LockPoisoned)? .get(key) .cloned(); if let Some(actor) = actor { if actor.is_available() { actor.stop()?; } else if actor.is_alive() { return Err(RuntimeError::ActorUnavailable(key.clone())); } self.inner .actors .lock() .map_err(|_| RuntimeError::LockPoisoned)? .remove(key); } self.inner .services .lock() .map_err(|_| RuntimeError::LockPoisoned)? .remove(key); Ok(()) } /// Returns counts without exposing internal Engine or Store state. pub fn stats(&self) -> Result { let registered_services = self .inner .services .lock() .map_err(|_| RuntimeError::LockPoisoned)? .len(); let running_actors = self .inner .actors .lock() .map_err(|_| RuntimeError::LockPoisoned)? .values() .filter(|actor| actor.is_available()) .count(); Ok(RuntimeStats { registered_services, running_actors, }) } /// Stops every Actor and returns the first stop failure, if any. pub fn stop_all(&self) -> Result<(), RuntimeError> { let _lifecycle = self .inner .lifecycle .lock() .map_err(|_| RuntimeError::LockPoisoned)?; let actors = self .inner .actors .lock() .map_err(|_| RuntimeError::LockPoisoned)? .values() .cloned() .collect::>(); let mut first_error = None; for actor in actors { if !actor.is_available() { continue; } if let Err(error) = actor.stop() && first_error.is_none() { first_error = Some(error); } } self.inner .actors .lock() .map_err(|_| RuntimeError::LockPoisoned)? .clear(); match first_error { Some(error) => Err(error), None => Ok(()), } } } /// Cloneable command handle for one resident Actor. /// /// Clones share the same bounded mailbox and status flags; they do not clone /// the Wasmtime Store or Component Instance. #[derive(Clone)] pub struct ActorHandle { key: ServiceKey, limits: Arc, sender: SyncSender, status: Arc, _ticker: Arc, } impl ActorHandle { /// Returns the service revision served by this Actor. pub fn key(&self) -> &ServiceKey { &self.key } /// Sends an invocation and waits up to the remaining service deadline. pub fn invoke(&self, input: Vec) -> Result, RuntimeError> { if !self.is_available() { return Err(RuntimeError::ActorUnavailable(self.key.clone())); } if input.len() > self.limits.max_input_bytes { return Err(RuntimeError::InputTooLarge { service: self.key.clone(), limit: self.limits.max_input_bytes, }); } let (response_sender, response_receiver) = mpsc::sync_channel(1); let deadline = Instant::now() + self.limits.deadline(); let command = ActorCommand::Invoke { input, deadline, response_sender, }; match self.sender.try_send(command) { Ok(()) => {} Err(TrySendError::Full(_)) => { return Err(RuntimeError::ActorOverloaded(self.key.clone())); } Err(TrySendError::Disconnected(_)) => { return Err(RuntimeError::ActorUnavailable(self.key.clone())); } } match response_receiver.recv_timeout(deadline.saturating_duration_since(Instant::now())) { Ok(result) => result, Err(mpsc::RecvTimeoutError::Timeout) => Err(RuntimeError::DeadlineExceeded { service: self.key.clone(), deadline: self.limits.deadline(), }), Err(mpsc::RecvTimeoutError::Disconnected) => { Err(RuntimeError::ActorStopped(self.key.clone())) } } } /// Prevents new calls, asks the worker to stop, and waits for acknowledgement. pub fn stop(&self) -> Result<(), RuntimeError> { if !self.status.accepting.swap(false, Ordering::AcqRel) { return Err(RuntimeError::ActorUnavailable(self.key.clone())); } let (response_sender, response_receiver) = mpsc::sync_channel(1); match self.sender.try_send(ActorCommand::Stop { response_sender }) { Ok(()) => {} Err(TrySendError::Full(_)) => { self.status.accepting.store(true, Ordering::Release); return Err(RuntimeError::ActorOverloaded(self.key.clone())); } Err(TrySendError::Disconnected(_)) => { self.status.alive.store(false, Ordering::Release); return Err(RuntimeError::ActorUnavailable(self.key.clone())); } } response_receiver .recv_timeout(self.limits.deadline()) .map_err(|_| RuntimeError::ActorStopped(self.key.clone())) } fn is_available(&self) -> bool { self.status.alive.load(Ordering::Acquire) && self.status.accepting.load(Ordering::Acquire) } fn is_alive(&self) -> bool { self.status.alive.load(Ordering::Acquire) } } enum ActorCommand { Invoke { input: Vec, deadline: Instant, response_sender: SyncSender, RuntimeError>>, }, Stop { response_sender: SyncSender<()>, }, } struct ActorStatus { accepting: AtomicBool, alive: AtomicBool, } impl ActorStatus { fn new() -> Self { Self { accepting: AtomicBool::new(false), alive: AtomicBool::new(false), } } fn mark_ready(&self) { self.alive.store(true, Ordering::Release); self.accepting.store(true, Ordering::Release); } fn mark_stopped(&self) { self.accepting.store(false, Ordering::Release); self.alive.store(false, Ordering::Release); } } pub(crate) struct HostState { store_limits: StoreLimits, wasi: WasiCtx, wasi_table: ResourceTable, pub(crate) capabilities: HostCapabilities, } impl HostState { fn new( limits: &ResourceLimits, capabilities: &[Capability], service_id: &str, kv_backend: Option>, ) -> Self { let store_limits = StoreLimitsBuilder::new() .memory_size(limits.memory_bytes) .table_elements(ACTOR_TABLE_ELEMENT_LIMIT) .instances(ACTOR_RESOURCE_COUNT_LIMIT) .tables(ACTOR_RESOURCE_COUNT_LIMIT) .memories(ACTOR_RESOURCE_COUNT_LIMIT) .trap_on_grow_failure(true) .build(); let mut wasi = WasiCtxBuilder::new(); wasi.allow_tcp(false) .allow_udp(false) .allow_ip_name_lookup(false) .socket_addr_check(|_, _| Box::pin(async { false })); Self { store_limits, wasi: wasi.build(), wasi_table: ResourceTable::new(), capabilities: HostCapabilities::new( capabilities, service_id, kv_backend, limits.deadline(), ), } } } impl HasData for HostState { type Data<'a> = &'a mut HostState; } impl WasiView for HostState { fn ctx(&mut self) -> WasiCtxView<'_> { WasiCtxView { ctx: &mut self.wasi, table: &mut self.wasi_table, } } } struct WasiIoTable; impl HasData for WasiIoTable { type Data<'a> = &'a mut ResourceTable; } struct ActorWorker { key: ServiceKey, limits: Arc, epoch_tick: Duration, store: Store, bindings: ServiceComponent, } impl ActorWorker { fn create( engine: Arc, epoch_tick: Duration, key: ServiceKey, service: RegisteredService, init_config: Vec, kv_backend: Option>, ) -> Result { let limits = Arc::new(service.manifest.limits.clone()); if init_config.len() > limits.max_input_bytes { return Err(RuntimeError::InputTooLarge { service: key, limit: limits.max_input_bytes, }); } let mut store = Store::new( &engine, HostState::new(&limits, &service.capabilities, key.id(), kv_backend), ); store.limiter(|state| &mut state.store_limits); let mut linker = Linker::new(&engine); add_restricted_wasi(&mut linker)?; CapabilityRegistry::add_to_linker(&mut linker, &service.capabilities)?; let epoch_ticks = ticks_for(limits.deadline(), epoch_tick); configure_call_budget(&mut store, &limits, epoch_ticks)?; let bindings = ServiceComponent::instantiate(&mut store, &service.component, &linker) .map_err(|error| RuntimeError::ActorInitialization(error.to_string()))?; configure_call_budget(&mut store, &limits, epoch_ticks)?; match bindings.call_init(&mut store, &init_config) { Ok(Ok(())) => {} Ok(Err(error)) => { return Err(component_error("init", error)); } Err(error) => { return Err(RuntimeError::ActorInitialization(error.to_string())); } } Ok(Self { key, limits, epoch_tick, store, bindings, }) } fn invoke(&mut self, input: Vec, deadline: Instant) -> Result, RuntimeError> { let remaining = deadline.saturating_duration_since(Instant::now()); if remaining.is_zero() { return Err(RuntimeError::DeadlineExceeded { service: self.key.clone(), deadline: self.limits.deadline(), }); } configure_call_budget( &mut self.store, &self.limits, ticks_for(remaining, self.epoch_tick), )?; match self.bindings.call_invoke(&mut self.store, &input) { Ok(Ok(output)) => { if output.len() > self.limits.max_output_bytes { return Err(RuntimeError::OutputTooLarge { service: self.key.clone(), limit: self.limits.max_output_bytes, }); } Ok(output) } Ok(Err(error)) => Err(component_error("invoke", error)), Err(error) => Err(RuntimeError::ActorFault { service: self.key.clone(), message: error.to_string(), }), } } } fn add_restricted_wasi(linker: &mut Linker) -> Result<(), RuntimeError> { let exit_options = cli::exit::LinkOptions::default(); link_wasi(clocks::wall_clock::add_to_linker::( linker, |state| state.clocks(), ))?; link_wasi(filesystem::preopens::add_to_linker::< HostState, WasiFilesystem, >(linker, |state| state.filesystem()))?; link_wasi( filesystem::types::add_to_linker::(linker, |state| { state.filesystem() }), )?; link_wasi(cli::exit::add_to_linker::( linker, &exit_options, |state| state.cli(), ))?; link_wasi(cli::environment::add_to_linker::( linker, |state| state.cli(), ))?; link_wasi(cli::stdin::add_to_linker::( linker, |state| state.cli(), ))?; link_wasi(cli::stdout::add_to_linker::( linker, |state| state.cli(), ))?; link_wasi(cli::stderr::add_to_linker::( linker, |state| state.cli(), ))?; link_wasi(io::error::add_to_linker::( linker, |state| state.ctx().table, ))?; link_wasi(io::streams::add_to_linker::( linker, |state| state.ctx().table, ))?; Ok(()) } fn link_wasi(result: wasmtime::Result<()>) -> Result<(), RuntimeError> { result.map_err(|error| RuntimeError::ActorInitialization(error.to_string())) } fn validate_component_imports( component: &Component, engine: &Engine, ) -> Result, RuntimeError> { // Import validation is deny-by-default. A host function is linked only // after its exact WIT identity or WASI family has passed this allowlist. let mut capabilities = Vec::new(); for (name, _) in component.component_type().imports(engine) { if let Some(capability) = CapabilityRegistry::resolve(name) { capabilities.push(capability); } else if !ALLOWED_WASI_IMPORTS .iter() .any(|allowed| name.starts_with(allowed)) { return Err(RuntimeError::UnsupportedImport(name.to_owned())); } } capabilities.sort_unstable(); capabilities.dedup(); Ok(capabilities) } fn spawn_actor( engine: Arc, ticker: Arc, epoch_tick: Duration, key: ServiceKey, service: RegisteredService, init_config: Vec, kv_backend: Option>, ) -> Result { let limits = Arc::new(service.manifest.limits.clone()); let (sender, receiver) = mpsc::sync_channel(limits.mailbox_capacity); let (ready_sender, ready_receiver) = mpsc::sync_channel(1); let thread_key = key.clone(); let status = Arc::new(ActorStatus::new()); let worker_status = Arc::clone(&status); let worker_ticker = Arc::clone(&ticker); thread::Builder::new() .name(format!("wasm-actor-{thread_key}")) .spawn(move || { let _ticker = worker_ticker; match ActorWorker::create(engine, epoch_tick, key, service, init_config, kv_backend) { Ok(mut worker) => { worker_status.mark_ready(); let _ = ready_sender.send(Ok(())); run_actor(&mut worker, receiver, &worker_status); } Err(error) => { worker_status.mark_stopped(); let _ = ready_sender.send(Err(error)); } } }) .map_err(RuntimeError::ActorThread)?; ready_receiver.recv().map_err(|_| { RuntimeError::ActorInitialization("actor thread stopped during initialization".to_owned()) })??; Ok(ActorHandle { key: thread_key, limits, sender, status, _ticker: ticker, }) } fn run_actor(worker: &mut ActorWorker, receiver: Receiver, status: &ActorStatus) { // This is the only loop that enters the Store. Multiple ActorHandle clones // therefore cannot invoke the same Component Instance concurrently. while let Ok(command) = receiver.recv() { match command { ActorCommand::Invoke { input, deadline, response_sender, } => { let result = if Instant::now() >= deadline { Err(RuntimeError::DeadlineExceeded { service: worker.key.clone(), deadline: worker.limits.deadline(), }) } else { worker.invoke(input, deadline) }; let fatal = is_fatal(&result); if fatal { status.mark_stopped(); } let _ = response_sender.send(result); if fatal { drain_after_failure(receiver, &worker.key); return; } } ActorCommand::Stop { response_sender } => { status.mark_stopped(); let _ = response_sender.send(()); return; } } } status.mark_stopped(); } fn drain_after_failure(receiver: Receiver, key: &ServiceKey) { while let Ok(command) = receiver.try_recv() { match command { ActorCommand::Invoke { response_sender, .. } => { let _ = response_sender.send(Err(RuntimeError::ActorUnavailable(key.clone()))); } ActorCommand::Stop { response_sender } => { let _ = response_sender.send(()); } } } } fn configure_call_budget( store: &mut Store, limits: &ResourceLimits, epoch_ticks: u64, ) -> Result<(), RuntimeError> { store .set_fuel(limits.fuel_per_call) .map_err(|error| RuntimeError::ActorInitialization(error.to_string()))?; store.set_epoch_deadline(epoch_ticks); Ok(()) } fn ticks_for(deadline: Duration, epoch_tick: Duration) -> u64 { let deadline_nanos = deadline.as_nanos(); let tick_nanos = epoch_tick.as_nanos().max(1); let ticks = deadline_nanos.div_ceil(tick_nanos).max(1); u64::try_from(ticks).unwrap_or(u64::MAX) } fn component_error(kind: &'static str, error: ServiceError) -> RuntimeError { let message = match error { ServiceError::InitFailed(message) | ServiceError::CallFailed(message) => message, }; RuntimeError::ComponentError { kind, message } } fn is_fatal(result: &Result, RuntimeError>) -> bool { matches!(result, Err(RuntimeError::ActorFault { .. })) }