use std::fs::File;
#[cfg(target_os = "linux")]
use std::os::fd::BorrowedFd;
use std::time::Duration;
use std::array::TryFromSliceError;
use std::collections::{HashSet, HashMap};
use std::rc::Rc;
use tracing::{debug, error, info, trace, warn, enabled, Level};
use super::*;
use crate::sharing::*;
use crate::event::*;
#[cfg(target_os = "linux")]
use std::os::fd::FromRawFd;
#[cfg(target_os = "windows")]
struct BorrowedFd<'a> { fd: &'a u32 }
pub mod abi;
pub mod rb;
mod events;
mod bpf;
use abi::*;
pub use rb::source::RingBufSessionBuilder;
pub use rb::{RingBufOptions, RingBufBuilder};
pub use rb::cpu_count;
#[derive(Default)]
pub(crate) struct CaptureEnvironmentOptions {
all_mmaps: bool,
}
impl CaptureEnvironmentOptions {
pub(crate) fn with_all_mmaps(&mut self) -> &mut Self {
self.all_mmaps = true;
self
}
pub(crate) fn all_mmaps(&self) -> bool {
self.all_mmaps
}
}
static EMPTY: &[u8] = &[];
#[derive(Default)]
pub struct AncillaryData {
cpu: u32,
attributes: Rc<perf_event_attr>,
}
impl AncillaryData {
pub fn cpu(&self) -> u32 {
self.cpu
}
pub fn config(&self) -> u64 {
self.attributes.config
}
pub fn event_type(&self) -> u32 {
self.attributes.event_type
}
pub fn sample_type(&self) -> u64 {
self.attributes.sample_type
}
pub fn read_format(&self) -> u64 {
self.attributes.read_format
}
pub fn non_sampled_id_offsets(&self) -> Option<SampleIdOffsets> {
self.attributes.non_sampled_id_offsets()
}
}
impl Clone for AncillaryData {
fn clone(&self) -> Self {
Self {
cpu: self.cpu,
attributes: self.attributes.clone(),
}
}
}
pub struct PerfData<'a> {
pub ancillary: AncillaryData,
pub raw_data: &'a [u8],
}
impl<'a> Default for PerfData<'a> {
fn default() -> Self {
Self {
ancillary: AncillaryData::default(),
raw_data: EMPTY,
}
}
}
impl<'a> PerfData<'a> {
fn has_format(
&self,
format: u64) -> bool {
self.ancillary.attributes.has_format(format)
}
fn has_read_format(
&self,
format: u64) -> bool {
self.ancillary.attributes.has_read_format(format)
}
fn regs_user_count(&self) -> usize {
self.ancillary.attributes.sample_regs_user.count_ones() as usize
}
fn non_sampled_id_offsets(&self) -> Option<SampleIdOffsets> {
self.ancillary.non_sampled_id_offsets()
}
fn read_format_size(&self) -> usize {
let mut size: usize = 0;
if self.has_read_format(abi::PERF_FORMAT_TOTAL_TIME_ENABLED) {
size += 8;
}
if self.has_read_format(abi::PERF_FORMAT_TOTAL_TIME_RUNNING) {
size += 8;
}
if self.has_read_format(abi::PERF_FORMAT_ID) {
size += 8;
}
if self.has_read_format(abi::PERF_FORMAT_GROUP) {
size += 8;
}
if self.has_read_format(abi::PERF_FORMAT_LOST) {
size += 8;
}
size
}
fn read_u64(
&self,
offset: usize) -> Result<u64, TryFromSliceError> {
let slice = self.raw_data[offset .. offset + 8].try_into()?;
Ok(u64::from_ne_bytes(slice))
}
fn read_u32(
&self,
offset: usize) -> Result<u32, TryFromSliceError> {
let slice = self.raw_data[offset .. offset + 4].try_into()?;
Ok(u32::from_ne_bytes(slice))
}
fn read_u16(
&self,
offset: usize) -> Result<u16, TryFromSliceError> {
let slice = self.raw_data[offset .. offset + 2].try_into()?;
Ok(u16::from_ne_bytes(slice))
}
}
pub struct PerfDataFile {
id: u64,
fd: i32,
}
impl PerfDataFile {
pub fn new(
id: u64,
fd: i32) -> Self {
Self {
id,
fd,
}
}
pub fn id(&self) -> u64 { self.id }
pub fn fd(&self) -> i32 { self.fd }
}
pub trait PerfDataSource {
fn enable(&mut self) -> IOResult<()>;
fn disable(&mut self) -> IOResult<()>;
fn target_pids(&self) -> Option<&[i32]>;
fn create_bpf_files(
&mut self,
event: Option<&Event>) -> IOResult<Vec<PerfDataFile>>;
fn add_event(
&mut self,
event: &Event) -> IOResult<()>;
fn begin_reading(&mut self);
fn read(
&mut self,
timeout: Duration) -> Option<PerfData<'_>>;
fn end_reading(&mut self);
fn more(&self) -> bool;
fn begin_reading_cpu(
&mut self,
_cpu: u32) -> Option<CpuReadContext> { None }
fn read_cpu(
&mut self,
_ctx: &CpuReadContext) -> Option<PerfData<'_>> { None }
fn end_reading_cpu(
&mut self,
_ctx: &CpuReadContext) { }
#[cfg(target_os = "linux")]
fn perf_fds(&self) -> Vec<(u32, BorrowedFd<'_>)> { Vec::new() }
fn take_in_process_writer(
&mut self) -> Option<rb::InProcessRingBufWriter> { None }
}
pub struct CpuReadContext {
cpu: u32,
reader_index: usize,
}
impl CpuReadContext {
pub(crate) fn new(
cpu: u32,
reader_index: usize) -> Self {
Self { cpu, reader_index }
}
pub(crate) const fn cpu(&self) -> u32 {
self.cpu
}
pub(crate) const fn reader_index(&self) -> usize {
self.reader_index
}
}
#[derive(Debug)]
pub enum ParseForCpuError {
UnknownCpu(u32),
Parse(TryFromSliceError),
}
impl std::fmt::Display for ParseForCpuError {
fn fmt(
&self,
f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
ParseForCpuError::UnknownCpu(cpu) => write!(
f,
"no ring buffer for cpu={cpu}; valid CPUs come from perf_fds()"),
ParseForCpuError::Parse(e) => write!(f, "{e}"),
}
}
}
impl std::error::Error for ParseForCpuError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match self {
ParseForCpuError::UnknownCpu(_) => None,
ParseForCpuError::Parse(e) => Some(e),
}
}
}
impl From<TryFromSliceError> for ParseForCpuError {
fn from(e: TryFromSliceError) -> Self {
ParseForCpuError::Parse(e)
}
}
pub struct PerfSession {
source: Box<dyn PerfDataSource>,
source_enabled: bool,
events: HashMap<usize, Event>,
errors: Vec<anyhow::Error>,
event_error_callback: Option<Box<dyn Fn(&Event, &anyhow::Error)>>,
ip_field: DataFieldRef,
pid_field: DataFieldRef,
tid_field: DataFieldRef,
time_field: DataFieldRef,
address_field: DataFieldRef,
id_field: DataFieldRef,
stream_id_field: DataFieldRef,
cpu_field: DataFieldRef,
period_field: DataFieldRef,
read_field: DataFieldRef,
callchain_field: DataFieldRef,
raw_field: DataFieldRef,
branch_stack_field: DataFieldRef,
regs_user_field: DataFieldRef,
stack_user_field: DataFieldRef,
cgroup_field: DataFieldRef,
read_timeout: Duration,
capture_env_options: CaptureEnvironmentOptions,
cpu_profile_event: Event,
cswitch_profile_event: Event,
lost_event: Event,
comm_event: Event,
exit_event: Event,
fork_event: Event,
mmap_event: Event,
lost_samples_event: Event,
cswitch_event: Event,
soft_page_fault_event: Event,
hard_page_fault_event: Event,
drop_event: Event,
bpf_events: HashMap<u64, Writable<Event>>,
bpf_map_files: Vec<File>,
ancillary: Writable<AncillaryData>,
misc_field: DataFieldRef,
data_type_field: DataFieldRef,
}
impl Drop for PerfSession {
fn drop(&mut self) {
self.drop_event.process(
EMPTY,
EMPTY,
&mut self.errors);
self.log_errors(&self.drop_event);
}
}
impl PerfSession {
pub fn new(
source: Box<dyn PerfDataSource>) -> Self {
debug!("PerfSession::new: creating new PerfSession");
unsafe {
let mut limit = libc::rlimit {
rlim_cur: 0,
rlim_max: 0,
};
if libc::getrlimit(libc::RLIMIT_NOFILE, &mut limit) == 0 {
let old_cur = limit.rlim_cur;
limit.rlim_cur = limit.rlim_max;
let result = libc::setrlimit(libc::RLIMIT_NOFILE, &limit);
if result == 0 {
debug!(
"PerfSession::new: increased rlimit NOFILE from {} to {}",
old_cur, limit.rlim_max
);
} else {
if enabled!(Level::WARN) {
let error = std::io::Error::last_os_error();
warn!(
"PerfSession::new: setrlimit failed, attempted to increase from {} to {}, error={}",
old_cur, limit.rlim_max, error
);
}
}
}
}
let session = Self {
source,
source_enabled: false,
events: HashMap::new(),
errors: Vec::new(),
event_error_callback: None,
ip_field: DataFieldRef::new(),
pid_field: DataFieldRef::new(),
tid_field: DataFieldRef::new(),
time_field: DataFieldRef::new(),
address_field: DataFieldRef::new(),
id_field: DataFieldRef::new(),
stream_id_field: DataFieldRef::new(),
cpu_field: DataFieldRef::new(),
period_field: DataFieldRef::new(),
read_field: DataFieldRef::new(),
callchain_field: DataFieldRef::new(),
raw_field: DataFieldRef::new(),
branch_stack_field: DataFieldRef::new(),
regs_user_field: DataFieldRef::new(),
stack_user_field: DataFieldRef::new(),
cgroup_field: DataFieldRef::new(),
bpf_events: HashMap::new(),
bpf_map_files: Vec::new(),
read_timeout: Duration::from_millis(15),
capture_env_options: CaptureEnvironmentOptions::default(),
cpu_profile_event: Event::new(0, "__cpu_profile".into()),
cswitch_profile_event: Event::new(0, "__cswitch_profile".into()),
soft_page_fault_event: Event::new(0, "__soft_page_fault".into()),
hard_page_fault_event: Event::new(0, "__hard_page_fault".into()),
lost_event: events::lost(),
comm_event: events::comm(),
exit_event: events::exit(),
fork_event: events::fork(),
mmap_event: events::mmap(),
lost_samples_event: events::lost_samples(),
cswitch_event: events::cswitch(),
drop_event: Event::new(0, "__session_drop".into()),
ancillary: Writable::new(AncillaryData::default()),
misc_field: DataFieldRef::new(),
data_type_field: DataFieldRef::new(),
};
session.data_type_field.update(0, 4);
session.misc_field.update(4, 2);
session
}
pub fn set_event_error_callback(
&mut self,
callback: impl Fn(&Event, &anyhow::Error) + 'static) {
self.event_error_callback = Some(Box::new(callback));
}
pub fn ancillary_data(&self) -> ReadOnly<AncillaryData> {
self.ancillary.read_only()
}
pub fn cpu_profile_event(&mut self) -> &mut Event {
&mut self.cpu_profile_event
}
pub fn cswitch_profile_event(&mut self) -> &mut Event {
&mut self.cswitch_profile_event
}
pub fn soft_page_fault_event(&mut self) -> &mut Event {
&mut self.soft_page_fault_event
}
pub fn hard_page_fault_event(&mut self) -> &mut Event {
&mut self.hard_page_fault_event
}
pub fn lost_event(&mut self) -> &mut Event {
&mut self.lost_event
}
pub fn comm_event(&mut self) -> &mut Event {
&mut self.comm_event
}
pub fn exit_event(&mut self) -> &mut Event {
&mut self.exit_event
}
pub fn fork_event(&mut self) -> &mut Event {
&mut self.fork_event
}
pub fn mmap_event(&mut self) -> &mut Event {
&mut self.mmap_event
}
pub fn lost_samples_event(&mut self) -> &mut Event {
&mut self.lost_samples_event
}
pub fn cswitch_event(&mut self) -> &mut Event {
&mut self.cswitch_event
}
pub fn drop_event(&mut self) -> &mut Event {
&mut self.drop_event
}
pub fn misc_data_ref(&self) -> DataFieldRef {
self.misc_field.clone()
}
pub fn data_type_ref(&self) -> DataFieldRef {
self.data_type_field.clone()
}
pub fn ip_data_ref(&self) -> DataFieldRef {
self.ip_field.clone()
}
pub fn pid_field_ref(&self) -> DataFieldRef {
self.pid_field.clone()
}
pub fn tid_data_ref(&self) -> DataFieldRef {
self.tid_field.clone()
}
pub fn time_data_ref(&self) -> DataFieldRef {
self.time_field.clone()
}
pub fn address_data_ref(&self) -> DataFieldRef {
self.address_field.clone()
}
pub fn id_data_ref(&self) -> DataFieldRef {
self.id_field.clone()
}
pub fn stream_id_data_ref(&self) -> DataFieldRef {
self.stream_id_field.clone()
}
pub fn cpu_data_ref(&self) -> DataFieldRef {
self.cpu_field.clone()
}
pub fn period_data_ref(&self) -> DataFieldRef {
self.period_field.clone()
}
pub fn read_data_ref(&self) -> DataFieldRef {
self.read_field.clone()
}
pub fn callchain_data_ref(&self) -> DataFieldRef {
self.callchain_field.clone()
}
pub fn raw_data_ref(&self) -> DataFieldRef {
self.raw_field.clone()
}
pub fn branch_stack_data_ref(&self) -> DataFieldRef {
self.branch_stack_field.clone()
}
pub fn regs_user_data_ref(&self) -> DataFieldRef {
self.regs_user_field.clone()
}
pub fn stack_user_data_ref(&self) -> DataFieldRef {
self.stack_user_field.clone()
}
pub fn cgroup_data_ref(&self) -> DataFieldRef {
self.cgroup_field.clone()
}
pub fn set_read_timeout(
&mut self,
timeout: Duration) {
self.read_timeout = timeout;
}
pub fn perf_fds(&self) -> Vec<(u32, BorrowedFd<'_>)> {
self.source.perf_fds()
}
pub(crate) const fn capture_env_options_mut(&mut self) -> &mut CaptureEnvironmentOptions {
&mut self.capture_env_options
}
pub fn create_bpf_files(
&mut self,
event: Option<&Event>) -> IOResult<Vec<PerfDataFile>> {
self.source.create_bpf_files(event)
}
pub fn attach_to_bpf_map_path(
&mut self,
path: &str,
event: Event) -> IOResult<()> {
let path = std::ffi::CString::new(path)?;
let fd = bpf::bpf_get_map_fd_by_path(&path)?;
self.attach_to_bpf_map_fd(fd, event)
}
pub fn attach_to_bpf_map_id(
&mut self,
id: u32,
event: Event) -> IOResult<()> {
let fd = bpf::bpf_get_map_fd(id)?;
self.attach_to_bpf_map_fd(fd, event)
}
pub fn attach_to_bpf_map_fd(
&mut self,
fd: i32,
event: Event) -> IOResult<()> {
debug!(
"attach_to_bpf_map_fd: attaching event_name={}, event_id={}, fd={}",
event.name(), event.id(), fd
);
let bpf_files = match self.create_bpf_files(Some(&event)) {
Ok(bpf_files) => { bpf_files },
Err(err) => {
warn!(
"attach_to_bpf_map_fd failed: cannot create BPF files for event_name={}, fd={}, error={}",
event.name(), fd, err
);
unsafe {
libc::close(fd);
}
return Err(err);
}
};
let event = Writable::new(event);
for (i, perf_file) in bpf_files.iter().enumerate() {
match bpf::bpf_set_map_element(
fd,
i as u64,
(perf_file.fd()) as u64) {
Ok(()) => {
let file = unsafe {
File::from_raw_fd(fd)
};
self.bpf_map_files.push(file);
self.bpf_events.insert(perf_file.id(), event.clone());
trace!(
"attach_to_bpf_map_fd: attached perf_file, index={}, id={}, fd={}",
i, perf_file.id(), perf_file.fd()
);
},
Err(err) => {
warn!(
"attach_to_bpf_map_fd: bpf_set_map_element failed, index={}, fd={}, error={}",
i, fd, err
);
unsafe {
libc::close(fd);
}
return Err(err);
}
}
}
info!(
"attach_to_bpf_map_fd completed: event_name={}, fd={}, bpf_file_count={}",
event.borrow().name(), fd, bpf_files.len()
);
Ok(())
}
pub fn add_event(
&mut self,
event: Event) -> IOResult<()> {
debug!("add_event: adding event_name={}, event_id={}", event.name(), event.id());
self.source.add_event(&event)?;
self.events.insert(event.id(), event);
info!("PerfSession event added: event_count={}", self.events.len());
Ok(())
}
pub fn enable(&mut self) -> IOResult<()> {
debug!("PerfSession::enable: enabling data source");
self.source_enabled = true;
self.source.enable()?;
info!("PerfSession enabled");
Ok(())
}
pub fn disable(&mut self) -> IOResult<()> {
debug!("PerfSession::disable: disabling data source");
self.source_enabled = false;
self.source.disable()?;
info!("PerfSession disabled");
Ok(())
}
pub fn parse_all(&mut self) -> Result<(), TryFromSliceError> {
debug!("parse_all: starting to parse all available data");
let result = self.parse_until(|| false);
match &result {
Ok(_) => info!("parse_all completed successfully"),
Err(e) => warn!("parse_all failed: error={}", e),
}
result
}
pub fn parse_for_duration(
&mut self,
duration: Duration) -> Result<(), TryFromSliceError> {
debug!("parse_for_duration: parsing for {:?}", duration);
let now = std::time::Instant::now();
let result = self.parse_until(|| { now.elapsed() >= duration });
match &result {
Ok(_) => info!("parse_for_duration completed successfully: elapsed={:?}", now.elapsed()),
Err(e) => warn!("parse_for_duration failed: error={}", e),
}
result
}
pub fn parse_for_cpu(
&mut self,
cpu: u32) -> Result<(), ParseForCpuError> {
debug!("parse_for_cpu: parsing cpu={}", cpu);
let Some(ctx) = self.source.begin_reading_cpu(cpu) else {
warn!("parse_for_cpu: no ring buffer for cpu={}", cpu);
return Err(ParseForCpuError::UnknownCpu(cpu));
};
while self.parse_perf_data(None, Some(&ctx))? { }
self.source.end_reading_cpu(&ctx);
info!("parse_for_cpu completed: cpu={}", cpu);
Ok(())
}
fn log_errors(
&self,
event: &Event) {
for error in &self.errors {
if let Some(callback) = &self.event_error_callback {
callback(event, error);
} else {
error!("Event '{}': {}", event.name(), error);
}
}
}
fn parse_perf_data(
&mut self,
perf_data: Option<PerfData>,
cpu_ctx: Option<&CpuReadContext>) -> Result<bool, TryFromSliceError> {
let perf_data = perf_data.or_else(|| {
match cpu_ctx {
Some(ctx) => self.source.read_cpu(ctx),
None => self.source.read(self.read_timeout),
}
});
if perf_data.is_none() {
trace!("parse_perf_data: no data available");
return Ok(false);
}
let perf_data = perf_data.unwrap();
let header = abi::Header::from_slice(perf_data.raw_data)?;
trace!(
"parse_perf_data: processing record, entry_type={}, size={}",
header.entry_type, header.size
);
self.ancillary.write(|value| {
*value = perf_data.ancillary.clone();
});
if header.entry_type != abi::PERF_RECORD_SAMPLE {
match perf_data.non_sampled_id_offsets() {
Some(offsets) => {
let mut offset = header.size as usize - offsets.size;
if offsets.pid.is_some() {
offset += self.pid_field.update(offset, 4);
} else {
self.pid_field.reset();
}
if offsets.tid.is_some() {
offset += self.tid_field.update(offset, 4);
} else {
self.tid_field.reset();
}
if offsets.time.is_some() {
offset += self.time_field.update(offset, 8);
} else {
self.time_field.reset();
}
if offsets.id.is_some() {
offset += self.id_field.update(offset, 8);
} else {
self.id_field.reset();
}
if offsets.stream_id.is_some() {
offset += self.stream_id_field.update(offset, 8);
} else {
self.stream_id_field.reset();
}
if offsets.cpu.is_some() {
offset += self.cpu_field.update(offset, 8);
} else {
self.cpu_field.reset();
}
if offsets.identifier.is_some() {
self.id_field.update(offset, 8);
}
},
None => {
self.pid_field.reset();
self.tid_field.reset();
self.time_field.reset();
self.id_field.reset();
self.stream_id_field.reset();
self.cpu_field.reset();
}
}
}
match header.entry_type {
abi::PERF_RECORD_SAMPLE => {
let mut offset: usize = abi::Header::data_offset();
let mut id: Option<usize> = None;
if perf_data.has_format(abi::PERF_SAMPLE_IDENTIFIER) {
offset += self.id_field.update(offset, 8);
} else {
self.id_field.reset();
}
if perf_data.has_format(abi::PERF_SAMPLE_IP) {
offset += self.ip_field.update(offset, 8);
} else {
self.ip_field.reset();
}
if perf_data.has_format(abi::PERF_SAMPLE_TID) {
offset += self.pid_field.update(offset, 4);
offset += self.tid_field.update(offset, 4);
} else {
self.pid_field.reset();
self.tid_field.reset();
}
if perf_data.has_format(abi::PERF_SAMPLE_TIME) {
offset += self.time_field.update(offset, 8);
} else {
self.time_field.reset();
}
if perf_data.has_format(abi::PERF_SAMPLE_ADDR) {
offset += self.address_field.update(offset, 8);
} else {
self.address_field.reset();
}
if perf_data.has_format(abi::PERF_SAMPLE_ID) {
offset += self.id_field.update(offset, 8);
}
if perf_data.has_format(abi::PERF_SAMPLE_STREAM_ID) {
offset += self.stream_id_field.update(offset, 8);
} else {
self.stream_id_field.reset();
}
if perf_data.has_format(abi::PERF_SAMPLE_CPU) {
offset += self.cpu_field.update(offset, 8);
} else {
self.cpu_field.reset();
}
if perf_data.has_format(abi::PERF_SAMPLE_PERIOD) {
offset += self.period_field.update(offset, 8);
} else {
self.period_field.reset();
}
if perf_data.has_format(abi::PERF_SAMPLE_READ) {
let read_size = perf_data.read_format_size();
offset += self.read_field.update(offset, read_size);
} else {
self.read_field.reset();
}
if perf_data.has_format(abi::PERF_SAMPLE_CALLCHAIN) {
let count = perf_data.read_u64(offset)?;
let size = (count * 8) as usize;
offset += 8;
offset += self.callchain_field.update(offset, size);
} else {
self.callchain_field.reset();
}
if perf_data.has_format(abi::PERF_SAMPLE_RAW) {
let size = perf_data.read_u32(offset)? as usize;
let raw_field_size = 4 + size;
let raw_field_aligned_size = abi::align_to_perf_record(raw_field_size);
offset += 4;
id = Some(perf_data.read_u16(offset)? as usize);
offset += self.raw_field.update(offset, size);
offset += raw_field_aligned_size - raw_field_size;
} else {
self.raw_field.reset();
}
if perf_data.has_format(abi::PERF_SAMPLE_BRANCH_STACK) {
let count = perf_data.read_u64(offset)? as usize;
offset += 8;
let size = count * 24;
offset += self.branch_stack_field.update(offset, size);
} else {
self.branch_stack_field.reset();
}
if perf_data.has_format(abi::PERF_SAMPLE_REGS_USER) {
let abi = perf_data.read_u64(offset)?;
offset += 8;
let count = perf_data.regs_user_count();
let size = count * (abi * 4) as usize;
offset += self.regs_user_field.update(offset, size);
} else {
self.regs_user_field.reset();
}
if perf_data.has_format(abi::PERF_SAMPLE_STACK_USER) {
let size = perf_data.read_u64(offset)? as usize;
offset += 8;
if size > 0 {
let stack_start = offset;
offset += size;
let dyn_size = perf_data.read_u64(offset)? as usize;
offset += 8;
self.stack_user_field.update(stack_start, dyn_size);
} else {
self.stack_user_field.reset();
}
} else {
self.stack_user_field.reset();
}
const UNPARSED_SAMPLE_FIELDS_BEFORE_CGROUP: u64 =
abi::PERF_SAMPLE_WEIGHT |
abi::PERF_SAMPLE_DATA_SRC |
abi::PERF_SAMPLE_TRANSACTION |
abi::PERF_SAMPLE_REGS_INTR |
abi::PERF_SAMPLE_PHYS_ADDR;
if perf_data.has_format(abi::PERF_SAMPLE_CGROUP) {
let sample_type = perf_data.ancillary.attributes.sample_type;
debug_assert_eq!(
sample_type & UNPARSED_SAMPLE_FIELDS_BEFORE_CGROUP, 0,
"a builder enabled a sample field between STACK_USER and CGROUP; \
the cgroup parser needs explicit skip logic for it");
if (sample_type & UNPARSED_SAMPLE_FIELDS_BEFORE_CGROUP) != 0 {
warn!("Skipping PERF_SAMPLE_CGROUP: unparsed sample fields precede it");
self.cgroup_field.reset();
} else {
offset += self.cgroup_field.update(offset, 8);
}
} else {
self.cgroup_field.reset();
}
if offset > perf_data.raw_data.len() {
warn!("Truncated sample: offset={}, data_len={}", offset, perf_data.raw_data.len());
}
if perf_data.ancillary.config() == PERF_COUNT_SW_BPF_OUTPUT {
let full_data = perf_data.raw_data;
if let Some(id) = self.id_field.try_get_u64(full_data) {
if let Some(event) = self.bpf_events.get_mut(&id) {
let event_data = self.raw_field.get_data(full_data);
event.borrow_mut().process(
full_data,
event_data,
&mut self.errors);
}
if !self.errors.is_empty() {
if let Some(event) = self.bpf_events.get(&id) {
self.log_errors(&event.borrow());
}
}
}
} else if let Some(id) = &id {
if let Some(event) = self.events.get_mut(id) {
let full_data = perf_data.raw_data;
let event_data = self.raw_field.get_data(full_data);
event.process(
full_data,
event_data,
&mut self.errors);
}
if !self.errors.is_empty() {
if let Some(event) = self.events.get(id) {
self.log_errors(&event);
}
}
} else {
match perf_data.ancillary.event_type() {
PERF_TYPE_SOFTWARE => {
match perf_data.ancillary.config() {
PERF_COUNT_SW_CPU_CLOCK => {
self.cpu_profile_event.process(
perf_data.raw_data,
perf_data.raw_data,
&mut self.errors);
self.log_errors(&self.cpu_profile_event);
},
PERF_COUNT_SW_CONTEXT_SWITCHES => {
self.cswitch_profile_event.process(
perf_data.raw_data,
perf_data.raw_data,
&mut self.errors);
self.log_errors(&self.cswitch_profile_event);
},
PERF_COUNT_SW_PAGE_FAULTS_MIN => {
self.soft_page_fault_event.process(
perf_data.raw_data,
perf_data.raw_data,
&mut self.errors);
self.log_errors(&self.soft_page_fault_event);
},
PERF_COUNT_SW_PAGE_FAULTS_MAJ => {
self.hard_page_fault_event.process(
perf_data.raw_data,
perf_data.raw_data,
&mut self.errors);
self.log_errors(&self.hard_page_fault_event);
},
_ => { },
}
},
_ => { },
}
}
},
abi::PERF_RECORD_LOST => {
let offset = abi::Header::data_offset();
self.lost_event.process(
perf_data.raw_data,
&perf_data.raw_data[offset..],
&mut self.errors);
self.log_errors(&self.lost_event);
},
abi::PERF_RECORD_COMM => {
let offset = abi::Header::data_offset();
self.comm_event.process(
perf_data.raw_data,
&perf_data.raw_data[offset..],
&mut self.errors);
self.log_errors(&self.comm_event);
},
abi::PERF_RECORD_EXIT => {
let offset = abi::Header::data_offset();
self.exit_event.process(
perf_data.raw_data,
&perf_data.raw_data[offset..],
&mut self.errors);
self.log_errors(&self.exit_event);
},
abi::PERF_RECORD_FORK => {
let offset = abi::Header::data_offset();
self.fork_event.process(
perf_data.raw_data,
&perf_data.raw_data[offset..],
&mut self.errors);
self.log_errors(&self.fork_event);
},
abi::PERF_RECORD_MMAP2 => {
let offset = abi::Header::data_offset();
self.mmap_event.process(
perf_data.raw_data,
&perf_data.raw_data[offset..],
&mut self.errors);
self.log_errors(&self.mmap_event);
},
abi::PERF_RECORD_LOST_SAMPLES => {
let offset = abi::Header::data_offset();
self.lost_samples_event.process(
perf_data.raw_data,
&perf_data.raw_data[offset..],
&mut self.errors);
self.log_errors(&self.lost_samples_event);
},
abi::PERF_RECORD_SWITCH_CPU_WIDE => {
let offset = abi::Header::data_offset();
self.cswitch_event.process(
perf_data.raw_data,
&perf_data.raw_data[offset..],
&mut self.errors);
self.log_errors(&self.cswitch_event);
},
_ => {
},
}
self.errors.clear();
Ok(true)
}
pub fn parse_until(
&mut self,
should_stop: impl Fn() -> bool) -> Result<(), TryFromSliceError> {
loop {
let mut i: u32 = 0;
self.source.begin_reading();
while self.parse_perf_data(None, None)? {
if i >= 100 {
break;
}
i += 1;
}
self.source.end_reading();
if should_stop() || !self.source.more() {
break;
}
}
Ok(())
}
pub fn spawn_capture_environment(
&mut self) -> std::thread::JoinHandle<()> {
let writer = self.source.take_in_process_writer();
let mut pid_lookup = None;
if let Some(pids) = self.source.target_pids() {
let mut lookup = HashSet::new();
for pid in pids {
lookup.insert(*pid);
}
debug!("spawn_capture_environment: filtering by {} target PIDs", lookup.len());
pid_lookup = Some(lookup);
}
let captures_all = self.capture_env_options.all_mmaps();
std::thread::spawn(move || {
if let Some(mut writer) = writer {
debug!("capture_environment thread: starting");
Self::write_environment_comms(&pid_lookup, &mut writer);
Self::write_environment_modules(&pid_lookup, captures_all, &mut writer);
writer.flush();
info!("capture_environment thread: completed");
}
})
}
fn write_environment_comms(
pid_lookup: &Option<HashSet<i32>>,
writer: &mut rb::InProcessRingBufWriter) {
debug!(
"write_environment_comms: starting, pid_filter={}",
pid_lookup.is_some()
);
let attributes = RingBufBuilder::common_attributes();
let mut event_data: Vec<u8> = Vec::new();
let mut full_data: Vec<u8> = Vec::new();
let id_bytes = 0u64.to_ne_bytes();
let mut comm_count = 0;
procfs::iter_processes(|pid, path_buf| {
if let Some(pid_lookup) = &pid_lookup {
if !pid_lookup.contains(&(pid as i32)) {
return;
}
}
let comm = procfs::get_comm(path_buf)
.unwrap_or(String::new());
event_data.clear();
full_data.clear();
let pid_bytes = pid.to_ne_bytes();
event_data.extend_from_slice(&pid_bytes);
event_data.extend_from_slice(&pid_bytes);
event_data.extend_from_slice(comm.as_bytes());
event_data.push(0);
event_data.extend_from_slice(&pid_bytes);
event_data.extend_from_slice(&pid_bytes);
let capture_time = rb::perf_timestamp(&attributes);
event_data.extend_from_slice(&capture_time.to_ne_bytes());
event_data.extend_from_slice(&id_bytes);
abi::Header::write(PERF_RECORD_COMM, 0, &event_data, &mut full_data);
writer.write(&full_data);
comm_count += 1;
});
info!("write_environment_comms completed: comm_count={}", comm_count);
}
fn write_environment_modules(
pid_lookup: &Option<HashSet<i32>>,
captures_all: bool,
writer: &mut rb::InProcessRingBufWriter) {
debug!(
"write_environment_modules: starting, pid_filter={}",
pid_lookup.is_some()
);
let attributes = RingBufBuilder::common_attributes();
let mut event_data: Vec<u8> = Vec::new();
let mut full_data: Vec<u8> = Vec::new();
let gen_bytes = 0u64.to_ne_bytes();
let prot_bytes = 4u32.to_ne_bytes();
let flag_bytes = 0u32.to_ne_bytes();
let id_bytes = 0u64.to_ne_bytes();
let mut module_count = 0;
procfs::iter_modules(|pid, module| {
if !captures_all && !module.is_exec() {
return;
}
if let Some(pid_lookup) = &pid_lookup {
if !pid_lookup.contains(&(pid as i32)) {
return;
}
}
let path = module.path.unwrap_or("");
event_data.clear();
full_data.clear();
let pid_bytes = pid.to_ne_bytes();
event_data.extend_from_slice(&pid_bytes);
event_data.extend_from_slice(&pid_bytes);
event_data.extend_from_slice(&module.start_addr.to_ne_bytes());
event_data.extend_from_slice(&module.len().to_ne_bytes());
event_data.extend_from_slice(&module.offset.to_ne_bytes());
event_data.extend_from_slice(&module.dev_maj.to_ne_bytes());
event_data.extend_from_slice(&module.dev_min.to_ne_bytes());
event_data.extend_from_slice(&module.ino.to_ne_bytes());
event_data.extend_from_slice(&gen_bytes);
event_data.extend_from_slice(&prot_bytes);
event_data.extend_from_slice(&flag_bytes);
event_data.extend_from_slice(path.as_bytes());
event_data.push(0);
let capture_time = rb::perf_timestamp(&attributes);
event_data.extend_from_slice(&capture_time.to_ne_bytes());
event_data.extend_from_slice(&id_bytes);
abi::Header::write(PERF_RECORD_MMAP2, 0, &event_data, &mut full_data);
writer.write(&full_data);
module_count += 1;
});
info!("write_environment_modules completed: module_count={}", module_count);
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
struct MockData {
data: Vec<u8>,
entries: Vec<(usize, usize)>,
attr: Rc<perf_event_attr>,
index: usize,
cpu: u32,
}
impl MockData {
pub(crate) fn new(
sample_type: u64,
read_format: u64) -> Self {
let mut attr = perf_event_attr::default();
attr.sample_type = sample_type;
attr.read_format = read_format;
Self {
data: Vec::new(),
entries: Vec::new(),
attr: Rc::new(attr),
index: 0,
cpu: 0,
}
}
pub(crate) fn set_cpu(
&mut self,
cpu: u32) {
self.cpu = cpu;
}
pub(crate) fn push(
&mut self,
slice: &[u8]) {
let entry: (usize, usize) = (self.data.len(), slice.len());
self.entries.push(entry);
for byte in slice {
self.data.push(*byte);
}
}
}
impl PerfDataSource for MockData {
fn enable(&mut self) -> IOResult<()> { Ok(()) }
fn disable(&mut self) -> IOResult<()> { Ok(()) }
fn target_pids(&self) -> Option<&[i32]> { None }
fn create_bpf_files(
&mut self,
_event: Option<&Event>) -> IOResult<Vec<PerfDataFile>> {
Ok(Vec::new())
}
fn add_event(
&mut self,
_event: &Event) -> IOResult<()> { Ok(()) }
fn begin_reading(&mut self) { }
fn read(
&mut self,
_timeout: Duration) -> Option<PerfData<'_>> {
if !self.more() {
return None;
}
let entry = self.entries[self.index];
self.index += 1;
let start = entry.0;
let end = start + entry.1;
Some(PerfData {
ancillary: AncillaryData {
cpu: self.cpu,
attributes: self.attr.clone(),
},
raw_data: &self.data[start .. end],
})
}
fn end_reading(&mut self) { }
fn more(&self) -> bool {
self.index < self.entries.len()
}
fn begin_reading_cpu(
&mut self,
cpu: u32) -> Option<CpuReadContext> {
if cpu == self.cpu {
Some(CpuReadContext::new(cpu, 0))
} else {
None
}
}
fn read_cpu(
&mut self,
ctx: &CpuReadContext) -> Option<PerfData<'_>> {
if ctx.cpu() != self.cpu {
return None;
}
self.read(Duration::ZERO)
}
fn end_reading_cpu(
&mut self,
_ctx: &CpuReadContext) { }
}
#[test]
fn it_works() {
let mock = MockData::new(abi::PERF_SAMPLE_RAW, 0);
let mut e = Event::new(1, "test".into());
let format = e.format_mut();
format.add_field(
EventField::new(
"1".into(), "unsigned char".into(),
LocationType::Static, 0, 1));
format.add_field(
EventField::new(
"2".into(), "unsigned char".into(),
LocationType::Static, 1, 1));
format.add_field(
EventField::new(
"3".into(), "unsigned char".into(),
LocationType::Static, 2, 1));
let mut session = PerfSession::new(Box::new(mock));
let count = Arc::new(AtomicUsize::new(0));
let first = format.get_field_ref("1").unwrap();
let second = format.get_field_ref("2").unwrap();
let third = format.get_field_ref("3").unwrap();
e.add_callback(move |data| {
let format = data.format();
let event_data = data.event_data();
let a = format.get_data(first, event_data);
let b = format.get_data(second, event_data);
let c = format.get_data(third, event_data);
assert!(a[0] == 1u8);
assert!(b[0] == 2u8);
assert!(c[0] == 3u8);
count.fetch_add(1, Ordering::Relaxed);
Ok(())
});
session.add_event(e).unwrap();
}
#[test]
fn mock_data_sanity() {
let mut mock = MockData::new(0, 0);
let mut data: Vec<u8> = Vec::new();
data.push(1);
mock.push(data.as_slice());
data.clear();
data.push(2);
mock.push(data.as_slice());
data.clear();
data.push(3);
mock.push(data.as_slice());
data.clear();
drop(data);
let timeout = Duration::from_millis(100);
let first = mock.read(timeout).unwrap();
assert_eq!(1, first.raw_data[0]);
assert_eq!(1, first.raw_data.len());
let second = mock.read(timeout).unwrap();
assert_eq!(2, second.raw_data[0]);
assert_eq!(1, second.raw_data.len());
let third = mock.read(timeout).unwrap();
assert_eq!(3, third.raw_data[0]);
assert_eq!(1, third.raw_data.len());
assert!(mock.read(timeout).is_none());
assert!(!mock.more());
}
#[cfg(target_os = "linux")]
#[test]
fn perf_fds_empty_for_mock_source() {
let mock = MockData::new(0, 0);
let session = PerfSession::new(Box::new(mock));
assert!(session.perf_fds().is_empty());
}
#[test]
fn mock_data_perf_session() {
let count = Arc::new(AtomicUsize::new(0));
let sample_format =
abi::PERF_SAMPLE_TIME |
abi::PERF_SAMPLE_RAW;
let mut mock = MockData::new(sample_format, 0);
let mut perf_data = Vec::new();
let mut raw_data = Vec::new();
let mut event_data = Vec::new();
let id: u16 = 1;
let magic: u64 = 1234;
let time: u64 = 4321;
event_data.extend_from_slice(&id.to_ne_bytes());
event_data.extend_from_slice(&magic.to_ne_bytes());
Sample::write_time(time, &mut raw_data);
let raw_len = event_data.len() as u32;
let field_size = 4 + event_data.len();
let aligned_field_size = abi::align_to_perf_record(field_size);
raw_data.extend_from_slice(&raw_len.to_ne_bytes());
raw_data.extend_from_slice(event_data.as_slice());
let padding = aligned_field_size - field_size;
if padding > 0 {
raw_data.resize(raw_data.len() + padding, 0);
}
Header::write(abi::PERF_RECORD_SAMPLE, 0, raw_data.as_slice(), &mut perf_data);
mock.push(perf_data.as_slice());
perf_data.clear();
let mut session = PerfSession::new(Box::new(mock));
let mut e = Event::new(id as usize, "test".into());
let format = e.format_mut();
format.add_field(
EventField::new(
"common_type".into(), "unsigned short".into(),
LocationType::Static, 0, 2));
format.add_field(
EventField::new(
"magic".into(), "u64".into(),
LocationType::Static, 2, 8));
let callback_count = Arc::clone(&count);
let time_data = session.time_data_ref();
let magic_ref = format.get_field_ref("magic").unwrap();
e.add_callback(move |data| {
let full_data = data.full_data();
let format = data.format();
let event_data = data.event_data();
let read_time = time_data.try_get_u64(full_data).unwrap();
let read_magic = format.try_get_u64(magic_ref, event_data).unwrap();
assert_eq!(4321, read_time);
assert_eq!(1234, read_magic);
callback_count.fetch_add(1, Ordering::Relaxed);
Ok(())
});
session.add_event(e).unwrap();
session.parse_all().unwrap();
assert_eq!(count.load(Ordering::Relaxed), 1);
}
#[test]
fn mock_data_perf_session_raw_aligned_with_following_cpu() {
let count = Arc::new(AtomicUsize::new(0));
let sample_format =
abi::PERF_SAMPLE_RAW |
abi::PERF_SAMPLE_CPU;
let mut mock = MockData::new(sample_format, 0);
let mut perf_data = Vec::new();
let mut raw_data = Vec::new();
let mut event_data = Vec::new();
let id: u16 = 1;
let marker: u8 = 0xAB;
let cpu: u32 = 7;
event_data.extend_from_slice(&id.to_ne_bytes());
event_data.push(marker);
raw_data.extend_from_slice(&cpu.to_ne_bytes());
raw_data.extend_from_slice(&0u32.to_ne_bytes());
let raw_len = event_data.len() as u32;
let field_size = 4 + event_data.len();
let aligned_field_size = abi::align_to_perf_record(field_size);
raw_data.extend_from_slice(&raw_len.to_ne_bytes());
raw_data.extend_from_slice(event_data.as_slice());
let padding = aligned_field_size - field_size;
if padding > 0 {
raw_data.resize(raw_data.len() + padding, 0);
}
assert_eq!(16, raw_data.len());
Header::write(abi::PERF_RECORD_SAMPLE, 0, raw_data.as_slice(), &mut perf_data);
mock.push(perf_data.as_slice());
let mut session = PerfSession::new(Box::new(mock));
let mut e = Event::new(id as usize, "test".into());
let format = e.format_mut();
format.add_field(
EventField::new(
"common_type".into(), "unsigned short".into(),
LocationType::Static, 0, 2));
format.add_field(
EventField::new(
"marker".into(), "unsigned char".into(),
LocationType::Static, 2, 1));
let cpu_data = session.cpu_data_ref();
let callback_count = Arc::clone(&count);
let marker_ref = format.get_field_ref("marker").unwrap();
e.add_callback(move |data| {
let read_marker = data.format().get_data(marker_ref, data.event_data());
let read_cpu = cpu_data.try_get_u32(data.full_data()).unwrap();
assert_eq!(marker, read_marker[0]);
assert_eq!(cpu, read_cpu);
callback_count.fetch_add(1, Ordering::Relaxed);
Ok(())
});
session.add_event(e).unwrap();
session.parse_all().unwrap();
assert_eq!(1, count.load(Ordering::Relaxed));
}
fn push_sample(
mock: &mut MockData,
id: u16,
magic: u64,
time: u64) {
let mut perf_data = Vec::new();
let mut raw_data = Vec::new();
let mut event_data = Vec::new();
event_data.extend_from_slice(&id.to_ne_bytes());
event_data.extend_from_slice(&magic.to_ne_bytes());
Sample::write_time(time, &mut raw_data);
let raw_len = event_data.len() as u32;
let field_size = 4 + event_data.len();
let aligned_field_size = abi::align_to_perf_record(field_size);
raw_data.extend_from_slice(&raw_len.to_ne_bytes());
raw_data.extend_from_slice(event_data.as_slice());
let padding = aligned_field_size - field_size;
if padding > 0 {
raw_data.resize(raw_data.len() + padding, 0);
}
Header::write(abi::PERF_RECORD_SAMPLE, 0, raw_data.as_slice(), &mut perf_data);
mock.push(perf_data.as_slice());
}
#[test]
fn parse_for_cpu_drains_matching_cpu() {
let sample_format = abi::PERF_SAMPLE_TIME | abi::PERF_SAMPLE_RAW;
let id: u16 = 1;
let cpu: u32 = 3;
let mut mock = MockData::new(sample_format, 0);
mock.set_cpu(cpu);
push_sample(&mut mock, id, 1234, 4321);
let mut session = PerfSession::new(Box::new(mock));
let mut e = Event::new(id as usize, "test".into());
let format = e.format_mut();
format.add_field(
EventField::new(
"common_type".into(), "unsigned short".into(),
LocationType::Static, 0, 2));
format.add_field(
EventField::new(
"magic".into(), "u64".into(),
LocationType::Static, 2, 8));
let count = Arc::new(AtomicUsize::new(0));
let callback_count = Arc::clone(&count);
let magic_ref = format.get_field_ref("magic").unwrap();
e.add_callback(move |data| {
let read_magic = data.format()
.try_get_u64(magic_ref, data.event_data())
.unwrap();
assert_eq!(1234, read_magic);
callback_count.fetch_add(1, Ordering::Relaxed);
Ok(())
});
session.add_event(e).unwrap();
session.parse_for_cpu(cpu).unwrap();
assert_eq!(1, count.load(Ordering::Relaxed));
session.parse_for_cpu(cpu).unwrap();
assert_eq!(1, count.load(Ordering::Relaxed));
}
#[test]
fn parse_for_cpu_unknown_cpu_returns_error() {
let mock = MockData::new(0, 0);
let mut session = PerfSession::new(Box::new(mock));
match session.parse_for_cpu(7) {
Err(ParseForCpuError::UnknownCpu(cpu)) => assert_eq!(7, cpu),
other => panic!("expected UnknownCpu(7), got {other:?}"),
}
}
#[test]
fn mock_data_cgroup_field() {
let count = Arc::new(AtomicUsize::new(0));
let sample_format =
abi::PERF_SAMPLE_TIME |
abi::PERF_SAMPLE_CGROUP |
abi::PERF_SAMPLE_RAW;
let mut mock = MockData::new(sample_format, 0);
let mut perf_data = Vec::new();
let mut raw_data = Vec::new();
let id: u16 = 1;
let time: u64 = 9999;
let cgroup_id: u64 = 404228;
Sample::write_time(time, &mut raw_data);
let id_bytes = id.to_ne_bytes();
let field_size = 4 + id_bytes.len();
raw_data.extend_from_slice(&(id_bytes.len() as u32).to_ne_bytes());
raw_data.extend_from_slice(&id_bytes);
raw_data.resize(raw_data.len() + abi::align_to_perf_record(field_size) - field_size, 0);
raw_data.extend_from_slice(&cgroup_id.to_ne_bytes());
Header::write(abi::PERF_RECORD_SAMPLE, 0, raw_data.as_slice(), &mut perf_data);
mock.push(perf_data.as_slice());
let mut session = PerfSession::new(Box::new(mock));
let mut e = Event::new(id as usize, "test".into());
let cgroup_data = session.cgroup_data_ref();
let time_data = session.time_data_ref();
let callback_count = Arc::clone(&count);
e.add_callback(move |data| {
let full_data = data.full_data();
assert_eq!(Some(time), time_data.try_get_u64(full_data));
assert_eq!(Some(cgroup_id), cgroup_data.try_get_u64(full_data));
callback_count.fetch_add(1, Ordering::Relaxed);
Ok(())
});
session.add_event(e).unwrap();
session.parse_all().unwrap();
assert_eq!(1, count.load(Ordering::Relaxed));
}
}