use std::future::Future;
use std::time::Duration;
use crate::bundle::RuntimeBundle;
use crate::bus::{BusFault, BusHandle, BusOwner, ParticipantReadyToken};
use crate::identity::ParticipantId;
use crate::participant::api::Participant;
use crate::participant::bus_log::{self, BusLogTask};
use crate::participant::clock::simulation::SimulationClock;
use crate::participant::clock::{ClockMode, ClockReading, ClockSource, TimeUnsynchronized};
use crate::participant::context::{SetupContext, TimelineRetention};
use crate::participant::managed::{ManagedTaskExit, ManagedTaskPolicy, ManagedTasks};
use crate::participant::runtime_performance::{RuntimePerformance, RuntimePerformancePublisher};
use crate::participant::scheduler::simulation::SimulationClockHandle;
use crate::participant::scheduler::{AnyStepScheduler, StepSchedule};
use super::ShutdownController;
use super::query::QuerySurface;
use super::startup::PreparedRun;
use super::teardown::{
ShutdownDeadline, Teardown, TeardownReport, abandon_setup, abandon_startup, combine,
};
#[allow(
dead_code,
reason = "compiled in every profile because a domain module never asks which profile it is in; its only consumer is a module one profile declares"
)]
pub(crate) enum BusLease {
Owned(BusOwner),
Borrowed,
}
#[derive(Debug)]
pub(crate) enum LoopExit {
ShutdownRequested,
ManagedTaskFaulted(ManagedTaskExit),
BusFaulted(BusFault),
ClockDisciplineLost(TimeUnsynchronized),
ResetFailed(anyhow::Error),
StepFailed(anyhow::Error),
QueryDispatchFailed(anyhow::Error),
}
impl LoopExit {
pub(crate) fn into_result(self) -> crate::Result<()> {
match self {
Self::ShutdownRequested => Ok(()),
Self::ManagedTaskFaulted(exit) => Err(ParticipantFault::ManagedTask(exit).into()),
Self::BusFaulted(fault) => Err(ParticipantFault::Bus(fault).into()),
Self::ClockDisciplineLost(reason) => {
Err(ParticipantFault::Clock(ClockDisciplineLost { reason }).into())
}
Self::ResetFailed(error) => Err(ParticipantFault::Reset(error).into()),
Self::StepFailed(error) => Err(ParticipantFault::Step(error).into()),
Self::QueryDispatchFailed(error) => Err(ParticipantFault::Query(error).into()),
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, thiserror::Error)]
#[error("clock discipline lost: {reason}")]
pub(crate) struct ClockDisciplineLost {
pub(crate) reason: TimeUnsynchronized,
}
#[derive(Debug)]
pub(crate) enum ParticipantFault {
ManagedTask(ManagedTaskExit),
Bus(BusFault),
Clock(ClockDisciplineLost),
Reset(anyhow::Error),
Step(anyhow::Error),
Query(anyhow::Error),
}
impl std::fmt::Display for ParticipantFault {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::ManagedTask(error) => error.fmt(formatter),
Self::Bus(error) => write!(formatter, "bus transport failed: {error}"),
Self::Clock(error) => error.fmt(formatter),
Self::Reset(error) => write!(formatter, "reset failed: {error}"),
Self::Step(error) => write!(formatter, "step failed: {error}"),
Self::Query(error) => write!(formatter, "query dispatch failed: {error}"),
}
}
}
impl std::error::Error for ParticipantFault {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match self {
Self::ManagedTask(error) => Some(error),
Self::Bus(error) => Some(error),
Self::Clock(error) => Some(error),
Self::Reset(error) | Self::Step(error) | Self::Query(error) => Some(error.as_ref()),
}
}
}
pub(crate) enum RunnerClock<C: ClockSource> {
Delegated(C),
Simulation(SimulationClock),
}
impl<C: ClockSource> ClockSource for RunnerClock<C> {
fn read(&self) -> ClockReading {
match self {
Self::Delegated(clock) => clock.read(),
Self::Simulation(clock) => clock.read(),
}
}
}
pub(crate) fn runner_clock<C: ClockSource>(
scheduler: &AnyStepScheduler,
clock: Option<C>,
) -> Result<RunnerClock<C>, ClockDisciplineLost> {
match scheduler {
AnyStepScheduler::Simulation(simulation) => {
Ok(RunnerClock::Simulation(simulation.simulation_clock()))
}
AnyStepScheduler::Real(_) | AnyStepScheduler::Disabled => match clock {
Some(clock) => Ok(RunnerClock::Delegated(clock)),
None => Err(ClockDisciplineLost {
reason: TimeUnsynchronized::ClockFault,
}),
},
}
}
pub(crate) struct RunnerTasks {
pub(crate) simulation_clock: Option<SimulationClockHandle>,
pub(crate) bus_log: BusLogTask,
pub(crate) query_reply_delay: Option<Duration>,
}
pub(crate) struct StartInputs<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<RuntimeBundle>,
pub(crate) config: R::Config,
pub(crate) clock: RunnerClock<C>,
pub(crate) scheduler: AnyStepScheduler,
pub(crate) schedule: Option<StepSchedule>,
pub(crate) clock_mode: ClockMode,
pub(crate) tasks: RunnerTasks,
}
pub(crate) enum StartOutcome<T> {
Ready(T),
Terminal {
result: crate::Result<()>,
deadline: ShutdownDeadline,
session: BusLease,
},
}
fn startup_terminal<T>(
primary: crate::Result<()>,
report: TeardownReport,
deadline: ShutdownDeadline,
session: BusLease,
) -> StartOutcome<T> {
StartOutcome::Terminal {
result: combine(primary, report),
deadline,
session,
}
}
async fn startup_teardown<T, R>(
managed_tasks: ManagedTasks,
participant: &R,
api: &R::Api,
state: &mut R::State,
shutdown_grace: Duration,
primary: crate::Result<()>,
session: BusLease,
) -> StartOutcome<T>
where
R: Participant,
{
let deadline = ShutdownDeadline::from_now(shutdown_grace);
let report = Teardown {
managed_tasks,
deadline,
}
.run(participant, api, state)
.await;
startup_terminal(primary, report, deadline, session)
}
pub(crate) async fn run<R, C, S>(
prepared: PreparedRun<R, C>,
shutdown: &mut ShutdownController<S>,
) -> crate::Result<()>
where
R: Participant,
C: ClockSource,
S: Future<Output = ()>,
{
let PreparedRun {
bus,
session,
participant_id,
shutdown_grace,
bundle,
config,
clock_mode,
clock,
query_reply_delay,
} = prepared;
let (bus_logs, bus_log_task) = bus_log::attach(bus.clone());
let schedule = R::__step_schedule();
let now = if clock_mode == ClockMode::Real {
let reading = clock
.as_ref()
.map(ClockSource::read)
.unwrap_or(ClockReading::Unsynchronized(TimeUnsynchronized::ClockFault));
match reading {
ClockReading::Synchronized(_) => reading.instant(),
ClockReading::Unsynchronized(reason) => {
let result = close_session_with_result(
Err(ClockDisciplineLost { reason }.into()),
session,
ShutdownDeadline::from_now(shutdown_grace),
)
.await;
bus_logs.shutdown();
return result;
}
}
} else {
None
};
let (scheduler, clock_handle) =
match AnyStepScheduler::for_clock_mode(clock_mode, schedule, now) {
Ok(value) => value,
Err(error) => {
let result = close_session_with_result(
Err(error),
session,
ShutdownDeadline::from_now(shutdown_grace),
)
.await;
bus_logs.shutdown();
return result;
}
};
let effective_clock = match runner_clock(&scheduler, clock) {
Ok(clock) => clock,
Err(error) => {
let result = close_session_with_result(
Err(error.into()),
session,
ShutdownDeadline::from_now(shutdown_grace),
)
.await;
bus_logs.shutdown();
return result;
}
};
let start = Runner::<R, C>::start(
StartInputs {
bus,
session,
participant_id,
shutdown_grace,
bundle,
config,
clock: effective_clock,
scheduler,
schedule,
clock_mode,
tasks: RunnerTasks {
simulation_clock: clock_handle,
bus_log: bus_log_task,
query_reply_delay,
},
},
shutdown,
)
.await;
let result = match start {
StartOutcome::Ready(runner) => runner.run(shutdown).await,
StartOutcome::Terminal {
result,
deadline,
session,
} => close_session_with_result(result, session, deadline).await,
};
bus_logs.shutdown();
result
}
pub(crate) struct Runner<R: Participant, C: ClockSource> {
pub(crate) participant: R,
pub(crate) api: R::Api,
pub(crate) state: R::State,
pub(crate) bus: BusHandle,
pub(crate) session: BusLease,
pub(crate) clock: RunnerClock<C>,
pub(crate) scheduler: AnyStepScheduler,
pub(crate) schedule: Option<StepSchedule>,
pub(crate) clock_mode: ClockMode,
pub(crate) timeline_retentions: Vec<TimelineRetention>,
pub(crate) queries: Option<QuerySurface<R>>,
pub(crate) runtime_performance_publisher: RuntimePerformancePublisher,
pub(crate) runtime_performance: RuntimePerformance,
pub(crate) managed_tasks: ManagedTasks,
pub(crate) ready: Option<ParticipantReadyToken>,
pub(crate) shutdown_grace: Duration,
}
impl<R: Participant, C: ClockSource> Runner<R, C> {
pub(crate) async fn start<S>(
inputs: StartInputs<R, C>,
shutdown: &mut ShutdownController<S>,
) -> StartOutcome<Self>
where
S: Future<Output = ()>,
{
let StartInputs {
bus,
session,
participant_id,
shutdown_grace,
bundle,
config,
clock,
scheduler,
schedule,
clock_mode,
tasks,
} = inputs;
let mut ctx = SetupContext::<R>::new(bus.clone(), bundle, participant_id.clone());
ctx.spawn_managed_with(
"bus-log-drain",
ManagedTaskPolicy::Finite,
tasks.bus_log.run(),
);
if let Some(handle) = tasks.simulation_clock {
ctx.spawn_managed(
"simulation-clock-ingest",
simulation_clock_feed(bus.clone(), handle),
);
}
let participant = R::__new();
let setup = participant.setup(&mut ctx, config);
let (mut state, api) = match tokio::select! {
biased;
_ = shutdown.wait() => {
let deadline = ShutdownDeadline::from_now(shutdown_grace);
let report = abandon_startup(ctx.take_managed_tasks(), deadline).await;
return startup_terminal(Ok(()), report, deadline, session);
}
fault = bus.wait_for_fatal() => {
let deadline = ShutdownDeadline::from_now(shutdown_grace);
let error = abandon_setup(
ctx.take_managed_tasks(),
ParticipantFault::Bus(fault).into(),
deadline,
).await;
return StartOutcome::Terminal { result: Err(error), deadline, session };
}
result = setup => result,
} {
Ok(pair) => pair,
Err(error) => {
let deadline = ShutdownDeadline::from_now(shutdown_grace);
let error = abandon_setup(ctx.take_managed_tasks(), error, deadline).await;
return StartOutcome::Terminal {
result: Err(error),
deadline,
session,
};
}
};
let mut managed_tasks = ctx.take_managed_tasks();
let timeline_retentions = ctx.take_timeline_retentions();
let query_registrations = ctx.take_query_registrations();
let query_reply_delay = tasks.query_reply_delay;
let mut queries = match tokio::select! {
biased;
_ = shutdown.wait() => {
return startup_teardown(
managed_tasks,
&participant,
&api,
&mut state,
shutdown_grace,
Ok(()),
session,
).await;
}
fault = bus.wait_for_fatal() => {
return startup_teardown(
managed_tasks,
&participant,
&api,
&mut state,
shutdown_grace,
Err(ParticipantFault::Bus(fault).into()),
session,
).await;
}
result = QuerySurface::declare(
&bus,
query_registrations,
&mut managed_tasks,
query_reply_delay,
) => result,
} {
Ok(queries) => queries,
Err(error) => {
return startup_teardown(
managed_tasks,
&participant,
&api,
&mut state,
shutdown_grace,
Err(error),
session,
)
.await;
}
};
if let Some(exit) = managed_tasks.try_next_unexpected_exit() {
if let Some(queries) = queries.take() {
queries.close();
}
return startup_teardown(
managed_tasks,
&participant,
&api,
&mut state,
shutdown_grace,
Err(exit.into()),
session,
)
.await;
}
if shutdown.is_requested() {
if let Some(queries) = queries.take() {
queries.close();
}
return startup_teardown(
managed_tasks,
&participant,
&api,
&mut state,
shutdown_grace,
Ok(()),
session,
)
.await;
}
if let crate::bus::BusTerminal::Fatal(fault) = bus.terminal() {
if let Some(queries) = queries.take() {
queries.close();
}
return startup_teardown(
managed_tasks,
&participant,
&api,
&mut state,
shutdown_grace,
Err(ParticipantFault::Bus(fault).into()),
session,
)
.await;
}
let ready = match &session {
BusLease::Borrowed => {
if shutdown.is_requested() {
if let Some(queries) = queries.take() {
queries.close();
}
return startup_teardown(
managed_tasks,
&participant,
&api,
&mut state,
shutdown_grace,
Ok(()),
session,
)
.await;
}
None
}
BusLease::Owned(owner) => Some(tokio::select! {
biased;
_ = shutdown.wait() => {
if let Some(queries) = queries.take() {
queries.close();
}
return startup_teardown(
managed_tasks,
&participant,
&api,
&mut state,
shutdown_grace,
Ok(()),
session,
).await;
}
fault = bus.wait_for_fatal() => {
if let Some(queries) = queries.take() {
queries.close();
}
return startup_teardown(
managed_tasks,
&participant,
&api,
&mut state,
shutdown_grace,
Err(ParticipantFault::Bus(fault).into()),
session,
).await;
}
exit = managed_tasks.next_unexpected_exit() => {
if let Some(queries) = queries.take() {
queries.close();
}
return startup_teardown(
managed_tasks,
&participant,
&api,
&mut state,
shutdown_grace,
Err(exit.into()),
session,
).await;
}
result = owner.declare_participant_ready() => match result {
Ok(token) => token,
Err(error) => {
if let Some(queries) = queries.take() {
queries.close();
}
return startup_teardown(
managed_tasks,
&participant,
&api,
&mut state,
shutdown_grace,
Err(error.into()),
session,
).await;
}
},
}),
};
tracing::info!(
target: "phoxal.runtime",
id = R::ID,
participant = %participant_id,
"runtime ready"
);
StartOutcome::Ready(Self {
participant,
api,
state,
bus: bus.clone(),
session,
clock,
scheduler,
schedule,
clock_mode,
timeline_retentions,
queries,
runtime_performance_publisher: RuntimePerformancePublisher::attach(bus),
runtime_performance: RuntimePerformance::new(schedule),
managed_tasks,
ready,
shutdown_grace,
})
}
pub(crate) async fn run<S>(mut self, shutdown: &mut ShutdownController<S>) -> crate::Result<()>
where
S: Future<Output = ()>,
{
let exit = self.main_loop(shutdown).await;
let primary = exit.into_result();
let report = self.finish().await;
combine(primary, report)
}
async fn finish(self) -> TeardownReport {
let Self {
participant,
api,
mut state,
queries,
managed_tasks,
session,
ready,
shutdown_grace,
..
} = self;
drop(ready);
if let Some(queries) = queries {
queries.close();
}
let deadline = ShutdownDeadline::from_now(shutdown_grace);
let mut report = Teardown {
managed_tasks,
deadline,
}
.run(&participant, &api, &mut state)
.await;
let close_report = close_session(session, deadline).await;
report.bus_close = close_report.bus_close;
report
}
}
async fn close_session(session: BusLease, deadline: ShutdownDeadline) -> TeardownReport {
let BusLease::Owned(owner) = session else {
return TeardownReport::default();
};
let close = owner.close_until(deadline.instant()).await;
if close.is_clean() {
TeardownReport::default()
} else {
TeardownReport {
bus_close: Some(close),
..TeardownReport::default()
}
}
}
pub(crate) async fn close_session_with_result<T>(
primary: crate::Result<T>,
session: BusLease,
deadline: ShutdownDeadline,
) -> crate::Result<T> {
combine(primary, close_session(session, deadline).await)
}
async fn simulation_clock_feed(bus: BusHandle, handle: SimulationClockHandle) -> crate::Result<()> {
super::event_loop::simulation_clock_feed(bus, handle).await
}