diff --git a/crates/wasmeld-runtime/src/capability/kv.rs b/crates/wasmeld-runtime/src/capability/kv.rs new file mode 100644 index 0000000..3dd216c --- /dev/null +++ b/crates/wasmeld-runtime/src/capability/kv.rs @@ -0,0 +1,145 @@ +//! `wasmeld:kv/store@0.1.0` host implementation and backend boundary. + +use std::{fmt, sync::Arc, time::Duration}; + +use thiserror::Error; +use wasmtime::component::Linker; + +use crate::{RuntimeError, runtime::HostState}; + +wasmtime::component::bindgen!({ + path: "../../wit/kv", + world: "kv-host", +}); + +pub(crate) const INTERFACE: &str = "wasmeld:kv/store@0.1.0"; +pub(crate) const PACKAGE: &str = "wasmeld:kv"; +pub(crate) const NAME: &str = "store"; +pub(crate) const VERSION: &str = "0.1.0"; + +/// Maximum UTF-8 byte length accepted for one KV key. +pub const KV_MAX_KEY_BYTES: usize = 256; +/// Maximum byte length accepted for one KV value. +pub const KV_MAX_VALUE_BYTES: usize = 64 * 1024; + +/// Synchronous storage boundary used by the Wasmtime host implementation. +/// +/// Implementations must bound their own blocking time because calls execute on +/// the Component Actor thread. +pub trait KvBackend: fmt::Debug + Send + Sync { + /// Returns the value stored under a service-scoped key. + fn get( + &self, + service_id: &str, + key: &str, + timeout: Duration, + ) -> Result>, KvBackendError>; + /// Creates or replaces a service-scoped value. + fn set( + &self, + service_id: &str, + key: &str, + value: &[u8], + timeout: Duration, + ) -> Result<(), KvBackendError>; + /// Removes a service-scoped value if it exists. + fn delete(&self, service_id: &str, key: &str, timeout: Duration) -> Result<(), KvBackendError>; +} + +/// Storage failures returned by a [`KvBackend`]. +#[derive(Debug, Error)] +pub enum KvBackendError { + #[error("KV backend is busy")] + Busy, + #[error("KV backend timed out")] + Timeout, + #[error("KV backend operation failed: {0}")] + Operation(String), +} + +pub(crate) struct KvState { + service_id: String, + backend: Arc, + timeout: Duration, +} + +impl KvState { + pub(crate) fn new(service_id: &str, backend: Arc, timeout: Duration) -> Self { + Self { + service_id: service_id.to_owned(), + backend, + timeout, + } + } + + fn get(&self, key: &str) -> Result>, wasmeld::kv::store::KvError> { + validate_key(key)?; + self.backend + .get(&self.service_id, key, self.timeout) + .map_err(operation_failed) + .and_then(|value| match value { + Some(value) if value.len() > KV_MAX_VALUE_BYTES => { + Err(wasmeld::kv::store::KvError::ValueTooLarge( + u64::try_from(value.len()).unwrap_or(u64::MAX), + )) + } + value => Ok(value), + }) + } + + fn set(&self, key: &str, value: &[u8]) -> Result<(), wasmeld::kv::store::KvError> { + validate_key(key)?; + if value.len() > KV_MAX_VALUE_BYTES { + return Err(wasmeld::kv::store::KvError::ValueTooLarge( + u64::try_from(value.len()).unwrap_or(u64::MAX), + )); + } + self.backend + .set(&self.service_id, key, value, self.timeout) + .map_err(operation_failed) + } + + fn delete(&self, key: &str) -> Result<(), wasmeld::kv::store::KvError> { + validate_key(key)?; + self.backend + .delete(&self.service_id, key, self.timeout) + .map_err(operation_failed) + } +} + +impl wasmeld::kv::store::Host for HostState { + fn get(&mut self, key: String) -> Result>, wasmeld::kv::store::KvError> { + self.capabilities.kv().get(&key) + } + + fn set(&mut self, key: String, value: Vec) -> Result<(), wasmeld::kv::store::KvError> { + self.capabilities.kv().set(&key, &value) + } + + fn delete(&mut self, key: String) -> Result<(), wasmeld::kv::store::KvError> { + self.capabilities.kv().delete(&key) + } +} + +pub(crate) fn add_to_linker(linker: &mut Linker) -> Result<(), RuntimeError> { + wasmeld::kv::store::add_to_linker::(linker, |state| state) + .map_err(|error| RuntimeError::ActorInitialization(error.to_string())) +} + +fn validate_key(key: &str) -> Result<(), wasmeld::kv::store::KvError> { + if key.is_empty() { + return Err(wasmeld::kv::store::KvError::InvalidKey( + "key must not be empty".to_owned(), + )); + } + if key.len() > KV_MAX_KEY_BYTES { + return Err(wasmeld::kv::store::KvError::InvalidKey(format!( + "key exceeds the {KV_MAX_KEY_BYTES}-byte limit" + ))); + } + Ok(()) +} + +fn operation_failed(error: KvBackendError) -> wasmeld::kv::store::KvError { + wasmeld::kv::store::KvError::OperationFailed(error.to_string()) +} diff --git a/crates/wasmeld-runtime/src/capability/mod.rs b/crates/wasmeld-runtime/src/capability/mod.rs index 83d51d1..8701e52 100644 --- a/crates/wasmeld-runtime/src/capability/mod.rs +++ b/crates/wasmeld-runtime/src/capability/mod.rs @@ -1,22 +1,37 @@ //! Versioned host capabilities available to WebAssembly Components. mod clock; +mod kv; mod registry; +pub use kv::{KV_MAX_KEY_BYTES, KV_MAX_VALUE_BYTES, KvBackend, KvBackendError}; pub use registry::CapabilityDescriptor; pub(crate) use registry::{Capability, CapabilityRegistry}; /// Per-Actor state owned by enabled host capabilities. pub(crate) struct HostCapabilities { clock: Option, + kv: Option, } impl HostCapabilities { - pub(crate) fn new(capabilities: &[Capability]) -> Self { + pub(crate) fn new( + capabilities: &[Capability], + service_id: &str, + kv_backend: Option>, + deadline: std::time::Duration, + ) -> Self { Self { clock: capabilities .contains(&Capability::MonotonicClock) .then(clock::ClockState::new), + kv: capabilities.contains(&Capability::KvStore).then(|| { + kv::KvState::new( + service_id, + kv_backend.expect("KV Components are rejected when no backend is configured"), + deadline, + ) + }), } } @@ -26,4 +41,10 @@ impl HostCapabilities { .expect("the Clock interface is linked only when its state is enabled") .now() } + + fn kv(&self) -> &kv::KvState { + self.kv + .as_ref() + .expect("the KV interface is linked only when its state is enabled") + } } diff --git a/crates/wasmeld-runtime/src/capability/registry.rs b/crates/wasmeld-runtime/src/capability/registry.rs index eee337a..d003b80 100644 --- a/crates/wasmeld-runtime/src/capability/registry.rs +++ b/crates/wasmeld-runtime/src/capability/registry.rs @@ -2,7 +2,7 @@ use wasmtime::component::Linker; -use super::clock; +use super::{clock, kv}; use crate::{RuntimeError, runtime::HostState}; /// Stable description of one versioned WIT interface used by a Component. @@ -53,6 +53,7 @@ impl CapabilityDescriptor { #[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)] pub(crate) enum Capability { MonotonicClock, + KvStore, } impl Capability { @@ -64,12 +65,16 @@ impl Capability { clock::NAME, clock::VERSION, ), + Self::KvStore => { + CapabilityDescriptor::new(kv::INTERFACE, kv::PACKAGE, kv::NAME, kv::VERSION) + } } } pub(crate) fn add_to_linker(self, linker: &mut Linker) -> Result<(), RuntimeError> { match self { Self::MonotonicClock => clock::add_to_linker(linker), + Self::KvStore => kv::add_to_linker(linker), } } } @@ -82,6 +87,7 @@ impl CapabilityRegistry { pub(crate) fn resolve(interface: &str) -> Option { match interface { clock::INTERFACE => Some(Capability::MonotonicClock), + kv::INTERFACE => Some(Capability::KvStore), _ => None, } } @@ -116,5 +122,10 @@ mod tests { CapabilityRegistry::resolve("wasmeld:clock/other@0.1.0"), None ); + assert_eq!( + CapabilityRegistry::resolve("wasmeld:kv/store@0.1.0"), + Some(Capability::KvStore) + ); + assert_eq!(CapabilityRegistry::resolve("wasmeld:kv/store@0.2.0"), None); } } diff --git a/crates/wasmeld-runtime/src/error.rs b/crates/wasmeld-runtime/src/error.rs index b3d7661..4b22831 100644 --- a/crates/wasmeld-runtime/src/error.rs +++ b/crates/wasmeld-runtime/src/error.rs @@ -30,6 +30,9 @@ pub enum RuntimeError { #[error("component imports unsupported capability {0}")] UnsupportedImport(String), + #[error("component requires unavailable capability {0}")] + CapabilityUnavailable(String), + #[error("runtime creation failed: {0}")] RuntimeCreation(#[source] wasmtime::Error), diff --git a/crates/wasmeld-runtime/src/lib.rs b/crates/wasmeld-runtime/src/lib.rs index 9b29bdb..448494a 100644 --- a/crates/wasmeld-runtime/src/lib.rs +++ b/crates/wasmeld-runtime/src/lib.rs @@ -11,7 +11,9 @@ mod error; mod manifest; mod runtime; -pub use capability::CapabilityDescriptor; +pub use capability::{ + CapabilityDescriptor, KV_MAX_KEY_BYTES, KV_MAX_VALUE_BYTES, KvBackend, KvBackendError, +}; pub use error::RuntimeError; pub use manifest::{ResourceLimits, ServiceKey, ServiceManifest}; pub use runtime::{ActorHandle, Runtime, RuntimeConfig, RuntimeStats}; diff --git a/crates/wasmeld-runtime/src/runtime.rs b/crates/wasmeld-runtime/src/runtime.rs index d0b5fc8..7cc064e 100644 --- a/crates/wasmeld-runtime/src/runtime.rs +++ b/crates/wasmeld-runtime/src/runtime.rs @@ -30,7 +30,7 @@ use wasmtime_wasi::{ }; use crate::{ - CapabilityDescriptor, RuntimeError, ServiceKey, ServiceManifest, + CapabilityDescriptor, KvBackend, RuntimeError, ServiceKey, ServiceManifest, bindings::{ServiceComponent, ServiceError}, capability::{Capability, CapabilityRegistry, HostCapabilities}, manifest::ResourceLimits, @@ -61,6 +61,8 @@ pub struct RuntimeConfig { 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 { @@ -68,6 +70,7 @@ impl Default for RuntimeConfig { Self { epoch_tick: DEFAULT_EPOCH_TICK, max_wasm_stack: DEFAULT_MAX_WASM_STACK, + kv_backend: None, } } } @@ -117,6 +120,7 @@ struct RuntimeInner { actors: Mutex>, lifecycle: Mutex<()>, ticker: Arc, + kv_backend: Option>, } struct EpochTicker { @@ -189,6 +193,7 @@ impl Runtime { actors: Mutex::new(HashMap::new()), lifecycle: Mutex::new(()), ticker, + kv_backend: config.kv_backend, }), }) } @@ -207,6 +212,11 @@ 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)?; + 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 @@ -298,6 +308,7 @@ impl Runtime { key.clone(), service, init_config, + self.inner.kv_backend.clone(), )?; self.inner @@ -580,7 +591,12 @@ pub(crate) struct HostState { } impl HostState { - fn new(limits: &ResourceLimits, capabilities: &[Capability]) -> Self { + 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) @@ -600,7 +616,12 @@ impl HostState { store_limits, wasi: wasi.build(), wasi_table: ResourceTable::new(), - capabilities: HostCapabilities::new(capabilities), + capabilities: HostCapabilities::new( + capabilities, + service_id, + kv_backend, + limits.deadline(), + ), } } } @@ -639,6 +660,7 @@ impl ActorWorker { 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 { @@ -648,7 +670,10 @@ impl ActorWorker { }); } - let mut store = Store::new(&engine, HostState::new(&limits, &service.capabilities)); + 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); @@ -797,6 +822,7 @@ fn spawn_actor( 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); @@ -811,7 +837,7 @@ fn spawn_actor( .spawn(move || { let _ticker = worker_ticker; - match ActorWorker::create(engine, epoch_tick, key, service, init_config) { + match ActorWorker::create(engine, epoch_tick, key, service, init_config, kv_backend) { Ok(mut worker) => { worker_status.mark_ready(); let _ = ready_sender.send(Ok(())); diff --git a/crates/wasmeld-runtime/tests/runtime_components.rs b/crates/wasmeld-runtime/tests/runtime_components.rs index e8aae78..70ddb15 100644 --- a/crates/wasmeld-runtime/tests/runtime_components.rs +++ b/crates/wasmeld-runtime/tests/runtime_components.rs @@ -1,13 +1,17 @@ use std::{ + collections::HashMap, fs, path::{Path, PathBuf}, process::Command, - sync::OnceLock, + sync::{Arc, Mutex, OnceLock}, thread, time::{Duration, Instant}, }; -use wasmeld_runtime::{ResourceLimits, Runtime, RuntimeConfig, RuntimeError, ServiceManifest}; +use wasmeld_runtime::{ + KV_MAX_VALUE_BYTES, KvBackend, KvBackendError, ResourceLimits, Runtime, RuntimeConfig, + RuntimeError, ServiceManifest, +}; static COMPONENTS_BUILT: OnceLock<()> = OnceLock::new(); @@ -188,6 +192,97 @@ fn explicit_clock_capability_is_linked() { runtime.stop(&key).unwrap(); } +#[test] +fn kv_capability_is_service_scoped_and_shared_across_revisions() { + let backend = Arc::new(MemoryKvBackend::default()); + let runtime = Runtime::new(RuntimeConfig { + kv_backend: Some(backend), + ..RuntimeConfig::default() + }) + .expect("runtime should start"); + let first_key = register_component_revision( + &runtime, + "shared-kv", + "0.1.0", + "kv_probe_component.wasm", + test_limits(), + ); + let second_key = register_component_revision( + &runtime, + "shared-kv", + "0.2.0", + "kv_probe_component.wasm", + test_limits(), + ); + let isolated_key = register_component_revision( + &runtime, + "isolated-kv", + "0.1.0", + "kv_probe_component.wasm", + test_limits(), + ); + + let first = runtime.start(&first_key, Vec::new()).unwrap(); + let second = runtime.start(&second_key, Vec::new()).unwrap(); + let isolated = runtime.start(&isolated_key, Vec::new()).unwrap(); + + assert_eq!(first.invoke(b"set:answer:forty-two".to_vec()).unwrap(), b""); + assert_eq!(second.invoke(b"get:answer".to_vec()).unwrap(), b"forty-two"); + assert_eq!(isolated.invoke(b"get:answer".to_vec()).unwrap(), b""); + assert_eq!( + runtime.capabilities(&first_key).unwrap()[0].interface(), + "wasmeld:kv/store@0.1.0" + ); + + runtime.stop_all().unwrap(); +} + +#[test] +fn kv_capability_requires_a_backend_and_enforces_limits() { + let artifact = component_artifact("kv_probe_component.wasm"); + let manifest = ServiceManifest { + id: "kv-without-backend".to_owned(), + revision: "0.1.0".to_owned(), + component: artifact.display().to_string(), + world: component_world("kv_probe_component.wasm").to_owned(), + limits: test_limits(), + }; + let bytes = fs::read(&artifact).unwrap(); + let runtime = Runtime::new(RuntimeConfig::default()).unwrap(); + assert!(matches!( + runtime.register(manifest, &bytes), + Err(RuntimeError::CapabilityUnavailable(interface)) + if interface == "wasmeld:kv/store@0.1.0" + )); + + let runtime = Runtime::new(RuntimeConfig { + kv_backend: Some(Arc::new(MemoryKvBackend::default())), + ..RuntimeConfig::default() + }) + .unwrap(); + let mut limits = test_limits(); + limits.max_input_bytes = 128 * 1024; + let key = register_component(&runtime, "kv-limits", "kv_probe_component.wasm", limits); + let actor = runtime.start(&key, Vec::new()).unwrap(); + + let long_key = format!("get:{}", "k".repeat(257)); + assert!(matches!( + actor.invoke(long_key.into_bytes()), + Err(RuntimeError::ComponentError { message, .. }) + if message.contains("InvalidKey") + )); + + let mut oversized = b"set:key:".to_vec(); + oversized.resize(oversized.len() + KV_MAX_VALUE_BYTES + 1, b'x'); + assert!(matches!( + actor.invoke(oversized), + Err(RuntimeError::ComponentError { message, .. }) + if message.contains("ValueTooLarge") + )); + + runtime.stop(&key).unwrap(); +} + #[test] fn registration_rejects_non_whitelisted_wasi_imports() { let runtime = Runtime::new(RuntimeConfig::default()).expect("runtime should start"); @@ -225,11 +320,21 @@ fn register_component( service_id: &str, artifact_name: &str, limits: ResourceLimits, +) -> wasmeld_runtime::ServiceKey { + register_component_revision(runtime, service_id, "0.1.0", artifact_name, limits) +} + +fn register_component_revision( + runtime: &Runtime, + service_id: &str, + revision: &str, + artifact_name: &str, + limits: ResourceLimits, ) -> wasmeld_runtime::ServiceKey { let artifact = component_artifact(artifact_name); let manifest = ServiceManifest { id: service_id.to_owned(), - revision: "0.1.0".to_owned(), + revision: revision.to_owned(), component: artifact.display().to_string(), world: component_world(artifact_name).to_owned(), limits, @@ -256,6 +361,7 @@ fn build_components_once() { "components/fault/Cargo.toml", "components/spin/Cargo.toml", "components/clock-probe/Cargo.toml", + "components/kv-probe/Cargo.toml", "components/wasi-clock-probe/Cargo.toml", ] { let status = Command::new("rustup") @@ -323,9 +429,58 @@ fn component_world(artifact_name: &str) -> &'static str { "fault_component.wasm" => "component:fault/fault-component@0.1.0", "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", "wasi_clock_probe_component.wasm" => { "component:wasi-clock-probe/wasi-clock-probe-component@0.1.0" } other => panic!("no WIT world registered for {other}"), } } + +#[derive(Debug, Default)] +struct MemoryKvBackend { + entries: Mutex>>, +} + +impl KvBackend for MemoryKvBackend { + fn get( + &self, + service_id: &str, + key: &str, + _timeout: Duration, + ) -> Result>, KvBackendError> { + Ok(self + .entries + .lock() + .unwrap() + .get(&(service_id.to_owned(), key.to_owned())) + .cloned()) + } + + fn set( + &self, + service_id: &str, + key: &str, + value: &[u8], + _timeout: Duration, + ) -> Result<(), KvBackendError> { + self.entries + .lock() + .unwrap() + .insert((service_id.to_owned(), key.to_owned()), value.to_vec()); + Ok(()) + } + + fn delete( + &self, + service_id: &str, + key: &str, + _timeout: Duration, + ) -> Result<(), KvBackendError> { + self.entries + .lock() + .unwrap() + .remove(&(service_id.to_owned(), key.to_owned())); + Ok(()) + } +}