fix: filter CPU profile samples to reviewed executable symbols (#8248)

This commit is contained in:
Chris
2026-09-30 01:06:24 +08:00
committed by GitHub
parent 1e7065101d
commit 5851d9eb54
3 changed files with 121 additions and 9 deletions
Generated
+2
View File
@@ -9553,6 +9553,7 @@ dependencies = [
"metrics-util", "metrics-util",
"mime_guess", "mime_guess",
"moka", "moka",
"object 0.40.0",
"opentelemetry 0.33.0", "opentelemetry 0.33.0",
"opentelemetry-proto", "opentelemetry-proto",
"opentelemetry_sdk 0.33.0", "opentelemetry_sdk 0.33.0",
@@ -9632,6 +9633,7 @@ dependencies = [
"snap", "snap",
"socket2", "socket2",
"subtle", "subtle",
"symbolic-demangle",
"sysinfo", "sysinfo",
"temp-env", "temp-env",
"tempfile", "tempfile",
+5 -1
View File
@@ -74,7 +74,7 @@ connect-e2e-short-credentials = []
# Builds the dedicated rustfs-cli-e2e target with a build-time public enrollment root. # Builds the dedicated rustfs-cli-e2e target with a build-time public enrollment root.
offline-enrollment-e2e-root = [] offline-enrollment-e2e-root = []
rio-v2 = ["rustfs-ecstore/rio-v2"] rio-v2 = ["rustfs-ecstore/rio-v2"]
pyroscope = ["rustfs-obs/pyroscope", "dep:pyroscope"] pyroscope = ["rustfs-obs/pyroscope", "dep:pyroscope", "dep:object", "dep:symbolic-demangle"]
# Tokio runtime telemetry. Requires `--cfg tokio_unstable`; use `make build-profiling`. # Tokio runtime telemetry. Requires `--cfg tokio_unstable`; use `make build-profiling`.
dial9 = ["rustfs-obs/dial9"] dial9 = ["rustfs-obs/dial9"]
hotpath = [ hotpath = [
@@ -379,6 +379,10 @@ libsystemd.workspace = true
[target.'cfg(any(all(target_os = "macos", any(target_arch = "x86_64", target_arch = "aarch64")), all(target_os = "linux", target_env = "gnu", any(target_arch = "x86_64", target_arch = "aarch64"))))'.dependencies] [target.'cfg(any(all(target_os = "macos", any(target_arch = "x86_64", target_arch = "aarch64")), all(target_os = "linux", target_env = "gnu", any(target_arch = "x86_64", target_arch = "aarch64"))))'.dependencies]
pyroscope = { workspace = true, optional = true, features = ["backend-pprof-rs"] } pyroscope = { workspace = true, optional = true, features = ["backend-pprof-rs"] }
[target.'cfg(all(target_os = "linux", target_env = "gnu", any(target_arch = "x86_64", target_arch = "aarch64")))'.dependencies]
object = { workspace = true, optional = true, features = ["read", "std", "elf"] }
symbolic-demangle = { version = "=13.9.0", default-features = false, features = ["rust"], optional = true }
[dev-dependencies] [dev-dependencies]
uuid = { workspace = true, features = ["v4", "v5", "fast-rng", "macro-diagnostics"] } uuid = { workspace = true, features = ["v4", "v5", "fast-rng", "macro-diagnostics"] }
serial_test = { workspace = true } serial_test = { workspace = true }
+114 -8
View File
@@ -649,9 +649,12 @@ pub async fn export_cpu_profile(
) )
))] ))]
mod local_cpu { mod local_cpu {
use std::collections::HashMap; use std::collections::{HashMap, HashSet};
use std::time::{Duration, Instant}; use std::time::{Duration, Instant};
#[cfg(all(target_os = "linux", target_env = "gnu"))]
use object::{Object as _, ObjectSymbol as _};
use pyroscope::backend::{BackendConfig, PprofConfig, ReportData, pprof_backend}; use pyroscope::backend::{BackendConfig, PprofConfig, ReportData, pprof_backend};
use sha2::{Digest as _, Sha256}; use sha2::{Digest as _, Sha256};
use tokio_util::sync::CancellationToken; use tokio_util::sync::CancellationToken;
@@ -665,6 +668,10 @@ mod local_cpu {
const MAX_UNIQUE_SYMBOLS: usize = 4_096; const MAX_UNIQUE_SYMBOLS: usize = 4_096;
const MAX_OUTPUT_SYMBOLS: usize = 256; const MAX_OUTPUT_SYMBOLS: usize = 256;
const MAX_SYMBOL_BYTES: usize = 4_096; const MAX_SYMBOL_BYTES: usize = 4_096;
#[cfg(all(target_os = "linux", target_env = "gnu"))]
const MAX_CATALOG_BYTES: usize = 32 * 1024 * 1024;
#[cfg(all(target_os = "linux", target_env = "gnu"))]
const MAX_CATALOG_SYMBOLS: usize = 200_000;
pub(super) fn collect( pub(super) fn collect(
request: &ProfileCaptureRequest, request: &ProfileCaptureRequest,
@@ -679,6 +686,8 @@ mod local_cpu {
return Err(ProfileError::LimitExceeded); return Err(ProfileError::LimitExceeded);
} }
let reviewed_symbols = reviewed_executable_symbols(cancel)?;
let started = Instant::now(); let started = Instant::now();
let deadline = started.checked_add(request.duration).ok_or(ProfileError::LimitExceeded)?; let deadline = started.checked_add(request.duration).ok_or(ProfileError::LimitExceeded)?;
let mut backend = pprof_backend(PprofConfig { sample_rate }, BackendConfig::default()) let mut backend = pprof_backend(PprofConfig { sample_rate }, BackendConfig::default())
@@ -694,7 +703,7 @@ mod local_cpu {
return Err(ProfileError::SourceUnavailable); return Err(ProfileError::SourceUnavailable);
}; };
let mut accumulator = Accumulator::new(request.nonce); let mut accumulator = Accumulator::new(request.nonce, reviewed_symbols.as_ref());
for report in reports { for report in reports {
for (stack, count) in report.data { for (stack, count) in report.data {
accumulator.record_stack(stack.frames.iter().filter_map(|frame| frame.name.as_deref()), count)?; accumulator.record_stack(stack.frames.iter().filter_map(|frame| frame.name.as_deref()), count)?;
@@ -718,18 +727,89 @@ mod local_cpu {
} }
} }
struct Accumulator { // The packaged Connect catalog is derived from this executable's ELF text symbols.
// Match that exact symbol set before signing; shared-library and runtime frames
// cannot be reviewed by a catalog bound to the RustFS executable alone.
fn reviewed_executable_symbols(cancel: &CancellationToken) -> Result<Option<HashSet<String>>, ProfileError> {
#[cfg(all(target_os = "linux", target_env = "gnu"))]
{
return load_executable_symbols(cancel).map(Some);
}
#[cfg(not(all(target_os = "linux", target_env = "gnu")))]
{
check_cancel(cancel)?;
Ok(None)
}
}
#[cfg(all(target_os = "linux", target_env = "gnu"))]
fn load_executable_symbols(cancel: &CancellationToken) -> Result<HashSet<String>, ProfileError> {
check_cancel(cancel)?;
let path = std::env::current_exe().map_err(|_| ProfileError::SourceUnavailable)?;
let executable = std::fs::File::open(path).map_err(|_| ProfileError::SourceUnavailable)?;
let size = executable.metadata().map_err(|_| ProfileError::SourceUnavailable)?.len();
if size == 0 || size > 1_073_741_824 {
return Err(ProfileError::SourceUnavailable);
}
let cache = object::read::ReadCache::new(executable);
let file = object::File::parse(&cache).map_err(|_| ProfileError::SourceUnavailable)?;
if file.format() != object::BinaryFormat::Elf || file.symbol_table().is_none() {
return Err(ProfileError::SourceUnavailable);
}
let mut symbols = HashSet::new();
let mut symbol_bytes = 0_usize;
for (index, symbol) in file.symbols().chain(file.dynamic_symbols()).enumerate() {
if index % 1_024 == 0 {
check_cancel(cancel)?;
}
if symbol.kind() != object::SymbolKind::Text || symbol.is_undefined() {
continue;
}
if let Some(name) = symbol.name_bytes().ok().and_then(catalog_symbol_name)
&& !symbols.contains(&name)
{
symbol_bytes = symbol_bytes
.checked_add(name.len() + 3)
.ok_or(ProfileError::SourceUnavailable)?;
if symbol_bytes > MAX_CATALOG_BYTES || symbols.len() >= MAX_CATALOG_SYMBOLS {
return Err(ProfileError::SourceUnavailable);
}
symbols.insert(name);
}
}
if symbols.is_empty() {
return Err(ProfileError::SourceUnavailable);
}
Ok(symbols)
}
#[cfg(all(target_os = "linux", target_env = "gnu"))]
fn catalog_symbol_name(raw: &[u8]) -> Option<String> {
let name = symbolic_demangle::demangle(&String::from_utf8_lossy(raw)).into_owned();
if name.is_empty()
|| name.len() > MAX_SYMBOL_BYTES
|| matches!(name.as_str(), "Unknown" | "<unresolved>")
|| name.bytes().any(|byte| byte < 32 || byte == 127)
{
return None;
}
Some(name)
}
struct Accumulator<'a> {
nonce: [u8; 32], nonce: [u8; 32],
reviewed_symbols: Option<&'a HashSet<String>>,
samples: HashMap<String, u64>, samples: HashMap<String, u64>,
stack_records: usize, stack_records: usize,
accepted_sample_count: u64, accepted_sample_count: u64,
dropped_sample_count: u64, dropped_sample_count: u64,
} }
impl Accumulator { impl<'a> Accumulator<'a> {
fn new(nonce: [u8; 32]) -> Self { fn new(nonce: [u8; 32], reviewed_symbols: Option<&'a HashSet<String>>) -> Self {
Self { Self {
nonce, nonce,
reviewed_symbols,
samples: HashMap::new(), samples: HashMap::new(),
stack_records: 0, stack_records: 0,
accepted_sample_count: 0, accepted_sample_count: 0,
@@ -756,6 +836,9 @@ mod local_cpu {
.into_iter() .into_iter()
.find(|symbol| !symbol.is_empty() && symbol.len() <= MAX_SYMBOL_BYTES) .find(|symbol| !symbol.is_empty() && symbol.len() <= MAX_SYMBOL_BYTES)
.unwrap_or("<unresolved>"); .unwrap_or("<unresolved>");
if self.reviewed_symbols.is_some_and(|reviewed| !reviewed.contains(symbol)) {
return self.drop_samples(accepted_count);
}
let symbol_id = symbol_id(&self.nonce, symbol); let symbol_id = symbol_id(&self.nonce, symbol);
if !self.samples.contains_key(&symbol_id) && self.samples.len() >= MAX_UNIQUE_SYMBOLS { if !self.samples.contains_key(&symbol_id) && self.samples.len() >= MAX_UNIQUE_SYMBOLS {
return self.drop_samples(accepted_count); return self.drop_samples(accepted_count);
@@ -820,12 +903,12 @@ mod local_cpu {
#[test] #[test]
fn summary_uses_nonce_bound_ids_and_excludes_raw_symbols() { fn summary_uses_nonce_bound_ids_and_excludes_raw_symbols() {
let raw_symbol = "rustfs::storage::disk::read_object"; let raw_symbol = "rustfs::storage::disk::read_object";
let mut first = Accumulator::new([7; 32]); let mut first = Accumulator::new([7; 32], None);
first first
.record_stack([raw_symbol].into_iter(), 9) .record_stack([raw_symbol].into_iter(), 9)
.expect("first stack should be recorded"); .expect("first stack should be recorded");
let first = first.finish(10_000).expect("first summary should be produced"); let first = first.finish(10_000).expect("first summary should be produced");
let mut second = Accumulator::new([8; 32]); let mut second = Accumulator::new([8; 32], None);
second second
.record_stack([raw_symbol].into_iter(), 9) .record_stack([raw_symbol].into_iter(), 9)
.expect("second stack should be recorded"); .expect("second stack should be recorded");
@@ -841,7 +924,7 @@ mod local_cpu {
#[test] #[test]
fn summary_bounds_samples_and_output_symbols() { fn summary_bounds_samples_and_output_symbols() {
let mut accumulator = Accumulator::new([3; 32]); let mut accumulator = Accumulator::new([3; 32], None);
for index in 0..=MAX_OUTPUT_SYMBOLS { for index in 0..=MAX_OUTPUT_SYMBOLS {
let symbol = format!("rustfs::bounded::{index}"); let symbol = format!("rustfs::bounded::{index}");
accumulator accumulator
@@ -867,6 +950,29 @@ mod local_cpu {
Err(ProfileError::Cancelled) Err(ProfileError::Cancelled)
)); ));
} }
#[test]
fn unreviewed_symbol_counts_are_dropped_without_exporting_their_ids() {
let reviewed = HashSet::from(["rustfs::known".to_owned()]);
let mut accumulator = Accumulator::new([4; 32], Some(&reviewed));
accumulator
.record_stack(["rustfs::known"].into_iter(), 3)
.expect("reviewed stack");
accumulator
.record_stack(["libc::unreviewed"].into_iter(), 2)
.expect("unreviewed stack");
let summary = accumulator.finish(10_000).expect("summary");
assert_eq!(summary.samples.len(), 1);
assert_eq!(summary.samples[0].sample_count, 3);
assert_eq!(summary.dropped_sample_count, 2);
assert_ne!(summary.samples[0].symbol_id, symbol_id(&[4; 32], "libc::unreviewed"));
let mut only_unreviewed = Accumulator::new([4; 32], Some(&reviewed));
only_unreviewed
.record_stack(["libc::unreviewed"].into_iter(), 2)
.expect("unreviewed stack");
assert!(matches!(only_unreviewed.finish(10_000), Err(ProfileError::SourceUnavailable)));
}
} }
} }