use std::path::{Path, PathBuf};
use std::process::Stdio;
use std::time::Duration;
use anyhow::{Context, Result};
use config::Config;
use ipc::SessionState;
use store::Store;
use tokio::process::Command;
use crate::cli::DaemonCmd;
use crate::commands::{broker, load, open_store};
use crate::output;
pub(crate) async fn dispatch(
cmd: DaemonCmd,
device: Option<&str>,
config_path: Option<PathBuf>,
) -> Result<Option<String>> {
match cmd {
DaemonCmd::Start { foreground: true } => {
let cfg = load(config_path)?;
let db = open_store(&cfg).await?;
start_foreground(cfg, device.map(str::to_owned), db).await?;
Ok(None)
}
DaemonCmd::Start { foreground: false } => {
let cfg = load(config_path.clone())?;
Ok(Some(start_background(&cfg, device, config_path).await?))
}
DaemonCmd::Stop => {
let cfg = load(config_path)?;
Ok(Some(stop(&cfg, device).await?))
}
DaemonCmd::Status => {
let cfg = load(config_path)?;
Ok(Some(broker::run_status(&cfg, device, "daemon").await?))
}
DaemonCmd::Install { system } => {
let addr = match device {
Some(d) => d.to_owned(),
None => default_device(config_path.clone(), system)?,
};
Ok(Some(install(&addr, config_path.as_deref(), system)?))
}
DaemonCmd::Uninstall { system } => Ok(Some(uninstall(system)?)),
}
}
fn default_device(config_path: Option<PathBuf>, system: bool) -> Result<String> {
if system && config_path.is_none() {
let home = service::invoking_home().context("resolving config for --system install")?;
let explicit = home.join(".config/imsg/imsg.toml");
return Ok(load(Some(explicit))?.device.address().to_owned());
}
Ok(load(config_path)?.device.address().to_owned())
}
async fn start_foreground(cfg: Config, device: Option<String>, store: Store) -> Result<()> {
store.set_meta("daemon_enabled", "true").await?;
output::line("daemon starting — Ctrl-C to stop")?;
let addr = device.clone().unwrap_or_else(|| cfg.device.address().to_owned());
tokio::spawn(announce_when_connected(addr));
imsg_broker::run_daemon(cfg, device, store).await
}
async fn announce_when_connected(addr: String) {
loop {
match broker::query_state(&addr).await {
Some(SessionState::Active) => {
let _ = output::line(&format!("daemon for {addr}: connected"));
return;
}
Some(SessionState::Failed) => return,
_ => {}
}
tokio::time::sleep(Duration::from_millis(250)).await;
}
}
async fn start_background(
cfg: &Config,
device: Option<&str>,
config_path: Option<PathBuf>,
) -> Result<String> {
let addr = device.unwrap_or_else(|| cfg.device.address()).to_owned();
match broker::query_persistent(&addr).await {
Some(true) => return Ok(format!("daemon for {addr}: already running")),
Some(false) => {
return Err(anyhow::anyhow!(
"an ephemeral broker for {addr} is already using this socket — wait for it to \
idle out (or check `imsg broker status`) before starting the daemon"
));
}
None => {}
}
let log_path = config::daemon_log_path(&addr);
let mut child = spawn_detached(&addr, config_path.as_deref(), &log_path).await?;
broker::connect_retry(
&addr,
&mut child,
&log_path,
cfg.broker.readiness_wait(),
cfg.broker.readiness_poll(),
)
.await?;
Ok(format!("daemon for {addr}: started (log: {})", log_path.display()))
}
async fn spawn_detached(
addr: &str,
config_path: Option<&Path>,
log_path: &Path,
) -> Result<tokio::process::Child> {
let exe = std::env::current_exe().context("resolving current executable path")?;
if let Some(parent) = log_path.parent() {
tokio::fs::create_dir_all(parent).await.context("creating daemon log directory")?;
}
let mut open_opts = tokio::fs::OpenOptions::new();
open_opts.create(true).write(true).truncate(true);
#[cfg(unix)]
open_opts.mode(0o600); let log_file = open_opts
.open(log_path)
.await
.with_context(|| format!("opening daemon log: {}", log_path.display()))?
.into_std()
.await;
let mut cmd = Command::new(exe);
cmd.args(["daemon", "start", "--foreground", "--device", addr]);
if let Some(p) = config_path {
cmd.args(["--config", p.to_str().context("config path is not valid UTF-8")?]);
}
cmd.stdin(Stdio::null());
cmd.stdout(Stdio::from(log_file.try_clone().context("duplicating daemon log handle")?));
cmd.stderr(Stdio::from(log_file));
#[cfg(unix)]
cmd.process_group(0);
cmd.spawn().context("spawning detached daemon subprocess")
}
async fn stop(cfg: &Config, device: Option<&str>) -> Result<String> {
broker::run_stop(cfg, device).await
}
const fn level(system: bool) -> service::ServiceLevel {
if system {
service::ServiceLevel::System
} else {
service::ServiceLevel::User
}
}
fn install(addr: &str, config_path: Option<&Path>, system: bool) -> Result<String> {
let lvl = level(system);
service::install(Some(addr), config_path, lvl).context("installing daemon service")?;
Ok(format!("daemon service installed ({lvl:?})"))
}
fn uninstall(system: bool) -> Result<String> {
let lvl = level(system);
service::uninstall(lvl).context("uninstalling daemon service")?;
Ok(format!("daemon service uninstalled ({lvl:?})"))
}