use super::{
error::{
WasmtimePluginRuntimeError, WasmtimePluginRuntimeFailure, WasmtimePluginRuntimeFailureKind,
WasmtimeWorkspaceObserveTaskEnqueueError,
},
execution::{EpochTicker, GuestExecutionGuard},
schedule::{
WasmtimePluginRevocationCleanup, WasmtimePluginScheduledUpdate,
WasmtimePluginUpdateScheduleError, WasmtimeWorkspaceObserveTaskScheduleError,
WasmtimeWorkspaceObserveTaskScheduleErrorParts,
},
store::WasmtimeStoreState,
};
use crate::plugin::{
PluginBufferHandle, PluginEffectBatchLimit, PluginHostImportBatches, PluginHostImportError,
PluginHostImportSession, PluginIdentity, PluginInstanceState, PluginOperationalExport,
PluginRevocationReason, PluginRevocationReport, PluginRuntimeLimitField, PluginRuntimeSpec,
PluginViewHandle, PluginWorkspaceIoBudget, PluginWorkspaceObserveTaskCancellationReport,
PluginWorkspaceObserveTaskEnqueueReport, SealedPluginWorkspaceObserveTaskBatch,
ValidatedPluginRuntimeLimits, abi::bindings::alma::editor::host as wit_host,
};
use std::{
fmt::{Debug, Formatter},
num::NonZeroUsize,
time::Duration,
};
use wasmtime::{
Config, Engine, Store, StoreLimits, StoreLimitsBuilder, WasmBacktraceDetails,
component::{Component, HasSelf, Instance, Linker, Resource},
};
const MAX_COMPONENT_CORE_INSTANCES: usize = 8;
const EPOCH_TICK_INTERVAL: Duration = Duration::from_millis(1);
#[derive(Debug)]
pub struct WasmtimePluginRuntime {
engine: Engine,
epoch_ticker: EpochTicker,
}
impl WasmtimePluginRuntime {
pub fn new() -> Result<Self, WasmtimePluginRuntimeError> {
let mut config = Config::new();
let _config = config.consume_fuel(true);
let _config = config.epoch_interruption(true);
let _config = config.wasm_backtrace(false);
let _config = config.wasm_backtrace_details(WasmBacktraceDetails::Disable);
let engine = Engine::new(&config).map_err(|source| WasmtimePluginRuntimeError::Engine {
failure: WasmtimePluginRuntimeFailure::from_wasmtime(
WasmtimePluginRuntimeFailureKind::Engine,
&source,
),
})?;
let epoch_ticker =
EpochTicker::start(engine.clone(), EPOCH_TICK_INTERVAL).map_err(|source| {
WasmtimePluginRuntimeError::RuntimeTimer {
failure: WasmtimePluginRuntimeFailure::from_io(
WasmtimePluginRuntimeFailureKind::RuntimeTimer,
&source,
),
}
})?;
Ok(Self {
engine,
epoch_ticker,
})
}
pub fn instantiate(
&self,
spec: &PluginRuntimeSpec,
) -> Result<WasmtimePluginInstance, WasmtimePluginRuntimeError> {
let component =
Component::from_binary(&self.engine, spec.component_bytes()).map_err(|source| {
WasmtimePluginRuntimeError::Component {
identity: spec.identity_proof().clone(),
byte_len: spec.component_bytes().len(),
failure: WasmtimePluginRuntimeFailure::from_wasmtime(
WasmtimePluginRuntimeFailureKind::Component,
&source,
),
}
})?;
let limits = spec.validated_limits();
let store_limits = store_limits(spec.identity_proof(), limits)?;
let mut store = Store::new(&self.engine, WasmtimeStoreState::new(store_limits, limits));
store.limiter(|state| &mut state.limits);
let mut linker = Linker::<WasmtimeStoreState>::new(&self.engine);
wit_host::add_to_linker::<_, HasSelf<_>>(&mut linker, |state| state).map_err(|source| {
WasmtimePluginRuntimeError::Instantiate {
identity: spec.identity_proof().clone(),
failure: WasmtimePluginRuntimeFailure::from_wasmtime(
WasmtimePluginRuntimeFailureKind::Linker,
&source,
),
}
})?;
let instance = linker
.instantiate(&mut store, &component)
.map_err(|source| WasmtimePluginRuntimeError::Instantiate {
identity: spec.identity_proof().clone(),
failure: WasmtimePluginRuntimeFailure::from_wasmtime(
WasmtimePluginRuntimeFailureKind::Instantiate,
&source,
),
})?;
let mut instance = WasmtimePluginInstance {
identity: spec.identity_proof().clone(),
limits,
epoch_ticker: self.epoch_ticker.clone(),
store,
instance,
};
instance.call_lifecycle_export(PluginOperationalExport::Load)?;
Ok(instance)
}
}
pub struct WasmtimePluginInstance {
pub(super) identity: PluginIdentity,
pub(super) limits: ValidatedPluginRuntimeLimits,
pub(super) epoch_ticker: EpochTicker,
pub(super) store: Store<WasmtimeStoreState>,
pub(super) instance: Instance,
}
impl WasmtimePluginInstance {
#[must_use]
pub fn identity(&self) -> &str {
self.identity_proof().as_str()
}
#[must_use]
pub const fn identity_proof(&self) -> &PluginIdentity {
&self.identity
}
pub fn update(
&mut self,
state: &PluginInstanceState,
) -> Result<PluginHostImportBatches, WasmtimePluginRuntimeError> {
if state.identity_proof() != &self.identity {
return Err(WasmtimePluginRuntimeError::InstanceStateMismatch {
runtime_identity: self.identity.clone(),
state_identity: state.identity_proof().clone(),
});
}
let active_update = state
.active_update()
.map_err(WasmtimePluginRuntimeError::SchedulingDenied)?;
let active_update = active_update.into_parts();
let workspace_budget = PluginWorkspaceIoBudget::from_max_message_byte_limit(
self.limits.max_message_byte_limit(),
);
let session = PluginHostImportSession::from_snapshot(
active_update.snapshot,
PluginEffectBatchLimit::default(),
workspace_budget,
);
let resources = self
.store
.data_mut()
.begin_update(
session,
active_update.update_handles.view(),
active_update.update_handles.buffer(),
)
.map_err(|source| WasmtimePluginRuntimeError::UpdateSetup {
identity: self.identity.clone(),
source,
})?;
let borrows = resources.borrows();
let result = self.call_update_export(borrows.view, borrows.buffer);
let session = self.store.data_mut().end_update().ok_or_else(|| {
WasmtimePluginRuntimeError::UpdateSetup {
identity: self.identity.clone(),
source: PluginHostImportError::SessionMissing,
}
})?;
let import_error = self.store.data_mut().take_import_rejection();
match (result, import_error) {
(Ok(()), None) => {
self.store
.data_mut()
.commit_workspace_observe_tasks_started_in_update();
Ok(session.seal())
}
(Ok(()), Some(rejection)) => {
let _invalidated = self
.store
.data_mut()
.invalidate_workspace_observe_tasks_started_in_update();
Err(WasmtimePluginRuntimeError::HostImport {
identity: self.identity.clone(),
export: PluginOperationalExport::Update,
import: rejection.import,
source: rejection.source,
}
.with_discard(session.discard()))
}
(Err(source), import_error) => {
let _invalidated = self
.store
.data_mut()
.invalidate_workspace_observe_tasks_started_in_update();
let source = import_error.map_or(source, |rejection| {
WasmtimePluginRuntimeError::HostImport {
identity: self.identity.clone(),
export: PluginOperationalExport::Update,
import: rejection.import,
source: rejection.source,
}
});
Err(source.with_discard(session.discard()))
}
}
}
pub fn update_and_enqueue_workspace_observe_tasks(
&mut self,
state: &PluginInstanceState,
) -> Result<WasmtimePluginScheduledUpdate, WasmtimePluginUpdateScheduleError> {
let update = self.update(state)?.into_owner_and_workspace_observe_tasks();
match self.enqueue_workspace_observe_tasks(state, update.workspace_observe_tasks) {
Ok(scheduled_workspace_observe_tasks) => Ok(WasmtimePluginScheduledUpdate::new(
update.owner_batches,
scheduled_workspace_observe_tasks,
)),
Err(source) => Err(WasmtimeWorkspaceObserveTaskScheduleError::new(
update.owner_batches,
source,
)
.into()),
}
}
pub fn retry_workspace_observe_task_schedule(
&mut self,
state: &PluginInstanceState,
error: WasmtimeWorkspaceObserveTaskScheduleError,
) -> Result<WasmtimePluginScheduledUpdate, WasmtimeWorkspaceObserveTaskScheduleError> {
let WasmtimeWorkspaceObserveTaskScheduleErrorParts {
owner_batches,
enqueue_error,
} = error.into_parts();
match self.enqueue_workspace_observe_tasks(state, enqueue_error.into_batch()) {
Ok(scheduled_workspace_observe_tasks) => Ok(WasmtimePluginScheduledUpdate::new(
owner_batches,
scheduled_workspace_observe_tasks,
)),
Err(source) => Err(WasmtimeWorkspaceObserveTaskScheduleError::new(
owner_batches,
source,
)),
}
}
pub fn enqueue_workspace_observe_tasks(
&mut self,
state: &PluginInstanceState,
batch: SealedPluginWorkspaceObserveTaskBatch,
) -> Result<PluginWorkspaceObserveTaskEnqueueReport, WasmtimeWorkspaceObserveTaskEnqueueError>
{
let batch = self.ensure_active_instance_state_for_task_enqueue(state, batch)?;
self.store
.data_mut()
.workspace_observe_tasks
.enqueue(batch)
.map_err(WasmtimeWorkspaceObserveTaskEnqueueError::from_queue)
}
pub fn execute_all_workspace_observe_tasks(
&mut self,
state: &PluginInstanceState,
filesystem: &crate::fs_utils::FilesystemConfig,
) -> Result<Vec<crate::plugin::PluginWorkspaceObserveTaskCompletion>, WasmtimePluginRuntimeError>
{
self.ensure_active_instance_state(state)?;
Ok(self
.store
.data_mut()
.workspace_observe_tasks
.execute_all(filesystem))
}
pub fn execute_workspace_observe_tasks_up_to(
&mut self,
state: &PluginInstanceState,
filesystem: &crate::fs_utils::FilesystemConfig,
max_tasks: NonZeroUsize,
) -> Result<Vec<crate::plugin::PluginWorkspaceObserveTaskCompletion>, WasmtimePluginRuntimeError>
{
self.ensure_active_instance_state(state)?;
Ok(self
.store
.data_mut()
.workspace_observe_tasks
.execute_up_to(filesystem, max_tasks))
}
pub fn execute_next_workspace_observe_task(
&mut self,
state: &PluginInstanceState,
filesystem: &crate::fs_utils::FilesystemConfig,
) -> Result<
Option<crate::plugin::PluginWorkspaceObserveTaskCompletion>,
WasmtimePluginRuntimeError,
> {
self.ensure_active_instance_state(state)?;
Ok(self
.store
.data_mut()
.workspace_observe_tasks
.execute_next(filesystem))
}
pub fn revoke_state(
&mut self,
state: &mut PluginInstanceState,
reason: PluginRevocationReason,
) -> Result<WasmtimePluginRevocationCleanup, WasmtimePluginRuntimeError> {
if state.identity_proof() != &self.identity {
return Err(WasmtimePluginRuntimeError::InstanceStateMismatch {
runtime_identity: self.identity.clone(),
state_identity: state.identity_proof().clone(),
});
}
let revocation = state.revoke(reason);
let workspace_observe_tasks = self
.store
.data_mut()
.cancel_workspace_observe_tasks(revocation.identity_proof());
Ok(WasmtimePluginRevocationCleanup {
revocation,
workspace_observe_tasks,
})
}
pub fn cancel_workspace_observe_tasks_for_revocation(
&mut self,
revocation: &PluginRevocationReport,
) -> Result<PluginWorkspaceObserveTaskCancellationReport, WasmtimePluginRuntimeError> {
if revocation.identity_proof() != &self.identity {
return Err(WasmtimePluginRuntimeError::InstanceStateMismatch {
runtime_identity: self.identity.clone(),
state_identity: revocation.identity_proof().clone(),
});
}
Ok(self
.store
.data_mut()
.cancel_workspace_observe_tasks(revocation.identity_proof()))
}
#[must_use]
pub fn cancel_workspace_observe_tasks_for_shutdown(
&mut self,
) -> Vec<PluginWorkspaceObserveTaskCancellationReport> {
self.store.data_mut().cancel_all_workspace_observe_tasks()
}
fn ensure_active_instance_state(
&self,
state: &PluginInstanceState,
) -> Result<(), WasmtimePluginRuntimeError> {
if state.identity_proof() != &self.identity {
return Err(WasmtimePluginRuntimeError::InstanceStateMismatch {
runtime_identity: self.identity.clone(),
state_identity: state.identity_proof().clone(),
});
}
state
.host_context()
.map(|_host| ())
.map_err(WasmtimePluginRuntimeError::SchedulingDenied)
}
fn ensure_active_instance_state_for_task_enqueue(
&self,
state: &PluginInstanceState,
batch: SealedPluginWorkspaceObserveTaskBatch,
) -> Result<SealedPluginWorkspaceObserveTaskBatch, WasmtimeWorkspaceObserveTaskEnqueueError>
{
if state.identity_proof() != &self.identity {
return Err(
WasmtimeWorkspaceObserveTaskEnqueueError::instance_state_mismatch(
self.identity.clone(),
state.identity_proof().clone(),
batch,
),
);
}
if batch.identity_proof() != &self.identity {
return Err(
WasmtimeWorkspaceObserveTaskEnqueueError::batch_identity_mismatch(
self.identity.clone(),
batch.identity_proof().clone(),
batch,
),
);
}
match state.host_context() {
Ok(_host) => Ok(batch),
Err(source) => Err(WasmtimeWorkspaceObserveTaskEnqueueError::scheduling_denied(
source, batch,
)),
}
}
fn call_lifecycle_export(
&mut self,
export: PluginOperationalExport,
) -> Result<(), WasmtimePluginRuntimeError> {
let guard = self.arm_guest_export(export)?;
let export_name = export.as_str();
let func = self
.instance
.get_typed_func::<(), ()>(&mut self.store, export_name)
.map_err(|source| WasmtimePluginRuntimeError::Export {
identity: self.identity.clone(),
export,
failure: WasmtimePluginRuntimeFailure::from_wasmtime(
WasmtimePluginRuntimeFailureKind::Export,
&source,
),
})?;
let result = func.call(&mut self.store, ());
guard.ensure_result(
&self.identity,
result,
WasmtimePluginRuntimeFailureKind::GuestTrap,
)?;
let result = func.post_return(&mut self.store);
guard.ensure_result(
&self.identity,
result,
WasmtimePluginRuntimeFailureKind::PostReturn,
)?;
guard.ensure_not_expired(&self.identity)?;
Ok(())
}
fn call_update_export(
&mut self,
view: Resource<PluginViewHandle>,
buffer: Resource<PluginBufferHandle>,
) -> Result<(), WasmtimePluginRuntimeError> {
let export = PluginOperationalExport::Update;
let guard = self.arm_guest_export(export)?;
let func = self
.instance
.get_typed_func::<(Resource<PluginViewHandle>, Resource<PluginBufferHandle>), ()>(
&mut self.store,
export.as_str(),
)
.map_err(|source| WasmtimePluginRuntimeError::Export {
identity: self.identity.clone(),
export,
failure: WasmtimePluginRuntimeFailure::from_wasmtime(
WasmtimePluginRuntimeFailureKind::Export,
&source,
),
})?;
let result = func.call(&mut self.store, (view, buffer));
guard.ensure_result(
&self.identity,
result,
WasmtimePluginRuntimeFailureKind::GuestTrap,
)?;
let result = func.post_return(&mut self.store);
guard.ensure_result(
&self.identity,
result,
WasmtimePluginRuntimeFailureKind::PostReturn,
)?;
guard.ensure_not_expired(&self.identity)?;
Ok(())
}
fn arm_guest_export(
&mut self,
export: PluginOperationalExport,
) -> Result<GuestExecutionGuard, WasmtimePluginRuntimeError> {
self.store
.set_fuel(self.limits.fuel_per_update())
.map_err(|source| WasmtimePluginRuntimeError::Fuel {
identity: self.identity.clone(),
failure: WasmtimePluginRuntimeFailure::from_wasmtime(
WasmtimePluginRuntimeFailureKind::Fuel,
&source,
),
})?;
let guard = GuestExecutionGuard::new(export, self.limits.timeout_ms(), &self.epoch_ticker);
self.store.set_epoch_deadline(guard.epoch_ticks());
self.store.epoch_deadline_trap();
Ok(guard)
}
}
impl Debug for WasmtimePluginInstance {
fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("WasmtimePluginInstance")
.field("identity", &self.identity)
.field("limits", &self.limits)
.finish_non_exhaustive()
}
}
fn store_limits(
identity: &PluginIdentity,
limits: ValidatedPluginRuntimeLimits,
) -> Result<StoreLimits, WasmtimePluginRuntimeError> {
let memory_size = usize::try_from(limits.max_memory_bytes()).map_err(|_source| {
WasmtimePluginRuntimeError::Limit {
identity: identity.clone(),
field: PluginRuntimeLimitField::MaxMemoryBytes,
value: limits.max_memory_bytes(),
}
})?;
Ok(StoreLimitsBuilder::new()
.memory_size(memory_size)
.instances(MAX_COMPONENT_CORE_INSTANCES)
.memories(1)
.tables(8)
.trap_on_grow_failure(true)
.build())
}