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;
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::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)?;
while !shutdown.load(Ordering::Relaxed) {
match events_listener.accept() {
Ok((stream, _)) => {
if let Err(e) = hub.add(stream) {
eprintln!("event subscriber error: {e}");
}
}
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}");
}
}
Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => {
thread::sleep(Duration::from_millis(50));
}
Err(e) => return Err(e.into()),
}
}
cleanup(&home)?;
Ok(())
}
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 } => {
let mut store = store.lock().map_err(lock_err)?;
let id = store.post_event(&payload)?;
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![],
})
}
}
}
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));
Ok(())
}
fn lock_err<E: std::fmt::Display>(e: E) -> Error {
Error::msg(format!("daemon store lock poisoned: {e}"))
}