diff --git a/bins/validator-node/src/wasm_executor.rs b/bins/validator-node/src/wasm_executor.rs index 5ec25adc..f12de48d 100644 --- a/bins/validator-node/src/wasm_executor.rs +++ b/bins/validator-node/src/wasm_executor.rs @@ -91,6 +91,7 @@ impl WasmChallengeExecutor { validator_id: "validator".to_string(), restart_id: String::new(), config_version: 0, + ..Default::default() }; let mut instance = self @@ -181,6 +182,7 @@ impl WasmChallengeExecutor { validator_id: "validator".to_string(), restart_id: String::new(), config_version: 0, + ..Default::default() }; let mut instance = self diff --git a/crates/wasm-runtime-interface/src/lib.rs b/crates/wasm-runtime-interface/src/lib.rs index 34227647..c00f5707 100644 --- a/crates/wasm-runtime-interface/src/lib.rs +++ b/crates/wasm-runtime-interface/src/lib.rs @@ -23,11 +23,13 @@ pub use exec::{ ExecError, ExecHostFunction, ExecHostFunctions, ExecPolicy, ExecRequest, ExecResponse, ExecState, }; -pub use network::{NetworkHostFunctions, NetworkState, NetworkStateError}; +pub use network::{ + NetworkHostFunctions, NetworkState, NetworkStateError, HOST_GET_TIMESTAMP, HOST_LOG_MESSAGE, +}; pub use storage::{ InMemoryStorageBackend, NoopStorageBackend, StorageAuditEntry, StorageAuditLogger, StorageBackend, StorageDeleteRequest, StorageGetRequest, StorageGetResponse, StorageHostConfig, - StorageHostError, StorageHostState, StorageHostStatus, StorageOperation, + StorageHostError, StorageHostFunctions, StorageHostState, StorageHostStatus, StorageOperation, StorageProposeWriteRequest, StorageProposeWriteResponse, }; @@ -43,7 +45,7 @@ pub use runtime::{ }; pub use storage::{ HOST_STORAGE_ALLOC, HOST_STORAGE_DELETE, HOST_STORAGE_GET, HOST_STORAGE_GET_RESULT, - HOST_STORAGE_NAMESPACE, HOST_STORAGE_PROPOSE_WRITE, + HOST_STORAGE_NAMESPACE, HOST_STORAGE_PROPOSE_WRITE, HOST_STORAGE_SET, }; pub use time::{TimeError, TimeHostFunction, TimeHostFunctions, TimeMode, TimePolicy, TimeState}; diff --git a/crates/wasm-runtime-interface/src/network.rs b/crates/wasm-runtime-interface/src/network.rs index 5b726362..b3245e0b 100644 --- a/crates/wasm-runtime-interface/src/network.rs +++ b/crates/wasm-runtime-interface/src/network.rs @@ -13,7 +13,7 @@ use std::io::Read; use std::net::IpAddr; use std::sync::Arc; use std::time::{Duration, Instant}; -use tracing::{info, warn}; +use tracing::{error, info, warn}; use trust_dns_resolver::config::{ResolverConfig, ResolverOpts}; use trust_dns_resolver::proto::rr::RecordType; use trust_dns_resolver::Resolver; @@ -21,6 +21,12 @@ use wasmtime::{Caller, Linker, Memory}; use crate::runtime::{HostFunctionRegistrar, RuntimeState, WasmRuntimeError}; +pub const HOST_LOG_MESSAGE: &str = "log_message"; +pub const HOST_GET_TIMESTAMP: &str = "get_timestamp"; + +const DEFAULT_RESPONSE_BUF_SIZE: i32 = 65536; +const DEFAULT_DNS_BUF_SIZE: i32 = 4096; + #[derive(Debug, thiserror::Error)] pub enum NetworkStateError { #[error("network policy invalid: {0}")] @@ -86,10 +92,9 @@ impl HostFunctionRegistrar for NetworkHostFunctions { |mut caller: Caller, req_ptr: i32, req_len: i32, - resp_ptr: i32, - resp_len: i32| + resp_ptr: i32| -> i32 { - handle_http_get(&mut caller, req_ptr, req_len, resp_ptr, resp_len) + handle_http_get(&mut caller, req_ptr, req_len, resp_ptr) }, ) .map_err(|err| WasmRuntimeError::HostFunction(err.to_string()))?; @@ -104,7 +109,8 @@ impl HostFunctionRegistrar for NetworkHostFunctions { req_ptr: i32, req_len: i32, resp_ptr: i32, - resp_len: i32| + resp_len: i32, + _extra: i32| -> i32 { handle_http_post(&mut caller, req_ptr, req_len, resp_ptr, resp_len) }, @@ -120,15 +126,32 @@ impl HostFunctionRegistrar for NetworkHostFunctions { |mut caller: Caller, req_ptr: i32, req_len: i32, - resp_ptr: i32, - resp_len: i32| + resp_ptr: i32| -> i32 { - handle_dns_request(&mut caller, req_ptr, req_len, resp_ptr, resp_len) + handle_dns_request(&mut caller, req_ptr, req_len, resp_ptr) }, ) .map_err(|err| WasmRuntimeError::HostFunction(err.to_string()))?; } + linker + .func_wrap( + HOST_FUNCTION_NAMESPACE, + HOST_LOG_MESSAGE, + |mut caller: Caller, level: i32, msg_ptr: i32, msg_len: i32| { + handle_log_message(&mut caller, level, msg_ptr, msg_len); + }, + ) + .map_err(|err| WasmRuntimeError::HostFunction(err.to_string()))?; + + linker + .func_wrap( + HOST_FUNCTION_NAMESPACE, + HOST_GET_TIMESTAMP, + |caller: Caller| -> i64 { handle_get_timestamp(&caller) }, + ) + .map_err(|err| WasmRuntimeError::HostFunction(err.to_string()))?; + Ok(()) } } @@ -147,11 +170,6 @@ pub struct NetworkState { } impl NetworkState { - /// Create network state for a single WASM execution. - /// - /// The policy is validated and then enforced by every host call. All host - /// functions share the same counters and audit logger, providing a single - /// enforcement surface for HTTP and DNS access. pub fn new( policy: NetworkPolicy, audit_logger: Option>, @@ -578,8 +596,8 @@ fn handle_http_get( req_ptr: i32, req_len: i32, resp_ptr: i32, - resp_len: i32, ) -> i32 { + let resp_len = DEFAULT_RESPONSE_BUF_SIZE; let enforcement = "http_get"; let request_bytes = match read_memory(caller, req_ptr, req_len) { Ok(bytes) => bytes, @@ -678,8 +696,8 @@ fn handle_dns_request( req_ptr: i32, req_len: i32, resp_ptr: i32, - resp_len: i32, ) -> i32 { + let resp_len = DEFAULT_DNS_BUF_SIZE; let enforcement = "dns_resolve"; let request_bytes = match read_memory(caller, req_ptr, req_len) { Ok(bytes) => bytes, @@ -716,6 +734,34 @@ fn handle_dns_request( write_result(caller, resp_ptr, resp_len, result) } +fn handle_log_message(caller: &mut Caller, level: i32, msg_ptr: i32, msg_len: i32) { + let msg = match read_memory(caller, msg_ptr, msg_len) { + Ok(bytes) => String::from_utf8_lossy(&bytes).into_owned(), + Err(err) => { + warn!( + challenge_id = %caller.data().challenge_id, + error = %err, + "log_message: failed to read message from wasm memory" + ); + return; + } + }; + + let challenge_id = caller.data().challenge_id.clone(); + match level { + 0 => info!(challenge_id = %challenge_id, "[wasm] {}", msg), + 1 => warn!(challenge_id = %challenge_id, "[wasm] {}", msg), + _ => error!(challenge_id = %challenge_id, "[wasm] {}", msg), + } +} + +fn handle_get_timestamp(caller: &Caller) -> i64 { + if let Some(ts) = caller.data().fixed_timestamp_ms { + return ts; + } + chrono::Utc::now().timestamp_millis() +} + fn resolve_dns( resolver: &Resolver, request: &DnsRequest, diff --git a/crates/wasm-runtime-interface/src/runtime.rs b/crates/wasm-runtime-interface/src/runtime.rs index 909d0adc..ca21d3c7 100644 --- a/crates/wasm-runtime-interface/src/runtime.rs +++ b/crates/wasm-runtime-interface/src/runtime.rs @@ -1,7 +1,11 @@ use crate::bridge::{self, BridgeError, EvalRequest, EvalResponse}; use crate::exec::{ExecPolicy, ExecState}; +use crate::storage::{ + InMemoryStorageBackend, StorageBackend, StorageHostConfig, StorageHostFunctions, + StorageHostState, +}; use crate::time::{TimePolicy, TimeState}; -use crate::{NetworkAuditLogger, NetworkPolicy, NetworkState}; +use crate::{NetworkAuditLogger, NetworkHostFunctions, NetworkPolicy, NetworkState}; use std::sync::Arc; use std::time::Instant; use thiserror::Error; @@ -103,6 +107,12 @@ pub struct InstanceConfig { pub restart_id: String, /// Configuration version for hot-restarts. pub config_version: u64, + /// Storage host function configuration. + pub storage_host_config: StorageHostConfig, + /// Storage backend implementation. + pub storage_backend: Arc, + /// Fixed timestamp for deterministic consensus execution. + pub fixed_timestamp_ms: Option, } impl Default for InstanceConfig { @@ -117,6 +127,9 @@ impl Default for InstanceConfig { validator_id: "unknown".to_string(), restart_id: String::new(), config_version: 0, + storage_host_config: StorageHostConfig::default(), + storage_backend: Arc::new(InMemoryStorageBackend::new()), + fixed_timestamp_ms: None, } } } @@ -140,6 +153,10 @@ pub struct RuntimeState { pub restart_id: String, /// Configuration version for hot-restarts. pub config_version: u64, + /// Storage host state for key-value operations. + pub storage_state: StorageHostState, + /// Fixed timestamp in milliseconds for deterministic consensus execution. + pub fixed_timestamp_ms: Option, limits: StoreLimits, } @@ -155,6 +172,8 @@ impl RuntimeState { validator_id: String, restart_id: String, config_version: u64, + storage_state: StorageHostState, + fixed_timestamp_ms: Option, limits: StoreLimits, ) -> Self { Self { @@ -167,6 +186,8 @@ impl RuntimeState { validator_id, restart_id, config_version, + storage_state, + fixed_timestamp_ms, limits, } } @@ -175,6 +196,10 @@ impl RuntimeState { self.network_state.reset_counters(); } + pub fn reset_storage_counters(&mut self) { + self.storage_state.reset_counters(); + } + pub fn reset_exec_counters(&mut self) { self.exec_state.reset_counters(); } @@ -242,6 +267,11 @@ impl WasmRuntime { instance_config.validator_id.clone(), ) .map_err(|err| WasmRuntimeError::HostFunction(err.to_string()))?; + let storage_state = StorageHostState::new( + instance_config.challenge_id.clone(), + instance_config.storage_host_config.clone(), + Arc::clone(&instance_config.storage_backend), + ); let exec_state = ExecState::new( instance_config.exec_policy.clone(), instance_config.challenge_id.clone(), @@ -262,6 +292,8 @@ impl WasmRuntime { instance_config.validator_id.clone(), instance_config.restart_id.clone(), instance_config.config_version, + storage_state, + instance_config.fixed_timestamp_ms, limits.build(), ); let mut store = Store::new(&self.engine, runtime_state); @@ -277,6 +309,13 @@ impl WasmRuntime { store.limiter(|state| &mut state.limits); let mut linker = Linker::new(&self.engine); + + let network_host_fns = NetworkHostFunctions::all(); + network_host_fns.register(&mut linker)?; + + let storage_host_fns = StorageHostFunctions::new(); + storage_host_fns.register(&mut linker)?; + if let Some(registrar) = registrar { registrar.register(&mut linker)?; } @@ -436,6 +475,22 @@ impl ChallengeInstance { self.store.data_mut().reset_network_counters(); } + pub fn reset_storage_state(&mut self) { + self.store.data_mut().reset_storage_counters(); + } + + pub fn storage_bytes_read(&self) -> u64 { + self.store.data().storage_state.bytes_read + } + + pub fn storage_bytes_written(&self) -> u64 { + self.store.data().storage_state.bytes_written + } + + pub fn storage_operations_count(&self) -> u32 { + self.store.data().storage_state.operations_count + } + pub fn challenge_id(&self) -> &str { &self.store.data().challenge_id } diff --git a/crates/wasm-runtime-interface/src/storage.rs b/crates/wasm-runtime-interface/src/storage.rs index 8c184949..d8da109c 100644 --- a/crates/wasm-runtime-interface/src/storage.rs +++ b/crates/wasm-runtime-interface/src/storage.rs @@ -5,7 +5,8 @@ //! //! # Host Functions //! -//! - `storage_get(key_ptr, key_len) -> i64` - Read from storage +//! - `storage_get(key_ptr, key_len, value_ptr) -> i32` - Read from storage +//! - `storage_set(key_ptr, key_len, value_ptr, value_len) -> i32` - Write to storage //! - `storage_propose_write(key_ptr, key_len, value_ptr, value_len) -> i64` - Propose a write //! - `storage_delete(key_ptr, key_len) -> i32` - Delete from storage (requires consensus) //! @@ -20,17 +21,19 @@ //! - Not found: returns 0 //! - Error: returns negative status code -#![allow(dead_code, unused_variables, unused_imports)] - use serde::{Deserialize, Serialize}; use sha2::{Digest, Sha256}; use std::collections::HashMap; use std::sync::{Arc, RwLock}; use thiserror::Error; -use tracing::{debug, info, warn}; +use tracing::warn; +use wasmtime::{Caller, Linker, Memory}; + +use crate::runtime::{HostFunctionRegistrar, RuntimeState, WasmRuntimeError}; pub const HOST_STORAGE_NAMESPACE: &str = "platform_storage"; pub const HOST_STORAGE_GET: &str = "storage_get"; +pub const HOST_STORAGE_SET: &str = "storage_set"; pub const HOST_STORAGE_PROPOSE_WRITE: &str = "storage_propose_write"; pub const HOST_STORAGE_DELETE: &str = "storage_delete"; pub const HOST_STORAGE_GET_RESULT: &str = "storage_get_result"; @@ -217,6 +220,7 @@ pub struct StorageDeleteRequest { pub struct StorageHostState { pub config: StorageHostConfig, pub challenge_id: String, + pub backend: Arc, pub pending_results: HashMap>, pub next_result_id: u32, pub bytes_read: u64, @@ -225,10 +229,15 @@ pub struct StorageHostState { } impl StorageHostState { - pub fn new(challenge_id: String, config: StorageHostConfig) -> Self { + pub fn new( + challenge_id: String, + config: StorageHostConfig, + backend: Arc, + ) -> Self { Self { config, challenge_id, + backend, pending_results: HashMap::new(), next_result_id: 1, bytes_read: 0, @@ -381,6 +390,7 @@ pub struct StorageAuditEntry { #[serde(rename_all = "snake_case")] pub enum StorageOperation { Get, + Set, ProposeWrite, Delete, } @@ -395,6 +405,202 @@ impl StorageAuditLogger for NoopStorageAuditLogger { fn record(&self, _entry: StorageAuditEntry) {} } +#[derive(Clone, Debug)] +pub struct StorageHostFunctions; + +impl StorageHostFunctions { + pub fn new() -> Self { + Self + } +} + +impl Default for StorageHostFunctions { + fn default() -> Self { + Self::new() + } +} + +impl HostFunctionRegistrar for StorageHostFunctions { + fn register(&self, linker: &mut Linker) -> Result<(), WasmRuntimeError> { + linker + .func_wrap( + HOST_STORAGE_NAMESPACE, + HOST_STORAGE_GET, + |mut caller: Caller, + key_ptr: i32, + key_len: i32, + value_ptr: i32| + -> i32 { + handle_storage_get(&mut caller, key_ptr, key_len, value_ptr) + }, + ) + .map_err(|err| WasmRuntimeError::HostFunction(err.to_string()))?; + + linker + .func_wrap( + HOST_STORAGE_NAMESPACE, + HOST_STORAGE_SET, + |mut caller: Caller, + key_ptr: i32, + key_len: i32, + value_ptr: i32, + value_len: i32| + -> i32 { + handle_storage_set(&mut caller, key_ptr, key_len, value_ptr, value_len) + }, + ) + .map_err(|err| WasmRuntimeError::HostFunction(err.to_string()))?; + + Ok(()) + } +} + +fn handle_storage_get( + caller: &mut Caller, + key_ptr: i32, + key_len: i32, + value_ptr: i32, +) -> i32 { + let key = match read_wasm_memory(caller, key_ptr, key_len) { + Ok(bytes) => bytes, + Err(err) => { + warn!(error = %err, "storage_get: failed to read key from wasm memory"); + return StorageHostStatus::InternalError.to_i32(); + } + }; + + let storage = &caller.data().storage_state; + if let Err(err) = storage.config.validate_key(&key) { + warn!(error = %err, "storage_get: key validation failed"); + return StorageHostStatus::from(err).to_i32(); + } + + let challenge_id = storage.challenge_id.clone(); + let backend = Arc::clone(&storage.backend); + + let value = match backend.get(&challenge_id, &key) { + Ok(Some(v)) => v, + Ok(None) => return 0, + Err(err) => { + warn!(error = %err, "storage_get: backend read failed"); + return StorageHostStatus::from(err).to_i32(); + } + }; + + caller.data_mut().storage_state.bytes_read += value.len() as u64; + caller.data_mut().storage_state.operations_count += 1; + + if let Err(err) = write_wasm_memory(caller, value_ptr, &value) { + warn!(error = %err, "storage_get: failed to write value to wasm memory"); + return StorageHostStatus::InternalError.to_i32(); + } + + value.len() as i32 +} + +fn handle_storage_set( + caller: &mut Caller, + key_ptr: i32, + key_len: i32, + value_ptr: i32, + value_len: i32, +) -> i32 { + let key = match read_wasm_memory(caller, key_ptr, key_len) { + Ok(bytes) => bytes, + Err(err) => { + warn!(error = %err, "storage_set: failed to read key from wasm memory"); + return StorageHostStatus::InternalError.to_i32(); + } + }; + + let value = match read_wasm_memory(caller, value_ptr, value_len) { + Ok(bytes) => bytes, + Err(err) => { + warn!(error = %err, "storage_set: failed to read value from wasm memory"); + return StorageHostStatus::InternalError.to_i32(); + } + }; + + let storage = &caller.data().storage_state; + if let Err(err) = storage.config.validate_key(&key) { + warn!(error = %err, "storage_set: key validation failed"); + return StorageHostStatus::from(err).to_i32(); + } + if let Err(err) = storage.config.validate_value(&value) { + warn!(error = %err, "storage_set: value validation failed"); + return StorageHostStatus::from(err).to_i32(); + } + + if storage.config.require_consensus && !storage.config.allow_direct_writes { + warn!("storage_set: direct writes require consensus or allow_direct_writes"); + return StorageHostStatus::ConsensusRequired.to_i32(); + } + + let challenge_id = storage.challenge_id.clone(); + let backend = Arc::clone(&storage.backend); + + match backend.propose_write(&challenge_id, &key, &value) { + Ok(_proposal_id) => { + caller.data_mut().storage_state.bytes_written += value.len() as u64; + caller.data_mut().storage_state.operations_count += 1; + StorageHostStatus::Success.to_i32() + } + Err(err) => { + warn!(error = %err, "storage_set: backend write failed"); + StorageHostStatus::from(err).to_i32() + } + } +} + +fn read_wasm_memory( + caller: &mut Caller, + ptr: i32, + len: i32, +) -> Result, String> { + if ptr < 0 || len < 0 { + return Err("negative pointer/length".to_string()); + } + let ptr = ptr as usize; + let len = len as usize; + let memory = get_memory(caller).ok_or_else(|| "memory export not found".to_string())?; + let data = memory.data(caller); + let end = ptr + .checked_add(len) + .ok_or_else(|| "pointer overflow".to_string())?; + if end > data.len() { + return Err("memory read out of bounds".to_string()); + } + Ok(data[ptr..end].to_vec()) +} + +fn write_wasm_memory( + caller: &mut Caller, + ptr: i32, + bytes: &[u8], +) -> Result<(), String> { + if ptr < 0 { + return Err("negative pointer".to_string()); + } + let ptr = ptr as usize; + let memory = get_memory(caller).ok_or_else(|| "memory export not found".to_string())?; + let end = ptr + .checked_add(bytes.len()) + .ok_or_else(|| "pointer overflow".to_string())?; + let data = memory.data_mut(caller); + if end > data.len() { + return Err("memory write out of bounds".to_string()); + } + data[ptr..end].copy_from_slice(bytes); + Ok(()) +} + +fn get_memory(caller: &mut Caller) -> Option { + let memory_export = caller.data().memory_export.clone(); + caller + .get_export(&memory_export) + .and_then(|export| export.into_memory()) +} + #[cfg(test)] mod tests { use super::*; @@ -443,8 +649,12 @@ mod tests { #[test] fn test_storage_host_state() { - let mut state = - StorageHostState::new("challenge-1".to_string(), StorageHostConfig::default()); + let backend = Arc::new(InMemoryStorageBackend::new()); + let mut state = StorageHostState::new( + "challenge-1".to_string(), + StorageHostConfig::default(), + backend, + ); let id1 = state.store_result(b"result1".to_vec()); let id2 = state.store_result(b"result2".to_vec());