#![cfg(feature = "comms")]
use std::path::{Path, PathBuf};
use std::process::{Child, Command};
use std::time::{Duration, Instant};
use basemind::comms::client::CommsClient;
use basemind::comms::ids::AgentId;
use basemind::comms::singleton::{CommsPaths, comms_socket_path, probe_alive};
const BIN: &str = env!("CARGO_BIN_EXE_basemind");
struct Daemon {
child: Child,
comms_dir: PathBuf,
socket: PathBuf,
}
impl Daemon {
fn start(comms_dir: &Path) -> Self {
let socket = comms_socket_path(comms_dir);
let child = Command::new(BIN)
.args(["comms", "daemon"])
.env("BASEMIND_COMMS_DIR", comms_dir)
.env("BASEMIND_DATA_HOME", comms_dir)
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.expect("spawn comms daemon");
let daemon = Self {
child,
comms_dir: comms_dir.to_path_buf(),
socket,
};
let deadline = Instant::now() + Duration::from_secs(10);
while Instant::now() < deadline {
if probe_alive(&daemon.socket) {
return daemon;
}
std::thread::sleep(Duration::from_millis(50));
}
panic!("comms daemon did not become ready");
}
fn socket(&self) -> &Path {
&self.socket
}
fn stop(self) {
let socket = self.socket.clone();
drop(self);
let deadline = Instant::now() + Duration::from_secs(10);
while Instant::now() < deadline {
if !probe_alive(&socket) {
return;
}
std::thread::sleep(Duration::from_millis(50));
}
panic!("comms daemon did not release its socket after stop");
}
}
impl Drop for Daemon {
fn drop(&mut self) {
let _ = Command::new(BIN)
.args(["comms", "stop"])
.env("BASEMIND_COMMS_DIR", &self.comms_dir)
.output();
if self.child.try_wait().ok().flatten().is_none() {
std::thread::sleep(Duration::from_millis(200));
if self.child.try_wait().ok().flatten().is_none() {
let _ = self.child.kill();
}
}
let _ = self.child.wait();
}
}
async fn connect(socket: &Path, agent: &str, root: &Path) -> CommsClient {
let paths = CommsPaths {
comms_dir: socket.parent().expect("socket parent").to_path_buf(),
socket_path: socket.to_path_buf(),
};
CommsClient::connect(
&paths,
AgentId::parse(agent).expect("agent id"),
None,
Some(root.to_path_buf()),
)
.await
.unwrap_or_else(|e| panic!("connect {agent}: {e}"))
}
fn git(args: &[&str], cwd: &Path) {
let out = Command::new("git")
.args(args)
.current_dir(cwd)
.output()
.expect("run git");
assert!(
out.status.success(),
"git {args:?} failed: {}",
String::from_utf8_lossy(&out.stderr)
);
}
fn init_git_repo(main: &Path) {
std::fs::create_dir_all(main).expect("mkdir main");
git(&["init", "-q", "-b", "main"], main);
git(&["config", "user.email", "t@example.com"], main);
git(&["config", "user.name", "Test"], main);
std::fs::write(main.join("a.rs"), b"pub fn alpha() {}\n").expect("write a.rs");
git(&["add", "."], main);
git(&["commit", "-qm", "init"], main);
}
fn init_bulk_git_repo(main: &Path, n_files: usize) {
std::fs::create_dir_all(main.join("src")).expect("mkdir src");
git(&["init", "-q", "-b", "main"], main);
git(&["config", "user.email", "t@example.com"], main);
git(&["config", "user.name", "Test"], main);
for i in 0..n_files {
let body = format!(
"pub fn f{i}() -> u32 {{ {i} }}\npub struct S{i};\nimpl S{i} {{ pub fn m{i}(&self) -> u32 {{ f{i}() }} }}\n"
);
std::fs::write(main.join("src").join(format!("m{i}.rs")), body).expect("write src file");
}
git(&["add", "."], main);
git(&["commit", "-qm", "bulk"], main);
}
fn stress_knob(var: &str, default: usize) -> usize {
std::env::var(var)
.ok()
.and_then(|v| v.parse().ok())
.filter(|n| *n > 0)
.unwrap_or(default)
}
#[tokio::test(flavor = "multi_thread")]
async fn concurrent_cold_rescans_open_the_workspace_once_and_all_succeed() {
const RACERS: usize = 6;
let tmp = tempfile::tempdir().expect("tempdir");
let comms_dir = tmp.path().join("comms");
let repo = tmp.path().join("repo");
init_git_repo(&repo);
let daemon = Daemon::start(&comms_dir);
let socket = daemon.socket().to_path_buf();
let mut clients = Vec::with_capacity(RACERS);
for c in 0..RACERS {
clients.push(connect(&socket, &format!("agent-race-{c}"), &repo).await);
}
let mut tasks = Vec::with_capacity(RACERS);
for (c, mut client) in clients.into_iter().enumerate() {
let repo = repo.clone();
tasks.push(tokio::spawn(async move {
client
.rescan(repo, None, true, false)
.await
.map(|_| ())
.map_err(|e| format!("racer {c}: {e}"))
}));
}
for task in tasks {
match task.await {
Ok(Ok(())) => {}
Ok(Err(message)) => panic!("a cold-open racer must succeed, not fail on the lock: {message}"),
Err(join) => panic!("a cold-open racer panicked: {join}"),
}
}
daemon.stop();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[ignore = "stress: many concurrent clients hammer one daemon; run with --ignored"]
async fn stress_many_concurrent_sessions_read_and_write_through_one_daemon() {
let clients = stress_knob("BASEMIND_STRESS_CLIENTS", 8);
let iters = stress_knob("BASEMIND_STRESS_ITERS", 8);
let files = stress_knob("BASEMIND_STRESS_FILES", 150);
let tmp = tempfile::tempdir().expect("tempdir");
let comms_dir = tmp.path().join("comms");
let repo = tmp.path().join("repo");
init_bulk_git_repo(&repo, files);
let daemon = Daemon::start(&comms_dir);
let socket = daemon.socket().to_path_buf();
let mut boot = connect(&socket, "agent-stress-boot", &repo).await;
let repo_id = boot.list_workspaces().await.expect("list workspaces")[0]
.repo_id
.clone()
.expect("git workspace has a repo id");
drop(boot);
let mut tasks = Vec::with_capacity(clients);
for c in 0..clients {
let socket = socket.clone();
let repo = repo.clone();
let repo_id = repo_id.clone();
tasks.push(tokio::spawn(async move {
let agent = format!("agent-stress-{c}");
let mut client = connect(&socket, &agent, &repo).await;
for _ in 0..iters {
client
.rescan(repo.clone(), None, true, false)
.await
.map_err(|e| format!("{agent} rescan: {e}"))?;
let workspaces = client
.list_workspaces()
.await
.map_err(|e| format!("{agent} list_workspaces: {e}"))?;
if workspaces.is_empty() {
return Err(format!("{agent}: registry lost the workspace mid-storm"));
}
let name = format!("wt-{c}");
let _ = client
.claim_worktree(repo_id.clone(), name.clone(), agent.clone())
.await;
let _ = client.release_worktree(repo_id.clone(), name, agent.clone()).await;
}
Ok::<(), String>(())
}));
}
let join_all = async {
let mut outcomes = Vec::with_capacity(tasks.len());
for task in tasks {
outcomes.push(task.await);
}
outcomes
};
let outcomes = tokio::time::timeout(Duration::from_secs(180), join_all)
.await
.expect("all stress clients must finish within 180s (no daemon deadlock)");
for (i, outcome) in outcomes.into_iter().enumerate() {
match outcome {
Ok(Ok(())) => {}
Ok(Err(message)) => panic!("stress client {i} failed: {message}"),
Err(join) => panic!("stress client {i} panicked: {join}"),
}
}
let mut after = connect(&socket, "agent-stress-after", &repo).await;
let workspaces = after.list_workspaces().await.expect("post-storm list_workspaces");
assert_eq!(
workspaces.len(),
1,
"exactly one workspace must remain registered after the storm, got {}",
workspaces.len()
);
let report = after
.rescan(repo.clone(), None, true, false)
.await
.expect("post-storm rescan");
assert!(
report.scanned >= 1,
"a post-storm rescan must still do real work (index not torn), got scanned={}",
report.scanned
);
drop(after);
daemon.stop();
}
#[test]
fn comms_stop_terminates_the_daemon_without_an_external_kill() {
let tmp = tempfile::tempdir().expect("tempdir");
let comms_dir = tmp.path().join("comms");
std::fs::create_dir_all(&comms_dir).expect("mkdir comms");
let socket = comms_socket_path(&comms_dir);
let mut child = Command::new(BIN)
.args(["comms", "daemon"])
.env("BASEMIND_COMMS_DIR", &comms_dir)
.env("BASEMIND_DATA_HOME", &comms_dir)
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.expect("spawn comms daemon");
let ready_by = Instant::now() + Duration::from_secs(10);
while Instant::now() < ready_by && !probe_alive(&socket) {
std::thread::sleep(Duration::from_millis(50));
}
assert!(probe_alive(&socket), "daemon did not become ready");
let stop = Command::new(BIN)
.args(["comms", "stop"])
.env("BASEMIND_COMMS_DIR", &comms_dir)
.output()
.expect("run comms stop");
assert!(
stop.status.success(),
"comms stop failed: {}",
String::from_utf8_lossy(&stop.stderr)
);
let exit_by = Instant::now() + Duration::from_secs(10);
let mut exited = false;
while Instant::now() < exit_by {
if child.try_wait().expect("try_wait").is_some() {
exited = true;
break;
}
std::thread::sleep(Duration::from_millis(50));
}
if !exited {
let _ = child.kill();
let _ = child.wait();
panic!("daemon did not self-terminate after `comms stop` within 10s (the #34 no-op bug)");
}
let _ = child.wait();
assert!(
!probe_alive(&socket),
"the socket must be released once the daemon self-terminates"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn registry_and_worktree_claim_survive_a_daemon_restart() {
let tmp = tempfile::tempdir().expect("tempdir");
let comms_dir = tmp.path().join("comms");
let repo = tmp.path().join("repo");
init_git_repo(&repo);
let daemon = Daemon::start(&comms_dir);
let socket = daemon.socket().to_path_buf();
let mut alice = connect(&socket, "agent-alice", &repo).await;
let workspaces = alice.list_workspaces().await.expect("list workspaces");
assert_eq!(workspaces.len(), 1, "Hello cwd auto-registers exactly one workspace");
let repo_id = workspaces[0].repo_id.clone().expect("a git workspace has a repo id");
let claimed = alice
.claim_worktree(repo_id.clone(), "(main)".to_string(), "agent-alice".to_string())
.await
.expect("alice claim");
assert!(claimed, "alice takes the previously-unclaimed (main) worktree");
drop(alice);
daemon.stop();
let daemon = Daemon::start(&comms_dir);
let socket = daemon.socket().to_path_buf();
let mut bob = connect(&socket, "agent-bob", &repo).await;
let workspaces = bob.list_workspaces().await.expect("list workspaces after restart");
assert_eq!(
workspaces.len(),
1,
"the registered workspace survives the daemon restart"
);
assert_eq!(
workspaces[0].repo_id.as_deref(),
Some(repo_id.as_str()),
"the same repo id reloads from the snapshot"
);
let worktrees = bob
.list_worktrees(repo_id.clone())
.await
.expect("list worktrees after restart");
let main = worktrees
.iter()
.find(|w| w.name == "(main)")
.expect("(main) worktree present after restart");
assert_eq!(
main.claimed_by.as_deref(),
Some("agent-alice"),
"populate_git preserves the reloaded claim when Hello re-enumerates the repo"
);
let bob_won = bob
.claim_worktree(repo_id.clone(), "(main)".to_string(), "agent-bob".to_string())
.await
.expect("bob claim after restart");
assert!(
!bob_won,
"the surviving claim blocks a second claimant across the restart"
);
let released = bob
.release_worktree(repo_id.clone(), "(main)".to_string(), "agent-alice".to_string())
.await
.expect("release alice's surviving claim");
assert!(released, "the reloaded claim is releasable by its original holder");
let bob_won = bob
.claim_worktree(repo_id.clone(), "(main)".to_string(), "agent-bob".to_string())
.await
.expect("bob claim after release");
assert!(bob_won, "with the claim released, the worktree is claimable again");
drop(bob);
daemon.stop();
}
#[tokio::test(flavor = "multi_thread")]
async fn should_converge_on_one_live_daemon_when_two_processes_race_a_cold_bind() {
let tmp = tempfile::tempdir().expect("tempdir");
let comms_dir = tmp.path().join("comms");
let socket = comms_socket_path(&comms_dir);
let spawn_one = || {
Command::new(BIN)
.args(["comms", "daemon"])
.env("BASEMIND_COMMS_DIR", &comms_dir)
.env("BASEMIND_DATA_HOME", &comms_dir)
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.expect("spawn comms daemon")
};
let mut child_a = spawn_one();
let mut child_b = spawn_one();
let deadline = Instant::now() + Duration::from_secs(20);
let mut alive = false;
while Instant::now() < deadline {
if probe_alive(&socket) {
alive = true;
break;
}
std::thread::sleep(Duration::from_millis(50));
}
assert!(alive, "one of the two racing daemons must come up serving");
let mut client = connect(&socket, "agent-race", tmp.path()).await;
let report = client.status().await.expect("status against the winning daemon");
assert!(
report.pid > 0,
"the winning daemon reports a real pid, got {}",
report.pid
);
drop(client);
let _ = Command::new(BIN)
.args(["comms", "stop"])
.env("BASEMIND_COMMS_DIR", &comms_dir)
.output();
for child in [&mut child_a, &mut child_b] {
if child.try_wait().ok().flatten().is_none() {
std::thread::sleep(Duration::from_millis(300));
if child.try_wait().ok().flatten().is_none() {
let _ = child.kill();
}
}
let _ = child.wait();
}
let deadline = Instant::now() + Duration::from_secs(10);
while Instant::now() < deadline {
if !probe_alive(&socket) {
return;
}
std::thread::sleep(Duration::from_millis(50));
}
panic!("socket still answers after both racing daemons were reaped");
}
#[tokio::test(flavor = "multi_thread")]
async fn should_reflect_branch_creation_and_worktree_add_remove_on_fresh_connects() {
let tmp = tempfile::tempdir().expect("tempdir");
let comms_dir = tmp.path().join("comms");
let repo = tmp.path().join("repo");
init_git_repo(&repo);
let daemon = Daemon::start(&comms_dir);
let socket = daemon.socket().to_path_buf();
let mut client = connect(&socket, "agent-reg", &repo).await;
let workspaces = client.list_workspaces().await.expect("list workspaces");
assert_eq!(workspaces.len(), 1, "Hello cwd auto-registers exactly one workspace");
let repo_id = workspaces[0].repo_id.clone().expect("a git workspace has a repo id");
let branches = client.list_branches(repo_id.clone()).await.expect("list branches");
assert!(
branches.iter().any(|b| b.name == "main"),
"the initial checkout's branch is enumerated, got {branches:?}"
);
drop(client);
git(&["branch", "feature"], &repo);
let mut client = connect(&socket, "agent-reg-2", &repo).await;
let branches = client
.list_branches(repo_id.clone())
.await
.expect("list branches after branch create");
assert!(
branches.iter().any(|b| b.name == "feature"),
"the new branch is enumerated after a fresh Hello re-populates, got {branches:?}"
);
drop(client);
let linked_path = tmp.path().join("linked-wt");
git(
&[
"worktree",
"add",
"-b",
"wt-feature",
linked_path.to_str().expect("utf8 path"),
"feature",
],
&repo,
);
let mut client = connect(&socket, "agent-reg-3", &repo).await;
let worktrees = client
.list_worktrees(repo_id.clone())
.await
.expect("list worktrees after add");
let linked = worktrees
.iter()
.find(|w| w.path == linked_path.canonicalize().expect("canonicalize linked path"));
assert!(
linked.is_some(),
"the newly linked worktree is enumerated after a fresh Hello, got {worktrees:?}"
);
assert_eq!(
linked.expect("checked above").branch.as_deref(),
Some("wt-feature"),
"the linked worktree's checked-out branch is recorded"
);
drop(client);
git(&["worktree", "remove", linked_path.to_str().expect("utf8 path")], &repo);
let mut client = connect(&socket, "agent-reg-4", &repo).await;
let worktrees = client
.list_worktrees(repo_id.clone())
.await
.expect("list worktrees after remove");
assert!(
!worktrees
.iter()
.any(|w| w.path == linked_path.canonicalize().unwrap_or(linked_path.clone())),
"the removed worktree is pruned from the registry after a fresh Hello, got {worktrees:?}"
);
drop(client);
daemon.stop();
}
#[cfg(all(feature = "comms", unix))]
#[tokio::test(flavor = "multi_thread")]
#[ignore = "watchdog fires on a ~30s timer; run explicitly with --ignored"]
async fn should_self_terminate_when_its_socket_is_reclaimed_by_another_daemon() {
let tmp = tempfile::tempdir().expect("tempdir");
let comms_dir = tmp.path().join("comms");
let socket = comms_socket_path(&comms_dir);
let mut child_a = Command::new(BIN)
.args(["comms", "daemon"])
.env("BASEMIND_COMMS_DIR", &comms_dir)
.env("BASEMIND_DATA_HOME", &comms_dir)
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.expect("spawn daemon A");
let deadline = Instant::now() + Duration::from_secs(10);
while Instant::now() < deadline && !probe_alive(&socket) {
std::thread::sleep(Duration::from_millis(50));
}
assert!(probe_alive(&socket), "daemon A must come up before we orphan it");
std::fs::remove_file(&socket).expect("unlink daemon A's socket");
let mut child_b = Command::new(BIN)
.args(["comms", "daemon"])
.env("BASEMIND_COMMS_DIR", &comms_dir)
.env("BASEMIND_DATA_HOME", &comms_dir)
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.expect("spawn daemon B");
let deadline = Instant::now() + Duration::from_secs(10);
while Instant::now() < deadline && !probe_alive(&socket) {
std::thread::sleep(Duration::from_millis(50));
}
assert!(probe_alive(&socket), "daemon B must come up on the rebound socket");
let watchdog_deadline = Instant::now() + Duration::from_secs(90);
let mut a_exited = false;
while Instant::now() < watchdog_deadline {
if child_a.try_wait().ok().flatten().is_some() {
a_exited = true;
break;
}
std::thread::sleep(Duration::from_millis(500));
}
assert!(
a_exited,
"daemon A must self-terminate once its socket is reclaimed by daemon B (orphan watchdog)"
);
let _ = child_a.wait();
let _ = Command::new(BIN)
.args(["comms", "stop"])
.env("BASEMIND_COMMS_DIR", &comms_dir)
.output();
if child_b.try_wait().ok().flatten().is_none() {
std::thread::sleep(Duration::from_millis(300));
if child_b.try_wait().ok().flatten().is_none() {
let _ = child_b.kill();
}
}
let _ = child_b.wait();
}
#[tokio::test(flavor = "multi_thread")]
async fn should_rescan_successfully_and_ignore_a_stale_legacy_in_repo_index() {
let tmp = tempfile::tempdir().expect("tempdir");
let comms_dir = tmp.path().join("comms");
let repo = tmp.path().join("repo");
init_git_repo(&repo);
#[derive(serde::Serialize)]
struct LegacyIndex {
schema_ver: u16,
files: std::collections::BTreeMap<String, ()>,
doc_files: std::collections::BTreeMap<String, ()>,
}
let legacy_dir = repo.join(".basemind");
std::fs::create_dir_all(&legacy_dir).expect("mkdir legacy .basemind");
let legacy_index = LegacyIndex {
schema_ver: 21,
files: std::collections::BTreeMap::new(),
doc_files: std::collections::BTreeMap::new(),
};
let legacy_bytes = rmp_serde::to_vec_named(&legacy_index).expect("encode legacy index");
let legacy_index_path = legacy_dir.join("index.msgpack");
std::fs::write(&legacy_index_path, &legacy_bytes).expect("write legacy index.msgpack");
let legacy_bytes_before = std::fs::read(&legacy_index_path).expect("read back legacy index.msgpack");
let daemon = Daemon::start(&comms_dir);
let socket = daemon.socket().to_path_buf();
let mut client = connect(&socket, "agent-migrate", &repo).await;
let report = client
.rescan(repo.clone(), None, true, false)
.await
.expect("rescan a repo carrying a stale legacy in-repo index");
assert!(
report.scanned >= 1,
"rescan must actually scan the repo's files, got scanned={}",
report.scanned
);
let legacy_bytes_after = std::fs::read(&legacy_index_path).expect("re-read legacy index.msgpack");
assert_eq!(
legacy_bytes_after, legacy_bytes_before,
"the in-repo legacy .basemind/index.msgpack must be left byte-for-byte untouched"
);
let workspace_key = basemind::store::workspace_key(&repo);
let workspace_dir = comms_dir.join("cache").join("workspaces").join(&workspace_key);
assert!(
workspace_dir.exists(),
"the global cache must gain a workspace dir for this repo at {}",
workspace_dir.display()
);
drop(client);
daemon.stop();
}
#[tokio::test(flavor = "multi_thread")]
async fn should_drain_cleanly_with_an_in_flight_rescan_and_reload_registry_after_restart() {
let tmp = tempfile::tempdir().expect("tempdir");
let comms_dir = tmp.path().join("comms");
let repo = tmp.path().join("repo");
init_git_repo(&repo);
let daemon = Daemon::start(&comms_dir);
let socket = daemon.socket().to_path_buf();
let mut client = connect(&socket, "agent-drain", &repo).await;
let workspaces = client.list_workspaces().await.expect("list workspaces");
let repo_id = workspaces[0].repo_id.clone().expect("git workspace has a repo id");
let rescan_repo = repo.clone();
let rescan_task = tokio::spawn(async move { client.rescan(rescan_repo, None, true, false).await });
tokio::time::sleep(Duration::from_millis(100)).await;
let comms_dir_for_stop = comms_dir.clone();
let stop_output = tokio::task::spawn_blocking(move || {
Command::new(BIN)
.args(["comms", "stop"])
.env("BASEMIND_COMMS_DIR", &comms_dir_for_stop)
.output()
})
.await
.expect("join stop task")
.expect("run the `comms stop` CLI invocation");
assert!(
stop_output.status.success(),
"the `comms stop` RPC must succeed even with a rescan in flight, stderr: {}",
String::from_utf8_lossy(&stop_output.stderr)
);
let rescan_outcome = tokio::time::timeout(Duration::from_secs(15), rescan_task)
.await
.expect("in-flight rescan must resolve within the timeout, not hang");
match rescan_outcome {
Ok(Ok(report)) => {
assert!(
report.scanned >= 1,
"a rescan that completed must report real work, got scanned={}",
report.scanned
);
}
Ok(Err(client_error)) => {
let _ = client_error;
}
Err(join_error) => panic!("rescan task must not panic, got: {join_error}"),
}
daemon.stop();
let daemon = Daemon::start(&comms_dir);
let socket = daemon.socket().to_path_buf();
let mut client = connect(&socket, "agent-drain-2", &repo).await;
let workspaces = client.list_workspaces().await.expect("list workspaces after restart");
assert_eq!(
workspaces.len(),
1,
"the registered workspace survives the drain + restart"
);
assert_eq!(
workspaces[0].repo_id.as_deref(),
Some(repo_id.as_str()),
"the same repo id reloads from the snapshot after the drain"
);
drop(client);
daemon.stop();
}