use super::*;
use crate::meerkat_machine_types::{
MeerkatMachineFieldlessRuntimeInternalInput,
canonical_meerkat_machine_runtime_internal_fieldless_input_variant_manifest,
canonical_meerkat_machine_runtime_internal_input_variant_manifest,
};
#[derive(Debug, Clone)]
pub(crate) struct DslTransitionEffects {
effects: Vec<dsl::MeerkatMachineEffect>,
}
#[derive(Debug, Clone)]
pub(crate) struct SessionRegistrationRefusal {
reason: dsl::SessionRegistrationRejectReasonKind,
registered_runtime_epoch_id: Option<dsl::RuntimeEpochId>,
attempted_runtime_epoch_id: Option<dsl::RuntimeEpochId>,
}
impl std::fmt::Display for SessionRegistrationRefusal {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let format_epoch = |epoch: &Option<dsl::RuntimeEpochId>| match epoch {
Some(epoch) => epoch.0.clone(),
None => "none".to_string(),
};
match self.reason {
dsl::SessionRegistrationRejectReasonKind::RuntimeEpochConflict => write!(
f,
"session registration refused: runtime epoch conflict (registered {}, attempted {})",
format_epoch(&self.registered_runtime_epoch_id),
format_epoch(&self.attempted_runtime_epoch_id)
),
}
}
}
impl DslTransitionEffects {
fn new(effects: Vec<dsl::MeerkatMachineEffect>) -> Self {
Self { effects }
}
pub(crate) fn as_slice(&self) -> &[dsl::MeerkatMachineEffect] {
&self.effects
}
pub(crate) fn session_registration_refusal(&self) -> Option<SessionRegistrationRefusal> {
self.effects.iter().find_map(|effect| match effect {
dsl::MeerkatMachineEffect::SessionRegistrationRejected {
reason,
registered_runtime_epoch_id,
attempted_runtime_epoch_id,
..
} => Some(SessionRegistrationRefusal {
reason: *reason,
registered_runtime_epoch_id: registered_runtime_epoch_id.clone(),
attempted_runtime_epoch_id: attempted_runtime_epoch_id.clone(),
}),
_ => None,
})
}
}
pub(crate) fn apply_dsl_transition_on_authority(
authority: &crate::driver::ephemeral::SharedIngressDslAuthority,
input: dsl::MeerkatMachineInput,
context: &str,
) -> Result<DslTransitionEffects, String> {
MeerkatMachine::reject_raw_fieldless_runtime_internal_dsl_input(&input)?;
let mut authority = authority
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
dsl::MeerkatMachineMutator::apply(&mut *authority, input)
.map(|transition| DslTransitionEffects::new(transition.into_effects()))
.map_err(|err| dsl_authority::map_error(err, context))
}
impl std::ops::Deref for DslTransitionEffects {
type Target = [dsl::MeerkatMachineEffect];
fn deref(&self) -> &Self::Target {
self.as_slice()
}
}
impl MeerkatMachine {
pub(super) async fn stage_session_runtime_internal_dsl_transition(
&self,
session_id: &SessionId,
input: MeerkatMachineFieldlessRuntimeInternalInput,
) -> Result<StagedSessionDslInput, String> {
let authority = self.session_dsl_authority(session_id).await?;
Self::stage_runtime_internal_dsl_transition_on_authority(&authority, input)
}
pub(super) fn stage_runtime_internal_dsl_transition_on_authority(
authority: &crate::driver::ephemeral::SharedIngressDslAuthority,
input: MeerkatMachineFieldlessRuntimeInternalInput,
) -> Result<StagedSessionDslInput, String> {
let variant = input.input_variant();
if !canonical_meerkat_machine_runtime_internal_input_variant_manifest().contains(&variant) {
return Err(format!(
"runtime-internal input {variant:?} is absent from the typed production manifest"
));
}
if !canonical_meerkat_machine_runtime_internal_fieldless_input_variant_manifest()
.contains(&variant)
{
return Err(format!(
"runtime-internal input {variant:?} is absent from the typed fieldless manifest"
));
}
if !input.requires_typed_runtime_internal_stager() {
return Err(format!(
"fieldless runtime-internal input {variant:?} is owned by {:?}, not the typed runtime-internal stager",
input.authority()
));
}
Self::stage_dsl_transition_on_authority_after_typed_gate(
authority,
input.dsl_input(),
variant.as_str(),
)
}
pub(super) fn stage_runtime_owner_dsl_transition_on_authority(
authority: &crate::driver::ephemeral::SharedIngressDslAuthority,
input: MeerkatMachineFieldlessRuntimeInternalInput,
) -> Result<StagedSessionDslInput, String> {
let mut authority = authority
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
Self::stage_runtime_owner_dsl_transition_on_locked_authority(&mut authority, input)
}
pub(super) fn stage_runtime_owner_dsl_transition_on_locked_authority(
authority: &mut dsl::MeerkatMachineAuthority,
input: MeerkatMachineFieldlessRuntimeInternalInput,
) -> Result<StagedSessionDslInput, String> {
let variant = input.input_variant();
if !canonical_meerkat_machine_runtime_internal_input_variant_manifest().contains(&variant) {
return Err(format!(
"runtime-owner input {variant:?} is absent from the typed production manifest"
));
}
if !canonical_meerkat_machine_runtime_internal_fieldless_input_variant_manifest()
.contains(&variant)
{
return Err(format!(
"runtime-owner input {variant:?} is absent from the typed fieldless manifest"
));
}
if input.authority()
!= crate::meerkat_machine_types::MeerkatMachineFieldlessRuntimeInternalAuthority::RuntimeOwner
{
return Err(format!(
"fieldless runtime-internal input {variant:?} is owned by {:?}, not RuntimeOwner",
input.authority()
));
}
Self::stage_dsl_transition_on_locked_authority_after_typed_gate(
authority,
input.dsl_input(),
variant.as_str(),
)
}
pub(super) async fn stage_session_dsl_input(
&self,
session_id: &SessionId,
input: dsl::MeerkatMachineInput,
context: &str,
) -> Result<dsl::MeerkatMachineAuthoritySnapshot, String> {
self.stage_session_dsl_transition(session_id, input, context)
.await
.map(|staged| staged.previous_snapshot)
}
pub(super) async fn stage_session_dsl_transition(
&self,
session_id: &SessionId,
input: dsl::MeerkatMachineInput,
context: &str,
) -> Result<StagedSessionDslInput, String> {
let sessions = self.sessions.read().await;
let entry = sessions.get(session_id).ok_or_else(|| {
RuntimeDriverError::NotReady {
state: RuntimeState::Destroyed,
}
.to_string()
})?;
if let Some(error) = entry.dsl_mutation_blocked_by_unregister(session_id) {
return Err(error.to_string());
}
Self::stage_dsl_transition_on_authority(&entry.dsl_authority, input, context)
}
pub(super) fn stage_dsl_transition_on_authority(
authority: &crate::driver::ephemeral::SharedIngressDslAuthority,
input: dsl::MeerkatMachineInput,
context: &str,
) -> Result<StagedSessionDslInput, String> {
Self::reject_raw_fieldless_runtime_internal_dsl_input(&input)?;
Self::stage_dsl_transition_on_authority_after_typed_gate(authority, input, context)
}
pub(super) fn stage_dsl_transition_on_locked_authority(
authority: &mut dsl::MeerkatMachineAuthority,
input: dsl::MeerkatMachineInput,
context: &str,
) -> Result<StagedSessionDslInput, String> {
Self::reject_raw_fieldless_runtime_internal_dsl_input(&input)?;
Self::stage_dsl_transition_on_locked_authority_after_typed_gate(authority, input, context)
}
fn stage_dsl_transition_on_authority_after_typed_gate(
authority: &crate::driver::ephemeral::SharedIngressDslAuthority,
input: dsl::MeerkatMachineInput,
context: &str,
) -> Result<StagedSessionDslInput, String> {
let mut authority = authority
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
Self::stage_dsl_transition_on_locked_authority_after_typed_gate(
&mut authority,
input,
context,
)
}
fn stage_dsl_transition_on_locked_authority_after_typed_gate(
authority: &mut dsl::MeerkatMachineAuthority,
input: dsl::MeerkatMachineInput,
context: &str,
) -> Result<StagedSessionDslInput, String> {
let previous_snapshot = authority.snapshot();
let effects = dsl::MeerkatMachineMutator::apply(authority, input)
.map(|transition| DslTransitionEffects::new(transition.into_effects()))
.map_err(|err| dsl_authority::map_error(err, context))?;
let committed_snapshot = authority.snapshot();
Ok(StagedSessionDslInput {
previous_snapshot,
committed_snapshot,
effects,
})
}
fn resolve_routed_entry_runtime_epoch(
input: &mut dsl::MeerkatMachineInput,
entry_epoch_id: &meerkat_core::RuntimeEpochId,
) {
let declared_epoch = match input {
dsl::MeerkatMachineInput::PrepareBindings {
runtime_epoch_id, ..
}
| dsl::MeerkatMachineInput::Ingest {
runtime_epoch_id, ..
} => runtime_epoch_id,
_ => return,
};
if declared_epoch.is_none() {
*declared_epoch = Some(dsl::RuntimeEpochId::from_domain(entry_epoch_id));
}
}
pub(super) fn reject_raw_fieldless_runtime_internal_dsl_input(
input: &dsl::MeerkatMachineInput,
) -> Result<(), String> {
MeerkatMachineFieldlessRuntimeInternalInput::reject_raw_dsl_input(input)
}
pub(super) async fn apply_session_dsl_input(
&self,
session_id: &SessionId,
input: dsl::MeerkatMachineInput,
context: &str,
) -> Result<(dsl::MeerkatMachineAuthoritySnapshot, DslTransitionEffects), String> {
self.apply_session_dsl_input_with_dispatch_failure(
session_id,
input,
context,
CommittedEffectDispatchFailure::PreserveCommittedDslState,
)
.await
}
pub(super) async fn apply_session_dsl_input_with_dispatch_failure(
&self,
session_id: &SessionId,
input: dsl::MeerkatMachineInput,
context: &str,
dispatch_failure: CommittedEffectDispatchFailure,
) -> Result<(dsl::MeerkatMachineAuthoritySnapshot, DslTransitionEffects), String> {
Self::reject_raw_fieldless_runtime_internal_dsl_input(&input)?;
let sessions = self.sessions.read().await;
let entry = sessions.get(session_id).ok_or_else(|| {
RuntimeDriverError::NotReady {
state: RuntimeState::Destroyed,
}
.to_string()
})?;
if let Some(error) = entry.dsl_mutation_blocked_by_unregister(session_id) {
return Err(error.to_string());
}
let (previous_snapshot, effects) = {
let mut authority = entry
.dsl_authority
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let previous_snapshot = authority.snapshot();
let effects = dsl::MeerkatMachineMutator::apply(&mut *authority, input)
.map(|transition| DslTransitionEffects::new(transition.into_effects()))
.map_err(|err| dsl_authority::map_error(err, context))?;
(previous_snapshot, effects)
};
drop(sessions);
if let Err(error) = self.dispatch_routed_signals_from_effects(&effects).await {
let CommittedEffectDispatchFailure::PreserveCommittedDslState = dispatch_failure;
return Err(format!(
"DSL authority ({context}): committed effect dispatch failed: {error}"
));
}
Ok((previous_snapshot, effects))
}
pub(super) async fn apply_routed_session_dsl_input(
&self,
session_id: &SessionId,
mut input: dsl::MeerkatMachineInput,
context: &str,
) -> Result<
(dsl::MeerkatMachineAuthoritySnapshot, DslTransitionEffects),
dsl_authority::DslTransitionRefusal,
> {
if let Err(reason) = Self::reject_raw_fieldless_runtime_internal_dsl_input(&input) {
return Err(dsl_authority::DslTransitionRefusal::other(
"routed_raw_internal_input_rejected",
reason,
));
}
let sessions = self.sessions.read().await;
let entry = sessions.get(session_id).ok_or_else(|| {
dsl_authority::DslTransitionRefusal::other(
"session_authority_unavailable",
RuntimeDriverError::NotReady {
state: RuntimeState::Destroyed,
}
.to_string(),
)
})?;
if let Some(error) = entry.dsl_mutation_blocked_by_unregister(session_id) {
return Err(dsl_authority::DslTransitionRefusal::other(
"unregister_finalization_pending",
error.to_string(),
));
}
Self::resolve_routed_entry_runtime_epoch(&mut input, &entry.epoch_id);
let (previous_snapshot, effects) = {
let mut authority = entry
.dsl_authority
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let previous_snapshot = authority.snapshot();
let effects = dsl::MeerkatMachineMutator::apply(&mut *authority, input)
.map(|transition| DslTransitionEffects::new(transition.into_effects()))
.map_err(|err| dsl_authority::refusal(err, context))?;
(previous_snapshot, effects)
};
drop(sessions);
if let Err(error) = self.dispatch_routed_signals_from_effects(&effects).await {
return Err(dsl_authority::DslTransitionRefusal::other(
"committed_effect_dispatch_failed",
format!("DSL authority ({context}): committed effect dispatch failed: {error}"),
));
}
Ok((previous_snapshot, effects))
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::BTreeSet;
fn seam_routed_inputs_declaring_runtime_epoch() -> BTreeSet<String> {
let seam = meerkat_machine_schema::meerkat_mob_seam_composition();
let machine = meerkat_machine_schema::catalog::dsl::dsl_meerkat_machine();
let consumer = crate::generated::meerkat_mob_seam::producers::meerkat_instance_id();
seam.routes
.iter()
.filter(|route| route.to.machine == consumer)
.filter_map(|route| match &route.to.input_variant {
meerkat_machine_schema::RouteVariantId::Input(variant) => Some(variant.as_str()),
meerkat_machine_schema::RouteVariantId::Signal(_) => None,
})
.filter(|variant| {
machine.inputs.variants.iter().any(|declared| {
declared.name.as_str() == *variant
&& declared
.fields
.iter()
.any(|field| field.name.as_str() == "runtime_epoch_id")
})
})
.map(str::to_owned)
.collect()
}
fn routed_runtime_epoch(
input: &dsl::MeerkatMachineInput,
) -> Option<&Option<dsl::RuntimeEpochId>> {
match input {
dsl::MeerkatMachineInput::PrepareBindings {
runtime_epoch_id, ..
}
| dsl::MeerkatMachineInput::Ingest {
runtime_epoch_id, ..
} => Some(runtime_epoch_id),
_ => None,
}
}
fn epochless_routed_inputs() -> Vec<(&'static str, dsl::MeerkatMachineInput)> {
vec![
(
"PrepareBindings",
dsl::MeerkatMachineInput::PrepareBindings {
agent_runtime_id: dsl::AgentRuntimeId("mob-runtime".into()),
fence_token: dsl::FenceToken(1),
generation: Some(dsl::Generation(0)),
runtime_epoch_id: None,
session_id: dsl::SessionId("routed-session".into()),
},
),
(
"Ingest",
dsl::MeerkatMachineInput::Ingest {
session_id: dsl::SessionId("routed-session".into()),
runtime_id: dsl::AgentRuntimeId("mob-runtime".into()),
fence_token: dsl::FenceToken(1),
generation: Some(dsl::Generation(0)),
runtime_epoch_id: None,
work_id: dsl::WorkId("work-1".into()),
origin: dsl::WorkOrigin::Ingest,
},
),
]
}
#[test]
fn routed_epoch_resolution_covers_every_seam_input_declaring_the_field() {
let declared = seam_routed_inputs_declaring_runtime_epoch();
let covered: BTreeSet<String> = epochless_routed_inputs()
.iter()
.map(|(name, _)| (*name).to_owned())
.collect();
assert_eq!(
covered,
BTreeSet::from(["Ingest".to_owned(), "PrepareBindings".to_owned()]),
"the routed epoch-bearing fixture set is the gate's own anchor and must stay explicit"
);
assert_eq!(
declared, covered,
"meerkat_mob_seam routes an input declaring `runtime_epoch_id` that this gate does \
not cover; fill it in MeerkatMachine::resolve_routed_entry_runtime_epoch \
(meerkat-runtime/src/meerkat_machine/dsl_effects.rs) and add it here"
);
let entry_epoch_id = meerkat_core::RuntimeEpochId::new();
let expected = Some(dsl::RuntimeEpochId::from_domain(&entry_epoch_id));
for (name, mut input) in epochless_routed_inputs() {
MeerkatMachine::resolve_routed_entry_runtime_epoch(&mut input, &entry_epoch_id);
assert_eq!(
routed_runtime_epoch(&input),
Some(&expected),
"routed `{name}` must carry the entry epoch after resolution"
);
}
}
#[test]
fn routed_epoch_resolution_leaves_a_stated_epoch_alone() {
let stated = Some(dsl::RuntimeEpochId("stated-epoch".into()));
let entry_epoch_id = meerkat_core::RuntimeEpochId::new();
for (name, mut input) in epochless_routed_inputs() {
match &mut input {
dsl::MeerkatMachineInput::PrepareBindings {
runtime_epoch_id, ..
}
| dsl::MeerkatMachineInput::Ingest {
runtime_epoch_id, ..
} => *runtime_epoch_id = stated.clone(),
_ => panic!("routed fixture `{name}` declares no runtime epoch field"),
}
MeerkatMachine::resolve_routed_entry_runtime_epoch(&mut input, &entry_epoch_id);
assert_eq!(
routed_runtime_epoch(&input),
Some(&stated),
"routed `{name}` must keep the epoch its producer stated"
);
}
}
}