#![expect(unused_crate_dependencies, reason = "Didn't remove them yet")] use std::{ fs::{self, File, OpenOptions}, io::Write, path::{Path, PathBuf}, time::{Duration, Instant}, }; use eyre::{Context, Result, bail, eyre}; use fs4::fs_std::FileExt; use tokio::time::sleep; use crate::{ aclient::{database::ClientSqlite, record::sqlite_store::SqliteStore, settings::Settings}, api::{ DAEMON_VERSION, server::{control::ControlService, history::HistoryService}, }, daemon::Daemon, }; pub mod aclient; pub mod api; pub(crate) mod daemon; pub(crate) mod events; pub(crate) mod server; /// Boot the daemon. /// /// This creates a daemon, /// starts the gRPC server with services, and runs the event loop. pub async fn boot(settings: Settings, store: SqliteStore, history_db: ClientSqlite) -> Result<()> { let pidfile_path = PathBuf::from(&settings.daemon.pidfile_path); let _pidfile_guard = PidfileGuard::acquire(&pidfile_path)?; let mut daemon = Daemon::builder(settings.clone()) .store(store) .history_db(history_db.clone()) .build()?; let handle = { let handle = daemon.handle(); // Spawn signal handler to emit ShutdownRequested on Ctrl+C/SIGTERM let signal_handle = handle.clone(); tokio::spawn(async move { shutdown_signal().await; tracing::info!("received shutdown signal"); signal_handle.shutdown(); }); handle }; let history_service = HistoryService::new(handle.clone(), history_db).await?; let control_service = ControlService::new(handle.clone()); server::run_grpc_server( &settings, history_service.into_server(), control_service.into_server(), handle, )?; daemon.run_event_loop().await?; tracing::info!("daemon shut down complete"); Ok(()) } /// Wait for a shutdown signal (Ctrl+C or SIGTERM). #[cfg(unix)] async fn shutdown_signal() { let mut term = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()) .expect("failed to register sigterm handler"); let mut int = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::interrupt()) .expect("failed to register sigint handler"); tokio::select! { _ = term.recv() => {}, _ = int.recv() => {}, } } struct PidfileGuard { file: File, } impl PidfileGuard { fn acquire(path: &Path) -> Result { let mut file = open_lock_file(path)?; if !file.try_lock_exclusive()? { bail!( "daemon already running (pidfile lock busy at {})", path.display() ); } file.set_len(0) .wrap_err_with(|| format!("could not truncate daemon pidfile {}", path.display()))?; writeln!(file, "{}", std::process::id()) .and_then(|()| writeln!(file, "{DAEMON_VERSION}")) .wrap_err_with(|| format!("could not write daemon pidfile {}", path.display()))?; Ok(Self { file }) } } impl Drop for PidfileGuard { fn drop(&mut self) { drop(self.file.unlock()); } } fn open_lock_file(path: &Path) -> Result { if let Some(parent) = path.parent() { fs::create_dir_all(parent) .wrap_err_with(|| format!("could not create lock directory {}", parent.display()))?; } OpenOptions::new() .read(true) .write(true) .create(true) .truncate(false) .open(path) .wrap_err_with(|| format!("could not open lock file {}", path.display())) } async fn wait_for_lock(path: &Path, timeout: Duration) -> Result { const LOCK_POLL: Duration = Duration::from_millis(20); let file = open_lock_file(path)?; let start = Instant::now(); loop { match file.try_lock_exclusive() { Ok(true) => return Ok(file), Ok(false) => { if start.elapsed() >= timeout { bail!("timed out waiting for lock at {}", path.display()); } sleep(LOCK_POLL).await; } Err(err) => { return Err(eyre!("could not lock {}: {err}", path.display())); } } } } async fn wait_for_pidfile_available(path: &Path, timeout: Duration) -> Result<()> { let file = wait_for_lock(path, timeout).await?; file.unlock() .wrap_err_with(|| format!("failed to unlock {}", path.display()))?; Ok(()) }