From 6e122893393b0417260be0e07f55ba71d8a1acd5 Mon Sep 17 00:00:00 2001 From: houseme Date: Sat, 16 May 2026 19:09:04 +0800 Subject: [PATCH] fix(runtime): finalize issue 2941 profiling cleanup (#2983) * perf(runtime): narrow profiling support and upgrade starshard * style(notify): normalize starshard imports * perf(ecstore): reduce list_path_raw coordination overhead * docs(scripts): add issue 2941 perf capture workflow * fix(runtime): finalize issue 2941 profiling cleanup * build(deps): bump quick-xml to 0.40.0 * chore(scripts): untrack local perf capture guide * fix(scripts): honor label in perf capture output --- Cargo.lock | 40 +- Cargo.toml | 4 +- build-rustfs.sh | 2 +- .../ecstore/src/cache_value/metacache_set.rs | 31 +- crates/notify/src/rule_engine.rs | 4 +- crates/notify/src/rules/subscriber_index.rs | 4 +- rustfs/Cargo.toml | 8 +- rustfs/src/admin/handlers/profile.rs | 57 +- rustfs/src/embedded.rs | 3 +- rustfs/src/main.rs | 15 +- rustfs/src/profiling.rs | 153 +---- rustfs/src/profiling/allocator.rs | 529 ------------------ scripts/run_issue_2941_perf_capture.sh | 290 ++++++++++ 13 files changed, 416 insertions(+), 724 deletions(-) delete mode 100644 rustfs/src/profiling/allocator.rs create mode 100755 scripts/run_issue_2941_perf_capture.sh diff --git a/Cargo.lock b/Cargo.lock index d212c3a4b..9d3171ee9 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -761,9 +761,9 @@ dependencies = [ [[package]] name = "async-rs" -version = "0.8.6" +version = "0.8.7" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5f1cd61fd1b13805da591787d1af1b78e1650b2677f6365a72f00e55bebf0976" +checksum = "53bf71bee8a75907b6e3c81c5476efa7fcbb34df6e12d30b706888abded72091" dependencies = [ "async-compat", "async-global-executor", @@ -3532,9 +3532,9 @@ dependencies = [ [[package]] name = "dial9-tokio-telemetry" -version = "0.3.8" +version = "0.3.9" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4186df8377ec0b72938f205610b7dc11ecb39901a4cfdf6aeb23517e796e74e3" +checksum = "5a9bc5e50c103d4916a4a1b95784d9cc515d96bb2d093c8ea37c2cc0d85d92dd" dependencies = [ "arc-swap", "bon", @@ -4745,9 +4745,6 @@ dependencies = [ "allocator-api2", "equivalent", "foldhash 0.2.0", - "rayon", - "serde", - "serde_core", ] [[package]] @@ -6856,7 +6853,7 @@ version = "5.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "51e219e79014df21a225b1860a479e2dcd7cbd9130f4defd4bd0e191ea31d67d" dependencies = [ - "base64 0.21.7", + "base64 0.22.1", "chrono", "getrandom 0.2.17", "http 1.4.0", @@ -8366,6 +8363,15 @@ name = "quick-xml" version = "0.39.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "cdcc8dd4e2f670d309a5f0e83fe36dfdc05af317008fea29144da1a2ac858e5e" +dependencies = [ + "memchr", +] + +[[package]] +name = "quick-xml" +version = "0.40.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2474bd2e5029e7ccb6abb2ba48cf2383a333851dedf495901544281590c7da7f" dependencies = [ "encoding_rs", "memchr", @@ -9280,7 +9286,6 @@ dependencies = [ "aws-config", "aws-sdk-s3", "axum", - "backtrace", "base64 0.22.1", "base64-simd", "bytes", @@ -9359,7 +9364,6 @@ dependencies = [ "sha2 0.11.0", "shadow-rs", "socket2", - "starshard", "subtle", "sysinfo", "temp-env", @@ -9545,7 +9549,7 @@ dependencies = [ "parking_lot 0.12.5", "path-absolutize", "pin-project-lite", - "quick-xml 0.39.4", + "quick-xml 0.40.1", "rand 0.10.1", "ratelimit", "reed-solomon-erasure", @@ -9857,7 +9861,7 @@ dependencies = [ "hashbrown 0.17.1", "metrics", "percent-encoding", - "quick-xml 0.39.4", + "quick-xml 0.40.1", "rayon", "rustc-hash", "rustfs-config", @@ -9993,7 +9997,7 @@ dependencies = [ "md5 0.8.0", "percent-encoding", "proptest", - "quick-xml 0.39.4", + "quick-xml 0.40.1", "rcgen", "regex", "russh", @@ -11358,12 +11362,12 @@ dependencies = [ [[package]] name = "starshard" -version = "1.1.0" +version = "2.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6b3a2034ea62d2981c3bdeb21002f07707952ff3bd4594aa39f86ae38ea27dc6" +checksum = "472e5a677707be0fbe7bd851ca25ed7e7757e83dc07a05447cdb6fbdd27db5e2" dependencies = [ "async-trait", - "hashbrown 0.16.1", + "hashbrown 0.17.1", "rayon", "rustc-hash", "serde", @@ -13097,9 +13101,9 @@ checksum = "d6bbff5f0aada427a1e5a6da5f1f98158182f26556f345ac9e04d36d0ebed650" [[package]] name = "winnow" -version = "1.0.2" +version = "1.0.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2ee1708bef14716a11bae175f579062d4554d95be2c6829f518df847b7b3fdd0" +checksum = "0592e1c9d151f854e6fd382574c3a0855250e1d9b2f99d9281c6e6391af352f1" dependencies = [ "memchr", ] diff --git a/Cargo.toml b/Cargo.toml index c0a47b14a..5cd3eb4be 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -154,7 +154,7 @@ byteorder = "1.5.0" flatbuffers = "25.12.19" form_urlencoded = "1.2.2" prost = "0.14.3" -quick-xml = "0.39.4" +quick-xml = "0.40.0" rmp = { version = "0.8.15" } rmp-serde = { version = "1.3.1" } serde = { version = "1.0.228", features = ["derive"] } @@ -267,7 +267,7 @@ smallvec = { version = "1.15.1", features = ["serde"] } smartstring = "1.0.1" snafu = "0.9.0" snap = "1.1.1" -starshard = { version = "1.1.0", features = ["rayon", "async", "serde"] } +starshard = { version = "2.2.0", features = ["rayon", "async", "serde"] } strum = { version = "0.28.0", features = ["derive"] } sysinfo = "0.39.1" temp-env = "0.3.6" diff --git a/build-rustfs.sh b/build-rustfs.sh index dfb33b021..ea6d06b3a 100755 --- a/build-rustfs.sh +++ b/build-rustfs.sh @@ -217,7 +217,7 @@ setup_rust_environment() { # Set up environment variables for musl targets if [[ "$PLATFORM" == *"musl"* ]]; then print_message $YELLOW "Setting up environment for musl target..." - export RUSTFLAGS="'--cfg tokio_unstable -C target-feature=-crt-static'" + export RUSTFLAGS="--cfg tokio_unstable -C target-feature=-crt-static" # For cargo-zigbuild, set up additional environment variables if command -v cargo-zigbuild &> /dev/null; then diff --git a/crates/ecstore/src/cache_value/metacache_set.rs b/crates/ecstore/src/cache_value/metacache_set.rs index d69a8b596..10561443f 100644 --- a/crates/ecstore/src/cache_value/metacache_set.rs +++ b/crates/ecstore/src/cache_value/metacache_set.rs @@ -19,9 +19,10 @@ use futures::future::join_all; use metrics::counter; use rustfs_filemeta::{MetaCacheEntries, MetaCacheEntry, MetacacheReader, is_io_eof}; use std::{ + collections::VecDeque, future::Future, pin::Pin, - sync::{Arc, Mutex}, + sync::{Arc, OnceLock}, time::Duration, }; use tokio::io::AsyncRead; @@ -112,9 +113,9 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d let mut jobs: Vec>> = Vec::new(); let mut readers = Vec::with_capacity(opts.disks.len()); - let fds = opts.fallback_disks.iter().flatten().cloned().collect::>(); + let fds = opts.fallback_disks.iter().flatten().cloned().collect::>(); let max_disk_failures = opts.disks.len().saturating_sub(opts.min_disks); - let producer_errs = Arc::new(Mutex::new(vec![None; opts.disks.len()])); + let producer_errs: Arc<[OnceLock]> = (0..opts.disks.len()).map(|_| OnceLock::new()).collect::>().into(); let cancel_rx = CancellationToken::new(); @@ -137,14 +138,14 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d return Ok(()); } TestReaderBehavior::ProducerError(err) => { - producer_errs_clone.lock().expect("producer error mutex poisoned")[disk_idx] = Some(err.clone()); + record_producer_error(&producer_errs_clone, disk_idx, &err); return Err(err); } TestReaderBehavior::PartialThenTimeout(entries) => { let mut wr = wr; let mut out = rustfs_filemeta::MetacacheWriter::new(&mut wr); let err = DiskError::Timeout; - producer_errs_clone.lock().expect("producer error mutex poisoned")[disk_idx] = Some(err.clone()); + record_producer_error(&producer_errs_clone, disk_idx, &err); let _ = out.write(&entries).await; drop(out); return Err(err); @@ -187,8 +188,7 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d while need_fallback { let mut disk_op = None; - while !fds_clone.is_empty() { - let disk = fds_clone.remove(0); + while let Some(disk) = fds_clone.pop_front() { if disk.is_online().await { disk_op = Some(disk); break; @@ -198,7 +198,7 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d let Some(disk) = disk_op else { warn!("list_path_raw: fallback disk is none"); let err = last_err.unwrap_or(DiskError::DiskNotFound); - producer_errs_clone.lock().expect("producer error mutex poisoned")[disk_idx] = Some(err.clone()); + record_producer_error(&producer_errs_clone, disk_idx, &err); return Err(err); }; @@ -281,7 +281,7 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d // info!("read entry disk: {}, name: {}", i, entry.name); entry } else { - if let Some(err) = producer_errs.lock().expect("producer error mutex poisoned")[i].clone() { + if let Some(err) = producer_error(&producer_errs, i) { has_err += 1; errs[i] = Some(err); continue; @@ -293,7 +293,7 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d } } PeekOutcome::Error(err) => { - if let Some(err) = producer_errs.lock().expect("producer error mutex poisoned")[i].clone() { + if let Some(err) = producer_error(&producer_errs, i) { has_err += 1; errs[i] = Some(err); continue; @@ -510,10 +510,21 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d Ok(()) } +#[inline] +fn record_producer_error(producer_errs: &[OnceLock], idx: usize, err: &DiskError) { + let _ = producer_errs[idx].set(err.clone()); +} + +#[inline] +fn producer_error(producer_errs: &[OnceLock], idx: usize) -> Option { + producer_errs[idx].get().cloned() +} + #[cfg(test)] mod tests { use super::*; use rustfs_filemeta::MetacacheWriter; + use std::sync::Mutex; #[tokio::test] async fn list_path_raw_empty_disks_returns_read_quorum() { diff --git a/crates/notify/src/rule_engine.rs b/crates/notify/src/rule_engine.rs index c2e05d719..42bcacf7a 100644 --- a/crates/notify/src/rule_engine.rs +++ b/crates/notify/src/rule_engine.rs @@ -16,7 +16,7 @@ use crate::rules::{RulesMap, TargetIdSet}; use percent_encoding::percent_decode_str; use rustfs_s3_common::EventName; use rustfs_targets::arn::TargetID; -use starshard::AsyncShardedHashMap; +use starshard::{AsyncShardedHashMap, DEFAULT_SHARDS, SnapshotMode}; use std::sync::Arc; use tracing::info; @@ -37,7 +37,7 @@ pub struct NotifyRuleEngine { impl NotifyRuleEngine { pub fn new() -> Self { Self { - bucket_rules_map: Arc::new(AsyncShardedHashMap::new(0)), + bucket_rules_map: Arc::new(AsyncShardedHashMap::with_snapshot_mode(DEFAULT_SHARDS, SnapshotMode::Cached)), } } diff --git a/crates/notify/src/rules/subscriber_index.rs b/crates/notify/src/rules/subscriber_index.rs index 1f45efba2..3b57bc8f3 100644 --- a/crates/notify/src/rules/subscriber_index.rs +++ b/crates/notify/src/rules/subscriber_index.rs @@ -15,7 +15,7 @@ use crate::rules::{BucketRulesSnapshot, BucketSnapshotRef, DynRulesContainer}; use arc_swap::ArcSwap; use rustfs_s3_common::EventName; -use starshard::ShardedHashMap; +use starshard::{DEFAULT_SHARDS, ShardedHashMap}; use std::fmt; use std::sync::Arc; @@ -46,7 +46,7 @@ impl SubscriberIndex { /// Returns a new instance of SubscriberIndex. pub fn new(empty_rules: Arc) -> Self { Self { - inner: ShardedHashMap::new(64), + inner: ShardedHashMap::new(DEFAULT_SHARDS), empty_rules, } } diff --git a/rustfs/Cargo.toml b/rustfs/Cargo.toml index b7484e205..7d7f5413d 100644 --- a/rustfs/Cargo.toml +++ b/rustfs/Cargo.toml @@ -176,13 +176,7 @@ libmimalloc-sys = { version = "0.1.47", features = ["extended"] } -# Only enable pprof-based profiling on non-Windows targets. -[target.'cfg(all(not(target_os = "windows"), not(all(target_os = "linux", target_env = "gnu", target_arch = "x86_64"))))'.dependencies] -starshard = { workspace = true } -backtrace = { workspace = true } -rand = { workspace = true } -pprof = { workspace = true } - +# Only enable pprof-based profiling on linux + gnu + x86_64. [target.'cfg(all(target_os = "linux", target_env = "gnu", target_arch = "x86_64"))'.dependencies] tikv-jemallocator = { workspace = true } tikv-jemalloc-ctl = { workspace = true } diff --git a/rustfs/src/admin/handlers/profile.rs b/rustfs/src/admin/handlers/profile.rs index 94c042d0f..242821b94 100644 --- a/rustfs/src/admin/handlers/profile.rs +++ b/rustfs/src/admin/handlers/profile.rs @@ -15,8 +15,11 @@ use crate::admin::{auth::validate_admin_request, router::Operation}; use crate::auth::{check_key_valid, get_session_token}; use crate::server::RemoteAddr; +#[cfg(all(target_os = "linux", target_env = "gnu", target_arch = "x86_64"))] +use http::HeaderMap; +use http::StatusCode; +#[cfg(all(target_os = "linux", target_env = "gnu", target_arch = "x86_64"))] use http::header::CONTENT_TYPE; -use http::{HeaderMap, StatusCode}; use matchit::Params; use rustfs_policy::policy::action::{Action, AdminAction}; use s3s::{Body, S3Request, S3Response, S3Result, s3_error}; @@ -49,14 +52,29 @@ impl Operation for TriggerProfileCPU { authorize_profile_request(&req).await?; info!("Triggering CPU profile dump via S3 request..."); - let dur = std::time::Duration::from_secs(60); - match crate::profiling::dump_cpu_pprof_for(dur).await { - Ok(path) => { - let mut header = HeaderMap::new(); - header.insert(CONTENT_TYPE, "text/html".parse().unwrap()); - Ok(S3Response::with_headers((StatusCode::OK, Body::from(path.display().to_string())), header)) + #[cfg(not(all(target_os = "linux", target_env = "gnu", target_arch = "x86_64")))] + { + return Ok(S3Response::new(( + StatusCode::NOT_IMPLEMENTED, + Body::from( + crate::profiling::dump_cpu_pprof_for(std::time::Duration::from_secs(0)) + .await + .unwrap_err(), + ), + ))); + } + + #[cfg(all(target_os = "linux", target_env = "gnu", target_arch = "x86_64"))] + { + let dur = std::time::Duration::from_secs(60); + match crate::profiling::dump_cpu_pprof_for(dur).await { + Ok(path) => { + let mut header = HeaderMap::new(); + header.insert(CONTENT_TYPE, "text/html".parse().unwrap()); + Ok(S3Response::with_headers((StatusCode::OK, Body::from(path.display().to_string())), header)) + } + Err(e) => Err(s3s::s3_error!(InternalError, "{}", format!("Failed to dump CPU profile: {e}"))), } - Err(e) => Err(s3s::s3_error!(InternalError, "{}", format!("Failed to dump CPU profile: {e}"))), } } } @@ -68,13 +86,24 @@ impl Operation for TriggerProfileMemory { authorize_profile_request(&req).await?; info!("Triggering Memory profile dump via S3 request..."); - match crate::profiling::dump_memory_pprof_now().await { - Ok(path) => { - let mut header = HeaderMap::new(); - header.insert(CONTENT_TYPE, "text/html".parse().unwrap()); - Ok(S3Response::with_headers((StatusCode::OK, Body::from(path.display().to_string())), header)) + #[cfg(not(all(target_os = "linux", target_env = "gnu", target_arch = "x86_64")))] + { + return Ok(S3Response::new(( + StatusCode::NOT_IMPLEMENTED, + Body::from(crate::profiling::dump_memory_pprof_now().await.unwrap_err()), + ))); + } + + #[cfg(all(target_os = "linux", target_env = "gnu", target_arch = "x86_64"))] + { + match crate::profiling::dump_memory_pprof_now().await { + Ok(path) => { + let mut header = HeaderMap::new(); + header.insert(CONTENT_TYPE, "text/html".parse().unwrap()); + Ok(S3Response::with_headers((StatusCode::OK, Body::from(path.display().to_string())), header)) + } + Err(e) => Err(s3s::s3_error!(InternalError, "{}", format!("Failed to dump Memory profile: {e}"))), } - Err(e) => Err(s3s::s3_error!(InternalError, "{}", format!("Failed to dump Memory profile: {e}"))), } } } diff --git a/rustfs/src/embedded.rs b/rustfs/src/embedded.rs index e845f29cb..914a52330 100644 --- a/rustfs/src/embedded.rs +++ b/rustfs/src/embedded.rs @@ -51,6 +51,7 @@ use crate::config::Config; use crate::init::{add_bucket_notification_configuration, init_buffer_profile_system, init_kms_system}; use crate::server::{init_event_notifier, shutdown_event_notifier, start_audit_system, start_http_server, stop_audit_system}; use rustfs_common::{GlobalReadiness, SystemStage, set_global_addr}; +use rustfs_config::ENV_RUSTFS_ALLOW_INSECURE_DEFAULT_CREDENTIALS; use rustfs_credentials::init_global_action_credentials; use rustfs_ecstore::store::init_lock_clients; use rustfs_ecstore::{ @@ -317,7 +318,7 @@ impl RustFSServerBuilder { )); } - let allow_insecure_defaults = get_env_bool(rustfs_config::ENV_RUSTFS_ALLOW_INSECURE_DEFAULT_CREDENTIALS, false); + let allow_insecure_defaults = get_env_bool(ENV_RUSTFS_ALLOW_INSECURE_DEFAULT_CREDENTIALS, false); if !config.default_credentials_allowed_for_addr(server_addr, allow_insecure_defaults) { return Err(ServerError::Init( "default root credentials are not allowed on non-loopback listeners; set access_key and secret_key to non-default values, bind to loopback, or set RUSTFS_ALLOW_INSECURE_DEFAULT_CREDENTIALS=true for local development only" diff --git a/rustfs/src/main.rs b/rustfs/src/main.rs index e1f4e2f52..3e188e1e2 100644 --- a/rustfs/src/main.rs +++ b/rustfs/src/main.rs @@ -34,6 +34,7 @@ use rustfs::server::{ start_audit_system, start_http_server, stop_audit_system, wait_for_shutdown, }; use rustfs_common::{GlobalReadiness, SystemStage, set_global_addr}; +use rustfs_config::ENV_RUSTFS_ALLOW_INSECURE_DEFAULT_CREDENTIALS; use rustfs_credentials::init_global_action_credentials; use rustfs_ecstore::store::init_lock_clients; use rustfs_ecstore::{ @@ -76,15 +77,7 @@ const ENV_HEAL_ENABLED_DEPRECATED: &str = "RUSTFS_ENABLE_HEAL"; #[global_allocator] static GLOBAL: tikv_jemallocator::Jemalloc = tikv_jemallocator::Jemalloc; -#[cfg(all( - not(target_os = "windows"), - not(all(target_os = "linux", target_env = "gnu", target_arch = "x86_64")) -))] -#[global_allocator] -static GLOBAL: rustfs::profiling::allocator::TracingAllocator = - rustfs::profiling::allocator::TracingAllocator::new(mimalloc::MiMalloc); - -#[cfg(target_os = "windows")] +#[cfg(not(all(target_os = "linux", target_env = "gnu", target_arch = "x86_64")))] #[global_allocator] static GLOBAL: mimalloc::MiMalloc = mimalloc::MiMalloc; @@ -145,7 +138,7 @@ const DEFAULT_CREDENTIALS_WARNING_MESSAGE: &str = "Detected default root credent const DEFAULT_CREDENTIALS_ERROR_MESSAGE: &str = "Default root credentials are not allowed on non-loopback listeners; set RUSTFS_ACCESS_KEY and RUSTFS_SECRET_KEY to non-default values, bind to loopback, or set RUSTFS_ALLOW_INSECURE_DEFAULT_CREDENTIALS=true for local development only"; fn allow_insecure_default_credentials() -> bool { - get_env_bool(rustfs_config::ENV_RUSTFS_ALLOW_INSECURE_DEFAULT_CREDENTIALS, false) + get_env_bool(ENV_RUSTFS_ALLOW_INSECURE_DEFAULT_CREDENTIALS, false) } async fn async_main() -> Result<()> { @@ -845,7 +838,7 @@ mod tests { for message in [DEFAULT_CREDENTIALS_WARNING_MESSAGE, DEFAULT_CREDENTIALS_ERROR_MESSAGE] { assert!(message.contains(rustfs_config::ENV_RUSTFS_ACCESS_KEY)); assert!(message.contains(rustfs_config::ENV_RUSTFS_SECRET_KEY)); - assert!(message.contains(rustfs_config::ENV_RUSTFS_ALLOW_INSECURE_DEFAULT_CREDENTIALS)); + assert!(message.contains(ENV_RUSTFS_ALLOW_INSECURE_DEFAULT_CREDENTIALS)); assert!(!message.contains(rustfs_credentials::DEFAULT_ACCESS_KEY)); assert!(!message.contains(rustfs_credentials::DEFAULT_SECRET_KEY)); } diff --git a/rustfs/src/profiling.rs b/rustfs/src/profiling.rs index 29688e33d..4770cdeea 100644 --- a/rustfs/src/profiling.rs +++ b/rustfs/src/profiling.rs @@ -12,154 +12,53 @@ // See the License for the specific language governing permissions and // limitations under the License. -#[cfg(all( - not(target_os = "windows"), - not(all(target_os = "linux", target_env = "gnu", target_arch = "x86_64")) -))] -pub mod allocator; - -#[cfg(target_os = "windows")] -mod windows_impl { +#[cfg(not(all(target_os = "linux", target_env = "gnu", target_arch = "x86_64")))] +mod unsupported_impl { use std::path::PathBuf; use std::time::Duration; use tracing::info; pub async fn init_from_env() { - info!("Profiling initialization skipped on Windows platform (not supported)"); + let target_env = option_env!("CARGO_CFG_TARGET_ENV").unwrap_or("unknown"); + info!( + target_os = std::env::consts::OS, + target_env, + target_arch = std::env::consts::ARCH, + "Profiling initialization skipped on unsupported platform" + ); } /// Stop all background profiling tasks pub fn shutdown_profiling() { - info!("profiling: shutdown called on Windows platform (no-op)"); + let target_env = option_env!("CARGO_CFG_TARGET_ENV").unwrap_or("unknown"); + info!( + target_os = std::env::consts::OS, + target_env, + target_arch = std::env::consts::ARCH, + "profiling: shutdown called on unsupported platform (no-op)" + ); } pub async fn dump_cpu_pprof_for(_duration: Duration) -> Result { - Err("CPU profiling is not supported on Windows platform".to_string()) + Err(unsupported_message("CPU profiling")) } pub async fn dump_memory_pprof_now() -> Result { - Err("Memory profiling is not supported on Windows platform".to_string()) + Err(unsupported_message("Memory profiling")) } -} -#[cfg(target_os = "windows")] -pub use windows_impl::{dump_cpu_pprof_for, dump_memory_pprof_now, init_from_env, shutdown_profiling}; -#[cfg(all( - not(target_os = "windows"), - not(all(target_os = "linux", target_env = "gnu", target_arch = "x86_64")) -))] -mod generic_impl { - use super::allocator; - use rustfs_config::{ - DEFAULT_ENABLE_PROFILING, DEFAULT_MEM_INTERVAL_SECS, DEFAULT_MEM_PERIODIC, DEFAULT_OUTPUT_DIR, ENV_ENABLE_PROFILING, - ENV_MEM_INTERVAL_SECS, ENV_MEM_PERIODIC, ENV_OUTPUT_DIR, - }; - use rustfs_utils::{get_env_bool, get_env_str, get_env_u64}; - use std::fs::create_dir_all; - use std::path::PathBuf; - use std::sync::OnceLock; - use std::time::Duration; - use tokio::time::sleep; - use tokio_util::sync::CancellationToken; - use tracing::{debug, error, info, warn}; - - // Global cancellation token for periodic profiling tasks - static PROFILING_CANCEL_TOKEN: OnceLock = OnceLock::new(); - - fn get_platform_info() -> (String, String, String) { - ( - std::env::consts::OS.to_string(), - option_env!("CARGO_CFG_TARGET_ENV").unwrap_or("unknown").to_string(), - std::env::consts::ARCH.to_string(), + fn unsupported_message(feature: &str) -> String { + let target_env = option_env!("CARGO_CFG_TARGET_ENV").unwrap_or("unknown"); + format!( + "{feature} is only supported on linux x86_64 gnu. target_os={}, target_env={target_env}, target_arch={}", + std::env::consts::OS, + std::env::consts::ARCH ) } - - fn output_dir() -> PathBuf { - let dir = get_env_str(ENV_OUTPUT_DIR, DEFAULT_OUTPUT_DIR); - let p = PathBuf::from(dir); - if let Err(e) = create_dir_all(&p) { - warn!("profiling: create output dir {} failed: {}, fallback to current dir", p.display(), e); - return PathBuf::from("."); - } - p - } - - fn ts() -> String { - jiff::Zoned::now().strftime("%Y%m%dT%H%M%S").to_string() - } - - pub async fn init_from_env() { - let enabled = get_env_bool(ENV_ENABLE_PROFILING, DEFAULT_ENABLE_PROFILING); - if !enabled { - debug!("profiling: disabled by env"); - return; - } - - allocator::set_enabled(true); - info!("profiling: Memory profiling enabled (mimalloc + tracing)"); - - // Initialize cancellation token - let token = PROFILING_CANCEL_TOKEN.get_or_init(CancellationToken::new).clone(); - - // Memory periodic dump - let mem_periodic = get_env_bool(ENV_MEM_PERIODIC, DEFAULT_MEM_PERIODIC); - let mem_interval = Duration::from_secs(get_env_u64(ENV_MEM_INTERVAL_SECS, DEFAULT_MEM_INTERVAL_SECS)); - if mem_periodic { - start_memory_periodic(mem_interval, token).await; - } - } - - async fn start_memory_periodic(interval: Duration, token: CancellationToken) { - info!(?interval, "start periodic memory pprof dump"); - tokio::spawn(async move { - loop { - tokio::select! { - _ = token.cancelled() => { - info!("periodic memory profiling task cancelled"); - break; - } - _ = sleep(interval) => { - let out = output_dir().join(format!("mem_profile_periodic_{}.pb", ts())); - match allocator::dump_profile(&out) { - Ok(_) => info!("periodic memory profile dumped to {}", out.display()), - Err(e) => error!("periodic mem dump failed: {}", e), - } - } - } - } - }); - } - - /// Stop all background profiling tasks - pub fn shutdown_profiling() { - if let Some(token) = PROFILING_CANCEL_TOKEN.get() { - token.cancel(); - } - allocator::set_enabled(false); - } - - pub async fn dump_cpu_pprof_for(_duration: Duration) -> Result { - let (target_os, target_env, target_arch) = get_platform_info(); - let msg = format!( - "CPU profiling is not supported on this platform. target_os={target_os}, target_env={target_env}, target_arch={target_arch}" - ); - Err(msg) - } - - pub async fn dump_memory_pprof_now() -> Result { - let out = output_dir().join(format!("mem_profile_{}.pb", ts())); - allocator::dump_profile(&out).map(|_| { - info!("Memory profile exported: {}", out.display()); - out - }) - } } -#[cfg(all( - not(target_os = "windows"), - not(all(target_os = "linux", target_env = "gnu", target_arch = "x86_64")) -))] -pub use generic_impl::{dump_cpu_pprof_for, dump_memory_pprof_now, init_from_env, shutdown_profiling}; +#[cfg(not(all(target_os = "linux", target_env = "gnu", target_arch = "x86_64")))] +pub use unsupported_impl::{dump_cpu_pprof_for, dump_memory_pprof_now, init_from_env, shutdown_profiling}; #[cfg(all(target_os = "linux", target_env = "gnu", target_arch = "x86_64"))] mod linux_impl { diff --git a/rustfs/src/profiling/allocator.rs b/rustfs/src/profiling/allocator.rs deleted file mode 100644 index edf4632ca..000000000 --- a/rustfs/src/profiling/allocator.rs +++ /dev/null @@ -1,529 +0,0 @@ -// Copyright 2024 RustFS Team -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. - -use backtrace::Backtrace; -use pprof::protos::Message; -use rand::RngExt; -use starshard::ShardedHashMap; -use std::alloc::{GlobalAlloc, Layout}; -use std::cell::Cell; -use std::collections::HashMap; -use std::fs::File; -use std::hash::{Hash, Hasher}; -use std::io::Write; -use std::path::Path; -use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; -use std::sync::{Arc, LazyLock, Weak}; - -type AllocatorShardedMap = ShardedHashMap>)>; -type LazyAllocatorShardedMap = LazyLock; - -type AllocatorSampleHashMap = HashMap<*const Vec, (i64, i64, Arc>)>; -/// A wrapper around a GlobalAlloc that samples allocations and records stack traces. -pub struct TracingAllocator { - inner: A, -} - -// Thread-local reentrancy guard to prevent infinite recursion when recording allocations -thread_local! { - static REENTRANCY_GUARD: Cell = const { Cell::new(false) }; -} - -// Global configuration -static SAMPLE_RATE: AtomicUsize = AtomicUsize::new(512 * 1024); // Default: sample every 512KB on average -static ENABLED: AtomicBool = AtomicBool::new(false); - -// Global storage for profile data -// Map: Address (usize) -> (Size (usize), StackTrace (Arc>)) -// We store the Arc to keep the stack trace alive as long as the allocation is live. -static LIVE_ALLOCATIONS: LazyAllocatorShardedMap = LazyLock::new(|| ShardedHashMap::new(64)); - -// Cache for deduplicating stack traces. -// Map: StackHash (u64) -> Weak> -// We use Weak references so that unused stack traces can be dropped when all referring allocations are freed. -static STACK_CACHE: LazyLock>>> = LazyLock::new(|| ShardedHashMap::new(64)); - -impl TracingAllocator { - pub const fn new(inner: A) -> Self { - Self { inner } - } -} - -// Public configuration functions -#[allow(dead_code)] -pub fn set_sample_rate(rate: usize) { - SAMPLE_RATE.store(rate, Ordering::Relaxed); -} - -pub fn set_enabled(enabled: bool) { - // Force initialization of LazyLocks before enabling profiling to avoid recursion during init. - // Accessing them is enough to trigger initialization. - let _ = &*LIVE_ALLOCATIONS; - let _ = &*STACK_CACHE; - - ENABLED.store(enabled, Ordering::Relaxed); -} - -fn should_sample(size: usize) -> bool { - if !ENABLED.load(Ordering::Relaxed) { - return false; - } - - let rate = SAMPLE_RATE.load(Ordering::Relaxed); - if rate == 0 { - return true; - } - - // Use a fresh RNG each time. - let mut rng = rand::rng(); - rng.random_range(0..rate) < size -} - -// Internal function, assumes guard is already held -fn record_alloc(ptr: *mut u8, size: usize) { - // Capture stack trace - let bt = Backtrace::new_unresolved(); - let mut frames = Vec::new(); - for frame in bt.frames() { - frames.push(frame.symbol_address() as usize); - } - - // Calculate hash of the stack trace - let mut hasher = std::collections::hash_map::DefaultHasher::new(); - frames.hash(&mut hasher); - let stack_hash = hasher.finish(); - - // Deduplicate stack trace using STACK_CACHE - let stack_arc = if let Some(weak) = STACK_CACHE.get(&stack_hash) { - if let Some(arc) = weak.upgrade() { - arc - } else { - // Entry exists but is dead, replace it - let arc = Arc::new(frames); - STACK_CACHE.insert(stack_hash, Arc::downgrade(&arc)); - arc - } - } else { - // New entry - let arc = Arc::new(frames); - STACK_CACHE.insert(stack_hash, Arc::downgrade(&arc)); - arc - }; - - // Store the allocation info with the Arc - LIVE_ALLOCATIONS.insert(ptr as usize, (size, stack_arc)); -} - -// Internal function, assumes guard is already held -fn record_dealloc(ptr: *mut u8) { - // Remove from live allocations. - // The Arc> will be dropped. - // If it was the last reference, the Vec is freed. - // The Weak pointer in STACK_CACHE remains but becomes upgrade-able to None. - LIVE_ALLOCATIONS.remove(&(ptr as usize)); -} - -/// Dump the current profile to a pprof protobuf file -pub fn dump_profile(path: &Path) -> Result<(), String> { - // Prevent reentrancy during dump - if REENTRANCY_GUARD.replace(true) { - return Err("Reentrancy detected during dump".to_string()); - } - - // Perform a lazy cleanup of the cache during dump - cleanup_cache(); - - let result = dump_profile_inner(path); - - REENTRANCY_GUARD.set(false); - result -} - -// Clean up dead entries from STACK_CACHE -fn cleanup_cache() { - // We collect dead keys first to avoid locking issues during iteration if any - let mut dead_keys = Vec::new(); - - // Note: This iteration might be slow if the cache is huge, but dump_profile is infrequent. - for entry in STACK_CACHE.iter() { - let (key, weak) = entry; - if weak.upgrade().is_none() { - dead_keys.push(key); - } - } - - for key in dead_keys { - STACK_CACHE.remove(&key); - } -} - -fn dump_profile_inner(path: &Path) -> Result<(), String> { - use pprof::protos as pb; - - let mut profile = pb::Profile::default(); - - // Basic metadata - profile.string_table.push("".to_string()); // 0: empty - profile.string_table.push("alloc_objects".to_string()); // 1 - profile.string_table.push("count".to_string()); // 2 - profile.string_table.push("alloc_space".to_string()); // 3 - profile.string_table.push("bytes".to_string()); // 4 - - let sample_type_count = pb::ValueType { - ty: 1, // "alloc_objects" - unit: 2, // "count" - ..Default::default() - }; - let sample_type_bytes = pb::ValueType { - ty: 3, // "alloc_space" - unit: 4, // "bytes" - ..Default::default() - }; - profile.sample_type = vec![sample_type_count, sample_type_bytes]; - - // Helper to get string ID - let mut string_map: HashMap = HashMap::new(); - string_map.insert("".to_string(), 0); - string_map.insert("alloc_objects".to_string(), 1); - string_map.insert("count".to_string(), 2); - string_map.insert("alloc_space".to_string(), 3); - string_map.insert("bytes".to_string(), 4); - - let mut get_string_id = |s: String| -> i64 { - if let Some(&id) = string_map.get(&s) { - id - } else { - let id = profile.string_table.len() as i64; - profile.string_table.push(s.clone()); - string_map.insert(s, id); - id - } - }; - - // Helper to get location ID - let mut location_map: HashMap = HashMap::new(); // addr -> loc_id - let mut function_map: HashMap = HashMap::new(); // addr -> func_id - - // Collect samples - // Aggregate by Stack Trace Pointer (deduplication via Arc pointer) - // Map: Arc pointer -> (Count, Bytes, Arc>) - let mut aggregated_samples: AllocatorSampleHashMap = HashMap::new(); - - // Step 1: Collect data from LIVE_ALLOCATIONS while holding the lock (implicitly via iter) - // We do NOT perform symbol resolution here to avoid deadlocks. - for entry in LIVE_ALLOCATIONS.iter() { - let (_ptr, (size, stack_arc)) = entry; - let stack_arc_clone = stack_arc.clone(); - let key = Arc::as_ptr(&stack_arc_clone); - - let agg = aggregated_samples.entry(key).or_insert_with(|| (0, 0, stack_arc_clone)); - agg.0 += 1; - agg.1 += size as i64; - } - // LIVE_ALLOCATIONS lock is released here as the iterator is dropped. - - // Step 2: Process samples and resolve symbols (outside of LIVE_ALLOCATIONS lock) - for (_key, (count, bytes, frames)) in aggregated_samples { - let mut sample = pb::Sample { - value: vec![count, bytes], - ..Default::default() - }; - - // Process frames - for &addr in frames.iter() { - let loc_id = if let Some(&id) = location_map.get(&addr) { - id - } else { - // Resolve symbol - // This might take time and locks, but we are safe now. - let mut func_name = "unknown".to_string(); - let mut file_name = "unknown".to_string(); - let mut line_no = 0; - - backtrace::resolve(addr as *mut std::ffi::c_void, |symbol| { - if let Some(name) = symbol.name() { - func_name = name.to_string(); - } - if let Some(filename) = symbol.filename() { - file_name = filename.to_string_lossy().to_string(); - } - if let Some(line) = symbol.lineno() { - line_no = line as i64; - } - }); - - // Create Function - let func_id = if let Some(&id) = function_map.get(&addr) { - id - } else { - let id = (profile.function.len() + 1) as u64; - let name_id = get_string_id(func_name); - let file_id = get_string_id(file_name); - - let func = pb::Function { - id, - name: name_id, - system_name: name_id, - filename: file_id, - start_line: 0, - ..Default::default() - }; - profile.function.push(func); - function_map.insert(addr, id); - id - }; - - // Create Location - let id = (profile.location.len() + 1) as u64; - let line = pb::Line { - function_id: func_id, - line: line_no, - ..Default::default() - }; - let loc = pb::Location { - id, - mapping_id: 0, - address: addr as u64, - line: vec![line], - is_folded: false, - ..Default::default() - }; - profile.location.push(loc); - location_map.insert(addr, id); - id - }; - sample.location_id.push(loc_id); - } - profile.sample.push(sample); - } - - // Write to file - let mut buf = Vec::with_capacity(1024 * 1024); - profile.write_to_vec(&mut buf).map_err(|e| format!("encode failed: {e}"))?; - - let mut f = File::create(path).map_err(|e| format!("create file failed: {e}"))?; - f.write_all(&buf).map_err(|e| format!("write file failed: {e}"))?; - - Ok(()) -} - -// Helper to handle sampling logic -#[inline(always)] -fn handle_alloc_sampling(ptr: *mut u8, size: usize) { - if !ptr.is_null() { - // Check reentrancy guard BEFORE calling should_sample - if !REENTRANCY_GUARD.replace(true) { - if should_sample(size) { - record_alloc(ptr, size); - } - REENTRANCY_GUARD.set(false); - } - } -} - -// Helper to handle dealloc logic -#[inline(always)] -fn handle_dealloc_sampling(ptr: *mut u8) { - if !REENTRANCY_GUARD.replace(true) { - record_dealloc(ptr); - REENTRANCY_GUARD.set(false); - } -} - -// SAFETY: This allocator wrapper preserves the `GlobalAlloc` contract by -// delegating all allocation operations to `inner` with the exact caller-provided -// layouts and pointers. Sampling records metadata only and does not take -// ownership of allocation memory. -#[allow(unsafe_code)] -unsafe impl GlobalAlloc for TracingAllocator { - unsafe fn alloc(&self, layout: Layout) -> *mut u8 { - // SAFETY: Delegating to inner allocator. - let ptr = unsafe { self.inner.alloc(layout) }; - handle_alloc_sampling(ptr, layout.size()); - ptr - } - - unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) { - handle_dealloc_sampling(ptr); - // SAFETY: Delegating to inner allocator. - unsafe { self.inner.dealloc(ptr, layout) }; - } - - unsafe fn alloc_zeroed(&self, layout: Layout) -> *mut u8 { - // SAFETY: Delegating to inner allocator. - let ptr = unsafe { self.inner.alloc_zeroed(layout) }; - handle_alloc_sampling(ptr, layout.size()); - ptr - } - - unsafe fn realloc(&self, ptr: *mut u8, layout: Layout, new_size: usize) -> *mut u8 { - handle_dealloc_sampling(ptr); - - // SAFETY: Delegating to inner allocator. - let new_ptr = unsafe { self.inner.realloc(ptr, layout, new_size) }; - - handle_alloc_sampling(new_ptr, new_size); - new_ptr - } -} - -#[cfg(test)] -// SAFETY: Tests call the unsafe `GlobalAlloc` methods with layouts created by -// `Layout::from_size_align` and deallocate each returned pointer with the same -// layout. -#[allow(unsafe_code)] -mod tests { - use super::*; - use serial_test::serial; - use std::alloc::System; - use std::thread; - use tempfile::NamedTempFile; - - // Use System allocator for testing - static TEST_ALLOCATOR: TracingAllocator = TracingAllocator::new(System); - - #[test] - #[serial] - fn test_basic_allocation_tracking() { - // Enable profiling and force sampling (rate = 1 means sample everything) - set_enabled(true); - set_sample_rate(1); - - unsafe { - let layout = Layout::from_size_align(1024, 8).unwrap(); - let ptr = TEST_ALLOCATOR.alloc(layout); - assert!(!ptr.is_null()); - - // Verify allocation is recorded - assert!(LIVE_ALLOCATIONS.get(&(ptr as usize)).is_some()); - - TEST_ALLOCATOR.dealloc(ptr, layout); - - // Verify allocation is removed - assert!(LIVE_ALLOCATIONS.get(&(ptr as usize)).is_none()); - } - - // Reset - set_enabled(false); - } - - #[test] - #[serial] - fn test_reentrancy_guard() { - set_enabled(true); - set_sample_rate(1); - - // Manually set guard to simulate reentrancy - REENTRANCY_GUARD.set(true); - - unsafe { - let layout = Layout::from_size_align(128, 8).unwrap(); - let ptr = TEST_ALLOCATOR.alloc(layout); - - // Should NOT be recorded because guard was true - assert!(LIVE_ALLOCATIONS.get(&(ptr as usize)).is_none()); - - TEST_ALLOCATOR.dealloc(ptr, layout); - } - - REENTRANCY_GUARD.set(false); - set_enabled(false); - } - - #[test] - #[serial] - fn test_sampling_logic() { - set_enabled(true); - // Set a high rate so small allocations are unlikely to be sampled - set_sample_rate(1_000_000); - - let mut sampled_count = 0; - let iterations = 100; - - unsafe { - let layout = Layout::from_size_align(8, 8).unwrap(); - for _ in 0..iterations { - let ptr = TEST_ALLOCATOR.alloc(layout); - if LIVE_ALLOCATIONS.get(&(ptr as usize)).is_some() { - sampled_count += 1; - } - TEST_ALLOCATOR.dealloc(ptr, layout); - } - } - - // With high sample rate and small size, sampled count should be low (likely 0) - // This is probabilistic, but 0 is very likely. - assert!(sampled_count < iterations); - - set_enabled(false); - } - - #[test] - #[serial] - fn test_profile_dump() { - set_enabled(true); - // Use a larger sample rate to avoid capturing too much noise from the test runner - // and ensure we only capture our large allocation. - set_sample_rate(1024 * 1024); - - unsafe { - // Allocate a large enough chunk to likely be sampled (2MB > 1MB rate) - let layout = Layout::from_size_align(2 * 1024 * 1024, 8).unwrap(); - let ptr = TEST_ALLOCATOR.alloc(layout); - - let file = NamedTempFile::new().unwrap(); - let path = file.path(); - - let result = dump_profile(path); - assert!(result.is_ok()); - - let metadata = std::fs::metadata(path).unwrap(); - assert!(metadata.len() > 0); - - TEST_ALLOCATOR.dealloc(ptr, layout); - } - set_enabled(false); - } - - #[test] - #[serial] - fn test_concurrent_allocations() { - set_enabled(true); - set_sample_rate(1); - - let threads: Vec<_> = (0..10) - .map(|_| { - thread::spawn(|| { - unsafe { - let layout = Layout::from_size_align(64, 8).unwrap(); - for _ in 0..100 { - let ptr = TEST_ALLOCATOR.alloc(layout); - // Just ensure no panic/crash - TEST_ALLOCATOR.dealloc(ptr, layout); - } - } - }) - }) - .collect(); - - for t in threads { - t.join().unwrap(); - } - - // After all threads join and dealloc, map should be empty (ignoring other potential allocations in test runner) - // Note: In a real test runner, other tests might be running, so we can't assert empty. - // But we verified no crashes. - set_enabled(false); - } -} diff --git a/scripts/run_issue_2941_perf_capture.sh b/scripts/run_issue_2941_perf_capture.sh new file mode 100755 index 000000000..e91f3886c --- /dev/null +++ b/scripts/run_issue_2941_perf_capture.sh @@ -0,0 +1,290 @@ +#!/usr/bin/env bash +set -euo pipefail + +PROJECT_ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +LABEL="${LABEL:-issue-2941}" +DURATION_SECS="${DURATION_SECS:-60}" +PERF_FREQ="${PERF_FREQ:-99}" +OUT_DIR="${OUT_DIR:-}" +RUSTFS_PID="${RUSTFS_PID:-}" +CONTAINER_NAME="${CONTAINER_NAME:-}" +ENDPOINT="${ENDPOINT:-http://127.0.0.1:9000}" +PERF_MODE="${PERF_MODE:-auto}" # auto|on|off +SUDO_CMD="${SUDO_CMD:-}" # example: sudo + +usage() { + cat <<'USAGE' +Usage: + scripts/run_issue_2941_perf_capture.sh [options] + +Options: + --label artifact label prefix + --duration sample duration in seconds (default: 60) + --out-dir artifact output directory + --pid rustfs pid; auto-detect if omitted + --container docker container name/id for extra stats + --endpoint rustfs endpoint for health probes (default: http://127.0.0.1:9000) + --perf whether to run perf record (default: auto) + --perf-freq perf sample frequency (default: 99) + --sudo-cmd optional prefix for privileged perf, e.g. "sudo" + -h, --help show help + +Environment: + LABEL + DURATION_SECS + PERF_FREQ + OUT_DIR + RUSTFS_PID + CONTAINER_NAME + ENDPOINT + PERF_MODE + SUDO_CMD + +Examples: + scripts/run_issue_2941_perf_capture.sh --label musl-baseline --container rustfs + scripts/run_issue_2941_perf_capture.sh --label glibc-test --pid 12345 --perf on --sudo-cmd sudo +USAGE +} + +log() { + printf '[INFO] %s\n' "$*" +} + +warn() { + printf '[WARN] %s\n' "$*" >&2 +} + +require_arg() { + local option="$1" + local value="${2-}" + if [[ $# -lt 2 || -z "${value}" || "${value}" == --* ]]; then + warn "missing value for ${option}" + usage + exit 1 + fi +} + +parse_args() { + while [[ $# -gt 0 ]]; do + case "$1" in + --label) require_arg "$1" "${2-}"; LABEL="$2"; shift 2 ;; + --duration) require_arg "$1" "${2-}"; DURATION_SECS="$2"; shift 2 ;; + --out-dir) require_arg "$1" "${2-}"; OUT_DIR="$2"; shift 2 ;; + --pid) require_arg "$1" "${2-}"; RUSTFS_PID="$2"; shift 2 ;; + --container) require_arg "$1" "${2-}"; CONTAINER_NAME="$2"; shift 2 ;; + --endpoint) require_arg "$1" "${2-}"; ENDPOINT="$2"; shift 2 ;; + --perf) require_arg "$1" "${2-}"; PERF_MODE="$2"; shift 2 ;; + --perf-freq) require_arg "$1" "${2-}"; PERF_FREQ="$2"; shift 2 ;; + --sudo-cmd) require_arg "$1" "${2-}"; SUDO_CMD="$2"; shift 2 ;; + -h|--help) usage; exit 0 ;; + *) + warn "unknown argument: $1" + usage + exit 1 + ;; + esac + done +} + +finalize_defaults() { + if [[ -z "${OUT_DIR}" ]]; then + OUT_DIR="${PROJECT_ROOT}/target/perf/${LABEL}-$(date +%Y%m%d-%H%M%S)}" + fi +} + +command_exists() { + command -v "$1" >/dev/null 2>&1 +} + +write_cmd_output() { + local out_file="$1" + shift + if "$@" >"$out_file" 2>&1; then + return 0 + fi + warn "command failed, see ${out_file}" + return 1 +} + +resolve_pid() { + if [[ -n "${RUSTFS_PID}" ]]; then + printf '%s\n' "${RUSTFS_PID}" + return + fi + + if [[ -n "${CONTAINER_NAME}" ]] && command_exists docker; then + local pid + pid="$(docker inspect --format '{{.State.Pid}}' "${CONTAINER_NAME}" 2>/dev/null || true)" + if [[ -n "${pid}" && "${pid}" != "0" ]]; then + printf '%s\n' "${pid}" + return + fi + fi + + pgrep -n rustfs || true +} + +snapshot_proc() { + local pid="$1" + local prefix="$2" + [[ -n "${pid}" ]] || return 0 + + [[ -r "/proc/${pid}/status" ]] && cp "/proc/${pid}/status" "${OUT_DIR}/${prefix}.proc-status.txt" || true + [[ -r "/proc/${pid}/io" ]] && cp "/proc/${pid}/io" "${OUT_DIR}/${prefix}.proc-io.txt" || true + [[ -r "/proc/${pid}/sched" ]] && cp "/proc/${pid}/sched" "${OUT_DIR}/${prefix}.proc-sched.txt" || true + [[ -r "/proc/${pid}/smaps_rollup" ]] && cp "/proc/${pid}/smaps_rollup" "${OUT_DIR}/${prefix}.proc-smaps-rollup.txt" || true + [[ -r "/proc/${pid}/limits" ]] && cp "/proc/${pid}/limits" "${OUT_DIR}/${prefix}.proc-limits.txt" || true + + if command_exists ps; then + ps -p "${pid}" -o pid,ppid,stat,pcpu,pmem,rss,vsz,etime,args >"${OUT_DIR}/${prefix}.ps.txt" 2>&1 || true + ps -L -p "${pid}" -o pid,tid,psr,pcpu,stat,wchan:32,comm >"${OUT_DIR}/${prefix}.threads.txt" 2>&1 || true + fi + + if command_exists top; then + if [[ "$(uname -s)" == "Linux" ]]; then + top -H -b -n 1 -p "${pid}" >"${OUT_DIR}/${prefix}.top.txt" 2>&1 || true + else + top -l 1 -pid "${pid}" >"${OUT_DIR}/${prefix}.top.txt" 2>&1 || true + fi + fi +} + +capture_host_info() { + write_cmd_output "${OUT_DIR}/uname.txt" uname -a || true + command_exists lscpu && write_cmd_output "${OUT_DIR}/lscpu.txt" lscpu || true + command_exists free && write_cmd_output "${OUT_DIR}/free.txt" free -h || true + command_exists df && write_cmd_output "${OUT_DIR}/df.txt" df -h || true + command_exists mount && write_cmd_output "${OUT_DIR}/mount.txt" mount || true +} + +capture_endpoint_info() { + if command_exists curl; then + curl -fsS "${ENDPOINT}/health" >"${OUT_DIR}/health.txt" 2>&1 || true + curl -fsS "${ENDPOINT}/health/ready" >"${OUT_DIR}/health-ready.txt" 2>&1 || true + fi +} + +capture_container_info() { + [[ -n "${CONTAINER_NAME}" ]] || return 0 + command_exists docker || return 0 + + docker inspect "${CONTAINER_NAME}" >"${OUT_DIR}/docker-inspect.json" 2>&1 || true + docker logs --tail 500 "${CONTAINER_NAME}" >"${OUT_DIR}/docker-logs-tail.txt" 2>&1 || true + docker stats --no-stream --format '{{json .}}' "${CONTAINER_NAME}" >"${OUT_DIR}/docker-stats-once.jsonl" 2>&1 || true +} + +sample_container_stats_loop() { + [[ -n "${CONTAINER_NAME}" ]] || return 0 + command_exists docker || return 0 + + local out_file="${OUT_DIR}/docker-stats-loop.jsonl" + : >"${out_file}" + local end_ts=$((SECONDS + DURATION_SECS)) + while (( SECONDS < end_ts )); do + docker stats --no-stream --format '{{json .}}' "${CONTAINER_NAME}" >>"${out_file}" 2>/dev/null || true + sleep 1 + done +} + +sample_pidstat() { + local pid="$1" + [[ -n "${pid}" ]] || return 0 + command_exists pidstat || { + echo "pidstat unavailable" >"${OUT_DIR}/pidstat.txt" + return 0 + } + + pidstat -durwh -p "${pid}" 1 "${DURATION_SECS}" >"${OUT_DIR}/pidstat.txt" 2>&1 || true +} + +sample_perf() { + local pid="$1" + [[ -n "${pid}" ]] || return 0 + [[ "${PERF_MODE}" == "off" ]] && return 0 + command_exists perf || { + echo "perf unavailable" >"${OUT_DIR}/perf-record.log" + [[ "${PERF_MODE}" == "on" ]] && warn "perf requested but not installed" + return 0 + } + + local perf_data="${OUT_DIR}/perf.data" + local perf_log="${OUT_DIR}/perf-record.log" + local perf_report="${OUT_DIR}/perf-report.txt" + + local -a prefix=() + if [[ -n "${SUDO_CMD}" ]]; then + read -r -a prefix <<<"${SUDO_CMD}" + fi + + if "${prefix[@]}" perf record -F "${PERF_FREQ}" -g -p "${pid}" -o "${perf_data}" -- sleep "${DURATION_SECS}" \ + >"${perf_log}" 2>&1; then + "${prefix[@]}" perf report --stdio -i "${perf_data}" >"${perf_report}" 2>&1 || true + else + if [[ "${PERF_MODE}" == "on" ]]; then + warn "perf record failed; see ${perf_log}" + fi + fi +} + +capture_version_info() { + local pid="$1" + if [[ -n "${pid}" && -x "/proc/${pid}/exe" ]]; then + readlink "/proc/${pid}/exe" >"${OUT_DIR}/binary-path.txt" 2>&1 || true + "/proc/${pid}/exe" --help >"${OUT_DIR}/binary-help.txt" 2>&1 || true + fi +} + +main() { + parse_args "$@" + finalize_defaults + mkdir -p "${OUT_DIR}" + + local pid + pid="$(resolve_pid)" + if [[ -z "${pid}" ]]; then + warn "failed to detect rustfs pid automatically" + else + log "using rustfs pid=${pid}" + fi + + cat >"${OUT_DIR}/capture-meta.txt" </dev/null || true) +git_head=$(git -C "${PROJECT_ROOT}" rev-parse HEAD 2>/dev/null || true) +EOF + + capture_host_info + capture_endpoint_info + capture_container_info + capture_version_info "${pid}" + snapshot_proc "${pid}" "start" + + local bg_pids=() + sample_pidstat "${pid}" & + bg_pids+=($!) + sample_container_stats_loop & + bg_pids+=($!) + sample_perf "${pid}" & + bg_pids+=($!) + + for bg_pid in "${bg_pids[@]}"; do + wait "${bg_pid}" || true + done + + snapshot_proc "${pid}" "end" + capture_endpoint_info + + log "issue-2941 perf capture artifacts written to ${OUT_DIR}" + find "${OUT_DIR}" -maxdepth 1 -type f | sort +} + +main "$@"