Files
wasmeld/crates/wasmeld-runtime/src/runtime.rs
T
Maofeng b63cf8c9eb feat(runtime): add resident Wasmtime actor runtime
- define versioned service and clock WIT contracts\n- enforce import allowlists, fuel, epoch, memory, I/O, and mailbox limits\n- keep one Store and Component Instance resident per serial Actor\n- add echo, counter, fault, spin, and capability probe components\n- cover lifecycle, concurrency, sandbox, and fault recovery behavior
2026-07-27 04:59:33 +08:00

932 lines
28 KiB
Rust

//! 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::{
RuntimeError, ServiceKey, ServiceManifest,
bindings::{ServiceComponent, ServiceError, clock as clock_bindings},
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 MONOTONIC_CLOCK_IMPORT: &str = "wasmeld:clock/monotonic-clock@0.1.0";
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,
}
impl Default for RuntimeConfig {
fn default() -> Self {
Self {
epoch_tick: DEFAULT_EPOCH_TICK,
max_wasm_stack: DEFAULT_MAX_WASM_STACK,
}
}
}
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<RuntimeInner>,
}
/// 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<Engine>,
epoch_tick: Duration,
services: Mutex<HashMap<ServiceKey, RegisteredService>>,
actors: Mutex<HashMap<ServiceKey, ActorHandle>>,
lifecycle: Mutex<()>,
ticker: Arc<EpochTicker>,
}
struct EpochTicker {
shutdown: mpsc::Sender<()>,
worker: Mutex<Option<thread::JoinHandle<()>>>,
}
impl EpochTicker {
fn start(engine: Arc<Engine>, tick: Duration) -> Result<Arc<Self>, 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<Component>,
capabilities: Vec<HostCapability>,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum HostCapability {
MonotonicClock,
}
impl Runtime {
/// Creates the shared Wasmtime Engine and epoch interruption worker.
pub fn new(config: RuntimeConfig) -> Result<Self, RuntimeError> {
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,
}),
})
}
/// 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<ServiceKey, RuntimeError> {
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)?;
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<ServiceKey, RuntimeError> {
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<u8>,
) -> Result<ActorHandle, RuntimeError> {
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
.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<u8>) -> Result<Vec<u8>, RuntimeError> {
let actor = self
.inner
.actors
.lock()
.map_err(|_| RuntimeError::LockPoisoned)?
.get(key)
.cloned()
.ok_or_else(|| RuntimeError::ActorUnavailable(key.clone()))?;
actor.invoke(input)
}
/// 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<u8>,
) -> Result<ActorHandle, RuntimeError> {
let key = self.register(manifest, component_bytes)?;
self.start(&key, init_config)
}
/// Returns counts without exposing internal Engine or Store state.
pub fn stats(&self) -> Result<RuntimeStats, RuntimeError> {
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::<Vec<_>>();
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<ResourceLimits>,
sender: SyncSender<ActorCommand>,
status: Arc<ActorStatus>,
_ticker: Arc<EpochTicker>,
}
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<u8>) -> Result<Vec<u8>, 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<u8>,
deadline: Instant,
response_sender: SyncSender<Result<Vec<u8>, 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);
}
}
struct HostState {
store_limits: StoreLimits,
wasi: WasiCtx,
wasi_table: ResourceTable,
clock_origin: Instant,
}
impl HostState {
fn new(limits: &ResourceLimits) -> 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(),
clock_origin: Instant::now(),
}
}
}
impl HasData for HostState {
type Data<'a> = &'a mut HostState;
}
impl clock_bindings::wasmeld::clock::monotonic_clock::Host for HostState {
fn now(&mut self) -> u64 {
u64::try_from(self.clock_origin.elapsed().as_nanos()).unwrap_or(u64::MAX)
}
}
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<ResourceLimits>,
epoch_tick: Duration,
store: Store<HostState>,
bindings: ServiceComponent,
}
impl ActorWorker {
fn create(
engine: Arc<Engine>,
epoch_tick: Duration,
key: ServiceKey,
service: RegisteredService,
init_config: Vec<u8>,
) -> Result<Self, RuntimeError> {
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));
store.limiter(|state| &mut state.store_limits);
let mut linker = Linker::new(&engine);
add_restricted_wasi(&mut linker)?;
add_host_capabilities(&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<u8>, deadline: Instant) -> Result<Vec<u8>, 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<HostState>) -> Result<(), RuntimeError> {
let exit_options = cli::exit::LinkOptions::default();
link_wasi(clocks::wall_clock::add_to_linker::<HostState, WasiClocks>(
linker,
|state| state.clocks(),
))?;
link_wasi(filesystem::preopens::add_to_linker::<
HostState,
WasiFilesystem,
>(linker, |state| state.filesystem()))?;
link_wasi(
filesystem::types::add_to_linker::<HostState, WasiFilesystem>(linker, |state| {
state.filesystem()
}),
)?;
link_wasi(cli::exit::add_to_linker::<HostState, WasiCli>(
linker,
&exit_options,
|state| state.cli(),
))?;
link_wasi(cli::environment::add_to_linker::<HostState, WasiCli>(
linker,
|state| state.cli(),
))?;
link_wasi(cli::stdin::add_to_linker::<HostState, WasiCli>(
linker,
|state| state.cli(),
))?;
link_wasi(cli::stdout::add_to_linker::<HostState, WasiCli>(
linker,
|state| state.cli(),
))?;
link_wasi(cli::stderr::add_to_linker::<HostState, WasiCli>(
linker,
|state| state.cli(),
))?;
link_wasi(io::error::add_to_linker::<HostState, WasiIoTable>(
linker,
|state| state.ctx().table,
))?;
link_wasi(io::streams::add_to_linker::<HostState, WasiIoTable>(
linker,
|state| state.ctx().table,
))?;
Ok(())
}
fn link_wasi(result: wasmtime::Result<()>) -> Result<(), RuntimeError> {
result.map_err(|error| RuntimeError::ActorInitialization(error.to_string()))
}
fn add_host_capabilities(
linker: &mut Linker<HostState>,
capabilities: &[HostCapability],
) -> Result<(), RuntimeError> {
for capability in capabilities {
match capability {
HostCapability::MonotonicClock => link_wasi(
clock_bindings::wasmeld::clock::monotonic_clock::add_to_linker::<
HostState,
HostState,
>(linker, |state| state),
)?,
}
}
Ok(())
}
fn validate_component_imports(
component: &Component,
engine: &Engine,
) -> Result<Vec<HostCapability>, 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 name == MONOTONIC_CLOCK_IMPORT {
capabilities.push(HostCapability::MonotonicClock);
} else if !ALLOWED_WASI_IMPORTS
.iter()
.any(|allowed| name.starts_with(allowed))
{
return Err(RuntimeError::UnsupportedImport(name.to_owned()));
}
}
Ok(capabilities)
}
fn spawn_actor(
engine: Arc<Engine>,
ticker: Arc<EpochTicker>,
epoch_tick: Duration,
key: ServiceKey,
service: RegisteredService,
init_config: Vec<u8>,
) -> Result<ActorHandle, RuntimeError> {
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) {
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<ActorCommand>, 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<ActorCommand>, 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<HostState>,
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<Vec<u8>, RuntimeError>) -> bool {
matches!(result, Err(RuntimeError::ActorFault { .. }))
}