From 951fc48eba9089e0d67a8985abcf71bf80e4aa20 Mon Sep 17 00:00:00 2001 From: can1357 Date: Sun, 1 Feb 2026 22:04:15 +0100 Subject: [PATCH] feat(work): added work scheduling profiler with circular buffer sampling and flamegraph visualization - Added work scheduling profiler with always-on circular buffer sampling to track CPU scheduling patterns and task execution times. - Added `getWorkProfile()` native function to retrieve profiling data including folded stacks, markdown summary, SVG flamegraph, and metrics for the last N seconds. - Added work profile debug menu item to visualize CPU scheduling patterns via flamegraph in the coding agent. - Added work profile support to report bundles with folded stacks, summary, and flamegraph visualization. - Updated `launch_blocking()` and `launch_async()` functions to accept task name tags for profiling instrumentation. - Added dependencies on inferno, smallvec, and heapless crates for profiling infrastructure. --- Cargo.lock | 79 +++++ crates/pi-natives/Cargo.toml | 3 + crates/pi-natives/src/clipboard.rs | 37 +- crates/pi-natives/src/find.rs | 2 +- crates/pi-natives/src/grep.rs | 4 +- crates/pi-natives/src/html.rs | 2 +- crates/pi-natives/src/image.rs | 17 +- crates/pi-natives/src/keys.rs | 2 +- crates/pi-natives/src/shell.rs | 2 +- crates/pi-natives/src/work.rs | 322 ++++++++++++++++-- packages/coding-agent/CHANGELOG.md | 4 + packages/coding-agent/src/debug/index.ts | 38 +++ .../coding-agent/src/debug/report-bundle.ts | 18 + packages/natives/CHANGELOG.md | 6 + packages/natives/src/native.ts | 2 + packages/natives/src/work/index.ts | 11 + packages/natives/src/work/types.ts | 31 ++ 17 files changed, 520 insertions(+), 60 deletions(-) create mode 100644 packages/natives/src/work/index.ts create mode 100644 packages/natives/src/work/types.ts diff --git a/Cargo.lock b/Cargo.lock index 405f38773..c8a817377 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -340,6 +340,12 @@ version = "1.25.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c8efb64bd706a16a1bdde310ae86b351e4d21550d98d056f22f8a7f7a2183fec" +[[package]] +name = "byteorder" +version = "1.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1fd0f2584146f6f2ef48085050886acf353beff7305ebd1ae69500e27c67f64b" + [[package]] name = "byteorder-lite" version = "0.1.0" @@ -1051,6 +1057,15 @@ dependencies = [ "zerocopy", ] +[[package]] +name = "hash32" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "47d60b12902ba28e2730cd37e95b8c9223af2808df9e902d4df49588d1470606" +dependencies = [ + "byteorder", +] + [[package]] name = "hashbrown" version = "0.15.5" @@ -1073,6 +1088,17 @@ dependencies = [ "foldhash 0.2.0", ] +[[package]] +name = "heapless" +version = "0.9.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2af2455f757db2b292a9b1768c4b70186d443bcb3b316252d6b540aec1cd89ed" +dependencies = [ + "hash32", + "serde_core", + "stable_deref_trait", +] + [[package]] name = "heck" version = "0.5.0" @@ -1301,6 +1327,22 @@ dependencies = [ "hashbrown 0.16.1", ] +[[package]] +name = "inferno" +version = "0.12.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d35223c50fdd26419a4ccea2c73be68bd2b29a3d7d6123ffe101c17f4c20a52a" +dependencies = [ + "ahash", + "itoa", + "log", + "num-format", + "once_cell", + "quick-xml", + "rgb", + "str_stack", +] + [[package]] name = "intl-memoizer" version = "0.5.3" @@ -1335,6 +1377,12 @@ dependencies = [ "either", ] +[[package]] +name = "itoa" +version = "1.0.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "92ecc6618181def0457392ccd0ee51198e065e016d1d527a7ac1b6dc7c1f09d2" + [[package]] name = "js-sys" version = "0.3.85" @@ -1633,6 +1681,16 @@ dependencies = [ "num-traits", ] +[[package]] +name = "num-format" +version = "0.4.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a652d9771a63711fd3c3deb670acfbe5c30a4072e664d7a3bf5a9e1056ac72c3" +dependencies = [ + "arrayvec", + "itoa", +] + [[package]] name = "num-integer" version = "0.1.46" @@ -1942,9 +2000,11 @@ dependencies = [ "grep-matcher", "grep-regex", "grep-searcher", + "heapless", "html-to-markdown-rs", "ignore", "image", + "inferno", "libc", "napi", "napi-build", @@ -1953,6 +2013,7 @@ dependencies = [ "parking_lot", "phf 0.13.1", "rayon", + "smallvec", "syntect", "sysinfo", "tokio", @@ -2212,6 +2273,15 @@ version = "0.8.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7a2d987857b319362043e95f5353c0535c1f58eec5336fdfcf626430af7def58" +[[package]] +name = "rgb" +version = "0.8.52" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0c6a884d2998352bb4daf0183589aec883f16a6da1f4dde84d8e2e9a5409a1ce" +dependencies = [ + "bytemuck", +] + [[package]] name = "rlimit" version = "0.10.2" @@ -2351,6 +2421,9 @@ name = "smallvec" version = "1.15.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "67b1b7a3b5fe4f1376887184045fcf45c69e92af734b7aaddc05fb777b6fbd03" +dependencies = [ + "serde", +] [[package]] name = "socket2" @@ -2368,6 +2441,12 @@ version = "1.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6ce2be8dc25455e1f91df71bfa12ad37d7af1092ae736f3a6cd0e37bc7810596" +[[package]] +name = "str_stack" +version = "0.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9091b6114800a5f2141aee1d1b9d6ca3592ac062dc5decb3764ec5895a47b4eb" + [[package]] name = "string_cache" version = "0.9.0" diff --git a/crates/pi-natives/Cargo.toml b/crates/pi-natives/Cargo.toml index d3e05a302..f1fde76a3 100644 --- a/crates/pi-natives/Cargo.toml +++ b/crates/pi-natives/Cargo.toml @@ -29,6 +29,7 @@ grep-matcher = "0.1" globset = "0.4" ignore = "0.4" rayon = "1.10" +inferno = { version = "0.12", default-features = false } image = { version = "0.25", default-features = false, features = [ "png", "jpeg", @@ -46,6 +47,8 @@ syntect = { version = "5.3", default-features = false, features = [ ] } html-to-markdown-rs = { version = "2.24", default-features = false } phf = { version = "0.13", features = ["macros"] } +smallvec = { version = "1.15.1", features = ["serde", "write", "union", "specialization"] } +heapless = { version = "0.9.2", features = ["serde", "nightly"] } [target.'cfg(unix)'.dependencies] libc = "0.2" diff --git a/crates/pi-natives/src/clipboard.rs b/crates/pi-natives/src/clipboard.rs index 45f1261fa..a10f33b9a 100644 --- a/crates/pi-natives/src/clipboard.rs +++ b/crates/pi-natives/src/clipboard.rs @@ -58,7 +58,7 @@ fn encode_png(image: ImageData<'_>) -> Result> { /// Returns an error if clipboard access fails. #[napi(js_name = "copyToClipboard")] pub async fn copy_to_clipboard(text: String) -> Result<()> { - launch_blocking(move || -> Result<()> { + launch_blocking("clipboard.copy", move || -> Result<()> { let mut clipboard = Clipboard::new() .map_err(|err| Error::from_reason(format!("Failed to access clipboard: {err}")))?; clipboard @@ -79,22 +79,23 @@ pub async fn copy_to_clipboard(text: String) -> Result<()> { /// Returns an error if clipboard access fails or image encoding fails. #[napi(js_name = "readImageFromClipboard")] pub async fn read_image_from_clipboard() -> Result> { - let result = launch_blocking(move || -> Result> { - let mut clipboard = Clipboard::new() - .map_err(|err| Error::from_reason(format!("Failed to access clipboard: {err}")))?; - match clipboard.get_image() { - Ok(image) => { - let bytes = encode_png(image)?; - Ok(Some(ClipboardImage { - data: Uint8Array::from(bytes), - mime_type: "image/png".to_string(), - })) - }, - Err(ClipboardError::ContentNotAvailable) => Ok(None), - Err(err) => Err(Error::from_reason(format!("Failed to read clipboard image: {err}"))), - } - }) - .wait() - .await?; + let result = + launch_blocking("clipboard.read_image", move || -> Result> { + let mut clipboard = Clipboard::new() + .map_err(|err| Error::from_reason(format!("Failed to access clipboard: {err}")))?; + match clipboard.get_image() { + Ok(image) => { + let bytes = encode_png(image)?; + Ok(Some(ClipboardImage { + data: Uint8Array::from(bytes), + mime_type: "image/png".to_string(), + })) + }, + Err(ClipboardError::ContentNotAvailable) => Ok(None), + Err(err) => Err(Error::from_reason(format!("Failed to read clipboard image: {err}"))), + } + }) + .wait() + .await?; Ok(result) } diff --git a/crates/pi-natives/src/find.rs b/crates/pi-natives/src/find.rs index 78b9f7aae..c90f3539c 100644 --- a/crates/pi-natives/src/find.rs +++ b/crates/pi-natives/src/find.rs @@ -317,7 +317,7 @@ pub async fn find( let mentions_node_modules = pattern.contains("node_modules"); let sort_by_mtime = sort_by_mtime.unwrap_or(false); - launch_blocking(move || { + launch_blocking("find", move || { let cancelled = AtomicBool::new(false); let config = FindConfig { root: search_path, diff --git a/crates/pi-natives/src/grep.rs b/crates/pi-natives/src/grep.rs index e173d4140..dad40f608 100644 --- a/crates/pi-natives/src/grep.rs +++ b/crates/pi-natives/src/grep.rs @@ -1034,7 +1034,7 @@ pub async fn grep( ThreadsafeFunction, >, ) -> Result { - launch_blocking(move || grep_sync(options, on_match.as_ref())) + launch_blocking("grep", move || grep_sync(options, on_match.as_ref())) .wait() .await } @@ -1156,7 +1156,7 @@ fn fuzzy_find_sync(options: FuzzyFindOptions) -> Result { /// Matching file and directory entries. #[napi(js_name = "fuzzyFind")] pub async fn fuzzy_find(options: FuzzyFindOptions) -> Result { - launch_blocking(move || fuzzy_find_sync(options)) + launch_blocking("fuzzy_find", move || fuzzy_find_sync(options)) .wait() .await } diff --git a/crates/pi-natives/src/html.rs b/crates/pi-natives/src/html.rs index da6a4aba5..d6c81e5a7 100644 --- a/crates/pi-natives/src/html.rs +++ b/crates/pi-natives/src/html.rs @@ -31,7 +31,7 @@ pub async fn html_to_markdown( let clean_content = options.clean_content.unwrap_or(false); let skip_images = options.skip_images.unwrap_or(false); - launch_blocking(move || { + launch_blocking("html_to_markdown", move || { let conversion_opts = ConversionOptions { skip_images, preprocessing: PreprocessingOptions { diff --git a/crates/pi-natives/src/image.rs b/crates/pi-natives/src/image.rs index d7633cf7e..ac2fa7184 100644 --- a/crates/pi-natives/src/image.rs +++ b/crates/pi-natives/src/image.rs @@ -62,7 +62,7 @@ impl PhotonImage { #[napi(factory, js_name = "parse")] pub async fn parse(bytes: Uint8Array) -> Result { let bytes = bytes.as_ref().to_vec(); - let img = launch_blocking(move || -> Result { + let img = launch_blocking("image.decode", move || -> Result { let reader = ImageReader::new(Cursor::new(bytes)) .with_guessed_format() .map_err(|e| Error::from_reason(format!("Failed to detect image format: {e}")))?; @@ -104,9 +104,10 @@ impl PhotonImage { #[napi(js_name = "encode")] pub async fn encode(&self, format: u8, quality: u8) -> Result { let img = Arc::clone(&self.img); - let buffer: Vec = launch_blocking(move || encode_image(&img, format, quality)) - .wait() - .await?; + let buffer: Vec = + launch_blocking("image.encode", move || encode_image(&img, format, quality)) + .wait() + .await?; Ok(Uint8Array::from(buffer)) } @@ -115,9 +116,11 @@ impl PhotonImage { #[napi(js_name = "resize")] pub async fn resize(&self, width: u32, height: u32, filter: SamplingFilter) -> Result { let img = Arc::clone(&self.img); - let resized = launch_blocking(move || Ok(img.resize_exact(width, height, filter.into()))) - .wait() - .await?; + let resized = launch_blocking("image.resize", move || { + Ok(img.resize_exact(width, height, filter.into())) + }) + .wait() + .await?; Ok(Self { img: Arc::new(resized) }) } } diff --git a/crates/pi-natives/src/keys.rs b/crates/pi-natives/src/keys.rs index d05ac3883..4f94c6e96 100644 --- a/crates/pi-natives/src/keys.rs +++ b/crates/pi-natives/src/keys.rs @@ -685,7 +685,7 @@ fn matches_key_inner(bytes: &[u8], key_id: &str, kitty_protocol_active: bool) -> // ctrl+symbol legacy mapping (layout dependent) if let Some(legacy_ctrl) = ctrl_symbol_to_byte(ch) - && bytes == &[legacy_ctrl] + && bytes == [legacy_ctrl] { return true; } diff --git a/crates/pi-natives/src/shell.rs b/crates/pi-natives/src/shell.rs index 3623a8ec3..bce69e6bf 100644 --- a/crates/pi-natives/src/shell.rs +++ b/crates/pi-natives/src/shell.rs @@ -554,7 +554,7 @@ async fn run_shell_command( } } - let reader_handle = launch_async(async move { + let reader_handle = launch_async("shell.read_output", async move { read_output(reader_file, on_chunk).await; Ok(()) }); diff --git a/crates/pi-natives/src/work.rs b/crates/pi-natives/src/work.rs index 30a4fae82..8a469c269 100644 --- a/crates/pi-natives/src/work.rs +++ b/crates/pi-natives/src/work.rs @@ -4,32 +4,30 @@ //! Runs CPU-bound or blocking Rust work on a shared Rayon thread pool instead //! of Tokio's limited blocking workers. //! -//! # Example -//! ```ignore -//! use pi_natives::work::launch_task; -//! -//! # async fn demo() -> napi::Result<()> { -//! let handle = launch_task(|| Ok(42)); -//! let value = handle.wait().await?; -//! assert_eq!(value, 42); -//! # Ok(()) -//! # } -//! ``` -//! -//! # Architecture -//! ```text -//! JS async -> N-API -> launch_task -> Rayon thread pool -//! ``` +//! # Profiling +//! Samples are always collected into a circular buffer. Call +//! `get_work_profile()` to retrieve the last N seconds of data. use std::{ + cell::RefCell, + cmp::Reverse, + collections::HashMap, panic::{AssertUnwindSafe, catch_unwind}, sync::LazyLock, + time::Instant, }; use napi::{Error, Result}; +use napi_derive::napi; +use parking_lot::Mutex; use rayon::{ThreadPool, ThreadPoolBuilder}; +use smallvec::{SmallVec, smallvec}; use tokio::{sync::oneshot, task::JoinHandle}; +// ───────────────────────────────────────────────────────────────────────────── +// Work Handle +// ───────────────────────────────────────────────────────────────────────────── + /// Handle for a scheduled blocking task. pub enum WorkHandle { Blocking(oneshot::Receiver>), @@ -38,9 +36,6 @@ pub enum WorkHandle { impl WorkHandle { /// Await completion of the scheduled work. - /// - /// # Errors - /// Returns an error if the task panics or the channel is cancelled. pub async fn wait(self) -> Result { match self { Self::Blocking(receiver) => match receiver.await { @@ -63,34 +58,303 @@ impl WorkHandle { } } -/// Schedule blocking work on the shared Rayon pool. +// ───────────────────────────────────────────────────────────────────────────── +// Work Profiler - Always-on circular buffer +// ───────────────────────────────────────────────────────────────────────────── + +/// Maximum samples to keep (roughly 60s at high activity). +const MAX_SAMPLES: usize = 10_000; + +/// Process start time for relative timestamps. +static PROCESS_START: LazyLock = LazyLock::new(Instant::now); + +/// Circular buffer of profiling samples. +static PROFILE_BUFFER: LazyLock> = + LazyLock::new(|| Mutex::new(CircularBuffer::new(MAX_SAMPLES))); + +thread_local! { + /// Thread-local stack of active regions. + static REGION_STACK: RefCell> = const { RefCell::new(Vec::new()) }; +} + +/// A single profiling sample with timing data. +#[derive(Clone)] +struct ProfileSample { + /// Stack of region names (from root to leaf). + stack: SmallVec<[&'static str; 2]>, + /// Duration in microseconds. + duration_us: u64, + /// Timestamp (microseconds since process start). + timestamp_us: u64, +} + +/// Circular buffer for samples. +struct CircularBuffer { + samples: Vec, + capacity: usize, + write_pos: usize, + count: usize, +} + +impl CircularBuffer { + fn new(capacity: usize) -> Self { + Self { samples: Vec::with_capacity(capacity), capacity, write_pos: 0, count: 0 } + } + + fn push(&mut self, sample: ProfileSample) { + if self.samples.len() < self.capacity { + self.samples.push(sample); + } else { + self.samples[self.write_pos] = sample; + } + self.write_pos = (self.write_pos + 1) % self.capacity; + self.count = self.count.saturating_add(1); + } + + fn get_since(&self, cutoff_us: u64) -> Vec { + self + .samples + .iter() + .filter(|s| s.timestamp_us >= cutoff_us) + .cloned() + .collect() + } +} + +/// RAII guard that records timing when dropped. +pub struct ProfileGuard { + region: &'static str, + start: Instant, +} + +impl ProfileGuard { + #[inline] + fn new(region: &'static str) -> Self { + REGION_STACK.with(|stack| stack.borrow_mut().push(region)); + Self { region, start: Instant::now() } + } +} + +impl Drop for ProfileGuard { + fn drop(&mut self) { + let duration = self.start.elapsed(); + let duration_us = duration.as_micros() as u64; + let timestamp_us = PROCESS_START.elapsed().as_micros() as u64; + + REGION_STACK.with(|stack| { + let mut stack = stack.borrow_mut(); + let sample = ProfileSample { stack: stack.clone().into(), duration_us, timestamp_us }; + + if stack.last() == Some(&self.region) { + stack.pop(); + } + + PROFILE_BUFFER.lock().push(sample); + }); + } +} + +/// Start a profiling region. Returns a guard that records timing on drop. +#[inline] +pub fn profile_region(region: &'static str) -> ProfileGuard { + ProfileGuard::new(region) +} + +// ───────────────────────────────────────────────────────────────────────────── +// Work Profile Results +// ───────────────────────────────────────────────────────────────────────────── + +/// Profiling results returned to JavaScript. +#[napi(object)] +#[derive(Clone)] +pub struct WorkProfile { + /// Folded stack format for flamegraph tools. + pub folded: String, + /// Markdown summary of profiling results. + pub summary: String, + /// SVG flamegraph (if generation succeeded). + pub svg: Option, + /// Total profiled duration in milliseconds. + pub total_ms: f64, + /// Number of samples collected. + pub sample_count: u32, +} + +fn generate_folded(samples: &[ProfileSample]) -> String { + let mut aggregated: HashMap = HashMap::new(); + + for sample in samples { + if sample.stack.is_empty() { + continue; + } + let key = sample.stack.join(";"); + *aggregated.entry(key).or_insert(0) += sample.duration_us; + } + + let mut sorted: Vec<_> = aggregated.into_iter().collect(); + sorted.sort_by_key(|x| Reverse(x.1)); + + let mut output = String::new(); + for (stack, count) in sorted { + output.push_str(&stack); + output.push(' '); + output.push_str(&count.to_string()); + output.push('\n'); + } + + output +} + +fn generate_summary(samples: &[ProfileSample], window_ms: f64) -> String { + let mut by_region: HashMap<&'static str, (u64, usize)> = HashMap::new(); + + for sample in samples { + if let Some(®ion) = sample.stack.last() { + let entry = by_region.entry(region).or_insert((0, 0)); + entry.0 += sample.duration_us; + entry.1 += 1; + } + } + + let mut sorted: Vec<_> = by_region.into_iter().collect(); + sorted.sort_by_key(|x| Reverse((x.1).0)); + + let total_us: u64 = sorted.iter().map(|(_, (us, _))| us).sum(); + let total_ms = total_us as f64 / 1000.0; + + let mut lines = vec![ + "# Work Profile Summary".to_string(), + String::new(), + format!("Window: {window_ms:.1}ms"), + format!("Total work time: {total_ms:.1}ms"), + format!("Samples: {}", samples.len()), + String::new(), + "## Time by Region".to_string(), + String::new(), + "| Region | Time (ms) | % | Calls |".to_string(), + "|--------|-----------|---|-------|".to_string(), + ]; + + for (region, (time_us, count)) in sorted { + let time_ms = time_us as f64 / 1000.0; + let pct = if total_us > 0 { + (time_us as f64 / total_us as f64) * 100.0 + } else { + 0.0 + }; + lines.push(format!("| {region} | {time_ms:.2} | {pct:.1}% | {count} |")); + } + + lines.join("\n") +} + +fn generate_svg(folded: &str) -> Option { + use inferno::flamegraph::{self, Options}; + + let mut options = Options::default(); + options.title = "Work Profile".to_string(); + options.count_name = "μs".to_string(); + options.min_width = 0.1; + + let mut svg_output = Vec::new(); + let reader = std::io::Cursor::new(folded.as_bytes()); + + match flamegraph::from_reader(&mut options, reader, &mut svg_output) { + Ok(()) => String::from_utf8(svg_output).ok(), + Err(_) => None, + } +} + +// ───────────────────────────────────────────────────────────────────────────── +// N-API Exports +// ───────────────────────────────────────────────────────────────────────────── + +/// Get work profile data from the last N seconds. /// -/// # Errors -/// The returned handle resolves to an error if the task panics or is cancelled. -pub fn launch_blocking(work: F) -> WorkHandle +/// Always-on profiling - no need to start/stop. Just call this to get +/// recent activity. +#[napi] +pub fn get_work_profile(last_seconds: f64) -> WorkProfile { + let window_us = (last_seconds * 1_000_000.0) as u64; + let now_us = PROCESS_START.elapsed().as_micros() as u64; + let cutoff_us = now_us.saturating_sub(window_us); + + let samples = PROFILE_BUFFER.lock().get_since(cutoff_us); + + let folded = generate_folded(&samples); + let summary = generate_summary(&samples, last_seconds * 1000.0); + let svg = if folded.is_empty() { + None + } else { + generate_svg(&folded) + }; + + let total_us: u64 = samples.iter().map(|s| s.duration_us).sum(); + + WorkProfile { + folded, + summary, + svg, + total_ms: total_us as f64 / 1000.0, + sample_count: samples.len() as u32, + } +} + +// ───────────────────────────────────────────────────────────────────────────── +// Work Scheduling +// ───────────────────────────────────────────────────────────────────────────── + +/// Schedule blocking work on the shared Rayon pool with a profiling tag. +pub fn launch_blocking(tag: &'static str, work: F) -> WorkHandle where F: FnOnce() -> Result + Send + 'static, T: Send + 'static, { let (sender, receiver) = oneshot::channel(); + let submit_time = Instant::now(); + POOL.spawn(move || { + // Record queue wait time + let wait_us = submit_time.elapsed().as_micros() as u64; + let timestamp_us = PROCESS_START.elapsed().as_micros() as u64; + PROFILE_BUFFER.lock().push(ProfileSample { + stack: smallvec![tag, "queue_wait"], + duration_us: wait_us, + timestamp_us, + }); + + // Execute with profiling + let guard = profile_region(tag); let result = catch_unwind(AssertUnwindSafe(work)) .unwrap_or_else(|_| Err(Error::from_reason("Rayon task panicked"))); + drop(guard); + let _ = sender.send(result); }); + WorkHandle::Blocking(receiver) } -/// Schedule non-blocking async work on the Tokio runtime. -/// -/// # Errors -/// The returned handle resolves to an error if the task panics or is cancelled. -pub fn launch_async(work: Fut) -> WorkHandle +/// Schedule non-blocking async work on the Tokio runtime with a profiling tag. +pub fn launch_async(tag: &'static str, work: Fut) -> WorkHandle where Fut: Future> + Send + 'static, T: Send + 'static, { - WorkHandle::Async(tokio::spawn(work)) + WorkHandle::Async(tokio::spawn(async move { + let start = Instant::now(); + let result = work.await; + let duration_us = start.elapsed().as_micros() as u64; + let timestamp_us = PROCESS_START.elapsed().as_micros() as u64; + + PROFILE_BUFFER.lock().push(ProfileSample { + stack: smallvec![tag], + duration_us, + timestamp_us, + }); + + result + })) } static POOL: LazyLock = LazyLock::new(|| { diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 5ed66653f..a21b6eda7 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -1,6 +1,10 @@ # Changelog ## [Unreleased] +### Added + +- Added work scheduling profiler to debug menu for analyzing CPU scheduling patterns over the last 30 seconds +- Added support for work profile data in report bundles including folded stacks, summary, and flamegraph visualization ## [10.0.0] - 2026-02-01 ### Added diff --git a/packages/coding-agent/src/debug/index.ts b/packages/coding-agent/src/debug/index.ts index e91c462ea..c2307323c 100644 --- a/packages/coding-agent/src/debug/index.ts +++ b/packages/coding-agent/src/debug/index.ts @@ -4,6 +4,7 @@ * Provides tools for debugging, bug report generation, and system diagnostics. */ import * as fs from "node:fs/promises"; +import { getWorkProfile } from "@oh-my-pi/pi-natives/work"; import { Container, Loader, type SelectItem, SelectList, Spacer, Text } from "@oh-my-pi/pi-tui"; import { getSessionsDir } from "../config"; import { DynamicBorder } from "../modes/components/dynamic-border"; @@ -17,6 +18,7 @@ import { collectSystemInfo, formatSystemInfo } from "./system-info"; const DEBUG_MENU_ITEMS: SelectItem[] = [ { value: "open-artifacts", label: "Open: artifact folder", description: "Open session artifacts in file manager" }, { value: "performance", label: "Report: performance issue", description: "Profile CPU, reproduce, then bundle" }, + { value: "work", label: "Profile: work scheduling", description: "Open flamegraph of last 30s" }, { value: "dump", label: "Report: dump session", description: "Create report bundle immediately" }, { value: "memory", label: "Report: memory issue", description: "Heap snapshot + bundle" }, { value: "logs", label: "View: recent logs", description: "Show last 50 log entries" }, @@ -69,6 +71,9 @@ export class DebugSelectorComponent extends Container { case "performance": await this.handlePerformanceReport(); break; + case "work": + await this.handleWorkReport(); + break; case "dump": await this.handleDumpReport(); break; @@ -162,6 +167,39 @@ export class DebugSelectorComponent extends Container { this.ctx.ui.requestRender(); } + private async handleWorkReport(): Promise { + try { + const workProfile = getWorkProfile(30); + + if (!workProfile.svg) { + this.ctx.showWarning(`No work profile data (${workProfile.sampleCount} samples)`); + return; + } + + // Write SVG to temp file and open in browser + const tmpPath = `/tmp/work-profile-${Date.now()}.svg`; + await Bun.write(tmpPath, workProfile.svg); + + const openCmd = + process.platform === "darwin" + ? ["open", tmpPath] + : process.platform === "win32" + ? ["cmd", "/c", "start", "", tmpPath] + : ["xdg-open", tmpPath]; + + Bun.spawn(openCmd, { stdout: "ignore", stderr: "ignore" }).unref(); + + this.ctx.chatContainer.addChild(new Spacer(1)); + this.ctx.chatContainer.addChild( + new Text(theme.fg("dim", `Opened flamegraph (${workProfile.sampleCount} samples)`), 1, 0), + ); + } catch (err) { + this.ctx.showError(`Failed to open profile: ${err instanceof Error ? err.message : String(err)}`); + } + + this.ctx.ui.requestRender(); + } + private async handleDumpReport(): Promise { const loader = new Loader( this.ctx.ui, diff --git a/packages/coding-agent/src/debug/report-bundle.ts b/packages/coding-agent/src/debug/report-bundle.ts index 144214aab..072ea386f 100644 --- a/packages/coding-agent/src/debug/report-bundle.ts +++ b/packages/coding-agent/src/debug/report-bundle.ts @@ -6,6 +6,7 @@ import * as fs from "node:fs/promises"; import * as os from "node:os"; import * as path from "node:path"; +import type { WorkProfile } from "@oh-my-pi/pi-natives/work"; import { isEnoent } from "@oh-my-pi/pi-utils"; import type { CpuProfile, HeapSnapshot } from "./profiler"; import { collectSystemInfo, sanitizeEnv } from "./system-info"; @@ -42,6 +43,8 @@ export interface ReportBundleOptions { cpuProfile?: CpuProfile; /** Heap snapshot (for memory reports) */ heapSnapshot?: HeapSnapshot; + /** Work profile (for work scheduling reports) */ + workProfile?: WorkProfile; } export interface ReportBundleResult { @@ -63,6 +66,9 @@ export interface ReportBundleResult { * - profile.cpuprofile: CPU profile (performance report only) * - profile.md: Markdown CPU profile (performance report only) * - heap.heapsnapshot: Heap snapshot (memory report only) + * - work.folded: Work profile folded stacks (work report only) + * - work.md: Work profile summary (work report only) + * - work.svg: Work profile flamegraph (work report only) */ export async function createReportBundle(options: ReportBundleOptions): Promise { const reportsDir = getReportsDir(); @@ -131,6 +137,18 @@ export async function createReportBundle(options: ReportBundleOptions): Promise< files.push("heap.heapsnapshot"); } + // Work profile + if (options.workProfile) { + data["work.folded"] = options.workProfile.folded; + files.push("work.folded"); + data["work.md"] = options.workProfile.summary; + files.push("work.md"); + if (options.workProfile.svg) { + data["work.svg"] = options.workProfile.svg; + files.push("work.svg"); + } + } + // Write archive await Bun.Archive.write(outputPath, data, { compress: "gzip" }); diff --git a/packages/natives/CHANGELOG.md b/packages/natives/CHANGELOG.md index e5af551be..2827b9d3f 100644 --- a/packages/natives/CHANGELOG.md +++ b/packages/natives/CHANGELOG.md @@ -1,11 +1,17 @@ # Changelog ## [Unreleased] + ### Breaking Changes - Changed `executionId` parameter type from `string` to `number` in `abortShellExecution()` and `ShellExecuteOptions` - Removed `sessionKey` field from `ShellExecuteOptions` +### Added + +- Added `getWorkProfile()` function to retrieve work scheduling profiling data from a circular buffer of recent activity +- Added `WorkProfile` type with folded stack format, markdown summary, SVG flamegraph, and sample metrics for profiling results + ## [9.8.0] - 2026-02-01 ### Breaking Changes diff --git a/packages/natives/src/native.ts b/packages/natives/src/native.ts index e22bf8a87..b941cf15f 100644 --- a/packages/natives/src/native.ts +++ b/packages/natives/src/native.ts @@ -24,6 +24,7 @@ import "./ps/types"; import "./shell/types"; import "./system-info/types"; import "./text/types"; +import "./work/types"; export type { NativeBindings, TsFunc } from "./bindings"; @@ -184,6 +185,7 @@ function validateNative(bindings: NativeBindings, source: string): void { checkFn("killTree"); checkFn("listDescendants"); checkFn("getSystemInfo"); + checkFn("getWorkProfile"); if (missing.length) { throw new Error( diff --git a/packages/natives/src/work/index.ts b/packages/natives/src/work/index.ts new file mode 100644 index 000000000..242d1853c --- /dev/null +++ b/packages/natives/src/work/index.ts @@ -0,0 +1,11 @@ +/** + * Work scheduling profiling via native instrumentation. + * + * Always-on profiling - samples are collected into a circular buffer. + * Call `getWorkProfile()` to retrieve recent activity. + */ + +import { native } from "../native"; + +export type { WorkProfile } from "./types"; +export const { getWorkProfile } = native; diff --git a/packages/natives/src/work/types.ts b/packages/natives/src/work/types.ts new file mode 100644 index 000000000..9161543d5 --- /dev/null +++ b/packages/natives/src/work/types.ts @@ -0,0 +1,31 @@ +/** + * Types for work scheduling profiling. + */ + +/** + * Profiling results from work scheduling instrumentation. + */ +export interface WorkProfile { + /** Folded stack format for flamegraph tools. */ + folded: string; + /** Markdown summary of profiling results. */ + summary: string; + /** SVG flamegraph (if generation succeeded). */ + svg: string | null; + /** Total work time in milliseconds. */ + totalMs: number; + /** Number of samples collected. */ + sampleCount: number; +} + +declare module "../bindings" { + interface NativeBindings { + /** + * Get work profile data from the last N seconds. + * + * Always-on profiling - samples are collected into a circular buffer. + * Call this to retrieve recent activity. + */ + getWorkProfile(lastSeconds: number): WorkProfile; + } +}