#![cfg(unix)]
use std::path::{Path, PathBuf};
use std::process::{Child, Command, ExitStatus, Stdio};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Barrier};
use std::time::{Duration, Instant};
use mermaid_cli::runtime::{NewApproval, NewTask, RuntimeStore, TaskStatus};
static SEQ: AtomicU64 = AtomicU64::new(0);
const TIMEOUT: Duration = Duration::from_secs(20);
fn fresh_base(tag: &str) -> PathBuf {
let base = std::env::temp_dir().join(format!(
"mermaid_daemon_it_{tag}_{}_{}",
std::process::id(),
SEQ.fetch_add(1, Ordering::Relaxed)
));
let _ = std::fs::remove_dir_all(&base);
std::fs::create_dir_all(&base).expect("create temp base dir");
base
}
fn data_paths(base: &Path) -> (PathBuf, PathBuf) {
let dir = base.join("mermaid");
(dir.join("runtime.sqlite3"), dir.join("mermaidd.sock"))
}
fn spawn_daemon(base: &Path, stderr_log: &Path) -> Child {
Command::new(env!("CARGO_BIN_EXE_mermaidd"))
.env("XDG_DATA_HOME", base)
.env("HOME", base) .env_remove("MERMAID_DAEMON_ENABLE_TCP") .stdout(Stdio::null())
.stderr(Stdio::from(
std::fs::File::create(stderr_log).expect("create stderr log"),
))
.spawn()
.expect("spawn mermaidd")
}
struct DaemonGuard(Child);
impl Drop for DaemonGuard {
fn drop(&mut self) {
let _ = self.0.kill();
let _ = self.0.wait();
}
}
struct DirGuard(PathBuf);
impl Drop for DirGuard {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.0);
}
}
fn wait_until(mut cond: impl FnMut() -> bool, timeout: Duration) -> bool {
let start = Instant::now();
loop {
if cond() {
return true;
}
if start.elapsed() >= timeout {
return false;
}
std::thread::sleep(Duration::from_millis(50));
}
}
fn wait_for_exit(child: &mut Child, timeout: Duration) -> Option<ExitStatus> {
let start = Instant::now();
loop {
match child.try_wait() {
Ok(Some(status)) => return Some(status),
Ok(None) => {},
Err(_) => return None,
}
if start.elapsed() >= timeout {
return None;
}
std::thread::sleep(Duration::from_millis(50));
}
}
fn read(path: &Path) -> String {
std::fs::read_to_string(path).unwrap_or_default()
}
fn make_approval(store: &RuntimeStore, action: &str) -> String {
store
.approvals()
.create(NewApproval {
task_id: None,
proposed_action: action.to_string(),
risk_classification: "Write".to_string(),
policy_decision: "Ask".to_string(),
args_summary: None,
checkpoint_id: None,
pending_action_json: None,
})
.expect("create approval")
.id
}
#[test]
fn approval_claim_has_exactly_one_winner_under_thread_contention() {
let base = fresh_base("claim");
let _dir = DirGuard(base.clone());
let (db, _sock) = data_paths(&base);
let id = {
let store = RuntimeStore::open(&db).expect("open store");
make_approval(&store, "write_file secret.txt")
};
const THREADS: usize = 8;
let barrier = Arc::new(Barrier::new(THREADS));
let handles: Vec<_> = (0..THREADS)
.map(|_| {
let db = db.clone();
let id = id.clone();
let barrier = Arc::clone(&barrier);
std::thread::spawn(move || {
let store = RuntimeStore::open(&db).expect("open store in thread");
barrier.wait(); store
.approvals()
.claim(&id)
.expect("claim must not error — busy_timeout serializes writers")
})
})
.collect();
let winners = handles
.into_iter()
.map(|h| h.join().expect("claimer thread panicked"))
.filter(|&won| won)
.count();
assert_eq!(
winners, 1,
"exactly one concurrent claim may win the #118 race, got {winners}"
);
let store = RuntimeStore::open(&db).expect("reopen store");
let appr = store
.approvals()
.get(&id)
.expect("get approval")
.expect("approval exists");
assert_eq!(
appr.user_decision.as_deref(),
Some("approving"),
"the single winner must hold the intermediate claim"
);
}
#[test]
#[ignore = "spawns a real mermaidd; run with: cargo test --test daemon_integration -- --ignored"]
fn daemon_singleton_flock_rejects_a_second_start() {
let base = fresh_base("flock");
let _dir = DirGuard(base.clone());
let (_db, sock) = data_paths(&base);
let a_err = base.join("a.stderr");
let b_err = base.join("b.stderr");
let _a = DaemonGuard(spawn_daemon(&base, &a_err));
assert!(
wait_until(|| sock.exists(), TIMEOUT),
"daemon A never bound its socket; stderr:\n{}",
read(&a_err)
);
let mut b = DaemonGuard(spawn_daemon(&base, &b_err));
let status = match wait_for_exit(&mut b.0, TIMEOUT) {
Some(status) => status,
None => panic!(
"daemon B never exited — the #131 singleton flock failed to reject a second start"
),
};
assert!(
!status.success(),
"daemon B must exit non-zero while A holds the lock; got {status:?}"
);
let b_stderr = read(&b_err);
assert!(
b_stderr.contains("lock held") || b_stderr.to_lowercase().contains("already"),
"B should report the singleton-lock rejection (#131); its stderr was:\n{b_stderr}"
);
}
#[test]
#[ignore = "spawns a real mermaidd; run with: cargo test --test daemon_integration -- --ignored"]
fn daemon_startup_reconciles_running_task_and_gcs_old_archived_row() {
let base = fresh_base("recover");
let _dir = DirGuard(base.clone());
let (db, sock) = data_paths(&base);
let err = base.join("d.stderr");
let (running_task, old_appr, fresh_appr) = {
let store = RuntimeStore::open(&db).expect("open store for seeding");
let task = store
.tasks()
.create(NewTask::new("stuck task", "/tmp/project", "test/model").daemon_owned())
.expect("create task");
store
.tasks()
.update_status(&task.id, TaskStatus::Running, None)
.expect("mark running");
let old = make_approval(&store, "old archived action");
let fresh = make_approval(&store, "fresh archived action");
store
.approvals()
.archive(&[old.clone(), fresh.clone()], "test")
.expect("archive both");
(task.id, old, fresh)
};
{
let conn = rusqlite::Connection::open(&db).expect("open db to backdate");
let changed = conn
.execute(
"UPDATE approvals SET archived_at = ?2 WHERE id = ?1",
rusqlite::params![old_appr, "2000-01-01T00:00:00+00:00"],
)
.expect("backdate archived_at");
assert_eq!(changed, 1, "backdate must hit exactly the old approval");
}
let _daemon = DaemonGuard(spawn_daemon(&base, &err));
assert!(
wait_until(|| sock.exists(), TIMEOUT),
"daemon never became ready (socket never bound); stderr:\n{}",
read(&err)
);
let store = RuntimeStore::open(&db).expect("reopen store to verify");
let task = store
.tasks()
.get(&running_task)
.expect("get task")
.expect("task row exists");
assert_eq!(
task.status,
TaskStatus::Failed,
"#120: a task left Running by a crashed daemon must be reconciled to Failed on restart"
);
assert!(
store.approvals().get(&old_appr).expect("get old").is_none(),
"#130: the long-archived approval must be GC'd on startup"
);
assert!(
store
.approvals()
.get(&fresh_appr)
.expect("get fresh")
.is_some(),
"#130: a freshly-archived approval must be kept"
);
}