use crate::daemon::{DaemonCommand, DaemonState};
use crate::sessions::SessionCommand;
use std::io;
use std::path::Path;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicUsize};
use std::sync::mpsc;
use std::thread;
use tracing::debug;
use tracing::info;
use tracing::warn;
pub(crate) struct CoreOptions {
pub acl: Option<Arc<crate::server::acl::SharedAcl>>,
pub config_watchers: bool,
pub auto_exit_wake_path: Option<String>,
}
pub(crate) struct DaemonCore {
pub daemon_tx: mpsc::Sender<DaemonCommand>,
pub shutdown: Arc<AtomicBool>,
pub global_lag: Arc<AtomicUsize>,
pub conn_count: Arc<AtomicUsize>,
pub cmd_handle: thread::JoinHandle<()>,
}
fn spawn_power_event_forwarder(
daemon_tx: mpsc::Sender<DaemonCommand>,
power_rx: crossbeam_channel::Receiver<choreo_power_events::SuspendEvent>,
) {
let _ = thread::Builder::new()
.name("power-events".into())
.spawn(move || {
for event in power_rx.iter() {
debug!(?event, "power event received; forwarding to command loop");
if daemon_tx.send(DaemonCommand::PowerEvent(event)).is_err() {
info!("daemon command loop gone; stopping power-event forwarder");
break;
}
}
});
}
pub(crate) fn start_daemon_core(state: DaemonState, opts: CoreOptions) -> io::Result<DaemonCore> {
let mut state = state;
let (daemon_tx, daemon_rx) = mpsc::channel::<DaemonCommand>();
state.daemon_tx = daemon_tx.clone();
if let Some(acl) = opts.acl {
let acl_path = acl.path().to_path_buf();
state.acl = Some(acl.clone());
if let Some(acl_dir) = acl_path.parent() {
let mut acl_watcher = crate::config_watch::ConfigWatcher::new(acl_dir.to_path_buf());
let acl_rx = acl_watcher.subscribe(
acl_path
.file_name()
.map(|n| n.to_string_lossy().into_owned())
.as_deref()
.unwrap_or("authorized_clients.toml"),
);
acl_watcher.spawn();
crate::server::acl::spawn_acl_watcher(daemon_tx.clone(), acl_rx);
} else {
warn!(
path = %acl_path.display(),
"ACL path has no parent directory; ACL hot-reload disabled"
);
}
}
let overlay_rx = match state.catalog_paths.overlay.parent().map(Path::to_path_buf) {
_ if !opts.config_watchers => {
warn!("config watchers disabled; config-file auto-reload disabled");
crossbeam_channel::never()
}
Some(config_dir) => {
let mut config_watcher = crate::config_watch::ConfigWatcher::new(config_dir);
let overlay_rx = config_watcher.subscribe(crate::catalog::USER_OVERLAY_NAME);
let accounts_rx = config_watcher.subscribe(crate::accounts::ACCOUNTS_TOML_NAME);
config_watcher.spawn();
crate::accounts::spawn_accounts_watcher(daemon_tx.clone(), accounts_rx);
overlay_rx
}
None => {
warn!("config directory not resolvable; config-file auto-reload disabled");
crossbeam_channel::never()
}
};
let maintenance_tx = crate::catalog::spawn_catalog_maintenance(
daemon_tx.clone(),
state.db.clone(),
state.catalog_paths.clone(),
overlay_rx,
);
state.maintenance_tx = Some(maintenance_tx);
let power_monitor = choreo_power_events::PowerMonitor::best_effort();
let power_rx = power_monitor.events().clone();
spawn_power_event_forwarder(daemon_tx.clone(), power_rx);
let shutdown = Arc::new(AtomicBool::new(false));
let global_lag = Arc::clone(&state.global_lag);
let conn_count = Arc::new(AtomicUsize::new(0));
let auto_exit_wake_path = opts.auto_exit_wake_path.clone();
let auto_exit_conn_count = Arc::clone(&conn_count);
let auto_exit_shutdown = Arc::clone(&shutdown);
let cmd_handle = thread::spawn(move || {
loop {
match daemon_rx.recv() {
Ok(DaemonCommand::LastClientDisconnected) => {
if let Some(wake_path) = &auto_exit_wake_path
&& auto_exit_conn_count.load(std::sync::atomic::Ordering::Relaxed) == 0
{
info!("auto-exit: last client disconnected; beginning graceful shutdown");
auto_exit_shutdown.store(true, std::sync::atomic::Ordering::SeqCst);
let _ = choreo_proto::connect_unix(wake_path);
}
}
Ok(DaemonCommand::Shutdown) => {
info!("command loop: shutdown command received; beginning teardown");
break;
}
Ok(cmd) => state.handle_command(cmd),
Err(mpsc::RecvError) => {
info!("command loop: all daemon command senders dropped");
break;
}
}
}
let active_sessions = std::mem::take(&mut state.active_sessions);
info!(
active_sessions = active_sessions.len(),
"command loop teardown: signalling session threads"
);
for entry in active_sessions.values() {
let _ = entry.cmd_tx.send(SessionCommand::Shutdown);
}
let joiners: Vec<_> = active_sessions
.into_iter()
.map(|(session_id, entry)| {
std::thread::spawn(move || {
crate::sessions::join_session_shutdown(entry.handle, session_id)
})
})
.collect();
for joiner in joiners {
let _ = joiner.join();
}
info!("command loop teardown: session threads drained");
state.mcp_manager.shutdown_all();
info!("command loop teardown: complete");
});
Ok(DaemonCore {
daemon_tx,
shutdown,
global_lag,
conn_count,
cmd_handle,
})
}