use std::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,
})
}
#[cfg(feature = "controller")]
fn start_controller(&self) {
if let Some(controller) = &self.controller {
controller.run();
}
}
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) -> super::handle::SupervisorHandle {
self.owner.core().start();
#[cfg(feature = "controller")]
self.start_controller();
let handle = super::handle::SupervisorHandle::new(Arc::clone(&self.owner));
#[cfg(feature = "controller")]
let handle = handle.with_controller(self.controller.clone());
handle
}
pub async fn run(&self, tasks: Vec<TaskSpec>) -> Result<(), RuntimeError> {
#[cfg(feature = "controller")]
self.start_controller();
self.owner.core().run(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::time::Duration;
#[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();
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");
}
}