use crate::{
DaemonConfig,
daemon::event::{DaemonEvent, DaemonEventSender},
hook::DaemonHook,
service::ServiceManager,
};
use ::transport::uds::accept_loop;
use anyhow::Result;
use model::ProviderRegistry;
use std::{
path::{Path, PathBuf},
sync::Arc,
};
use tokio::sync::{RwLock, broadcast, mpsc, oneshot};
use wcore::Runtime;
pub(crate) mod builder;
pub mod event;
mod protocol;
#[derive(Clone)]
pub struct Daemon {
pub runtime: Arc<RwLock<Arc<Runtime<ProviderRegistry, DaemonHook>>>>,
pub(crate) config_dir: PathBuf,
pub(crate) event_tx: DaemonEventSender,
}
impl Daemon {
pub async fn start(config_dir: &Path) -> Result<DaemonHandle> {
let config_path = config_dir.join("walrus.toml");
let config = DaemonConfig::load(&config_path)?;
tracing::info!("loaded configuration from {}", config_path.display());
let (event_tx, event_rx) = mpsc::unbounded_channel::<DaemonEvent>();
let (daemon, service_manager) =
Daemon::build(&config, config_dir, event_tx.clone()).await?;
let (shutdown_tx, _) = broadcast::channel::<()>(1);
let shutdown_event_tx = event_tx.clone();
let mut shutdown_rx = shutdown_tx.subscribe();
tokio::spawn(async move {
let _ = shutdown_rx.recv().await;
let _ = shutdown_event_tx.send(DaemonEvent::Shutdown);
});
for (name, agent) in &config.agents {
if agent.heartbeat.interval == 0 {
continue;
}
let agent_name = name.clone();
let heartbeat_tx = event_tx.clone();
let mut heartbeat_shutdown = shutdown_tx.subscribe();
let interval_secs = agent.heartbeat.interval * 60;
tokio::spawn(async move {
let mut tick = tokio::time::interval(std::time::Duration::from_secs(interval_secs));
tick.tick().await; loop {
tokio::select! {
_ = tick.tick() => {
let event = DaemonEvent::Heartbeat {
agent: agent_name.clone(),
};
if heartbeat_tx.send(event).is_err() {
break;
}
}
_ = heartbeat_shutdown.recv() => break,
}
}
});
tracing::info!(
"heartbeat timer started for '{}' (interval: {}m)",
name,
agent.heartbeat.interval,
);
}
let d = daemon.clone();
let event_loop_join = tokio::spawn(async move {
d.handle_events(event_rx).await;
});
Ok(DaemonHandle {
config,
event_tx,
shutdown_tx,
daemon,
event_loop_join: Some(event_loop_join),
service_manager,
})
}
}
pub struct DaemonHandle {
pub config: DaemonConfig,
pub event_tx: DaemonEventSender,
pub shutdown_tx: broadcast::Sender<()>,
#[allow(unused)]
daemon: Daemon,
event_loop_join: Option<tokio::task::JoinHandle<()>>,
service_manager: Option<ServiceManager>,
}
impl DaemonHandle {
pub async fn wait_until_ready(&self) -> Result<()> {
Ok(())
}
pub async fn shutdown(mut self) -> Result<()> {
if let Some(ref mut sm) = self.service_manager {
sm.shutdown_all().await;
}
let _ = self.shutdown_tx.send(());
if let Some(join) = self.event_loop_join.take() {
join.await?;
}
Ok(())
}
}
pub fn setup_socket(
shutdown_tx: &broadcast::Sender<()>,
event_tx: &DaemonEventSender,
) -> Result<(&'static Path, tokio::task::JoinHandle<()>)> {
let resolved_path: &'static Path = &wcore::paths::SOCKET_PATH;
if let Some(parent) = resolved_path.parent() {
std::fs::create_dir_all(parent)?;
}
if resolved_path.exists() {
std::fs::remove_file(resolved_path)?;
}
let listener = tokio::net::UnixListener::bind(resolved_path)?;
tracing::info!("daemon listening on {}", resolved_path.display());
let socket_shutdown = bridge_shutdown(shutdown_tx.subscribe());
let socket_tx = event_tx.clone();
let join = tokio::spawn(accept_loop(
listener,
move |msg, reply| {
let _ = socket_tx.send(DaemonEvent::Message { msg, reply });
},
socket_shutdown,
));
Ok((resolved_path, join))
}
pub fn setup_tcp(
shutdown_tx: &broadcast::Sender<()>,
event_tx: &DaemonEventSender,
) -> Result<(tokio::task::JoinHandle<()>, u16)> {
let (std_listener, addr) = transport::tcp::bind()?;
let listener = tokio::net::TcpListener::from_std(std_listener)?;
tracing::info!("daemon listening on tcp://{addr}");
let tcp_shutdown = bridge_shutdown(shutdown_tx.subscribe());
let tcp_tx = event_tx.clone();
let join = tokio::spawn(transport::tcp::accept_loop(
listener,
move |msg, reply| {
let _ = tcp_tx.send(DaemonEvent::Message { msg, reply });
},
tcp_shutdown,
));
Ok((join, addr.port()))
}
pub fn bridge_shutdown(mut rx: broadcast::Receiver<()>) -> oneshot::Receiver<()> {
let (otx, orx) = oneshot::channel();
tokio::spawn(async move {
let _ = rx.recv().await;
let _ = otx.send(());
});
orx
}