use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader};
use tokio::net::{UnixListener, UnixStream};
use tokio::sync::Semaphore;
use tokio::task::{JoinHandle, JoinSet};
use hallouminate_config::{self, Config};
use super::dispatch::dispatch;
use super::heartbeat::TaskName;
use super::ipc::{DaemonRequest, DaemonResponse};
use super::socket::daemon_socket_path;
use super::state::{DaemonState, WorkClass};
use super::watchdog;
#[derive(Debug, Default, Clone)]
pub struct DaemonArgs {
pub config: Option<PathBuf>,
}
pub const IDLE_READ_TIMEOUT: Duration = Duration::from_secs(30);
const MAX_REQUEST_LINE_BYTES: u64 = 4 * 1024 * 1024;
const MAX_CONCURRENT_CONNECTIONS: usize = 64;
const SHUTDOWN_DRAIN_TIMEOUT: Duration = Duration::from_secs(30);
pub async fn run_daemon(cfg: Config, args: DaemonArgs) -> anyhow::Result<()> {
let xdg_path = args
.config
.clone()
.unwrap_or_else(hallouminate_config::xdg_config_path);
let socket_path = daemon_socket_path()?;
serve_with_config(cfg, Some(xdg_path), &socket_path).await
}
async fn serve_with_config(
cfg: Config,
xdg_path: Option<PathBuf>,
socket_path: &Path,
) -> anyhow::Result<()> {
let now_unix = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
if let watchdog::BootDecision::Backoff {
retry_after_secs,
backoff_secs,
recent_trips,
} = watchdog::check_boot_backoff(
&socket_path.with_file_name("watchdog-trips"),
cfg.daemon.boot_backoff_floor_secs,
cfg.daemon.boot_backoff_cap_secs,
now_unix,
) {
tracing::error!(
target: "hallouminate::daemon",
retry_after_secs,
backoff_secs,
recent_trips,
"watchdog trip backoff active; refusing to start yet",
);
std::process::exit(watchdog::BOOT_BACKOFF_EXIT_CODE);
}
prepare_socket_dir(socket_path).await?;
let lock_path = lock_path_for(socket_path);
let lock = acquire_single_instance(&lock_path)?;
let state = DaemonState::open_with_socket(cfg, xdg_path, socket_path.to_path_buf()).await?;
remove_stale_socket(socket_path).await;
let watcher_enabled = super::watch::spawn_corpus_watcher(&state).is_some();
{
let sup = state.supervisor().clone();
let factory_state = state.clone();
sup.spawn(TaskName::WatcherPump, move || {
let state = factory_state.clone();
async move {
match super::watch::spawn_corpus_watcher(&state) {
Some(handle) => handle.join().await,
None => std::future::pending::<()>().await,
}
}
});
}
spawn_signal_handlers(&state);
spawn_idle_exit(&state, state.baseline().daemon.idle_exit_secs);
{
let sup = state.supervisor().clone();
let factory_state = state.clone();
sup.spawn(TaskName::CatchUp, move || {
super::dispatch::catch_up_index(factory_state.clone())
});
}
spawn_watchdog_when_armed(&state, watcher_enabled, socket_path);
let (result, shutdown_deadline): (anyhow::Result<()>, Instant) =
match serve_on_listener(&state, socket_path, IDLE_READ_TIMEOUT).await {
Ok(deadline) => (Ok(()), deadline),
Err(error) => (Err(error), Instant::now() + SHUTDOWN_DRAIN_TIMEOUT),
};
state.shutdown_token().cancel();
let maintenance = state.take_maintenance_task().await;
finish_shutdown(maintenance, lock, socket_path, shutdown_deadline).await;
result
}
pub fn spawn_signal_handlers(state: &DaemonState) {
let token = state.shutdown_token().clone();
let sigterm = match tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()) {
Ok(s) => s,
Err(e) => {
tracing::warn!(target: "hallouminate::daemon", error = %e, "failed to install SIGTERM handler");
return;
}
};
let sigterm = Arc::new(tokio::sync::Mutex::new(sigterm));
state.supervisor().spawn(TaskName::Signal, move || {
let token = token.clone();
let sigterm = Arc::clone(&sigterm);
async move {
let mut sigterm = sigterm.lock().await;
tokio::select! {
_ = tokio::signal::ctrl_c() => {
tracing::info!(target: "hallouminate::daemon", "received SIGINT; shutting down");
}
_ = sigterm.recv() => {
tracing::info!(target: "hallouminate::daemon", "received SIGTERM; shutting down");
}
}
token.cancel();
}
});
}
fn spawn_idle_exit(state: &DaemonState, idle_exit_secs: u64) {
if idle_exit_secs == 0 {
return;
}
let sup = state.supervisor().clone();
let factory_state = state.clone();
sup.spawn(TaskName::IdleExit, move || {
let state = factory_state.clone();
async move {
let cancel = state.shutdown_token().clone();
loop {
let secs = state.secs_until_idle(idle_exit_secs).max(1);
tokio::select! {
biased;
_ = cancel.cancelled() => break,
_ = sleep_with_idle_heartbeat(&state, Duration::from_secs(secs)) => {}
}
state.heartbeat().bump(TaskName::IdleExit);
if state.should_idle_exit(idle_exit_secs) {
tracing::info!(
target: "hallouminate::daemon",
idle_secs = idle_exit_secs,
"daemon idle-exit; exiting so the OS reclaims all memory",
);
state.shutdown_token().cancel();
break;
}
}
}
});
}
async fn sleep_with_idle_heartbeat(state: &DaemonState, total: Duration) {
const CHUNK: Duration = Duration::from_secs(60);
let mut remaining = total;
while remaining > CHUNK {
tokio::time::sleep(CHUNK).await;
state.heartbeat().bump(TaskName::IdleExit);
remaining -= CHUNK;
}
tokio::time::sleep(remaining).await;
}
fn spawn_watchdog_when_armed(state: &DaemonState, watcher_enabled: bool, socket_path: &Path) {
let daemon = &state.baseline().daemon;
let stall_secs = daemon.watchdog_stall_secs;
let mut candidates: Vec<TaskName> = Vec::new();
if daemon.maintenance_interval_secs != 0 {
candidates.push(TaskName::Maintenance);
}
if watcher_enabled {
candidates.push(TaskName::WatcherPump);
}
if daemon.idle_exit_secs != 0 {
candidates.push(TaskName::IdleExit);
}
if stall_secs == 0 || candidates.is_empty() {
tracing::info!(
target: "hallouminate::daemon",
stall_secs,
candidate_count = candidates.len(),
"watchdog disabled (no stall window or no monitorable tasks)",
);
return;
}
let poll = Duration::from_secs((stall_secs / 4).clamp(1, 60));
let stall = Duration::from_secs(stall_secs);
let state = state.clone();
let trip_path = socket_path.with_file_name("watchdog-trips");
tokio::spawn(async move {
let shutdown = state.shutdown_token().clone();
let heartbeat = state.heartbeat().clone();
loop {
if candidates.iter().all(|task| heartbeat.epoch(*task) > 0) {
break;
}
tokio::select! {
biased;
_ = shutdown.cancelled() => return,
_ = tokio::time::sleep(poll) => {}
}
}
let watchdog = watchdog::Watchdog::spawn(
heartbeat,
candidates,
stall,
poll,
trip_path,
Box::new(|_| std::process::abort()),
);
shutdown.cancelled().await;
watchdog.stop();
});
}
async fn finish_shutdown(
maintenance: Option<JoinHandle<()>>,
lock: std::fs::File,
socket_path: &Path,
deadline: Instant,
) {
if let Some(mut task) = maintenance {
tracing::info!(
target: "hallouminate::daemon",
"draining periodic maintenance before releasing daemon resources",
);
let remaining = deadline.saturating_duration_since(Instant::now());
match tokio::time::timeout(remaining, &mut task).await {
Ok(Ok(())) => {}
Ok(Err(error)) => {
tracing::warn!(
target: "hallouminate::daemon",
error = %error,
"periodic maintenance task exited unexpectedly during shutdown",
);
}
Err(_elapsed) => {
tracing::warn!(
target: "hallouminate::daemon",
timeout_secs = remaining.as_secs_f64(),
"shutdown drain timed out; aborting periodic maintenance",
);
task.abort();
drop(task.await);
}
}
}
cleanup(lock, socket_path).await;
}
async fn cleanup(lock: std::fs::File, socket_path: &Path) {
let _ = tokio::fs::remove_file(socket_path).await;
drop(lock);
}
async fn prepare_socket_dir(socket_path: &Path) -> anyhow::Result<()> {
if let Some(parent) = socket_path.parent()
&& !parent.as_os_str().is_empty()
{
tokio::fs::create_dir_all(parent)
.await
.map_err(|e| anyhow::anyhow!("create socket parent dir {}: {e}", parent.display()))?;
use std::os::unix::fs::PermissionsExt;
let perms = std::fs::Permissions::from_mode(0o700);
tokio::fs::set_permissions(parent, perms)
.await
.map_err(|e| {
anyhow::anyhow!(
"failed to set owner-only permissions (0o700) on socket parent dir {}: {e}",
parent.display(),
)
})?;
}
Ok(())
}
async fn remove_stale_socket(socket_path: &Path) {
if let Err(e) = tokio::fs::remove_file(socket_path).await
&& e.kind() != std::io::ErrorKind::NotFound
{
tracing::warn!(
target: "hallouminate::daemon",
socket = %socket_path.display(),
error = %e,
"failed to remove stale socket before bind; bind may fail with address-in-use",
);
}
}
pub async fn serve(state: &DaemonState, socket_path: &Path) -> anyhow::Result<()> {
serve_with_idle_timeout(state, socket_path, IDLE_READ_TIMEOUT).await
}
pub async fn serve_with_idle_timeout(
state: &DaemonState,
socket_path: &Path,
idle_timeout: Duration,
) -> anyhow::Result<()> {
prepare_socket_dir(socket_path).await?;
let lock_path = lock_path_for(socket_path);
let lock = acquire_single_instance(&lock_path)?;
remove_stale_socket(socket_path).await;
let watcher = super::watch::spawn_corpus_watcher(state);
spawn_idle_exit(state, state.baseline().daemon.idle_exit_secs);
let (result, shutdown_deadline): (anyhow::Result<()>, Instant) =
match serve_on_listener(state, socket_path, idle_timeout).await {
Ok(deadline) => (Ok(()), deadline),
Err(error) => (Err(error), Instant::now() + SHUTDOWN_DRAIN_TIMEOUT),
};
drop(watcher);
state.shutdown_token().cancel();
let maintenance = state.take_maintenance_task().await;
finish_shutdown(maintenance, lock, socket_path, shutdown_deadline).await;
result
}
async fn serve_on_listener(
state: &DaemonState,
socket_path: &Path,
idle_timeout: Duration,
) -> anyhow::Result<Instant> {
let listener = UnixListener::bind(socket_path).map_err(|e| {
tracing::error!(
target: "hallouminate::daemon",
socket = %socket_path.display(),
error = %e,
"failed to bind daemon socket",
);
anyhow::anyhow!("bind {}: {e}", socket_path.display())
})?;
use std::os::unix::fs::PermissionsExt;
let perms = std::fs::Permissions::from_mode(0o600);
tokio::fs::set_permissions(socket_path, perms)
.await
.map_err(|e| {
anyhow::anyhow!(
"failed to set owner-only permissions (0o600) on socket {}: {e}",
socket_path.display(),
)
})?;
tracing::info!(
target: "hallouminate::daemon",
socket = %socket_path.display(),
"daemon listening"
);
let shutdown = state.shutdown_token().clone();
let semaphore = Arc::new(Semaphore::new(MAX_CONCURRENT_CONNECTIONS));
let mut handlers: JoinSet<()> = JoinSet::new();
loop {
let (stream, _addr) = tokio::select! {
_ = shutdown.cancelled() => {
tracing::info!(target: "hallouminate::daemon", "shutdown requested; stopping accept loop");
break;
}
accepted = listener.accept() => match accepted {
Ok(pair) => pair,
Err(e) => {
tracing::warn!(target: "hallouminate::daemon", error = %e, "accept error");
continue;
}
},
};
let permit = tokio::select! {
_ = shutdown.cancelled() => {
tracing::info!(target: "hallouminate::daemon", "shutdown requested; stopping accept loop");
break;
}
acquired = Arc::clone(&semaphore).acquire_owned() => match acquired {
Ok(permit) => permit,
Err(_closed) => break,
},
};
let state = state.clone();
let conn = state.enter_connection(WorkClass::External);
handlers.spawn(async move {
let _conn = conn;
let _permit = permit;
if let Err(e) = handle_connection(state, stream, idle_timeout).await {
tracing::warn!(
target: "hallouminate::daemon",
error = %e,
"connection handler errored"
);
}
});
}
let shutdown_deadline = Instant::now() + SHUTDOWN_DRAIN_TIMEOUT;
drain_handlers(&mut handlers, shutdown_deadline).await;
Ok(shutdown_deadline)
}
async fn drain_handlers(handlers: &mut JoinSet<()>, deadline: Instant) {
if handlers.is_empty() {
return;
}
let pending = handlers.len();
tracing::info!(
target: "hallouminate::daemon",
pending,
"draining in-flight connection handlers before releasing daemon resources",
);
let remaining = deadline.saturating_duration_since(Instant::now());
let drained = tokio::time::timeout(remaining, async {
while handlers.join_next().await.is_some() {}
})
.await;
if drained.is_err() {
tracing::warn!(
target: "hallouminate::daemon",
timeout_secs = remaining.as_secs_f64(),
"shutdown drain timed out; aborting remaining in-flight handlers",
);
handlers.abort_all();
while handlers.join_next().await.is_some() {}
}
}
async fn handle_connection(
state: DaemonState,
stream: UnixStream,
idle_timeout: Duration,
) -> anyhow::Result<()> {
let peer_uid = peer_credential_uid(&stream);
let effective_uid = rustix::process::geteuid().as_raw();
let (read_half, mut write_half) = stream.into_split();
let mut reader = BufReader::new(read_half).take(MAX_REQUEST_LINE_BYTES);
let mut line = String::new();
let n = match tokio::time::timeout(idle_timeout, reader.read_line(&mut line)).await {
Ok(res) => res?,
Err(_) => {
tracing::debug!(
target: "hallouminate::daemon",
timeout_secs = idle_timeout.as_secs_f64(),
"connection idle timeout waiting for request line; closing",
);
return Ok(());
}
};
if n == 0 {
return Ok(());
}
let response = if !line.ends_with('\n') {
tracing::warn!(
target: "hallouminate::daemon",
cap_bytes = MAX_REQUEST_LINE_BYTES,
"request line exceeded the size cap; returning structured error",
);
DaemonResponse::invalid_params(format!(
"request line exceeds {MAX_REQUEST_LINE_BYTES}-byte cap"
))
} else {
match serde_json::from_str::<DaemonRequest>(line.trim_end()) {
Ok(req) => match authorize_peer(peer_uid, effective_uid, &req.payload) {
Some(denied) => denied,
None => dispatch(&state, req).await,
},
Err(e) => DaemonResponse::invalid_params(format!("invalid request: {e}")),
}
};
state.touch_activity(WorkClass::External);
let mut text = serde_json::to_string(&response)?;
text.push('\n');
let write_result = tokio::time::timeout(idle_timeout, async {
write_half.write_all(text.as_bytes()).await?;
write_half.flush().await
})
.await;
match write_result {
Ok(res) => res?,
Err(_) => {
tracing::debug!(
target: "hallouminate::daemon",
timeout_secs = idle_timeout.as_secs_f64(),
"connection idle timeout writing response; closing",
);
}
}
Ok(())
}
fn lock_path_for(socket_path: &Path) -> PathBuf {
let mut s = socket_path.as_os_str().to_os_string();
s.push(".lock");
PathBuf::from(s)
}
fn acquire_single_instance(lock_path: &Path) -> anyhow::Result<std::fs::File> {
use std::fs::OpenOptions;
use std::os::unix::fs::OpenOptionsExt;
use rustix::fs::{FlockOperation, flock};
let file = OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.mode(0o600)
.open(lock_path)
.map_err(|e| anyhow::anyhow!("open lockfile {}: {e}", lock_path.display()))?;
if let Err(errno) = flock(&file, FlockOperation::NonBlockingLockExclusive) {
return Err(anyhow::anyhow!(
"another hallouminate daemon already holds {} ({})",
lock_path.display(),
std::io::Error::from(errno)
));
}
Ok(file)
}
fn authorize_peer(
peer_uid: Option<u32>,
effective_uid: u32,
payload: &super::ipc::DaemonRequestPayload,
) -> Option<DaemonResponse> {
if !is_mutating_payload(payload) {
return None;
}
match peer_uid {
Some(uid) if uid == effective_uid => None,
Some(uid) => Some(DaemonResponse::invalid_params(format!(
"peer uid {uid} is not authorized for mutating requests (daemon uid {effective_uid})"
))),
None => Some(DaemonResponse::invalid_params(
"peer credentials unavailable; refusing mutating request",
)),
}
}
fn is_mutating_payload(payload: &super::ipc::DaemonRequestPayload) -> bool {
use super::ipc::DaemonRequestPayload;
matches!(
payload,
DaemonRequestPayload::AddMarkdown(_)
| DaemonRequestPayload::DeleteMarkdown(_)
| DaemonRequestPayload::Index(_)
| DaemonRequestPayload::Shutdown
)
}
fn peer_credential_uid(stream: &UnixStream) -> Option<u32> {
use std::os::fd::AsRawFd;
raw_peer_uid(stream.as_raw_fd())
}
#[cfg(target_os = "linux")]
fn raw_peer_uid(fd: std::os::fd::RawFd) -> Option<u32> {
let mut cred: libc::ucred = unsafe { std::mem::zeroed() };
let mut len = std::mem::size_of::<libc::ucred>() as libc::socklen_t;
let ret = unsafe {
libc::getsockopt(
fd,
libc::SOL_SOCKET,
libc::SO_PEERCRED,
&mut cred as *mut libc::ucred as *mut libc::c_void,
&mut len,
)
};
if ret == 0 { Some(cred.uid) } else { None }
}
#[cfg(any(
target_os = "macos",
target_os = "ios",
target_os = "freebsd",
target_os = "netbsd",
target_os = "openbsd",
target_os = "dragonfly"
))]
fn raw_peer_uid(fd: std::os::fd::RawFd) -> Option<u32> {
let mut uid: libc::uid_t = 0;
let mut gid: libc::gid_t = 0;
let ret = unsafe { libc::getpeereid(fd, &mut uid, &mut gid) };
if ret == 0 { Some(uid) } else { None }
}
#[cfg(not(any(
target_os = "linux",
target_os = "macos",
target_os = "ios",
target_os = "freebsd",
target_os = "netbsd",
target_os = "openbsd",
target_os = "dragonfly"
)))]
fn raw_peer_uid(_fd: std::os::fd::RawFd) -> Option<u32> {
None
}
#[cfg(test)]
mod tests {
use super::*;
use std::path::PathBuf;
#[test]
fn lock_path_appends_dot_lock_suffix() {
let sock = PathBuf::from("/tmp/hallouminate/daemon.sock");
assert_eq!(
lock_path_for(&sock),
PathBuf::from("/tmp/hallouminate/daemon.sock.lock"),
);
}
#[tokio::test]
async fn remove_stale_socket_tolerates_missing_file() {
let dir = std::env::temp_dir().join(format!("hallouminate-test-{}", std::process::id()));
let missing = dir.join("never-existed.sock");
assert!(!missing.exists());
remove_stale_socket(&missing).await;
assert!(!missing.exists());
}
#[tokio::test]
async fn remove_stale_socket_unlinks_existing_file() {
let dir = std::env::temp_dir().join(format!(
"hallouminate-test-{}-{}",
std::process::id(),
"stale"
));
std::fs::create_dir_all(&dir).expect("create temp dir");
let stale = dir.join("daemon.sock");
std::fs::write(&stale, b"").expect("create stale socket stand-in");
assert!(stale.exists());
remove_stale_socket(&stale).await;
assert!(!stale.exists(), "stale socket must be removed before bind");
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn authorize_peer_allows_same_uid_mutating_request() {
use super::super::ipc::{AddMarkdownRequest, DaemonRequestPayload};
let payload = DaemonRequestPayload::AddMarkdown(AddMarkdownRequest::default());
assert!(
authorize_peer(Some(501), 501, &payload).is_none(),
"same-uid peer must be authorized for a mutating request"
);
}
#[test]
fn authorize_peer_rejects_different_uid_mutating_request() {
use super::super::ipc::{AddMarkdownRequest, DaemonRequestPayload};
let payload = DaemonRequestPayload::AddMarkdown(AddMarkdownRequest::default());
let response = authorize_peer(Some(999), 501, &payload)
.expect("different-uid peer must be rejected for a mutating request");
match response {
DaemonResponse::Err { kind, message } => {
assert_eq!(
kind,
super::super::ipc::ErrorKind::InvalidParams,
"{message}"
);
assert!(
message.contains("999") && message.contains("501"),
"{message}"
);
}
DaemonResponse::Ok { result } => {
panic!("unauthorized mutating request must error; got Ok({result:?})")
}
}
}
#[test]
fn authorize_peer_allows_different_uid_read_only_request() {
use super::super::ipc::DaemonRequestPayload;
assert!(
authorize_peer(Some(999), 501, &DaemonRequestPayload::Ping).is_none(),
"read-only requests must stay unrestricted regardless of peer uid"
);
}
#[test]
fn authorize_peer_fails_closed_when_peer_uid_unknown() {
use super::super::ipc::{AddMarkdownRequest, DaemonRequestPayload};
let payload = DaemonRequestPayload::AddMarkdown(AddMarkdownRequest::default());
assert!(
authorize_peer(None, 501, &payload).is_some(),
"unresolvable peer credentials must fail closed for a mutating request"
);
}
fn acquire_released_lock(lock_path: &Path) -> std::fs::File {
let deadline = Instant::now() + Duration::from_secs(5);
loop {
match acquire_single_instance(lock_path) {
Ok(file) => return file,
Err(e) if Instant::now() >= deadline => {
panic!("replacement daemon never acquired released lock: {e}")
}
Err(_) => std::thread::sleep(Duration::from_millis(1)),
}
}
}
#[tokio::test]
async fn finish_shutdown_drains_maintenance_before_releasing_lock() {
let tmp = tempfile::tempdir().expect("tempdir");
let socket_path = tmp.path().join("daemon.sock");
let lock_path = lock_path_for(&socket_path);
let lock = acquire_single_instance(&lock_path).expect("first daemon lock");
let (release, wait_for_release) = tokio::sync::oneshot::channel::<()>();
let maintenance = tokio::spawn(async move {
wait_for_release.await.expect("release maintenance");
});
let cleanup_socket_path = socket_path.clone();
let shutdown = tokio::spawn(async move {
finish_shutdown(
Some(maintenance),
lock,
&cleanup_socket_path,
Instant::now() + Duration::from_secs(1),
)
.await;
});
tokio::task::yield_now().await;
assert!(
acquire_single_instance(&lock_path).is_err(),
"replacement daemon must not acquire the lock while maintenance is running",
);
release.send(()).expect("finish maintenance");
shutdown.await.expect("finish shutdown");
let replacement = acquire_released_lock(&lock_path);
drop(replacement);
}
#[tokio::test(start_paused = true)]
async fn idle_exit_epoch_advances_at_least_every_60s_during_a_long_sleep() {
let tmp = tempfile::tempdir().expect("tempdir");
let mut cfg = Config::default();
cfg.embeddings.enabled = false;
cfg.storage.ground_dir = tmp.path().to_string_lossy().into_owned();
let state = DaemonState::open(cfg, None).await.expect("open");
let sleep_task = tokio::spawn({
let state = state.clone();
async move {
sleep_with_idle_heartbeat(&state, Duration::from_secs(900)).await;
}
});
tokio::task::yield_now().await;
let mut previous = state.heartbeat().epoch(TaskName::IdleExit);
for _ in 0..14 {
tokio::time::advance(Duration::from_secs(60)).await;
tokio::task::yield_now().await;
let current = state.heartbeat().epoch(TaskName::IdleExit);
assert!(
current > previous,
"IdleExit epoch must advance at least every 60s during a long idle-exit sleep"
);
previous = current;
}
sleep_task.await.expect("idle-exit sleep task");
}
}