#[derive(Debug, thiserror::Error)]
pub enum VMError {
#[error("coroutine {coro_id} faulted: {fault}")]
Fault {
coro_id: usize,
fault: Fault,
},
#[error("max sessions ({max}) exceeded")]
TooManySessions {
max: usize,
},
#[error("max coroutines ({max}) exceeded")]
TooManyCoroutines {
max: usize,
},
#[error("session {0} not found")]
SessionNotFound(SessionId),
#[error("effect handler error: {0}")]
HandlerError(String),
#[error("persistence error: {0}")]
PersistenceError(String),
#[error("invalid concurrency level: {n}")]
InvalidConcurrency {
n: usize,
},
#[error("invalid VM config: {reason}")]
InvalidConfig {
reason: String,
},
#[error("thread pool build failed: {message}")]
ThreadPoolBuild {
message: String,
},
#[error("invalid code image: {reason}")]
InvalidCodeImage {
reason: String,
},
}
pub(crate) enum CoroUpdate {
AdvancePc,
SetPc(PC),
Block(BlockReason),
AdvancePcBlock(BlockReason),
Halt,
AdvancePcWriteReg { reg: u16, val: Value },
}
pub(crate) enum TypeUpdate {
Advance(LocalTypeR),
AdvanceWithOriginal(LocalTypeR, LocalTypeR),
Remove,
}
pub(crate) fn resolve_type_update(
cont: &LocalTypeR,
original: &LocalTypeR,
ep: &Endpoint,
) -> (LocalTypeR, Option<(Endpoint, TypeUpdate)>) {
let (resolved, new_scope) = unfold_if_var_with_scope(cont, original);
let update = if let Some(mu) = new_scope {
Some((
ep.clone(),
TypeUpdate::AdvanceWithOriginal(resolved.clone(), mu),
))
} else {
Some((ep.clone(), TypeUpdate::Advance(resolved.clone())))
};
(resolved, update)
}
pub(crate) struct StepPack {
pub(crate) coro_update: CoroUpdate,
pub(crate) type_update: Option<(Endpoint, TypeUpdate)>,
pub(crate) events: Vec<ObsEvent>,
}
#[derive(Clone, Copy)]
pub(crate) struct GuardAcquireInput<'a> {
pub coro_idx: usize,
pub endpoint: &'a Endpoint,
pub role: &'a str,
pub sid: SessionId,
pub layer: &'a str,
pub dst: u16,
}
#[derive(Clone, Copy)]
pub(crate) struct GuardReleaseInput<'a> {
pub coro_idx: usize,
pub endpoint: &'a Endpoint,
pub role: &'a str,
pub sid: SessionId,
pub layer: &'a str,
pub evidence: u16,
}
pub(crate) enum ExecOutcome {
Continue,
Blocked(BlockReason),
Halted,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct VM<I = (), G = (), P = NoopPersistence, Nu = DefaultVerificationModel>
where
P: PersistenceModel,
{
config: VMConfig,
code: Option<Program>,
programs: Vec<Program>,
identity_model: PhantomData<I>,
guard_model: PhantomData<G>,
persistence_model: PhantomData<P>,
persistent: P::PState,
verification: Nu,
#[serde(default)]
communication_consumption: DefaultCommunicationConsumption,
#[serde(default)]
communication_consumption_artifacts: Vec<CommunicationConsumptionArtifact>,
coroutines: Vec<Coroutine>,
sessions: SessionStore,
arena: Arena,
resource_states: BTreeMap<ScopeId, ResourceState>,
sched: Scheduler,
monitor: SessionMonitor,
obs_trace: Vec<ObsEvent>,
role_symbols: SymbolTable,
label_symbols: SymbolTable,
clock: SimClock,
next_coro_id: usize,
next_session_id: SessionId,
paused_roles: BTreeSet<String>,
guard_layer: InMemoryGuardLayer,
effect_trace: Vec<EffectTraceEntry>,
next_effect_id: u64,
output_condition_checks: Vec<OutputConditionCheck>,
crashed_sites: BTreeSet<SiteId>,
partitioned_edges: BTreeSet<(SiteId, SiteId)>,
corrupted_edges: BTreeMap<(SiteId, SiteId), CorruptionType>,
timed_out_sites: BTreeMap<SiteId, u64>,
last_sched_step: Option<SchedStepDebug>,
handler_identity_anchor: Option<String>,
}
pub type VMState<I = (), G = (), P = NoopPersistence, Nu = DefaultVerificationModel> =
VM<I, G, P, Nu>;
impl<I, G, P, Nu> VM<I, G, P, Nu>
where
P: PersistenceModel,
{
#[must_use]
pub fn new_with_models(config: VMConfig) -> Self
where
P::PState: Default,
Nu: VerificationModel + Default,
{
config.assert_invariants();
let tick_duration = config.tick_duration;
let communication_replay_mode = config.communication_replay_mode;
let sched = Scheduler::new(config.sched_policy.clone());
let mut guard_resources = BTreeMap::new();
for layer in &config.guard_layers {
guard_resources.insert(layer.id.clone(), Value::Unit);
}
Self {
config,
code: None,
programs: Vec::new(),
identity_model: PhantomData,
guard_model: PhantomData,
persistence_model: PhantomData,
persistent: P::PState::default(),
verification: Nu::default(),
communication_consumption: DefaultCommunicationConsumption::new(
communication_replay_mode,
),
communication_consumption_artifacts: Vec::new(),
coroutines: Vec::new(),
sessions: SessionStore::new(),
arena: Arena::default(),
resource_states: BTreeMap::new(),
sched,
monitor: SessionMonitor::default(),
obs_trace: Vec::new(),
role_symbols: SymbolTable::new(),
label_symbols: SymbolTable::new(),
clock: SimClock::new(tick_duration),
next_coro_id: 0,
next_session_id: 0,
paused_roles: BTreeSet::new(),
guard_layer: InMemoryGuardLayer {
resources: guard_resources
.into_iter()
.map(|(k, v)| (LayerId(k), v))
.collect(),
},
effect_trace: Vec::new(),
next_effect_id: 0,
output_condition_checks: Vec::new(),
crashed_sites: BTreeSet::new(),
partitioned_edges: BTreeSet::new(),
corrupted_edges: BTreeMap::new(),
timed_out_sites: BTreeMap::new(),
last_sched_step: None,
handler_identity_anchor: None,
}
}
#[must_use]
pub fn persistent_state(&self) -> &P::PState {
&self.persistent
}
pub fn persistent_state_mut(&mut self) -> &mut P::PState {
&mut self.persistent
}
fn apply_open_delta(&mut self, sid: SessionId) -> Result<(), String> {
let delta = P::open_delta(sid);
P::apply(&mut self.persistent, &delta)
}
fn apply_close_delta(&mut self, sid: SessionId) -> Result<(), String> {
let delta = P::close_delta(sid);
P::apply(&mut self.persistent, &delta)
}
fn apply_invoke_delta(&mut self, sid: SessionId, action: &str) -> Result<(), String> {
if let Some(delta) = P::invoke_delta(sid, action) {
P::apply(&mut self.persistent, &delta)?;
}
Ok(())
}
#[must_use]
pub fn bridge_guard_layer_for_participant<B>(
&self,
bridge: &B,
participant: &I::ParticipantId,
) -> LayerId
where
I: IdentityModel,
G: GuardLayer,
B: IdentityGuardBridge<I, G>,
{
bridge.guard_layer_for_participant(participant)
}
#[must_use]
pub fn bridge_verifying_key_for_participant<B>(
&self,
bridge: &B,
participant: &I::ParticipantId,
) -> Nu::VerifyingKey
where
I: IdentityModel,
Nu: VerificationModel,
B: IdentityVerificationBridge<I, Nu>,
{
bridge.verification_key_for_participant(participant)
}
}