use std::{future::Future, sync::Arc};
use crate::core::{RuntimeOwner, SupervisorConfig, SupervisorCore, builder::SupervisorBuilder};
use crate::{error::RuntimeError, subscribers::Subscribe, tasks::TaskSpec};
pub struct Supervisor {
owner: Arc<RuntimeOwner>,
#[cfg(feature = "controller")]
controller: Option<Arc<crate::controller::Controller>>,
}
impl std::fmt::Debug for Supervisor {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Supervisor")
.field("core", self.owner.core())
.finish_non_exhaustive()
}
}
impl Supervisor {
pub(super) fn from_parts(
core: Arc<SupervisorCore>,
#[cfg(feature = "controller")] controller: Option<Arc<crate::controller::Controller>>,
) -> Arc<Self> {
Arc::new(Self {
owner: RuntimeOwner::new(core),
#[cfg(feature = "controller")]
controller,
})
}
pub fn new(cfg: SupervisorConfig, subscribers: Vec<Arc<dyn Subscribe>>) -> Arc<Self> {
Self::builder(cfg).with_subscribers(subscribers).build()
}
pub fn builder(cfg: SupervisorConfig) -> SupervisorBuilder {
SupervisorBuilder::new(cfg)
}
#[must_use = "use the returned runtime handle to manage or shut down the supervisor"]
pub fn serve(&self) -> Result<super::handle::SupervisorHandle, RuntimeError> {
self.owner.core().start()?;
let handle = super::handle::SupervisorHandle::new(Arc::clone(&self.owner));
#[cfg(feature = "controller")]
let handle = handle.with_controller(self.controller.clone());
Ok(handle)
}
pub async fn run(&self, tasks: Vec<TaskSpec>) -> Result<(), RuntimeError> {
self.owner.core().run(tasks).await
}
pub async fn run_until<F>(&self, tasks: Vec<TaskSpec>, shutdown: F) -> Result<(), RuntimeError>
where
F: Future<Output = ()>,
{
self.owner.core().run_until(tasks, shutdown).await
}
pub async fn run_with_os_signals(&self, tasks: Vec<TaskSpec>) -> Result<(), RuntimeError> {
self.owner.core().run_with_os_signals(tasks).await
}
#[must_use = "use the returned runtime configuration"]
pub fn runtime_config(&self) -> &SupervisorConfig {
self.owner.core().runtime_config()
}
#[must_use = "use the returned task defaults"]
pub fn task_defaults(&self) -> &crate::TaskDefaults {
self.owner.core().task_defaults()
}
#[cfg(test)]
pub(crate) fn core(&self) -> &Arc<SupervisorCore> {
self.owner.core()
}
#[cfg(all(test, feature = "controller"))]
pub(crate) fn owner(&self) -> &Arc<RuntimeOwner> {
&self.owner
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::{
pin::pin,
task::{Context, Poll, Waker},
time::Duration,
};
fn poll_ready_without_runtime<F: Future>(future: F) -> F::Output {
let mut future = pin!(future);
let mut context = Context::from_waker(Waker::noop());
match future.as_mut().poll(&mut context) {
Poll::Ready(output) => output,
Poll::Pending => panic!("startup failure must resolve without an active runtime"),
}
}
#[test]
fn serve_without_tokio_returns_typed_error_and_can_retry() {
let supervisor = Supervisor::new(SupervisorConfig::default(), vec![]);
assert!(!supervisor.core().drop_domain().is_started());
assert!(matches!(
supervisor.serve(),
Err(RuntimeError::TokioRuntimeUnavailable)
));
assert!(!supervisor.core().drop_domain().is_started());
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("test runtime");
runtime.block_on(async {
let handle = supervisor.serve().expect("retry inside Tokio must start");
assert!(!supervisor.core().drop_domain().is_started());
handle.shutdown().await.expect("runtime must shut down");
assert!(!supervisor.core().drop_domain().is_started());
});
}
#[test]
fn failed_static_start_does_not_consume_single_shot_lifecycle() {
let supervisor = Supervisor::new(SupervisorConfig::default(), vec![]);
assert!(!supervisor.core().drop_domain().is_started());
for _ in 0..2 {
assert!(matches!(
poll_ready_without_runtime(supervisor.run(vec![])),
Err(RuntimeError::TokioRuntimeUnavailable)
));
assert!(!supervisor.core().drop_domain().is_started());
}
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("test runtime");
runtime.block_on(async {
supervisor
.run(vec![])
.await
.expect("retry inside Tokio must retain the run lifecycle");
assert!(!supervisor.core().drop_domain().is_started());
});
}
#[test]
fn nonempty_static_without_tokio_keeps_destructor_isolation_dormant() {
let supervisor = Supervisor::new(SupervisorConfig::default(), vec![]);
let task = crate::TaskFn::arc(|_ctx| async { Ok(()) });
assert!(matches!(
poll_ready_without_runtime(supervisor.run(vec![TaskSpec::once("outside-tokio", task)])),
Err(RuntimeError::TokioRuntimeUnavailable)
));
assert!(!supervisor.core().drop_domain().is_started());
}
#[tokio::test]
async fn last_public_owner_drop_releases_the_runtime_core() {
let supervisor = Supervisor::new(SupervisorConfig::default(), vec![]);
let weak = Arc::downgrade(supervisor.core());
let handle = supervisor.serve().expect("runtime startup");
drop(supervisor);
assert!(weak.upgrade().is_some(), "the live handle owns the runtime");
drop(handle);
tokio::time::timeout(Duration::from_secs(2), async {
while weak.upgrade().is_some() {
tokio::task::yield_now().await;
}
})
.await
.expect("last-owner Drop must not leave a core ownership cycle");
}
}