use crate::config::{resolve_global_xbp_root_dir, resolve_worktree_watch_config};
use chrono::Utc;
use colored::Colorize;
use serde::Serialize;
use std::collections::BTreeMap;
use std::fs;
use std::path::{Path, PathBuf};
use std::thread;
use std::time::Duration;
use sysinfo::{Pid, ProcessRefreshKind, ProcessesToUpdate, System, UpdateKind};
#[derive(Debug, Clone)]
pub struct ResourceMonitorOptions {
pub json: bool,
pub watch_seconds: Option<u64>,
pub samples: u32,
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct ResourceReport {
captured_at: String,
this_process: ProcessRow,
xbp_processes: Vec<ProcessRow>,
totals: ProcessTotals,
disk_cache: DiskCacheReport,
worktree_watch: WorktreeWatchResourceSummary,
notes: Vec<String>,
}
#[derive(Debug, Serialize, Clone)]
#[serde(rename_all = "camelCase")]
struct ProcessRow {
pid: u32,
role: String,
name: String,
cpu_percent: f64,
rss_mb: f64,
virtual_mb: f64,
threads: Option<usize>,
cmdline: String,
exe: Option<String>,
}
#[derive(Debug, Serialize, Default)]
#[serde(rename_all = "camelCase")]
struct ProcessTotals {
process_count: usize,
cpu_percent: f64,
rss_mb: f64,
virtual_mb: f64,
}
#[derive(Debug, Serialize, Default)]
#[serde(rename_all = "camelCase")]
struct DiskCacheReport {
root: String,
total_mb: f64,
file_count: u64,
directories: BTreeMap<String, DirSize>,
}
#[derive(Debug, Serialize, Default, Clone)]
#[serde(rename_all = "camelCase")]
struct DirSize {
mb: f64,
files: u64,
}
#[derive(Debug, Serialize, Default)]
#[serde(rename_all = "camelCase")]
struct WorktreeWatchResourceSummary {
parent_cover: Option<String>,
repo_rank_count: usize,
tier_counts: BTreeMap<String, u64>,
schedule: BTreeMap<String, u64>,
estimated_fs_watches_max: u64,
spool_mutations_mb: f64,
spool_mutations_files: u64,
}
pub async fn run_resource_monitor(opts: ResourceMonitorOptions) -> Result<(), String> {
let samples = opts.samples.max(1);
let interval = opts.watch_seconds.unwrap_or(0);
if interval == 0 {
let report = collect_report(samples)?;
print_report(&report, opts.json);
return Ok(());
}
loop {
let report = collect_report(samples)?;
if opts.json {
print_report(&report, true);
} else {
print!("\x1b[2J\x1b[H");
print_report(&report, false);
println!(
"\n{} refresh every {}s · Ctrl+C to stop",
"·".bright_black(),
interval
);
}
tokio::time::sleep(Duration::from_secs(interval)).await;
}
}
fn collect_report(cpu_samples: u32) -> Result<ResourceReport, String> {
let mut notes = Vec::new();
let processes = sample_xbp_processes(cpu_samples, &mut notes);
let this_pid = std::process::id();
let this_process = processes
.iter()
.find(|p| p.pid == this_pid)
.cloned()
.unwrap_or_else(|| ProcessRow {
pid: this_pid,
role: "cli".into(),
name: "xbp".into(),
cpu_percent: 0.0,
rss_mb: 0.0,
virtual_mb: 0.0,
threads: None,
cmdline: std::env::args().collect::<Vec<_>>().join(" "),
exe: std::env::current_exe()
.ok()
.map(|p| p.display().to_string()),
});
let totals = ProcessTotals {
process_count: processes.len(),
cpu_percent: round1(processes.iter().map(|p| p.cpu_percent).sum()),
rss_mb: round1(processes.iter().map(|p| p.rss_mb).sum()),
virtual_mb: round1(processes.iter().map(|p| p.virtual_mb).sum()),
};
let disk_cache = measure_disk_cache();
let worktree_watch = summarize_worktree_watch(&disk_cache);
if processes.len() <= 1 {
notes.push(
"Only this CLI sample is running — start worktree-watch with `xbp worktree-watch start --detach` to monitor the long-lived watcher."
.into(),
);
}
if totals.rss_mb > 200.0 {
notes.push(
"Aggregate RSS is high — prefer a release build (`cargo build -p xbp --release`) and tighter worktree_watch knobs (max_commit_checks_per_tick / max_fs_watches).".into(),
);
}
Ok(ResourceReport {
captured_at: Utc::now().to_rfc3339(),
this_process,
xbp_processes: processes,
totals,
disk_cache,
worktree_watch,
notes,
})
}
fn sample_xbp_processes(samples: u32, notes: &mut Vec<String>) -> Vec<ProcessRow> {
let mut system = System::new();
let refresh = ProcessRefreshKind::everything()
.with_cmd(UpdateKind::Always)
.with_exe(UpdateKind::Always)
.with_memory();
let rounds = samples.max(2);
for i in 0..rounds {
system.refresh_processes_specifics(ProcessesToUpdate::All, true, refresh);
if i + 1 < rounds {
thread::sleep(Duration::from_millis(250));
}
}
let mut rows = Vec::new();
for (pid, proc) in system.processes() {
if !process_is_xbp(proc) {
continue;
}
let pid_u = pid.as_u32();
let cmd_parts: Vec<String> = proc
.cmd()
.iter()
.map(|s| s.to_string_lossy().into_owned())
.collect();
let cmdline = if cmd_parts.is_empty() {
proc.name().to_string_lossy().into_owned()
} else {
cmd_parts.join(" ")
};
let role = classify_role(&cmd_parts, &cmdline);
let rss = proc.memory() as f64 / (1024.0 * 1024.0);
let virt = proc.virtual_memory() as f64 / (1024.0 * 1024.0);
rows.push(ProcessRow {
pid: pid_u,
role,
name: proc.name().to_string_lossy().into_owned(),
cpu_percent: round1(f64::from(proc.cpu_usage())),
rss_mb: round1(rss),
virtual_mb: round1(virt),
threads: None, cmdline: truncate_middle(&cmdline, 140),
exe: proc
.exe()
.map(|p| p.display().to_string()),
});
}
rows.sort_by(|a, b| {
b.rss_mb
.partial_cmp(&a.rss_mb)
.unwrap_or(std::cmp::Ordering::Equal)
.then_with(|| a.pid.cmp(&b.pid))
});
if rows.is_empty() {
notes.push("No xbp processes matched (unexpected — at least this CLI should appear).".into());
}
rows
}
fn process_is_xbp(proc: &sysinfo::Process) -> bool {
let name = proc.name().to_string_lossy().to_ascii_lowercase();
if name == "xbp" || name == "xbp.exe" {
return true;
}
if let Some(exe) = proc.exe() {
let s = exe.to_string_lossy().to_ascii_lowercase().replace('\\', "/");
if s.ends_with("/xbp.exe") || s.ends_with("/xbp") {
return true;
}
}
false
}
fn classify_role(cmd: &[String], joined: &str) -> String {
let lower = joined.to_ascii_lowercase();
if lower.contains("worktree-watch") {
if lower.contains(" tray") || lower.contains("\\tray") || lower.contains(" tray") {
return "worktree-watch tray".into();
}
if lower.contains("--parent") {
return "worktree-watch parent".into();
}
return "worktree-watch".into();
}
if std::env::var_os("PORT_XBP_MCP").is_some()
|| lower.contains("mcp") && lower.contains("serve")
{
return "mcp".into();
}
if std::env::var_os("PORT_XBP_API").is_some() || lower.contains(" api") {
if cmd.iter().any(|c| c.contains("api")) {
return "api".into();
}
}
if lower.contains(" resource") {
return "resource-monitor".into();
}
"cli".into()
}
fn measure_disk_cache() -> DiskCacheReport {
let root = resolve_global_xbp_root_dir();
let mut directories = BTreeMap::new();
let mut total_bytes: u64 = 0;
let mut total_files: u64 = 0;
if root.is_dir() {
if let Ok(rd) = fs::read_dir(&root) {
for ent in rd.flatten() {
let path = ent.path();
let name = ent.file_name().to_string_lossy().into_owned();
let size = dir_size(&path);
total_bytes = total_bytes.saturating_add(size.bytes);
total_files = total_files.saturating_add(size.files);
directories.insert(
name,
DirSize {
mb: round1(size.bytes as f64 / (1024.0 * 1024.0)),
files: size.files,
},
);
}
}
}
DiskCacheReport {
root: root.display().to_string(),
total_mb: round1(total_bytes as f64 / (1024.0 * 1024.0)),
file_count: total_files,
directories,
}
}
struct RawSize {
bytes: u64,
files: u64,
}
fn dir_size(path: &Path) -> RawSize {
if path.is_file() {
let bytes = fs::metadata(path).map(|m| m.len()).unwrap_or(0);
return RawSize { bytes, files: 1 };
}
let mut bytes = 0u64;
let mut files = 0u64;
let mut stack = vec![path.to_path_buf()];
let mut visited = 0u32;
const MAX_ENTRIES: u32 = 200_000; while let Some(dir) = stack.pop() {
let Ok(rd) = fs::read_dir(&dir) else {
continue;
};
for ent in rd.flatten() {
visited += 1;
if visited > MAX_ENTRIES {
return RawSize { bytes, files };
}
let p = ent.path();
if p.is_dir() {
stack.push(p);
} else if let Ok(meta) = ent.metadata() {
bytes = bytes.saturating_add(meta.len());
files = files.saturating_add(1);
}
}
}
RawSize { bytes, files }
}
fn summarize_worktree_watch(disk: &DiskCacheReport) -> WorktreeWatchResourceSummary {
let cfg = resolve_worktree_watch_config();
let mut tier_counts: BTreeMap<String, u64> = BTreeMap::new();
for r in &cfg.repo_ranks {
let tier = r
.tier
.as_deref()
.unwrap_or("unknown")
.to_ascii_lowercase();
*tier_counts.entry(tier).or_default() += 1;
}
let mut schedule = BTreeMap::new();
schedule.insert("idleAfterSeconds".into(), cfg.idle_after_seconds());
schedule.insert("coldAfterSeconds".into(), cfg.cold_after_seconds());
schedule.insert("hotCommitCheckSeconds".into(), cfg.hot_commit_check_seconds());
schedule.insert(
"warmCommitCheckSeconds".into(),
cfg.warm_commit_check_seconds(),
);
schedule.insert(
"coldCommitCheckSeconds".into(),
cfg.cold_commit_check_seconds(),
);
schedule.insert("parentRescanSeconds".into(), cfg.parent_rescan_seconds());
schedule.insert("maxHotRepos".into(), cfg.max_hot_repos() as u64);
schedule.insert("maxFsWatches".into(), cfg.max_fs_watches() as u64);
schedule.insert(
"maxCommitChecksPerTick".into(),
cfg.max_commit_checks_per_tick() as u64,
);
let mutations = disk
.directories
.get("mutations")
.cloned()
.unwrap_or_default();
WorktreeWatchResourceSummary {
parent_cover: detect_parent_cover(),
repo_rank_count: cfg.repo_ranks.len(),
tier_counts,
schedule,
estimated_fs_watches_max: cfg.max_fs_watches() as u64,
spool_mutations_mb: mutations.mb,
spool_mutations_files: mutations.files,
}
}
fn detect_parent_cover() -> Option<String> {
let home = resolve_global_xbp_root_dir();
let state_hint = home.join("worktree-watch-parent.json");
if state_hint.is_file() {
if let Ok(raw) = fs::read_to_string(&state_hint) {
if let Ok(v) = serde_json::from_str::<serde_json::Value>(&raw) {
if let Some(p) = v.get("parent").and_then(|x| x.as_str()) {
return Some(p.to_string());
}
if let Some(p) = v.get("parentRoot").and_then(|x| x.as_str()) {
return Some(p.to_string());
}
}
}
}
None
}
fn print_report(report: &ResourceReport, json: bool) {
if json {
match serde_json::to_string_pretty(report) {
Ok(s) => println!("{s}"),
Err(e) => eprintln!("json encode failed: {e}"),
}
return;
}
println!(
"{} {}",
"xbp resource".bright_cyan().bold(),
report.captured_at.bright_black()
);
println!();
println!("{}", "processes".bright_white().bold());
println!(
" {} {} process(es) · CPU {:.1}% · RSS {:.1} MiB · virt {:.0} MiB",
"Σ".bright_magenta(),
report.totals.process_count,
report.totals.cpu_percent,
report.totals.rss_mb,
report.totals.virtual_mb
);
println!();
println!(
" {:<8} {:<22} {:>7} {:>9} {:>10} {}",
"PID".bright_black(),
"ROLE".bright_black(),
"CPU%".bright_black(),
"RSS_MiB".bright_black(),
"VIRT_MiB".bright_black(),
"CMD".bright_black()
);
for p in &report.xbp_processes {
let self_mark = if p.pid == report.this_process.pid {
"*".bright_green().to_string()
} else {
" ".into()
};
println!(
"{self_mark} {:<7} {:<22} {:>7.1} {:>9.1} {:>10.0} {}",
p.pid,
truncate_middle(&p.role, 22),
p.cpu_percent,
p.rss_mb,
p.virtual_mb,
p.cmdline.bright_black()
);
}
println!(
" {} this CLI sample",
"*".bright_green()
);
println!();
println!("{}", "disk cache (~/.xbp)".bright_white().bold());
println!(
" root {} · {:.1} MiB · {} file(s)",
report.disk_cache.root.bright_black(),
report.disk_cache.total_mb,
report.disk_cache.file_count
);
if !report.disk_cache.directories.is_empty() {
println!(
" {:<22} {:>10} {:>8}",
"DIR".bright_black(),
"MiB".bright_black(),
"FILES".bright_black()
);
let mut dirs: Vec<_> = report.disk_cache.directories.iter().collect();
dirs.sort_by(|a, b| {
b.1.mb
.partial_cmp(&a.1.mb)
.unwrap_or(std::cmp::Ordering::Equal)
});
for (name, size) in dirs.into_iter().take(16) {
println!(" {:<22} {:>10.1} {:>8}", name, size.mb, size.files);
}
}
println!();
let ww = &report.worktree_watch;
println!("{}", "worktree-watch (config + spool)".bright_white().bold());
if let Some(p) = &ww.parent_cover {
println!(" parent cover {}", p.cyan());
}
println!(
" ranks {} · spool mutations {:.1} MiB / {} files · max FS watches {}",
ww.repo_rank_count, ww.spool_mutations_mb, ww.spool_mutations_files, ww.estimated_fs_watches_max
);
if !ww.tier_counts.is_empty() {
let tiers: Vec<String> = ww
.tier_counts
.iter()
.map(|(k, v)| format!("{k}={v}"))
.collect();
println!(" tiers {}", tiers.join(" · "));
}
if !ww.schedule.is_empty() {
println!(
" schedule hot/warm/cold HEAD {}/{}/{}s · cold_after {}s · commit_budget {}",
ww.schedule.get("hotCommitCheckSeconds").copied().unwrap_or(0),
ww.schedule
.get("warmCommitCheckSeconds")
.copied()
.unwrap_or(0),
ww.schedule
.get("coldCommitCheckSeconds")
.copied()
.unwrap_or(0),
ww.schedule.get("coldAfterSeconds").copied().unwrap_or(0),
ww.schedule
.get("maxCommitChecksPerTick")
.copied()
.unwrap_or(0),
);
}
println!();
if !report.notes.is_empty() {
println!("{}", "notes".bright_white().bold());
for n in &report.notes {
println!(" {} {}", "·".bright_yellow(), n);
}
}
}
fn round1(v: f64) -> f64 {
(v * 10.0).round() / 10.0
}
fn truncate_middle(s: &str, max: usize) -> String {
if s.chars().count() <= max {
return s.to_string();
}
if max < 5 {
return s.chars().take(max).collect();
}
let keep = max.saturating_sub(1) / 2;
let front: String = s.chars().take(keep).collect();
let back: String = s.chars().rev().take(keep).collect::<String>().chars().rev().collect();
format!("{front}…{back}")
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn classify_role_detects_worktree_parent() {
let cmd = vec![
"xbp.exe".into(),
"worktree-watch".into(),
"start".into(),
"--parent".into(),
r"C:\Users\me\Documents\GitHub".into(),
];
let joined = cmd.join(" ");
assert_eq!(classify_role(&cmd, &joined), "worktree-watch parent");
}
#[test]
fn dir_size_counts_file() {
let dir = std::env::temp_dir().join(format!("xbp-resource-test-{}", std::process::id()));
let _ = fs::remove_dir_all(&dir);
fs::create_dir_all(&dir).unwrap();
fs::write(dir.join("a.txt"), b"hello-world").unwrap();
let s = dir_size(&dir);
assert_eq!(s.files, 1);
assert!(s.bytes >= 11);
let _ = fs::remove_dir_all(&dir);
}
}