#[cfg(all(feature = "comms", any(unix, windows)))]
use anyhow::Context;
use anyhow::Result;
pub(crate) fn cmd_statusline(root: Option<&std::path::Path>) -> Result<()> {
if let Some(root) = root {
println!("{}", render_repo_statusline(root));
return Ok(());
}
#[cfg(all(feature = "comms", any(unix, windows)))]
{
use basemind::comms::client::CommsClient;
use basemind::comms::ids::AgentId;
use basemind::comms::singleton;
let line = (|| -> Option<String> {
let paths = singleton::resolve_paths().ok()?;
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.ok()?;
runtime.block_on(async move {
let agent = AgentId::parse("basemind-statusline").ok()?;
let mut client = CommsClient::connect(&paths, agent, None, None).await.ok()?;
let hot = client.accessed_paths().await.ok()?;
Some(format_statusline(&hot))
})
})();
if let Some(line) = line {
println!("{line}");
}
}
Ok(())
}
const BRAND: &str = "\x1b[38;2;249;115;22m";
const CYAN: &str = "\x1b[38;5;51m";
const MAGENTA: &str = "\x1b[38;5;201m";
const LABEL: &str = "\x1b[38;5;255m";
const SEP: &str = "\x1b[38;5;240m";
const BOLD: &str = "\x1b[1m";
const RESET: &str = "\x1b[0m";
const BRAND_GLYPH: &str = "◆";
fn brand_mark() -> String {
format!("{BRAND}{BRAND_GLYPH}{RESET} {BOLD}{BRAND}basemind{RESET}")
}
fn render_repo_statusline(root: &std::path::Path) -> String {
use basemind::store::{read_status_sidecar, workspace_cache_dir};
let basemind_dir = workspace_cache_dir(root);
let Some(status) = read_status_sidecar(&basemind_dir) else {
return format!(
"{} {SEP}│{RESET} {LABEL}no index — run:{RESET} {BOLD}{CYAN}basemind scan{RESET}",
brand_mark()
);
};
let age = format_scan_age(status.scanned_unix);
let (calls, saved) = telemetry_today(&basemind_dir);
let mut out = format!(
"{} {BOLD}{CYAN}{}{RESET} {LABEL}files{RESET} {SEP}·{RESET} {BOLD}{CYAN}{age}{RESET}",
brand_mark(),
fmt_count(status.file_count as u64),
);
out.push_str(&format!(
" {SEP}│{RESET} {BOLD}{MAGENTA}{}{RESET} {LABEL}calls{RESET} {SEP}·{RESET} {BOLD}{MAGENTA}{}{RESET} {LABEL}saved{RESET}",
fmt_count(calls),
fmt_count(saved),
));
out
}
fn format_scan_age(scanned_unix: i64) -> String {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs() as i64)
.unwrap_or(0);
let delta = now - scanned_unix;
if scanned_unix <= 0 || delta < 0 {
return "never".to_string();
}
if delta < 60 {
format!("{delta}s ago")
} else if delta < 3_600 {
format!("{}m ago", delta / 60)
} else if delta < 86_400 {
format!("{}h ago", delta / 3_600)
} else {
format!("{}d ago", delta / 86_400)
}
}
#[derive(serde::Deserialize)]
struct StatuslineTelemetryRow {
ts_micros: i64,
#[serde(default)]
est_tokens_saved: u64,
}
fn telemetry_today(basemind_dir: &std::path::Path) -> (u64, u64) {
use std::io::{BufRead, BufReader};
const TAIL_ROWS: usize = 2_000;
const DAY_MICROS: i64 = 24 * 3_600 * 1_000_000;
let now_micros = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| i64::try_from(d.as_micros()).unwrap_or(i64::MAX))
.unwrap_or(0);
let cutoff = now_micros.saturating_sub(DAY_MICROS);
let Ok(file) = std::fs::File::open(basemind_dir.join("telemetry.jsonl")) else {
return (0, 0);
};
let mut tail: std::collections::VecDeque<StatuslineTelemetryRow> =
std::collections::VecDeque::with_capacity(TAIL_ROWS);
for line in BufReader::new(file).lines().map_while(Result::ok) {
if line.trim().is_empty() {
continue;
}
if let Ok(row) = serde_json::from_str::<StatuslineTelemetryRow>(&line) {
if tail.len() == TAIL_ROWS {
tail.pop_front();
}
tail.push_back(row);
}
}
let mut calls = 0u64;
let mut saved = 0u64;
for row in tail.iter().filter(|r| r.ts_micros >= cutoff) {
calls += 1;
saved = saved.saturating_add(row.est_tokens_saved);
}
(calls, saved)
}
fn fmt_count(n: u64) -> String {
if n < 1_000 {
format!("{n}")
} else if n < 10_000 {
format!("{}.{}k", n / 1_000, (n * 10 / 1_000) % 10)
} else if n < 1_000_000 {
format!("{}k", n / 1_000)
} else {
format!("{}M", n / 1_000_000)
}
}
#[cfg(all(feature = "comms", any(unix, windows)))]
fn format_statusline(workspaces: &[basemind::comms::workspace_pool::AccessedWorkspace]) -> String {
if workspaces.is_empty() {
return "bm: idle".to_string();
}
const MAX_NAMES: usize = 3;
let names: Vec<&str> = workspaces
.iter()
.take(MAX_NAMES)
.map(|w| w.root.file_name().and_then(|n| n.to_str()).unwrap_or("?"))
.collect();
let mut label = names.join(" · ");
if workspaces.len() > MAX_NAMES {
label.push_str(&format!(" +{}", workspaces.len() - MAX_NAMES));
}
format!("bm: {label} · {} hot", workspaces.len())
}
#[cfg(all(feature = "comms", any(unix, windows)))]
pub(crate) fn cmd_comms(action: crate::CommsLifecycleCmd, json: bool) -> Result<()> {
match action {
crate::CommsLifecycleCmd::Daemon => basemind::cli::comms_daemon::run(),
crate::CommsLifecycleCmd::Start => cmd_comms_start(),
crate::CommsLifecycleCmd::Stop { all: true } => cmd_comms_stop_all(json),
crate::CommsLifecycleCmd::Stop { all: false } => cmd_comms_lifecycle_rpc(CommsRpc::Stop, json),
crate::CommsLifecycleCmd::Status => cmd_comms_lifecycle_rpc(CommsRpc::Status, json),
crate::CommsLifecycleCmd::Doctor => cmd_comms_doctor(json),
}
}
#[cfg(all(feature = "comms", any(unix, windows)))]
enum CommsRpc {
Stop,
Status,
}
#[cfg(all(feature = "comms", any(unix, windows)))]
const HTTP_READY_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
#[cfg(all(feature = "comms", any(unix, windows)))]
pub(crate) fn cmd_daemon(action: crate::DaemonCmd, json: bool) -> Result<()> {
match action {
crate::DaemonCmd::Ensure => cmd_daemon_ensure(json),
}
}
#[cfg(all(feature = "comms", any(unix, windows)))]
fn cmd_daemon_ensure(json: bool) -> Result<()> {
use basemind::comms::http_frontend;
use basemind::comms::singleton;
let paths = singleton::resolve_paths().context("resolve comms paths")?;
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.context("build tokio runtime")?;
let addr = runtime.block_on(async move {
singleton::ensure_daemon(&paths)
.await
.map_err(|e| anyhow::anyhow!("ensure comms daemon: {e}"))?;
http_frontend::await_http_ready(&paths.comms_dir, HTTP_READY_TIMEOUT)
.await
.context("wait for streamable-HTTP MCP transport")
})?;
let url = http_frontend::base_url(&addr);
if json {
println!("{{\"ready\":true,\"addr\":\"{addr}\",\"url\":\"{url}\"}}");
} else {
println!("{url}");
}
Ok(())
}
#[cfg(all(feature = "comms", any(unix, windows)))]
fn cmd_comms_start() -> Result<()> {
use basemind::comms::singleton;
let paths = singleton::resolve_paths().context("resolve comms paths")?;
let socket_path = paths.socket_path.clone();
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.context("build tokio runtime")?;
runtime.block_on(async move {
singleton::ensure_daemon(&paths)
.await
.map_err(|e| anyhow::anyhow!("ensure comms daemon: {e}"))
})?;
println!("comms daemon is running ({})", socket_path.display());
Ok(())
}
#[cfg(all(feature = "comms", any(unix, windows)))]
fn cmd_comms_lifecycle_rpc(rpc: CommsRpc, json: bool) -> Result<()> {
use basemind::comms::client::CommsClient;
use basemind::comms::singleton;
let paths = singleton::resolve_paths().context("resolve comms paths")?;
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.context("build tokio runtime")?;
runtime.block_on(async move {
let root = std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from("."));
let agent = basemind::comms::identity::cli_agent_id(&root);
let mut client = CommsClient::connect(&paths, agent, None, None)
.await
.map_err(|e| anyhow::anyhow!("connect to comms daemon: {e}"))?;
match rpc {
CommsRpc::Stop => {
client.stop().await.map_err(|e| anyhow::anyhow!("stop: {e}"))?;
if json {
println!("{{\"stopped\":true}}");
} else {
println!("comms daemon stopping");
}
}
CommsRpc::Status => {
let status = client.status().await.map_err(|e| anyhow::anyhow!("status: {e}"))?;
if json {
println!(
"{}",
serde_json::to_string(&status).map_err(|e| anyhow::anyhow!("serialize status: {e}"))?
);
} else {
println!(
"pid={} version={} build={} proto={} uptime={}s threads={} subscribers={}",
status.pid,
status.version,
if status.build_id.is_empty() {
"unreported"
} else {
&status.build_id
},
status.proto_ver,
status.uptime_secs,
status.threads,
status.subscribers,
);
let ours = basemind::version::build_id();
if !status.build_id.is_empty() && status.build_id != ours {
println!(
" WARNING: this daemon is running a DIFFERENT build of the same version \
(daemon {} vs this binary {}).",
status.build_id, ours
);
println!(
" It will keep answering with its own code — a version check cannot see this. \
Restart it with `basemind comms stop` to pick up the current binary."
);
}
}
}
}
Ok::<(), anyhow::Error>(())
})?;
Ok(())
}
#[cfg(all(feature = "comms", any(unix, windows)))]
fn now_unix() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs() as i64)
.unwrap_or(0)
}
#[cfg(all(feature = "comms", any(unix, windows)))]
fn cmd_comms_doctor(json: bool) -> Result<()> {
use basemind::daemon_lock;
let daemons = daemon_lock::live_daemons();
let ceiling = daemon_lock::max_live_daemons();
let now = now_unix();
if json {
let items: Vec<serde_json::Value> = daemons
.iter()
.map(|record| {
serde_json::json!({
"pid": record.pid,
"kind": record.kind,
"dir": record.dir,
"version": record.version,
"uptime_secs": (now - record.started_unix).max(0),
})
})
.collect();
let report = serde_json::json!({
"count": daemons.len(),
"ceiling": ceiling,
"over_ceiling": daemons.len() > ceiling,
"daemons": items,
});
println!("{report}");
return Ok(());
}
if daemons.is_empty() {
println!("no live basemind daemons");
return Ok(());
}
println!("{} live daemon(s) (ceiling {ceiling}):", daemons.len());
for record in &daemons {
println!(
" pid={} kind={} version={} uptime={}s dir={}",
record.pid,
record.kind,
record.version,
(now - record.started_unix).max(0),
record.dir.display(),
);
}
if daemons.len() > ceiling {
println!(
"WARNING: {} daemons exceed the ceiling of {ceiling}; run `basemind comms stop --all` to reclaim",
daemons.len(),
);
}
Ok(())
}
#[cfg(all(feature = "comms", any(unix, windows)))]
fn cmd_comms_stop_all(json: bool) -> Result<()> {
use basemind::comms::singleton;
use basemind::daemon_lock::{self, DaemonKind};
let daemons = daemon_lock::live_daemons_of(DaemonKind::Comms);
for record in &daemons {
singleton::request_stop(&singleton::comms_socket_path(&record.dir));
}
if json {
println!("{{\"stopped\":{}}}", daemons.len());
} else if daemons.is_empty() {
println!("no live basemind daemons to stop");
} else {
println!("asked {} daemon(s) to stop", daemons.len());
}
Ok(())
}
#[cfg(all(test, feature = "comms", any(unix, windows)))]
mod statusline_tests {
use std::path::PathBuf;
use basemind::comms::workspace_pool::AccessedWorkspace;
fn ws(root: &str) -> AccessedWorkspace {
AccessedWorkspace {
root: PathBuf::from(root),
key: "k".to_string(),
idle_secs: 0,
}
}
#[test]
fn empty_hot_set_reads_idle() {
assert_eq!(super::format_statusline(&[]), "bm: idle");
}
#[test]
fn lists_workspace_basenames_and_the_hot_count() {
let hot = [ws("/repos/web"), ws("/repos/api")];
assert_eq!(super::format_statusline(&hot), "bm: web · api · 2 hot");
}
#[test]
fn caps_the_name_list_with_an_overflow_marker() {
let hot = [ws("/a/one"), ws("/a/two"), ws("/a/three"), ws("/a/four"), ws("/a/five")];
assert_eq!(super::format_statusline(&hot), "bm: one · two · three +2 · 5 hot");
}
}