use std::fs;
use std::io::{BufRead, BufReader, Write};
use std::os::unix::net::{UnixListener, UnixStream};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::thread;
use std::time::{Duration, Instant};
use crate::daemon::idle::{
idle_after, idle_stop_after, max_load, should_idle_flush, should_idle_stop, system_is_idle,
};
use crate::daemon::notify::{event_name, EventHub, Notice};
use crate::daemon::paths::{daemon_dir, events_socket_path, pid_path, socket_path};
use crate::daemon::protocol::{
decode_request, encode_response, ok_empty, MessageDto, Request, Response,
};
use crate::daemon::www;
use crate::envelope::DEFAULT_SENDER;
use crate::error::{Error, Result};
use crate::home::UnifierHome;
use crate::store::HotStore;
use crate::tick::TickStartOutcome;
pub fn run(home: UnifierHome) -> Result<()> {
let shutdown = Arc::new(AtomicBool::new(false));
home.ensure()?;
fs::create_dir_all(daemon_dir(&home))?;
let events_sock = events_socket_path(&home);
if events_sock.exists() {
fs::remove_file(&events_sock)?;
}
let events_listener = UnixListener::bind(&events_sock)?;
events_listener.set_nonblocking(true)?;
let hub = Arc::new(EventHub::default());
let sock = socket_path(&home);
if sock.exists() {
fs::remove_file(&sock)?;
}
let listener = UnixListener::bind(&sock)?;
listener.set_nonblocking(true)?;
let store = Arc::new(std::sync::Mutex::new(HotStore::load(&home)?));
write_pid(&home)?;
let http_activity = Arc::new(AtomicBool::new(false));
let http_port = match www::spawn(home.clone(), Arc::clone(&shutdown), Arc::clone(&http_activity))
{
Ok(port) => {
eprintln!("www listening on http://127.0.0.1:{port}");
Some(port)
}
Err(e) => {
eprintln!("www server disabled: {e}");
None
}
};
let _ = http_port;
let idle_after = idle_after();
let idle_stop_after = idle_stop_after();
let max_load = max_load();
let mut last_activity = Instant::now();
let mut last_event_gc = Instant::now();
while !shutdown.load(Ordering::Relaxed) {
if !home.path().is_dir() {
let _ = cleanup(&home);
return Ok(());
}
if http_activity.swap(false, Ordering::Relaxed) {
last_activity = Instant::now();
}
match events_listener.accept() {
Ok((stream, _)) => {
if let Err(e) = hub.add(stream) {
eprintln!("event subscriber error: {e}");
}
last_activity = Instant::now();
}
Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => {}
Err(e) => return Err(e.into()),
}
match listener.accept() {
Ok((stream, _)) => {
let home = home.clone();
let store = Arc::clone(&store);
let shutdown = Arc::clone(&shutdown);
let hub = Arc::clone(&hub);
if let Err(e) = handle_client(stream, &home, &store, &shutdown, &hub) {
eprintln!("daemon client error: {e}");
}
last_activity = Instant::now();
}
Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => {
let _ = crate::namespace::reap_dead(&home);
if last_event_gc.elapsed() >= Duration::from_secs(30) {
if let Ok(mut store) = store.lock() {
if let Err(e) = store.expire_events(&home) {
eprintln!("event gc error: {e}");
}
}
last_event_gc = Instant::now();
}
maybe_idle_flush(&home, &store, last_activity, idle_after, max_load);
if maybe_idle_stop(
&home,
&store,
&hub,
&shutdown,
last_activity,
idle_stop_after,
) {
break;
}
thread::sleep(Duration::from_millis(50));
}
Err(e) => return Err(e.into()),
}
}
cleanup(&home)?;
Ok(())
}
fn maybe_idle_flush(
home: &UnifierHome,
store: &Arc<std::sync::Mutex<HotStore>>,
last_activity: Instant,
idle_after: Duration,
max_load: f64,
) {
if !should_idle_flush(
last_activity,
Instant::now(),
idle_after,
system_is_idle(max_load),
true,
false,
) {
return;
}
let Ok(mut store) = store.lock() else {
return;
};
if store.active_tick.is_some() || !store.is_dirty() {
return;
}
if let Err(e) = store.flush(home) {
eprintln!("idle flush error: {e}");
}
}
fn maybe_idle_stop(
home: &UnifierHome,
store: &Arc<std::sync::Mutex<HotStore>>,
hub: &Arc<EventHub>,
shutdown: &Arc<AtomicBool>,
last_activity: Instant,
idle_stop_after: Duration,
) -> bool {
let tick_active = store
.lock()
.map(|s| s.active_tick.is_some())
.unwrap_or(true);
if !should_idle_stop(
last_activity,
Instant::now(),
idle_stop_after,
tick_active,
hub.subscriber_count(),
) {
return false;
}
if let Ok(mut store) = store.lock() {
if store.is_dirty() {
if let Err(e) = store.flush(home) {
eprintln!("idle stop flush error: {e}");
}
}
}
shutdown.store(true, Ordering::Relaxed);
true
}
fn handle_client(
mut stream: UnixStream,
home: &UnifierHome,
store: &Arc<std::sync::Mutex<HotStore>>,
shutdown: &Arc<AtomicBool>,
hub: &Arc<EventHub>,
) -> Result<()> {
let mut reader = BufReader::new(stream.try_clone()?);
let mut line = String::new();
reader.read_line(&mut line)?;
if line.trim().is_empty() {
return Ok(());
}
let request = decode_request(&line).map_err(|e| Error::msg(format!("invalid request: {e}")))?;
let response = match dispatch(home, store, shutdown, hub, request) {
Ok(resp) => resp,
Err(e) => Response::Err {
error: e.to_string(),
},
};
stream.write_all(encode_response(&response)?.as_bytes())?;
stream.flush()?;
Ok(())
}
fn dispatch(
home: &UnifierHome,
store: &Arc<std::sync::Mutex<HotStore>>,
shutdown: &Arc<AtomicBool>,
hub: &Arc<EventHub>,
request: Request,
) -> Result<Response> {
match request {
Request::Ping => Ok(ok_empty()),
Request::Shutdown => {
shutdown.store(true, Ordering::Relaxed);
let mut store = store.lock().map_err(lock_err)?;
if store.is_dirty() && store.active_tick.is_none() {
store.flush(home)?;
} else if store.active_tick.is_some() {
return Err(Error::msg(
"cannot shutdown with active tick; run tick end first",
));
}
Ok(Response::Ok {
value: None,
uuid: None,
found: None,
dirty: Some(false),
tick: None,
queued: None,
messages: vec![],
})
}
Request::Flush => {
let mut store = store.lock().map_err(lock_err)?;
let was_dirty = store.is_dirty();
if store.active_tick.is_some() {
return Err(Error::msg("cannot flush while a tick is active"));
}
if was_dirty {
store.flush(home)?;
}
Ok(Response::Ok {
value: None,
uuid: None,
found: None,
dirty: Some(was_dirty),
tick: None,
queued: None,
messages: vec![],
})
}
Request::Put { key, value } => {
let mut store = store.lock().map_err(lock_err)?;
store.put_key(&key, &value)?;
Ok(ok_empty())
}
Request::Get { key } => {
let store = store.lock().map_err(lock_err)?;
match store.get_key(&key)? {
Some(value) => Ok(Response::Ok {
value: Some(value),
uuid: None,
found: None,
dirty: None,
tick: None,
queued: None,
messages: vec![],
}),
None => Ok(Response::Err {
error: format!("key not found: {key}"),
}),
}
}
Request::Del { key } => {
let mut store = store.lock().map_err(lock_err)?;
if store.delete_key(&key)? {
Ok(ok_empty())
} else {
Ok(Response::Err {
error: format!("key not found: {key}"),
})
}
}
Request::Send {
from,
recipient,
message,
} => {
let from = from.unwrap_or_else(|| DEFAULT_SENDER.to_string());
let mut store = store.lock().map_err(lock_err)?;
let id = store.send_from(&from, &recipient, &message)?;
hub.broadcast(&Notice::mailbox(id, &from, &recipient));
Ok(Response::Ok {
value: None,
uuid: Some(id),
found: None,
dirty: None,
tick: None,
queued: None,
messages: vec![],
})
}
Request::Cron { schedule, message } => {
let mut store = store.lock().map_err(lock_err)?;
let id = store.post_cron(&schedule, &message)?;
Ok(Response::Ok {
value: None,
uuid: Some(id),
found: None,
dirty: None,
tick: None,
queued: None,
messages: vec![],
})
}
Request::Poll { recipient } => {
let store = store.lock().map_err(lock_err)?;
let messages = store.poll_mailbox(&recipient)?;
Ok(messages_response(messages))
}
Request::PollCron => {
let store = store.lock().map_err(lock_err)?;
let messages = store.poll_cron()?;
Ok(messages_response(messages))
}
Request::List { path } => {
let store = store.lock().map_err(lock_err)?;
let messages = store.list_dir(home, &path)?;
Ok(messages_response(messages))
}
Request::Ack { id_or_path } => {
let mut store = store.lock().map_err(lock_err)?;
let found = store.ack(home, &id_or_path)?;
Ok(Response::Ok {
value: None,
uuid: None,
found: Some(found),
dirty: None,
tick: None,
queued: None,
messages: vec![],
})
}
Request::TickStart { label } => {
let mut store = store.lock().map_err(lock_err)?;
match store.tick_start(&label)? {
TickStartOutcome::Started { tick } => Ok(Response::Ok {
value: None,
uuid: None,
found: None,
dirty: None,
tick: Some(tick),
queued: None,
messages: vec![],
}),
TickStartOutcome::Queued { position, .. } => {
eprintln!("tick queue: queued start {label:?} at position {position}");
Ok(Response::Ok {
value: None,
uuid: None,
found: None,
dirty: None,
tick: None,
queued: Some(position),
messages: vec![],
})
}
}
}
Request::TickEnd => {
let mut store = store.lock().map_err(lock_err)?;
let tick = store.tick_end(home)?;
Ok(Response::Ok {
value: None,
uuid: None,
found: None,
dirty: None,
tick: Some(tick),
queued: None,
messages: vec![],
})
}
Request::TickStatus => {
let store = store.lock().map_err(lock_err)?;
let status = store.tick_status();
Ok(Response::Ok {
value: Some(format!(
"committed={} active={:?} queued={} locks={:?}",
status.committed_tick, status.active_tick, status.queued, status.locked_keys
)),
uuid: None,
found: None,
dirty: None,
tick: status.active_tick,
queued: Some(status.queued),
messages: vec![],
})
}
Request::TickLock { key } => {
let mut store = store.lock().map_err(lock_err)?;
store.tick_lock(&key)?;
Ok(ok_empty())
}
Request::TickUnlock { key } => {
let mut store = store.lock().map_err(lock_err)?;
let found = store.tick_unlock(&key)?;
Ok(Response::Ok {
value: None,
uuid: None,
found: Some(found),
dirty: None,
tick: None,
queued: None,
messages: vec![],
})
}
Request::Event { payload, ttl } => {
let mut store = store.lock().map_err(lock_err)?;
let id = store.post_event(&payload, ttl)?;
hub.broadcast(&Notice::event(id, event_name(&payload)));
Ok(Response::Ok {
value: None,
uuid: Some(id),
found: None,
dirty: None,
tick: None,
queued: None,
messages: vec![],
})
}
Request::AgentMessage { from, to, payload } => {
let mut store = store.lock().map_err(lock_err)?;
let id = store.send_agent_message(&from, &to, &payload)?;
hub.broadcast(&Notice::mailbox(id, &from, &to));
Ok(Response::Ok {
value: None,
uuid: Some(id),
found: None,
dirty: None,
tick: None,
queued: None,
messages: vec![],
})
}
Request::WebStatus => {
let value = match www::base_url(home) {
Some(url) => url,
None => return Err(Error::msg("web server is not listening")),
};
Ok(Response::Ok {
value: Some(value),
uuid: None,
found: None,
dirty: None,
tick: None,
queued: None,
messages: vec![],
})
}
Request::WebList => {
let entries = www::list(home)?;
let lines: Vec<String> = entries
.into_iter()
.map(|e| {
let url = www::entry_url(home, &e.name).unwrap_or_default();
format!("{}\t{}\t{}\t{}", e.name, e.content_type, e.bytes, url)
})
.collect();
Ok(Response::Ok {
value: Some(lines.join("\n")),
uuid: None,
found: None,
dirty: None,
tick: None,
queued: None,
messages: vec![],
})
}
Request::WebRm { name } => {
let found = www::remove(home, &name)?;
Ok(Response::Ok {
value: None,
uuid: None,
found: Some(found),
dirty: None,
tick: None,
queued: None,
messages: vec![],
})
}
}
}
fn messages_response(messages: Vec<crate::postbox::Message>) -> Response {
Response::Ok {
value: None,
uuid: None,
found: None,
dirty: None,
tick: None,
queued: None,
messages: messages.into_iter().map(MessageDto::from).collect(),
}
}
fn write_pid(home: &UnifierHome) -> Result<()> {
let pid = std::process::id();
fs::write(pid_path(home), format!("{pid}\n"))?;
Ok(())
}
fn cleanup(home: &UnifierHome) -> Result<()> {
let _ = fs::remove_file(pid_path(home));
let _ = fs::remove_file(socket_path(home));
let _ = fs::remove_file(events_socket_path(home));
let _ = fs::remove_file(crate::daemon::paths::http_port_path(home));
Ok(())
}
fn lock_err<E: std::fmt::Display>(e: E) -> Error {
Error::msg(format!("daemon store lock poisoned: {e}"))
}