unifier-cli 0.3.0

Filesystem postbox for inter-process communication via a Unix tree
Documentation
#![cfg(unix)]

use assert_cmd::Command;
use predicates::prelude::*;
use std::path::PathBuf;
use std::thread;
use std::time::Duration;
use unifier::daemon::{run_server, Client, Request};
use unifier::home::UnifierHome;

struct DaemonGuard {
    _tmp: tempfile::TempDir,
    home: UnifierHome,
    handle: Option<thread::JoinHandle<()>>,
}

impl DaemonGuard {
    fn start() -> Self {
        let tmp = tempfile::tempdir().unwrap();
        let home = UnifierHome::resolve(Some(tmp.path().to_path_buf()), None).unwrap();
        let run_home = home.clone();
        let handle = thread::spawn(move || {
            let _ = run_server(run_home);
        });
        wait_for_socket(home.path());
        Self {
            _tmp: tmp,
            home,
            handle: Some(handle),
        }
    }

    fn home_path(&self) -> PathBuf {
        self.home.path().to_path_buf()
    }

    fn home_arg(&self) -> String {
        self.home.path().display().to_string()
    }
}

impl Drop for DaemonGuard {
    fn drop(&mut self) {
        if Client::is_running(&self.home) {
            if let Ok(mut client) = Client::connect(&self.home) {
                let _ = client.request(Request::Shutdown);
            }
        }
        if let Some(handle) = self.handle.take() {
            let _ = handle.join();
        }
    }
}

fn wait_for_socket(home: &std::path::Path) {
    let sock = home.join(".daemon/unifier.sock");
    let events = home.join(".daemon/events.sock");
    for _ in 0..100 {
        if sock.exists() && events.exists() {
            return;
        }
        thread::sleep(Duration::from_millis(50));
    }
    panic!("daemon sockets did not appear");
}

#[test]
fn hot_put_get_without_flush_leaves_disk_clean() {
    let daemon = DaemonGuard::start();
    let home = daemon.home_path();
    let home_s = daemon.home_arg();

    Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "put", "app/theme", "dark"])
        .assert()
        .success();

    Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "get", "app/theme"])
        .assert()
        .success()
        .stdout("dark\n");

    assert!(!home.join("keys/app/theme").exists());

    Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "daemon", "flush"])
        .assert()
        .success()
        .stdout(predicate::str::contains("flushed"));

    assert!(home.join("keys/app/theme").is_file());
}

#[test]
fn daemon_status_reports_running() {
    let daemon = DaemonGuard::start();
    let home_s = daemon.home_arg();

    Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "daemon", "status"])
        .assert()
        .success()
        .stdout(predicate::str::contains("running"));
}

#[test]
fn no_daemon_flag_uses_filesystem() {
    let daemon = DaemonGuard::start();
    let home = daemon.home_path();
    let home_s = daemon.home_arg();

    Command::cargo_bin("unifier")
        .unwrap()
        .args([
            "--home",
            &home_s,
            "--no-daemon",
            "put",
            "direct/key",
            "value",
        ])
        .assert()
        .success();

    assert!(home.join("keys/direct/key").is_file());
}

#[test]
fn stop_flushes_dirty_state() {
    let daemon = DaemonGuard::start();
    let home = daemon.home_path();
    let home_s = daemon.home_arg();

    Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "put", "persist/me", "yes"])
        .assert()
        .success();

    assert!(!home.join("keys/persist/me").exists());

    Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "daemon", "stop"])
        .assert()
        .success();

    assert!(home.join("keys/persist/me").is_file());

    Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "daemon", "status"])
        .assert()
        .success()
        .stdout(predicate::str::contains("stopped"));
}

#[test]
fn mailbox_message_notifies_event_socket() {
    use std::io::{BufRead, BufReader};
    use std::os::unix::net::UnixStream;
    use unifier::daemon::{events_socket_path, Notice};
    use unifier::envelope::Envelope;

    let daemon = DaemonGuard::start();
    let events_path = events_socket_path(&daemon.home);
    let listener = UnixStream::connect(&events_path).expect("connect events.sock");
    listener
        .set_read_timeout(Some(Duration::from_secs(2)))
        .unwrap();
    thread::sleep(Duration::from_millis(120));

    let home_s = daemon.home_arg();
    let assert = Command::cargo_bin("unifier")
        .unwrap()
        .args([
            "--home",
            &home_s,
            "message",
            "--from",
            "myagent",
            "youragent",
            r#"{"hello":"world"}"#,
        ])
        .assert()
        .success();
    let id = String::from_utf8(assert.get_output().stdout.clone())
        .unwrap()
        .trim()
        .to_string();

    let mut reader = BufReader::new(listener);
    let mut line = String::new();
    reader.read_line(&mut line).expect("read wakeup notice");
    let notice = Notice::from_json(&line).unwrap();
    match notice {
        Notice::Mailbox { id: nid, from, to } => {
            assert_eq!(nid.hyphenated().to_string(), id);
            assert_eq!(from, "myagent");
            assert_eq!(to, "youragent");
        }
        other => panic!("expected mailbox notice, got {other:?}"),
    }

    let poll = Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "poll", "youragent"])
        .assert()
        .success();
    let body = String::from_utf8(poll.get_output().stdout.clone()).unwrap();
    let env = Envelope::from_json(body.lines().nth(1).unwrap()).unwrap();
    assert_eq!(env.from, "myagent");
    assert_eq!(env.to, "youragent");
    assert_eq!(env.id.hyphenated().to_string(), id);
    assert_eq!(env.payload["hello"], "world");
}

#[test]
fn gradle_style_auto_starts_daemon() {
    let tmp = tempfile::tempdir().unwrap();
    let home_s = tmp.path().display().to_string();

    Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "put", "auto/start", "yes"])
        .assert()
        .success();

    Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "daemon", "status"])
        .assert()
        .success()
        .stdout(predicate::str::contains("running"));

    Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "daemon", "stop"])
        .assert()
        .success();
}

#[test]
fn tick_staging_via_daemon() {
    let daemon = DaemonGuard::start();
    let home_s = daemon.home_arg();

    Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "tick", "start", "job1"])
        .assert()
        .success()
        .stdout(predicate::str::contains("tick 1 started"));

    Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "put", "state/x", "draft"])
        .assert()
        .success();

    assert!(!daemon.home_path().join("keys/state/x").exists());

    Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "tick", "end"])
        .assert()
        .success()
        .stdout(predicate::str::contains("tick 1 committed"));

    assert!(daemon.home_path().join("ticks/1/meta.json").is_file());
}

#[test]
fn namespace_prefixes_keys_via_daemon() {
    let daemon = DaemonGuard::start();
    let home_s = daemon.home_arg();

    Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "namespace", "set", "agent"])
        .assert()
        .success()
        .stdout("agent\n");

    Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "put", "app/theme", "dark"])
        .assert()
        .success();

    Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "get", "app/theme"])
        .assert()
        .success()
        .stdout("dark\n");

    assert!(!daemon.home_path().join("keys/app/theme").exists());

    Command::cargo_bin("unifier")
        .unwrap()
        .args(["--home", &home_s, "daemon", "flush"])
        .assert()
        .success();

    assert!(daemon.home_path().join("keys/agent/app/theme").is_file());
    assert!(!daemon.home_path().join("keys/app/theme").exists());
}