use anyhow::{Context, Result};
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use std::fs::{self, File, OpenOptions};
use std::io::{BufRead, BufReader, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct BuildEvent {
pub ts: DateTime<Utc>,
pub crate_name: String,
#[serde(default, skip_serializing_if = "String::is_empty")]
pub root: String,
#[serde(default)]
pub version: String,
pub result: EventResult,
pub elapsed_ms: u64,
#[serde(default)]
pub compile_time_ms: u64,
pub size: u64,
#[serde(default)]
pub cache_key: String,
#[serde(default)]
pub schema: u32,
#[serde(default, skip_serializing_if = "String::is_empty")]
pub session_id: String,
#[serde(default)]
pub key_ms: u64,
#[serde(default)]
pub key_hash_hits: u64,
#[serde(default)]
pub key_hash_misses: u64,
#[serde(default)]
pub key_hash_bytes: u64,
#[serde(default)]
pub lookup_ms: u64,
#[serde(default)]
pub restore_ms: u64,
#[serde(default)]
pub store_ms: u64,
#[serde(default)]
pub store_output_blobs: u32,
#[serde(default)]
pub store_duplicate_blobs: u32,
#[serde(default)]
pub store_new_blobs: u32,
#[serde(default)]
pub compiler_runs: u32,
#[serde(default)]
pub preprocessor_runs: u32,
#[serde(default)]
pub probe_runs: u32,
#[serde(default)]
pub reflinked_bytes: u64,
#[serde(default)]
pub hardlinked_bytes: u64,
#[serde(default)]
pub copied_bytes: u64,
#[serde(default)]
pub store_reflinked_bytes: u64,
#[serde(default)]
pub store_hardlinked_bytes: u64,
#[serde(default)]
pub store_copied_bytes: u64,
#[serde(default, skip_serializing_if = "String::is_empty")]
pub passthrough_reason: String,
#[serde(default, skip_serializing_if = "String::is_empty")]
pub store_error: String,
#[serde(default, skip_serializing_if = "is_false")]
pub fallback: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub exit_code: Option<i32>,
#[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
pub key_fields: std::collections::BTreeMap<String, String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub key_diff: Vec<String>,
#[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
pub key_externs: std::collections::BTreeMap<String, String>,
#[serde(default, skip_serializing_if = "is_false")]
pub key_externs_recorded: bool,
}
fn is_false(value: &bool) -> bool {
!*value
}
#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum EventResult {
LocalHit,
PrefetchHit,
RemoteHit,
Dup,
Miss,
Error,
Passthrough,
Skipped,
}
impl std::fmt::Display for EventResult {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
EventResult::LocalHit => write!(f, "local_hit"),
EventResult::PrefetchHit => write!(f, "prefetch_hit"),
EventResult::RemoteHit => write!(f, "remote_hit"),
EventResult::Dup => write!(f, "dup"),
EventResult::Miss => write!(f, "miss"),
EventResult::Error => write!(f, "error"),
EventResult::Passthrough => write!(f, "passthrough"),
EventResult::Skipped => write!(f, "skipped"),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct BuildSummaryEvent {
pub ts: DateTime<Utc>,
pub schema: u32,
#[serde(default)]
pub session_id: String,
#[serde(default)]
pub root: String,
#[serde(default)]
pub plan_source: String,
#[serde(default)]
pub plan_id: String,
#[serde(default)]
pub closure_reason: String,
#[serde(default)]
pub started_at_ms: u64,
#[serde(default)]
pub last_activity_ms: u64,
#[serde(default)]
pub candidate_keys: u64,
#[serde(default)]
pub downloaded_keys: u64,
#[serde(default)]
pub downloaded_bytes: u64,
#[serde(default)]
pub used_keys: u64,
#[serde(default)]
pub used_bytes: u64,
#[serde(default)]
pub demanded_keys: u64,
#[serde(default)]
pub demanded_candidate_keys: u64,
#[serde(default)]
pub cancelled: bool,
#[serde(default)]
pub list_requests: u64,
#[serde(default)]
pub list_duration_ms: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HeartbeatEvent {
pub event: String,
pub ts: DateTime<Utc>,
pub crate_name: String,
#[serde(default, skip_serializing_if = "String::is_empty")]
pub root: String,
pub pid: u32,
pub elapsed_s: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub typical_s: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub eta_s: Option<u64>,
#[serde(default)]
pub schema: u32,
}
pub const HEARTBEAT_EVENT_TAG: &str = "heartbeat";
pub const HEARTBEAT_SCHEMA: u32 = 1;
#[derive(Debug, Clone)]
pub enum EventRecord {
Build(Box<BuildEvent>),
Heartbeat(HeartbeatEvent),
}
fn parse_event_line(line: &str) -> Option<EventRecord> {
if let Ok(hb) = serde_json::from_str::<HeartbeatEvent>(line) {
if hb.event == HEARTBEAT_EVENT_TAG {
return Some(EventRecord::Heartbeat(hb));
}
return None;
}
serde_json::from_str::<BuildEvent>(line)
.ok()
.map(|e| EventRecord::Build(Box::new(e)))
}
fn append_log_line(event_log_path: &Path, line: String) -> Result<()> {
if let Some(parent) = event_log_path.parent() {
fs::create_dir_all(parent)?;
}
let lock = open_log_lock(event_log_path).context("opening event log lock")?;
lock.lock().context("locking event log")?;
let mut file = OpenOptions::new()
.create(true)
.append(true)
.open(event_log_path)
.context("opening event log")?;
let mut bytes = line.into_bytes();
bytes.push(b'\n');
let write_result = file.write_all(&bytes).context("writing event to log");
lock.unlock().context("unlocking event log")?;
write_result
}
pub fn log_event(event_log_path: &Path, event: &BuildEvent) -> Result<()> {
append_log_line(
event_log_path,
serde_json::to_string(event).context("serializing event")?,
)
}
pub fn log_heartbeat(event_log_path: &Path, event: &HeartbeatEvent) -> Result<()> {
append_log_line(
event_log_path,
serde_json::to_string(event).context("serializing heartbeat")?,
)
}
const TYPICAL_WINDOW: usize = 20;
pub fn typical_compile_ms(event_log_path: &Path, crate_name: &str, root: &str) -> Option<u64> {
let events = read_events(event_log_path).ok()?;
let samples: Vec<u64> = events
.iter()
.filter(|e| {
e.crate_name == crate_name
&& e.root == root
&& e.compile_time_ms > 0
&& matches!(e.result, EventResult::Miss | EventResult::Dup)
})
.map(|e| e.compile_time_ms)
.collect();
let window = &samples[samples.len().saturating_sub(TYPICAL_WINDOW)..];
let center = median(window)?;
let n = window.len() as f64;
let mean = window.iter().sum::<u64>() as f64 / n;
let sigma = (window
.iter()
.map(|&s| {
let d = s as f64 - mean;
d * d
})
.sum::<f64>()
/ n)
.sqrt();
let kept: Vec<u64> = window
.iter()
.copied()
.filter(|&s| sigma == 0.0 || (s as f64 - center as f64).abs() <= 3.0 * sigma)
.collect();
median(&kept)
}
fn median(samples: &[u64]) -> Option<u64> {
if samples.is_empty() {
return None;
}
let mut sorted = samples.to_vec();
sorted.sort_unstable();
Some(sorted[(sorted.len() - 1) / 2])
}
pub fn read_events(event_log_path: &Path) -> Result<Vec<BuildEvent>> {
if !event_log_path.exists() {
return Ok(Vec::new());
}
let lock = open_log_lock(event_log_path).context("opening event log lock")?;
lock.lock_shared().context("locking event log for read")?;
let file = File::open(event_log_path).context("opening event log")?;
let reader = BufReader::new(&file);
let mut events = Vec::new();
for line in reader.lines() {
let line = line?;
if line.trim().is_empty() {
continue;
}
match serde_json::from_str::<BuildEvent>(&line) {
Ok(event) => events.push(event),
Err(e) => {
tracing::debug!("skipping invalid event line: {}", e);
}
}
}
lock.unlock().context("unlocking event log")?;
Ok(events)
}
pub fn read_events_since(event_log_path: &Path, since: DateTime<Utc>) -> Result<Vec<BuildEvent>> {
let all = read_events(event_log_path)?;
Ok(all.into_iter().filter(|e| e.ts >= since).collect())
}
#[cfg(windows)]
fn get_file_identity(file: &std::fs::File) -> Option<(u32, u32, u32)> {
use std::os::windows::io::AsRawHandle;
use windows_sys::Win32::Storage::FileSystem::{
BY_HANDLE_FILE_INFORMATION, GetFileInformationByHandle,
};
let handle = file.as_raw_handle();
let mut info: BY_HANDLE_FILE_INFORMATION = unsafe { std::mem::zeroed() };
let ok = unsafe { GetFileInformationByHandle(handle as _, &mut info) };
if ok != 0 {
Some((
info.dwVolumeSerialNumber,
info.nFileIndexHigh,
info.nFileIndexLow,
))
} else {
None
}
}
pub struct EventTailer {
path: PathBuf,
position: u64,
file: Option<File>,
}
impl EventTailer {
pub fn new(path: PathBuf) -> Self {
let file = File::open(&path).ok();
let position = file
.as_ref()
.and_then(|file| file.metadata().ok())
.map(|m| m.len())
.unwrap_or(0);
EventTailer {
path,
position,
file,
}
}
pub fn from_start(path: PathBuf) -> Self {
let file = File::open(&path).ok();
EventTailer {
path,
position: 0,
file,
}
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn poll(&mut self) -> Result<Vec<BuildEvent>> {
Ok(self
.poll_records()?
.into_iter()
.filter_map(|r| match r {
EventRecord::Build(e) => Some(*e),
EventRecord::Heartbeat(_) => None,
})
.collect())
}
pub fn poll_records(&mut self) -> Result<Vec<EventRecord>> {
if self.file.is_none() {
match File::open(&self.path) {
Ok(file) => self.file = Some(file),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
Err(e) => return Err(e.into()),
}
}
let mut rotated = false;
if let Some(file) = &self.file {
#[cfg(unix)]
{
use std::os::unix::fs::MetadataExt;
if let Ok(m1) = file.metadata()
&& let Ok(m2) = std::fs::metadata(&self.path)
&& (m1.dev() != m2.dev() || m1.ino() != m2.ino())
{
rotated = true;
}
}
#[cfg(windows)]
{
if let Some(id1) = get_file_identity(file)
&& let Ok(f2) = std::fs::File::open(&self.path)
&& let Some(id2) = get_file_identity(&f2)
&& id1 != id2
{
rotated = true;
}
}
}
if rotated {
match File::open(&self.path) {
Ok(file) => {
self.file = Some(file);
self.position = 0;
}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
self.file = None;
self.position = 0;
return Ok(Vec::new());
}
Err(e) => return Err(e.into()),
}
}
let file = self.file.as_mut().unwrap();
let file_len = file.metadata()?.len();
if file_len < self.position {
self.position = 0;
}
if file_len <= self.position {
return Ok(Vec::new());
}
file.seek(SeekFrom::Start(self.position))?;
let reader = BufReader::new(file);
let mut records = Vec::new();
let mut bytes_read = 0u64;
for line in reader.lines() {
let line = line?;
bytes_read += line.len() as u64 + 1; if line.trim().is_empty() {
continue;
}
if let Some(record) = parse_event_line(&line) {
records.push(record);
}
}
self.position += bytes_read;
Ok(records)
}
}
fn rotate_log_impl(
log_path: &Path,
max_size: u64,
keep_lines: usize,
log_label: &str,
) -> Result<()> {
if !log_path.exists() {
return Ok(());
}
if let Ok(meta) = fs::metadata(log_path)
&& meta.len() <= max_size
{
return Ok(());
}
let lock = open_log_lock(log_path).context("opening log lock")?;
lock.lock().context("locking log for rotation")?;
let res = (|| -> Result<()> {
let meta = fs::metadata(log_path)?;
if meta.len() <= max_size {
return Ok(());
}
if let Some(parent) = log_path.parent()
&& let Some(file_prefix) = log_path.file_name().and_then(|n| n.to_str())
{
let _ = crate::atomic::cleanup_temp_files(
parent,
file_prefix,
std::time::Duration::from_secs(300),
);
}
let content = fs::read_to_string(log_path)?;
let lines: Vec<&str> = content.lines().collect();
let keep_from = lines.len().saturating_sub(keep_lines);
let mut kept: Vec<&str> = lines[keep_from..].to_vec();
let mut total_bytes: u64 = kept.iter().map(|line| line.len() as u64 + 1).sum();
while total_bytes > max_size && kept.len() > 1 {
let removed = kept.remove(0);
total_bytes -= (removed.len() + 1) as u64;
}
let output = kept.join("\n") + "\n";
crate::atomic::atomic_replace(log_path, output.as_bytes())
.context("writing and replacing log file atomically")?;
tracing::info!(
"rotated {}: kept {} of {} lines",
log_label,
kept.len(),
lines.len()
);
Ok(())
})();
let _ = lock.unlock();
res
}
pub fn rotate_if_needed(event_log_path: &Path, max_size: u64, keep_lines: usize) -> Result<()> {
rotate_log_impl(event_log_path, max_size, keep_lines, "event log")
}
pub fn log_summary(summary_log_path: &Path, event: &BuildSummaryEvent) -> Result<()> {
if let Some(parent) = summary_log_path.parent() {
fs::create_dir_all(parent)?;
}
let lock = open_log_lock(summary_log_path).context("opening summary log lock")?;
lock.lock().context("locking summary log")?;
let mut file = OpenOptions::new()
.create(true)
.append(true)
.open(summary_log_path)
.context("opening summary log")?;
let line = serde_json::to_string(event).context("serializing summary event")?;
let mut bytes = line.into_bytes();
bytes.push(b'\n');
file.write_all(&bytes)
.context("writing summary event to log")?;
lock.unlock().context("unlocking summary log")?;
Ok(())
}
pub fn read_summaries(summary_log_path: &Path) -> Result<Vec<BuildSummaryEvent>> {
if !summary_log_path.exists() {
return Ok(Vec::new());
}
let lock = open_log_lock(summary_log_path).context("opening summary log lock")?;
lock.lock_shared().context("locking summary log for read")?;
let file = File::open(summary_log_path).context("opening summary log")?;
let reader = BufReader::new(&file);
let mut events = Vec::new();
for line in reader.lines() {
let line = line?;
if line.trim().is_empty() {
continue;
}
match serde_json::from_str::<BuildSummaryEvent>(&line) {
Ok(event) => events.push(event),
Err(e) => {
tracing::debug!("skipping invalid summary line: {}", e);
}
}
}
lock.unlock().context("unlocking summary log")?;
Ok(events)
}
use crate::daemon::TransferEvent;
pub fn log_transfer(transfer_log_path: &Path, event: &TransferEvent) -> Result<()> {
if let Some(parent) = transfer_log_path.parent() {
fs::create_dir_all(parent)?;
}
let lock = open_log_lock(transfer_log_path).context("opening transfer log lock")?;
lock.lock().context("locking transfer log")?;
let mut file = OpenOptions::new()
.create(true)
.append(true)
.open(transfer_log_path)
.context("opening transfer log")?;
let line = serde_json::to_string(event).context("serializing transfer event")?;
let mut bytes = line.into_bytes();
bytes.push(b'\n');
file.write_all(&bytes)
.context("writing transfer event to log")?;
lock.unlock().context("unlocking transfer log")?;
Ok(())
}
pub fn read_transfers(transfer_log_path: &Path) -> Result<Vec<TransferEvent>> {
if !transfer_log_path.exists() {
return Ok(Vec::new());
}
let lock = open_log_lock(transfer_log_path).context("opening transfer log lock")?;
lock.lock_shared()
.context("locking transfer log for read")?;
let file = File::open(transfer_log_path).context("opening transfer log")?;
let reader = BufReader::new(&file);
let mut events = Vec::new();
for line in reader.lines() {
let line = line?;
if line.trim().is_empty() {
continue;
}
match serde_json::from_str::<TransferEvent>(&line) {
Ok(event) => events.push(event),
Err(e) => {
tracing::debug!("skipping invalid transfer line: {}", e);
}
}
}
lock.unlock().context("unlocking transfer log")?;
Ok(events)
}
fn open_log_lock(log_path: &Path) -> Result<File> {
OpenOptions::new()
.create(true)
.truncate(false)
.read(true)
.write(true)
.open(sidecar_lock_path(log_path))
.context("opening log lock")
}
fn sidecar_lock_path(log_path: &Path) -> PathBuf {
let mut path = log_path.as_os_str().to_owned();
path.push(".lock");
PathBuf::from(path)
}
pub fn read_transfers_since(transfer_log_path: &Path, since_ts: u64) -> Result<Vec<TransferEvent>> {
let all = read_transfers(transfer_log_path)?;
Ok(all
.into_iter()
.filter(|e| e.timestamp >= since_ts)
.collect())
}
pub fn rotate_transfers_if_needed(
transfer_log_path: &Path,
max_size: u64,
keep_lines: usize,
) -> Result<()> {
rotate_log_impl(transfer_log_path, max_size, keep_lines, "transfer log")
}
#[allow(dead_code)]
pub fn clear_events(event_log_path: &Path) -> Result<()> {
if !event_log_path.exists() {
return Ok(());
}
let lock = open_log_lock(event_log_path).context("opening event log lock")?;
lock.lock().context("locking event log for clearing")?;
let res = fs::write(event_log_path, "");
let _ = lock.unlock();
res.context("clearing event log")
}
pub struct EventStats {
#[allow(dead_code)]
pub total: usize,
pub local_hits: usize,
pub prefetch_hits: usize,
pub remote_hits: usize,
pub dups: usize,
pub misses: usize,
pub errors: usize,
pub total_size: u64,
pub total_elapsed_ms: u64,
pub hit_elapsed_ms: u64,
pub miss_elapsed_ms: u64,
pub hit_compile_time_ms: u64,
pub miss_compile_time_ms: u64,
pub total_key_ms: u64,
pub total_lookup_ms: u64,
pub total_restore_ms: u64,
pub total_store_ms: u64,
pub store_output_blobs: u32,
pub store_duplicate_blobs: u32,
pub store_new_blobs: u32,
pub reflinked_bytes: u64,
pub hardlinked_bytes: u64,
pub copied_bytes: u64,
pub store_reflinked_bytes: u64,
pub store_hardlinked_bytes: u64,
pub store_copied_bytes: u64,
pub store_failures: usize,
}
pub fn compute_stats(events: &[BuildEvent]) -> EventStats {
let mut stats = EventStats {
total: events.len(),
local_hits: 0,
prefetch_hits: 0,
remote_hits: 0,
dups: 0,
misses: 0,
errors: 0,
total_size: 0,
total_elapsed_ms: 0,
hit_elapsed_ms: 0,
miss_elapsed_ms: 0,
hit_compile_time_ms: 0,
miss_compile_time_ms: 0,
total_key_ms: 0,
total_lookup_ms: 0,
total_restore_ms: 0,
total_store_ms: 0,
store_output_blobs: 0,
store_duplicate_blobs: 0,
store_new_blobs: 0,
reflinked_bytes: 0,
hardlinked_bytes: 0,
copied_bytes: 0,
store_reflinked_bytes: 0,
store_hardlinked_bytes: 0,
store_copied_bytes: 0,
store_failures: 0,
};
for event in events {
match event.result {
EventResult::LocalHit => {
stats.local_hits += 1;
stats.hit_elapsed_ms += event.elapsed_ms;
stats.hit_compile_time_ms += event.compile_time_ms;
}
EventResult::PrefetchHit => {
stats.prefetch_hits += 1;
stats.hit_elapsed_ms += event.elapsed_ms;
stats.hit_compile_time_ms += event.compile_time_ms;
}
EventResult::RemoteHit => {
stats.remote_hits += 1;
stats.hit_elapsed_ms += event.elapsed_ms;
stats.hit_compile_time_ms += event.compile_time_ms;
}
EventResult::Dup => {
stats.dups += 1;
stats.miss_elapsed_ms += event.elapsed_ms;
stats.miss_compile_time_ms += if event.compile_time_ms > 0 {
event.compile_time_ms
} else {
event.elapsed_ms
};
}
EventResult::Miss => {
stats.misses += 1;
stats.miss_elapsed_ms += event.elapsed_ms;
stats.miss_compile_time_ms += if event.compile_time_ms > 0 {
event.compile_time_ms
} else {
event.elapsed_ms
};
}
EventResult::Error => stats.errors += 1,
EventResult::Passthrough | EventResult::Skipped => continue,
}
if matches!(event.result, EventResult::Miss | EventResult::Dup)
&& !event.store_error.is_empty()
{
stats.store_failures += 1;
}
stats.total_size += event.size;
stats.total_elapsed_ms += event.elapsed_ms;
stats.total_key_ms += event.key_ms;
stats.total_lookup_ms += event.lookup_ms;
stats.total_restore_ms += event.restore_ms;
stats.total_store_ms += event.store_ms;
stats.store_output_blobs += event.store_output_blobs;
stats.store_duplicate_blobs += event.store_duplicate_blobs;
stats.store_new_blobs += event.store_new_blobs;
stats.reflinked_bytes += event.reflinked_bytes;
stats.hardlinked_bytes += event.hardlinked_bytes;
stats.copied_bytes += event.copied_bytes;
stats.store_reflinked_bytes += event.store_reflinked_bytes;
stats.store_hardlinked_bytes += event.store_hardlinked_bytes;
stats.store_copied_bytes += event.store_copied_bytes;
}
stats
}
#[cfg(test)]
impl BuildEvent {
pub(crate) fn new_for_test(crate_name: &str, result: EventResult) -> BuildEvent {
BuildEvent::test_event(crate_name, result, 0, 0, 0, "")
}
#[cfg(test)]
fn test_event(
crate_name: &str,
result: EventResult,
elapsed_ms: u64,
compile_time_ms: u64,
size: u64,
cache_key: &str,
) -> BuildEvent {
BuildEvent {
ts: Utc::now(),
crate_name: crate_name.to_string(),
root: "/work/tree".to_string(),
version: "0.0.0".to_string(),
session_id: String::new(),
result,
elapsed_ms,
compile_time_ms,
size,
cache_key: cache_key.to_string(),
schema: 8,
key_ms: 0,
key_hash_hits: 0,
key_hash_misses: 0,
key_hash_bytes: 0,
lookup_ms: 0,
restore_ms: 0,
store_ms: 0,
store_output_blobs: 0,
store_duplicate_blobs: 0,
store_new_blobs: 0,
compiler_runs: 0,
preprocessor_runs: 0,
probe_runs: 0,
reflinked_bytes: 0,
hardlinked_bytes: 0,
copied_bytes: 0,
store_reflinked_bytes: 0,
store_hardlinked_bytes: 0,
store_copied_bytes: 0,
passthrough_reason: String::new(),
store_error: String::new(),
fallback: false,
exit_code: None,
key_fields: Default::default(),
key_diff: Vec::new(),
key_externs: Default::default(),
key_externs_recorded: false,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn test_event(
crate_name: &str,
result: EventResult,
elapsed_ms: u64,
compile_time_ms: u64,
size: u64,
cache_key: &str,
) -> BuildEvent {
BuildEvent::test_event(
crate_name,
result,
elapsed_ms,
compile_time_ms,
size,
cache_key,
)
}
#[test]
fn test_log_and_read_events() {
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("events.jsonl");
let event = BuildEvent {
ts: Utc::now(),
crate_name: "serde".to_string(),
root: "/work/tree".to_string(),
version: "1.0.210".to_string(),
session_id: String::new(),
result: EventResult::LocalHit,
elapsed_ms: 2,
compile_time_ms: 250,
size: 3145728,
cache_key: "abc123".to_string(),
schema: 8,
key_ms: 0,
key_hash_hits: 0,
key_hash_misses: 0,
key_hash_bytes: 0,
lookup_ms: 0,
restore_ms: 0,
store_ms: 0,
store_output_blobs: 0,
store_duplicate_blobs: 0,
store_new_blobs: 0,
compiler_runs: 0,
preprocessor_runs: 0,
probe_runs: 0,
reflinked_bytes: 0,
hardlinked_bytes: 0,
copied_bytes: 0,
store_reflinked_bytes: 0,
store_hardlinked_bytes: 0,
store_copied_bytes: 0,
passthrough_reason: String::new(),
store_error: String::new(),
fallback: false,
exit_code: None,
key_fields: Default::default(),
key_diff: Vec::new(),
key_externs: Default::default(),
key_externs_recorded: false,
};
log_event(&log_path, &event).unwrap();
log_event(&log_path, &event).unwrap();
let events = read_events(&log_path).unwrap();
assert_eq!(events.len(), 2);
assert_eq!(events[0].crate_name, "serde");
assert_eq!(events[0].result, EventResult::LocalHit);
}
fn test_heartbeat(crate_name: &str, elapsed_s: u64) -> HeartbeatEvent {
HeartbeatEvent {
event: HEARTBEAT_EVENT_TAG.to_string(),
ts: Utc::now(),
crate_name: crate_name.to_string(),
root: "/work/tree".to_string(),
pid: 4242,
elapsed_s,
typical_s: Some(471),
eta_s: Some(211),
schema: HEARTBEAT_SCHEMA,
}
}
#[test]
fn heartbeat_lines_are_invisible_to_build_event_readers() {
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("events.jsonl");
log_event(
&log_path,
&test_event("gkrust", EventResult::Miss, 471_000, 471_000, 1024, "k1"),
)
.unwrap();
log_heartbeat(&log_path, &test_heartbeat("gkrust", 260)).unwrap();
log_event(
&log_path,
&test_event("serde", EventResult::LocalHit, 2, 250, 64, "k2"),
)
.unwrap();
let builds = read_events(&log_path).unwrap();
assert_eq!(
builds.len(),
2,
"BuildEvent readers must skip the heartbeat line"
);
let hb_line = serde_json::to_string(&test_heartbeat("gkrust", 260)).unwrap();
assert!(
serde_json::from_str::<BuildEvent>(&hb_line).is_err(),
"heartbeat must not deserialize as BuildEvent"
);
let build_line =
serde_json::to_string(&test_event("gkrust", EventResult::Miss, 1, 1, 1, "k")).unwrap();
assert!(
serde_json::from_str::<HeartbeatEvent>(&build_line).is_err(),
"BuildEvent must not deserialize as heartbeat"
);
}
#[test]
fn pre_609_events_deserialize_without_extern_digests() {
let line = r#"{"ts":"2026-07-01T00:00:00Z","crate_name":"serde",
"result":"miss","elapsed_ms":1,"size":2,"cache_key":"k","schema":11,
"key_fields":{"externs":"aaaa"}}"#;
let event: BuildEvent = serde_json::from_str(line).expect("schema 11 must still parse");
assert_eq!(event.schema, 11);
assert!(event.key_externs.is_empty());
assert_eq!(
event.key_fields.get("externs").map(String::as_str),
Some("aaaa")
);
}
#[test]
fn empty_extern_digests_are_not_serialized() {
let event = test_event("serde", EventResult::Miss, 1, 1, 1, "k");
let line = serde_json::to_string(&event).unwrap();
assert!(
!line.contains("key_externs"),
"empty key_externs must not be written: {line}"
);
}
#[test]
fn tailer_poll_records_yields_heartbeats_and_builds() {
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("events.jsonl");
let mut tailer = EventTailer::from_start(log_path.clone());
log_heartbeat(&log_path, &test_heartbeat("gkrust", 30)).unwrap();
log_event(
&log_path,
&test_event("gkrust", EventResult::Miss, 60_000, 60_000, 1024, "k1"),
)
.unwrap();
let records = tailer.poll_records().unwrap();
assert_eq!(records.len(), 2);
assert!(matches!(&records[0], EventRecord::Heartbeat(h) if h.elapsed_s == 30));
assert!(matches!(&records[1], EventRecord::Build(e) if e.result == EventResult::Miss));
log_heartbeat(&log_path, &test_heartbeat("gkrust", 60)).unwrap();
assert!(
tailer.poll().unwrap().is_empty(),
"builds-only poll must skip a heartbeat-only append"
);
}
#[test]
fn typical_compile_ms_is_a_robust_per_crate_median() {
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("events.jsonl");
assert_eq!(
typical_compile_ms(&log_path, "gkrust", "/work/tree"),
None,
"no history → no estimate"
);
for ms in [100_000u64, 101_000, 99_000, 100_500, 100_200] {
log_event(
&log_path,
&test_event("gkrust", EventResult::Miss, ms, ms, 1024, "k"),
)
.unwrap();
}
log_event(
&log_path,
&test_event("gkrust", EventResult::Miss, 900_000, 900_000, 1024, "k"),
)
.unwrap();
log_event(
&log_path,
&test_event("serde", EventResult::Miss, 2_000, 2_000, 64, "k"),
)
.unwrap();
log_event(
&log_path,
&test_event("gkrust", EventResult::LocalHit, 5, 0, 1024, "k"),
)
.unwrap();
let typical = typical_compile_ms(&log_path, "gkrust", "/work/tree").unwrap();
assert!(
(99_000..=101_000).contains(&typical),
"median must sit in the cluster and shed the 900s outlier, got {typical}"
);
assert_eq!(
typical_compile_ms(&log_path, "serde", "/work/tree"),
Some(2_000)
);
}
#[test]
fn test_event_tailer() {
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("events.jsonl");
let mut tailer = EventTailer::from_start(log_path.clone());
assert_eq!(tailer.poll().unwrap().len(), 0);
let event = test_event("tokio", EventResult::Miss, 5000, 4800, 8388608, "def456");
log_event(&log_path, &event).unwrap();
let new_events = tailer.poll().unwrap();
assert_eq!(new_events.len(), 1);
assert_eq!(tailer.poll().unwrap().len(), 0);
log_event(&log_path, &event).unwrap();
let new_events = tailer.poll().unwrap();
assert_eq!(new_events.len(), 1);
}
#[test]
fn test_event_rotation() {
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("events.jsonl");
for i in 0..100 {
let event = test_event(
&format!("crate_{i}"),
EventResult::LocalHit,
1,
25,
1024,
&format!("key_{i}"),
);
log_event(&log_path, &event).unwrap();
}
rotate_if_needed(&log_path, 10000, 10).unwrap();
let events = read_events(&log_path).unwrap();
assert_eq!(events.len(), 10);
assert_eq!(events[0].crate_name, "crate_90");
rotate_if_needed(&log_path, 200, 10).unwrap();
let events = read_events(&log_path).unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].crate_name, "crate_99");
}
#[test]
fn test_event_result_display() {
assert_eq!(EventResult::LocalHit.to_string(), "local_hit");
assert_eq!(EventResult::PrefetchHit.to_string(), "prefetch_hit");
assert_eq!(EventResult::RemoteHit.to_string(), "remote_hit");
assert_eq!(EventResult::Dup.to_string(), "dup");
assert_eq!(EventResult::Miss.to_string(), "miss");
assert_eq!(EventResult::Error.to_string(), "error");
assert_eq!(EventResult::Passthrough.to_string(), "passthrough");
assert_eq!(EventResult::Skipped.to_string(), "skipped");
}
#[test]
fn test_read_events_nonexistent_file() {
let events = read_events(Path::new("/nonexistent/events.jsonl")).unwrap();
assert!(events.is_empty());
}
#[test]
fn test_read_events_with_invalid_lines() {
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("events.jsonl");
let event = test_event("valid", EventResult::Miss, 100, 90, 1024, "key");
log_event(&log_path, &event).unwrap();
use std::io::Write;
let mut f = OpenOptions::new().append(true).open(&log_path).unwrap();
writeln!(f, "this is not json").unwrap();
writeln!(f, "{{}}").unwrap();
let events = read_events(&log_path).unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].crate_name, "valid");
}
#[test]
fn test_read_events_since() {
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("events.jsonl");
let mut old_event = test_event("old", EventResult::Miss, 100, 80, 1024, "key1");
old_event.ts = Utc::now() - chrono::Duration::hours(2);
let new_event = test_event("new", EventResult::LocalHit, 10, 250, 512, "key2");
log_event(&log_path, &old_event).unwrap();
log_event(&log_path, &new_event).unwrap();
let since = Utc::now() - chrono::Duration::hours(1);
let events = read_events_since(&log_path, since).unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].crate_name, "new");
}
#[test]
fn test_compute_stats() {
let events = vec![
test_event("a", EventResult::LocalHit, 10, 300, 100, "k1"),
test_event("b", EventResult::PrefetchHit, 5, 250, 150, "k1b"),
test_event("c", EventResult::RemoteHit, 50, 900, 200, "k2"),
test_event("dup", EventResult::Dup, 700, 650, 400, "kdup"),
test_event("d", EventResult::Miss, 1000, 950, 500, "k3"),
test_event("e", EventResult::Error, 5, 0, 0, "k4"),
test_event("f", EventResult::Skipped, 0, 0, 0, "k5"),
test_event("g", EventResult::Passthrough, 25, 0, 0, ""),
];
let stats = compute_stats(&events);
assert_eq!(stats.total, 8);
assert_eq!(stats.local_hits, 1);
assert_eq!(stats.prefetch_hits, 1);
assert_eq!(stats.remote_hits, 1);
assert_eq!(stats.dups, 1);
assert_eq!(stats.misses, 1);
assert_eq!(stats.errors, 1);
assert_eq!(stats.total_size, 1350);
assert_eq!(stats.total_elapsed_ms, 1770);
assert_eq!(stats.hit_elapsed_ms, 65);
assert_eq!(stats.miss_elapsed_ms, 1700);
assert_eq!(stats.hit_compile_time_ms, 1450);
assert_eq!(stats.miss_compile_time_ms, 1600);
}
#[test]
fn test_compute_stats_empty() {
let stats = compute_stats(&[]);
assert_eq!(stats.total, 0);
assert_eq!(stats.local_hits, 0);
}
#[test]
fn test_clear_events() {
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("events.jsonl");
let event = test_event("test", EventResult::Miss, 100, 80, 1024, "key");
log_event(&log_path, &event).unwrap();
assert!(!read_events(&log_path).unwrap().is_empty());
clear_events(&log_path).unwrap();
assert!(read_events(&log_path).unwrap().is_empty());
}
#[test]
fn test_clear_events_nonexistent() {
clear_events(Path::new("/nonexistent/events.jsonl")).unwrap();
}
#[test]
fn test_rotate_skips_small_file() {
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("events.jsonl");
let event = test_event("test", EventResult::Miss, 100, 80, 1024, "key");
log_event(&log_path, &event).unwrap();
let size_before = fs::metadata(&log_path).unwrap().len();
rotate_if_needed(&log_path, 1_000_000, 10).unwrap();
let size_after = fs::metadata(&log_path).unwrap().len();
assert_eq!(size_before, size_after);
}
#[test]
fn test_rotate_nonexistent() {
rotate_if_needed(Path::new("/nonexistent/events.jsonl"), 100, 10).unwrap();
}
#[test]
fn test_rotate_transfers_trims_to_keep_lines_when_oversized() {
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("transfers.jsonl");
let body: String = (0..100).map(|i| format!("line {i}\n")).collect();
fs::write(&log_path, body).unwrap();
rotate_transfers_if_needed(&log_path, 100, 10).unwrap();
let kept = fs::read_to_string(&log_path).unwrap();
let lines: Vec<&str> = kept.lines().collect();
assert_eq!(lines.len(), 10, "should keep the last 10 lines");
assert_eq!(lines[0], "line 90", "keeps the tail");
assert_eq!(lines[9], "line 99");
}
#[test]
fn test_rotate_transfers_skips_small_and_nonexistent() {
rotate_transfers_if_needed(Path::new("/nonexistent/transfers.jsonl"), 100, 10).unwrap();
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("transfers.jsonl");
fs::write(&log_path, "a\nb\n").unwrap();
rotate_transfers_if_needed(&log_path, 1_000_000, 1).unwrap();
assert_eq!(fs::read_to_string(&log_path).unwrap(), "a\nb\n");
}
#[test]
fn test_event_tailer_handles_truncation() {
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("events.jsonl");
let event = test_event("test", EventResult::Miss, 100, 80, 1024, "key");
for _ in 0..10 {
log_event(&log_path, &event).unwrap();
}
let mut tailer = EventTailer::from_start(log_path.clone());
assert_eq!(tailer.poll().unwrap().len(), 10);
fs::write(&log_path, "").unwrap();
log_event(&log_path, &event).unwrap();
let events = tailer.poll().unwrap();
assert_eq!(events.len(), 1);
}
#[test]
fn test_read_transfers_missing_file_is_empty() {
let got = read_transfers(Path::new("/nonexistent/transfers.jsonl")).unwrap();
assert!(got.is_empty());
}
#[test]
fn test_read_transfers_skips_blank_and_invalid_lines() {
let dir = tempfile::tempdir().unwrap();
let log = dir.path().join("transfers.jsonl");
fs::write(&log, "\n \nnot json at all\n{ partial: \n").unwrap();
let got = read_transfers(&log).unwrap();
assert!(got.is_empty(), "invalid transfer lines are skipped");
}
#[test]
fn test_event_tailer_handles_rename_rotation() {
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("events.jsonl");
let event = test_event("test", EventResult::Miss, 100, 80, 1024, "key");
for _ in 0..5 {
log_event(&log_path, &event).unwrap();
}
let mut tailer = EventTailer::from_start(log_path.clone());
assert_eq!(tailer.poll().unwrap().len(), 5);
rotate_if_needed(&log_path, 1500, 2).unwrap();
for _ in 0..10 {
log_event(&log_path, &event).unwrap();
}
let events = tailer.poll().unwrap();
assert_eq!(events.len(), 2 + 10);
}
#[test]
fn test_concurrent_log_append_and_rotate() {
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
let dir = tempfile::tempdir().unwrap();
let log_path = dir.path().join("events.jsonl");
let first_event = test_event("0", EventResult::LocalHit, 10, 10, 100, "key0");
log_event(&log_path, &first_event).unwrap();
let running = Arc::new(AtomicBool::new(true));
let appender_done = Arc::new(AtomicBool::new(false));
let log_path_clone = log_path.clone();
let running_clone = running.clone();
let rotator = std::thread::spawn(move || {
while running_clone.load(Ordering::Relaxed) {
let _ = rotate_if_needed(&log_path_clone, 65536, 100);
std::thread::sleep(std::time::Duration::from_millis(20));
}
});
let log_path_clone2 = log_path.clone();
let appender_done_clone = appender_done.clone();
let appender = std::thread::spawn(move || {
for i in 1..=200 {
let event = test_event(&i.to_string(), EventResult::LocalHit, 10, 10, 100, "key");
log_event(&log_path_clone2, &event).expect("log_event failed");
std::thread::sleep(std::time::Duration::from_millis(3));
}
appender_done_clone.store(true, Ordering::Relaxed);
});
let log_path_clone3 = log_path.clone();
let running_clone3 = running.clone();
let appender_done_clone2 = appender_done.clone();
let tailer_thread = std::thread::spawn(move || {
let mut tailer = EventTailer::from_start(log_path_clone3);
let mut polled_ids = std::collections::HashSet::new();
while running_clone3.load(Ordering::Relaxed)
|| !appender_done_clone2.load(Ordering::Relaxed)
{
if let Ok(events) = tailer.poll() {
for event in events {
if let Ok(id) = event.crate_name.parse::<usize>() {
polled_ids.insert(id);
}
}
}
std::thread::sleep(std::time::Duration::from_millis(2));
}
if let Ok(events) = tailer.poll() {
for event in events {
if let Ok(id) = event.crate_name.parse::<usize>() {
polled_ids.insert(id);
}
}
}
polled_ids
});
appender.join().unwrap();
running.store(false, Ordering::Relaxed);
rotator.join().unwrap();
let polled_ids = tailer_thread.join().unwrap();
let content = fs::read_to_string(&log_path).unwrap();
assert!(content.ends_with("\n"), "log must end with a newline");
assert!(
!content.ends_with("\n\n"),
"log must not end with double newlines"
);
let mut ids = Vec::new();
for line in content.lines() {
if line.trim().is_empty() {
continue;
}
let event: BuildEvent =
serde_json::from_str(line).expect("each line must be valid JSON");
let id: usize = event
.crate_name
.parse()
.expect("crate name must be parsed as id");
ids.push(id);
}
assert!(!ids.is_empty(), "at least some events must survive");
assert!(ids.len() > 1, "multiple events must survive");
let start = ids[0];
let end = ids[ids.len() - 1];
let expected: Vec<usize> = (start..=end).collect();
assert_eq!(
ids, expected,
"there must be no missing events or gaps in the surviving log: got {:?}",
ids
);
let mut actual_polled: Vec<usize> = polled_ids.into_iter().collect();
actual_polled.sort();
assert!(
!actual_polled.is_empty(),
"EventTailer must have observed some events"
);
let start_id = actual_polled[0];
let end_id = *actual_polled.last().unwrap();
let expected_sequence: Vec<usize> = (start_id..=end_id).collect();
assert_eq!(
actual_polled, expected_sequence,
"EventTailer must not miss any events in its observed sequence: got {:?}",
actual_polled
);
assert_eq!(
end_id,
200,
"EventTailer stopped following the log across rotation: never observed the \
final event (saw {} ids, max {:?})",
actual_polled.len(),
actual_polled.last()
);
}
}