unifier-cli 0.2.0

Filesystem postbox for inter-process communication via a Unix tree
Documentation
//! Hot daemon server: owns a [`HotStore`] and serves requests over a Unix socket.

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}"))
}