use std::time::Duration;
use crate::participant::api::Participant;
use crate::participant::managed::ManagedTasks;
#[derive(Clone, Copy, Debug)]
pub(crate) struct ShutdownDeadline(tokio::time::Instant);
impl ShutdownDeadline {
pub(crate) fn from_now(grace: Duration) -> Self {
Self(tokio::time::Instant::now() + grace)
}
pub(crate) fn instant(self) -> tokio::time::Instant {
self.0
}
pub(crate) fn remaining(self) -> Duration {
self.0
.saturating_duration_since(tokio::time::Instant::now())
}
}
#[derive(Debug, Default)]
pub(crate) struct TeardownReport {
pub(crate) shutdown_error: Option<anyhow::Error>,
pub(crate) shutdown_timed_out: bool,
pub(crate) unjoined_tasks: Vec<String>,
pub(crate) task_errors: Vec<anyhow::Error>,
pub(crate) bus_close: Option<phoxal_bus::BusCloseReport>,
}
impl TeardownReport {
pub(crate) fn is_clean(&self) -> bool {
self.shutdown_error.is_none()
&& !self.shutdown_timed_out
&& self.unjoined_tasks.is_empty()
&& self.task_errors.is_empty()
&& self
.bus_close
.as_ref()
.is_none_or(phoxal_bus::BusCloseReport::is_clean)
}
}
impl std::fmt::Display for TeardownReport {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let mut first = true;
let mut item = |label: &str, detail: &dyn std::fmt::Display| -> std::fmt::Result {
if !first {
write!(formatter, "; ")?;
}
first = false;
write!(formatter, "{label}={detail}")
};
if let Some(error) = &self.shutdown_error {
item("shutdown", error)?;
}
if self.shutdown_timed_out {
item("shutdown-timeout", &"true")?;
}
if !self.unjoined_tasks.is_empty() {
item("unjoined-tasks", &format_args!("{:?}", self.unjoined_tasks))?;
}
if !self.task_errors.is_empty() {
let errors: Vec<_> = self.task_errors.iter().map(ToString::to_string).collect();
item("task-errors", &format_args!("{errors:?}"))?;
}
if let Some(report) = &self.bus_close
&& !report.is_clean()
{
item("bus-close-report", report)?;
}
if first {
formatter.write_str("clean")
} else {
Ok(())
}
}
}
#[derive(Debug)]
pub(crate) struct TerminalError {
pub(crate) primary: Option<anyhow::Error>,
pub(crate) teardown: TeardownReport,
}
impl std::fmt::Display for TerminalError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match &self.primary {
Some(primary) => write!(formatter, "{primary}; teardown: {}", self.teardown),
None => write!(formatter, "teardown failed: {}", self.teardown),
}
}
}
impl std::error::Error for TerminalError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
self.primary
.as_deref()
.map(|error| error as _)
.or_else(|| self.teardown.first_error())
}
}
impl TeardownReport {
fn first_error(&self) -> Option<&(dyn std::error::Error + 'static)> {
self.shutdown_error
.as_ref()
.map(|error| error.as_ref() as _)
.or_else(|| self.task_errors.first().map(|error| error.as_ref() as _))
.or_else(|| {
self.bus_close
.as_ref()
.filter(|report| !report.is_clean())
.map(|report| report as &(dyn std::error::Error + 'static))
})
.or_else(|| {
(!self.unjoined_tasks.is_empty())
.then_some(&UNJOINED_TASKS_ERROR as &(dyn std::error::Error + 'static))
})
.or_else(|| {
self.shutdown_timed_out
.then_some(&SHUTDOWN_TIMEOUT_ERROR as &(dyn std::error::Error + 'static))
})
}
}
static SHUTDOWN_TIMEOUT_ERROR: ShutdownTimeoutError = ShutdownTimeoutError;
static UNJOINED_TASKS_ERROR: UnjoinedTasksError = UnjoinedTasksError;
#[derive(Debug)]
struct ShutdownTimeoutError;
impl std::fmt::Display for ShutdownTimeoutError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str("participant shutdown hook exceeded its bounded grace")
}
}
impl std::error::Error for ShutdownTimeoutError {}
#[derive(Debug)]
struct UnjoinedTasksError;
impl std::fmt::Display for UnjoinedTasksError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str("managed tasks remained unjoined after bounded reaping")
}
}
impl std::error::Error for UnjoinedTasksError {}
pub(crate) fn combine<T>(primary: crate::Result<T>, teardown: TeardownReport) -> crate::Result<T> {
if teardown.is_clean() {
return primary;
}
Err(TerminalError {
primary: primary.err(),
teardown,
}
.into())
}
pub(crate) struct Teardown {
pub(crate) managed_tasks: ManagedTasks,
pub(crate) deadline: ShutdownDeadline,
}
impl Teardown {
pub(crate) async fn run<R>(
self,
participant: &R,
api: &R::Api,
state: &mut R::State,
) -> TeardownReport
where
R: Participant,
{
let Teardown {
mut managed_tasks,
deadline,
} = self;
let mut report = TeardownReport::default();
let remaining = deadline.remaining();
match tokio::time::timeout(remaining, participant.shutdown(api, state)).await {
Ok(Ok(())) => {}
Ok(Err(error)) => {
report.shutdown_error = Some(error);
}
Err(_elapsed) => {
report.shutdown_timed_out = true;
}
}
managed_tasks.cancel();
let task_report = managed_tasks.join_until(deadline.instant()).await;
report.unjoined_tasks = task_report.unjoined;
report.task_errors = task_report.failures;
tracing::info!(target: "phoxal.runtime", id = R::ID, "runtime stopped");
report
}
}
pub(crate) async fn abandon_setup(
mut managed_tasks: ManagedTasks,
error: anyhow::Error,
deadline: ShutdownDeadline,
) -> anyhow::Error {
managed_tasks.cancel();
let report = task_report(managed_tasks.join_until(deadline.instant()).await);
if report.is_clean() {
error
} else {
TerminalError {
primary: Some(error),
teardown: report,
}
.into()
}
}
pub(crate) async fn abandon_startup(
mut managed_tasks: ManagedTasks,
deadline: ShutdownDeadline,
) -> TeardownReport {
managed_tasks.cancel();
task_report(managed_tasks.join_until(deadline.instant()).await)
}
fn task_report(shutdown: crate::participant::managed::ManagedTaskShutdown) -> TeardownReport {
TeardownReport {
unjoined_tasks: shutdown.unjoined,
task_errors: shutdown.failures,
..TeardownReport::default()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::participant::managed::ManagedTaskPolicy;
use crate::prelude::*;
use std::error::Error;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
#[derive(Clone, Default)]
struct HookTrace {
called: Arc<AtomicBool>,
completed: Arc<AtomicBool>,
}
#[phoxal::service(id = "hanging-shutdown", state = HookTrace)]
struct HangingShutdown;
impl Participant for HangingShutdown {
async fn setup(
&self,
_ctx: &mut SetupContext<Self>,
_config: Self::Config,
) -> crate::Result<(Self::State, Self::Api)> {
Ok((HookTrace::default(), ()))
}
async fn shutdown(&self, _api: &Self::Api, state: &mut Self::State) -> crate::Result<()> {
state.called.store(true, Ordering::Relaxed);
std::future::pending::<()>().await;
state.completed.store(true, Ordering::Relaxed);
Ok(())
}
}
#[phoxal::service(id = "failing-shutdown", state = HookTrace)]
struct FailingShutdown;
impl Participant for FailingShutdown {
async fn setup(
&self,
_ctx: &mut SetupContext<Self>,
_config: Self::Config,
) -> crate::Result<(Self::State, Self::Api)> {
Ok((HookTrace::default(), ()))
}
async fn shutdown(&self, _api: &Self::Api, state: &mut Self::State) -> crate::Result<()> {
state.called.store(true, Ordering::Relaxed);
anyhow::bail!("could not park the wheels")
}
}
async fn pending_managed_task(name: &str) -> (ManagedTasks, Arc<AtomicBool>) {
let cancelled = Arc::new(AtomicBool::new(false));
let observed = Arc::clone(&cancelled);
let started = Arc::new(AtomicBool::new(false));
let running = Arc::clone(&started);
let mut managed = ManagedTasks::default();
managed.spawn(name, ManagedTaskPolicy::Critical, async move {
struct OnCancel(Arc<AtomicBool>);
impl Drop for OnCancel {
fn drop(&mut self) {
self.0.store(true, Ordering::Relaxed);
}
}
let _guard = OnCancel(observed);
running.store(true, Ordering::Relaxed);
std::future::pending::<()>().await;
});
while !started.load(Ordering::Relaxed) {
tokio::task::yield_now().await;
}
(managed, cancelled)
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_hanging_shutdown_hook_is_bounded_by_the_grace_deadline() {
let trace = HookTrace::default();
let mut state = trace.clone();
let began = std::time::Instant::now();
Teardown {
managed_tasks: ManagedTasks::default(),
deadline: ShutdownDeadline::from_now(Duration::from_millis(150)),
}
.run(&HangingShutdown, &(), &mut state)
.await;
let elapsed = began.elapsed();
assert!(trace.called.load(Ordering::Relaxed), "the hook must run");
assert!(
elapsed < Duration::from_millis(700),
"teardown must return at the grace deadline, took {elapsed:?}"
);
assert!(
!trace.completed.load(Ordering::Relaxed),
"the timed-out hook is dropped, not awaited to completion"
);
}
#[tokio::test]
async fn a_failing_shutdown_hook_does_not_abort_teardown() {
let trace = HookTrace::default();
let (managed_tasks, cancelled) = pending_managed_task("after-a-failing-hook").await;
let mut state = trace.clone();
let report = Teardown {
managed_tasks,
deadline: ShutdownDeadline::from_now(Duration::from_secs(5)),
}
.run(&FailingShutdown, &(), &mut state)
.await;
assert!(trace.called.load(Ordering::Relaxed), "the hook must run");
assert!(
cancelled.load(Ordering::Relaxed),
"the work after the failing hook must still happen"
);
assert_eq!(
report
.shutdown_error
.as_ref()
.map(ToString::to_string)
.as_deref(),
Some("could not park the wheels"),
"cleanup failures must remain structured evidence"
);
}
#[tokio::test]
async fn a_failed_setup_cancels_its_tasks_and_keeps_its_error() {
let (managed, cancelled) = pending_managed_task("spawned-in-setup").await;
let returned = abandon_setup(
managed,
anyhow::anyhow!("the serial port was not there"),
ShutdownDeadline::from_now(Duration::from_secs(5)),
)
.await;
assert!(
cancelled.load(Ordering::Relaxed),
"a task spawned before the failure must not outlive it"
);
assert_eq!(
format!("{returned}"),
"the serial port was not there",
"cleanup must not mask why setup failed"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn managed_tasks_are_cancelled_and_joined_within_the_same_deadline() {
let (managed_tasks, cancelled) = pending_managed_task("sensor-loop").await;
let mut state = HookTrace::default();
let began = std::time::Instant::now();
let report = Teardown {
managed_tasks,
deadline: ShutdownDeadline::from_now(Duration::from_millis(150)),
}
.run(&HangingShutdown, &(), &mut state)
.await;
let elapsed = began.elapsed();
assert!(
cancelled.load(Ordering::Relaxed)
|| report.unjoined_tasks == vec!["sensor-loop".to_string()],
"a task that cannot observe cancellation before the shared deadline must be reported"
);
assert!(
elapsed < Duration::from_millis(700),
"one deadline covers the hook and the joining, took {elapsed:?}"
);
}
#[test]
fn cleanup_failure_changes_a_clean_terminal_result_without_losing_primary_fault() {
let clean = combine::<()>(
Ok(()),
TeardownReport {
shutdown_error: Some(anyhow::anyhow!("park failed").context("shutdown hook")),
..TeardownReport::default()
},
)
.expect_err("cleanup failure must fail an otherwise clean run");
assert!(
format!("{clean}").contains("shutdown=shutdown hook"),
"unexpected cleanup rendering: {clean}"
);
let clean_terminal = clean
.downcast_ref::<TerminalError>()
.expect("cleanup failure must retain the terminal structure");
assert!(
clean_terminal.source().is_some(),
"a clean-exit cleanup failure must expose an Error source"
);
assert!(
clean_terminal
.teardown
.shutdown_error
.as_ref()
.is_some_and(|error| error.chain().count() >= 2),
"cleanup source chains must remain inspectable"
);
let primary = anyhow::anyhow!("step failed").context("step transition");
let combined = combine::<()>(
Err(primary),
TeardownReport {
shutdown_error: Some(anyhow::anyhow!("park failed").context("shutdown hook")),
..TeardownReport::default()
},
)
.expect_err("cleanup evidence must remain attached to a primary fault");
assert!(format!("{combined}").contains("step transition"));
let terminal = combined
.downcast_ref::<TerminalError>()
.expect("primary and cleanup must remain in a TerminalError");
assert_eq!(
terminal.primary.as_ref().map(ToString::to_string),
Some("step transition".to_string())
);
assert!(terminal.teardown.shutdown_error.is_some());
assert!(terminal.source().is_some());
}
#[test]
fn bus_close_report_stays_structured_in_terminal_evidence() {
let result = combine::<()>(
Ok(()),
TeardownReport {
bus_close: Some(phoxal_bus::BusCloseReport {
transport_error_count: 3,
transport_errors: vec!["first failure".to_string()],
transport_errors_truncated: 2,
..phoxal_bus::BusCloseReport::default()
}),
..TeardownReport::default()
},
)
.expect_err("transport close evidence must fail an otherwise clean run");
let terminal = result
.downcast_ref::<TerminalError>()
.expect("close evidence must retain the terminal structure");
let close = terminal
.teardown
.bus_close
.as_ref()
.expect("the structured close report must remain attached");
assert_eq!(close.transport_error_count, 3);
assert_eq!(close.transport_errors, ["first failure"]);
assert!(format!("{result}").contains("3 transport failures"));
}
}