#![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");
let tick = home.join(".daemon/tick.sock");
let http = home.join(".daemon/http.port");
for _ in 0..100 {
if sock.exists() && events.exists() && tick.exists() && http.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 missing_home_stops_daemon() {
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());
assert!(Client::is_running(&home));
drop(tmp);
for _ in 0..100 {
if handle.is_finished() {
let _ = handle.join();
return;
}
thread::sleep(Duration::from_millis(50));
}
panic!("daemon did not exit after home directory was removed");
}
#[test]
fn idle_stop_exits_without_subscribers() {
std::env::set_var("UNIFIER_IDLE_STOP_SECS", "1");
std::env::set_var("UNIFIER_IDLE_FLUSH_SECS", "0");
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());
for _ in 0..80 {
if handle.is_finished() {
let _ = handle.join();
std::env::remove_var("UNIFIER_IDLE_STOP_SECS");
std::env::remove_var("UNIFIER_IDLE_FLUSH_SECS");
return;
}
thread::sleep(Duration::from_millis(50));
}
std::env::remove_var("UNIFIER_IDLE_STOP_SECS");
std::env::remove_var("UNIFIER_IDLE_FLUSH_SECS");
let _ = Client::connect(&home).and_then(|mut c| c.request(Request::Shutdown));
let _ = handle.join();
panic!("daemon did not idle-stop within timeout");
}
#[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());
}
#[test]
fn serve_pipes_html_and_http_fetches_it() {
use std::io::{Read, Write};
use std::net::TcpStream;
let daemon = DaemonGuard::start();
let home_s = daemon.home_arg();
let assert = Command::cargo_bin("unifier")
.unwrap()
.args([
"--home",
&home_s,
"serve",
"jan-report",
"--wrap",
"--title",
"Nightly",
])
.write_stdin("<p>hello from jan</p>")
.assert()
.success();
let url = String::from_utf8(assert.get_output().stdout.clone())
.unwrap()
.trim()
.to_string();
assert!(url.starts_with("http://127.0.0.1:"), "{url}");
assert!(url.ends_with("/jan-report"), "{url}");
Command::cargo_bin("unifier")
.unwrap()
.args(["--home", &home_s, "web", "list"])
.assert()
.success()
.stdout(predicate::str::contains("jan-report"));
Command::cargo_bin("unifier")
.unwrap()
.args(["--home", &home_s, "daemon", "status"])
.assert()
.success()
.stdout(predicate::str::contains("www: http://127.0.0.1:"));
let host_port = url.trim_start_matches("http://").split('/').next().unwrap();
let mut stream = TcpStream::connect(host_port).unwrap();
stream
.write_all(b"GET /jan-report HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n")
.unwrap();
let mut resp = String::new();
stream.read_to_string(&mut resp).unwrap();
assert!(resp.contains("HTTP/1.1 200"), "{resp}");
assert!(resp.contains("Unifier"), "{resp}");
assert!(resp.contains("hello from jan"), "{resp}");
assert!(resp.contains("Nightly"), "{resp}");
let mut stream = TcpStream::connect(host_port).unwrap();
stream
.write_all(b"GET / HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n")
.unwrap();
let mut index = String::new();
stream.read_to_string(&mut index).unwrap();
assert!(index.contains("jan-report"), "{index}");
Command::cargo_bin("unifier")
.unwrap()
.args([
"--home",
&home_s,
"put",
"gifts/ideas/ada-lovelace",
r#"{"signal":80,"person":"Ada"}"#,
])
.assert()
.success();
let mut stream = TcpStream::connect(host_port).unwrap();
stream
.write_all(
b"GET /keys/gifts/ideas/ada-lovelace HTTP/1.1\r\nHost: localhost\r\nAccept: application/json\r\nConnection: close\r\n\r\n",
)
.unwrap();
let mut key_resp = String::new();
stream.read_to_string(&mut key_resp).unwrap();
assert!(key_resp.contains("HTTP/1.1 200"), "{key_resp}");
assert!(key_resp.contains(r#""signal":80"#) || key_resp.contains(r#""signal": 80"#), "{key_resp}");
assert!(key_resp.contains("application/json"), "{key_resp}");
let mut stream = TcpStream::connect(host_port).unwrap();
stream
.write_all(b"GET /keys/gifts/ HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n")
.unwrap();
let mut list_resp = String::new();
stream.read_to_string(&mut list_resp).unwrap();
assert!(list_resp.contains("HTTP/1.1 200"), "{list_resp}");
assert!(list_resp.contains("gifts/ideas/ada-lovelace"), "{list_resp}");
Command::cargo_bin("unifier")
.unwrap()
.args([
"--home",
&home_s,
"web",
"key-url",
"gifts/ideas/ada-lovelace",
])
.assert()
.success()
.stdout(predicate::str::contains("/keys/gifts/ideas/ada-lovelace"));
Command::cargo_bin("unifier")
.unwrap()
.args(["--home", &home_s, "web", "rm", "jan-report"])
.assert()
.success();
Command::cargo_bin("unifier")
.unwrap()
.args(["--home", &home_s, "web", "list"])
.assert()
.success()
.stdout("");
}
#[test]
fn event_persists_expiry_envelope() {
let daemon = DaemonGuard::start();
let home_s = daemon.home_arg();
let home = daemon.home_path();
let assert = Command::cargo_bin("unifier")
.unwrap()
.args([
"--home",
&home_s,
"event",
r#"{"name":"keep"}"#,
"--ttl",
"0",
])
.assert()
.success();
let id = String::from_utf8(assert.get_output().stdout.clone())
.unwrap()
.trim()
.to_string();
Command::cargo_bin("unifier")
.unwrap()
.args(["--home", &home_s, "daemon", "flush"])
.assert()
.success();
let path = home.join("events").join(format!("{id}.json"));
let body = std::fs::read_to_string(&path).unwrap();
let v: serde_json::Value = serde_json::from_str(&body).unwrap();
assert_eq!(v["payload"]["name"], "keep");
assert!(v.get("created_at").is_some());
assert!(v.get("expires_at").is_none() || v["expires_at"].is_null());
}
#[test]
fn tick_control_socket_drives_lifecycle_and_broadcasts() {
use std::io::{BufRead, BufReader};
use std::os::unix::net::UnixStream;
use unifier::daemon::{events_socket_path, tick_socket_path, Notice, Response};
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));
assert!(tick_socket_path(&daemon.home).exists());
let tick_req = |req: Request| {
let mut tick = Client::connect_tick(&daemon.home).expect("connect tick.sock");
tick.request(req).unwrap()
};
let started = tick_req(Request::TickStart {
label: "jan-driver".into(),
});
match started {
Response::Ok {
tick: Some(1),
phase: Some(ref p),
label: Some(ref l),
..
} => {
assert_eq!(p, "start");
assert_eq!(l, "jan-driver");
}
other => panic!("unexpected start response: {other:?}"),
}
let phased = tick_req(Request::TickPhase {
phase: "sense".into(),
});
match phased {
Response::Ok {
tick: Some(1),
phase: Some(ref p),
..
} => assert_eq!(p, "sense"),
other => panic!("unexpected phase response: {other:?}"),
}
let rejected = tick_req(Request::Put {
key: "x".into(),
value: "y".into(),
});
match rejected {
Response::Err { error } => {
assert!(error.contains("tick.sock"), "{error}");
}
other => panic!("expected reject, got {other:?}"),
}
tick_req(Request::TickEnd);
let mut reader = BufReader::new(listener);
let mut notices = Vec::new();
for _ in 0..3 {
let mut line = String::new();
reader.read_line(&mut line).expect("read notice");
notices.push(Notice::from_json(&line).unwrap());
}
assert_eq!(
notices[0],
Notice::tick(1, "start", Some("jan-driver".into()))
);
assert_eq!(
notices[1],
Notice::tick(1, "sense", Some("jan-driver".into()))
);
assert_eq!(
notices[2],
Notice::tick(1, "end", Some("jan-driver".into()))
);
}
#[test]
fn daemon_status_lists_tick_socket() {
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("tick:"))
.stdout(predicate::str::contains("tick.sock"));
}
#[test]
fn poll_ack_drains_mailbox_over_one_connection() {
let daemon = DaemonGuard::start();
let home_s = daemon.home_arg();
for n in ["1", "2"] {
Command::cargo_bin("unifier")
.unwrap()
.args([
"--home",
&home_s,
"message",
"--from",
"scout",
"worker",
&format!(r#"{{"n":{n}}}"#),
])
.assert()
.success();
}
Command::cargo_bin("unifier")
.unwrap()
.args(["--home", &home_s, "poll", "worker", "--ack"])
.assert()
.success()
.stdout(predicate::str::contains(r#""n":1"#))
.stdout(predicate::str::contains(r#""n":2"#));
Command::cargo_bin("unifier")
.unwrap()
.args(["--home", &home_s, "poll", "worker", "--ack"])
.assert()
.success()
.stdout("");
}
#[test]
fn spawned_daemon_detaches_and_flushes_on_sigterm() {
let tmp = tempfile::tempdir().unwrap();
let home_s = tmp.path().display().to_string();
Command::cargo_bin("unifier")
.unwrap()
.args(["--home", &home_s, "put", "survive/me", "hello"])
.assert()
.success();
let pid_path = tmp.path().join(".daemon/unifier.pid");
let pid: i32 = std::fs::read_to_string(&pid_path)
.unwrap()
.trim()
.parse()
.unwrap();
let stat = std::fs::read_to_string(format!("/proc/{pid}/stat")).unwrap();
let sid: i32 = stat.split_whitespace().nth(5).unwrap().parse().unwrap();
assert_eq!(sid, pid, "daemon should be its own session leader");
unsafe { libc::kill(pid, libc::SIGTERM) };
for _ in 0..100 {
if !PathBuf::from(format!("/proc/{pid}/stat")).exists() {
break;
}
thread::sleep(Duration::from_millis(50));
}
assert!(tmp.path().join("keys/survive/me").exists());
Command::cargo_bin("unifier")
.unwrap()
.args(["--home", &home_s, "--no-daemon", "get", "survive/me"])
.assert()
.success()
.stdout(predicate::str::contains("hello"));
}
#[test]
fn client_connection_serves_sequential_requests() {
let daemon = DaemonGuard::start();
let mut client = Client::connect(&daemon.home).unwrap();
client
.request(Request::Put {
key: "a".into(),
value: "1".into(),
})
.unwrap();
let resp = client.request(Request::Get { key: "a".into() }).unwrap();
match resp {
unifier::daemon::Response::Ok { value, .. } => assert_eq!(value.as_deref(), Some("1")),
other => panic!("unexpected response: {other:?}"),
}
client.request(Request::Ping).unwrap();
}