use std::path::{Path, PathBuf};
use std::process::Output;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::{Duration, Instant};
use tempfile::TempDir;
mod common;
use common::{build_kache, isolated_config_path, kache_binary};
fn rustc_path() -> String {
std::env::var("RUSTC").unwrap_or_else(|_| "rustc".to_string())
}
struct Client {
_cache: TempDir,
cache_dir: PathBuf,
command_seq: AtomicUsize,
}
impl Client {
fn new(shared_folder: &Path) -> Self {
let cache = TempDir::new().unwrap();
let cache_dir = cache.path().to_path_buf();
std::fs::write(
isolated_config_path(&cache_dir),
format!(
"[cache.remote]\n\
type = \"filesystem\"\n\
path = \"{}\"\n\
prefix = \"artifacts\"\n",
shared_folder.display().to_string().replace('\\', "\\\\"),
),
)
.unwrap();
Client {
_cache: cache,
cache_dir,
command_seq: AtomicUsize::new(0),
}
}
fn kache(&self) -> std::process::Command {
let mut cmd = std::process::Command::new(kache_binary());
cmd.env("KACHE_CACHE_DIR", &self.cache_dir)
.env("KACHE_CONFIG", isolated_config_path(&self.cache_dir))
.env("KACHE_LOG", "off")
.env("KACHE_DAEMON_IDLE_TIMEOUT", "60")
.env_remove("KACHE_SOCKET_PATH")
.env_remove("RUSTC_WRAPPER")
.env_remove("CARGO_BUILD_RUSTC_WRAPPER");
cmd
}
fn run_within(&self, args: &[&str], deadline: Duration) -> Output {
let seq = self.command_seq.fetch_add(1, Ordering::Relaxed);
let out_path = self.cache_dir.join(format!("cmd-{seq}.out"));
let err_path = self.cache_dir.join(format!("cmd-{seq}.err"));
let stdout = std::fs::File::create(&out_path).expect("creating stdout capture");
let stderr = std::fs::File::create(&err_path).expect("creating stderr capture");
let mut child = self
.kache()
.args(args)
.stdin(std::process::Stdio::null())
.stdout(stdout)
.stderr(stderr)
.spawn()
.expect("failed to spawn kache");
let start = Instant::now();
let status = loop {
match child.try_wait().expect("polling kache") {
Some(status) => break status,
None if start.elapsed() >= deadline => {
let _ = child.kill();
let _ = child.wait();
panic!(
"`kache {}` did not finish within {deadline:?} (kunobi-ninja/kache#704).\n\
stdout: {}\nstderr: {}",
args.join(" "),
std::fs::read_to_string(&out_path).unwrap_or_default(),
std::fs::read_to_string(&err_path).unwrap_or_default(),
);
}
None => std::thread::sleep(Duration::from_millis(50)),
}
};
Output {
status,
stdout: std::fs::read(&out_path).unwrap_or_default(),
stderr: std::fs::read(&err_path).unwrap_or_default(),
}
}
fn run(&self, args: &[&str]) -> Output {
self.run_within(args, Duration::from_secs(90))
}
fn compile(&self, src: &Path, out_dir: &Path) {
let out_dir = out_dir.display().to_string();
let src = src.display().to_string();
let rustc = rustc_path();
let output = self.run(&[
&rustc,
"--crate-name",
"fsremote",
"--crate-type",
"lib",
"--edition",
"2021",
"--emit=link",
"--out-dir",
&out_dir,
&src,
]);
assert!(
output.status.success(),
"kache rustc failed.\nstderr: {}",
String::from_utf8_lossy(&output.stderr),
);
}
fn start_daemon(&self) {
let output = self.run(&["daemon", "start"]);
assert!(
output.status.success(),
"kache daemon start failed.\nstdout: {}\nstderr: {}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr),
);
}
fn stop_daemon_and_wait(&self) {
let output = self.run(&["daemon", "stop"]);
assert!(
output.status.success(),
"kache daemon stop failed.\nstdout: {}\nstderr: {}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr),
);
let run_lock_path = self.cache_dir.join("daemon.run.lock");
let run_lock = std::fs::OpenOptions::new()
.create(true)
.write(true)
.truncate(false)
.open(&run_lock_path)
.expect("opening daemon run lock probe");
let deadline = Instant::now() + Duration::from_secs(45);
loop {
match run_lock.try_lock() {
Ok(()) => {
run_lock.unlock().expect("releasing daemon run lock probe");
return;
}
Err(std::fs::TryLockError::Error(error)) => {
panic!(
"failed to probe daemon lifetime lock {}: {error}",
run_lock_path.display()
);
}
Err(error @ std::fs::TryLockError::WouldBlock) if Instant::now() >= deadline => {
panic!(
"daemon lifetime lock {} remained held after its drain phase: {error}",
run_lock_path.display()
);
}
Err(std::fs::TryLockError::WouldBlock) => {
std::thread::sleep(Duration::from_millis(100));
}
}
}
}
fn sync(&self, args: &[&str]) -> Output {
let mut argv = vec!["sync"];
argv.extend_from_slice(args);
let output = self.run(&argv);
assert!(
output.status.success(),
"kache sync {args:?} failed.\nstdout: {}\nstderr: {}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr),
);
let stderr = String::from_utf8_lossy(&output.stderr);
assert!(
!stderr
.lines()
.any(|line| line.trim_end().ends_with("failed)")),
"kache sync {args:?} reported failed transfers despite exiting successfully.\n\
stdout: {}\nstderr: {stderr}",
String::from_utf8_lossy(&output.stdout),
);
output
}
fn results(&self) -> Vec<String> {
let output = self.run(&["report", "--format", "json", "--since", "1h"]);
assert!(output.status.success(), "kache report failed");
let report: serde_json::Value =
serde_json::from_slice(&output.stdout).expect("report should be valid json");
report["all_events"]
.as_array()
.map(|events| {
events
.iter()
.filter_map(|e| e["result"].as_str().map(str::to_string))
.collect()
})
.unwrap_or_default()
}
}
impl Drop for Client {
fn drop(&mut self) {
let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
self.run_within(&["daemon", "stop"], Duration::from_secs(30))
}));
}
}
#[test]
fn two_independent_clients_share_hits_through_one_folder() {
build_kache();
let shared = TempDir::new().unwrap();
let sources = TempDir::new().unwrap();
let src = sources.path().join("lib.rs");
std::fs::write(&src, "pub fn answer() -> u32 { 42 }\n").unwrap();
let alpha = Client::new(shared.path());
let out_a = TempDir::new().unwrap();
alpha.compile(&src, out_a.path());
assert!(
alpha.results().iter().any(|r| r == "miss" || r == "dup"),
"client A's first compile must be a real compile: {:?}",
alpha.results()
);
alpha.stop_daemon_and_wait();
let push = alpha.sync(&["--push"]);
let manifests = shared.path().join("artifacts/v3/manifests");
assert!(
manifests.is_dir(),
"the push must have written the v3 object layout into the shared folder.\n\
stdout: {}\nstderr: {}",
String::from_utf8_lossy(&push.stdout),
String::from_utf8_lossy(&push.stderr),
);
let beta = Client::new(shared.path());
let out_b = TempDir::new().unwrap();
beta.start_daemon();
beta.compile(&src, out_b.path());
let beta_results = beta.results();
assert!(
beta_results.iter().any(|r| r == "remote_hit"),
"client B must pull A's artifact from the shared folder on miss: {beta_results:?}"
);
assert!(
out_b.path().join("libfsremote.rlib").is_file(),
"client B must end up with the artifact materialized"
);
}
#[test]
fn a_fresh_client_can_seed_itself_from_the_shared_folder() {
build_kache();
let shared = TempDir::new().unwrap();
let sources = TempDir::new().unwrap();
let src = sources.path().join("lib.rs");
std::fs::write(&src, "pub fn seeded() -> u32 { 11 }\n").unwrap();
let alpha = Client::new(shared.path());
let out_a = TempDir::new().unwrap();
alpha.compile(&src, out_a.path());
alpha.stop_daemon_and_wait();
alpha.sync(&["--push"]);
let beta = Client::new(shared.path());
let out_b = TempDir::new().unwrap();
beta.sync(&["--pull", "--all"]);
beta.compile(&src, out_b.path());
let beta_results = beta.results();
assert!(
beta_results.iter().any(|r| r == "local_hit"),
"after --pull --all the artifact is local, so the compile is a local hit: {beta_results:?}"
);
assert!(out_b.path().join("libfsremote.rlib").is_file());
}
#[test]
fn capturing_kache_output_through_pipes_does_not_hang() {
build_kache();
let shared = TempDir::new().unwrap();
let client = Client::new(shared.path());
let sources = TempDir::new().unwrap();
let src = sources.path().join("lib.rs");
std::fs::write(&src, "pub fn piped() -> u32 { 3 }\n").unwrap();
let out_dir = TempDir::new().unwrap();
let (tx, rx) = std::sync::mpsc::channel();
let cache_dir = client.cache_dir.clone();
let src_path = src.display().to_string();
let out_path = out_dir.path().display().to_string();
let rustc = rustc_path();
std::thread::spawn(move || {
let result = std::process::Command::new(kache_binary())
.env("KACHE_CACHE_DIR", &cache_dir)
.env("KACHE_CONFIG", cache_dir.join("config.toml"))
.env("KACHE_LOG", "off")
.env("KACHE_DAEMON_IDLE_TIMEOUT", "60")
.env_remove("KACHE_SOCKET_PATH")
.env_remove("RUSTC_WRAPPER")
.env_remove("CARGO_BUILD_RUSTC_WRAPPER")
.args([
&rustc,
"--crate-name",
"piped",
"--crate-type",
"lib",
"--edition",
"2021",
"--emit=link",
"--out-dir",
&out_path,
&src_path,
])
.output();
let _ = tx.send(result.map(|o| o.status.success()));
});
match rx.recv_timeout(Duration::from_secs(120)) {
Ok(Ok(true)) => {}
Ok(Ok(false)) => panic!("the piped invocation failed"),
Ok(Err(e)) => panic!("could not run kache: {e}"),
Err(_) => panic!(
"a piped `Command::output()` around kache did not return within 120s — \
the auto-started daemon is holding the caller's pipe open \
(kunobi-ninja/kache#704)"
),
}
}
#[cfg(unix)]
#[test]
fn dropping_a_client_stops_its_daemon() {
build_kache();
let shared = TempDir::new().unwrap();
let socket = {
let client = Client::new(shared.path());
client.start_daemon();
let socket = client.cache_dir.join("daemon.sock");
assert!(
socket.exists(),
"the daemon should have published its socket at {}",
socket.display()
);
socket
};
for _ in 0..50 {
if !socket.exists() {
return;
}
std::thread::sleep(Duration::from_millis(100));
}
panic!(
"daemon socket {} still present after the client was dropped — \
Drop cleanup did not stop the daemon",
socket.display()
);
}
#[test]
fn the_shared_folder_holds_objects_only() {
build_kache();
let shared = TempDir::new().unwrap();
let sources = TempDir::new().unwrap();
let src = sources.path().join("lib.rs");
std::fs::write(&src, "pub fn only_objects() -> u32 { 7 }\n").unwrap();
let client = Client::new(shared.path());
let out = TempDir::new().unwrap();
client.compile(&src, out.path());
client.stop_daemon_and_wait();
client.sync(&["--push"]);
let mut files = Vec::new();
collect_files(shared.path(), &mut files);
assert!(
!files.is_empty(),
"the push should have written objects to the shared folder"
);
for file in &files {
let name = file.file_name().unwrap_or_default().to_string_lossy();
assert!(
!name.contains("index.db"),
"a SQLite index must never live on the shared folder: {}",
file.display()
);
assert!(
!name.ends_with("-wal") && !name.ends_with("-shm"),
"SQLite WAL sidecars must never live on the shared folder: {}",
file.display()
);
#[cfg(unix)]
{
use std::os::unix::fs::MetadataExt;
let meta = std::fs::metadata(file).unwrap();
assert_eq!(
meta.nlink(),
1,
"shared-folder object must not be a hardlink: {}",
file.display()
);
}
}
}
fn collect_files(dir: &Path, out: &mut Vec<PathBuf>) {
let Ok(entries) = std::fs::read_dir(dir) else {
return;
};
for entry in entries.flatten() {
let path = entry.path();
if path.is_dir() {
collect_files(&path, out);
} else {
out.push(path);
}
}
}