use std::future::Future;
use std::time::Duration;
use crate::participant::api::Participant;
use crate::participant::clock::real::RealClock;
use crate::participant::clock::{ClockReading, ClockSource, TimeUnsynchronized};
use crate::participant::config::ParticipantConfig;
use crate::participant::launch::SupervisedLaunch;
use crate::participant::scheduler::AnyStepScheduler;
use phoxal_bundle::ParticipantClock;
use phoxal_bundle::ParticipantRuntimeInputs;
use phoxal_bus::{BusConfig, BusHandle, BusOwner};
use phoxal_runtime_contract::identity::{ParticipantArtifactId, ParticipantId};
use phoxal_runtime_contract::metadata::ParticipantContract;
use phoxal_runtime_contract::origin::ExecutionOrigin;
use phoxal_runtime_contract::version::FrameworkVersion;
use super::ShutdownController;
use super::inputs::{participant_config, participant_inputs_for_launch};
use super::lifecycle::{self, BusLease};
pub(crate) struct PreparedRun<R: Participant, C: ClockSource> {
pub(crate) bus: BusHandle,
pub(crate) session: BusLease,
pub(crate) participant_id: ParticipantId,
pub(crate) shutdown_grace: Duration,
pub(crate) bundle: Option<ParticipantRuntimeInputs>,
pub(crate) config: R::Config,
pub(crate) clock_mode: ParticipantClock,
pub(crate) clock: Option<C>,
pub(crate) query_reply_delay: Option<Duration>,
}
pub(crate) async fn run_supervised<R, S>(launch: SupervisedLaunch, shutdown: S) -> crate::Result<()>
where
R: Participant,
S: Future<Output = ()>,
{
let mut shutdown = ShutdownController::new(shutdown);
let shutdown_grace = launch.shutdown_grace;
let (bundle, config, clock, clock_mode) = validate_launch::<R>(&launch)?;
tracing::info!(
target: "phoxal.runtime",
endpoints = ?launch.connect_endpoints,
"connecting to the bus"
);
let (owner, bus) = tokio::select! {
biased;
_ = shutdown.wait() => return Ok(()),
result = BusOwner::open(BusConfig::for_participant(
launch.execution_id,
launch.participant_id.clone(),
launch.connect_endpoints.clone(),
)) => result?,
};
lifecycle::run(
PreparedRun::<R, RealClock> {
bus,
session: BusLease::Owned(owner),
participant_id: launch.participant_id,
shutdown_grace,
bundle: Some(bundle),
config,
clock_mode,
clock,
query_reply_delay: None,
},
&mut shutdown,
)
.await
}
pub(crate) fn validate_launch<R>(
launch: &SupervisedLaunch,
) -> crate::Result<(
ParticipantRuntimeInputs,
R::Config,
Option<RealClock>,
ParticipantClock,
)>
where
R: Participant,
{
let bundle = participant_inputs_for_launch(&launch.bundle_root, &launch.participant_id)?;
let compiled_artifact_id = ParticipantArtifactId::new(R::ID).map_err(|error| {
anyhow::anyhow!(
"binary '{}' carries an invalid compiled artifact id: {error}",
R::ID
)
})?;
if bundle.artifact().contract().id != compiled_artifact_id {
anyhow::bail!(
"selected artifact '{}' does not match this binary's compiled artifact id '{}'",
bundle.artifact().contract().id,
R::ID
);
}
if bundle.artifact().contract().kind != R::KIND {
anyhow::bail!(
"selected artifact '{}' has kind {:?}, but this binary declares {:?}",
bundle.artifact().contract().id,
bundle.artifact().contract().kind,
R::KIND
);
}
let process_config_schema: serde_json::Value = serde_json::from_str(R::Config::SCHEMA_JSON)
.map_err(|error| {
anyhow::anyhow!(
"binary '{}' carries an invalid compiled config schema: {error}",
R::ID
)
})?;
let expected_contract = ParticipantContract {
framework: FrameworkVersion::CURRENT,
id: compiled_artifact_id,
kind: R::KIND,
requirement: R::REQUIREMENT,
config_schema: process_config_schema,
};
if bundle.artifact().contract() != &expected_contract {
anyhow::bail!(
"selected artifact '{}' contract does not match this binary's compiled participant contract",
bundle.artifact().contract().id
);
}
let clock_mode = bundle.participant().clock();
let config = participant_config::<R::Config>(bundle.participant().config())?;
let clock = clock_for_mode(clock_mode, launch.execution_origin)?;
validate_clock_inputs::<R, _>(clock_mode, clock.as_ref())?;
Ok((bundle, config, clock, clock_mode))
}
pub(crate) fn clock_for_mode(
clock_mode: ParticipantClock,
execution_origin: Option<ExecutionOrigin>,
) -> crate::Result<Option<RealClock>> {
if clock_mode == ParticipantClock::Real {
let origin = execution_origin.ok_or(TimeUnsynchronized::MissingOrigin)?;
Ok(Some(RealClock::new(origin)?))
} else {
Ok(None)
}
}
pub(crate) fn validate_clock_inputs<R, C>(
clock_mode: ParticipantClock,
clock: Option<&C>,
) -> crate::Result<()>
where
R: Participant,
C: ClockSource,
{
let now = if clock_mode == ParticipantClock::Real {
let reading = clock
.map(ClockSource::read)
.unwrap_or(ClockReading::Unsynchronized(
TimeUnsynchronized::MissingOrigin,
));
match reading {
ClockReading::Synchronized(_) => reading.instant(),
ClockReading::Unsynchronized(reason) => {
return Err(lifecycle::ClockDisciplineLost { reason }.into());
}
}
} else {
None
};
AnyStepScheduler::validate_clock_mode(clock_mode, R::__step_schedule(), now)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn execution_origin_is_required_only_for_real_clock_mode() {
let origin = ExecutionOrigin::mint();
let missing_real = clock_for_mode(ParticipantClock::Real, None)
.expect_err("real mode without an origin must fail before bus startup");
assert!(matches!(
missing_real.downcast_ref::<TimeUnsynchronized>(),
Some(TimeUnsynchronized::MissingOrigin)
));
assert!(
clock_for_mode(ParticipantClock::Real, Some(origin))
.expect("a current-host real origin is valid")
.is_some()
);
assert!(
clock_for_mode(ParticipantClock::Simulation, None)
.expect("simulation mode does not need a host origin")
.is_none()
);
assert!(
clock_for_mode(ParticipantClock::Simulation, Some(origin))
.expect("simulation mode ignores a supplied host origin")
.is_none()
);
assert!(
clock_for_mode(ParticipantClock::Clockless, None)
.expect("clockless mode does not need a host origin")
.is_none()
);
}
}