mod mixed;
mod scroll;
use crate::{
config::{McPaths, read_settings},
instructions::discover_agents_with_additional_markdown,
model_catalog::{ModelCatalogEntry, cached_model_context_window, write_catalog_cache},
output::{
ActivityEvent, ActivityId, ActivityKind, ActivityMetadata, ActivityStatus, OutputEvent,
},
providers::{ProviderEvent, StreamParser},
sessions::SessionManager,
skills::{SkillRoots, discover_skills},
tools::{ToolRuntime, process::terminate_child_tree_and_wait},
tui::state::MissionControlState,
};
use ratatui::{
Terminal, TerminalOptions, Viewport,
backend::{Backend, CrosstermBackend, TestBackend},
buffer::Buffer,
layout::{Position, Rect, Size},
};
use serde::Serialize;
use serde_json::{Value, json};
use std::{
cell::{Cell, RefCell},
collections::{BTreeMap, HashMap},
fs,
hint::black_box,
io::{self, Write},
path::{Path, PathBuf},
process::{Command, Stdio},
rc::Rc,
time::{Duration, Instant},
};
use tempfile::TempDir;
const RESPONSES_SSE: &str = include_str!("../../scripts/fixtures/cpu_profiling/responses.sse");
const CHAT_COMPLETIONS_SSE: &str =
include_str!("../../scripts/fixtures/cpu_profiling/chat_completions.sse");
const RENDERING_TRANSCRIPT: &str =
include_str!("../../scripts/fixtures/cpu_profiling/rendering_transcript.md");
const DEFAULT_ITERATIONS: u64 = 25;
thread_local! {
static EMIT_RENDER_SUMMARIES: Cell<bool> = const { Cell::new(true) };
}
#[derive(Debug)]
struct ProfileResult {
scenario: &'static str,
iterations: u64,
elapsed: Duration,
primary_metric: &'static str,
metric_value: u64,
artifact: &'static str,
counters: Vec<(&'static str, u64)>,
}
impl ProfileResult {
fn print(&self) {
let counters = self
.counters
.iter()
.map(|(key, value)| format!("{key}:{value}"))
.collect::<Vec<_>>()
.join(",");
println!(
"profile_cpu scenario={} iterations={} elapsed_ms={} primary_metric={} metric_value={} artifact={} counters={}",
self.scenario,
self.iterations,
self.elapsed.as_millis(),
self.primary_metric,
self.metric_value,
self.artifact,
counters
);
}
}
fn profile_iterations() -> u64 {
std::env::var("PROFILE_CPU_ITERATIONS")
.ok()
.and_then(|value| value.parse::<u64>().ok())
.filter(|value| *value > 0)
.unwrap_or(DEFAULT_ITERATIONS)
}
fn profile_min_duration() -> Duration {
std::env::var("PROFILE_CPU_MIN_SECONDS")
.ok()
.and_then(|value| value.parse::<u64>().ok())
.map(Duration::from_secs)
.unwrap_or_default()
}
fn run_profile(
scenario: &'static str,
primary_metric: &'static str,
artifact: &'static str,
mut run_once: impl FnMut() -> u64,
) -> ProfileResult {
let min_iterations = profile_iterations();
let min_duration = profile_min_duration();
let started = Instant::now();
let mut iterations = 0;
let mut metric_value = 0u64;
while iterations < min_iterations || started.elapsed() < min_duration {
metric_value = metric_value.saturating_add(run_once());
iterations += 1;
}
let elapsed = started.elapsed();
ProfileResult {
scenario,
iterations,
elapsed,
primary_metric,
metric_value,
artifact,
counters: vec![("per_iteration", metric_value / iterations.max(1))],
}
}
fn comma_env(name: &str) -> Option<Vec<String>> {
std::env::var(name).ok().and_then(|value| {
let values = value
.split(',')
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
.collect::<Vec<_>>();
(!values.is_empty()).then_some(values)
})
}
fn parse_usize_env(name: &str, default: usize, min: usize, max: usize) -> usize {
std::env::var(name)
.ok()
.and_then(|value| value.parse::<usize>().ok())
.map(|value| value.clamp(min, max))
.unwrap_or(default)
}
fn hash_root_label(root: &Path) -> String {
use sha2::{Digest, Sha256};
let mut hasher = Sha256::new();
hasher.update(root.to_string_lossy().as_bytes());
crate::hex::lower_hex(hasher.finalize())[..16].to_string()
}
#[derive(Debug, Clone, Copy, Serialize)]
struct FindDurationStats {
min: f64,
p50: f64,
p95: f64,
max: f64,
}
fn duration_stats(durations: &[Duration]) -> FindDurationStats {
let mut millis = durations
.iter()
.map(|duration| duration.as_secs_f64() * 1000.0)
.collect::<Vec<_>>();
millis.sort_by(f64::total_cmp);
FindDurationStats {
min: round_millis(millis.first().copied().unwrap_or_default()),
p50: round_millis(percentile(&millis, 0.50)),
p95: round_millis(percentile(&millis, 0.95)),
max: round_millis(millis.last().copied().unwrap_or_default()),
}
}
fn percentile(sorted: &[f64], percentile: f64) -> f64 {
if sorted.is_empty() {
return 0.0;
}
let rank = (percentile * sorted.len() as f64).ceil() as usize;
sorted[rank.saturating_sub(1).min(sorted.len() - 1)]
}
fn round_millis(value: f64) -> f64 {
(value * 1000.0).round() / 1000.0
}
#[derive(Debug, Clone, Copy, Default)]
struct ResourceSnapshot {
rss_kb: Option<i64>,
fd_count: Option<i64>,
thread_count: Option<i64>,
}
impl ResourceSnapshot {
fn capture() -> Self {
Self {
rss_kb: current_rss_kb(),
fd_count: current_fd_count(),
thread_count: current_thread_count(),
}
}
fn delta(self, after: Self) -> ResourceDelta {
ResourceDelta {
rss_kb_delta: option_delta(self.rss_kb, after.rss_kb),
fd_delta: option_delta(self.fd_count, after.fd_count),
thread_delta: option_delta(self.thread_count, after.thread_count),
}
}
}
#[derive(Debug, Clone, Copy, Default, Serialize)]
struct ResourceDelta {
rss_kb_delta: Option<i64>,
fd_delta: Option<i64>,
thread_delta: Option<i64>,
}
fn option_delta(before: Option<i64>, after: Option<i64>) -> Option<i64> {
Some(after? - before?)
}
fn current_rss_kb() -> Option<i64> {
let output = Command::new("ps")
.args(["-o", "rss=", "-p", &std::process::id().to_string()])
.output()
.ok()?;
if !output.status.success() {
return None;
}
String::from_utf8_lossy(&output.stdout)
.trim()
.parse::<i64>()
.ok()
}
fn current_fd_count() -> Option<i64> {
fs::read_dir("/dev/fd").ok()?.count().try_into().ok()
}
#[cfg(target_os = "linux")]
fn current_thread_count() -> Option<i64> {
fs::read_dir("/proc/self/task")
.ok()?
.count()
.try_into()
.ok()
}
#[cfg(all(unix, not(target_os = "linux")))]
fn current_thread_count() -> Option<i64> {
let output = Command::new("ps")
.args(["-M", "-p", &std::process::id().to_string()])
.output()
.ok()?;
if !output.status.success() {
return None;
}
let lines = String::from_utf8_lossy(&output.stdout)
.lines()
.filter(|line| !line.trim().is_empty())
.count();
lines.checked_sub(1)?.try_into().ok()
}
#[cfg(not(unix))]
fn current_thread_count() -> Option<i64> {
None
}
fn metadata_u64(metadata: &Value, key: &str) -> u64 {
metadata
.get(key)
.and_then(Value::as_u64)
.unwrap_or_default()
}
use find::collect_churn_dirs;
use render::{
assert_render_summaries, run_render_scenario, seed_render_heavy_transcript_state,
seed_render_streaming_state,
};
mod content_search;
mod cpu;
mod find;
mod render;