feat(runtime): supervise resident component resources

Detect the resident actor export, route Host events through the bounded Actor mailbox, enforce fuel, deadlines, payload limits, and translate WIT effects into typed runtime values.

Add a revision-scoped ResidentSession with deny-by-default TCP, UDP, and Unix endpoint policies, non-reused resource IDs, stream lifecycle and backpressure state, atomic effect validation, timers, subscriptions, extension sources, and shutdown handling.

Cover passive rejection, real Component ABI round trips, resource ownership, endpoint policy, and driver operation validation with integration tests.
This commit is contained in:
Maofeng
2026-07-30 07:40:08 +08:00
parent 949b8d6cdb
commit 572fed47b4
8 changed files with 1983 additions and 7 deletions
+8
View File
@@ -5,3 +5,11 @@ wasmtime::component::bindgen!({
path: "../../wit/service",
world: "service-component",
});
/// Optional exports implemented only by active resident Components.
pub(crate) mod resident {
wasmtime::component::bindgen!({
path: "../../wit/resident",
world: "resident-component",
});
}
+12
View File
@@ -68,6 +68,18 @@ pub enum RuntimeError {
#[error("actor for {0} stopped before returning a result")]
ActorStopped(ServiceKey),
#[error("service revision {0} does not export the resident actor interface")]
NotResident(ServiceKey),
#[error("resident event for {service} exceeds the {limit}-byte payload limit")]
ResidentEventTooLarge { service: ServiceKey, limit: usize },
#[error("resident event for {service} returned more than {limit} effects")]
TooManyResidentEffects { service: ServiceKey, limit: usize },
#[error("invalid resident effect: {0}")]
InvalidResidentEffect(String),
#[error("input for {service} exceeds the {limit}-byte limit")]
InputTooLarge { service: ServiceKey, limit: usize },
+11
View File
@@ -9,6 +9,8 @@ mod bindings;
mod capability;
mod error;
mod manifest;
mod resident;
mod resident_host;
mod runtime;
pub use capability::{
@@ -16,5 +18,14 @@ pub use capability::{
};
pub use error::RuntimeError;
pub use manifest::{ResourceLimits, ServiceKey, ServiceManifest};
pub use resident::{
ComponentExecution, ResidentEffect, ResidentEvent, ResidentLimits, ResourceId,
StreamCloseReason,
};
pub use resident_host::{
NetworkScope, ResidentEndpoint, ResidentHostError, ResidentOperation, ResidentPolicy,
ResidentResourceInfo, ResidentResourceKind, ResidentResourceMetadata, ResidentResourceState,
ResidentSession,
};
pub use runtime::{ActorHandle, Runtime, RuntimeConfig, RuntimeStats};
pub use wasmeld_package::SERVICE_WORLD;
+6 -2
View File
@@ -4,7 +4,7 @@ use std::{fmt, time::Duration};
use serde::{Deserialize, Serialize};
use crate::error::RuntimeError;
use crate::{error::RuntimeError, resident::ResidentLimits};
/// Immutable identity of one service revision.
#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
@@ -103,6 +103,9 @@ pub struct ResourceLimits {
pub max_input_bytes: usize,
/// Maximum response payload size.
pub max_output_bytes: usize,
/// Limits for active events and Host-owned resources.
#[serde(default)]
pub resident: ResidentLimits,
}
impl Default for ResourceLimits {
@@ -114,6 +117,7 @@ impl Default for ResourceLimits {
mailbox_capacity: 16,
max_input_bytes: 16 * 1024,
max_output_bytes: 16 * 1024,
resident: ResidentLimits::default(),
}
}
}
@@ -151,7 +155,7 @@ impl ResourceLimits {
));
}
Ok(())
self.resident.validate()
}
/// Returns the configured call deadline.
+317
View File
@@ -0,0 +1,317 @@
//! Host-owned event and effect protocol for active resident Components.
//!
//! A [`ResidentEvent`] is always produced by a Host driver and delivered
//! through the same bounded Actor mailbox as ordinary service invocations.
//! The Component never receives an operating-system descriptor. It can only
//! reference opaque [`ResourceId`] values and return [`ResidentEffect`] values
//! for the owning driver to validate and apply.
use serde::{Deserialize, Serialize};
use crate::{RuntimeError, bindings::resident::exports::wasmeld::resident::actor as wit};
pub(crate) const RESIDENT_INTERFACE: &str = "wasmeld:resident/actor@0.1.0";
/// Component execution model discovered from its exported WIT interfaces.
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum ComponentExecution {
/// Retains memory but runs only for `invoke` mailbox commands.
Service,
/// Also accepts Host events and returns effects for Host-owned resources.
Resident,
}
/// Opaque Host resource identity scoped to one resident Actor revision.
///
/// IDs are never reused while the Actor is alive. `0` is reserved so an
/// uninitialized ID cannot accidentally address a real resource.
#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
#[serde(transparent)]
pub struct ResourceId(u64);
impl ResourceId {
/// Creates a valid opaque resource identity.
pub fn new(value: u64) -> Option<Self> {
(value != 0).then_some(Self(value))
}
/// Returns the WIT-compatible integer representation.
pub const fn get(self) -> u64 {
self.0
}
}
/// Why a Host-owned stream stopped accepting events and effects.
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum StreamCloseReason {
PeerClosed,
HostClosed,
IdleTimeout,
ProtocolError,
TransportError,
}
/// One bounded event delivered to a resident Component.
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum ResidentEvent {
Timer {
timer_id: ResourceId,
scheduled_at_ns: u64,
},
StreamOpened {
endpoint_id: ResourceId,
stream_id: ResourceId,
peer: Option<String>,
},
StreamData {
stream_id: ResourceId,
bytes: Vec<u8>,
},
StreamWritable {
stream_id: ResourceId,
},
StreamHalfClosed {
stream_id: ResourceId,
},
StreamClosed {
stream_id: ResourceId,
reason: StreamCloseReason,
},
Datagram {
endpoint_id: ResourceId,
peer: String,
bytes: Vec<u8>,
},
Message {
subscription_id: ResourceId,
message_id: u64,
bytes: Vec<u8>,
},
Source {
source_id: ResourceId,
event: String,
payload: Vec<u8>,
},
Shutdown,
}
/// A requested operation returned by a resident Component.
///
/// Effects are inert until the Host supervisor checks resource ownership,
/// payload limits, current stream state, and driver policy.
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum ResidentEffect {
WriteStream {
stream_id: ResourceId,
bytes: Vec<u8>,
},
CloseStream {
stream_id: ResourceId,
},
PauseStream {
stream_id: ResourceId,
},
ResumeStream {
stream_id: ResourceId,
},
SendDatagram {
endpoint_id: ResourceId,
peer: String,
bytes: Vec<u8>,
},
ArmTimer {
timer_id: ResourceId,
delay_ms: u64,
interval_ms: Option<u64>,
},
CancelTimer {
timer_id: ResourceId,
},
AcknowledgeMessage {
subscription_id: ResourceId,
message_id: u64,
},
SourceCommand {
source_id: ResourceId,
command: String,
payload: Vec<u8>,
},
}
/// Execution limits specific to active resident Components.
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct ResidentLimits {
/// Maximum number of Host resources owned by one Actor revision.
pub max_resources: usize,
/// Maximum effects accepted from one event handler call.
pub max_effects_per_event: usize,
/// Maximum bytes in one stream read or write.
pub max_stream_chunk_bytes: usize,
/// Maximum bytes in one received or sent datagram.
pub max_datagram_bytes: usize,
/// Minimum delay accepted for Component-created timers.
pub min_timer_delay_ms: u64,
}
impl Default for ResidentLimits {
fn default() -> Self {
Self {
max_resources: 1024,
max_effects_per_event: 64,
max_stream_chunk_bytes: 64 * 1024,
max_datagram_bytes: 64 * 1024,
min_timer_delay_ms: 1,
}
}
}
impl ResidentLimits {
pub(crate) fn validate(&self) -> Result<(), RuntimeError> {
if self.max_resources == 0 {
return Err(invalid_limit("max_resources"));
}
if self.max_effects_per_event == 0 {
return Err(invalid_limit("max_effects_per_event"));
}
if self.max_stream_chunk_bytes == 0 {
return Err(invalid_limit("max_stream_chunk_bytes"));
}
if self.max_datagram_bytes == 0 {
return Err(invalid_limit("max_datagram_bytes"));
}
if self.min_timer_delay_ms == 0 {
return Err(invalid_limit("min_timer_delay_ms"));
}
Ok(())
}
}
pub(crate) fn into_wit_event(event: ResidentEvent) -> wit::Event {
match event {
ResidentEvent::Timer {
timer_id,
scheduled_at_ns,
} => wit::Event::Timer(wit::TimerFired {
timer_id: timer_id.get(),
scheduled_at_ns,
}),
ResidentEvent::StreamOpened {
endpoint_id,
stream_id,
peer,
} => wit::Event::StreamOpened(wit::StreamOpened {
endpoint_id: endpoint_id.get(),
stream_id: stream_id.get(),
peer,
}),
ResidentEvent::StreamData { stream_id, bytes } => {
wit::Event::StreamData(wit::StreamChunk {
stream_id: stream_id.get(),
bytes,
})
}
ResidentEvent::StreamWritable { stream_id } => wit::Event::StreamWritable(stream_id.get()),
ResidentEvent::StreamHalfClosed { stream_id } => {
wit::Event::StreamHalfClosed(stream_id.get())
}
ResidentEvent::StreamClosed { stream_id, reason } => {
wit::Event::StreamClosed(wit::StreamClosed {
stream_id: stream_id.get(),
reason: match reason {
StreamCloseReason::PeerClosed => wit::CloseReason::PeerClosed,
StreamCloseReason::HostClosed => wit::CloseReason::HostClosed,
StreamCloseReason::IdleTimeout => wit::CloseReason::IdleTimeout,
StreamCloseReason::ProtocolError => wit::CloseReason::ProtocolError,
StreamCloseReason::TransportError => wit::CloseReason::TransportError,
},
})
}
ResidentEvent::Datagram {
endpoint_id,
peer,
bytes,
} => wit::Event::Datagram(wit::Datagram {
endpoint_id: endpoint_id.get(),
peer,
bytes,
}),
ResidentEvent::Message {
subscription_id,
message_id,
bytes,
} => wit::Event::Message(wit::Message {
subscription_id: subscription_id.get(),
message_id,
bytes,
}),
ResidentEvent::Source {
source_id,
event,
payload,
} => wit::Event::Source(wit::SourceEvent {
source_id: source_id.get(),
kind: event,
payload,
}),
ResidentEvent::Shutdown => wit::Event::Shutdown,
}
}
pub(crate) fn from_wit_effects(
effects: Vec<wit::Effect>,
) -> Result<Vec<ResidentEffect>, RuntimeError> {
effects.into_iter().map(from_wit_effect).collect()
}
fn from_wit_effect(effect: wit::Effect) -> Result<ResidentEffect, RuntimeError> {
Ok(match effect {
wit::Effect::WriteStream(effect) => ResidentEffect::WriteStream {
stream_id: resource_id(effect.stream_id)?,
bytes: effect.bytes,
},
wit::Effect::CloseStream(stream_id) => ResidentEffect::CloseStream {
stream_id: resource_id(stream_id)?,
},
wit::Effect::PauseStream(stream_id) => ResidentEffect::PauseStream {
stream_id: resource_id(stream_id)?,
},
wit::Effect::ResumeStream(stream_id) => ResidentEffect::ResumeStream {
stream_id: resource_id(stream_id)?,
},
wit::Effect::SendDatagram(effect) => ResidentEffect::SendDatagram {
endpoint_id: resource_id(effect.endpoint_id)?,
peer: effect.peer,
bytes: effect.bytes,
},
wit::Effect::ArmTimer(effect) => ResidentEffect::ArmTimer {
timer_id: resource_id(effect.timer_id)?,
delay_ms: effect.delay_ms,
interval_ms: effect.interval_ms,
},
wit::Effect::CancelTimer(timer_id) => ResidentEffect::CancelTimer {
timer_id: resource_id(timer_id)?,
},
wit::Effect::AcknowledgeMessage(effect) => ResidentEffect::AcknowledgeMessage {
subscription_id: resource_id(effect.subscription_id)?,
message_id: effect.message_id,
},
wit::Effect::SourceCommand(effect) => ResidentEffect::SourceCommand {
source_id: resource_id(effect.source_id)?,
command: effect.command,
payload: effect.payload,
},
})
}
fn resource_id(value: u64) -> Result<ResourceId, RuntimeError> {
ResourceId::new(value)
.ok_or_else(|| RuntimeError::InvalidResidentEffect("resource ID 0 is reserved".to_owned()))
}
fn invalid_limit(field: &str) -> RuntimeError {
RuntimeError::InvalidManifest(format!("limits.resident.{field} must be greater than zero"))
}
File diff suppressed because it is too large Load Diff
+306 -3
View File
@@ -31,10 +31,15 @@ use wasmtime_wasi::{
};
use crate::{
CapabilityDescriptor, KvBackend, RuntimeError, ServiceKey, ServiceManifest,
bindings::{ServiceComponent, ServiceError},
CapabilityDescriptor, ComponentExecution, KvBackend, ResidentEffect, ResidentEvent,
RuntimeError, ServiceKey, ServiceManifest,
bindings::{
ServiceComponent, ServiceError,
resident::{ResidentComponent, exports::wasmeld::resident::actor::ResidentError},
},
capability::{Capability, CapabilityRegistry, HostCapabilities},
manifest::ResourceLimits,
resident::{RESIDENT_INTERFACE, from_wit_effects, into_wit_event},
};
const DEFAULT_EPOCH_TICK: Duration = Duration::from_millis(5);
@@ -171,6 +176,7 @@ struct RegisteredService {
manifest: ServiceManifest,
component: Arc<Component>,
capabilities: Vec<Capability>,
execution: ComponentExecution,
}
impl Runtime {
@@ -231,6 +237,7 @@ impl Runtime {
let component = Component::new(&self.inner.engine, component_bytes.as_ref())
.map_err(RuntimeError::ComponentCompilation)?;
let capabilities = validate_component_imports(&component, &self.inner.engine)?;
let execution = component_execution(&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(),
@@ -253,6 +260,7 @@ impl Runtime {
manifest,
component: Arc::new(component),
capabilities,
execution,
},
);
@@ -372,6 +380,34 @@ impl Runtime {
.collect())
}
/// Returns whether a registered Component is passive or event-driven.
pub fn execution(&self, key: &ServiceKey) -> Result<ComponentExecution, RuntimeError> {
self.inner
.services
.lock()
.map_err(|_| RuntimeError::LockPoisoned)?
.get(key)
.map(|service| service.execution)
.ok_or_else(|| RuntimeError::ServiceNotRegistered(key.clone()))
}
/// Delivers one Host event through the Actor's bounded serial mailbox.
pub fn dispatch_event(
&self,
key: &ServiceKey,
event: ResidentEvent,
) -> Result<Vec<ResidentEffect>, RuntimeError> {
let actor = self
.inner
.actors
.lock()
.map_err(|_| RuntimeError::LockPoisoned)?
.get(key)
.cloned()
.ok_or_else(|| RuntimeError::ActorUnavailable(key.clone()))?;
actor.dispatch_event(event)
}
/// Stops one Actor and releases its Store and Component Instance.
pub fn stop(&self, key: &ServiceKey) -> Result<(), RuntimeError> {
let _lifecycle = self
@@ -528,6 +564,7 @@ impl Runtime {
#[derive(Clone)]
pub struct ActorHandle {
key: ServiceKey,
execution: ComponentExecution,
limits: Arc<ResourceLimits>,
sender: SyncSender<ActorCommand>,
status: Arc<ActorStatus>,
@@ -540,6 +577,15 @@ impl ActorHandle {
&self.key
}
/// Returns the execution model detected from the Component exports.
pub fn execution(&self) -> ComponentExecution {
self.execution
}
pub(crate) fn resident_resource_limit(&self) -> usize {
self.limits.resident.max_resources
}
/// Sends an invocation and waits up to the remaining service deadline.
pub fn invoke(&self, input: Vec<u8>) -> Result<Vec<u8>, RuntimeError> {
if !self.is_available() {
@@ -583,6 +629,45 @@ impl ActorHandle {
}
}
/// Sends a Host-owned resource event and waits for validated effects.
pub fn dispatch_event(
&self,
event: ResidentEvent,
) -> Result<Vec<ResidentEffect>, RuntimeError> {
if !self.is_available() {
return Err(RuntimeError::ActorUnavailable(self.key.clone()));
}
validate_resident_event(&self.key, &self.limits, &event)?;
let (response_sender, response_receiver) = mpsc::sync_channel(1);
let deadline = Instant::now() + self.limits.deadline();
let command = ActorCommand::ResidentEvent {
event,
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) {
@@ -622,6 +707,11 @@ enum ActorCommand {
deadline: Instant,
response_sender: SyncSender<Result<Vec<u8>, RuntimeError>>,
},
ResidentEvent {
event: ResidentEvent,
deadline: Instant,
response_sender: SyncSender<Result<Vec<ResidentEffect>, RuntimeError>>,
},
Stop {
response_sender: SyncSender<()>,
},
@@ -719,6 +809,7 @@ struct ActorWorker {
epoch_tick: Duration,
store: Store<HostState>,
bindings: ServiceComponent,
resident: Option<ResidentComponent>,
}
impl ActorWorker {
@@ -750,8 +841,18 @@ impl ActorWorker {
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)
let instance = linker
.instantiate(&mut store, &service.component)
.map_err(|error| RuntimeError::ActorInitialization(error.to_string()))?;
let bindings = ServiceComponent::new(&mut store, &instance)
.map_err(|error| RuntimeError::ActorInitialization(error.to_string()))?;
let resident = match service.execution {
ComponentExecution::Service => None,
ComponentExecution::Resident => Some(
ResidentComponent::new(&mut store, &instance)
.map_err(|error| RuntimeError::ActorInitialization(error.to_string()))?,
),
};
configure_call_budget(&mut store, &limits, epoch_ticks)?;
@@ -771,6 +872,7 @@ impl ActorWorker {
epoch_tick,
store,
bindings,
resident,
})
}
@@ -807,6 +909,54 @@ impl ActorWorker {
}),
}
}
fn dispatch_event(
&mut self,
event: ResidentEvent,
deadline: Instant,
) -> Result<Vec<ResidentEffect>, 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(),
});
}
let resident = self
.resident
.as_ref()
.ok_or_else(|| RuntimeError::NotResident(self.key.clone()))?;
configure_call_budget(
&mut self.store,
&self.limits,
ticks_for(remaining, self.epoch_tick),
)?;
match resident
.wasmeld_resident_actor()
.call_handle_event(&mut self.store, &into_wit_event(event))
{
Ok(Ok(effects)) => {
if effects.len() > self.limits.resident.max_effects_per_event {
return Err(RuntimeError::TooManyResidentEffects {
service: self.key.clone(),
limit: self.limits.resident.max_effects_per_event,
});
}
let effects = from_wit_effects(effects)?;
validate_resident_effects(&self.key, &self.limits, &effects)?;
Ok(effects)
}
Ok(Err(ResidentError::EventFailed(message))) => Err(RuntimeError::ComponentError {
kind: "resident-event",
message,
}),
Err(error) => Err(RuntimeError::ActorFault {
service: self.key.clone(),
message: error.to_string(),
}),
}
}
}
fn add_restricted_wasi(linker: &mut Linker<HostState>) -> Result<(), RuntimeError> {
@@ -883,6 +1033,18 @@ fn validate_component_imports(
Ok(capabilities)
}
fn component_execution(component: &Component, engine: &Engine) -> ComponentExecution {
if component
.component_type()
.exports(engine)
.any(|(name, _)| name == RESIDENT_INTERFACE)
{
ComponentExecution::Resident
} else {
ComponentExecution::Service
}
}
fn spawn_actor(
engine: Arc<Engine>,
ticker: Arc<EpochTicker>,
@@ -892,6 +1054,7 @@ fn spawn_actor(
init_config: Vec<u8>,
kv_backend: Option<Arc<dyn KvBackend>>,
) -> Result<ActorHandle, RuntimeError> {
let execution = service.execution;
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);
@@ -925,6 +1088,7 @@ fn spawn_actor(
Ok(ActorHandle {
key: thread_key,
execution,
limits,
sender,
status,
@@ -961,6 +1125,29 @@ fn run_actor(worker: &mut ActorWorker, receiver: Receiver<ActorCommand>, status:
return;
}
}
ActorCommand::ResidentEvent {
event,
deadline,
response_sender,
} => {
let result = if Instant::now() >= deadline {
Err(RuntimeError::DeadlineExceeded {
service: worker.key.clone(),
deadline: worker.limits.deadline(),
})
} else {
worker.dispatch_event(event, deadline)
};
let fatal = is_fatal_resident(&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(());
@@ -980,6 +1167,11 @@ fn drain_after_failure(receiver: Receiver<ActorCommand>, key: &ServiceKey) {
} => {
let _ = response_sender.send(Err(RuntimeError::ActorUnavailable(key.clone())));
}
ActorCommand::ResidentEvent {
response_sender, ..
} => {
let _ = response_sender.send(Err(RuntimeError::ActorUnavailable(key.clone())));
}
ActorCommand::Stop { response_sender } => {
let _ = response_sender.send(());
}
@@ -987,6 +1179,113 @@ fn drain_after_failure(receiver: Receiver<ActorCommand>, key: &ServiceKey) {
}
}
fn validate_resident_event(
key: &ServiceKey,
limits: &ResourceLimits,
event: &ResidentEvent,
) -> Result<(), RuntimeError> {
let (payload_len, limit) = match event {
ResidentEvent::StreamData { bytes, .. } => {
(bytes.len(), limits.resident.max_stream_chunk_bytes)
}
ResidentEvent::Datagram { bytes, .. } => (bytes.len(), limits.resident.max_datagram_bytes),
ResidentEvent::Message { bytes, .. } => (bytes.len(), limits.max_input_bytes),
ResidentEvent::Source { event, payload, .. } => (
event.len().saturating_add(payload.len()),
limits.max_input_bytes,
),
ResidentEvent::Timer { .. }
| ResidentEvent::StreamOpened { .. }
| ResidentEvent::StreamWritable { .. }
| ResidentEvent::StreamHalfClosed { .. }
| ResidentEvent::StreamClosed { .. }
| ResidentEvent::Shutdown => return Ok(()),
};
if payload_len > limit {
Err(RuntimeError::ResidentEventTooLarge {
service: key.clone(),
limit,
})
} else {
Ok(())
}
}
fn validate_resident_effects(
key: &ServiceKey,
limits: &ResourceLimits,
effects: &[ResidentEffect],
) -> Result<(), RuntimeError> {
for effect in effects {
match effect {
ResidentEffect::WriteStream { bytes, .. }
if bytes.len() > limits.resident.max_stream_chunk_bytes =>
{
return Err(invalid_resident_effect(
key,
"stream write",
limits.resident.max_stream_chunk_bytes,
));
}
ResidentEffect::SendDatagram { peer, bytes, .. } => {
if peer.is_empty() {
return Err(RuntimeError::InvalidResidentEffect(format!(
"{key} returned a datagram without a peer"
)));
}
if bytes.len() > limits.resident.max_datagram_bytes {
return Err(invalid_resident_effect(
key,
"datagram",
limits.resident.max_datagram_bytes,
));
}
}
ResidentEffect::ArmTimer {
delay_ms,
interval_ms,
..
} => {
let minimum = limits.resident.min_timer_delay_ms;
if *delay_ms < minimum || interval_ms.is_some_and(|value| value < minimum) {
return Err(RuntimeError::InvalidResidentEffect(format!(
"{key} returned a timer below the {minimum} ms minimum"
)));
}
}
ResidentEffect::SourceCommand {
command, payload, ..
} => {
if command.is_empty() {
return Err(RuntimeError::InvalidResidentEffect(format!(
"{key} returned an empty source command"
)));
}
if command.len().saturating_add(payload.len()) > limits.max_output_bytes {
return Err(invalid_resident_effect(
key,
"source command",
limits.max_output_bytes,
));
}
}
ResidentEffect::WriteStream { .. }
| ResidentEffect::CloseStream { .. }
| ResidentEffect::PauseStream { .. }
| ResidentEffect::ResumeStream { .. }
| ResidentEffect::CancelTimer { .. }
| ResidentEffect::AcknowledgeMessage { .. } => {}
}
}
Ok(())
}
fn invalid_resident_effect(key: &ServiceKey, kind: &str, limit: usize) -> RuntimeError {
RuntimeError::InvalidResidentEffect(format!(
"{key} returned a {kind} payload above the {limit}-byte limit"
))
}
fn configure_call_budget(
store: &mut Store<HostState>,
limits: &ResourceLimits,
@@ -1017,3 +1316,7 @@ fn component_error(kind: &'static str, error: ServiceError) -> RuntimeError {
fn is_fatal(result: &Result<Vec<u8>, RuntimeError>) -> bool {
matches!(result, Err(RuntimeError::ActorFault { .. }))
}
fn is_fatal_resident(result: &Result<Vec<ResidentEffect>, RuntimeError>) -> bool {
matches!(result, Err(RuntimeError::ActorFault { .. }))
}
@@ -1,6 +1,7 @@
use std::{
collections::HashMap,
fs,
net::SocketAddr,
path::{Path, PathBuf},
process::Command,
sync::{Arc, Mutex, OnceLock},
@@ -9,8 +10,11 @@ use std::{
};
use wasmeld_runtime::{
KV_MAX_VALUE_BYTES, KvBackend, KvBackendError, ResourceLimits, Runtime, RuntimeConfig,
RuntimeError, ServiceManifest,
ComponentExecution, KV_MAX_VALUE_BYTES, KvBackend, KvBackendError, NetworkScope,
ResidentEffect, ResidentEndpoint, ResidentEvent, ResidentHostError, ResidentLimits,
ResidentOperation, ResidentPolicy, ResidentResourceKind, ResidentResourceMetadata,
ResidentSession, ResourceId, ResourceLimits, Runtime, RuntimeConfig, RuntimeError,
ServiceManifest, StreamCloseReason,
};
static COMPONENTS_BUILT: OnceLock<()> = OnceLock::new();
@@ -192,6 +196,253 @@ fn explicit_clock_capability_is_linked() {
runtime.stop(&key).unwrap();
}
#[test]
fn resident_component_handles_host_events_in_its_actor_mailbox() {
let runtime = Runtime::new(RuntimeConfig::default()).expect("runtime should start");
let (key, actor) = start_component(&runtime, "resident-probe", "resident_probe_component.wasm");
let stream_id = resource_id(1);
let endpoint_id = resource_id(2);
let timer_id = resource_id(3);
let subscription_id = resource_id(4);
let source_id = resource_id(5);
assert_eq!(
runtime.execution(&key).unwrap(),
ComponentExecution::Resident
);
assert_eq!(
actor
.dispatch_event(ResidentEvent::StreamData {
stream_id,
bytes: b"stream".to_vec(),
})
.unwrap(),
vec![ResidentEffect::WriteStream {
stream_id,
bytes: b"stream".to_vec(),
}]
);
assert_eq!(
runtime
.dispatch_event(
&key,
ResidentEvent::Datagram {
endpoint_id,
peer: "127.0.0.1:9000".to_owned(),
bytes: b"datagram".to_vec(),
},
)
.unwrap(),
vec![ResidentEffect::SendDatagram {
endpoint_id,
peer: "127.0.0.1:9000".to_owned(),
bytes: b"datagram".to_vec(),
}]
);
assert_eq!(
actor
.dispatch_event(ResidentEvent::Timer {
timer_id,
scheduled_at_ns: 42,
})
.unwrap(),
vec![ResidentEffect::ArmTimer {
timer_id,
delay_ms: 10,
interval_ms: None,
}]
);
assert_eq!(
actor
.dispatch_event(ResidentEvent::Message {
subscription_id,
message_id: 99,
bytes: b"message".to_vec(),
})
.unwrap(),
vec![ResidentEffect::AcknowledgeMessage {
subscription_id,
message_id: 99,
}]
);
assert_eq!(
actor
.dispatch_event(ResidentEvent::Source {
source_id,
event: "flush".to_owned(),
payload: b"source".to_vec(),
})
.unwrap(),
vec![ResidentEffect::SourceCommand {
source_id,
command: "flush".to_owned(),
payload: b"source".to_vec(),
}]
);
assert_eq!(decode_counter(actor.invoke(Vec::new()).unwrap()), 5);
runtime.stop(&key).unwrap();
}
#[test]
fn passive_component_rejects_resident_events_without_stopping() {
let runtime = Runtime::new(RuntimeConfig::default()).expect("runtime should start");
let (key, actor) = start_component(&runtime, "passive-echo", "echo_component.wasm");
assert_eq!(
runtime.execution(&key).unwrap(),
ComponentExecution::Service
);
assert!(matches!(
actor.dispatch_event(ResidentEvent::Shutdown),
Err(RuntimeError::NotResident(service)) if service == key
));
assert_eq!(
actor.invoke(b"still-running".to_vec()).unwrap(),
b"still-running"
);
runtime.stop(&key).unwrap();
}
#[test]
fn resident_session_owns_resources_and_validates_driver_operations() {
let runtime = Runtime::new(RuntimeConfig::default()).expect("runtime should start");
let (key, actor) = start_component(
&runtime,
"resident-session",
"resident_probe_component.wasm",
);
let policy = ResidentPolicy {
tcp_listen: NetworkScope::Loopback,
udp_bind: NetworkScope::Loopback,
unix_listen_roots: vec![std::env::temp_dir()],
};
let mut session = ResidentSession::new(actor, policy).unwrap();
let tcp_id = session
.register_endpoint(ResidentEndpoint::Tcp {
name: "ingress".to_owned(),
bind: socket_addr("127.0.0.1:7000"),
})
.unwrap();
let udp_id = session
.register_endpoint(ResidentEndpoint::Udp {
name: "discovery".to_owned(),
bind: socket_addr("127.0.0.1:7001"),
})
.unwrap();
let unix_id = session
.register_endpoint(ResidentEndpoint::Unix {
name: "local-ingress".to_owned(),
path: std::env::temp_dir().join("wasmeld-resident-test.sock"),
})
.unwrap();
let timer_id = session.register_timer("heartbeat").unwrap();
let subscription_id = session.register_message_subscription("jobs").unwrap();
let source_id = session
.register_source("watcher", "fs-watch@0.1.0")
.unwrap();
let (stream_id, opened) = session
.accept_stream(tcp_id, Some("127.0.0.1:51000".to_owned()))
.unwrap();
assert!(opened.is_empty());
assert_eq!(
session.stream_data(stream_id, b"stream".to_vec()).unwrap(),
vec![ResidentOperation::WriteStream {
stream_id,
bytes: b"stream".to_vec(),
}]
);
assert_eq!(
session
.datagram(udp_id, "127.0.0.1:51001".to_owned(), b"datagram".to_vec(),)
.unwrap(),
vec![ResidentOperation::SendDatagram {
endpoint_id: udp_id,
peer: "127.0.0.1:51001".to_owned(),
bytes: b"datagram".to_vec(),
}]
);
assert_eq!(
session.timer_fired(timer_id, 1).unwrap(),
vec![ResidentOperation::ArmTimer {
timer_id,
delay_ms: 10,
interval_ms: None,
}]
);
assert!(matches!(
session.resource_info(timer_id).unwrap().metadata,
ResidentResourceMetadata::Timer { armed: true }
));
assert_eq!(
session
.message(subscription_id, 7, b"work".to_vec())
.unwrap(),
vec![ResidentOperation::AcknowledgeMessage {
subscription_id,
message_id: 7,
}]
);
assert_eq!(
session
.source_event(source_id, "flush".to_owned(), b"state".to_vec())
.unwrap(),
vec![ResidentOperation::SourceCommand {
source_id,
command: "flush".to_owned(),
payload: b"state".to_vec(),
}]
);
assert!(matches!(
session.datagram(tcp_id, "127.0.0.1:1".to_owned(), Vec::new()),
Err(ResidentHostError::WrongResourceKind { .. })
));
session
.stream_closed(stream_id, StreamCloseReason::PeerClosed)
.unwrap();
assert!(matches!(
session.resource_info(stream_id),
Err(ResidentHostError::UnknownResource { .. })
));
assert_eq!(session.owner(), &key);
assert_eq!(
session.resource_info(tcp_id).unwrap().kind,
ResidentResourceKind::TcpEndpoint
);
assert_eq!(
session.resource_info(unix_id).unwrap().kind,
ResidentResourceKind::UnixEndpoint
);
session.shutdown().unwrap();
assert!(matches!(
session.register_timer("late"),
Err(ResidentHostError::ShuttingDown(_))
));
runtime.stop(&key).unwrap();
}
#[test]
fn resident_endpoint_policy_is_deny_by_default() {
let runtime = Runtime::new(RuntimeConfig::default()).expect("runtime should start");
let (key, actor) =
start_component(&runtime, "resident-policy", "resident_probe_component.wasm");
let mut session = ResidentSession::new(actor, ResidentPolicy::default()).unwrap();
assert!(matches!(
session.register_endpoint(ResidentEndpoint::Tcp {
name: "public".to_owned(),
bind: socket_addr("0.0.0.0:8080"),
}),
Err(ResidentHostError::EndpointDenied(_))
));
runtime.stop(&key).unwrap();
}
#[test]
fn kv_capability_is_service_scoped_and_shared_across_revisions() {
let backend = Arc::new(MemoryKvBackend::default());
@@ -379,6 +630,7 @@ fn build_components_once() {
"components/spin/Cargo.toml",
"components/clock-probe/Cargo.toml",
"components/kv-probe/Cargo.toml",
"components/resident-probe/Cargo.toml",
"components/wasi-clock-probe/Cargo.toml",
] {
let status = Command::new("rustup")
@@ -421,9 +673,18 @@ fn test_limits() -> ResourceLimits {
mailbox_capacity: 16,
max_input_bytes: 16 * 1024,
max_output_bytes: 16 * 1024,
resident: ResidentLimits::default(),
}
}
fn resource_id(value: u64) -> ResourceId {
ResourceId::new(value).expect("test resource id must be non-zero")
}
fn socket_addr(value: &str) -> SocketAddr {
value.parse().expect("test socket address must be valid")
}
fn decode_counter(bytes: Vec<u8>) -> u64 {
let array: [u8; 8] = bytes
.try_into()
@@ -447,6 +708,9 @@ fn component_world(artifact_name: &str) -> &'static str {
"spin_component.wasm" => "component:spin/spin-component@0.1.0",
"clock_probe_component.wasm" => "component:clock-probe/clock-probe-component@0.1.0",
"kv_probe_component.wasm" => "component:kv-probe/kv-probe-component@0.1.0",
"resident_probe_component.wasm" => {
"component:resident-probe/resident-probe-component@0.1.0"
}
"wasi_clock_probe_component.wasm" => {
"component:wasi-clock-probe/wasi-clock-probe-component@0.1.0"
}