diff options
| author | Benedikt Peetz <benedikt.peetz@b-peetz.de> | 2026-07-20 19:30:40 +0200 |
|---|---|---|
| committer | Benedikt Peetz <benedikt.peetz@b-peetz.de> | 2026-07-20 19:30:40 +0200 |
| commit | 966a80c4199a49898cc7d8641012d520ce6b2efa (patch) | |
| tree | 51029ff75842090fd1eecbea97b6f7c447e3dea9 /crates/daemon/src/main.rs | |
| parent | chore(server): Remove warnings (diff) | |
| download | atuin-966a80c4199a49898cc7d8641012d520ce6b2efa.zip | |
chore: Commit
Diffstat (limited to '')
| -rw-r--r-- | crates/daemon/src/main.rs | 167 |
1 files changed, 159 insertions, 8 deletions
diff --git a/crates/daemon/src/main.rs b/crates/daemon/src/main.rs index 46d9ad4d..59d4c7ff 100644 --- a/crates/daemon/src/main.rs +++ b/crates/daemon/src/main.rs @@ -1,15 +1,35 @@ -#![expect(unused_crate_dependencies, reason = "Didn't remove them yet")] +#![expect(unused_crate_dependencies)] + +use std::{ + fs::{self, File, OpenOptions}, + io::Write, + path::{Path, PathBuf}, + time::{Duration, Instant}, +}; use clap::Parser; -use eyre::{Result, WrapErr}; -use std::path::PathBuf; -use turtle_daemon::aclient::{ - database::ClientSqlite, record::sqlite_store::SqliteStore, settings::Settings, +use eyre::WrapErr; +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::{control::ControlService, history::HistoryService}, + daemon::Daemon, }; +pub(crate) mod aclient; +pub(crate) mod api; +pub(crate) mod daemon; +pub(crate) mod events; +pub(crate) mod server; + +const DAEMON_VERSION: &str = env!("CARGO_PKG_VERSION"); + #[derive(Parser, Debug)] #[command(infer_subcommands = true)] -pub(crate) enum Cmd { +enum Cmd { /// Start the daemon server Start { /// Also write daemon logs to the console (useful for debugging) @@ -28,8 +48,139 @@ async fn main() -> Result<()> { let sqlite_store = SqliteStore::new(record_store_path, settings.local_timeout).await?; match Cmd::parse() { - Cmd::Start { show_logs, .. } => { - turtle_daemon::boot(settings, sqlite_store, history_db).await + Cmd::Start { show_logs, .. } => boot(settings, sqlite_store, history_db).await, + } +} + +/// Boot the daemon. +/// +/// This creates a daemon, +/// starts the gRPC server with services, and runs the event loop. +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<Self> { + 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<File> { + 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<File> { + 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(()) } |
