use crate::transport::prelude::*;
use crate::transport::{ListenerOptions, TokioListener, TokioStream, socket_name};
use anyhow::{Context, Result};
use kache_core::{PrefetchDisposition, PrefetchPlan};
use serde::{Deserialize, Serialize};
use std::collections::{HashMap, HashSet};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, AtomicU32, AtomicU64, Ordering};
use std::sync::{Arc, Mutex, OnceLock};
use std::time::{Duration, Instant};
use tokio::io::{AsyncBufRead, AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::sync::{Notify, RwLock};
use crate::config::Config;
use crate::events;
use crate::store::Store;
const KEY_CACHE_AUTHORITATIVE_MULTIPLIER: u64 = 5;
const KEY_CACHE_AUTHORITATIVE_MAX_AGE: Duration = Duration::from_secs(300);
const REMOTE_CHECK_WARMING_GRACE: Duration = Duration::from_millis(750);
fn key_cache_miss_is_authoritative(refresh_secs: u64, age: Option<Duration>) -> bool {
if refresh_secs == 0 {
return false;
}
let refresh_window =
Duration::from_secs(refresh_secs.saturating_mul(KEY_CACHE_AUTHORITATIVE_MULTIPLIER));
let authoritative_for = refresh_window.min(KEY_CACHE_AUTHORITATIVE_MAX_AGE);
matches!(age, Some(age) if age <= authoritative_for)
}
fn speculative_prefetch_disabled(prefetch_enabled: bool) -> bool {
!prefetch_enabled
}
fn should_start_speculative_prefetch(remote_configured: bool, prefetch_enabled: bool) -> bool {
remote_configured && prefetch_enabled
}
fn key_cache_periodic_refresh_disabled(refresh_secs: u64) -> bool {
refresh_secs == 0
}
const REMOTE_HEAD_FAILURE_THRESHOLD: u32 = 3;
const REMOTE_HEAD_DEGRADED_FOR: Duration = Duration::from_secs(45);
const DAEMON_START_TIMEOUT: Duration = Duration::from_secs(8);
const DAEMON_START_POLL_INTERVAL: Duration = Duration::from_millis(100);
const DAEMON_COORD_HEARTBEAT_INTERVAL: Duration = Duration::from_secs(2);
const DAEMON_CONFIG_WATCH_INTERVAL: Duration = Duration::from_secs(15);
const DAEMON_COORD_STALE_AFTER: Duration = Duration::from_secs(15);
const VERSION: &str = crate::VERSION;
const FILE_HASH_MEMORY_CACHE_CAP: usize = 4096;
pub fn build_epoch() -> u64 {
static BUILD_EPOCH: OnceLock<u64> = OnceLock::new();
*BUILD_EPOCH.get_or_init(|| {
std::env::current_exe()
.and_then(std::fs::metadata)
.and_then(|m| m.modified())
.ok()
.and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
.map(|d| d.as_secs())
.unwrap_or(0)
})
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
enum DaemonPhase {
Starting,
Ready,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
struct DaemonCoordState {
pid: u32,
build_epoch: u64,
phase: DaemonPhase,
updated_at_ms: u64,
}
#[derive(Debug, Clone)]
struct DaemonCoordFile {
path: PathBuf,
pid: u32,
build_epoch: u64,
}
impl DaemonCoordFile {
fn for_socket(socket_path: &Path) -> Self {
Self {
path: daemon_state_path(socket_path),
pid: std::process::id(),
build_epoch: build_epoch(),
}
}
fn write_phase(&self, phase: DaemonPhase) -> Result<()> {
let state = DaemonCoordState {
pid: self.pid,
build_epoch: self.build_epoch,
phase,
updated_at_ms: now_millis(),
};
write_json_atomically(&self.path, &state)
}
}
struct DaemonCoordGuard {
path: PathBuf,
}
struct SocketCleanupGuard {
path: PathBuf,
}
struct RemoteHealth {
head_probe_failures: AtomicU32,
head_probe_degraded_until_ms: AtomicU64,
suppressed_head_probes: AtomicU32,
}
impl RemoteHealth {
fn new() -> Self {
Self {
head_probe_failures: AtomicU32::new(0),
head_probe_degraded_until_ms: AtomicU64::new(0),
suppressed_head_probes: AtomicU32::new(0),
}
}
fn head_probe_is_degraded(&self) -> bool {
now_millis() < self.head_probe_degraded_until_ms.load(Ordering::Acquire)
}
fn note_head_probe_failure(&self, error: &str) {
let failures = self.head_probe_failures.fetch_add(1, Ordering::AcqRel) + 1;
if failures < REMOTE_HEAD_FAILURE_THRESHOLD {
if failures == 1 {
tracing::warn!(
"remote HEAD probe failed ({failures}/{REMOTE_HEAD_FAILURE_THRESHOLD} before degradation): {error}"
);
} else {
tracing::debug!(
"remote HEAD probe failed ({failures}/{REMOTE_HEAD_FAILURE_THRESHOLD} before degradation): {error}"
);
}
return;
}
let was_degraded = self.head_probe_is_degraded();
let degrade_until = now_millis() + REMOTE_HEAD_DEGRADED_FOR.as_millis() as u64;
self.head_probe_degraded_until_ms
.store(degrade_until, Ordering::Release);
self.suppressed_head_probes.store(0, Ordering::Release);
if !was_degraded {
tracing::warn!(
"remote HEAD probes degraded for {}s after {failures} consecutive failure(s); last error: {error}",
REMOTE_HEAD_DEGRADED_FOR.as_secs()
);
} else {
tracing::debug!("remote HEAD probe failed while degraded: {error}");
}
}
fn note_head_probe_success(&self) {
let failures = self.head_probe_failures.swap(0, Ordering::AcqRel);
let degraded_until = self.head_probe_degraded_until_ms.swap(0, Ordering::AcqRel);
let suppressed = self.suppressed_head_probes.swap(0, Ordering::AcqRel);
let now = now_millis();
if failures >= REMOTE_HEAD_FAILURE_THRESHOLD || degraded_until > now || suppressed > 0 {
tracing::info!(
"remote HEAD probes recovered after {failures} consecutive failure(s); suppressed {suppressed} probe(s) while degraded"
);
}
}
fn note_head_probe_suppressed(&self) {
self.suppressed_head_probes.fetch_add(1, Ordering::Relaxed);
}
}
impl DaemonCoordGuard {
fn new(path: PathBuf) -> Self {
Self { path }
}
}
impl Drop for DaemonCoordGuard {
fn drop(&mut self) {
let _ = std::fs::remove_file(&self.path);
}
}
impl Drop for SocketCleanupGuard {
fn drop(&mut self) {
let _ = std::fs::remove_file(&self.path);
}
}
fn daemon_state_path(socket_path: &Path) -> PathBuf {
socket_path.with_extension("state.json")
}
fn now_millis() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}
fn write_json_atomically<T: Serialize>(path: &Path, value: &T) -> Result<()> {
let parent = path
.parent()
.ok_or_else(|| anyhow::anyhow!("state file has no parent directory"))?;
std::fs::create_dir_all(parent)?;
let file_name = path
.file_name()
.ok_or_else(|| anyhow::anyhow!("state file has no file name"))?
.to_string_lossy();
let tmp_path = parent.join(format!("{file_name}.{}.tmp", std::process::id()));
let json = serde_json::to_vec(value)?;
std::fs::write(&tmp_path, json)?;
std::fs::rename(&tmp_path, path)?;
Ok(())
}
fn read_daemon_state(socket_path: &Path) -> Option<DaemonCoordState> {
let path = daemon_state_path(socket_path);
let bytes = std::fs::read(path).ok()?;
serde_json::from_slice(&bytes).ok()
}
fn daemon_state_is_recent(state: &DaemonCoordState) -> bool {
let age_ms = now_millis().saturating_sub(state.updated_at_ms);
age_ms <= DAEMON_COORD_STALE_AFTER.as_millis() as u64
}
fn client_epoch_is_newer(client_epoch: u64, daemon_epoch: u64) -> bool {
client_epoch > 0 && daemon_epoch > 0 && client_epoch > daemon_epoch
}
use crate::platform::is_process_alive as process_is_alive;
fn wait_for_run_lock_release(socket_path: &Path, timeout: Duration) -> Result<bool> {
let deadline = Instant::now() + timeout;
loop {
if !daemon_run_lock_is_held(socket_path)? {
return Ok(true);
}
if Instant::now() >= deadline {
return Ok(false);
}
std::thread::sleep(DAEMON_START_POLL_INTERVAL);
}
}
fn terminate_daemon_pid(pid: u32, socket_path: &Path) -> Result<bool> {
crate::platform::terminate_process(pid);
if wait_for_run_lock_release(socket_path, Duration::from_secs(1))? {
return Ok(true);
}
crate::platform::kill_process(pid);
wait_for_run_lock_release(socket_path, Duration::from_secs(1))
}
fn recover_unhealthy_daemon(socket_path: &Path, reason: &str) -> Result<bool> {
let run_lock_held = daemon_run_lock_is_held(socket_path)?;
if let Some(state) = read_daemon_state(socket_path) {
let state_recent = daemon_state_is_recent(&state);
if run_lock_held && process_is_alive(state.pid) {
tracing::info!(
socket = %socket_path.display(),
pid = state.pid,
?state.phase,
heartbeat_fresh = state_recent,
reason,
"terminating unhealthy daemon coordinator"
);
if !terminate_daemon_pid(state.pid, socket_path)? {
tracing::warn!(
socket = %socket_path.display(),
pid = state.pid,
heartbeat_fresh = state_recent,
reason,
"daemon process did not release run lock during recovery"
);
return Ok(false);
}
}
}
if daemon_run_lock_is_held(socket_path)? {
tracing::warn!(
socket = %socket_path.display(),
reason,
"daemon run lock still held and no recoverable coordinator state was found"
);
return Ok(false);
}
let _ = std::fs::remove_file(socket_path);
let _ = std::fs::remove_file(daemon_state_path(socket_path));
Ok(true)
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "snake_case")]
pub(crate) enum Request {
Upload(UploadJob),
Gc(GcRequest),
RemoteCheck(RemoteCheckRequest),
Stats(StatsRequest),
BatchRemoteCheck(BatchRemoteCheckRequest),
HashFiles(HashFilesRequest),
LocalLookup(LocalLookupRequest),
Prefetch(PrefetchRequest),
BuildStarted(BuildStartedRequest),
CompileStarted(CompileStartedRequest),
CompileFinished(CompileFinishedRequest),
Shutdown,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct UploadJob {
pub key: String,
pub entry_dir: String,
#[serde(default)]
pub crate_name: String,
#[serde(default)]
pub client_epoch: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct GcRequest {
pub max_age_hours: Option<u64>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct RemoteCheckRequest {
pub key: String,
pub entry_dir: String,
#[serde(default)]
pub crate_name: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct StatsRequest {
pub include_entries: bool,
pub sort_by: Option<String>,
pub event_hours: Option<u64>,
#[serde(default)]
pub client_epoch: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct BatchRemoteCheckRequest {
pub checks: Vec<RemoteCheckRequest>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct LocalLookupRequest {
pub key: String,
#[serde(default)]
pub client_epoch: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct LocalLookupReply {
pub outcome: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub meta: Option<crate::store::EntryMeta>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reason: Option<String>,
}
impl LocalLookupReply {
pub(crate) fn hit(meta: crate::store::EntryMeta) -> Self {
Self {
outcome: "hit".to_string(),
meta: Some(meta),
reason: None,
}
}
pub(crate) fn miss() -> Self {
Self {
outcome: "miss".to_string(),
meta: None,
reason: None,
}
}
pub(crate) fn fallback(reason: impl Into<String>) -> Self {
Self {
outcome: "fallback".to_string(),
meta: None,
reason: Some(reason.into()),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct HashFilesRequest {
pub files: Vec<HashFileRequest>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct HashFileRequest {
pub path: String,
pub size: i64,
pub mtime_ns: i64,
pub ctime_ns: i64,
#[serde(default)]
pub inode: i64,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct HashFileResult {
pub path: String,
pub size: i64,
pub mtime_ns: i64,
pub ctime_ns: i64,
#[serde(default)]
pub inode: i64,
#[serde(skip_serializing_if = "Option::is_none")]
pub hash: Option<String>,
#[serde(default)]
pub cache_hit: bool,
#[serde(default)]
pub bytes_hashed: u64,
#[serde(skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct PrefetchRequest {
pub keys: Vec<(String, String)>,
#[serde(default)]
pub warm_all: bool,
}
impl PrefetchRequest {
pub fn from_plan(plan: PrefetchPlan) -> Self {
Self {
warm_all: false,
keys: plan
.candidates
.into_iter()
.filter(|c| {
let ok = crate::cache_key::is_valid_cache_key(&c.cache_key)
&& crate::cache_key::is_valid_crate_name(&c.crate_name);
if !ok {
tracing::warn!(
cache_key = key_prefix(&c.cache_key),
cache_key_len = c.cache_key.len(),
"prefetch: dropping planner candidate with invalid cache_key/crate_name"
);
}
ok
})
.map(|candidate| (candidate.cache_key, candidate.crate_name))
.collect(),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct BuildStartedRequest {
#[serde(default)]
pub intent: kache_core::BuildIntent,
#[serde(default)]
pub client_epoch: u64,
#[serde(default)]
pub session_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct CompileStartedRequest {
pub crate_name: String,
#[serde(default)]
pub root: String,
pub pid: u32,
pub started_at_ms: u64,
#[serde(default)]
pub typical_ms: Option<u64>,
#[serde(default)]
pub client_epoch: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct CompileFinishedRequest {
pub pid: u32,
#[serde(default)]
pub started_at_ms: u64,
}
#[allow(dead_code)]
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct BatchResponse {
pub ok: bool,
pub results: Vec<Response>,
#[serde(skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct StatsResponse {
pub total_size: u64,
pub max_size: u64,
pub entry_count: usize,
pub entries: Option<Vec<StatsEntry>>,
pub events: EventStatsResponse,
#[serde(default)]
pub version: String,
#[serde(default)]
pub build_epoch: u64,
#[serde(default)]
pub pending_uploads: usize,
#[serde(default)]
pub active_downloads: usize,
#[serde(default)]
pub s3_concurrency_total: usize,
#[serde(default)]
pub s3_concurrency_used: usize,
#[serde(default)]
pub upload_queue_capacity: usize,
#[serde(default)]
pub uploads_completed: u64,
#[serde(default)]
pub uploads_failed: u64,
#[serde(default)]
pub uploads_skipped: u64,
#[serde(default)]
pub downloads_completed: u64,
#[serde(default)]
pub downloads_failed: u64,
#[serde(default)]
pub bytes_uploaded: u64,
#[serde(default)]
pub bytes_downloaded: u64,
#[serde(default)]
pub recent_transfers: Vec<TransferEvent>,
#[serde(default)]
pub prefetch: PrefetchStatsSnapshot,
#[serde(default)]
pub in_flight: Vec<InFlightEntry>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct InFlightEntry {
pub crate_name: String,
#[serde(default)]
pub root: String,
pub pid: u32,
pub elapsed_s: u64,
#[serde(default)]
pub typical_s: Option<u64>,
#[serde(default)]
pub eta_s: Option<u64>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq)]
pub struct PrefetchStatsSnapshot {
#[serde(default)]
pub downloads_completed: u64,
#[serde(default)]
pub bytes_downloaded: u64,
#[serde(default)]
pub keys_used: u64,
#[serde(default)]
pub keys_cancelled: u64,
#[serde(default)]
pub keys_over_budget: u64,
#[serde(default)]
pub cancelled: bool,
#[serde(default)]
pub plans_advisory: u64,
#[serde(default)]
pub plans_fallback: u64,
#[serde(default)]
pub last_plan_candidates: u64,
#[serde(default)]
pub dedup_join_waits: u64,
#[serde(default)]
pub dedup_join_wait_ms: u64,
#[serde(default)]
pub last_list_duration_ms: u64,
#[serde(default)]
pub last_list_key_count: u64,
#[serde(default)]
pub list_requests_total: u64,
#[serde(default)]
pub list_failures_total: u64,
#[serde(default)]
pub list_duration_ms_total: u64,
#[serde(default)]
pub list_keys_total: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct StatsEntry {
pub cache_key: String,
pub crate_name: String,
pub crate_type: String,
pub profile: String,
pub size: u64,
pub hit_count: u64,
pub created_at: String,
pub last_accessed: String,
pub content_hash: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct EventStatsResponse {
pub local_hits: usize,
#[serde(default)]
pub prefetch_hits: usize,
pub remote_hits: usize,
#[serde(default)]
pub dups: usize,
pub misses: usize,
pub errors: usize,
pub total_elapsed_ms: u64,
#[serde(default)]
pub hit_elapsed_ms: u64,
#[serde(default)]
pub miss_elapsed_ms: u64,
#[serde(default)]
pub hit_compile_time_ms: u64,
#[serde(default)]
pub miss_compile_time_ms: u64,
#[serde(default)]
pub store_output_blobs: u32,
#[serde(default)]
pub store_duplicate_blobs: u32,
#[serde(default)]
pub store_new_blobs: u32,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub(crate) struct Response {
pub ok: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub evicted: Option<usize>,
#[serde(default, skip_serializing_if = "is_false")]
pub skipped: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub found: Option<bool>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub prefetched: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub stats: Option<StatsResponse>,
#[serde(skip_serializing_if = "Option::is_none")]
pub batch_results: Option<Vec<Response>>,
#[serde(skip_serializing_if = "Option::is_none")]
pub hash_results: Option<Vec<HashFileResult>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub local_lookup: Option<LocalLookupReply>,
#[serde(skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
}
fn is_false(value: &bool) -> bool {
!*value
}
impl Response {
fn ok() -> Self {
Self {
ok: true,
evicted: None,
skipped: false,
found: None,
prefetched: None,
stats: None,
batch_results: None,
hash_results: None,
local_lookup: None,
error: None,
}
}
fn ok_evicted(n: usize) -> Self {
Self {
ok: true,
evicted: Some(n),
skipped: false,
found: None,
prefetched: None,
stats: None,
batch_results: None,
hash_results: None,
local_lookup: None,
error: None,
}
}
fn ok_gc_skipped() -> Self {
Self {
ok: true,
evicted: Some(0),
skipped: true,
found: None,
prefetched: None,
stats: None,
batch_results: None,
hash_results: None,
local_lookup: None,
error: None,
}
}
fn ok_stats(stats: StatsResponse) -> Self {
Self {
ok: true,
evicted: None,
skipped: false,
found: None,
prefetched: None,
stats: Some(stats),
batch_results: None,
hash_results: None,
local_lookup: None,
error: None,
}
}
fn ok_batch(results: Vec<Response>) -> Self {
Self {
ok: true,
evicted: None,
skipped: false,
found: None,
prefetched: None,
stats: None,
batch_results: Some(results),
hash_results: None,
local_lookup: None,
error: None,
}
}
fn ok_hash_results(results: Vec<HashFileResult>) -> Self {
Self {
ok: true,
evicted: None,
skipped: false,
found: None,
prefetched: None,
stats: None,
batch_results: None,
hash_results: Some(results),
local_lookup: None,
error: None,
}
}
fn found(val: bool) -> Self {
Self {
ok: true,
evicted: None,
skipped: false,
found: Some(val),
prefetched: None,
stats: None,
batch_results: None,
hash_results: None,
local_lookup: None,
error: None,
}
}
fn found_prefetched(val: bool, prefetched: bool) -> Self {
Self {
ok: true,
evicted: None,
skipped: false,
found: Some(val),
prefetched: Some(prefetched),
stats: None,
batch_results: None,
hash_results: None,
local_lookup: None,
error: None,
}
}
fn ok_local_lookup(reply: LocalLookupReply) -> Self {
Self {
local_lookup: Some(reply),
..Self::ok()
}
}
fn err(msg: impl Into<String>) -> Self {
Self {
ok: false,
evicted: None,
skipped: false,
found: None,
prefetched: None,
stats: None,
batch_results: None,
hash_results: None,
local_lookup: None,
error: Some(msg.into()),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "snake_case")]
pub enum TransferDirection {
Upload,
Download,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct TransferEvent {
#[serde(default = "default_transfer_schema")]
pub schema: u32,
pub crate_name: String,
pub direction: TransferDirection,
#[serde(default)]
pub format: String,
#[serde(default)]
pub cache_key: String,
#[serde(default)]
pub object_key: String,
pub compressed_bytes: u64,
pub elapsed_ms: u64,
#[serde(default)]
pub network_ms: u64,
#[serde(default)]
pub semaphore_wait_ms: u64,
#[serde(default)]
pub head_ms: u64,
#[serde(default)]
pub request_ms: u64,
#[serde(default)]
pub body_ms: u64,
#[serde(default)]
pub request_count: u32,
#[serde(default)]
pub original_bytes: u64,
#[serde(default)]
pub decompress_ms: u64,
#[serde(default)]
pub extract_ms: u64,
#[serde(default)]
pub disk_io_ms: u64,
#[serde(default)]
pub import_ms: u64,
#[serde(default)]
pub compression_ms: u64,
#[serde(default)]
pub head_checks_ms: u64,
#[serde(default)]
pub blobs_skipped: u32,
#[serde(default)]
pub blobs_total: u32,
pub ok: bool,
pub timestamp: u64,
}
const fn default_transfer_schema() -> u32 {
2
}
pub(crate) struct TransferCounters {
pub uploads_completed: std::sync::atomic::AtomicU64,
pub uploads_failed: std::sync::atomic::AtomicU64,
pub uploads_skipped: std::sync::atomic::AtomicU64,
pub downloads_completed: std::sync::atomic::AtomicU64,
pub downloads_failed: std::sync::atomic::AtomicU64,
pub bytes_uploaded: std::sync::atomic::AtomicU64,
pub bytes_downloaded: std::sync::atomic::AtomicU64,
}
impl TransferCounters {
fn new() -> Self {
Self {
uploads_completed: 0.into(),
uploads_failed: 0.into(),
uploads_skipped: 0.into(),
downloads_completed: 0.into(),
downloads_failed: 0.into(),
bytes_uploaded: 0.into(),
bytes_downloaded: 0.into(),
}
}
}
fn prefetch_concurrency_cap(s3_concurrency: u32) -> usize {
let total = s3_concurrency.max(1) as usize;
let reserve = (total / 4).clamp(1, 4).min(total.saturating_sub(1));
(total - reserve).max(1)
}
pub(crate) struct PrefetchStats {
pub downloads_completed: std::sync::atomic::AtomicU64,
pub bytes_downloaded: std::sync::atomic::AtomicU64,
pub keys_used: std::sync::atomic::AtomicU64,
pub keys_cancelled: std::sync::atomic::AtomicU64,
pub keys_over_budget: std::sync::atomic::AtomicU64,
pub plans_advisory: std::sync::atomic::AtomicU64,
pub plans_fallback: std::sync::atomic::AtomicU64,
pub last_plan_candidates: std::sync::atomic::AtomicU64,
pub dedup_join_waits: std::sync::atomic::AtomicU64,
pub dedup_join_wait_ms: std::sync::atomic::AtomicU64,
pub last_list_duration_ms: std::sync::atomic::AtomicU64,
pub last_list_key_count: std::sync::atomic::AtomicU64,
pub list_requests_total: std::sync::atomic::AtomicU64,
pub list_failures_total: std::sync::atomic::AtomicU64,
pub list_duration_ms_total: std::sync::atomic::AtomicU64,
pub list_keys_total: std::sync::atomic::AtomicU64,
}
impl PrefetchStats {
fn new() -> Self {
Self {
downloads_completed: 0.into(),
bytes_downloaded: 0.into(),
keys_used: 0.into(),
keys_cancelled: 0.into(),
keys_over_budget: 0.into(),
plans_advisory: 0.into(),
plans_fallback: 0.into(),
last_plan_candidates: 0.into(),
dedup_join_waits: 0.into(),
dedup_join_wait_ms: 0.into(),
last_list_duration_ms: 0.into(),
last_list_key_count: 0.into(),
list_requests_total: 0.into(),
list_failures_total: 0.into(),
list_duration_ms_total: 0.into(),
list_keys_total: 0.into(),
}
}
}
const RECENT_TRANSFERS_CAP: usize = 50;
#[derive(Debug)]
pub(crate) struct ActivePlan {
pub session_id: String,
pub plan_id: String,
pub plan_source: &'static str,
pub candidates: HashSet<String>,
pub demanded: HashSet<String>,
pub demanded_candidates: HashSet<String>,
pub downloaded: HashMap<String, u64>,
pub used: HashSet<String>,
pub cancelled: bool,
pub started_at_ms: u64,
pub last_activity_ms: u64,
pub list_requests_at_install: u64,
pub list_duration_ms_at_install: u64,
}
impl ActivePlan {
fn new(
session_id: String,
plan_id: String,
plan_source: &'static str,
candidates: HashSet<String>,
list_requests_at_install: u64,
list_duration_ms_at_install: u64,
) -> Self {
let now = epoch_ms();
Self {
session_id,
plan_id,
plan_source,
candidates,
demanded: HashSet::new(),
demanded_candidates: HashSet::new(),
downloaded: HashMap::new(),
used: HashSet::new(),
cancelled: false,
started_at_ms: now,
last_activity_ms: now,
list_requests_at_install,
list_duration_ms_at_install,
}
}
fn record_demand(&mut self, key: &str) -> bool {
self.last_activity_ms = epoch_ms();
if self.demanded.insert(key.to_string()) {
if self.candidates.contains(key) {
self.demanded_candidates.insert(key.to_string());
}
if self.downloaded.contains_key(key) {
self.used.insert(key.to_string());
}
}
if self.cancelled {
return false;
}
let downloaded_not_demanded = self
.downloaded
.keys()
.filter(|k| !self.demanded.contains(*k))
.count() as u64;
if should_cancel_prefetch(
self.demanded.len() as u64,
self.demanded_candidates.len() as u64,
downloaded_not_demanded,
) {
self.cancelled = true;
return true;
}
false
}
fn record_download(&mut self, key: &str, compressed_bytes: u64) {
self.last_activity_ms = epoch_ms();
self.downloaded.insert(key.to_string(), compressed_bytes);
if self.demanded.contains(key) {
self.used.insert(key.to_string());
}
}
fn used_bytes(&self) -> u64 {
self.used
.iter()
.filter_map(|k| self.downloaded.get(k))
.sum()
}
}
pub(crate) fn prefetch_key_budget_overflow(offered: usize, max_keys: u64) -> usize {
if max_keys == 0 {
return 0;
}
offered.saturating_sub(max_keys as usize)
}
pub(crate) fn prefetch_byte_budget_exhausted(max_bytes: u64, spent: u64) -> bool {
max_bytes > 0 && spent >= max_bytes
}
pub(crate) fn should_cancel_prefetch(
demanded: u64,
demanded_candidates: u64,
downloaded_not_demanded: u64,
) -> bool {
if demanded < 10 {
return false;
}
let upper_bound_hits = demanded_candidates + downloaded_not_demanded;
(upper_bound_hits as f64 / demanded as f64) < 0.3
}
fn epoch_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as u64
}
#[derive(Default)]
struct S3Index {
keys: HashSet<String>,
by_crate: HashMap<String, Vec<String>>,
}
pub(crate) struct S3KeyCache {
index: RwLock<Option<S3Index>>,
populated: AtomicBool,
last_populated: RwLock<Option<Instant>>,
}
impl S3KeyCache {
fn new() -> Self {
Self {
index: RwLock::new(None),
populated: AtomicBool::new(false),
last_populated: RwLock::new(None),
}
}
pub async fn age(&self) -> Option<Duration> {
let guard = self.last_populated.read().await;
guard.map(|t| t.elapsed())
}
pub async fn check(&self, key: &str) -> Option<bool> {
if !self.populated.load(Ordering::Acquire) {
return None;
}
let guard = self.index.read().await;
guard.as_ref().map(|i| i.keys.contains(key))
}
pub async fn keys_for_crate(&self, crate_name: &str) -> Vec<String> {
if !self.populated.load(Ordering::Acquire) {
return vec![];
}
let guard = self.index.read().await;
guard
.as_ref()
.and_then(|i| i.by_crate.get(crate_name))
.cloned()
.unwrap_or_default()
}
pub async fn populate(&self, keys: HashMap<String, String>) {
let mut by_crate: HashMap<String, Vec<String>> = HashMap::new();
for (cache_key, crate_name) in &keys {
by_crate
.entry(crate_name.clone())
.or_default()
.push(cache_key.clone());
}
let new_index = S3Index {
keys: keys.into_keys().collect(),
by_crate,
};
let mut guard = self.index.write().await;
*guard = Some(new_index);
drop(guard);
self.populated.store(true, Ordering::Release);
let mut ts = self.last_populated.write().await;
*ts = Some(Instant::now());
}
pub async fn insert(&self, key: String, crate_name: Option<&str>) {
let mut guard = self.index.write().await;
if let Some(index) = guard.as_mut() {
index.keys.insert(key.clone());
if let Some(name) = crate_name {
index
.by_crate
.entry(name.to_string())
.or_default()
.push(key);
}
}
}
pub async fn remove(&self, key: &str) {
let mut guard = self.index.write().await;
if let Some(index) = guard.as_mut() {
index.keys.remove(key);
for keys in index.by_crate.values_mut() {
keys.retain(|k| k != key);
}
}
}
}
pub(crate) struct Daemon {
config: Config,
store: OnceLock<Mutex<Store>>,
local_hit: OnceLock<crate::daemon_local::LocalHitService>,
remote_backend: tokio::sync::OnceCell<Arc<dyn crate::remote_backend::RemoteBackend>>,
key_cache: Arc<S3KeyCache>,
remote_health: Arc<RemoteHealth>,
s3_semaphore: Arc<tokio::sync::Semaphore>,
upload_tx: Mutex<Option<tokio::sync::mpsc::UnboundedSender<UploadJob>>>,
upload_queue_closed: AtomicBool,
pending_uploads: Arc<RwLock<HashSet<String>>>,
downloading: Arc<RwLock<HashMap<String, Arc<Notify>>>>,
warming_tx: tokio::sync::watch::Sender<bool>,
prefetched_keys: Arc<RwLock<HashSet<String>>>,
prefetch_cancel: tokio::sync::watch::Sender<bool>,
prefetch_stats: PrefetchStats,
prefetch_gate: Arc<tokio::sync::Semaphore>,
prefetch_used_keys: Arc<RwLock<HashSet<String>>>,
active_plan: Arc<std::sync::Mutex<Option<ActivePlan>>>,
in_flight_compiles: std::sync::Mutex<HashMap<u32, CompileStartedRequest>>,
version: String,
build_epoch: u64,
transfer_counters: TransferCounters,
recent_transfers: std::sync::Mutex<std::collections::VecDeque<TransferEvent>>,
file_hash_cache: Arc<Mutex<HashMap<FileHashCacheKey, String>>>,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
struct FileHashCacheKey {
path: String,
size: i64,
mtime_ns: i64,
ctime_ns: i64,
inode: i64,
}
impl Daemon {
pub fn new(config: Config) -> Self {
let permits = config.s3_concurrency.max(1) as usize;
let (warming_tx, _) = tokio::sync::watch::channel(false);
let (prefetch_cancel, _) = tokio::sync::watch::channel(false);
Self {
store: OnceLock::new(),
local_hit: OnceLock::new(),
s3_semaphore: Arc::new(tokio::sync::Semaphore::new(permits)),
remote_backend: tokio::sync::OnceCell::new(),
key_cache: Arc::new(S3KeyCache::new()),
remote_health: Arc::new(RemoteHealth::new()),
upload_tx: Mutex::new(None),
upload_queue_closed: AtomicBool::new(false),
pending_uploads: Arc::new(RwLock::new(HashSet::new())),
downloading: Arc::new(RwLock::new(HashMap::new())),
warming_tx,
prefetched_keys: Arc::new(RwLock::new(HashSet::new())),
prefetch_cancel,
prefetch_stats: PrefetchStats::new(),
prefetch_gate: Arc::new(tokio::sync::Semaphore::new(prefetch_concurrency_cap(
config.s3_concurrency,
))),
prefetch_used_keys: Arc::new(RwLock::new(HashSet::new())),
active_plan: Arc::new(std::sync::Mutex::new(None)),
in_flight_compiles: std::sync::Mutex::new(HashMap::new()),
version: VERSION.to_string(),
build_epoch: build_epoch(),
transfer_counters: TransferCounters::new(),
recent_transfers: std::sync::Mutex::new(std::collections::VecDeque::new()),
file_hash_cache: Arc::new(Mutex::new(HashMap::new())),
config,
}
}
fn store_lock(&self) -> Result<&Mutex<Store>> {
if let Some(store) = self.store.get() {
return Ok(store);
}
let store = Store::open(&self.config)?;
let _ = self.store.set(Mutex::new(store));
self.store
.get()
.ok_or_else(|| anyhow::anyhow!("daemon store failed to initialize"))
}
pub(crate) fn with_store<T>(&self, f: impl FnOnce(&Store) -> Result<T>) -> Result<T> {
let guard = self
.store_lock()?
.lock()
.map_err(|_| anyhow::anyhow!("daemon store mutex poisoned"))?;
f(&guard)
}
pub(crate) fn entry_dir_for(&self, cache_key: &str) -> PathBuf {
debug_assert!(
crate::cache_key::is_valid_cache_key(cache_key),
"entry_dir_for called with unvalidated cache_key"
);
self.config.store_dir().join(cache_key)
}
pub(crate) fn remote_config(&self) -> Option<&crate::config::RemoteConfig> {
self.config.remote.as_ref()
}
pub(crate) async fn key_cache_keys_for_crate(&self, crate_name: &str) -> Vec<String> {
self.key_cache.keys_for_crate(crate_name).await
}
async fn wait_for_warming(&self, timeout: Duration) -> bool {
let mut rx = self.warming_tx.subscribe();
if *rx.borrow() {
return true;
}
matches!(
tokio::time::timeout(timeout, rx.changed()).await,
Ok(Ok(()))
) || *rx.borrow()
}
fn signal_warming_complete(&self) {
self.warming_tx.send_replace(true);
}
fn push_transfer_event(&self, event: TransferEvent) {
if let Err(e) = events::log_transfer(&self.config.transfer_log_path(), &event) {
tracing::warn!("failed to log transfer event: {e}");
}
if let Ok(mut q) = self.recent_transfers.lock() {
if q.len() >= RECENT_TRANSFERS_CAP {
q.pop_front();
}
q.push_back(event);
}
}
pub fn set_upload_tx(&self, tx: tokio::sync::mpsc::UnboundedSender<UploadJob>) {
*self.upload_tx.lock().expect("upload queue mutex poisoned") = Some(tx);
self.upload_queue_closed.store(false, Ordering::Relaxed);
}
fn upload_tx(&self) -> Option<tokio::sync::mpsc::UnboundedSender<UploadJob>> {
self.upload_tx
.lock()
.expect("upload queue mutex poisoned")
.clone()
}
fn close_upload_queue(&self) {
self.upload_queue_closed.store(true, Ordering::Relaxed);
self.upload_tx
.lock()
.expect("upload queue mutex poisoned")
.take();
}
pub(crate) async fn get_remote_backend(
&self,
) -> Result<&Arc<dyn crate::remote_backend::RemoteBackend>> {
self.remote_backend
.get_or_try_init(|| async {
let remote = self
.config
.remote
.as_ref()
.ok_or_else(|| anyhow::anyhow!("no remote configured"))?;
crate::remote_backend::create_backend(remote, self.config.s3_pool_idle_secs).await
})
.await
}
#[cfg(test)]
pub fn handle_request_sync(&self, req: &Request) -> Response {
match req {
Request::Gc(gc) => self.handle_gc(gc),
Request::Stats(sr) => self.handle_stats(sr),
Request::HashFiles(req) => self.handle_hash_files(req),
Request::CompileStarted(req) => self.handle_compile_started(req.clone()),
Request::CompileFinished(req) => self.handle_compile_finished(req),
Request::Upload(_)
| Request::RemoteCheck(_)
| Request::BatchRemoteCheck(_)
| Request::LocalLookup(_)
| Request::Prefetch(_)
| Request::BuildStarted(_) => {
Response::err(
"upload/remote_check/batch/local_lookup/prefetch/build_started must be handled async",
)
}
Request::Shutdown => Response::ok(),
}
}
pub fn handle_stats(&self, req: &StatsRequest) -> Response {
let (total_size, entry_count, entries) = match self.with_store(|store| {
let total_size = store.total_size().unwrap_or(0);
let entry_count = store.entry_count().unwrap_or(0);
let entries = if req.include_entries {
let sort = req.sort_by.as_deref().unwrap_or("size");
store.list_entries(sort).ok().map(|list| {
list.into_iter()
.map(|e| StatsEntry {
cache_key: e.cache_key,
crate_name: e.crate_name,
crate_type: e.crate_type,
profile: e.profile,
size: e.size,
hit_count: e.hit_count,
created_at: e.created_at,
last_accessed: e.last_accessed,
content_hash: e.content_hash,
})
.collect()
})
} else {
None
};
Ok((total_size, entry_count, entries))
}) {
Ok(values) => values,
Err(e) => return Response::err(format!("store open failed: {e}")),
};
let hours = req.event_hours.unwrap_or(24);
let since = chrono::Utc::now() - chrono::Duration::hours(hours as i64);
let event_list =
events::read_events_since(&self.config.event_log_path(), since).unwrap_or_default();
let es = events::compute_stats(&event_list);
let pending_uploads = self
.pending_uploads
.try_read()
.map(|g| g.len())
.unwrap_or(0);
let active_downloads = self.downloading.try_read().map(|g| g.len()).unwrap_or(0);
let tc = &self.transfer_counters;
let ps = &self.prefetch_stats;
let s3_total = self.config.s3_concurrency.max(1) as usize;
let s3_used = s3_total - self.s3_semaphore.available_permits();
let recent_transfers = self
.recent_transfers
.try_lock()
.map(|q| q.iter().cloned().collect())
.unwrap_or_default();
let in_flight = self.in_flight_snapshot();
Response::ok_stats(StatsResponse {
total_size,
max_size: self.config.max_size,
entry_count,
entries,
events: EventStatsResponse {
local_hits: es.local_hits,
prefetch_hits: es.prefetch_hits,
remote_hits: es.remote_hits,
dups: es.dups,
misses: es.misses,
errors: es.errors,
total_elapsed_ms: es.total_elapsed_ms,
hit_elapsed_ms: es.hit_elapsed_ms,
miss_elapsed_ms: es.miss_elapsed_ms,
hit_compile_time_ms: es.hit_compile_time_ms,
miss_compile_time_ms: es.miss_compile_time_ms,
store_output_blobs: es.store_output_blobs,
store_duplicate_blobs: es.store_duplicate_blobs,
store_new_blobs: es.store_new_blobs,
},
version: self.version.clone(),
build_epoch: self.build_epoch,
pending_uploads,
active_downloads,
s3_concurrency_total: s3_total,
s3_concurrency_used: s3_used,
upload_queue_capacity: 0,
uploads_completed: tc.uploads_completed.load(Ordering::Relaxed),
uploads_failed: tc.uploads_failed.load(Ordering::Relaxed),
uploads_skipped: tc.uploads_skipped.load(Ordering::Relaxed),
downloads_completed: tc.downloads_completed.load(Ordering::Relaxed),
downloads_failed: tc.downloads_failed.load(Ordering::Relaxed),
bytes_uploaded: tc.bytes_uploaded.load(Ordering::Relaxed),
bytes_downloaded: tc.bytes_downloaded.load(Ordering::Relaxed),
recent_transfers,
prefetch: PrefetchStatsSnapshot {
downloads_completed: ps.downloads_completed.load(Ordering::Relaxed),
bytes_downloaded: ps.bytes_downloaded.load(Ordering::Relaxed),
keys_used: ps.keys_used.load(Ordering::Relaxed),
keys_cancelled: ps.keys_cancelled.load(Ordering::Relaxed),
keys_over_budget: ps.keys_over_budget.load(Ordering::Relaxed),
cancelled: *self.prefetch_cancel.borrow(),
plans_advisory: ps.plans_advisory.load(Ordering::Relaxed),
plans_fallback: ps.plans_fallback.load(Ordering::Relaxed),
last_plan_candidates: ps.last_plan_candidates.load(Ordering::Relaxed),
dedup_join_waits: ps.dedup_join_waits.load(Ordering::Relaxed),
dedup_join_wait_ms: ps.dedup_join_wait_ms.load(Ordering::Relaxed),
last_list_duration_ms: ps.last_list_duration_ms.load(Ordering::Relaxed),
last_list_key_count: ps.last_list_key_count.load(Ordering::Relaxed),
list_requests_total: ps.list_requests_total.load(Ordering::Relaxed),
list_failures_total: ps.list_failures_total.load(Ordering::Relaxed),
list_duration_ms_total: ps.list_duration_ms_total.load(Ordering::Relaxed),
list_keys_total: ps.list_keys_total.load(Ordering::Relaxed),
},
in_flight,
})
}
pub fn handle_compile_started(&self, req: CompileStartedRequest) -> Response {
if let Ok(mut map) = self.in_flight_compiles.lock() {
prune_in_flight(&mut map);
map.insert(req.pid, req);
}
Response::ok()
}
pub fn handle_compile_finished(&self, req: &CompileFinishedRequest) -> Response {
if let Ok(mut map) = self.in_flight_compiles.lock()
&& let Some(entry) = map.get(&req.pid)
&& (req.started_at_ms == 0 || entry.started_at_ms == req.started_at_ms)
{
map.remove(&req.pid);
}
Response::ok()
}
fn in_flight_snapshot(&self) -> Vec<InFlightEntry> {
let Ok(mut map) = self.in_flight_compiles.lock() else {
return Vec::new();
};
prune_in_flight(&mut map);
let now_ms = unix_ms();
let mut entries: Vec<InFlightEntry> = map
.values()
.map(|c| {
let elapsed_s = now_ms.saturating_sub(c.started_at_ms) / 1000;
let typical_s = c.typical_ms.map(|ms| ms.div_ceil(1000));
InFlightEntry {
crate_name: c.crate_name.clone(),
root: c.root.clone(),
pid: c.pid,
elapsed_s,
typical_s,
eta_s: typical_s.map(|t| t.saturating_sub(elapsed_s)),
}
})
.collect();
entries.sort_by_key(|e| std::cmp::Reverse(e.elapsed_s));
entries
}
pub fn handle_hash_files(&self, req: &HashFilesRequest) -> Response {
let mut results = Vec::with_capacity(req.files.len());
for file in &req.files {
let key = FileHashCacheKey {
path: file.path.clone(),
size: file.size,
mtime_ns: file.mtime_ns,
ctime_ns: file.ctime_ns,
inode: file.inode,
};
if let Ok(cache) = self.file_hash_cache.lock()
&& let Some(hash) = cache.get(&key).cloned()
{
results.push(HashFileResult {
path: file.path.clone(),
size: file.size,
mtime_ns: file.mtime_ns,
ctime_ns: file.ctime_ns,
inode: file.inode,
hash: Some(hash),
cache_hit: true,
bytes_hashed: 0,
error: None,
});
continue;
}
match std::fs::metadata(&file.path) {
Ok(metadata)
if i64::try_from(metadata.len()).unwrap_or(i64::MAX) == file.size
&& crate::cache_key::metadata_mtime_ns(&metadata) == file.mtime_ns
&& crate::cache_key::metadata_ctime_ns(&metadata) == file.ctime_ns
&& crate::cache_key::metadata_inode(&metadata) == file.inode => {}
Ok(_) => {
results.push(HashFileResult {
path: file.path.clone(),
size: file.size,
mtime_ns: file.mtime_ns,
ctime_ns: file.ctime_ns,
inode: file.inode,
hash: None,
cache_hit: false,
bytes_hashed: 0,
error: Some("file metadata changed before hashing".into()),
});
continue;
}
Err(e) => {
results.push(HashFileResult {
path: file.path.clone(),
size: file.size,
mtime_ns: file.mtime_ns,
ctime_ns: file.ctime_ns,
inode: file.inode,
hash: None,
cache_hit: false,
bytes_hashed: 0,
error: Some(e.to_string()),
});
continue;
}
}
let path = Path::new(&file.path);
let computed: anyhow::Result<(String, bool, u64)> =
match self.with_store(|store| Ok(store.file_hash_lookup(path))) {
Ok(crate::cache_key::FileHashLookup::Hit(hash)) => Ok((hash, true, 0)),
Ok(crate::cache_key::FileHashLookup::NeedsHash(fp)) => {
crate::cache_key::hash_file(path).map(|hash| {
let _ = self.with_store(|store| {
store.file_hash_record(&fp, &hash);
Ok(())
});
(hash, false, file.size.max(0) as u64)
})
}
Ok(crate::cache_key::FileHashLookup::Uncacheable) => {
crate::cache_key::hash_file(path)
.map(|hash| (hash, false, file.size.max(0) as u64))
}
Err(e) => Err(e),
};
match computed {
Ok((hash, cache_hit, bytes_hashed)) => {
if let Ok(mut cache) = self.file_hash_cache.lock() {
if cache.len() >= FILE_HASH_MEMORY_CACHE_CAP {
cache.clear();
}
cache.insert(key, hash.clone());
}
results.push(HashFileResult {
path: file.path.clone(),
size: file.size,
mtime_ns: file.mtime_ns,
ctime_ns: file.ctime_ns,
inode: file.inode,
hash: Some(hash),
cache_hit,
bytes_hashed,
error: None,
});
}
Err(e) => results.push(HashFileResult {
path: file.path.clone(),
size: file.size,
mtime_ns: file.mtime_ns,
ctime_ns: file.ctime_ns,
inode: file.inode,
hash: None,
cache_hit: false,
bytes_hashed: 0,
error: Some(e.to_string()),
}),
}
}
Response::ok_hash_results(results)
}
pub async fn handle_local_lookup(self: &Arc<Self>, req: &LocalLookupRequest) -> Response {
if !crate::cache_key::is_valid_cache_key(&req.key) {
return Response::err("invalid cache key");
}
let reply = tokio::time::timeout(crate::daemon_local::LOCAL_LOOKUP_DEADLINE, async {
if self.local_hit.get().is_none() {
let daemon = Arc::clone(self);
let _ = tokio::task::spawn_blocking(move || {
match crate::daemon_local::LocalHitService::new(&daemon.config) {
Ok(svc) => {
let _ = daemon.local_hit.set(svc);
}
Err(e) => tracing::warn!("local-hit service init failed: {e:#}"),
}
})
.await;
}
match self.local_hit.get() {
Some(service) => service.lookup(&req.key).await,
None => LocalLookupReply::fallback("service unavailable"),
}
})
.await
.unwrap_or_else(|_| LocalLookupReply::fallback("deadline exceeded"));
Response::ok_local_lookup(reply)
}
pub fn handle_gc(&self, req: &GcRequest) -> Response {
match self.run_gc(req.max_age_hours) {
Ok(stats) if stats.skipped => Response::ok_gc_skipped(),
Ok(stats) => Response::ok_evicted(stats.entries_evicted),
Err(e) => Response::err(format!("gc failed: {e}")),
}
}
pub async fn handle_upload(&self, job: &UploadJob) -> Response {
if self.config.remote_readonly {
tracing::debug!(
crate_name = job.crate_name,
key = key_prefix(&job.key),
"remote uploads disabled (read-only mode)"
);
return Response::ok();
}
if self.config.remote.is_none() {
return Response::err("no remote configured");
}
if let Some(tx) = self.upload_tx() {
{
let mut pending = self.pending_uploads.write().await;
if !pending.insert(job.key.clone()) {
return Response::ok(); }
}
return match tx.send(job.clone()) {
Ok(()) => Response::ok(),
Err(_) => {
self.pending_uploads.write().await.remove(&job.key);
Response::err("upload queue closed")
}
};
}
if self.upload_queue_closed.load(Ordering::Relaxed) {
return Response::err("upload queue closed");
}
let Ok(_permit) = self.s3_semaphore.acquire().await else {
return Response::err("remote semaphore closed");
};
self.do_upload(job).await
}
pub async fn do_upload(&self, job: &UploadJob) -> Response {
let key_short = key_prefix(&job.key);
if self.config.remote_readonly {
tracing::debug!(
crate_name = job.crate_name,
key = key_short,
"skipping upload (read-only mode)"
);
return Response::ok();
}
let Some(remote) = &self.config.remote else {
return Response::err("no remote configured");
};
let backend = match self.get_remote_backend().await {
Ok(b) => b,
Err(e) => {
tracing::warn!(
crate_name = job.crate_name,
key = key_short,
"remote backend init failed: {e:#}"
);
return Response::err(format!("remote backend init failed: {e:#}"));
}
};
let plan = crate::remote_plan::RemotePlanner::new(&self.config)
.plan(crate::remote_plan::RemoteWorkload::BackgroundUpload);
let layout = plan.layout(backend.as_ref(), remote);
let already_exists = layout
.exists_entry(&job.key, &job.crate_name)
.await
.unwrap_or(false);
if already_exists {
self.key_cache
.insert(job.key.clone(), Some(&job.crate_name))
.await;
self.transfer_counters
.uploads_skipped
.fetch_add(1, Ordering::Relaxed);
tracing::debug!(
crate_name = job.crate_name,
key = key_short,
"skipping upload — already in remote"
);
return Response::ok();
}
tracing::debug!(
crate_name = job.crate_name,
key = key_short,
remote = %remote.describe(),
"starting remote upload"
);
let entry_dir = PathBuf::from(&job.entry_dir);
let blobs_dir = self.config.store_dir().join("blobs");
let start = Instant::now();
match layout
.upload_entry(
&job.key,
&job.crate_name,
&entry_dir,
&blobs_dir,
self.config.compression_level,
)
.await
{
Ok(ul) => {
let elapsed_ms = start.elapsed().as_millis() as u64;
self.transfer_counters
.uploads_completed
.fetch_add(1, Ordering::Relaxed);
self.transfer_counters
.bytes_uploaded
.fetch_add(ul.transfer.compressed_bytes, Ordering::Relaxed);
self.push_transfer_event(TransferEvent {
schema: default_transfer_schema(),
crate_name: job.crate_name.clone(),
direction: TransferDirection::Upload,
format: ul.format.to_string(),
cache_key: job.key.clone(),
object_key: String::new(),
compressed_bytes: ul.transfer.compressed_bytes,
elapsed_ms,
network_ms: ul.transfer.network_ms,
semaphore_wait_ms: 0,
head_ms: 0,
request_ms: 0,
body_ms: 0,
request_count: 0,
original_bytes: 0,
decompress_ms: 0,
extract_ms: 0,
disk_io_ms: 0,
import_ms: 0,
compression_ms: ul.transfer.compression_ms,
head_checks_ms: ul.transfer.head_checks_ms,
blobs_skipped: 0,
blobs_total: 0,
ok: true,
timestamp: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs(),
});
self.key_cache
.insert(job.key.clone(), Some(&job.crate_name))
.await;
self.maybe_evict_after_upload();
Response::ok()
}
Err(e) => {
let elapsed_ms = start.elapsed().as_millis() as u64;
self.transfer_counters
.uploads_failed
.fetch_add(1, Ordering::Relaxed);
self.push_transfer_event(TransferEvent {
schema: default_transfer_schema(),
crate_name: job.crate_name.clone(),
direction: TransferDirection::Upload,
format: plan.transfer_format().to_string(),
cache_key: job.key.clone(),
object_key: String::new(),
compressed_bytes: 0,
elapsed_ms,
network_ms: 0,
semaphore_wait_ms: 0,
head_ms: 0,
request_ms: 0,
body_ms: 0,
request_count: 0,
original_bytes: 0,
decompress_ms: 0,
extract_ms: 0,
disk_io_ms: 0,
import_ms: 0,
compression_ms: 0,
head_checks_ms: 0,
blobs_skipped: 0,
blobs_total: 0,
ok: false,
timestamp: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs(),
});
tracing::warn!(
crate_name = job.crate_name,
key = key_short,
elapsed_ms,
"remote upload failed: {e:#}"
);
Response::err(format!("upload failed: {e:#}"))
}
}
}
pub async fn handle_remote_check(&self, req: &RemoteCheckRequest) -> Response {
let Some(remote) = &self.config.remote else {
return Response::err("no remote configured");
};
if !self.wait_for_warming(REMOTE_CHECK_WARMING_GRACE).await {
tracing::debug!(
"remote check: warming barrier timed out after {}ms, continuing with fallback path",
REMOTE_CHECK_WARMING_GRACE.as_millis()
);
}
{
let is_prefetched = self.prefetched_keys.read().await.contains(&req.key);
if is_prefetched {
if self
.prefetch_used_keys
.write()
.await
.insert(req.key.clone())
{
self.prefetch_stats
.keys_used
.fetch_add(1, Ordering::Relaxed);
}
}
let fire_cancel = {
let mut plan = self.active_plan.lock().unwrap_or_else(|p| p.into_inner());
match plan.as_mut() {
Some(p) => p.record_demand(&req.key),
None => false,
}
};
if fire_cancel {
let _ = self.prefetch_cancel.send(true);
let (demanded, hits) = {
let plan = self.active_plan.lock().unwrap_or_else(|p| p.into_inner());
plan.as_ref()
.map(|p| (p.demanded.len(), p.demanded_candidates.len()))
.unwrap_or((0, 0))
};
tracing::info!(
"adaptive prefetch cancel: {hits}/{demanded} demanded keys were plan candidates, cancelling remaining downloads"
);
}
}
let cn = &req.crate_name;
let mut needs_head_probe = false;
let mut head_ms = 0u64;
let mut semaphore_wait_ms = 0u64;
match self.key_cache.check(&req.key).await {
Some(false) => {
let authoritative = key_cache_miss_is_authoritative(
self.config.remote_key_cache_refresh_secs,
self.key_cache.age().await,
);
if authoritative {
tracing::debug!("key cache: {} not found (skipping remote)", &req.key);
return Response::found(false);
}
if self.remote_health.head_probe_is_degraded() {
self.remote_health.note_head_probe_suppressed();
tracing::debug!(
"key cache: {} not found but cache is stale and remote HEAD probes are degraded, treating as miss",
&req.key
);
return Response::found(false);
}
tracing::debug!(
"key cache: {} not found but cache is stale, falling through to HEAD",
&req.key
);
needs_head_probe = true;
}
Some(true) => {
tracing::debug!("key cache: {} found, skipping HEAD", &req.key);
}
None => {
if self.remote_health.head_probe_is_degraded() {
self.remote_health.note_head_probe_suppressed();
tracing::debug!(
"key cache unavailable and remote HEAD probes are degraded, treating {} as a miss",
&req.key
);
return Response::found(false);
}
needs_head_probe = true;
}
}
let backend = match self.get_remote_backend().await {
Ok(b) => b,
Err(e) => return Response::err(format!("remote backend init failed: {e}")),
};
let plan = crate::remote_plan::RemotePlanner::new(&self.config)
.plan(crate::remote_plan::RemoteWorkload::RestoreCheck);
let layout = plan.layout(backend.as_ref(), remote);
if needs_head_probe {
let semaphore_start = Instant::now();
let Ok(_permit) = self.s3_semaphore.acquire().await else {
return Response::err("remote semaphore closed");
};
semaphore_wait_ms += semaphore_start.elapsed().as_millis() as u64;
let head_start = Instant::now();
let exists = layout.exists_entry(&req.key, cn).await;
head_ms += head_start.elapsed().as_millis() as u64;
match exists {
Ok(false) => {
self.remote_health.note_head_probe_success();
return Response::found(false);
}
Ok(true) => {
self.remote_health.note_head_probe_success();
self.key_cache
.insert(req.key.clone(), Some(&req.crate_name))
.await;
}
Err(e) => {
tokio::time::sleep(Duration::from_millis(150)).await;
let retry_start = Instant::now();
let retried = layout.exists_entry(&req.key, cn).await;
head_ms += retry_start.elapsed().as_millis() as u64;
match retried {
Ok(true) => {
self.remote_health.note_head_probe_success();
self.key_cache
.insert(req.key.clone(), Some(&req.crate_name))
.await;
}
Ok(false) => {
self.remote_health.note_head_probe_success();
return Response::found(false);
}
Err(e2) => {
let error = format!(
"remote exists check failed (after retry): {e2}; first: {e}"
);
self.remote_health.note_head_probe_failure(&error);
return Response::found(false);
}
}
}
}
}
let mut claimed = true;
if let Some(mut notify) = claim_download(&self.downloading, &req.key).await {
tracing::debug!("already downloading {}, waiting for completion", &req.key);
let join_start = Instant::now();
let deadline = tokio::time::Instant::now() + DOWNLOAD_JOIN_BUDGET;
let entry_dir = PathBuf::from(&req.entry_dir);
claimed = false;
let found = loop {
let mut timed_out = false;
let mut adopt: Option<Arc<Notify>> = None;
{
let notified = notify.notified();
tokio::pin!(notified);
notified.as_mut().enable();
match self.downloading.read().await.get(&req.key).cloned() {
Some(cur) if Arc::ptr_eq(&cur, ¬ify) => {
timed_out = tokio::time::timeout_at(deadline, notified).await.is_err();
}
Some(cur) if tokio::time::Instant::now() < deadline => {
adopt = Some(cur);
}
Some(_) => timed_out = true,
None => {}
}
}
if let Some(cur) = adopt {
notify = cur;
continue;
}
if entry_dir.join("meta.json").exists() {
break true;
}
match claim_download(&self.downloading, &req.key).await {
None => {
claimed = true;
break false;
}
Some(next) => {
if timed_out {
tracing::warn!(
key = key_prefix(&req.key),
"download dedup wait exceeded {DOWNLOAD_JOIN_BUDGET:?}; \
proceeding without claim"
);
break false;
}
notify = next;
}
}
};
self.prefetch_stats
.dedup_join_waits
.fetch_add(1, Ordering::Relaxed);
self.prefetch_stats
.dedup_join_wait_ms
.fetch_add(join_start.elapsed().as_millis() as u64, Ordering::Relaxed);
if found {
let was_prefetched = self.prefetched_keys.read().await.contains(&req.key);
return Response::found_prefetched(true, was_prefetched);
}
}
let _dl_guard =
claimed.then(|| DownloadingGuard::new(self.downloading.clone(), req.key.clone()));
let semaphore_start = Instant::now();
let Ok(_permit) = self.s3_semaphore.acquire().await else {
return Response::err("remote semaphore closed");
};
semaphore_wait_ms += semaphore_start.elapsed().as_millis() as u64;
let entry_dir = PathBuf::from(&req.entry_dir);
let blobs_dir = self.config.store_dir().join("blobs");
let start = Instant::now();
let download_result = layout
.download_entry(&req.key, cn, &entry_dir, &blobs_dir)
.await;
match download_result {
Ok(dl) => {
let elapsed_ms = start.elapsed().as_millis() as u64;
let import_start = Instant::now();
let import_ms = if let Err(e) =
self.with_store(|store| store.import_restored_entry(&req.key))
{
tracing::warn!("failed to import downloaded entry {}: {e}", &req.key);
0
} else {
import_start.elapsed().as_millis() as u64
};
self.transfer_counters
.downloads_completed
.fetch_add(1, Ordering::Relaxed);
self.transfer_counters
.bytes_downloaded
.fetch_add(dl.compressed_bytes, Ordering::Relaxed);
self.push_transfer_event(TransferEvent {
schema: default_transfer_schema(),
crate_name: cn.to_string(),
direction: TransferDirection::Download,
format: dl.format.to_string(),
cache_key: req.key.clone(),
object_key: dl.object_key,
compressed_bytes: dl.compressed_bytes,
elapsed_ms,
network_ms: dl.network_ms,
semaphore_wait_ms,
head_ms,
request_ms: dl.request_ms,
body_ms: dl.body_ms,
request_count: dl.request_count,
original_bytes: dl.original_bytes,
decompress_ms: dl.decompress_ms,
extract_ms: dl.extract_ms,
disk_io_ms: dl.disk_io_ms,
import_ms,
compression_ms: 0,
head_checks_ms: 0,
blobs_skipped: dl.blobs_skipped,
blobs_total: dl.blobs_total,
ok: true,
timestamp: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs(),
});
Response::found(true)
}
Err(e)
if e.downcast_ref::<crate::remote_layout::EntryNotFound>()
.is_some() =>
{
tracing::debug!("remote GET 404 for {} — treating as miss", &req.key);
self.key_cache.remove(&req.key).await;
Response::found(false)
}
Err(e) => {
let elapsed_ms = start.elapsed().as_millis() as u64;
self.transfer_counters
.downloads_failed
.fetch_add(1, Ordering::Relaxed);
self.push_transfer_event(TransferEvent {
schema: default_transfer_schema(),
crate_name: cn.to_string(),
direction: TransferDirection::Download,
format: plan.transfer_format().to_string(),
cache_key: req.key.clone(),
object_key: String::new(),
compressed_bytes: 0,
elapsed_ms,
network_ms: 0,
semaphore_wait_ms,
head_ms,
request_ms: 0,
body_ms: 0,
request_count: 0,
original_bytes: 0,
decompress_ms: 0,
extract_ms: 0,
disk_io_ms: 0,
import_ms: 0,
compression_ms: 0,
head_checks_ms: 0,
blobs_skipped: 0,
blobs_total: 0,
ok: false,
timestamp: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs(),
});
Response::err(format!("remote download failed: {e}"))
}
}
}
pub async fn handle_batch_remote_check(
self: &Arc<Self>,
req: &BatchRemoteCheckRequest,
) -> Response {
let futures: Vec<_> = req
.checks
.iter()
.map(|check| self.handle_remote_check(check))
.collect();
let results = futures::future::join_all(futures).await;
Response::ok_batch(results)
}
pub async fn handle_prefetch(self: &Arc<Self>, req: &PrefetchRequest) -> Response {
if !self.config.prefetch_enabled {
tracing::debug!("prefetch request ignored: speculative prefetch disabled");
return Response::ok();
}
let Some(remote) = &self.config.remote else {
return Response::err("no remote configured");
};
if self.get_remote_backend().await.is_err() {
return Response::err("remote backend init failed");
}
let mut keys_to_fetch: Vec<(String, String, PathBuf)> = Vec::new();
let downloading_guard = self.downloading.read().await;
for (key, crate_name) in &req.keys {
if !crate::cache_key::is_valid_cache_key(key)
|| !crate::cache_key::is_valid_crate_name(crate_name)
{
tracing::warn!(
key = key_prefix(key),
"prefetch: skipping request key with invalid cache_key/crate_name"
);
continue;
}
let entry_dir = self.entry_dir_for(key);
if entry_dir.exists() {
continue;
}
if downloading_guard.contains_key(key) {
continue;
}
keys_to_fetch.push((key.clone(), crate_name.clone(), entry_dir));
}
drop(downloading_guard);
if req.warm_all
&& let Ok(backend) = self.get_remote_backend().await
&& let Ok(s3_keys) = crate::remote_plan::RemotePlanner::new(&self.config)
.plan(crate::remote_plan::RemoteWorkload::KeyDiscovery)
.layout(backend.as_ref(), remote)
.list_keys()
.await
{
for (key, crate_name) in s3_keys {
if !crate::cache_key::is_valid_cache_key(&key)
|| !crate::cache_key::is_valid_crate_name(&crate_name)
{
tracing::warn!(
key = key_prefix(&key),
"prefetch: skipping listing key with invalid cache_key/crate_name"
);
continue;
}
let entry_dir = self.entry_dir_for(&key);
if !entry_dir.exists() {
keys_to_fetch.push((key, crate_name, entry_dir));
}
}
}
let offered = keys_to_fetch.len();
let dropped_over_key_budget =
prefetch_key_budget_overflow(offered, self.config.prefetch_max_keys);
if dropped_over_key_budget > 0 {
keys_to_fetch.truncate(offered - dropped_over_key_budget);
}
let count = keys_to_fetch.len();
if count == 0 {
tracing::info!("prefetch: nothing to fetch");
return Response::ok();
}
if dropped_over_key_budget > 0 {
self.prefetch_stats
.keys_over_budget
.fetch_add(dropped_over_key_budget as u64, Ordering::Relaxed);
tracing::warn!(
offered,
admitted = count,
dropped = dropped_over_key_budget,
max_keys = self.config.prefetch_max_keys,
"prefetch: plan truncated by the key budget"
);
}
let daemon = Arc::clone(self);
let remote_config = remote.clone();
let cancel_rx = self.prefetch_cancel.subscribe();
tokio::spawn(async move {
let mut in_flight = futures::stream::FuturesUnordered::new();
let max_concurrent = prefetch_concurrency_cap(daemon.config.s3_concurrency);
let byte_budget = daemon.config.prefetch_max_bytes;
let bytes_at_start = daemon
.prefetch_stats
.bytes_downloaded
.load(Ordering::Relaxed);
let deadline = match daemon.config.prefetch_deadline_secs {
0 => None,
secs => Some(Instant::now() + Duration::from_secs(secs)),
};
let mut keys_iter = keys_to_fetch.into_iter().peekable();
while let Some((key, crate_name, entry_dir)) = keys_iter.next() {
if let Some(deadline) = deadline
&& Instant::now() >= deadline
{
let dropped = 1 + keys_iter.count() as u64;
daemon
.prefetch_stats
.keys_over_budget
.fetch_add(dropped, Ordering::Relaxed);
tracing::warn!(
dropped,
deadline_secs = daemon.config.prefetch_deadline_secs,
"prefetch: plan truncated by the time budget"
);
break;
}
{
let spent = daemon
.prefetch_stats
.bytes_downloaded
.load(Ordering::Relaxed)
.saturating_sub(bytes_at_start);
if prefetch_byte_budget_exhausted(byte_budget, spent) {
let dropped = 1 + keys_iter.count() as u64;
daemon
.prefetch_stats
.keys_over_budget
.fetch_add(dropped, Ordering::Relaxed);
tracing::warn!(
dropped,
spent_bytes = spent,
max_bytes = byte_budget,
in_flight = in_flight.len(),
"prefetch: plan truncated by the byte budget (soft: in-flight \
downloads still finish)"
);
break;
}
}
if *cancel_rx.borrow() {
tracing::info!("prefetch: cancelled by adaptive hit-rate check");
let cancelled = 1 + keys_iter.count() as u64;
daemon
.prefetch_stats
.keys_cancelled
.fetch_add(cancelled, Ordering::Relaxed);
break;
}
while in_flight.len() >= max_concurrent {
use futures::StreamExt;
in_flight.next().await;
}
let sem = daemon.s3_semaphore.clone();
let d = daemon.clone();
let remote_cfg = remote_config.clone();
let download_plan = crate::remote_plan::RemotePlanner::new(&d.config)
.plan(crate::remote_plan::RemoteWorkload::Prefetch);
in_flight.push(tokio::spawn(async move {
if entry_dir.exists() {
return;
}
let Ok(_gate) = d.prefetch_gate.clone().acquire_owned().await else {
tracing::warn!("prefetch: gate closed for {}", key);
return;
};
let semaphore_start = Instant::now();
let Ok(_permit) = sem.acquire().await else {
tracing::warn!("prefetch: semaphore closed for {}", key);
return;
};
let semaphore_wait_ms = semaphore_start.elapsed().as_millis() as u64;
if claim_download(&d.downloading, &key).await.is_some() {
tracing::debug!("prefetch: {} already claimed, skipping", key_prefix(&key));
return;
}
let _dl_guard = DownloadingGuard::new(d.downloading.clone(), key.clone());
if entry_dir.exists() {
return;
}
let backend = match d.get_remote_backend().await {
Ok(b) => b,
Err(_) => {
return;
}
};
let blobs_dir = d.config.store_dir().join("blobs");
let start = Instant::now();
let download_result = download_plan
.layout(backend.as_ref(), &remote_cfg)
.download_entry(&key, &crate_name, &entry_dir, &blobs_dir)
.await;
match download_result {
Ok(dl) => {
let elapsed_ms = start.elapsed().as_millis() as u64;
let import_start = Instant::now();
let import_ms = if let Err(e) =
d.with_store(|store| store.import_restored_entry(&key))
{
tracing::warn!("prefetch import failed for {}: {e}", key);
0
} else {
import_start.elapsed().as_millis() as u64
};
d.transfer_counters
.downloads_completed
.fetch_add(1, Ordering::Relaxed);
d.transfer_counters
.bytes_downloaded
.fetch_add(dl.compressed_bytes, Ordering::Relaxed);
d.prefetch_stats
.downloads_completed
.fetch_add(1, Ordering::Relaxed);
d.prefetch_stats
.bytes_downloaded
.fetch_add(dl.compressed_bytes, Ordering::Relaxed);
{
let mut plan =
d.active_plan.lock().unwrap_or_else(|p| p.into_inner());
if let Some(p) = plan.as_mut() {
p.record_download(&key, dl.compressed_bytes);
}
}
d.push_transfer_event(TransferEvent {
schema: default_transfer_schema(),
crate_name: crate_name.clone(),
direction: TransferDirection::Download,
format: dl.format.to_string(),
cache_key: key.clone(),
object_key: dl.object_key,
compressed_bytes: dl.compressed_bytes,
elapsed_ms,
network_ms: dl.network_ms,
semaphore_wait_ms,
head_ms: 0,
request_ms: dl.request_ms,
body_ms: dl.body_ms,
request_count: dl.request_count,
original_bytes: dl.original_bytes,
decompress_ms: dl.decompress_ms,
extract_ms: dl.extract_ms,
disk_io_ms: dl.disk_io_ms,
import_ms,
compression_ms: 0,
head_checks_ms: 0,
blobs_skipped: dl.blobs_skipped,
blobs_total: dl.blobs_total,
ok: true,
timestamp: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs(),
});
{
const MAX_PREFETCHED_KEYS: usize = 50_000;
let mut pf = d.prefetched_keys.write().await;
if pf.len() >= MAX_PREFETCHED_KEYS {
pf.clear();
d.prefetch_used_keys.write().await.clear();
}
pf.insert(key.clone());
}
}
Err(e) => {
let elapsed_ms = start.elapsed().as_millis() as u64;
d.transfer_counters
.downloads_failed
.fetch_add(1, Ordering::Relaxed);
d.push_transfer_event(TransferEvent {
schema: default_transfer_schema(),
crate_name: crate_name.clone(),
direction: TransferDirection::Download,
format: download_plan.transfer_format().to_string(),
cache_key: key.clone(),
object_key: String::new(),
compressed_bytes: 0,
elapsed_ms,
network_ms: 0,
semaphore_wait_ms,
head_ms: 0,
request_ms: 0,
body_ms: 0,
request_count: 0,
original_bytes: 0,
decompress_ms: 0,
extract_ms: 0,
disk_io_ms: 0,
import_ms: 0,
compression_ms: 0,
head_checks_ms: 0,
blobs_skipped: 0,
blobs_total: 0,
ok: false,
timestamp: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs(),
});
tracing::warn!("prefetch download failed for {}: {e}", key);
}
}
}));
}
use futures::StreamExt;
while in_flight.next().await.is_some() {}
tracing::info!("prefetch: completed {} downloads", count);
});
tracing::info!("prefetch: queued {} downloads", count);
Response::ok()
}
fn install_plan(
&self,
session_id: &str,
plan_id: &str,
plan_source: &'static str,
candidates: impl Iterator<Item = String>,
) {
let _ = self.prefetch_cancel.send(false);
let plan = ActivePlan::new(
session_id.to_string(),
plan_id.to_string(),
plan_source,
candidates.collect(),
self.prefetch_stats
.list_requests_total
.load(Ordering::Relaxed),
self.prefetch_stats
.list_duration_ms_total
.load(Ordering::Relaxed),
);
let prev = {
let mut slot = self.active_plan.lock().unwrap_or_else(|p| p.into_inner());
slot.replace(plan)
};
if let Some(prev) = prev {
self.emit_plan_summary(prev, "superseded");
}
}
pub(crate) fn finalize_inactive_plan(&self, inactivity_ms: u64) {
let prev = {
let mut slot = self.active_plan.lock().unwrap_or_else(|p| p.into_inner());
match slot.as_ref() {
Some(p) if epoch_ms().saturating_sub(p.last_activity_ms) >= inactivity_ms => {
slot.take()
}
_ => None,
}
};
if let Some(prev) = prev {
self.emit_plan_summary(prev, "inactivity");
}
}
fn emit_plan_summary(&self, plan: ActivePlan, closure_reason: &str) {
let used_bytes = plan.used_bytes();
let downloaded_bytes: u64 = plan.downloaded.values().sum();
let event = crate::events::BuildSummaryEvent {
ts: chrono::Utc::now(),
schema: 1,
session_id: plan.session_id,
root: String::new(),
plan_source: plan.plan_source.to_string(),
plan_id: plan.plan_id,
closure_reason: closure_reason.to_string(),
started_at_ms: plan.started_at_ms,
last_activity_ms: plan.last_activity_ms,
candidate_keys: plan.candidates.len() as u64,
downloaded_keys: plan.downloaded.len() as u64,
downloaded_bytes,
used_keys: plan.used.len() as u64,
used_bytes,
demanded_keys: plan.demanded.len() as u64,
demanded_candidate_keys: plan.demanded_candidates.len() as u64,
cancelled: plan.cancelled,
list_requests: self
.prefetch_stats
.list_requests_total
.load(Ordering::Relaxed)
.saturating_sub(plan.list_requests_at_install),
list_duration_ms: self
.prefetch_stats
.list_duration_ms_total
.load(Ordering::Relaxed)
.saturating_sub(plan.list_duration_ms_at_install),
};
let path = self.config.summary_log_path();
if let Err(e) = crate::events::log_summary(&path, &event) {
tracing::debug!("failed to write build summary: {e}");
}
}
pub async fn handle_build_started(self: &Arc<Self>, req: &BuildStartedRequest) -> Response {
let Some(_remote) = &self.config.remote else {
return Response::err("no remote configured");
};
if speculative_prefetch_disabled(self.config.prefetch_enabled) {
tracing::debug!("build-started: speculative prefetch disabled");
return Response::ok();
}
{
let prev = {
let mut slot = self.active_plan.lock().unwrap_or_else(|p| p.into_inner());
match slot.as_ref() {
Some(p) if p.session_id != req.session_id => slot.take(),
_ => None,
}
};
if let Some(prev) = prev {
self.emit_plan_summary(prev, "superseded");
}
}
match crate::planner_client::resolve_prefetch_plan(&req.intent).await {
Ok(Some(plan)) => {
let plan_id = plan.plan_id.clone();
let planner = plan.planner.clone();
match plan.disposition {
PrefetchDisposition::Execute if plan.candidates.is_empty() => {
tracing::warn!(
plan_id = ?plan_id,
planner = ?planner,
"build-started: planner returned execute with no candidates, falling back to local planning"
);
}
PrefetchDisposition::Execute => {
let prefetch_req = PrefetchRequest::from_plan(plan);
self.install_plan(
&req.session_id,
plan_id.as_deref().unwrap_or(""),
"advisory",
prefetch_req.keys.iter().map(|(k, _)| k.clone()),
);
let resp = self.handle_prefetch(&prefetch_req).await;
if resp.ok {
self.prefetch_stats
.plans_advisory
.fetch_add(1, Ordering::Relaxed);
self.prefetch_stats
.last_plan_candidates
.store(prefetch_req.keys.len() as u64, Ordering::Relaxed);
tracing::info!(
plan_id = ?plan_id,
planner = ?planner,
candidate_count = prefetch_req.keys.len(),
"build-started: using advisory planner plan"
);
return resp;
}
tracing::warn!(
plan_id = ?plan_id,
planner = ?planner,
"build-started: planner plan execution failed, falling back to local planning"
);
}
PrefetchDisposition::UseFallback => {
tracing::debug!(
plan_id = ?plan_id,
planner = ?planner,
"build-started: planner requested fallback to local planning"
);
}
PrefetchDisposition::DoNothing => {
tracing::info!(
plan_id = ?plan_id,
planner = ?planner,
"build-started: planner explicitly requested no prefetch"
);
return Response::ok();
}
}
}
Ok(None) => {}
Err(e) => {
tracing::warn!(
"build-started: planner lookup failed, falling back to local planning: {e}"
);
}
}
let fallback_plan =
match crate::fallback_planner::build_prefetch_plan(self, &req.intent).await {
Ok(plan) => plan,
Err(e) => return Response::err(format!("fallback planning failed: {e}")),
};
if fallback_plan.candidates.is_empty() {
tracing::debug!(
"build-started: nothing to prefetch ({} crate names checked)",
req.intent.crate_names.len()
);
return Response::ok();
}
tracing::info!(
"build-started: using fallback planner with {} candidates for {} crates",
fallback_plan.candidates.len(),
req.intent.crate_names.len()
);
self.prefetch_stats
.plans_fallback
.fetch_add(1, Ordering::Relaxed);
self.prefetch_stats
.last_plan_candidates
.store(fallback_plan.candidates.len() as u64, Ordering::Relaxed);
let prefetch_req = PrefetchRequest::from_plan(fallback_plan);
self.install_plan(
&req.session_id,
"",
"fallback",
prefetch_req.keys.iter().map(|(k, _)| k.clone()),
);
self.handle_prefetch(&prefetch_req).await
}
fn maybe_evict_after_upload(&self) {
let _ = self.with_store(|store| {
let _gc_lock = match store.try_gc_lock()? {
Some(lock) => lock,
None => {
tracing::debug!(
"gc.lock held by another GC; skipping upload-triggered eviction"
);
return Ok(());
}
};
let size = store.total_size()?;
if size > self.config.max_size {
tracing::info!(
"store size {} > max {}, running LRU eviction",
size,
self.config.max_size
);
let _ = store.evict();
}
Ok(())
});
}
pub fn run_gc(&self, max_age_hours: Option<u64>) -> Result<crate::store::GcStats> {
let start = Instant::now();
let _gc_lock = match self.with_store(|store| store.try_gc_lock())? {
Some(lock) => lock,
None => {
tracing::info!("gc.lock held by another GC; skipping this run");
return Ok(crate::store::GcStats {
skipped: true,
..Default::default()
});
}
};
let (dedup_stats, evict_stats, incremental_cleaned, orphan_stats) =
self.with_store(|store| {
let backfilled = store.backfill_content_hashes().unwrap_or(0);
if backfilled > 0 {
tracing::info!("backfilled {backfilled} content hashes");
}
let costs = store.backfill_compile_times().unwrap_or(0);
if costs > 0 {
tracing::info!("backfilled {costs} compile times");
}
let pruned = store
.prune_tombstones(crate::store::TOMBSTONE_RETENTION_DAYS)
.unwrap_or(0);
if let Ok((tracked, demanded)) = store.tombstone_stats()
&& tracked > 0
{
tracing::info!(
tracked,
demanded,
pruned,
demand_rate_pct = demanded * 100 / tracked.max(1),
"gc: post-eviction demand"
);
}
let dedup_stats = store.evict_duplicate_entries().unwrap_or_default();
if dedup_stats.entries_evicted > 0 {
tracing::info!("evicted {} duplicate entries", dedup_stats.entries_evicted);
}
let evict_stats = if let Some(hours) = max_age_hours {
store.evict_older_than(hours)?
} else {
store.evict()?
};
let incremental_cleaned = if self.config.clean_incremental {
store.clean_registered_incremental_dirs().unwrap_or(0)
} else {
0
};
let orphan_stats = store
.sweep_orphan_blobs(std::time::Duration::from_secs(3600))
.unwrap_or_default();
if orphan_stats.removed > 0 {
tracing::info!(
"swept {} of {} blobs as orphans ({} reclaimed)",
orphan_stats.removed,
orphan_stats.scanned,
crate::report::format_bytes(orphan_stats.bytes_reclaimed)
);
}
Ok((dedup_stats, evict_stats, incremental_cleaned, orphan_stats))
})?;
Self::clean_tool_version_caches(&self.config.cache_dir);
if incremental_cleaned > 0 {
tracing::info!("cleaned {incremental_cleaned} registered incremental dirs");
}
let stats = crate::store::GcStats {
entries_evicted: dedup_stats.entries_evicted + evict_stats.entries_evicted,
bytes_freed: dedup_stats.bytes_freed
+ evict_stats.bytes_freed
+ orphan_stats.bytes_reclaimed,
blobs_removed: dedup_stats.blobs_removed
+ evict_stats.blobs_removed
+ orphan_stats.removed,
duration_ms: start.elapsed().as_millis() as u64,
skipped: false,
entries_pinned: evict_stats.entries_pinned,
};
tracing::info!(
"gc complete: {} entries evicted, {} freed, {} blobs removed in {}ms",
stats.entries_evicted,
crate::report::format_bytes(stats.bytes_freed),
stats.blobs_removed,
stats.duration_ms,
);
let gc_stats_path = self.config.cache_dir.join("gc_stats.json");
let persisted = crate::report::GcStatsPersisted {
last_run: chrono::Utc::now().to_rfc3339(),
entries_evicted: stats.entries_evicted,
bytes_freed: stats.bytes_freed,
blobs_removed: stats.blobs_removed,
duration_ms: stats.duration_ms,
};
if let Ok(json) = serde_json::to_string_pretty(&persisted) {
let _ = std::fs::write(&gc_stats_path, json);
}
Ok(stats)
}
fn clean_tool_version_caches(cache_dir: &Path) {
let cutoff = std::time::SystemTime::now() - std::time::Duration::from_secs(7 * 24 * 3600);
let Ok(entries) = std::fs::read_dir(cache_dir) else {
return;
};
for entry in entries.flatten() {
let name = entry.file_name();
let name = name.to_string_lossy();
if (name.starts_with("rustc-ver-") || name.starts_with("linker-ver-"))
&& name.ends_with(".txt")
&& let Ok(meta) = entry.metadata()
&& let Ok(modified) = meta.modified()
&& modified < cutoff
{
let _ = std::fs::remove_file(entry.path());
}
}
}
}
pub fn run_server(config: &Config) -> Result<()> {
let socket_path = config.socket_path();
let lock_path = socket_path.with_extension("run.lock");
std::fs::create_dir_all(socket_path.parent().unwrap())?;
let lock_file = std::fs::OpenOptions::new()
.create(true)
.write(true)
.truncate(false)
.open(&lock_path)
.context("opening daemon run lock file")?;
if lock_file.try_lock().is_err() {
tracing::info!("another daemon holds the run lock, exiting");
return Ok(());
}
let _lock = lock_file;
let coord = DaemonCoordFile::for_socket(&socket_path);
coord
.write_phase(DaemonPhase::Starting)
.context("writing daemon coordinator state")?;
let _coord_guard = DaemonCoordGuard::new(coord.path.clone());
let rt = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()?;
rt.block_on(server_main(config, coord))
}
fn start_manifest_warming(daemon: &Arc<Daemon>) -> Option<tokio::task::JoinHandle<()>> {
if should_start_speculative_prefetch(
daemon.config.remote.is_some(),
daemon.config.prefetch_enabled,
) {
let manifest_daemon = daemon.clone();
Some(tokio::spawn(async move {
manifest_prefetch(&manifest_daemon).await;
manifest_daemon.signal_warming_complete();
}))
} else {
daemon.signal_warming_complete();
None
}
}
async fn server_main(config: &Config, coord: DaemonCoordFile) -> Result<()> {
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap())?;
let probe_name = socket_name(&socket_path)?;
match TokioStream::connect(probe_name).await {
Ok(_) => {
tracing::info!("another daemon is already running (socket is active), exiting cleanly",);
return Ok(());
}
Err(_) => {
let _ = std::fs::remove_file(&socket_path);
}
}
let bind_name = socket_name(&socket_path)?;
let listener = ListenerOptions::new()
.name(bind_name)
.create_tokio()
.context("binding local IPC socket")?;
let _socket_guard = SocketCleanupGuard {
path: socket_path.clone(),
};
coord
.write_phase(DaemonPhase::Ready)
.context("publishing daemon ready state")?;
tracing::info!("daemon listening on {}", socket_path.display());
#[cfg(target_os = "macos")]
let _ = crate::store::exclude_from_indexing(&config.cache_dir);
match crate::cache_fs::classify(&crate::cache_fs::probe(&config.cache_dir)) {
crate::cache_fs::CacheFsVerdict::NotLocal { name } => tracing::warn!(
cache_dir = %config.cache_dir.display(),
filesystem = %name,
"cache directory is not on host-local storage: the WAL index needs working \
file locking and a single writing machine, and can be corrupted on a shared \
or network mount. Set KACHE_CACHE_DIR to a local path; to share artifacts \
between machines use a remote cache instead."
),
verdict => tracing::debug!(
cache_dir = %config.cache_dir.display(),
?verdict,
"cache filesystem locality check"
),
}
let (buffer_tx, mut buffer_rx) = tokio::sync::mpsc::unbounded_channel::<UploadJob>();
let num_workers = (config.s3_concurrency as usize).max(1);
let (worker_tx, worker_rx) = tokio::sync::mpsc::channel::<UploadJob>(num_workers * 2);
let worker_rx = Arc::new(tokio::sync::Mutex::new(worker_rx));
let daemon_inner = Daemon::new(config.clone());
daemon_inner.set_upload_tx(buffer_tx);
let daemon = Arc::new(daemon_inner);
let enqueue_handle = tokio::spawn(async move {
while let Some(job) = buffer_rx.recv().await {
if worker_tx.send(job).await.is_err() {
break;
}
}
});
let mut upload_handles: Vec<tokio::task::JoinHandle<()>> = Vec::new();
for _ in 0..num_workers {
let rx = worker_rx.clone();
let d = daemon.clone();
upload_handles.push(tokio::spawn(async move {
while let Some(job) = rx.lock().await.recv().await {
let Ok(_permit) = d.s3_semaphore.acquire().await else {
d.transfer_counters
.uploads_failed
.fetch_add(1, Ordering::Relaxed);
tracing::error!("upload worker: semaphore closed, exiting");
break;
};
let resp = d.do_upload(&job).await;
d.pending_uploads.write().await.remove(&job.key);
if !resp.ok {
tracing::warn!(
"upload worker: {} failed: {}",
job.key,
resp.error.as_deref().unwrap_or("unknown")
);
}
}
}));
}
tracing::info!("started {} upload workers", num_workers);
let gc_daemon = daemon.clone();
let sweep_daemon = daemon.clone();
tokio::spawn(async move {
const SESSION_INACTIVITY_MS: u64 = 300_000;
let mut interval = tokio::time::interval(std::time::Duration::from_secs(60));
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
interval.tick().await;
sweep_daemon.finalize_inactive_plan(SESSION_INACTIVITY_MS);
}
});
let gc_handle = tokio::spawn(async move {
let mut interval = tokio::time::interval(std::time::Duration::from_secs(6 * 3600));
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
interval.tick().await;
tracing::info!("periodic GC sweep starting");
let gc = gc_daemon.clone();
match tokio::task::spawn_blocking(move || gc.run_gc(None)).await {
Ok(Ok(_)) => {}
Ok(Err(e)) => tracing::warn!("periodic GC failed: {e}"),
Err(e) => tracing::warn!("periodic GC task panicked: {e}"),
}
}
});
let cache_handle = if should_start_speculative_prefetch(
config.remote.is_some(),
config.prefetch_enabled,
) {
let cache_daemon = daemon.clone();
let refresh_secs = config.remote_key_cache_refresh_secs;
Some(tokio::spawn(async move {
let mut delay = std::time::Duration::from_secs(1);
for attempt in 1..=5 {
match populate_key_cache(&cache_daemon).await {
Ok(count) => {
tracing::info!("remote key cache populated: {count} keys");
break;
}
Err(e) => {
tracing::warn!(
"remote key cache population attempt {attempt}/5 failed: {e}"
);
if attempt < 5 {
tokio::time::sleep(delay).await;
delay *= 2;
}
}
}
}
if key_cache_periodic_refresh_disabled(refresh_secs) {
tracing::info!("remote key cache periodic refresh disabled");
return;
}
let mut interval = tokio::time::interval(std::time::Duration::from_secs(refresh_secs));
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
interval.tick().await; let mut consecutive_refresh_failures = 0u32;
loop {
interval.tick().await;
match populate_key_cache(&cache_daemon).await {
Ok(count) => {
if consecutive_refresh_failures > 0 {
tracing::info!(
"remote key cache refresh recovered after {consecutive_refresh_failures} failed attempt(s)"
);
consecutive_refresh_failures = 0;
}
tracing::debug!("remote key cache refreshed: {count} keys");
}
Err(e) => {
consecutive_refresh_failures += 1;
if should_warn_key_cache_refresh_failure(consecutive_refresh_failures) {
tracing::warn!(
"remote key cache refresh failed (attempt {consecutive_refresh_failures}): {e}"
);
} else {
tracing::debug!(
"remote key cache refresh failed (attempt {consecutive_refresh_failures}): {e}"
);
}
}
}
}
}))
} else {
None
};
let manifest_handle = start_manifest_warming(&daemon);
let migration_config = config.clone();
tokio::spawn(async move {
let result = tokio::task::spawn_blocking(move || {
if let Ok(store) = Store::open(&migration_config) {
store.migrate_to_blobs(|_, _| {})
} else {
Err(anyhow::anyhow!("failed to open store for migration"))
}
})
.await;
if let Ok(Ok(stats)) = result
&& stats.entries_migrated > 0
{
tracing::info!(
"background migration: migrated {} entries",
stats.entries_migrated,
);
}
});
let shutdown_flag = Arc::new(AtomicBool::new(false));
let heartbeat_coord = coord.clone();
let heartbeat_handle = tokio::spawn(async move {
let mut interval = tokio::time::interval(DAEMON_COORD_HEARTBEAT_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
interval.tick().await;
loop {
interval.tick().await;
if let Err(e) = heartbeat_coord.write_phase(DaemonPhase::Ready) {
tracing::debug!("daemon coordinator heartbeat failed: {e}");
}
}
});
let shutdown_notify = Arc::new(Notify::new());
let config_fingerprint = crate::config::config_file_fingerprint();
let config_watch_flag = Arc::clone(&shutdown_flag);
let config_watch_notify = Arc::clone(&shutdown_notify);
let config_watch_handle = tokio::spawn(async move {
let mut interval = tokio::time::interval(DAEMON_CONFIG_WATCH_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
interval.tick().await;
loop {
interval.tick().await;
if config_watch_flag.load(Ordering::Relaxed) {
break;
}
if crate::config::config_file_fingerprint() != config_fingerprint {
tracing::info!("config file changed on disk, scheduling restart to reload it");
config_watch_flag.store(true, Ordering::Relaxed);
config_watch_notify.notify_one();
break;
}
}
});
let idle_timeout = if config.daemon_idle_timeout_secs > 0 {
Some(Duration::from_secs(config.daemon_idle_timeout_secs))
} else {
None
};
accept_loop(
&listener,
&daemon,
&shutdown_flag,
&shutdown_notify,
idle_timeout,
shutdown_signal(),
)
.await;
gc_handle.abort();
if let Some(h) = cache_handle {
h.abort();
}
if let Some(h) = manifest_handle {
h.abort();
}
heartbeat_handle.abort();
config_watch_handle.abort();
daemon.close_upload_queue();
drop(daemon);
let _ = enqueue_handle.await;
let drain_deadline = tokio::time::sleep(Duration::from_secs(30));
tokio::pin!(drain_deadline);
for h in &mut upload_handles {
tokio::select! {
_ = h => {}
_ = &mut drain_deadline => {
tracing::warn!("upload drain timeout, aborting remaining workers");
break;
}
}
}
for h in upload_handles {
h.abort();
}
tracing::info!("daemon stopped");
Ok(())
}
const ACCEPT_LOOP_IDLE_TICK: Duration = Duration::from_secs(60);
const DOWNLOAD_JOIN_BUDGET: Duration = Duration::from_secs(30);
async fn claim_download(
downloading: &RwLock<HashMap<String, Arc<Notify>>>,
key: &str,
) -> Option<Arc<Notify>> {
use std::collections::hash_map::Entry;
match downloading.write().await.entry(key.to_string()) {
Entry::Occupied(e) => Some(e.get().clone()),
Entry::Vacant(v) => {
v.insert(Arc::new(Notify::new()));
None
}
}
}
struct DownloadingGuard {
map: Arc<RwLock<HashMap<String, Arc<Notify>>>>,
key: String,
}
impl DownloadingGuard {
fn new(map: Arc<RwLock<HashMap<String, Arc<Notify>>>>, key: String) -> Self {
Self { map, key }
}
}
impl Drop for DownloadingGuard {
fn drop(&mut self) {
let key = std::mem::take(&mut self.key);
if let Ok(mut g) = self.map.try_write() {
let notify = g.remove(&key);
drop(g);
if let Some(notify) = notify {
notify.notify_waiters();
}
return;
}
let map = self.map.clone();
if let Ok(handle) = tokio::runtime::Handle::try_current() {
handle.spawn(async move {
let notify = map.write().await.remove(&key);
if let Some(notify) = notify {
notify.notify_waiters();
}
});
}
}
}
async fn accept_loop(
listener: &TokioListener,
daemon: &Arc<Daemon>,
shutdown_flag: &Arc<AtomicBool>,
shutdown_notify: &Arc<Notify>,
idle_timeout: Option<Duration>,
shutdown_signal: impl std::future::Future<Output = ()>,
) {
tokio::pin!(shutdown_signal);
let mut last_activity = Instant::now();
const MAX_CONCURRENT_CONNECTIONS: usize = 128;
let conn_limiter = Arc::new(tokio::sync::Semaphore::new(MAX_CONCURRENT_CONNECTIONS));
loop {
if shutdown_flag.load(Ordering::Relaxed) {
tracing::info!("shutdown requested via protocol, draining...");
break;
}
if let Some(timeout) = idle_timeout
&& last_activity.elapsed() > timeout
{
tracing::info!("daemon idle for {:?}, shutting down", timeout);
break;
}
tokio::select! {
accept = listener.accept() => {
match accept {
Ok(stream) => {
last_activity = Instant::now();
let d = daemon.clone();
let flag = shutdown_flag.clone();
let notify = shutdown_notify.clone();
let limiter = conn_limiter.clone();
tokio::spawn(async move {
let _permit = limiter.acquire_owned().await.ok();
if let Err(e) = handle_connection(stream, &d, &flag, ¬ify).await {
if e.downcast_ref::<std::io::Error>()
.is_some_and(is_client_disconnect)
{
tracing::debug!("connection handler: client disconnected: {e}");
} else {
tracing::warn!("connection handler error: {e}");
}
}
});
}
Err(e) => {
tracing::warn!("accept error: {e}");
}
}
}
_ = shutdown_notify.notified() => {}
_ = tokio::time::sleep(ACCEPT_LOOP_IDLE_TICK) => {}
_ = &mut shutdown_signal => {
tracing::info!("shutdown signal received, draining...");
break;
}
}
}
}
async fn populate_key_cache(daemon: &Daemon) -> Result<usize> {
let backend = daemon.get_remote_backend().await?;
let remote = daemon
.config
.remote
.as_ref()
.ok_or_else(|| anyhow::anyhow!("no remote configured"))?;
let list_start = Instant::now();
daemon
.prefetch_stats
.list_requests_total
.fetch_add(1, Ordering::Relaxed);
let keys = match crate::remote_plan::RemotePlanner::new(&daemon.config)
.plan(crate::remote_plan::RemoteWorkload::KeyDiscovery)
.layout(backend.as_ref(), remote)
.list_keys()
.await
{
Ok(keys) => keys,
Err(e) => {
daemon
.prefetch_stats
.list_failures_total
.fetch_add(1, Ordering::Relaxed);
daemon
.prefetch_stats
.list_duration_ms_total
.fetch_add(list_start.elapsed().as_millis() as u64, Ordering::Relaxed);
return Err(e);
}
};
let list_elapsed_ms = list_start.elapsed().as_millis() as u64;
daemon
.prefetch_stats
.last_list_duration_ms
.store(list_elapsed_ms, Ordering::Relaxed);
daemon
.prefetch_stats
.last_list_key_count
.store(keys.len() as u64, Ordering::Relaxed);
daemon
.prefetch_stats
.list_duration_ms_total
.fetch_add(list_elapsed_ms, Ordering::Relaxed);
daemon
.prefetch_stats
.list_keys_total
.fetch_add(keys.len() as u64, Ordering::Relaxed);
let count = keys.len();
daemon.key_cache.populate(keys).await;
Ok(count)
}
async fn manifest_prefetch(daemon: &Arc<Daemon>) {
let Some(remote) = &daemon.config.remote else {
return;
};
let backend = match daemon.get_remote_backend().await {
Ok(b) => b,
Err(e) => {
tracing::warn!("manifest prefetch: remote backend init failed: {e}");
return;
}
};
if let Ok(namespace) = std::env::var("KACHE_NAMESPACE") {
let lock_path = std::path::Path::new("Cargo.lock");
if lock_path.exists() {
match shard_prefetch(daemon, backend, &remote.prefix, &namespace, lock_path).await {
Ok(n) => {
tracing::info!("shard prefetch: queued {n} keys from shards");
return;
}
Err(e) => {
tracing::warn!(
"shard prefetch failed, falling back to monolithic build manifest: {e}"
);
}
}
} else {
tracing::info!(
"KACHE_NAMESPACE set but no Cargo.lock found, falling back to monolithic build manifest"
);
}
}
monolithic_manifest_prefetch(daemon, backend.as_ref(), remote).await;
}
async fn shard_prefetch(
daemon: &Arc<Daemon>,
backend: &Arc<dyn crate::remote_backend::RemoteBackend>,
prefix: &str,
namespace: &str,
lock_path: &std::path::Path,
) -> anyhow::Result<usize> {
let deps = crate::shards::parse_cargo_lock(lock_path)?;
shard_prefetch_for_deps(daemon, backend, prefix, namespace, &deps).await
}
async fn shard_prefetch_for_deps(
daemon: &Arc<Daemon>,
backend: &Arc<dyn crate::remote_backend::RemoteBackend>,
prefix: &str,
namespace: &str,
deps: &[(String, String)],
) -> anyhow::Result<usize> {
let shard_set = crate::shards::compute_shards(namespace, deps);
tracing::info!(
"shard prefetch: {} deps -> {} shards for namespace '{namespace}'",
deps.len(),
shard_set.shards.len()
);
let mut handles = Vec::new();
for (hash, _entries) in &shard_set.shards {
let b = Arc::clone(backend);
let p = prefix.to_string();
let ns = namespace.to_string();
let h = hash.clone();
handles.push(tokio::spawn(async move {
crate::remote::download_shard(b.as_ref(), &p, &ns, &h).await
}));
}
let mut prefetch_keys: Vec<(String, String)> = Vec::new();
let mut shards_matched = 0usize;
for handle in handles {
match handle.await {
Ok(Ok(Some(shard))) => {
shards_matched += 1;
for entry in shard.entries {
prefetch_keys.push((entry.cache_key, entry.crate_name));
}
}
Ok(Ok(None)) => {} Ok(Err(e)) => tracing::warn!("shard download error: {e}"),
Err(e) => tracing::warn!("shard download task panicked: {e}"),
}
}
tracing::info!(
"shard prefetch: {shards_matched}/{} shards matched, {} keys to prefetch",
shard_set.shards.len(),
prefetch_keys.len()
);
if prefetch_keys.is_empty() {
return Ok(0);
}
let count = prefetch_keys.len();
let req = PrefetchRequest {
keys: prefetch_keys,
warm_all: false,
};
let resp = daemon.handle_prefetch(&req).await;
if !resp.ok {
anyhow::bail!(
"prefetch failed: {}",
resp.error.as_deref().unwrap_or("unknown")
);
}
Ok(count)
}
async fn monolithic_manifest_prefetch(
daemon: &Arc<Daemon>,
backend: &dyn crate::remote_backend::RemoteBackend,
remote: &crate::config::RemoteConfig,
) {
let manifest_key =
std::env::var("KACHE_MANIFEST_KEY").unwrap_or_else(|_| crate::cli::default_manifest_key());
let min_compile_ms: u64 = std::env::var("KACHE_MIN_COMPILE_MS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(1000);
let manifest = match crate::remote::download_manifest(backend, &remote.prefix, &manifest_key)
.await
{
Ok(m) => m,
Err(e) => {
tracing::info!("manifest prefetch: no manifest for '{manifest_key}' ({e}), skipping");
return;
}
};
let mut worth_prefetching: Vec<_> = manifest
.entries
.iter()
.filter(|e| e.compile_time_ms >= min_compile_ms)
.collect();
worth_prefetching.sort_by_key(|entry| std::cmp::Reverse(entry.compile_time_ms));
let skipped = manifest.entries.len() - worth_prefetching.len();
tracing::info!(
"manifest prefetch: {} entries, prefetching {} (skipped {} cheap crates < {}ms)",
manifest.entries.len(),
worth_prefetching.len(),
skipped,
min_compile_ms
);
if worth_prefetching.is_empty() {
return;
}
let prefetch_keys: Vec<(String, String)> = worth_prefetching
.iter()
.map(|e| (e.cache_key.clone(), e.crate_name.clone()))
.collect();
let req = PrefetchRequest {
keys: prefetch_keys,
warm_all: false,
};
let resp = daemon.handle_prefetch(&req).await;
if !resp.ok {
tracing::warn!(
"manifest prefetch failed: {}",
resp.error.as_deref().unwrap_or("unknown")
);
}
}
const MAX_REQUEST_FRAME_BYTES: usize = 8 * 1024 * 1024;
async fn read_bounded_line<R>(reader: &mut R, buf: &mut Vec<u8>) -> std::io::Result<Option<String>>
where
R: AsyncBufRead + Unpin,
{
buf.clear();
loop {
let available = reader.fill_buf().await?;
if available.is_empty() {
return Ok((!buf.is_empty()).then(|| decode_request_frame(buf)));
}
if let Some(pos) = available.iter().position(|&b| b == b'\n') {
buf.extend_from_slice(&available[..pos]);
std::pin::Pin::new(&mut *reader).consume(pos + 1);
return Ok(Some(decode_request_frame(buf)));
}
buf.extend_from_slice(available);
let consumed = available.len();
std::pin::Pin::new(&mut *reader).consume(consumed);
if buf.len() > MAX_REQUEST_FRAME_BYTES {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"request frame exceeds maximum size",
));
}
}
}
fn decode_request_frame(buf: &[u8]) -> String {
let mut s = String::from_utf8_lossy(buf).into_owned();
if s.ends_with('\r') {
s.pop();
}
s
}
async fn offload<F>(f: F) -> Response
where
F: FnOnce() -> Response + Send + 'static,
{
match tokio::task::spawn_blocking(f).await {
Ok(resp) => resp,
Err(e) => Response::err(format!("daemon handler task failed: {e}")),
}
}
async fn handle_connection(
stream: TokioStream,
daemon: &Arc<Daemon>,
shutdown_flag: &AtomicBool,
shutdown_notify: &Notify,
) -> Result<()> {
let mut reader = BufReader::new(&stream);
let mut frame = Vec::new();
loop {
let line = match read_bounded_line(&mut reader, &mut frame).await {
Ok(Some(l)) => l,
Ok(None) => break,
Err(e) if is_client_disconnect(&e) => {
tracing::debug!("client disconnected mid-read: {e}");
break;
}
Err(e) => return Err(e.into()),
};
let start = Instant::now();
let parsed = serde_json::from_str::<Request>(&line);
let client_epoch = match &parsed {
Ok(Request::Upload(job)) => job.client_epoch,
Ok(Request::Stats(req)) => req.client_epoch,
Ok(Request::BuildStarted(req)) => req.client_epoch,
Ok(Request::LocalLookup(req)) => req.client_epoch,
_ => 0,
};
let resp = match parsed {
Ok(Request::Upload(ref job)) => {
tracing::debug!(
crate_name = job.crate_name,
key = key_prefix(&job.key),
"handling upload request"
);
daemon.handle_upload(job).await
}
Ok(Request::Gc(req)) => {
let d = Arc::clone(daemon);
offload(move || d.handle_gc(&req)).await
}
Ok(Request::RemoteCheck(req)) => daemon.handle_remote_check(&req).await,
Ok(Request::LocalLookup(req)) => daemon.handle_local_lookup(&req).await,
Ok(Request::Stats(req)) => {
let d = Arc::clone(daemon);
offload(move || d.handle_stats(&req)).await
}
Ok(Request::BatchRemoteCheck(req)) => daemon.handle_batch_remote_check(&req).await,
Ok(Request::HashFiles(req)) => {
let d = Arc::clone(daemon);
offload(move || d.handle_hash_files(&req)).await
}
Ok(Request::Prefetch(req)) => daemon.handle_prefetch(&req).await,
Ok(Request::BuildStarted(req)) => daemon.handle_build_started(&req).await,
Ok(Request::CompileStarted(req)) => daemon.handle_compile_started(req),
Ok(Request::CompileFinished(req)) => daemon.handle_compile_finished(&req),
Ok(Request::Shutdown) => {
shutdown_flag.store(true, Ordering::Relaxed);
shutdown_notify.notify_one();
Response::ok()
}
Err(e) => {
tracing::warn!("invalid request from client: {e}");
Response::err(format!("invalid request: {e}"))
}
};
let elapsed = start.elapsed();
if client_epoch_is_newer(client_epoch, daemon.build_epoch)
&& !shutdown_flag.load(Ordering::Relaxed)
{
tracing::info!(
daemon_epoch = daemon.build_epoch,
client_epoch,
"client binary is newer than daemon, scheduling restart"
);
shutdown_flag.store(true, Ordering::Relaxed);
shutdown_notify.notify_one();
}
if !resp.ok {
tracing::warn!(
elapsed_ms = elapsed.as_millis() as u64,
error = resp.error.as_deref().unwrap_or("unknown"),
"request failed"
);
}
let mut resp_line = serde_json::to_string(&resp)?;
resp_line.push('\n');
if let Err(e) = (&stream).write_all(resp_line.as_bytes()).await {
tracing::debug!("response write failed (client likely closed): {e}");
break;
}
}
Ok(())
}
fn is_client_disconnect(e: &std::io::Error) -> bool {
matches!(
e.kind(),
std::io::ErrorKind::BrokenPipe | std::io::ErrorKind::ConnectionReset
) || e.raw_os_error() == Some(32) }
fn key_prefix(key: &str) -> &str {
let mut end = key.len().min(16);
while end > 0 && !key.is_char_boundary(end) {
end -= 1;
}
&key[..end]
}
fn send_retry_delay(attempt: u32, pid: u32) -> Duration {
let jitter = (u64::from(pid) * 7) % 50;
Duration::from_millis(100 * u64::from(attempt) + jitter)
}
fn should_warn_key_cache_refresh_failure(consecutive_refresh_failures: u32) -> bool {
consecutive_refresh_failures == 1 || consecutive_refresh_failures.is_multiple_of(10)
}
fn rotate_daemon_log_if_large(log_path: &Path) {
if std::fs::metadata(log_path).is_ok_and(|m| m.len() > 2 * 1024 * 1024) {
let _ = std::fs::write(log_path, b"--- log rotated ---\n");
}
}
use crate::platform::wait_for_shutdown as shutdown_signal;
pub fn send_upload_job(
config: &Config,
key: &str,
entry_dir: &Path,
crate_name: &str,
) -> Result<()> {
let socket_path = config.socket_path();
let req = Request::Upload(UploadJob {
key: key.to_string(),
entry_dir: entry_dir.to_string_lossy().into_owned(),
crate_name: crate_name.to_string(),
client_epoch: build_epoch(),
});
let key_short = key_prefix(key);
let try_send = |path: &Path| -> Result<()> { send_request_fire_and_forget(path, &req) };
match try_send(&socket_path) {
Ok(()) => return Ok(()),
Err(first_err) => {
tracing::debug!(
crate_name,
key = key_short,
"initial upload send failed, starting daemon: {first_err:#}",
);
match start_daemon_background() {
Ok(true) => {}
Ok(false) | Err(_) => {
tracing::warn!(
crate_name,
key = key_short,
"could not reach or start daemon, skipping upload"
);
return Ok(());
}
}
}
}
for attempt in 1..=3u32 {
match try_send(&socket_path) {
Ok(()) => return Ok(()),
Err(e) => {
if attempt < 3 {
let delay = send_retry_delay(attempt, std::process::id());
tracing::debug!(
crate_name,
key = key_short,
attempt,
"upload send retry {attempt}/3 failed, backoff {delay:?}: {e:#}",
);
std::thread::sleep(delay);
} else {
tracing::warn!(
crate_name,
key = key_short,
socket = %socket_path.display(),
"upload send failed after {attempt} retries: {e:#}",
);
}
}
}
}
Ok(()) }
pub struct GcRequestOutcome {
pub evicted: Option<usize>,
pub skipped: bool,
}
pub fn send_gc_request(config: &Config, max_age_hours: Option<u64>) -> Result<GcRequestOutcome> {
let socket_path = config.socket_path();
let req = Request::Gc(GcRequest { max_age_hours });
let try_send = |path: &Path| -> Result<Response> {
let resp_str = send_request(path, &req)?;
let resp: Response = serde_json::from_str(&resp_str)?;
Ok(resp)
};
match try_send(&socket_path) {
Ok(resp) => {
if resp.ok {
Ok(GcRequestOutcome {
evicted: resp.evicted,
skipped: resp.skipped,
})
} else {
anyhow::bail!("daemon GC error: {}", resp.error.unwrap_or_default());
}
}
Err(_) => {
if start_daemon_background()? {
let resp = try_send(&socket_path)?;
if resp.ok {
Ok(GcRequestOutcome {
evicted: resp.evicted,
skipped: resp.skipped,
})
} else {
anyhow::bail!("daemon GC error: {}", resp.error.unwrap_or_default());
}
} else {
anyhow::bail!("could not reach or start daemon");
}
}
}
}
pub struct RemoteCheckResult {
pub found: bool,
pub prefetched: bool,
}
fn remote_check_result_from_response_line(resp_str: &str) -> Option<RemoteCheckResult> {
match serde_json::from_str::<Response>(resp_str) {
Ok(resp) if resp.ok => resp.found.map(|found| RemoteCheckResult {
found,
prefetched: resp.prefetched.unwrap_or(false),
}),
Ok(resp) => {
tracing::warn!(
"remote check error: {}",
resp.error.as_deref().unwrap_or("unknown")
);
None
}
Err(e) => {
tracing::warn!("remote check response parse error: {e}");
None
}
}
}
pub fn send_remote_check(
config: &Config,
key: &str,
entry_dir: &Path,
crate_name: &str,
) -> Option<RemoteCheckResult> {
let socket_path = config.socket_path();
if !crate::transport::is_reachable(&socket_path) {
return None;
}
let req = Request::RemoteCheck(RemoteCheckRequest {
key: key.to_string(),
entry_dir: entry_dir.to_string_lossy().into_owned(),
crate_name: crate_name.to_string(),
});
match send_request_with_timeout(&socket_path, &req, std::time::Duration::from_secs(3)) {
Ok(resp_str) => remote_check_result_from_response_line(&resp_str),
Err(e) => {
tracing::debug!("remote check: daemon unreachable ({e})");
None
}
}
}
pub fn send_local_lookup(config: &Config, key: &str) -> Option<LocalLookupReply> {
let socket_path = config.socket_path();
if !crate::transport::is_reachable(&socket_path) {
return None;
}
let req = Request::LocalLookup(LocalLookupRequest {
key: key.to_string(),
client_epoch: build_epoch(),
});
let timeout = std::env::var("KACHE_LOCAL_HIT_TIMEOUT_MS")
.ok()
.and_then(|v| v.parse().ok())
.map(std::time::Duration::from_millis)
.unwrap_or(std::time::Duration::from_millis(250));
match send_request_with_timeout(&socket_path, &req, timeout) {
Ok(resp_str) => match serde_json::from_str::<Response>(&resp_str) {
Ok(resp) if resp.ok => resp.local_lookup,
_ => None,
},
Err(e) => {
tracing::debug!("local lookup: daemon unreachable ({e})");
None
}
}
}
pub fn send_hash_files_request(
socket_path: &Path,
files: Vec<HashFileRequest>,
) -> Result<Vec<HashFileResult>> {
if files.is_empty() {
return Ok(Vec::new());
}
if !socket_path.exists() {
anyhow::bail!("daemon socket does not exist: {}", socket_path.display());
}
let req = Request::HashFiles(HashFilesRequest { files });
let resp_str = send_request_with_timeout(socket_path, &req, std::time::Duration::from_secs(3))?;
hash_files_results_from_response_line(&resp_str)
}
fn hash_files_results_from_response_line(resp_str: &str) -> Result<Vec<HashFileResult>> {
let resp: Response = serde_json::from_str(resp_str)?;
if !resp.ok {
anyhow::bail!(
"daemon hash_files error: {}",
resp.error.unwrap_or_default()
);
}
Ok(resp.hash_results.unwrap_or_default())
}
#[allow(dead_code)]
pub fn send_prefetch(config: &Config, keys: &[(String, String)]) -> Result<()> {
let socket_path = config.socket_path();
let req = Request::Prefetch(PrefetchRequest {
keys: keys.to_vec(),
warm_all: false,
});
let try_send = |path: &Path| -> Result<()> { send_request_fire_and_forget(path, &req) };
match try_send(&socket_path) {
Ok(()) => return Ok(()),
Err(_) => match start_daemon_background() {
Ok(true) => {}
Ok(false) | Err(_) => {
tracing::warn!("could not reach or start daemon, skipping prefetch");
return Ok(());
}
},
}
for attempt in 1..=3u32 {
match try_send(&socket_path) {
Ok(()) => return Ok(()),
Err(e) => {
if attempt < 3 {
std::thread::sleep(send_retry_delay(attempt, std::process::id()));
} else {
tracing::warn!("prefetch send failed after {attempt} retries: {e}");
}
}
}
}
Ok(()) }
pub fn send_build_started(config: &Config, req: BuildStartedRequest) {
let socket_path = config.socket_path();
let crate_count = req.intent.crate_names.len();
let req = Request::BuildStarted(req);
match send_request_fire_and_forget(&socket_path, &req) {
Ok(()) => {
tracing::debug!("build-started hint sent for {} crates", crate_count);
}
Err(e) => {
tracing::debug!("build-started hint: daemon unreachable ({e}), skipping");
}
}
}
const IN_FLIGHT_MAX_AGE_MS: u64 = 6 * 60 * 60 * 1000;
fn unix_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}
fn prune_in_flight(map: &mut HashMap<u32, CompileStartedRequest>) {
let now = unix_ms();
map.retain(|&pid, c| {
now.saturating_sub(c.started_at_ms) <= IN_FLIGHT_MAX_AGE_MS && pid_alive(pid)
});
}
#[cfg(unix)]
fn pid_alive(pid: u32) -> bool {
let rc = unsafe { libc::kill(pid as libc::pid_t, 0) };
rc == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM)
}
#[cfg(not(unix))]
fn pid_alive(_pid: u32) -> bool {
true
}
pub fn send_compile_started(socket_path: &std::path::Path, req: CompileStartedRequest) {
if !crate::transport::is_reachable(socket_path) {
return;
}
let req = Request::CompileStarted(req);
if let Err(e) = send_request_fire_and_forget(socket_path, &req) {
tracing::debug!("compile-started: daemon unreachable ({e}), skipping");
}
}
pub fn send_compile_finished(socket_path: &std::path::Path, pid: u32, started_at_ms: u64) {
if !crate::transport::is_reachable(socket_path) {
return;
}
let req = Request::CompileFinished(CompileFinishedRequest { pid, started_at_ms });
if let Err(e) = send_request_fire_and_forget(socket_path, &req) {
tracing::debug!("compile-finished: daemon unreachable ({e}), skipping");
}
}
pub fn send_stats_request(
config: &Config,
include_entries: bool,
sort_by: Option<&str>,
event_hours: Option<u64>,
) -> Result<StatsResponse> {
let socket_path = config.socket_path();
let client_epoch = build_epoch();
let req = Request::Stats(StatsRequest {
include_entries,
sort_by: sort_by.map(String::from),
event_hours,
client_epoch,
});
let resp_str =
send_request_with_timeout(&socket_path, &req, std::time::Duration::from_secs(5))?;
let resp: Response = serde_json::from_str(&resp_str)?;
let stats = if resp.ok {
resp.stats
.ok_or_else(|| anyhow::anyhow!("stats response missing payload"))?
} else {
anyhow::bail!("daemon stats error: {}", resp.error.unwrap_or_default())
};
if client_epoch_is_newer(client_epoch, stats.build_epoch) {
tracing::info!(
daemon_epoch = stats.build_epoch,
client_epoch,
"stale daemon detected via stats request, restarting"
);
if restart_daemon_for_stale_client(config)?
&& let Ok(fresh_resp_str) =
send_request_with_timeout(&socket_path, &req, std::time::Duration::from_secs(3))
&& let Ok(fresh_resp) = serde_json::from_str::<Response>(&fresh_resp_str)
&& fresh_resp.ok
&& let Some(fresh_stats) = fresh_resp.stats
{
return Ok(fresh_stats);
}
}
Ok(stats)
}
pub fn send_shutdown_request(config: &Config) -> Result<()> {
let socket_path = config.socket_path();
match send_request_with_timeout(&socket_path, &Request::Shutdown, Duration::from_secs(5)) {
Ok(_) => {
eprintln!("daemon stopped");
Ok(())
}
Err(e) => {
if let Some(state) = read_daemon_state(&socket_path)
&& process_is_alive(state.pid)
{
tracing::info!(
pid = state.pid,
"socket unreachable, terminating daemon process"
);
crate::platform::terminate_process(state.pid);
if wait_for_run_lock_release(&socket_path, Duration::from_secs(3))? {
let _ = std::fs::remove_file(&socket_path);
eprintln!("daemon stopped (terminated stale process)");
return Ok(());
}
tracing::warn!(pid = state.pid, "daemon did not stop, force-killing");
crate::platform::kill_process(state.pid);
if wait_for_run_lock_release(&socket_path, Duration::from_secs(2))? {
let _ = std::fs::remove_file(&socket_path);
eprintln!("daemon stopped (killed stale process)");
return Ok(());
}
}
Err(e).context("connecting to daemon socket")
}
}
}
pub fn find_daemon_pids() -> Vec<u32> {
let own_pid = std::process::id();
#[cfg(unix)]
{
let output = match std::process::Command::new("pgrep")
.args(["-f", "kache daemon run"])
.output()
{
Ok(o) if o.status.success() => o,
_ => return Vec::new(),
};
String::from_utf8_lossy(&output.stdout)
.lines()
.filter_map(|l| l.trim().parse::<u32>().ok())
.filter(|&pid| pid != own_pid && process_is_alive(pid))
.collect()
}
#[cfg(windows)]
{
let output = match std::process::Command::new("tasklist")
.args(["/FI", "IMAGENAME eq kache.exe", "/FO", "CSV", "/NH"])
.output()
{
Ok(o) if o.status.success() => o,
_ => return Vec::new(),
};
String::from_utf8_lossy(&output.stdout)
.lines()
.filter_map(|line| {
let fields: Vec<&str> = line.split(',').collect();
fields.get(1)?.trim_matches('"').parse::<u32>().ok()
})
.filter(|&pid| pid != own_pid && process_is_alive(pid))
.collect()
}
}
pub fn force_recover(config: &Config) -> Result<()> {
let socket_path = config.socket_path();
let pids = find_daemon_pids();
if !pids.is_empty() {
tracing::info!(?pids, "killing lingering kache daemon processes");
for &pid in &pids {
crate::platform::terminate_process(pid);
}
std::thread::sleep(Duration::from_millis(500));
for &pid in &pids {
if process_is_alive(pid) {
tracing::warn!(pid, "graceful terminate did not land, force-killing");
crate::platform::kill_process(pid);
}
}
std::thread::sleep(Duration::from_millis(200));
}
let _ = std::fs::remove_file(&socket_path);
let _ = std::fs::remove_file(daemon_state_path(&socket_path));
let _ = std::fs::remove_file(socket_path.with_extension("lock"));
let _ = std::fs::remove_file(socket_path.with_extension("run.lock"));
Ok(())
}
pub fn restart(config: &Config) -> Result<bool> {
let socket_path = config.socket_path();
match crate::service::kickstart() {
Ok(true) => {
eprintln!("restarting daemon via service manager...");
if wait_for_socket_until(&socket_path, None, Duration::from_secs(10))? {
let responsive = send_stats_request(config, false, None, None).is_ok();
let pids = find_daemon_pids();
if responsive && pids.len() <= 1 {
eprintln!("daemon restarted");
return Ok(true);
}
tracing::warn!(
responsive,
daemon_pids = ?pids,
"service kickstart reported success but daemon isn't healthy; attempting nuclear recovery"
);
} else {
tracing::warn!(
"service kickstart completed but socket not ready; attempting nuclear recovery"
);
}
}
Ok(false) => {
}
Err(e) => {
tracing::warn!("service kickstart failed: {e:#}; attempting nuclear recovery");
}
}
let _ = send_shutdown_request(config);
force_recover(config)?;
match start_daemon_background()? {
true => {
eprintln!("daemon restarted");
Ok(true)
}
false => {
eprintln!("daemon did not start within timeout");
Ok(false)
}
}
}
fn restart_daemon_for_stale_client(config: &Config) -> Result<bool> {
let socket_path = config.socket_path();
let _ = send_request_with_timeout(&socket_path, &Request::Shutdown, Duration::from_secs(2));
for _ in 0..4 {
if !crate::transport::is_reachable(&socket_path) {
break;
}
std::thread::sleep(Duration::from_millis(100));
}
start_daemon_background()
}
fn send_request(socket_path: &Path, req: &Request) -> Result<String> {
send_request_with_timeout(socket_path, req, std::time::Duration::from_secs(30))
}
fn send_request_with_timeout(
socket_path: &Path,
req: &Request,
read_timeout: std::time::Duration,
) -> Result<String> {
#[cfg(windows)]
{
send_request_with_async_timeout(socket_path, req, read_timeout)
}
#[cfg(not(windows))]
{
send_request_with_socket_timeout(socket_path, req, read_timeout)
}
}
#[cfg(not(windows))]
fn send_request_with_socket_timeout(
socket_path: &Path,
req: &Request,
read_timeout: std::time::Duration,
) -> Result<String> {
use crate::transport::SyncStream;
use interprocess::local_socket::traits::Stream as _;
use std::io::{BufRead, Write};
let name = socket_name(socket_path)?;
let mut stream = SyncStream::connect(name)
.with_context(|| format!("connecting to daemon socket {}", socket_path.display()))?;
let _ = stream.set_recv_timeout(Some(read_timeout));
let _ = stream.set_send_timeout(Some(std::time::Duration::from_secs(5)));
let mut line = serde_json::to_string(req)?;
line.push('\n');
stream
.write_all(line.as_bytes())
.context("writing request to daemon")?;
stream.flush().context("flushing request to daemon")?;
let mut reader = std::io::BufReader::new(&stream);
let mut resp = String::new();
reader.read_line(&mut resp).with_context(|| {
format!(
"reading response from daemon (timeout {:?}, socket {})",
read_timeout,
socket_path.display()
)
})?;
Ok(resp)
}
#[cfg(windows)]
fn send_request_with_async_timeout(
socket_path: &Path,
req: &Request,
read_timeout: std::time::Duration,
) -> Result<String> {
let mut line = serde_json::to_string(req)?;
line.push('\n');
if tokio::runtime::Handle::try_current().is_ok() {
let socket_path = socket_path.to_path_buf();
std::thread::spawn(move || {
send_request_with_async_timeout_blocking(&socket_path, line, read_timeout)
})
.join()
.map_err(|_| anyhow::anyhow!("daemon client timeout thread panicked"))?
} else {
send_request_with_async_timeout_blocking(socket_path, line, read_timeout)
}
}
#[cfg(windows)]
fn send_request_with_async_timeout_blocking(
socket_path: &Path,
line: String,
read_timeout: std::time::Duration,
) -> Result<String> {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_io()
.enable_time()
.build()
.context("creating daemon client runtime")?;
runtime.block_on(async {
tokio::time::timeout(
read_timeout,
send_request_with_async_transport(socket_path, line, read_timeout),
)
.await
.with_context(|| {
format!(
"daemon request timed out after {:?} (socket {})",
read_timeout,
socket_path.display()
)
})?
})
}
#[cfg(windows)]
async fn send_request_with_async_transport(
socket_path: &Path,
line: String,
read_timeout: std::time::Duration,
) -> Result<String> {
let name = socket_name(socket_path)?;
let mut stream = TokioStream::connect(name)
.await
.with_context(|| format!("connecting to daemon socket {}", socket_path.display()))?;
stream
.write_all(line.as_bytes())
.await
.context("writing request to daemon")?;
stream.flush().await.context("flushing request to daemon")?;
let mut reader = BufReader::new(stream);
let mut resp = String::new();
reader.read_line(&mut resp).await.with_context(|| {
format!(
"reading response from daemon (timeout {:?}, socket {})",
read_timeout,
socket_path.display()
)
})?;
Ok(resp)
}
fn send_request_fire_and_forget(socket_path: &Path, req: &Request) -> Result<()> {
use crate::transport::SyncStream;
use interprocess::local_socket::traits::Stream as _;
use std::io::Write;
let name = socket_name(socket_path)?;
let mut stream = SyncStream::connect(name)
.with_context(|| format!("connecting to daemon socket {}", socket_path.display()))?;
let _ = stream.set_send_timeout(Some(std::time::Duration::from_secs(5)));
let mut line = serde_json::to_string(req)?;
line.push('\n');
stream
.write_all(line.as_bytes())
.context("writing request to daemon")?;
stream.flush().context("flushing request to daemon")?;
Ok(())
}
pub fn start_daemon_background() -> Result<bool> {
let config = Config::load()?;
let socket_path = config.socket_path();
let lock_path = socket_path.with_extension("lock");
let mut recovered_once = false;
for attempt in 0..2 {
std::fs::create_dir_all(socket_path.parent().unwrap())?;
let lock_file = std::fs::OpenOptions::new()
.create(true)
.write(true)
.truncate(false)
.open(&lock_path)
.context("opening daemon lock file")?;
let got_lock = lock_file.try_lock().is_ok();
if !got_lock {
tracing::debug!("daemon start already in progress, waiting for socket");
if wait_for_socket(&socket_path, None)? {
if recovered_once {
tracing::info!(
socket = %socket_path.display(),
"daemon startup recovered after retry"
);
}
return Ok(true);
}
if attempt == 0 {
tracing::info!(
socket = %socket_path.display(),
"daemon starter timed out without publishing a ready socket, retrying coordination"
);
std::thread::sleep(DAEMON_START_POLL_INTERVAL);
continue;
}
return Ok(false);
}
if crate::transport::is_reachable(&socket_path) {
let my_epoch = build_epoch();
let is_stale = send_request_with_timeout(
&socket_path,
&Request::Stats(StatsRequest {
include_entries: false,
sort_by: None,
event_hours: None,
client_epoch: my_epoch,
}),
Duration::from_secs(2),
)
.ok()
.and_then(|s| serde_json::from_str::<Response>(&s).ok())
.and_then(|r| r.stats)
.map(|s| client_epoch_is_newer(my_epoch, s.build_epoch))
.unwrap_or(false);
if !is_stale {
tracing::debug!("daemon already running");
return Ok(true);
}
tracing::info!("stale daemon detected, requesting shutdown before restart");
let _ =
send_request_with_timeout(&socket_path, &Request::Shutdown, Duration::from_secs(2));
if !wait_for_run_lock_release(&socket_path, Duration::from_secs(5))? {
tracing::info!(
socket = %socket_path.display(),
"stale daemon did not exit within timeout, attempting bounded recovery"
);
if attempt == 0
&& recover_unhealthy_daemon(
&socket_path,
"stale daemon did not exit after shutdown request",
)?
{
recovered_once = true;
continue;
}
return Ok(false);
}
}
if daemon_run_lock_is_held(&socket_path)? {
tracing::debug!(
socket = %socket_path.display(),
"daemon run lock already held, waiting for socket"
);
if wait_for_socket(&socket_path, None)? {
return Ok(true);
}
if attempt == 0
&& recover_unhealthy_daemon(
&socket_path,
"daemon run lock held but no ready socket became reachable",
)?
{
recovered_once = true;
continue;
}
return Ok(false);
}
let exe = std::env::current_exe().context("getting current executable path")?;
tracing::info!("auto-starting daemon");
let log_path = socket_path.with_extension("log");
rotate_daemon_log_if_large(&log_path);
let stderr_target = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(&log_path)
.map(std::process::Stdio::from)
.unwrap_or_else(|_| std::process::Stdio::null());
let mut child = std::process::Command::new(exe)
.args(["daemon", "run"])
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(stderr_target)
.spawn()
.context("spawning daemon process")?;
let ready = wait_for_socket(&socket_path, Some(&mut child))?;
if ready {
if recovered_once {
tracing::info!(
socket = %socket_path.display(),
"daemon started successfully after recovery"
);
} else {
tracing::info!("daemon started successfully");
}
return Ok(true);
}
if attempt == 0
&& recover_unhealthy_daemon(
&socket_path,
"daemon starter failed to publish a ready socket before timeout",
)?
{
recovered_once = true;
continue;
}
return Ok(false);
}
Ok(false)
}
fn daemon_run_lock_is_held(socket_path: &Path) -> Result<bool> {
let run_lock_path = socket_path.with_extension("run.lock");
let run_lock_file = std::fs::OpenOptions::new()
.create(true)
.write(true)
.truncate(false)
.open(&run_lock_path)
.context("opening daemon run lock probe file")?;
if run_lock_file.try_lock().is_ok() {
let _ = run_lock_file.unlock();
Ok(false)
} else {
Ok(true)
}
}
fn wait_for_socket(socket_path: &Path, child: Option<&mut std::process::Child>) -> Result<bool> {
wait_for_socket_until(socket_path, child, DAEMON_START_TIMEOUT)
}
fn wait_for_socket_until(
socket_path: &Path,
mut child: Option<&mut std::process::Child>,
timeout: Duration,
) -> Result<bool> {
let deadline = Instant::now() + timeout;
while Instant::now() < deadline {
if crate::transport::is_reachable(socket_path) {
return Ok(true);
}
if let Some(child_proc) = child.as_mut()
&& let Some(status) = child_proc
.try_wait()
.context("checking daemon process status")?
{
if status.success() {
tracing::debug!(
socket = %socket_path.display(),
?status,
"daemon starter exited cleanly before socket became ready, continuing to wait"
);
child = None;
continue;
}
tracing::warn!(
socket = %socket_path.display(),
?status,
"daemon exited before socket became ready"
);
return Ok(false);
}
std::thread::sleep(DAEMON_START_POLL_INTERVAL);
}
if crate::transport::is_reachable(socket_path) {
return Ok(true);
}
if let Some(child) = child.as_mut()
&& child
.try_wait()
.context("checking daemon process status after timeout")?
.is_none()
{
tracing::debug!(
socket = %socket_path.display(),
timeout_ms = timeout.as_millis(),
"daemon did not start within timeout, terminating starter process"
);
let _ = child.kill();
let _ = child.wait();
}
tracing::warn!(
socket = %socket_path.display(),
timeout_ms = timeout.as_millis(),
"daemon did not start within timeout"
);
Ok(false)
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::mpsc;
#[test]
fn should_cancel_prefetch_fires_on_low_candidate_share() {
assert!(should_cancel_prefetch(12, 1, 0));
}
#[test]
fn should_cancel_prefetch_holds_below_min_demands() {
assert!(!should_cancel_prefetch(9, 0, 0));
}
#[test]
fn should_cancel_prefetch_holds_when_plan_is_good() {
assert!(!should_cancel_prefetch(20, 15, 0));
}
#[test]
fn should_cancel_prefetch_counts_undmanded_downloads_as_potential_hits() {
assert!(should_cancel_prefetch(20, 2, 0));
assert!(!should_cancel_prefetch(20, 2, 8));
}
#[test]
fn active_plan_tracks_demand_download_and_use() {
let mut plan = ActivePlan::new(
"sess-1".into(),
"plan-1".into(),
"fallback",
["a", "b"].into_iter().map(String::from).collect(),
0,
0,
);
assert!(!plan.record_demand("a"));
assert_eq!(plan.demanded.len(), 1);
assert_eq!(plan.demanded_candidates.len(), 1);
assert!(plan.used.is_empty());
plan.record_download("a", 100);
assert!(plan.used.contains("a"));
plan.record_download("b", 50);
assert!(!plan.record_demand("b"));
assert!(plan.used.contains("b"));
assert_eq!(plan.used_bytes(), 150);
assert!(!plan.record_demand("a"));
assert_eq!(plan.demanded.len(), 2);
}
#[test]
fn active_plan_cancel_latch_fires_once() {
let mut plan = ActivePlan::new(
"sess-2".into(),
String::new(),
"advisory",
["only-candidate".to_string()].into_iter().collect(),
0,
0,
);
for i in 0..9 {
assert!(!plan.record_demand(&format!("k{i}")));
}
assert!(plan.record_demand("k9"));
assert!(plan.cancelled);
assert!(!plan.record_demand("k10"));
}
use crate::transport::{ListenerOptions, TokioListener, TokioStream, socket_name};
fn bind_listener(path: &Path) -> TokioListener {
let name = socket_name(path).expect("socket name");
ListenerOptions::new()
.name(name)
.create_tokio()
.expect("create_tokio listener")
}
async fn connect_stream(path: &Path) -> TokioStream {
let name = socket_name(path).expect("socket name");
TokioStream::connect(name).await.expect("connect")
}
fn bind_sync_listener(path: &Path) -> interprocess::local_socket::Listener {
let name = socket_name(path).expect("socket name");
ListenerOptions::new()
.name(name)
.create_sync()
.expect("create_sync listener")
}
fn spawn_quick_exit_child() -> std::process::Child {
#[cfg(unix)]
{
std::process::Command::new("sh")
.args(["-c", "exit 0"])
.spawn()
.unwrap()
}
#[cfg(windows)]
{
std::process::Command::new("cmd")
.args(["/c", "exit", "0"])
.spawn()
.unwrap()
}
}
fn spawn_blocking_child() -> std::process::Child {
#[cfg(unix)]
{
std::process::Command::new("sh")
.args(["-c", "sleep 30"])
.spawn()
.unwrap()
}
#[cfg(windows)]
{
std::process::Command::new("cmd")
.args(["/c", "ping", "-n", "31", "127.0.0.1"])
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.unwrap()
}
}
async fn client_roundtrip(socket_path: &Path, req: &Request) -> Response {
let mut stream = connect_stream(socket_path).await;
let mut line = serde_json::to_string(req).expect("serialize request");
line.push('\n');
stream
.write_all(line.as_bytes())
.await
.expect("write request");
let mut resp_line = String::new();
{
let mut reader = BufReader::new(&stream);
reader
.read_line(&mut resp_line)
.await
.expect("read response");
}
drop(stream);
serde_json::from_str(&resp_line).expect("parse response")
}
async fn one_shot_request(daemon: &Arc<Daemon>, socket_path: &Path, req: &Request) -> Response {
let listener = bind_listener(socket_path);
let server_daemon = daemon.clone();
let server = tokio::spawn(async move {
let stream = listener.accept().await.expect("accept");
handle_connection(
stream,
&server_daemon,
&AtomicBool::new(false),
&Notify::new(),
)
.await
.expect("handle_connection");
});
let resp = client_roundtrip(socket_path, req).await;
server.await.expect("join server task");
resp
}
#[test]
fn in_flight_registry_upserts_prunes_and_snapshots() {
let dir = tempfile::tempdir().unwrap();
let daemon = Daemon::new(test_config(dir.path()));
let now = unix_ms();
let pid = std::process::id();
daemon.handle_compile_started(CompileStartedRequest {
crate_name: "gkrust".into(),
root: "/w".into(),
pid,
started_at_ms: now.saturating_sub(10_000),
typical_ms: None,
client_epoch: 0,
});
daemon.handle_compile_started(CompileStartedRequest {
crate_name: "gkrust".into(),
root: "/w".into(),
pid,
started_at_ms: now.saturating_sub(10_000),
typical_ms: Some(471_000),
client_epoch: 0,
});
daemon.handle_compile_started(CompileStartedRequest {
crate_name: "ghost".into(),
root: "/w".into(),
pid: pid.wrapping_add(1),
started_at_ms: now.saturating_sub(IN_FLIGHT_MAX_AGE_MS + 60_000),
typical_ms: None,
client_epoch: 0,
});
let snapshot = daemon.in_flight_snapshot();
assert_eq!(
snapshot.len(),
1,
"ghost pruned, upsert deduped: {snapshot:?}"
);
let entry = &snapshot[0];
assert_eq!(entry.crate_name, "gkrust");
assert_eq!(entry.pid, pid);
assert!(entry.elapsed_s >= 10);
assert_eq!(entry.typical_s, Some(471));
assert_eq!(entry.eta_s, Some(471u64.saturating_sub(entry.elapsed_s)));
daemon.handle_compile_finished(&CompileFinishedRequest {
pid,
started_at_ms: 12345,
});
assert_eq!(daemon.in_flight_snapshot().len(), 1);
daemon.handle_compile_finished(&CompileFinishedRequest {
pid,
started_at_ms: now.saturating_sub(10_000),
});
assert!(daemon.in_flight_snapshot().is_empty());
}
#[test]
fn compile_started_wire_tags_and_stats_default() {
let req = Request::CompileStarted(CompileStartedRequest {
crate_name: "c".into(),
root: String::new(),
pid: 1,
started_at_ms: 2,
typical_ms: None,
client_epoch: 0,
});
let wire = serde_json::to_string(&req).unwrap();
assert!(wire.contains("\"compile_started\""), "{wire}");
let round: Request = serde_json::from_str(&wire).unwrap();
assert_eq!(round, req);
let mut old = serde_json::to_value(StatsResponse {
total_size: 0,
max_size: 0,
entry_count: 0,
entries: None,
events: EventStatsResponse {
local_hits: 0,
prefetch_hits: 0,
remote_hits: 0,
dups: 0,
misses: 0,
errors: 0,
total_elapsed_ms: 0,
hit_elapsed_ms: 0,
miss_elapsed_ms: 0,
hit_compile_time_ms: 0,
miss_compile_time_ms: 0,
store_output_blobs: 0,
store_duplicate_blobs: 0,
store_new_blobs: 0,
},
version: String::new(),
build_epoch: 0,
pending_uploads: 0,
active_downloads: 0,
s3_concurrency_total: 0,
s3_concurrency_used: 0,
upload_queue_capacity: 0,
uploads_completed: 0,
uploads_failed: 0,
uploads_skipped: 0,
downloads_completed: 0,
downloads_failed: 0,
bytes_uploaded: 0,
bytes_downloaded: 0,
recent_transfers: Vec::new(),
prefetch: PrefetchStatsSnapshot::default(),
in_flight: vec![InFlightEntry {
crate_name: "x".into(),
root: String::new(),
pid: 1,
elapsed_s: 1,
typical_s: None,
eta_s: None,
}],
})
.unwrap();
old.as_object_mut().unwrap().remove("in_flight");
let parsed: StatsResponse = serde_json::from_value(old).unwrap();
assert!(parsed.in_flight.is_empty());
}
#[tokio::test]
async fn test_shutdown_request_sets_flag_and_stores_notify_permit() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let listener = bind_listener(&socket_path);
let daemon = Arc::new(Daemon::new(config));
let shutdown_flag = Arc::new(AtomicBool::new(false));
let shutdown_notify = Arc::new(Notify::new());
let server_daemon = daemon.clone();
let server_flag = shutdown_flag.clone();
let server_notify = shutdown_notify.clone();
let server = tokio::spawn(async move {
let stream = listener.accept().await.expect("accept");
handle_connection(stream, &server_daemon, &server_flag, &server_notify)
.await
.expect("handle_connection");
});
let resp = client_roundtrip(&socket_path, &Request::Shutdown).await;
server.await.expect("join server task");
assert!(resp.ok, "stop request should return ok");
assert!(
shutdown_flag.load(Ordering::Relaxed),
"stop request must set the shutdown flag"
);
tokio::time::timeout(Duration::from_secs(1), shutdown_notify.notified())
.await
.expect("stop request must leave a notify permit (issue #288)");
}
#[tokio::test]
async fn test_accept_loop_breaks_promptly_on_stop_request() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let listener = bind_listener(&socket_path);
let daemon = Arc::new(Daemon::new(config));
let shutdown_flag = Arc::new(AtomicBool::new(false));
let shutdown_notify = Arc::new(Notify::new());
let client_socket = socket_path.clone();
let client =
tokio::spawn(async move { client_roundtrip(&client_socket, &Request::Shutdown).await });
let outcome = tokio::time::timeout(
Duration::from_secs(5),
accept_loop(
&listener,
&daemon,
&shutdown_flag,
&shutdown_notify,
None,
std::future::pending::<()>(),
),
)
.await;
assert!(
outcome.is_ok(),
"accept_loop did not break within 5s of a stop request (issue #288 regression)"
);
assert!(
shutdown_flag.load(Ordering::Relaxed),
"shutdown flag should be set after the stop request"
);
let resp = client.await.expect("join client task");
assert!(resp.ok, "stop request should return ok");
}
#[tokio::test]
async fn test_send_request_with_timeout_bounds_unresponsive_daemon() {
let dir = tempfile::tempdir().unwrap();
let socket_path = dir.path().join("daemon.sock");
let listener = bind_listener(&socket_path);
let server = tokio::spawn(async move {
let stream = listener.accept().await.expect("accept");
let mut request_line = String::new();
{
let mut reader = BufReader::new(&stream);
reader
.read_line(&mut request_line)
.await
.expect("read request");
}
assert!(request_line.contains("\"stats\""));
tokio::time::sleep(Duration::from_secs(1)).await;
drop(stream);
});
let req = Request::Stats(StatsRequest {
include_entries: false,
sort_by: None,
event_hours: None,
client_epoch: 0,
});
let client_socket_path = socket_path.clone();
let started = Instant::now();
let result = tokio::task::spawn_blocking(move || {
send_request_with_timeout(&client_socket_path, &req, Duration::from_millis(75))
})
.await
.expect("join client task");
assert!(result.is_err());
assert!(started.elapsed() < Duration::from_millis(750));
server.abort();
}
fn hold_run_lock_for_test(
socket_path: &Path,
hold_for: Duration,
) -> std::thread::JoinHandle<()> {
let run_lock_path = socket_path.with_extension("run.lock");
let (tx, rx) = mpsc::channel();
let handle = std::thread::spawn(move || {
let file = std::fs::OpenOptions::new()
.create(true)
.write(true)
.truncate(false)
.open(&run_lock_path)
.unwrap();
file.lock().unwrap();
tx.send(()).unwrap();
std::thread::sleep(hold_for);
let _ = file.unlock();
});
rx.recv().unwrap();
handle
}
fn test_config(dir: &Path) -> Config {
Config {
fallback: None,
key_salt: None,
cc_extra_allowlist_flags: Vec::new(),
local_only: false,
remote_readonly: false,
modified_input_guard: false,
local_hit_daemon: false,
windows_hardlink: false,
auto_gc: true,
storage_layout_advice: true,
heartbeat_secs: 30,
explain_miss: false,
path_only_env_vars: Vec::new(),
key_env_vars: Vec::new(),
base_dirs: Vec::new(),
cache_dir: dir.to_path_buf(),
max_size: 50 * 1024 * 1024, remote: None,
remote_error: None,
disabled: false,
cache_executables: false,
clean_incremental: false,
event_log_max_size: 10 * 1024 * 1024,
event_log_keep_lines: 1000,
compression_level: 3,
s3_concurrency: 16,
prefetch_enabled: crate::config::DEFAULT_PREFETCH_ENABLED,
remote_key_cache_refresh_secs: crate::config::DEFAULT_REMOTE_KEY_CACHE_REFRESH_SECS,
prefetch_max_keys: crate::config::DEFAULT_PREFETCH_MAX_KEYS,
prefetch_max_bytes: crate::config::DEFAULT_PREFETCH_MAX_BYTES,
prefetch_deadline_secs: crate::config::DEFAULT_PREFETCH_DEADLINE_SECS,
daemon_idle_timeout_secs: crate::config::DEFAULT_DAEMON_IDLE_TIMEOUT_SECS,
s3_pool_idle_secs: crate::config::DEFAULT_S3_POOL_IDLE_SECS,
}
}
#[test]
fn key_cache_authoritative_truth_table() {
assert!(key_cache_miss_is_authoritative(1, Some(Duration::ZERO)));
assert!(key_cache_miss_is_authoritative(
1,
Some(Duration::from_secs(5))
));
assert!(!key_cache_miss_is_authoritative(
1,
Some(Duration::from_secs(6))
));
assert!(key_cache_miss_is_authoritative(
60,
Some(Duration::from_secs(300))
));
assert!(!key_cache_miss_is_authoritative(
60,
Some(Duration::from_secs(301))
));
assert!(key_cache_miss_is_authoritative(
900,
Some(Duration::from_secs(300))
));
assert!(!key_cache_miss_is_authoritative(
900,
Some(Duration::from_secs(301))
));
assert!(!key_cache_miss_is_authoritative(0, Some(Duration::ZERO)));
assert!(!key_cache_miss_is_authoritative(60, None));
}
#[test]
fn speculative_prefetch_decision_truth_table() {
assert!(speculative_prefetch_disabled(false));
assert!(!speculative_prefetch_disabled(true));
assert!(should_start_speculative_prefetch(true, true));
assert!(!should_start_speculative_prefetch(false, true));
assert!(!should_start_speculative_prefetch(true, false));
assert!(!should_start_speculative_prefetch(false, false));
}
#[test]
fn key_cache_periodic_refresh_disabled_truth_table() {
assert!(key_cache_periodic_refresh_disabled(0));
assert!(!key_cache_periodic_refresh_disabled(1));
assert!(!key_cache_periodic_refresh_disabled(60));
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn server_main_binds_socket_and_handles_shutdown() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.daemon_idle_timeout_secs = 0;
let socket_path = config.socket_path();
let coord = DaemonCoordFile::for_socket(&socket_path);
let server_config = config.clone();
let server = tokio::spawn(async move { server_main(&server_config, coord).await });
let ready_socket = socket_path.clone();
let ready = tokio::task::spawn_blocking(move || {
wait_for_socket_until(&ready_socket, None, Duration::from_secs(5))
})
.await
.unwrap()
.unwrap();
assert!(ready, "server_main must bind its configured socket");
let shutdown_config = config.clone();
tokio::task::spawn_blocking(move || send_shutdown_request(&shutdown_config))
.await
.unwrap()
.unwrap();
let result = tokio::time::timeout(Duration::from_secs(10), server)
.await
.expect("server_main should stop after a shutdown request")
.expect("server_main task should not panic");
assert!(
result.is_ok(),
"server_main should exit cleanly: {result:?}"
);
assert!(
!socket_path.exists(),
"server_main should remove its socket during shutdown"
);
}
#[test]
fn test_request_upload_serde() {
let req = Request::Upload(UploadJob {
key: "abc123".into(),
entry_dir: "/tmp/store/abc123".into(),
crate_name: String::new(),
client_epoch: 0,
});
let json = serde_json::to_string(&req).unwrap();
let parsed: Request = serde_json::from_str(&json).unwrap();
assert_eq!(req, parsed);
assert!(json.contains("\"upload\""));
assert!(json.contains("\"key\":\"abc123\""));
}
#[test]
fn test_wait_for_socket_until_observes_late_socket() {
let dir = tempfile::tempdir().unwrap();
let socket_path = dir.path().join("daemon.sock");
let socket_path_bg = socket_path.clone();
let handle = std::thread::spawn(move || {
std::thread::sleep(Duration::from_millis(150));
let listener = bind_sync_listener(&socket_path_bg);
std::thread::sleep(Duration::from_millis(200));
drop(listener);
});
let ready = wait_for_socket_until(&socket_path, None, Duration::from_secs(1)).unwrap();
handle.join().unwrap();
assert!(ready);
}
#[test]
fn test_wait_for_socket_until_times_out_cleanly() {
let dir = tempfile::tempdir().unwrap();
let socket_path = dir.path().join("missing.sock");
let ready = wait_for_socket_until(&socket_path, None, Duration::from_millis(150)).unwrap();
assert!(!ready);
}
#[test]
fn test_wait_for_socket_until_ignores_clean_child_exit_if_socket_appears() {
let dir = tempfile::tempdir().unwrap();
let socket_path = dir.path().join("daemon.sock");
let socket_path_bg = socket_path.clone();
let handle = std::thread::spawn(move || {
std::thread::sleep(Duration::from_millis(150));
let listener = bind_sync_listener(&socket_path_bg);
std::thread::sleep(Duration::from_millis(200));
drop(listener);
});
let mut child = spawn_quick_exit_child();
let ready =
wait_for_socket_until(&socket_path, Some(&mut child), Duration::from_secs(1)).unwrap();
handle.join().unwrap();
assert!(ready);
}
#[test]
fn test_wait_for_socket_until_kills_stuck_child_after_timeout() {
let dir = tempfile::tempdir().unwrap();
let socket_path = dir.path().join("missing.sock");
let mut child = spawn_blocking_child();
let ready =
wait_for_socket_until(&socket_path, Some(&mut child), Duration::from_millis(150))
.unwrap();
assert!(!ready);
let status = child.try_wait().unwrap();
assert!(status.is_some());
}
#[test]
fn decode_request_frame_strips_trailing_carriage_return() {
assert_eq!(decode_request_frame(b"{\"x\":1}\r"), "{\"x\":1}");
assert_eq!(decode_request_frame(b"{\"x\":1}"), "{\"x\":1}");
assert_eq!(decode_request_frame(b""), "");
}
#[test]
fn is_client_disconnect_matches_disconnect_kinds() {
use std::io::{Error, ErrorKind};
assert!(is_client_disconnect(&Error::from(ErrorKind::BrokenPipe)));
assert!(is_client_disconnect(&Error::from(
ErrorKind::ConnectionReset
)));
assert!(is_client_disconnect(&Error::from_raw_os_error(32))); assert!(!is_client_disconnect(&Error::from(ErrorKind::NotFound)));
assert!(!is_client_disconnect(&Error::from(ErrorKind::TimedOut)));
}
#[test]
fn key_prefix_is_multibyte_safe() {
let hex = "0123456789abcdef".repeat(4);
assert_eq!(key_prefix(&hex), "0123456789abcdef");
assert_eq!(key_prefix("short"), "short");
assert_eq!(key_prefix(""), "");
let s = "アアアアアアアア"; let p = key_prefix(s);
assert!(s.starts_with(p));
assert!(p.len() <= 16);
}
#[test]
fn client_epoch_comparison_ignores_zero_and_detects_newer() {
assert!(!client_epoch_is_newer(0, 10));
assert!(!client_epoch_is_newer(10, 0));
assert!(!client_epoch_is_newer(10, 10));
assert!(!client_epoch_is_newer(9, 10));
assert!(client_epoch_is_newer(11, 10));
}
#[test]
fn send_retry_delay_uses_linear_backoff_and_pid_jitter() {
assert_eq!(send_retry_delay(1, 7), Duration::from_millis(100 + 49));
assert_eq!(send_retry_delay(3, 8), Duration::from_millis(300 + 6));
}
#[test]
fn key_cache_refresh_warning_cadence_is_first_and_every_tenth() {
assert!(should_warn_key_cache_refresh_failure(1));
assert!(!should_warn_key_cache_refresh_failure(2));
assert!(!should_warn_key_cache_refresh_failure(9));
assert!(should_warn_key_cache_refresh_failure(10));
assert!(should_warn_key_cache_refresh_failure(20));
}
#[test]
fn rotate_daemon_log_if_large_truncates_only_oversized_logs() {
let dir = tempfile::tempdir().unwrap();
let small = dir.path().join("small.log");
std::fs::write(&small, b"small log").unwrap();
rotate_daemon_log_if_large(&small);
assert_eq!(std::fs::read(&small).unwrap(), b"small log");
let large = dir.path().join("large.log");
std::fs::write(&large, vec![b'x'; 2 * 1024 * 1024 + 1]).unwrap();
rotate_daemon_log_if_large(&large);
assert_eq!(std::fs::read(&large).unwrap(), b"--- log rotated ---\n");
}
#[test]
fn daemon_state_path_uses_state_json_extension() {
assert_eq!(
daemon_state_path(Path::new("/tmp/kache/daemon.sock")),
Path::new("/tmp/kache/daemon.state.json")
);
}
#[test]
fn daemon_state_is_recent_distinguishes_fresh_from_stale() {
let fresh = DaemonCoordState {
pid: 1,
build_epoch: build_epoch(),
phase: DaemonPhase::Ready,
updated_at_ms: now_millis(),
};
assert!(daemon_state_is_recent(&fresh));
let stale = DaemonCoordState {
pid: 1,
build_epoch: build_epoch(),
phase: DaemonPhase::Ready,
updated_at_ms: now_millis()
.saturating_sub(DAEMON_COORD_STALE_AFTER.as_millis() as u64 * 2),
};
assert!(!daemon_state_is_recent(&stale));
}
#[test]
fn test_daemon_coord_state_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let socket_path = dir.path().join("daemon.sock");
let coord = DaemonCoordFile::for_socket(&socket_path);
coord.write_phase(DaemonPhase::Starting).unwrap();
let state = read_daemon_state(&socket_path).unwrap();
assert_eq!(state.pid, std::process::id());
assert_eq!(state.build_epoch, build_epoch());
assert_eq!(state.phase, DaemonPhase::Starting);
assert!(daemon_state_is_recent(&state));
}
#[test]
fn test_recover_unhealthy_daemon_cleans_stale_socket_and_state() {
let dir = tempfile::tempdir().unwrap();
let socket_path = dir.path().join("daemon.sock");
std::fs::write(&socket_path, b"stale").unwrap();
let state = DaemonCoordState {
pid: u32::MAX,
build_epoch: build_epoch(),
phase: DaemonPhase::Starting,
updated_at_ms: now_millis(),
};
write_json_atomically(&daemon_state_path(&socket_path), &state).unwrap();
assert!(recover_unhealthy_daemon(&socket_path, "test").unwrap());
assert!(!socket_path.exists());
assert!(read_daemon_state(&socket_path).is_none());
}
#[test]
fn test_recover_unhealthy_daemon_terminates_recent_recorded_pid() {
let dir = tempfile::tempdir().unwrap();
let socket_path = dir.path().join("daemon.sock");
std::fs::write(&socket_path, b"stale").unwrap();
let run_lock_handle = hold_run_lock_for_test(&socket_path, Duration::from_millis(150));
let mut child = spawn_blocking_child();
let state = DaemonCoordState {
pid: child.id(),
build_epoch: build_epoch(),
phase: DaemonPhase::Ready,
updated_at_ms: now_millis(),
};
write_json_atomically(&daemon_state_path(&socket_path), &state).unwrap();
assert!(recover_unhealthy_daemon(&socket_path, "test").unwrap());
run_lock_handle.join().unwrap();
assert_ne!(child.wait().unwrap().code(), Some(0));
assert!(!socket_path.exists());
assert!(read_daemon_state(&socket_path).is_none());
}
#[test]
fn test_recover_unhealthy_daemon_terminates_stale_recorded_pid() {
let dir = tempfile::tempdir().unwrap();
let socket_path = dir.path().join("daemon.sock");
std::fs::write(&socket_path, b"stale").unwrap();
let run_lock_handle = hold_run_lock_for_test(&socket_path, Duration::from_millis(150));
let mut child = spawn_blocking_child();
let state = DaemonCoordState {
pid: child.id(),
build_epoch: build_epoch(),
phase: DaemonPhase::Ready,
updated_at_ms: now_millis()
.saturating_sub(DAEMON_COORD_STALE_AFTER.as_millis() as u64 + 1),
};
write_json_atomically(&daemon_state_path(&socket_path), &state).unwrap();
assert!(recover_unhealthy_daemon(&socket_path, "test").unwrap());
run_lock_handle.join().unwrap();
assert_ne!(child.wait().unwrap().code(), Some(0));
assert!(!socket_path.exists());
assert!(read_daemon_state(&socket_path).is_none());
}
#[test]
fn test_recover_unhealthy_daemon_does_not_kill_pid_without_run_lock() {
let dir = tempfile::tempdir().unwrap();
let socket_path = dir.path().join("daemon.sock");
std::fs::write(&socket_path, b"stale").unwrap();
let mut child = spawn_blocking_child();
let state = DaemonCoordState {
pid: child.id(),
build_epoch: build_epoch(),
phase: DaemonPhase::Ready,
updated_at_ms: now_millis(),
};
write_json_atomically(&daemon_state_path(&socket_path), &state).unwrap();
assert!(recover_unhealthy_daemon(&socket_path, "test").unwrap());
assert!(child.try_wait().unwrap().is_none());
let _ = child.kill();
let _ = child.wait();
assert!(!socket_path.exists());
assert!(read_daemon_state(&socket_path).is_none());
}
#[test]
fn test_recover_unhealthy_daemon_refuses_held_lock_without_state() {
let dir = tempfile::tempdir().unwrap();
let socket_path = dir.path().join("daemon.sock");
let run_lock_handle = hold_run_lock_for_test(&socket_path, Duration::from_millis(150));
assert!(!recover_unhealthy_daemon(&socket_path, "test").unwrap());
run_lock_handle.join().unwrap();
}
#[test]
fn test_request_gc_serde() {
let req = Request::Gc(GcRequest {
max_age_hours: Some(168),
});
let json = serde_json::to_string(&req).unwrap();
let parsed: Request = serde_json::from_str(&json).unwrap();
assert_eq!(req, parsed);
assert!(json.contains("\"gc\""));
assert!(json.contains("\"max_age_hours\":168"));
}
#[test]
fn test_request_gc_null_age_serde() {
let req = Request::Gc(GcRequest {
max_age_hours: None,
});
let json = serde_json::to_string(&req).unwrap();
let parsed: Request = serde_json::from_str(&json).unwrap();
assert_eq!(req, parsed);
assert!(json.contains("\"max_age_hours\":null"));
}
#[test]
fn test_request_remote_check_serde() {
let req = Request::RemoteCheck(RemoteCheckRequest {
key: "abc123".into(),
entry_dir: "/tmp/store/abc123".into(),
crate_name: String::new(),
});
let json = serde_json::to_string(&req).unwrap();
let parsed: Request = serde_json::from_str(&json).unwrap();
assert_eq!(req, parsed);
assert!(json.contains("\"remote_check\""));
assert!(json.contains("\"key\":\"abc123\""));
assert!(json.contains("\"entry_dir\":\"/tmp/store/abc123\""));
}
#[test]
fn test_response_ok_serde() {
let resp = Response::ok();
let json = serde_json::to_string(&resp).unwrap();
assert_eq!(json, r#"{"ok":true}"#);
}
#[test]
fn test_response_ok_evicted_serde() {
let resp = Response::ok_evicted(5);
let json = serde_json::to_string(&resp).unwrap();
assert_eq!(json, r#"{"ok":true,"evicted":5}"#);
}
#[test]
fn test_response_gc_skipped_serde() {
let resp = Response::ok_gc_skipped();
let json = serde_json::to_string(&resp).unwrap();
assert_eq!(json, r#"{"ok":true,"evicted":0,"skipped":true}"#);
}
#[test]
fn test_response_found_true_serde() {
let resp = Response::found(true);
let json = serde_json::to_string(&resp).unwrap();
assert_eq!(json, r#"{"ok":true,"found":true}"#);
}
#[test]
fn test_response_found_false_serde() {
let resp = Response::found(false);
let json = serde_json::to_string(&resp).unwrap();
assert_eq!(json, r#"{"ok":true,"found":false}"#);
}
#[test]
fn test_response_found_prefetched_serde() {
let resp = Response::found_prefetched(true, true);
let json = serde_json::to_string(&resp).unwrap();
assert_eq!(json, r#"{"ok":true,"found":true,"prefetched":true}"#);
}
#[test]
fn test_stats_response_prefetch_field_is_backward_compatible() {
let old_json = r#"{"total_size":0,"max_size":0,"entry_count":0,"entries":null,
"events":{"local_hits":0,"prefetch_hits":0,"remote_hits":0,"dups":0,
"misses":0,"errors":0,"total_elapsed_ms":0,"hit_elapsed_ms":0,
"miss_elapsed_ms":0,"hit_compile_time_ms":0,"miss_compile_time_ms":0,
"store_output_blobs":0,"store_duplicate_blobs":0,"store_new_blobs":0}}"#;
let parsed: StatsResponse = serde_json::from_str(old_json).unwrap();
assert_eq!(parsed.prefetch, PrefetchStatsSnapshot::default());
let snap = PrefetchStatsSnapshot {
downloads_completed: 3,
bytes_downloaded: 1024,
keys_used: 2,
keys_cancelled: 1,
keys_over_budget: 5,
cancelled: true,
plans_advisory: 1,
plans_fallback: 4,
last_plan_candidates: 17,
dedup_join_waits: 2,
dedup_join_wait_ms: 250,
last_list_duration_ms: 42,
last_list_key_count: 9001,
list_requests_total: 7,
list_failures_total: 1,
list_duration_ms_total: 900,
list_keys_total: 63007,
};
let json = serde_json::to_string(&snap).unwrap();
let back: PrefetchStatsSnapshot = serde_json::from_str(&json).unwrap();
assert_eq!(back, snap);
}
#[test]
fn test_response_err_serde() {
let resp = Response::err("something broke");
let json = serde_json::to_string(&resp).unwrap();
let parsed: Response = serde_json::from_str(&json).unwrap();
assert!(!parsed.ok);
assert_eq!(parsed.error.as_deref(), Some("something broke"));
assert_eq!(parsed.evicted, None);
assert_eq!(parsed.found, None);
}
#[test]
fn test_invalid_request_json() {
let result = serde_json::from_str::<Request>(r#"{"bogus": 42}"#);
assert!(result.is_err());
}
#[tokio::test]
async fn test_key_cache_unpopulated_returns_none() {
let cache = S3KeyCache::new();
assert_eq!(cache.check("any_key").await, None);
}
#[tokio::test]
async fn test_key_cache_populate_and_check() {
let cache = S3KeyCache::new();
let mut keys = HashMap::new();
keys.insert("key_a".to_string(), "crate_a".to_string());
keys.insert("key_b".to_string(), "crate_b".to_string());
cache.populate(keys).await;
assert_eq!(cache.check("key_a").await, Some(true));
assert_eq!(cache.check("key_b").await, Some(true));
assert_eq!(cache.check("key_c").await, Some(false));
let crate_a_keys = cache.keys_for_crate("crate_a").await;
assert_eq!(crate_a_keys, vec!["key_a"]);
assert!(cache.keys_for_crate("unknown").await.is_empty());
}
#[tokio::test]
async fn test_key_cache_insert_after_populate() {
let cache = S3KeyCache::new();
cache.populate(HashMap::new()).await;
assert_eq!(cache.check("new_key").await, Some(false));
cache.insert("new_key".to_string(), Some("my_crate")).await;
assert_eq!(cache.check("new_key").await, Some(true));
let keys = cache.keys_for_crate("my_crate").await;
assert_eq!(keys, vec!["new_key"]);
}
#[tokio::test]
async fn test_key_cache_insert_before_populate_is_noop() {
let cache = S3KeyCache::new();
cache.insert("key".to_string(), Some("crate")).await;
assert_eq!(cache.check("key").await, None);
assert!(cache.keys_for_crate("crate").await.is_empty());
}
#[tokio::test]
async fn test_key_cache_views_stay_consistent_under_concurrency() {
use std::sync::Arc;
let cache = Arc::new(S3KeyCache::new());
let seed: HashMap<String, String> = (0..50)
.map(|i| (format!("seed_{i}"), format!("crate_{}", i % 5)))
.collect();
cache.populate(seed).await;
let mut tasks = Vec::new();
for r in 0..8 {
let c = cache.clone();
tasks.push(tokio::spawn(async move {
let mut m: HashMap<String, String> = (0..50)
.map(|i| (format!("seed_{i}"), format!("crate_{}", i % 5)))
.collect();
m.insert(format!("refresh_{r}"), "crate_r".to_string());
c.populate(m).await;
}));
}
for k in 0..8 {
let c = cache.clone();
tasks.push(tokio::spawn(async move {
c.insert(format!("up_{k}"), Some("crate_up")).await;
}));
}
for t in tasks {
t.await.unwrap();
}
assert_eq!(cache.check("seed_0").await, Some(true));
let guard = cache.index.read().await;
let idx = guard.as_ref().expect("populated");
let reverse_total: usize = idx.by_crate.values().map(Vec::len).sum();
assert_eq!(
idx.keys.len(),
reverse_total,
"forward set and reverse index must agree on key count"
);
for keys in idx.by_crate.values() {
for key in keys {
assert!(
idx.keys.contains(key),
"key {key} is in by_crate but missing from the forward set"
);
}
}
}
#[test]
fn test_handle_gc_empty_store() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Daemon::new(config);
let resp = daemon.handle_gc(&GcRequest {
max_age_hours: None,
});
assert!(resp.ok);
assert_eq!(resp.evicted, Some(0));
}
#[test]
fn test_handle_gc_reports_lock_skip() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let store = Store::open(&config).unwrap();
let _gc_lock = store.try_gc_lock().unwrap().expect("gc lock");
let daemon = Daemon::new(config);
let resp = daemon.handle_gc(&GcRequest {
max_age_hours: None,
});
assert!(resp.ok);
assert!(resp.skipped);
assert_eq!(resp.evicted, Some(0));
}
#[test]
fn test_handle_gc_with_max_age() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Daemon::new(config);
let resp = daemon.handle_gc(&GcRequest {
max_age_hours: Some(24),
});
assert!(resp.ok);
assert_eq!(resp.evicted, Some(0));
}
#[test]
fn test_handle_gc_evicts_entries() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
let src_file = dir.path().join("big.rlib");
std::fs::write(&src_file, vec![0u8; 200]).unwrap();
let store = Store::open(&config).unwrap();
store
.put(
"testkey",
"testcrate",
&["lib".into()],
&[],
"host",
"dev",
&[(src_file, "lib.rlib".into())],
"",
"",
)
.unwrap();
assert!(store.contains("testkey"));
assert!(store.total_size().unwrap() >= 200);
store.set_last_accessed_for_test("testkey", "-1 hour");
drop(store);
config.max_size = 100;
let daemon = Daemon::new(config);
let stats = daemon.run_gc(None).unwrap();
assert!(
stats.entries_evicted > 0,
"should have evicted at least 1 entry"
);
}
#[test]
fn test_upload_triggered_eviction_respects_gc_lock() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.max_size = 100;
let src_file = dir.path().join("big.rlib");
std::fs::write(&src_file, vec![0u8; 200]).unwrap();
let store = Store::open(&config).unwrap();
store
.put(
"upload_evict_key",
"testcrate",
&["lib".into()],
&[],
"host",
"dev",
&[(src_file, "lib.rlib".into())],
"",
"",
)
.unwrap();
store.set_last_accessed_for_test("upload_evict_key", "-1 hour");
let gc_lock = store.try_gc_lock().unwrap().expect("gc lock");
let daemon = Daemon::new(config);
daemon.maybe_evict_after_upload();
assert!(
store.contains("upload_evict_key"),
"upload-triggered eviction must skip while gc.lock is held"
);
drop(gc_lock);
daemon.maybe_evict_after_upload();
assert!(
!store.contains("upload_evict_key"),
"eviction should run once gc.lock is available"
);
}
#[test]
fn test_handle_request_sync_dispatches_gc() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Daemon::new(config);
let req = Request::Gc(GcRequest {
max_age_hours: None,
});
let resp = daemon.handle_request_sync(&req);
assert!(resp.ok);
assert_eq!(resp.evicted, Some(0));
}
#[tokio::test]
async fn offload_returns_the_handler_response() {
let resp = offload(Response::ok).await;
assert!(resp.ok);
}
#[tokio::test]
async fn offload_maps_a_handler_panic_to_an_error_response() {
let resp = offload(|| panic!("handler boom")).await;
assert!(!resp.ok, "a panicking handler must yield an error response");
assert!(
resp.error
.as_deref()
.unwrap_or_default()
.contains("task failed"),
"error should explain the handler task failed, got {:?}",
resp.error
);
}
#[test]
fn test_handle_request_sync_rejects_upload() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Daemon::new(config);
let req = Request::Upload(UploadJob {
key: "k".into(),
entry_dir: "/tmp".into(),
crate_name: String::new(),
client_epoch: 0,
});
let resp = daemon.handle_request_sync(&req);
assert!(!resp.ok);
assert!(resp.error.as_deref().unwrap().contains("async"));
}
#[test]
fn test_handle_request_sync_rejects_remote_check() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Daemon::new(config);
let req = Request::RemoteCheck(RemoteCheckRequest {
key: "k".into(),
entry_dir: "/tmp".into(),
crate_name: String::new(),
});
let resp = daemon.handle_request_sync(&req);
assert!(!resp.ok);
assert!(resp.error.as_deref().unwrap().contains("async"));
}
#[tokio::test]
async fn test_handle_upload_no_remote() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path()); let daemon = Daemon::new(config);
let job = UploadJob {
key: "k".into(),
entry_dir: "/tmp".into(),
crate_name: String::new(),
client_epoch: 0,
};
let resp = daemon.handle_upload(&job).await;
assert!(!resp.ok);
assert!(
resp.error
.as_deref()
.unwrap()
.contains("no remote configured")
);
}
#[tokio::test]
async fn test_handle_upload_remote_readonly() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote_readonly = true;
let daemon = Daemon::new(config);
let job = UploadJob {
key: "k".into(),
entry_dir: "/tmp".into(),
crate_name: String::new(),
client_epoch: 0,
};
let resp = daemon.handle_upload(&job).await;
assert!(resp.ok);
assert!(resp.error.is_none());
let resp_do = daemon.do_upload(&job).await;
assert!(resp_do.ok);
assert!(resp_do.error.is_none());
}
#[tokio::test]
async fn test_handle_remote_check_no_remote() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path()); let daemon = Daemon::new(config);
let req = RemoteCheckRequest {
key: "k".into(),
entry_dir: "/tmp".into(),
crate_name: String::new(),
};
let resp = daemon.handle_remote_check(&req).await;
assert!(!resp.ok);
assert!(
resp.error
.as_deref()
.unwrap()
.contains("no remote configured")
);
}
#[test]
fn test_run_gc_returns_count() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Daemon::new(config);
let stats = daemon.run_gc(None).unwrap();
assert_eq!(stats.entries_evicted, 0);
}
#[test]
fn test_run_gc_cleans_registered_incremental_dirs_once() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.clean_incremental = true;
let incremental_dir = dir.path().join("workspace/target/debug/incremental");
std::fs::create_dir_all(&incremental_dir).unwrap();
std::fs::write(incremental_dir.join("junk"), b"tmp").unwrap();
let store = Store::open(&config).unwrap();
store.remember_incremental_dir(&incremental_dir).unwrap();
drop(store);
let daemon = Daemon::new(config.clone());
let stats = daemon.run_gc(None).unwrap();
assert_eq!(stats.entries_evicted, 0);
assert!(!incremental_dir.exists());
std::fs::create_dir_all(&incremental_dir).unwrap();
std::fs::write(incremental_dir.join("junk"), b"tmp2").unwrap();
let stats = daemon.run_gc(None).unwrap();
assert_eq!(stats.entries_evicted, 0);
assert!(incremental_dir.exists());
}
#[test]
fn clean_tool_version_caches_removes_only_old_tool_version_txt() {
let dir = tempfile::tempdir().unwrap();
let old_rustc = dir.path().join("rustc-ver-old.txt");
let old_linker = dir.path().join("linker-ver-old.txt");
let fresh_rustc = dir.path().join("rustc-ver-fresh.txt");
let old_other = dir.path().join("other-ver-old.txt");
for path in [&old_rustc, &old_linker, &fresh_rustc, &old_other] {
std::fs::write(path, b"version").unwrap();
}
let old = filetime::FileTime::from_system_time(
std::time::SystemTime::now() - Duration::from_secs(8 * 24 * 3600),
);
for path in [&old_rustc, &old_linker, &old_other] {
filetime::set_file_mtime(path, old).unwrap();
}
Daemon::clean_tool_version_caches(dir.path());
assert!(!old_rustc.exists());
assert!(!old_linker.exists());
assert!(fresh_rustc.exists());
assert!(old_other.exists());
}
#[tokio::test]
async fn test_socket_gc_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let daemon = Arc::new(Daemon::new(config));
let resp = one_shot_request(
&daemon,
&socket_path,
&Request::Gc(GcRequest {
max_age_hours: None,
}),
)
.await;
assert!(resp.ok);
assert_eq!(resp.evicted, Some(0));
}
#[tokio::test]
async fn test_socket_remote_check_no_remote_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path()); let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let daemon = Arc::new(Daemon::new(config));
let resp = one_shot_request(
&daemon,
&socket_path,
&Request::RemoteCheck(RemoteCheckRequest {
key: "test_key".into(),
entry_dir: "/tmp/test".into(),
crate_name: String::new(),
}),
)
.await;
assert!(!resp.ok);
assert!(
resp.error
.as_deref()
.unwrap()
.contains("no remote configured")
);
}
#[cfg(unix)]
#[test]
fn test_stale_socket_cleanup() {
let dir = tempfile::tempdir().unwrap();
let socket_path = dir.path().join("daemon.sock");
std::fs::write(&socket_path, b"stale").unwrap();
assert!(socket_path.exists());
let result = std::os::unix::net::UnixStream::connect(&socket_path);
assert!(result.is_err());
std::fs::remove_file(&socket_path).unwrap();
assert!(!socket_path.exists());
}
#[test]
fn test_send_request_to_nonexistent_socket() {
let dir = tempfile::tempdir().unwrap();
let socket_path = dir.path().join("nonexistent.sock");
let req = Request::Gc(GcRequest {
max_age_hours: None,
});
let result = send_request(&socket_path, &req);
assert!(result.is_err());
}
#[test]
fn test_send_remote_check_unreachable_returns_none() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let result = send_remote_check(&config, "some_key", Path::new("/tmp/test"), "unknown");
assert!(result.is_none());
}
#[test]
fn remote_check_response_parser_handles_prefetched_error_and_malformed() {
let hit = serde_json::to_string(&Response::found_prefetched(true, true)).unwrap();
let result = remote_check_result_from_response_line(&hit).unwrap();
assert!(result.found);
assert!(result.prefetched);
let plain_hit = serde_json::to_string(&Response::found(true)).unwrap();
let result = remote_check_result_from_response_line(&plain_hit).unwrap();
assert!(result.found);
assert!(!result.prefetched);
let err = serde_json::to_string(&Response::err("remote down")).unwrap();
assert!(remote_check_result_from_response_line(&err).is_none());
assert!(remote_check_result_from_response_line("{not json").is_none());
}
#[test]
fn test_response_constructors() {
let ok = Response::ok();
assert!(ok.ok && ok.evicted.is_none() && ok.error.is_none() && ok.found.is_none());
assert!(ok.batch_results.is_none());
let evicted = Response::ok_evicted(3);
assert!(evicted.ok && evicted.evicted == Some(3));
let found_true = Response::found(true);
assert!(found_true.ok && found_true.found == Some(true));
let found_false = Response::found(false);
assert!(found_false.ok && found_false.found == Some(false));
let batch = Response::ok_batch(vec![Response::found(true), Response::found(false)]);
assert!(batch.ok && batch.batch_results.as_ref().unwrap().len() == 2);
let err = Response::err("oops");
assert!(!err.ok && err.error.as_deref() == Some("oops"));
}
#[test]
fn test_stats_request_serde() {
let req = Request::Stats(StatsRequest {
include_entries: true,
sort_by: Some("size".into()),
event_hours: Some(48),
client_epoch: 0,
});
let json = serde_json::to_string(&req).unwrap();
let parsed: Request = serde_json::from_str(&json).unwrap();
assert_eq!(req, parsed);
assert!(json.contains("\"stats\""));
assert!(json.contains("\"include_entries\":true"));
assert!(json.contains("\"sort_by\":\"size\""));
assert!(json.contains("\"event_hours\":48"));
}
#[test]
fn test_stats_response_serde() {
let stats = StatsResponse {
total_size: 1024,
max_size: 4096,
entry_count: 5,
entries: None,
events: EventStatsResponse {
local_hits: 10,
prefetch_hits: 0,
remote_hits: 2,
dups: 1,
misses: 3,
errors: 1,
total_elapsed_ms: 5000,
hit_elapsed_ms: 120,
miss_elapsed_ms: 4880,
hit_compile_time_ms: 22000,
miss_compile_time_ms: 9000,
store_output_blobs: 4,
store_duplicate_blobs: 1,
store_new_blobs: 3,
},
version: String::new(),
build_epoch: 0,
pending_uploads: 0,
active_downloads: 0,
s3_concurrency_total: 0,
s3_concurrency_used: 0,
upload_queue_capacity: 0,
uploads_completed: 0,
uploads_failed: 0,
uploads_skipped: 0,
downloads_completed: 0,
downloads_failed: 0,
bytes_uploaded: 0,
bytes_downloaded: 0,
recent_transfers: Vec::new(),
prefetch: PrefetchStatsSnapshot::default(),
in_flight: Vec::new(),
};
let resp = Response::ok_stats(stats.clone());
let json = serde_json::to_string(&resp).unwrap();
let parsed: Response = serde_json::from_str(&json).unwrap();
assert!(parsed.ok);
let parsed_stats = parsed.stats.unwrap();
assert_eq!(parsed_stats, stats);
}
#[test]
fn test_stats_response_with_entries() {
let stats = StatsResponse {
total_size: 2048,
max_size: 8192,
entry_count: 2,
entries: Some(vec![
StatsEntry {
cache_key: "abc123def456".into(),
crate_name: "serde".into(),
crate_type: "lib".into(),
profile: "release".into(),
size: 1024,
hit_count: 5,
created_at: "2025-01-01 00:00:00".into(),
last_accessed: "2025-06-01 12:00:00".into(),
content_hash: None,
},
StatsEntry {
cache_key: "789abc012def".into(),
crate_name: "tokio".into(),
crate_type: "lib".into(),
profile: "dev".into(),
size: 1024,
hit_count: 3,
created_at: "2025-02-01 00:00:00".into(),
last_accessed: "2025-05-15 08:00:00".into(),
content_hash: None,
},
]),
events: EventStatsResponse {
local_hits: 0,
prefetch_hits: 0,
remote_hits: 0,
dups: 0,
misses: 0,
errors: 0,
total_elapsed_ms: 0,
hit_elapsed_ms: 0,
miss_elapsed_ms: 0,
hit_compile_time_ms: 0,
miss_compile_time_ms: 0,
store_output_blobs: 0,
store_duplicate_blobs: 0,
store_new_blobs: 0,
},
version: String::new(),
build_epoch: 0,
pending_uploads: 0,
active_downloads: 0,
s3_concurrency_total: 0,
s3_concurrency_used: 0,
upload_queue_capacity: 0,
uploads_completed: 0,
uploads_failed: 0,
uploads_skipped: 0,
downloads_completed: 0,
downloads_failed: 0,
bytes_uploaded: 0,
bytes_downloaded: 0,
recent_transfers: Vec::new(),
prefetch: PrefetchStatsSnapshot::default(),
in_flight: Vec::new(),
};
let resp = Response::ok_stats(stats);
let json = serde_json::to_string(&resp).unwrap();
let parsed: Response = serde_json::from_str(&json).unwrap();
let entries = parsed.stats.unwrap().entries.unwrap();
assert_eq!(entries.len(), 2);
assert_eq!(entries[0].crate_name, "serde");
assert_eq!(entries[1].crate_name, "tokio");
}
#[test]
fn test_handle_stats_empty_store() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Daemon::new(config);
let resp = daemon.handle_stats(&StatsRequest {
include_entries: true,
sort_by: None,
event_hours: Some(24),
client_epoch: 0,
});
assert!(resp.ok);
let stats = resp.stats.unwrap();
assert_eq!(stats.total_size, 0);
assert_eq!(stats.entry_count, 0);
assert_eq!(stats.max_size, 50 * 1024 * 1024);
assert_eq!(stats.entries.unwrap().len(), 0);
assert_eq!(stats.events.local_hits, 0);
assert_eq!(stats.events.misses, 0);
}
#[test]
fn test_daemon_reuses_store_handle() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Daemon::new(config);
let first = daemon.store_lock().unwrap() as *const _;
let second = daemon.store_lock().unwrap() as *const _;
assert_eq!(first, second);
}
#[test]
fn test_handle_hash_files_uses_memory_cache() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Daemon::new(config);
let file = dir.path().join("large.rlib");
std::fs::write(&file, vec![7u8; 70 * 1024]).unwrap();
let metadata = std::fs::metadata(&file).unwrap();
let req = HashFilesRequest {
files: vec![HashFileRequest {
path: file.to_string_lossy().into_owned(),
size: i64::try_from(metadata.len()).unwrap(),
mtime_ns: crate::cache_key::metadata_mtime_ns(&metadata),
ctime_ns: crate::cache_key::metadata_ctime_ns(&metadata),
inode: crate::cache_key::metadata_inode(&metadata),
}],
};
let first = daemon.handle_hash_files(&req);
assert!(first.ok);
let first_result = &first.hash_results.as_ref().unwrap()[0];
assert!(first_result.hash.is_some());
assert!(!first_result.cache_hit);
assert!(first_result.bytes_hashed > 0);
let second = daemon.handle_hash_files(&req);
assert!(second.ok);
let second_result = &second.hash_results.as_ref().unwrap()[0];
assert_eq!(first_result.hash, second_result.hash);
assert!(second_result.cache_hit);
assert_eq!(second_result.bytes_hashed, 0);
}
#[test]
fn handle_hash_files_persistent_cache_hit_across_daemons() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let file = dir.path().join("big.rlib");
std::fs::write(&file, vec![3u8; 80 * 1024]).unwrap(); let metadata = std::fs::metadata(&file).unwrap();
let req = HashFilesRequest {
files: vec![HashFileRequest {
path: file.to_string_lossy().into_owned(),
size: i64::try_from(metadata.len()).unwrap(),
mtime_ns: crate::cache_key::metadata_mtime_ns(&metadata),
ctime_ns: crate::cache_key::metadata_ctime_ns(&metadata),
inode: crate::cache_key::metadata_inode(&metadata),
}],
};
let expected = crate::cache_key::hash_file(&file).unwrap();
let a = Daemon::new(config.clone());
let ra = a.handle_hash_files(&req);
let ra = &ra.hash_results.as_ref().unwrap()[0];
assert_eq!(ra.hash.as_deref(), Some(expected.as_str()));
assert!(!ra.cache_hit, "first hash is a persistent-cache miss");
assert!(ra.bytes_hashed > 0);
let b = Daemon::new(config);
let rb = b.handle_hash_files(&req);
let rb = &rb.hash_results.as_ref().unwrap()[0];
assert_eq!(rb.hash.as_deref(), Some(expected.as_str()));
assert!(rb.cache_hit, "second daemon must hit the persistent cache");
assert_eq!(rb.bytes_hashed, 0);
}
#[test]
fn handle_hash_files_rejects_changed_metadata_before_hashing() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Daemon::new(config);
let file = dir.path().join("input.bin");
std::fs::write(&file, b"stable bytes").unwrap();
let metadata = std::fs::metadata(&file).unwrap();
let resp = daemon.handle_hash_files(&HashFilesRequest {
files: vec![HashFileRequest {
path: file.to_string_lossy().into_owned(),
size: i64::try_from(metadata.len()).unwrap() + 1,
mtime_ns: crate::cache_key::metadata_mtime_ns(&metadata),
ctime_ns: crate::cache_key::metadata_ctime_ns(&metadata),
inode: crate::cache_key::metadata_inode(&metadata),
}],
});
assert!(resp.ok);
let result = &resp.hash_results.as_ref().unwrap()[0];
assert_eq!(result.hash, None);
assert_eq!(
result.error.as_deref(),
Some("file metadata changed before hashing")
);
}
#[test]
fn handle_hash_files_reports_hash_read_error() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Daemon::new(config);
let input_dir = dir.path().join("not-a-file");
std::fs::create_dir(&input_dir).unwrap();
let metadata = std::fs::metadata(&input_dir).unwrap();
let resp = daemon.handle_hash_files(&HashFilesRequest {
files: vec![HashFileRequest {
path: input_dir.to_string_lossy().into_owned(),
size: i64::try_from(metadata.len()).unwrap(),
mtime_ns: crate::cache_key::metadata_mtime_ns(&metadata),
ctime_ns: crate::cache_key::metadata_ctime_ns(&metadata),
inode: crate::cache_key::metadata_inode(&metadata),
}],
});
assert!(resp.ok);
let result = &resp.hash_results.as_ref().unwrap()[0];
assert_eq!(result.hash, None);
assert_eq!(result.bytes_hashed, 0);
assert!(
result
.error
.as_deref()
.unwrap_or_default()
.contains("reading"),
"got {:?}",
result.error
);
}
#[test]
fn test_handle_stats_with_store_entries() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let src_file = dir.path().join("lib.rlib");
std::fs::write(&src_file, vec![0u8; 100]).unwrap();
let store = Store::open(&config).unwrap();
store
.put(
"key1",
"mycrate",
&["lib".into()],
&[],
"host",
"dev",
&[(src_file, "lib.rlib".into())],
"",
"",
)
.unwrap();
drop(store);
let daemon = Daemon::new(config);
let resp = daemon.handle_stats(&StatsRequest {
include_entries: true,
sort_by: Some("size".into()),
event_hours: Some(24),
client_epoch: 0,
});
assert!(resp.ok);
let stats = resp.stats.unwrap();
assert_eq!(stats.entry_count, 1);
assert!(stats.total_size >= 100);
let entries = stats.entries.unwrap();
assert_eq!(entries.len(), 1);
assert_eq!(entries[0].crate_name, "mycrate");
}
#[test]
fn test_handle_request_sync_dispatches_stats() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Daemon::new(config);
let req = Request::Stats(StatsRequest {
include_entries: false,
sort_by: None,
event_hours: None,
client_epoch: 0,
});
let resp = daemon.handle_request_sync(&req);
assert!(resp.ok);
assert!(resp.stats.is_some());
}
#[test]
fn test_send_stats_request_unreachable() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let result = send_stats_request(&config, false, None, None);
assert!(result.is_err());
}
#[tokio::test]
async fn test_socket_stats_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let daemon = Arc::new(Daemon::new(config));
let resp = one_shot_request(
&daemon,
&socket_path,
&Request::Stats(StatsRequest {
include_entries: true,
sort_by: Some("size".into()),
event_hours: Some(24),
client_epoch: 0,
}),
)
.await;
assert!(resp.ok);
let stats = resp.stats.unwrap();
assert_eq!(stats.total_size, 0);
assert_eq!(stats.entry_count, 0);
assert!(stats.entries.unwrap().is_empty());
}
#[tokio::test]
async fn test_socket_local_lookup_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let store = Store::open(&config).unwrap();
let output_file = dir.path().join("out.rlib");
std::fs::write(&output_file, b"artifact-bytes").unwrap();
store
.put(
"0000000000000000000000000000000000000000000000000000000000000001",
"probe_crate",
&["lib".to_string()],
&[],
"x86_64-unknown-linux-gnu",
"dev",
&[(output_file, "libout.rlib".to_string())],
"cached stdout",
"",
)
.unwrap();
let index_db = crate::store::open_index_db(&config.index_db_path()).unwrap();
index_db
.execute(
"UPDATE entries SET last_accessed = datetime('now', '-1 hour')",
[],
)
.unwrap();
drop(store);
let daemon = Arc::new(Daemon::new(config));
let key = "0000000000000000000000000000000000000000000000000000000000000001";
let resp = one_shot_request(
&daemon,
&socket_path,
&Request::LocalLookup(LocalLookupRequest {
key: key.to_string(),
client_epoch: 0,
}),
)
.await;
assert!(resp.ok);
let reply = resp.local_lookup.expect("local_lookup payload");
assert_eq!(reply.outcome, "hit");
let meta = reply.meta.expect("hit carries meta");
assert_eq!(meta.cache_key, key);
assert_eq!(meta.stdout, "cached stdout");
assert_eq!(meta.files.len(), 1);
let (hits, recent): (i64, i64) = index_db
.query_row(
"SELECT hit_count, last_accessed >= datetime('now', '-60 seconds')
FROM entries WHERE cache_key = ?1",
[key],
|row| Ok((row.get(0)?, row.get(1)?)),
)
.unwrap();
assert_eq!(hits, 1, "hit must be accounted by the pin writer");
assert_eq!(recent, 1, "pin must refresh last_accessed before the reply");
let resp = one_shot_request(
&daemon,
&socket_path,
&Request::LocalLookup(LocalLookupRequest {
key: "0000000000000000000000000000000000000000000000000000000000000002".to_string(),
client_epoch: 0,
}),
)
.await;
assert!(resp.ok);
assert_eq!(resp.local_lookup.expect("payload").outcome, "miss");
}
#[tokio::test]
async fn test_socket_hash_files_roundtrip_hashes_a_real_file() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let file_path = dir.path().join("input.bin");
std::fs::write(&file_path, b"hash me please").unwrap();
let meta = std::fs::metadata(&file_path).unwrap();
let req = HashFileRequest {
path: file_path.to_string_lossy().into_owned(),
size: meta.len() as i64,
mtime_ns: crate::cache_key::metadata_mtime_ns(&meta),
ctime_ns: crate::cache_key::metadata_ctime_ns(&meta),
inode: crate::cache_key::metadata_inode(&meta),
};
let expected = blake3::hash(b"hash me please").to_hex().to_string();
let daemon = Arc::new(Daemon::new(config));
let resp = one_shot_request(
&daemon,
&socket_path,
&Request::HashFiles(HashFilesRequest { files: vec![req] }),
)
.await;
assert!(resp.ok, "hash-files request should succeed: {resp:?}");
let results = resp.hash_results.expect("hash_results present");
assert_eq!(results.len(), 1);
assert_eq!(results[0].hash.as_deref(), Some(expected.as_str()));
assert_eq!(results[0].error, None);
}
#[tokio::test]
async fn test_socket_hash_files_missing_file_reports_error_result() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let req = HashFileRequest {
path: dir
.path()
.join("does-not-exist")
.to_string_lossy()
.into_owned(),
size: 10,
mtime_ns: 0,
ctime_ns: 0,
inode: 0,
};
let daemon = Arc::new(Daemon::new(config));
let resp = one_shot_request(
&daemon,
&socket_path,
&Request::HashFiles(HashFilesRequest { files: vec![req] }),
)
.await;
assert!(resp.ok);
let results = resp.hash_results.expect("hash_results present");
assert_eq!(results.len(), 1);
assert!(results[0].hash.is_none());
assert!(results[0].error.is_some(), "missing file should error");
}
#[test]
fn send_hash_files_request_empty_is_ok_without_socket() {
let result = send_hash_files_request(Path::new("/nonexistent/socket"), Vec::new()).unwrap();
assert!(result.is_empty());
}
#[test]
fn send_hash_files_request_missing_socket_errors() {
let req = HashFileRequest {
path: "/some/file".into(),
size: 1,
mtime_ns: 0,
ctime_ns: 0,
inode: 0,
};
let err = send_hash_files_request(Path::new("/nonexistent/socket.sock"), vec![req])
.expect_err("missing socket -> error");
assert!(
err.to_string().contains("socket does not exist"),
"got: {err}"
);
}
#[test]
fn hash_files_response_parser_handles_results_error_and_malformed() {
let ok = Response::ok_hash_results(vec![HashFileResult {
path: "/tmp/a".into(),
size: 1,
mtime_ns: 2,
ctime_ns: 3,
inode: 4,
hash: Some("abc".into()),
cache_hit: false,
bytes_hashed: 1,
error: None,
}]);
let ok_json = serde_json::to_string(&ok).unwrap();
assert_eq!(
hash_files_results_from_response_line(&ok_json)
.unwrap()
.len(),
1
);
let err_json = serde_json::to_string(&Response::err("bad hash")).unwrap();
let err = hash_files_results_from_response_line(&err_json).unwrap_err();
assert!(err.to_string().contains("daemon hash_files error"));
let err = hash_files_results_from_response_line("{not json").unwrap_err();
assert!(err.to_string().contains("key must be a string"));
}
#[cfg(unix)]
#[tokio::test]
async fn send_hash_files_request_client_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let file_path = dir.path().join("input.bin");
std::fs::write(&file_path, b"hash me please").unwrap();
let meta = std::fs::metadata(&file_path).unwrap();
let req = HashFileRequest {
path: file_path.to_string_lossy().into_owned(),
size: meta.len() as i64,
mtime_ns: crate::cache_key::metadata_mtime_ns(&meta),
ctime_ns: crate::cache_key::metadata_ctime_ns(&meta),
inode: crate::cache_key::metadata_inode(&meta),
};
let expected = blake3::hash(b"hash me please").to_hex().to_string();
let listener = bind_listener(&socket_path);
let daemon = Arc::new(Daemon::new(config.clone()));
let server = tokio::spawn(async move {
let stream = listener.accept().await.expect("accept");
let _ =
handle_connection(stream, &daemon, &AtomicBool::new(false), &Notify::new()).await;
});
let sp = socket_path.clone();
let results = tokio::task::spawn_blocking(move || send_hash_files_request(&sp, vec![req]))
.await
.unwrap()
.expect("send_hash_files_request should succeed");
server.await.unwrap();
assert_eq!(results.len(), 1);
assert_eq!(results[0].hash.as_deref(), Some(expected.as_str()));
}
#[tokio::test]
async fn test_socket_build_started_roundtrip_without_remote() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let daemon = Arc::new(Daemon::new(config));
let resp = one_shot_request(
&daemon,
&socket_path,
&Request::BuildStarted(BuildStartedRequest {
intent: kache_core::BuildIntent {
crate_names: vec!["serde".into()],
namespace: Some("ns".into()),
cargo_lock_deps: vec![],
},
client_epoch: 0,
session_id: String::new(),
}),
)
.await;
assert!(!resp.ok);
assert!(resp.error.is_some());
}
#[tokio::test]
async fn test_socket_batch_remote_check_roundtrip_without_remote() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let daemon = Arc::new(Daemon::new(config));
let resp = one_shot_request(
&daemon,
&socket_path,
&Request::BatchRemoteCheck(BatchRemoteCheckRequest {
checks: vec![RemoteCheckRequest {
key: "k".into(),
entry_dir: "/tmp/k".into(),
crate_name: String::new(),
}],
}),
)
.await;
assert!(resp.batch_results.is_some() || resp.error.is_some());
}
fn seed_store_entry(config: &Config, cache_key: &str, crate_name: &str, dir: &Path) {
let store = Store::open(config).unwrap();
let src = dir.join(format!("{cache_key}-src"));
std::fs::create_dir_all(&src).unwrap();
let artifact = src.join("libfoo.rlib");
std::fs::write(&artifact, b"artifact bytes").unwrap();
store
.put(
cache_key,
crate_name,
&["lib".to_string()],
&[],
"x86_64-unknown-linux-gnu",
"debug",
&[(artifact, "libfoo.rlib".to_string())],
"",
"",
)
.unwrap();
}
#[tokio::test]
async fn test_socket_stats_roundtrip_with_populated_store() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
seed_store_entry(&config, "statskey1", "serde", dir.path());
let daemon = Arc::new(Daemon::new(config));
let resp = one_shot_request(
&daemon,
&socket_path,
&Request::Stats(StatsRequest {
include_entries: true,
sort_by: Some("size".into()),
event_hours: Some(24),
client_epoch: 0,
}),
)
.await;
assert!(resp.ok);
let stats = resp.stats.unwrap();
assert_eq!(stats.entry_count, 1);
assert!(stats.total_size > 0);
let entries = stats.entries.unwrap();
assert_eq!(entries.len(), 1);
assert_eq!(entries[0].crate_name, "serde");
}
#[tokio::test]
async fn test_send_stats_request_client_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
seed_store_entry(&config, "ckey1", "serde", dir.path());
let listener = bind_listener(&socket_path);
let daemon = Arc::new(Daemon::new(config.clone()));
let server = tokio::spawn(async move {
let stream = listener.accept().await.expect("accept");
handle_connection(stream, &daemon, &AtomicBool::new(false), &Notify::new())
.await
.expect("handle_connection");
});
let cfg = config.clone();
let stats = tokio::task::spawn_blocking(move || {
send_stats_request(&cfg, true, Some("size"), Some(24))
})
.await
.unwrap()
.expect("send_stats_request should succeed");
server.await.unwrap();
assert_eq!(stats.entry_count, 1);
assert_eq!(stats.entries.unwrap()[0].crate_name, "serde");
}
#[tokio::test]
async fn test_send_gc_request_client_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
seed_store_entry(&config, "gcc1", "serde", dir.path());
let listener = bind_listener(&socket_path);
let daemon = Arc::new(Daemon::new(config.clone()));
let server = tokio::spawn(async move {
let stream = listener.accept().await.expect("accept");
handle_connection(stream, &daemon, &AtomicBool::new(false), &Notify::new())
.await
.expect("handle_connection");
});
let cfg = config.clone();
let outcome = tokio::task::spawn_blocking(move || send_gc_request(&cfg, Some(0)))
.await
.unwrap()
.expect("send_gc_request should succeed");
server.await.unwrap();
assert!(!outcome.skipped);
assert!(outcome.evicted.is_some());
}
#[tokio::test]
async fn test_send_remote_check_client_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(crate::config::RemoteConfig::test_s3("test", "artifacts"));
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let listener = bind_listener(&socket_path);
let daemon = Arc::new(Daemon::new(config.clone()));
daemon.signal_warming_complete();
let mut keys = HashMap::new();
keys.insert("c".repeat(64), "othercrate".to_string());
daemon.key_cache.populate(keys).await;
let server = tokio::spawn(async move {
loop {
let stream = listener.accept().await.expect("accept");
let _ = handle_connection(stream, &daemon, &AtomicBool::new(false), &Notify::new())
.await;
}
});
let cfg = config.clone();
let missing = "d".repeat(64);
let result = tokio::task::spawn_blocking(move || {
send_remote_check(&cfg, &missing, Path::new("/tmp/test"), "crate")
})
.await
.unwrap();
server.abort();
let result = result.expect("authoritative miss yields a definitive result");
assert!(
!result.found,
"the missing key should round-trip as not found"
);
}
#[tokio::test]
async fn test_send_remote_check_error_response_yields_none() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path()); let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let listener = bind_listener(&socket_path);
let daemon = Arc::new(Daemon::new(config.clone()));
daemon.signal_warming_complete();
let server = tokio::spawn(async move {
loop {
let stream = listener.accept().await.expect("accept");
let _ = handle_connection(stream, &daemon, &AtomicBool::new(false), &Notify::new())
.await;
}
});
let cfg = config.clone();
let key = "e".repeat(64);
let result = tokio::task::spawn_blocking(move || {
send_remote_check(&cfg, &key, Path::new("/tmp/test"), "crate")
})
.await
.unwrap();
server.abort();
assert!(
result.is_none(),
"an error response (no remote) must yield None"
);
}
#[tokio::test]
async fn test_send_shutdown_request_client_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let listener = bind_listener(&socket_path);
let daemon = Arc::new(Daemon::new(config.clone()));
let server = tokio::spawn(async move {
let stream = listener.accept().await.expect("accept");
handle_connection(stream, &daemon, &AtomicBool::new(false), &Notify::new())
.await
.expect("handle_connection");
});
let cfg = config.clone();
let result = tokio::task::spawn_blocking(move || send_shutdown_request(&cfg))
.await
.unwrap();
server.await.unwrap();
assert!(result.is_ok(), "shutdown request should round-trip ok");
}
#[tokio::test]
async fn test_socket_gc_roundtrip_evicts_populated_store() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
seed_store_entry(&config, "gckey1", "tokio", dir.path());
let daemon = Arc::new(Daemon::new(config.clone()));
let resp = one_shot_request(
&daemon,
&socket_path,
&Request::Gc(GcRequest {
max_age_hours: Some(0),
}),
)
.await;
assert!(resp.ok, "gc should succeed: {resp:?}");
assert!(resp.evicted.is_some(), "gc reports an evicted count");
}
fn test_remote_config() -> crate::config::RemoteConfig {
crate::config::RemoteConfig::test_s3("bucket", "prefix")
}
fn test_remote_backend() -> Arc<dyn crate::remote_backend::RemoteBackend> {
Arc::new(crate::remote_backend::memory_backend())
}
fn test_manifest_object_key(cache_key: &str, crate_name: &str) -> String {
format!("prefix/v3/manifests/{crate_name}/{cache_key}.json")
}
fn test_pack_object_key(cache_key: &str, crate_name: &str) -> String {
format!("prefix/v3/packs/{crate_name}/{cache_key}.tar.zst")
}
fn test_build_manifest_object_key() -> String {
format!(
"prefix/_manifests/{}.json",
crate::cli::default_manifest_key()
)
}
async fn put_test_object(
backend: &Arc<dyn crate::remote_backend::RemoteBackend>,
key: &str,
body: &[u8],
) {
backend
.put(key, body.to_vec(), None)
.await
.expect("seed test remote object");
}
struct PutFailBackend;
#[async_trait::async_trait]
impl crate::remote_backend::RemoteBackend for PutFailBackend {
async fn head(&self, _key: &str) -> Result<bool> {
Ok(false)
}
async fn get(
&self,
_key: &str,
_max_bytes: Option<u64>,
) -> Result<Option<crate::remote_backend::GetObject>> {
Ok(None)
}
async fn put(&self, _key: &str, _body: Vec<u8>, _content_type: Option<&str>) -> Result<()> {
anyhow::bail!("injected PUT failure")
}
async fn list(&self, _prefix: &str) -> Result<Vec<String>> {
Ok(Vec::new())
}
fn describe(&self, key: &str) -> String {
format!("failure://test/{key}")
}
}
#[tokio::test]
async fn test_socket_remote_check_miss_with_injected_mock_client() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let client = test_remote_backend();
let daemon = Arc::new(Daemon::new(config));
assert!(
daemon.remote_backend.set(client).is_ok(),
"inject mock backend"
);
let resp = one_shot_request(
&daemon,
&socket_path,
&Request::RemoteCheck(RemoteCheckRequest {
key: "abc123def456".into(),
entry_dir: dir.path().join("entry").to_string_lossy().into_owned(),
crate_name: "serde".into(),
}),
)
.await;
assert!(resp.ok, "remote check should return a response: {resp:?}");
assert_eq!(resp.found, Some(false), "missing remote key -> found=false");
}
#[tokio::test]
async fn test_socket_prefetch_empty_keys_lists_remote_then_no_op() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let client = test_remote_backend();
let daemon = Arc::new(Daemon::new(config));
assert!(
daemon.remote_backend.set(client).is_ok(),
"inject mock backend"
);
let resp = one_shot_request(
&daemon,
&socket_path,
&Request::Prefetch(PrefetchRequest {
keys: Vec::new(),
warm_all: false,
}),
)
.await;
assert!(resp.ok, "prefetch over empty remote should be ok: {resp:?}");
}
#[tokio::test]
async fn test_do_upload_skips_when_entry_already_in_remote() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
let client = test_remote_backend();
put_test_object(
&client,
&test_manifest_object_key("abc123def456", "serde"),
b"{}",
)
.await;
let daemon = Arc::new(Daemon::new(config));
assert!(
daemon.remote_backend.set(client).is_ok(),
"inject mock backend"
);
let resp = daemon
.do_upload(&UploadJob {
key: "abc123def456".into(),
entry_dir: dir.path().join("entry").to_string_lossy().into_owned(),
crate_name: "serde".into(),
client_epoch: 0,
})
.await;
assert!(
resp.ok,
"already-present upload should be a no-op ok: {resp:?}"
);
}
#[tokio::test]
async fn test_do_upload_uploads_when_not_in_remote() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
config.prefetch_enabled = false;
seed_store_entry(&config, "upkey123", "serde", dir.path());
let entry_dir = config.store_dir().join("upkey123");
let client = test_remote_backend();
let daemon = Arc::new(Daemon::new(config));
assert!(
daemon.remote_backend.set(client.clone()).is_ok(),
"inject mock backend"
);
let resp = daemon
.do_upload(&UploadJob {
key: "upkey123".into(),
entry_dir: entry_dir.to_string_lossy().into_owned(),
crate_name: "serde".into(),
client_epoch: 0,
})
.await;
assert!(resp.ok, "upload of a new entry should succeed: {resp:?}");
assert!(
client
.head(&test_pack_object_key("upkey123", "serde"))
.await
.unwrap()
);
assert!(
client
.head(&test_manifest_object_key("upkey123", "serde"))
.await
.unwrap()
);
}
#[tokio::test]
async fn test_do_upload_records_failure_when_put_errors() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
seed_store_entry(&config, "upfail1", "serde", dir.path());
let entry_dir = config.store_dir().join("upfail1");
let client: Arc<dyn crate::remote_backend::RemoteBackend> = Arc::new(PutFailBackend);
let daemon = Arc::new(Daemon::new(config));
assert!(
daemon.remote_backend.set(client).is_ok(),
"inject mock backend"
);
let resp = daemon
.do_upload(&UploadJob {
key: "upfail1".into(),
entry_dir: entry_dir.to_string_lossy().into_owned(),
crate_name: "serde".into(),
client_epoch: 0,
})
.await;
assert!(!resp.ok, "a denied upload PUT must fail: {resp:?}");
assert_eq!(
daemon
.transfer_counters
.uploads_failed
.load(Ordering::Relaxed),
1
);
}
#[tokio::test]
async fn test_handle_build_started_falls_back_to_local_planning() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
let client = test_remote_backend();
let daemon = Arc::new(Daemon::new(config));
assert!(
daemon.remote_backend.set(client).is_ok(),
"inject mock backend"
);
let req = BuildStartedRequest {
intent: kache_core::BuildIntent {
crate_names: vec!["serde".into(), "tokio".into()],
namespace: None,
cargo_lock_deps: vec![],
},
client_epoch: 0,
session_id: String::new(),
};
let resp = daemon.handle_build_started(&req).await;
assert!(
resp.ok,
"fallback with nothing to prefetch should be ok: {resp:?}"
);
}
#[tokio::test]
async fn test_batch_remote_check_remote_path_with_injected_mock() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
let client = test_remote_backend();
let daemon = Arc::new(Daemon::new(config));
assert!(
daemon.remote_backend.set(client).is_ok(),
"inject mock backend"
);
let resp = daemon
.handle_batch_remote_check(&BatchRemoteCheckRequest {
checks: vec![
RemoteCheckRequest {
key: "aaaa1111bbbb2222".into(),
entry_dir: dir.path().join("a").to_string_lossy().into_owned(),
crate_name: "serde".into(),
},
RemoteCheckRequest {
key: "cccc3333dddd4444".into(),
entry_dir: dir.path().join("b").to_string_lossy().into_owned(),
crate_name: "tokio".into(),
},
],
})
.await;
assert!(resp.ok);
let results = resp.batch_results.expect("batch results present");
assert_eq!(results.len(), 2);
assert!(results.iter().all(|r| r.found == Some(false)));
}
#[tokio::test]
async fn test_remote_check_hit_then_download_failure() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
let client = test_remote_backend();
put_test_object(
&client,
&test_manifest_object_key("hit01hit02hit03", "serde"),
b"{}",
)
.await;
put_test_object(
&client,
&test_pack_object_key("hit01hit02hit03", "serde"),
b"not a pack",
)
.await;
let daemon = Arc::new(Daemon::new(config));
assert!(
daemon.remote_backend.set(client).is_ok(),
"inject mock backend"
);
let resp = daemon
.handle_remote_check(&RemoteCheckRequest {
key: "hit01hit02hit03".into(),
entry_dir: dir.path().join("entry").to_string_lossy().into_owned(),
crate_name: "serde".into(),
})
.await;
assert!(
!resp.ok,
"download failure should surface as an error: {resp:?}"
);
assert!(resp.error.is_some());
}
#[test]
fn test_prefetch_concurrency_cap_reserves_interactive_permits() {
assert_eq!(prefetch_concurrency_cap(16), 12); assert_eq!(prefetch_concurrency_cap(8), 6); assert_eq!(prefetch_concurrency_cap(4), 3); assert_eq!(prefetch_concurrency_cap(2), 1); assert_eq!(prefetch_concurrency_cap(1), 1); assert_eq!(prefetch_concurrency_cap(0), 1); assert_eq!(prefetch_concurrency_cap(64), 60); for n in 2..=64u32 {
assert!(
prefetch_concurrency_cap(n) < n as usize,
"pool {n}: prefetch must never be able to hold every permit"
);
}
}
#[tokio::test]
async fn test_remote_check_known_positive_get_404_is_clean_miss() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
let client = test_remote_backend();
let daemon = Arc::new(Daemon::new(config));
assert!(
daemon.remote_backend.set(client).is_ok(),
"inject mock backend"
);
let key = "gone01gone02gone03";
let mut keys = HashMap::new();
keys.insert(key.to_string(), "serde".to_string());
daemon.key_cache.populate(keys).await;
let resp = daemon
.handle_remote_check(&RemoteCheckRequest {
key: key.into(),
entry_dir: dir.path().join("entry").to_string_lossy().into_owned(),
crate_name: "serde".into(),
})
.await;
assert!(
resp.ok,
"GET 404 must be a clean miss, not an error: {resp:?}"
);
assert_eq!(resp.found, Some(false));
assert_eq!(daemon.key_cache.check(key).await, Some(false));
assert_eq!(
daemon
.transfer_counters
.downloads_failed
.load(Ordering::Relaxed),
0
);
}
fn build_entry_pack(key: &str, crate_name: &str) -> Vec<u8> {
let tmp = tempfile::tempdir().unwrap();
let cfg = test_config(tmp.path());
let store = Store::open(&cfg).unwrap();
let src = tmp.path().join("src");
std::fs::create_dir_all(&src).unwrap();
let artifact = src.join("libfoo.rlib");
std::fs::write(&artifact, b"real artifact bytes").unwrap();
store
.put(
key,
crate_name,
&["lib".to_string()],
&[],
"x86_64-unknown-linux-gnu",
"debug",
&[(artifact, "libfoo.rlib".to_string())],
"",
"",
)
.unwrap();
let entry_dir = store.entry_dir(key);
let meta: crate::store::EntryMeta =
serde_json::from_slice(&std::fs::read(entry_dir.join("meta.json")).unwrap()).unwrap();
crate::remote_layout::create_entry_pack_zstd(&entry_dir, &store.blobs_dir(), &meta, 3)
.unwrap()
}
#[tokio::test]
async fn test_remote_check_hit_downloads_and_imports() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
config.prefetch_enabled = false;
let key = "hitok01hitok02hi";
let pack = build_entry_pack(key, "serde");
let entry_dir = config.store_dir().join(key);
let client = test_remote_backend();
put_test_object(&client, &test_manifest_object_key(key, "serde"), b"{}").await;
put_test_object(&client, &test_pack_object_key(key, "serde"), &pack).await;
let daemon = Arc::new(Daemon::new(config.clone()));
assert!(
daemon.remote_backend.set(client).is_ok(),
"inject mock backend"
);
let resp = daemon
.handle_remote_check(&RemoteCheckRequest {
key: key.to_string(),
entry_dir: entry_dir.to_string_lossy().into_owned(),
crate_name: "serde".to_string(),
})
.await;
assert!(resp.ok, "hit+download should succeed: {resp:?}");
assert_eq!(resp.found, Some(true));
assert!(
config.store_dir().join(key).join("meta.json").exists(),
"entry should be imported into the local store"
);
}
#[tokio::test]
async fn test_handle_prefetch_disabled_ignores_explicit_keys() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
config.prefetch_enabled = false;
let key = "0123456789abcdef".repeat(4);
let pack = build_entry_pack(&key, "serde");
let client = test_remote_backend();
put_test_object(&client, &test_pack_object_key(&key, "serde"), &pack).await;
let daemon = Arc::new(Daemon::new(config.clone()));
assert!(daemon.remote_backend.set(client).is_ok());
let resp = daemon
.handle_prefetch(&PrefetchRequest {
keys: vec![(key.clone(), "serde".to_string())],
warm_all: false,
})
.await;
assert!(resp.ok);
assert!(!config.store_dir().join(&key).join("meta.json").exists());
assert_eq!(
daemon
.prefetch_stats
.downloads_completed
.load(Ordering::Relaxed),
0
);
}
#[tokio::test]
async fn test_handle_prefetch_explicit_key_downloads_in_background() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
let key = "abcdef0123456789".repeat(4);
let key = key.as_str();
let pack = build_entry_pack(key, "serde");
let client = test_remote_backend();
put_test_object(&client, &test_pack_object_key(key, "serde"), &pack).await;
let daemon = Arc::new(Daemon::new(config.clone()));
assert!(
daemon.remote_backend.set(client).is_ok(),
"inject mock backend"
);
let resp = daemon
.handle_prefetch(&PrefetchRequest {
keys: vec![(key.to_string(), "serde".to_string())],
warm_all: false,
})
.await;
assert!(resp.ok, "prefetch dispatch should be ok: {resp:?}");
let entry_meta = config.store_dir().join(key).join("meta.json");
let mut imported = false;
for _ in 0..100 {
if entry_meta.exists() {
imported = true;
break;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
assert!(
imported,
"background prefetch coordinator should download + import the entry"
);
}
#[tokio::test]
async fn test_handle_prefetch_records_a_failed_download() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
let key = "abcdef0123456789".repeat(4);
let key = key.as_str();
let client = test_remote_backend();
put_test_object(
&client,
&test_pack_object_key(key, "serde"),
b"not a valid pack",
)
.await;
let daemon = Arc::new(Daemon::new(config.clone()));
assert!(
daemon.remote_backend.set(client).is_ok(),
"inject mock backend"
);
let resp = daemon
.handle_prefetch(&PrefetchRequest {
keys: vec![(key.to_string(), "serde".to_string())],
warm_all: false,
})
.await;
assert!(
resp.ok,
"prefetch dispatch is ok even if downloads fail: {resp:?}"
);
let mut failed = false;
for _ in 0..100 {
if daemon
.transfer_counters
.downloads_failed
.load(Ordering::Relaxed)
>= 1
{
failed = true;
break;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
assert!(failed, "a garbage pack must record a failed download");
assert!(!config.store_dir().join(key).join("meta.json").exists());
}
#[tokio::test]
async fn test_populate_key_cache_lists_and_populates() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
let client = test_remote_backend();
put_test_object(&client, "prefix/v3/manifests/serde/key1aaaa.json", b"{}").await;
put_test_object(&client, "prefix/v3/manifests/tokio/key2bbbb.json", b"{}").await;
let daemon = Daemon::new(config);
assert!(
daemon.remote_backend.set(client).is_ok(),
"inject mock backend"
);
let count = populate_key_cache(&daemon)
.await
.expect("populate_key_cache should succeed");
assert_eq!(count, 2);
assert_eq!(daemon.key_cache.check("key1aaaa").await, Some(true));
}
#[tokio::test]
async fn test_monolithic_manifest_prefetch_downloads_and_filters() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
let remote = test_remote_config();
config.remote = Some(remote.clone());
let manifest = crate::remote::BuildManifest {
version: 3,
created: "2025-01-01T00:00:00Z".to_string(),
manifest_key: crate::cli::default_manifest_key(),
entries: vec![crate::remote::ManifestEntry {
cache_key: "cheapkey".to_string(),
crate_name: "cheap".to_string(),
compile_time_ms: 10, artifact_size: 100,
}],
};
let body = serde_json::to_vec(&manifest).unwrap();
let client = test_remote_backend();
put_test_object(&client, &test_build_manifest_object_key(), &body).await;
let daemon = Arc::new(Daemon::new(config));
monolithic_manifest_prefetch(&daemon, client.as_ref(), &remote).await;
}
#[tokio::test]
async fn test_monolithic_manifest_prefetch_skips_when_no_manifest() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
let remote = test_remote_config();
config.remote = Some(remote.clone());
let client = test_remote_backend();
let daemon = Arc::new(Daemon::new(config));
monolithic_manifest_prefetch(&daemon, client.as_ref(), &remote).await; }
#[tokio::test]
async fn test_monolithic_manifest_prefetch_dispatches_expensive_entries() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
let remote = test_remote_config();
config.remote = Some(remote.clone());
let key = "abcdef0123456789".repeat(4);
let manifest = crate::remote::BuildManifest {
version: 3,
created: "2025-01-01T00:00:00Z".to_string(),
manifest_key: crate::cli::default_manifest_key(),
entries: vec![crate::remote::ManifestEntry {
cache_key: key.clone(),
crate_name: "expensive".to_string(),
compile_time_ms: 5000, artifact_size: 100,
}],
};
let body = serde_json::to_vec(&manifest).unwrap();
let client = test_remote_backend();
put_test_object(&client, &test_build_manifest_object_key(), &body).await;
put_test_object(&client, &test_pack_object_key(&key, "expensive"), b"nope").await;
let daemon = Arc::new(Daemon::new(config));
assert!(
daemon.remote_backend.set(client.clone()).is_ok(),
"inject mock backend"
);
monolithic_manifest_prefetch(&daemon, client.as_ref(), &remote).await; }
#[tokio::test]
async fn test_shard_prefetch_all_shards_missing_returns_zero() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
let lock = dir.path().join("Cargo.lock");
std::fs::write(
&lock,
"version = 3\n\n[[package]]\nname = \"serde\"\nversion = \"1.0.0\"\n\n\
[[package]]\nname = \"tokio\"\nversion = \"1.0.0\"\n",
)
.unwrap();
let client = test_remote_backend();
let daemon = Arc::new(Daemon::new(config));
let count = shard_prefetch(&daemon, &client, "prefix", "ns", &lock)
.await
.expect("shard prefetch should succeed");
assert_eq!(count, 0, "no shards matched -> nothing queued");
}
#[test]
fn test_batch_remote_check_request_serde() {
let req = Request::BatchRemoteCheck(BatchRemoteCheckRequest {
checks: vec![
RemoteCheckRequest {
key: "key1".into(),
entry_dir: "/tmp/key1".into(),
crate_name: String::new(),
},
RemoteCheckRequest {
key: "key2".into(),
entry_dir: "/tmp/key2".into(),
crate_name: String::new(),
},
],
});
let json = serde_json::to_string(&req).unwrap();
let parsed: Request = serde_json::from_str(&json).unwrap();
assert_eq!(req, parsed);
assert!(json.contains("\"batch_remote_check\""));
assert!(json.contains("\"key1\""));
assert!(json.contains("\"key2\""));
}
#[test]
fn test_prefetch_request_serde() {
let req = Request::Prefetch(PrefetchRequest {
keys: vec![
("key_a".into(), "serde".into()),
("key_b".into(), "tokio".into()),
],
warm_all: false,
});
let json = serde_json::to_string(&req).unwrap();
let parsed: Request = serde_json::from_str(&json).unwrap();
assert_eq!(req, parsed);
assert!(json.contains("\"prefetch\""));
assert!(json.contains("\"key_a\""));
}
#[test]
fn test_hash_files_request_serde() {
let req = Request::HashFiles(HashFilesRequest {
files: vec![HashFileRequest {
path: "/tmp/libfoo.rlib".into(),
size: 123,
mtime_ns: 456,
ctime_ns: 789,
inode: 1011,
}],
});
let json = serde_json::to_string(&req).unwrap();
let parsed: Request = serde_json::from_str(&json).unwrap();
assert_eq!(req, parsed);
assert!(json.contains("\"hash_files\""));
}
#[test]
fn test_prefetch_request_empty_keys_serde() {
let req = Request::Prefetch(PrefetchRequest {
keys: vec![],
warm_all: false,
});
let json = serde_json::to_string(&req).unwrap();
let parsed: Request = serde_json::from_str(&json).unwrap();
assert_eq!(req, parsed);
}
#[test]
fn test_prefetch_request_from_plan() {
let valid_key = "a".repeat(64);
let plan = PrefetchPlan {
plan_id: Some("plan-1".into()),
planner: Some("fallback".into()),
disposition: PrefetchDisposition::Execute,
candidates: vec![
kache_core::PrefetchCandidate::new(valid_key.clone(), "serde".into()),
kache_core::PrefetchCandidate::new("../../../etc/passwd".into(), "serde".into()),
kache_core::PrefetchCandidate::new(valid_key.clone(), "../evil".into()),
],
};
let req = PrefetchRequest::from_plan(plan);
assert_eq!(req.keys, vec![(valid_key, "serde".into())]);
}
#[test]
fn test_batch_response_serde() {
let batch = BatchResponse {
ok: true,
results: vec![Response::found(true), Response::found(false)],
error: None,
};
let json = serde_json::to_string(&batch).unwrap();
let parsed: BatchResponse = serde_json::from_str(&json).unwrap();
assert_eq!(batch, parsed);
assert_eq!(parsed.results.len(), 2);
assert_eq!(parsed.results[0].found, Some(true));
assert_eq!(parsed.results[1].found, Some(false));
}
#[tokio::test]
async fn test_wait_for_warming_already_signaled() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Daemon::new(config);
daemon.signal_warming_complete();
let start = std::time::Instant::now();
assert!(daemon.wait_for_warming(Duration::from_millis(100)).await);
assert!(start.elapsed() < Duration::from_millis(500));
}
#[tokio::test]
async fn test_prefetch_disabled_remote_releases_warming_barrier() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(crate::config::RemoteConfig::test_s3("test", "artifacts"));
config.prefetch_enabled = false;
let daemon = Arc::new(Daemon::new(config));
assert!(
start_manifest_warming(&daemon).is_none(),
"prefetch-disabled startup must not spawn a warming task"
);
let start = std::time::Instant::now();
assert!(daemon.wait_for_warming(Duration::from_millis(100)).await);
assert!(
start.elapsed() < Duration::from_millis(500),
"prefetch-disabled exact checks must not pay the warming grace"
);
}
#[tokio::test]
async fn test_wait_for_warming_blocks_then_signals() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Arc::new(Daemon::new(config));
let d = daemon.clone();
tokio::spawn(async move {
tokio::time::sleep(Duration::from_millis(50)).await;
d.signal_warming_complete();
});
let start = std::time::Instant::now();
assert!(daemon.wait_for_warming(Duration::from_secs(5)).await);
let elapsed = start.elapsed();
assert!(elapsed >= Duration::from_millis(30));
assert!(elapsed < Duration::from_secs(1));
}
#[tokio::test]
async fn test_wait_for_warming_timeout() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Daemon::new(config);
let start = std::time::Instant::now();
assert!(!daemon.wait_for_warming(Duration::from_millis(100)).await);
let elapsed = start.elapsed();
assert!(elapsed >= Duration::from_millis(90));
assert!(elapsed < Duration::from_millis(500));
}
#[tokio::test]
async fn test_wait_for_warming_multiple_waiters() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Arc::new(Daemon::new(config));
let d1 = daemon.clone();
let d2 = daemon.clone();
let h1 = tokio::spawn(async move { d1.wait_for_warming(Duration::from_secs(5)).await });
let h2 = tokio::spawn(async move { d2.wait_for_warming(Duration::from_secs(5)).await });
tokio::time::sleep(Duration::from_millis(50)).await;
daemon.signal_warming_complete();
let (r1, r2) = tokio::join!(h1, h2);
assert!(r1.unwrap());
assert!(r2.unwrap());
}
#[test]
fn test_remote_health_degrades_after_threshold_and_recovers_on_success() {
let health = RemoteHealth::new();
health.note_head_probe_failure("boom-1");
health.note_head_probe_failure("boom-2");
assert!(!health.head_probe_is_degraded());
health.note_head_probe_failure("boom-3");
assert!(health.head_probe_is_degraded());
health.note_head_probe_suppressed();
health.note_head_probe_success();
assert!(!health.head_probe_is_degraded());
assert_eq!(health.head_probe_failures.load(Ordering::Acquire), 0);
assert_eq!(health.suppressed_head_probes.load(Ordering::Acquire), 0);
}
#[tokio::test]
async fn test_handle_remote_check_skips_head_when_probe_circuit_is_open() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(crate::config::RemoteConfig::test_s3("test", "artifacts"));
let daemon = Daemon::new(config);
daemon.signal_warming_complete();
daemon.remote_health.note_head_probe_failure("boom-1");
daemon.remote_health.note_head_probe_failure("boom-2");
daemon.remote_health.note_head_probe_failure("boom-3");
let req = RemoteCheckRequest {
key: "k".into(),
entry_dir: "/tmp/test".into(),
crate_name: "crate".into(),
};
let resp = daemon.handle_remote_check(&req).await;
assert!(resp.ok);
assert_eq!(resp.found, Some(false));
assert_eq!(
daemon
.remote_health
.suppressed_head_probes
.load(Ordering::Acquire),
1
);
}
#[tokio::test]
async fn test_handle_remote_check_authoritative_key_cache_skips_s3() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(crate::config::RemoteConfig::test_s3("test", "artifacts"));
let daemon = Daemon::new(config);
daemon.signal_warming_complete();
let present = "a".repeat(64);
let mut keys = HashMap::new();
keys.insert(present.clone(), "othercrate".to_string());
daemon.key_cache.populate(keys).await;
let missing = "b".repeat(64);
let req = RemoteCheckRequest {
key: missing,
entry_dir: "/tmp/test".into(),
crate_name: "crate".into(),
};
let resp = daemon.handle_remote_check(&req).await;
assert!(resp.ok);
assert_eq!(
resp.found,
Some(false),
"fresh key cache should authoritatively report the missing key as not found"
);
assert_eq!(
daemon
.remote_health
.suppressed_head_probes
.load(Ordering::Acquire),
0
);
}
#[tokio::test]
async fn test_handle_prefetch_no_remote() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path()); let daemon = Arc::new(Daemon::new(config));
let req = PrefetchRequest {
keys: vec![("k".into(), "mycrate".into())],
warm_all: false,
};
let resp = daemon.handle_prefetch(&req).await;
assert!(!resp.ok);
assert!(
resp.error
.as_deref()
.unwrap()
.contains("no remote configured")
);
}
#[tokio::test]
async fn test_prefetch_key_budget_truncates_and_reports() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
config.prefetch_max_keys = 1;
config.s3_concurrency = 2;
let keys = [
"1111111111111111".repeat(4),
"2222222222222222".repeat(4),
"3333333333333333".repeat(4),
];
let client = test_remote_backend();
for key in &keys {
put_test_object(&client, &test_manifest_object_key(key, "serde"), b"{}").await;
put_test_object(
&client,
&test_pack_object_key(key, "serde"),
&build_entry_pack(key, "serde"),
)
.await;
}
let daemon = Arc::new(Daemon::new(config.clone()));
assert!(
daemon.remote_backend.set(client).is_ok(),
"inject mock backend"
);
let _gate = daemon
.prefetch_gate
.clone()
.acquire_owned()
.await
.expect("gate permit");
let resp = daemon
.handle_prefetch(&PrefetchRequest {
keys: keys
.iter()
.map(|k| (k.clone(), "serde".to_string()))
.collect(),
warm_all: false,
})
.await;
assert!(resp.ok, "prefetch dispatch should be ok: {resp:?}");
assert_eq!(
daemon
.prefetch_stats
.keys_over_budget
.load(Ordering::Relaxed),
2,
"two of three candidates should be reported as dropped over budget"
);
}
#[test]
fn test_prefetch_key_budget_overflow() {
assert_eq!(prefetch_key_budget_overflow(10, 4), 6);
assert_eq!(prefetch_key_budget_overflow(4, 4), 0, "exactly at budget");
assert_eq!(prefetch_key_budget_overflow(3, 4), 0, "under budget");
assert_eq!(prefetch_key_budget_overflow(0, 4), 0, "empty plan");
assert_eq!(
prefetch_key_budget_overflow(10_000, 0),
0,
"0 disables the key budget"
);
}
#[test]
fn test_prefetch_byte_budget_exhausted() {
assert!(!prefetch_byte_budget_exhausted(1024, 0));
assert!(!prefetch_byte_budget_exhausted(1024, 1023));
assert!(
prefetch_byte_budget_exhausted(1024, 1024),
"a budget exactly met stops the next download"
);
assert!(prefetch_byte_budget_exhausted(1024, 4096), "overshot");
assert!(
!prefetch_byte_budget_exhausted(0, u64::MAX),
"0 disables the byte budget"
);
}
#[tokio::test]
async fn test_prefetch_deadline_stops_the_plan() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
config.prefetch_deadline_secs = 0;
let key = "4444444444444444".repeat(4);
let client = test_remote_backend();
put_test_object(&client, &test_manifest_object_key(&key, "serde"), b"{}").await;
put_test_object(
&client,
&test_pack_object_key(&key, "serde"),
&build_entry_pack(&key, "serde"),
)
.await;
let daemon = Arc::new(Daemon::new(config.clone()));
assert!(
daemon.remote_backend.set(client).is_ok(),
"inject mock backend"
);
let resp = daemon
.handle_prefetch(&PrefetchRequest {
keys: vec![(key.clone(), "serde".to_string())],
warm_all: false,
})
.await;
assert!(resp.ok, "prefetch dispatch should be ok: {resp:?}");
let entry_meta = config.store_dir().join(&key).join("meta.json");
let mut imported = false;
for _ in 0..100 {
if entry_meta.exists() {
imported = true;
break;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
assert!(
imported,
"prefetch_deadline_secs = 0 disables the deadline rather than dropping the plan"
);
}
#[tokio::test]
async fn test_empty_prefetch_request_does_not_warm_the_bucket() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
let key = "cccccccccccccccc".repeat(4);
let client = test_remote_backend();
put_test_object(&client, &test_manifest_object_key(&key, "serde"), b"{}").await;
put_test_object(
&client,
&test_pack_object_key(&key, "serde"),
&build_entry_pack(&key, "serde"),
)
.await;
let daemon = Arc::new(Daemon::new(config.clone()));
assert!(
daemon.remote_backend.set(client).is_ok(),
"inject mock backend"
);
let resp = daemon
.handle_prefetch(&PrefetchRequest {
keys: Vec::new(),
warm_all: false,
})
.await;
assert!(
resp.ok,
"an empty request is a no-op, not an error: {resp:?}"
);
tokio::time::sleep(Duration::from_millis(200)).await;
assert!(
daemon.downloading.read().await.is_empty(),
"an empty prefetch request must not claim any key"
);
assert!(
!config.store_dir().join(&key).join("meta.json").exists(),
"an empty prefetch request must not download anything"
);
assert_eq!(
daemon
.transfer_counters
.downloads_completed
.load(Ordering::Relaxed),
0
);
}
#[tokio::test]
async fn test_warm_all_prefetch_request_downloads_missing_keys() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
let key = "dddddddddddddddd".repeat(4);
let client = test_remote_backend();
put_test_object(&client, &test_manifest_object_key(&key, "serde"), b"{}").await;
put_test_object(
&client,
&test_pack_object_key(&key, "serde"),
&build_entry_pack(&key, "serde"),
)
.await;
let daemon = Arc::new(Daemon::new(config.clone()));
assert!(
daemon.remote_backend.set(client).is_ok(),
"inject mock backend"
);
let resp = daemon
.handle_prefetch(&PrefetchRequest {
keys: Vec::new(),
warm_all: true,
})
.await;
assert!(resp.ok, "warm_all dispatch should be ok: {resp:?}");
let entry_meta = config.store_dir().join(&key).join("meta.json");
let mut imported = false;
for _ in 0..100 {
if entry_meta.exists() {
imported = true;
break;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
assert!(
imported,
"warm_all should discover the key by listing and import it"
);
}
#[tokio::test]
async fn test_demand_does_not_wait_behind_unstarted_prefetch_candidates() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
config.s3_concurrency = 2;
let stalled_key = "aaaaaaaaaaaaaaaa".repeat(4);
let demanded_key = "bbbbbbbbbbbbbbbb".repeat(4);
let client = test_remote_backend();
for key in [&stalled_key, &demanded_key] {
put_test_object(&client, &test_manifest_object_key(key, "serde"), b"{}").await;
put_test_object(
&client,
&test_pack_object_key(key, "serde"),
&build_entry_pack(key, "serde"),
)
.await;
}
let daemon = Arc::new(Daemon::new(config.clone()));
assert!(
daemon.remote_backend.set(client).is_ok(),
"inject mock backend"
);
daemon.signal_warming_complete();
let _gate = daemon
.prefetch_gate
.clone()
.acquire_owned()
.await
.expect("gate permit");
let resp = daemon
.handle_prefetch(&PrefetchRequest {
keys: vec![
(stalled_key.clone(), "serde".to_string()),
(demanded_key.clone(), "serde".to_string()),
],
warm_all: false,
})
.await;
assert!(resp.ok, "prefetch dispatch should be ok: {resp:?}");
tokio::time::sleep(Duration::from_millis(100)).await;
assert!(
!daemon.downloading.read().await.contains_key(&demanded_key),
"an unstarted prefetch candidate must not be claimed in `downloading`"
);
let resp = tokio::time::timeout(
Duration::from_secs(5),
daemon.handle_remote_check(&RemoteCheckRequest {
key: demanded_key.clone(),
entry_dir: config
.store_dir()
.join(&demanded_key)
.to_string_lossy()
.into_owned(),
crate_name: "serde".into(),
}),
)
.await
.expect("demand must not block behind an unstarted prefetch candidate");
assert!(resp.ok, "demand download should succeed: {resp:?}");
assert_eq!(
resp.found,
Some(true),
"the demanded entry should have been downloaded"
);
}
#[tokio::test]
async fn test_handle_upload_with_queue_returns_immediately() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(crate::config::RemoteConfig::test_s3("test", "artifacts"));
let (tx, _rx) = tokio::sync::mpsc::unbounded_channel::<UploadJob>();
let daemon = Daemon::new(config);
daemon.set_upload_tx(tx);
let job = UploadJob {
key: "test_key".into(),
entry_dir: "/tmp/test".into(),
crate_name: String::new(),
client_epoch: 0,
};
let resp = daemon.handle_upload(&job).await;
assert!(resp.ok);
assert!(resp.error.is_none());
}
#[tokio::test]
async fn test_handle_upload_queue_closed() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(crate::config::RemoteConfig::test_s3("test", "artifacts"));
let (tx, rx) = tokio::sync::mpsc::unbounded_channel::<UploadJob>();
let daemon = Daemon::new(config);
daemon.set_upload_tx(tx);
drop(rx);
let job = UploadJob {
key: "k1".into(),
entry_dir: "/tmp/test".into(),
crate_name: String::new(),
client_epoch: 0,
};
let resp = daemon.handle_upload(&job).await;
assert!(!resp.ok);
assert!(resp.error.as_deref().unwrap().contains("queue closed"));
}
#[tokio::test]
async fn test_handle_upload_dedup() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(crate::config::RemoteConfig::test_s3("test", "artifacts"));
let (tx, _rx) = tokio::sync::mpsc::unbounded_channel::<UploadJob>();
let daemon = Daemon::new(config);
daemon.set_upload_tx(tx);
let job = UploadJob {
key: "same-key".into(),
entry_dir: "/tmp/test".into(),
crate_name: String::new(),
client_epoch: 0,
};
let resp1 = daemon.handle_upload(&job).await;
assert!(resp1.ok);
let resp2 = daemon.handle_upload(&job).await;
assert!(resp2.ok);
}
#[tokio::test]
async fn test_close_upload_queue_closes_buffer_with_daemon_clones_alive() {
let dir = tempfile::tempdir().unwrap();
let daemon = Arc::new(Daemon::new(test_config(dir.path())));
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<UploadJob>();
daemon.set_upload_tx(tx);
let worker_daemon = daemon.clone();
daemon.close_upload_queue();
let recv = tokio::time::timeout(Duration::from_millis(100), rx.recv())
.await
.expect("upload buffer should close promptly after close_upload_queue");
assert!(
recv.is_none(),
"upload buffer must close even while daemon clones remain alive"
);
drop(worker_daemon);
}
#[tokio::test]
async fn test_handle_upload_after_queue_close_rejects_without_direct_upload() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(crate::config::RemoteConfig::test_s3("test", "artifacts"));
let (tx, _rx) = tokio::sync::mpsc::unbounded_channel::<UploadJob>();
let daemon = Daemon::new(config);
daemon.set_upload_tx(tx);
daemon.close_upload_queue();
let job = UploadJob {
key: "late-key".into(),
entry_dir: "/tmp/test".into(),
crate_name: String::new(),
client_epoch: 0,
};
let resp = daemon.handle_upload(&job).await;
assert!(!resp.ok);
assert!(resp.error.as_deref().unwrap().contains("queue closed"));
}
#[test]
fn test_semaphore_created_with_config() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.s3_concurrency = 4;
let daemon = Daemon::new(config);
assert_eq!(daemon.s3_semaphore.available_permits(), 4);
}
#[test]
fn test_semaphore_min_one_permit() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.s3_concurrency = 0;
let daemon = Daemon::new(config);
assert_eq!(daemon.s3_semaphore.available_permits(), 1);
}
#[tokio::test]
async fn test_socket_prefetch_no_remote_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path()); let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let daemon = Arc::new(Daemon::new(config));
let resp = one_shot_request(
&daemon,
&socket_path,
&Request::Prefetch(PrefetchRequest {
keys: vec![("key1".into(), "mycrate".into())],
warm_all: false,
}),
)
.await;
assert!(!resp.ok);
assert!(
resp.error
.as_deref()
.unwrap()
.contains("no remote configured")
);
}
#[tokio::test]
async fn test_key_cache_age_none_before_populate() {
let cache = S3KeyCache::new();
assert!(cache.age().await.is_none());
}
#[tokio::test]
async fn test_key_cache_age_some_after_populate() {
let cache = S3KeyCache::new();
cache.populate(HashMap::new()).await;
let age = cache.age().await;
assert!(age.is_some());
assert!(age.unwrap() < Duration::from_secs(1));
}
#[test]
fn test_build_started_request_serde() {
let req = Request::BuildStarted(BuildStartedRequest {
intent: kache_core::BuildIntent {
crate_names: vec!["serde".into(), "tokio".into(), "anyhow".into()],
namespace: Some("x86_64/hash/release".into()),
cargo_lock_deps: vec![("serde".into(), "1.0.0".into())],
},
client_epoch: 0,
session_id: String::new(),
});
let json = serde_json::to_string(&req).unwrap();
let parsed: Request = serde_json::from_str(&json).unwrap();
assert_eq!(req, parsed);
assert!(json.contains("\"build_started\""));
assert!(json.contains("\"serde\""));
assert!(json.contains("\"tokio\""));
assert!(json.contains("x86_64/hash/release"));
}
#[test]
fn test_build_started_request_empty_serde() {
let req = Request::BuildStarted(BuildStartedRequest {
intent: kache_core::BuildIntent::default(),
client_epoch: 0,
session_id: String::new(),
});
let json = serde_json::to_string(&req).unwrap();
let parsed: Request = serde_json::from_str(&json).unwrap();
assert_eq!(req, parsed);
}
#[tokio::test]
async fn test_send_build_started_client_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let listener = bind_listener(&socket_path);
let daemon = Arc::new(Daemon::new(config.clone()));
let server = tokio::spawn(async move {
let stream = listener.accept().await.expect("accept");
let _ =
handle_connection(stream, &daemon, &AtomicBool::new(false), &Notify::new()).await;
});
let cfg = config.clone();
tokio::task::spawn_blocking(move || {
send_build_started(
&cfg,
BuildStartedRequest {
intent: kache_core::BuildIntent {
crate_names: vec!["serde".into()],
..Default::default()
},
client_epoch: 0,
session_id: String::new(),
},
)
})
.await
.unwrap();
server.await.unwrap();
}
#[tokio::test]
async fn test_send_upload_job_client_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let listener = bind_listener(&socket_path);
let daemon = Arc::new(Daemon::new(config.clone()));
let (tx, _rx) = tokio::sync::mpsc::unbounded_channel::<UploadJob>();
daemon.set_upload_tx(tx);
let server = tokio::spawn(async move {
let stream = listener.accept().await.expect("accept");
let _ =
handle_connection(stream, &daemon, &AtomicBool::new(false), &Notify::new()).await;
});
let cfg = config.clone();
let result = tokio::task::spawn_blocking(move || {
send_upload_job(&cfg, &"a".repeat(64), Path::new("/tmp/test"), "serde")
})
.await
.unwrap();
server.await.unwrap();
assert!(result.is_ok(), "upload job should send to a live daemon");
}
#[tokio::test]
async fn test_send_prefetch_client_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let socket_path = config.socket_path();
std::fs::create_dir_all(socket_path.parent().unwrap()).unwrap();
let listener = bind_listener(&socket_path);
let daemon = Arc::new(Daemon::new(config.clone()));
let server = tokio::spawn(async move {
let stream = listener.accept().await.expect("accept");
let _ =
handle_connection(stream, &daemon, &AtomicBool::new(false), &Notify::new()).await;
});
let cfg = config.clone();
let result = tokio::task::spawn_blocking(move || {
send_prefetch(&cfg, &[("a".repeat(64), "serde".to_string())])
})
.await
.unwrap();
server.await.unwrap();
assert!(result.is_ok(), "prefetch hint should send to a live daemon");
}
#[tokio::test]
async fn test_handle_build_started_no_remote() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path()); let daemon = Arc::new(Daemon::new(config));
let req = BuildStartedRequest {
intent: kache_core::BuildIntent {
crate_names: vec!["mycrate".into()],
..Default::default()
},
client_epoch: 0,
session_id: String::new(),
};
let resp = daemon.handle_build_started(&req).await;
assert!(!resp.ok);
assert!(
resp.error
.as_deref()
.unwrap()
.contains("no remote configured")
);
}
#[tokio::test]
async fn test_handle_build_started_prefetch_disabled_is_a_no_op() {
let dir = tempfile::tempdir().unwrap();
let mut config = test_config(dir.path());
config.remote = Some(test_remote_config());
config.prefetch_enabled = false;
let daemon = Arc::new(Daemon::new(config));
let resp = daemon
.handle_build_started(&BuildStartedRequest {
intent: kache_core::BuildIntent {
crate_names: vec!["serde".into(), "tokio".into()],
..Default::default()
},
client_epoch: 0,
session_id: "disabled-prefetch".into(),
})
.await;
assert!(resp.ok);
assert!(daemon.active_plan.lock().unwrap().is_none());
assert_eq!(
daemon.prefetch_stats.plans_advisory.load(Ordering::Relaxed),
0
);
assert_eq!(
daemon.prefetch_stats.plans_fallback.load(Ordering::Relaxed),
0
);
}
#[test]
fn test_handle_request_sync_rejects_build_started() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Daemon::new(config);
let req = Request::BuildStarted(BuildStartedRequest {
intent: kache_core::BuildIntent {
crate_names: vec!["c".into()],
..Default::default()
},
client_epoch: 0,
session_id: String::new(),
});
let resp = daemon.handle_request_sync(&req);
assert!(!resp.ok);
assert!(resp.error.as_deref().unwrap().contains("async"));
}
#[tokio::test]
async fn test_downloading_map_starts_empty() {
let dir = tempfile::tempdir().unwrap();
let config = test_config(dir.path());
let daemon = Daemon::new(config);
assert!(daemon.downloading.read().await.is_empty());
}
async fn park_on_claim(map: &RwLock<HashMap<String, Arc<Notify>>>, notify: &Notify, key: &str) {
let notified = notify.notified();
tokio::pin!(notified);
notified.as_mut().enable();
if map.read().await.contains_key(key) {
let _ = tokio::time::timeout(Duration::from_secs(10), notified).await;
}
}
#[tokio::test]
async fn downloading_guard_removes_key_via_runtime_when_lock_contended() {
let notify = Arc::new(Notify::new());
let mut keys = HashMap::new();
keys.insert("cache-key".to_string(), notify.clone());
let map = Arc::new(RwLock::new(keys));
let waiter = tokio::spawn({
let map = map.clone();
let notify = notify.clone();
async move {
park_on_claim(&map, ¬ify, "cache-key").await;
!map.read().await.contains_key("cache-key")
}
});
tokio::time::sleep(Duration::from_millis(20)).await;
let write_guard = map.write().await;
let guard = DownloadingGuard::new(map.clone(), "cache-key".to_string());
drop(guard);
assert!(write_guard.contains_key("cache-key"));
drop(write_guard);
let mut removed = false;
for _ in 0..20 {
if !map.read().await.contains_key("cache-key") {
removed = true;
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
assert!(removed, "contended drop should eventually remove the key");
let key_gone_at_wake = tokio::time::timeout(Duration::from_secs(5), waiter)
.await
.expect("waiter should be notified by the contended drop path")
.unwrap();
assert!(key_gone_at_wake, "wake must happen after the map removal");
}
#[tokio::test]
async fn waiter_wakes_promptly_and_reclaims_when_leader_fails() {
let map: Arc<RwLock<HashMap<String, Arc<Notify>>>> = Arc::new(RwLock::new(HashMap::new()));
assert!(
claim_download(&map, "k").await.is_none(),
"first claim is the leader"
);
let leader_guard = DownloadingGuard::new(map.clone(), "k".to_string());
let notify = claim_download(&map, "k")
.await
.expect("second claim is a waiter");
let waiter = tokio::spawn({
let map = map.clone();
async move {
let start = Instant::now();
park_on_claim(&map, ¬ify, "k").await;
let won = claim_download(&map, "k").await.is_none();
(start.elapsed(), won)
}
});
tokio::time::sleep(Duration::from_millis(50)).await; drop(leader_guard); let (elapsed, won) = waiter.await.unwrap();
assert!(won, "waiter should win the re-claim after leader failure");
assert!(
elapsed < Duration::from_secs(5),
"waiter should wake promptly, waited {elapsed:?}"
);
}
#[tokio::test]
async fn exactly_one_waiter_wins_reclaim_after_leader_failure() {
let map: Arc<RwLock<HashMap<String, Arc<Notify>>>> = Arc::new(RwLock::new(HashMap::new()));
assert!(claim_download(&map, "k").await.is_none());
let leader_guard = DownloadingGuard::new(map.clone(), "k".to_string());
let n1 = claim_download(&map, "k").await.unwrap();
let n2 = claim_download(&map, "k").await.unwrap();
let spawn_waiter = |notify: Arc<Notify>| {
let map = map.clone();
tokio::spawn(async move {
park_on_claim(&map, ¬ify, "k").await;
claim_download(&map, "k").await.is_none()
})
};
let w1 = spawn_waiter(n1);
let w2 = spawn_waiter(n2);
tokio::time::sleep(Duration::from_millis(50)).await; drop(leader_guard);
let (r1, r2) = tokio::join!(w1, w2);
let wins = usize::from(r1.unwrap()) + usize::from(r2.unwrap());
assert_eq!(wins, 1, "exactly one waiter must win the re-claim");
}
#[tokio::test]
async fn waiter_sees_meta_json_at_wake_on_leader_success() {
let dir = tempfile::tempdir().unwrap();
let entry_dir = dir.path().join("entry");
std::fs::create_dir_all(&entry_dir).unwrap();
let meta = entry_dir.join("meta.json");
let map: Arc<RwLock<HashMap<String, Arc<Notify>>>> = Arc::new(RwLock::new(HashMap::new()));
assert!(claim_download(&map, "k").await.is_none());
let leader_guard = DownloadingGuard::new(map.clone(), "k".to_string());
let notify = claim_download(&map, "k").await.unwrap();
let waiter = tokio::spawn({
let map = map.clone();
let meta = meta.clone();
async move {
park_on_claim(&map, ¬ify, "k").await;
meta.exists()
}
});
tokio::time::sleep(Duration::from_millis(50)).await; std::fs::write(&meta, "{}").unwrap(); drop(leader_guard); let found = tokio::time::timeout(Duration::from_secs(5), waiter)
.await
.expect("waiter should wake when the leader's guard drops")
.unwrap();
assert!(found, "waiter must observe meta.json at wake");
assert!(map.read().await.is_empty(), "claim fully released");
}
#[tokio::test]
async fn read_bounded_line_strips_and_handles_eof() {
let data = b"hello\nwith-cr\r\n\nlast"; let mut reader = BufReader::new(&data[..]);
let mut buf = Vec::new();
let r = |res: std::io::Result<Option<String>>| res.unwrap();
assert_eq!(
r(read_bounded_line(&mut reader, &mut buf).await).as_deref(),
Some("hello")
);
assert_eq!(
r(read_bounded_line(&mut reader, &mut buf).await).as_deref(),
Some("with-cr")
);
assert_eq!(
r(read_bounded_line(&mut reader, &mut buf).await).as_deref(),
Some("")
);
assert_eq!(
r(read_bounded_line(&mut reader, &mut buf).await).as_deref(),
Some("last")
);
assert_eq!(r(read_bounded_line(&mut reader, &mut buf).await), None);
}
#[tokio::test]
async fn read_bounded_line_rejects_oversized_frame() {
let big = vec![b'x'; MAX_REQUEST_FRAME_BYTES + 4096];
let mut reader = BufReader::new(&big[..]);
let mut buf = Vec::new();
let err = read_bounded_line(&mut reader, &mut buf).await.unwrap_err();
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
}
}