use miette::IntoDiagnostic;
use std::{io, process::ExitCode, time::Duration};
use tokio::runtime::Handle;
use tosub::Subsystem;
use tracing::{info, level_filters::LevelFilter};
use tracing_subscriber::{EnvFilter, Layer, fmt, layer::SubscriberExt, util::SubscriberInitExt};
#[tokio::main(flavor = "current_thread")]
async fn main() -> miette::Result<ExitCode> {
tracing_subscriber::registry()
.with(
fmt::Layer::new().with_writer(io::stderr).with_filter(
EnvFilter::builder()
.with_default_directive(LevelFilter::INFO.into())
.from_env_lossy(),
),
)
.init();
info!("Root runtime created.");
tosub::build_root("root")
.catch_signals()
.with_timeout(Duration::from_secs(1))
.start(run)
.await?;
info!("Root runtime completed.");
Ok(ExitCode::SUCCESS)
}
async fn run(root: Subsystem) -> miette::Result<()> {
root.spawn("1", run_subsys);
root.spawn("2", run_subsys);
root.spawn("3", run_subsys);
root.shutdown_requested().await;
Ok(())
}
async fn run_subsys(subsys: Subsystem) -> miette::Result<()> {
let s = subsys.clone();
Handle::current().spawn_blocking(move || run_subsys_on_other_runtime(s));
subsys.shutdown_requested().await;
Ok(())
}
fn run_subsys_on_other_runtime(s: Subsystem) -> Result<(), miette::Error> {
info!("Creating new runtime for subsystem '{}'", s.name());
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.into_diagnostic()?;
info!("Runtime for subsystem '{}' created.", s.name());
rt.block_on(do_something_on_other_runtime(s.clone()));
Ok(())
}
async fn do_something_on_other_runtime(subsys: Subsystem) {
info!(
"Hello from subsystem '{}' on runtime '{}'",
subsys.name(),
Handle::current().id()
);
subsys.shutdown_requested().await;
info!("Subsys '{}' stopped.", subsys.name());
}