#[cfg(target_os = "linux")]
use std::os::fd::BorrowedFd;
use tracing::{debug, info, trace, warn};
use super::*;
type BoxedBuilderHook = Box<dyn FnOnce(&mut RingBufSessionBuilder)>;
type BoxedSessionHook = Box<dyn FnOnce(&mut PerfSession)>;
struct RingBufSessionHook {
builder_hook: Option<BoxedBuilderHook>,
session_hook: Option<BoxedSessionHook>,
}
impl RingBufSessionHook {
pub(crate) fn new(
builder_hook: impl FnOnce(&mut RingBufSessionBuilder) + 'static,
session_hook: impl FnOnce(&mut PerfSession) + 'static) -> Self {
Self {
builder_hook: Some(Box::new(builder_hook)),
session_hook: Some(Box::new(session_hook)),
}
}
pub(crate) fn builder_hook(&mut self) -> Option<BoxedBuilderHook> {
self.builder_hook.take()
}
pub(crate) fn session_hook(&mut self) -> Option<BoxedSessionHook> {
self.session_hook.take()
}
}
pub struct RingBufSessionBuilder {
pages: usize,
target_pids: Option<Vec<i32>>,
target_cpus: Option<Vec<u16>>,
wakeup_watermark: Option<u32>,
kernel_builder: Option<RingBufBuilder<Kernel>>,
event_builder: Option<RingBufBuilder<Tracepoint>>,
profiling_builder: Option<RingBufBuilder<Profiling>>,
cswitch_builder: Option<RingBufBuilder<ContextSwitches>>,
soft_page_faults_builder: Option<RingBufBuilder<PageFaults>>,
hard_page_faults_builder: Option<RingBufBuilder<PageFaults>>,
bpf_builder: Option<RingBufBuilder<Bpf>>,
hooks: Option<Vec<RingBufSessionHook>>,
}
impl Default for RingBufSessionBuilder {
fn default() -> Self {
Self::new()
}
}
impl RingBufSessionBuilder {
pub fn new() -> Self {
Self {
pages: 1,
target_pids: None,
target_cpus: None,
wakeup_watermark: None,
kernel_builder: None,
event_builder: None,
profiling_builder: None,
cswitch_builder: None,
soft_page_faults_builder: None,
hard_page_faults_builder: None,
bpf_builder: None,
hooks: None,
}
}
pub fn with_target_pid(
&mut self,
pid: i32) -> Self {
let pids = match self.target_pids.take() {
Some(mut pids) => {
pids.push(pid);
Some(pids)
},
None => {
let mut pids = Vec::new();
pids.push(pid);
Some(pids)
},
};
Self {
pages: self.pages,
target_pids: pids,
target_cpus: self.target_cpus.take(),
wakeup_watermark: self.wakeup_watermark,
kernel_builder: self.kernel_builder.take(),
event_builder: self.event_builder.take(),
profiling_builder: self.profiling_builder.take(),
cswitch_builder: self.cswitch_builder.take(),
bpf_builder: self.bpf_builder.take(),
soft_page_faults_builder: self.soft_page_faults_builder.take(),
hard_page_faults_builder: self.hard_page_faults_builder.take(),
hooks: self.hooks.take(),
}
}
pub fn with_target_cpu(
&mut self,
cpu: u16) -> Self {
let cpus = match self.target_cpus.take() {
Some(mut cpus) => {
cpus.push(cpu);
Some(cpus)
},
None => {
let mut cpus = Vec::new();
cpus.push(cpu);
Some(cpus)
},
};
Self {
pages: self.pages,
target_pids: self.target_pids.take(),
target_cpus: cpus,
wakeup_watermark: self.wakeup_watermark,
kernel_builder: self.kernel_builder.take(),
event_builder: self.event_builder.take(),
profiling_builder: self.profiling_builder.take(),
cswitch_builder: self.cswitch_builder.take(),
bpf_builder: self.bpf_builder.take(),
soft_page_faults_builder: self.soft_page_faults_builder.take(),
hard_page_faults_builder: self.hard_page_faults_builder.take(),
hooks: self.hooks.take(),
}
}
pub fn with_page_count(
&mut self,
pages: usize) -> Self {
Self {
pages,
target_pids: self.target_pids.take(),
target_cpus: self.target_cpus.take(),
wakeup_watermark: self.wakeup_watermark,
kernel_builder: self.kernel_builder.take(),
event_builder: self.event_builder.take(),
profiling_builder: self.profiling_builder.take(),
cswitch_builder: self.cswitch_builder.take(),
bpf_builder: self.bpf_builder.take(),
soft_page_faults_builder: self.soft_page_faults_builder.take(),
hard_page_faults_builder: self.hard_page_faults_builder.take(),
hooks: self.hooks.take(),
}
}
pub fn with_wakeup_watermark(
&mut self,
bytes: u32) -> Self {
Self {
pages: self.pages,
target_pids: self.target_pids.take(),
target_cpus: self.target_cpus.take(),
wakeup_watermark: Some(bytes),
kernel_builder: self.kernel_builder.take(),
event_builder: self.event_builder.take(),
profiling_builder: self.profiling_builder.take(),
cswitch_builder: self.cswitch_builder.take(),
bpf_builder: self.bpf_builder.take(),
soft_page_faults_builder: self.soft_page_faults_builder.take(),
hard_page_faults_builder: self.hard_page_faults_builder.take(),
hooks: self.hooks.take(),
}
}
pub fn with_kernel_events(
&mut self,
builder: RingBufBuilder<Kernel>) -> Self {
Self {
pages: self.pages,
target_pids: self.target_pids.take(),
target_cpus: self.target_cpus.take(),
wakeup_watermark: self.wakeup_watermark,
kernel_builder: Some(builder),
event_builder: self.event_builder.take(),
profiling_builder: self.profiling_builder.take(),
cswitch_builder: self.cswitch_builder.take(),
bpf_builder: self.bpf_builder.take(),
soft_page_faults_builder: self.soft_page_faults_builder.take(),
hard_page_faults_builder: self.hard_page_faults_builder.take(),
hooks: self.hooks.take(),
}
}
pub fn take_kernel_events(
&mut self) -> Option<RingBufBuilder<Kernel>> {
self.kernel_builder.take()
}
pub fn replace_kernel_events(
&mut self,
builder: RingBufBuilder<Kernel>) -> Option<RingBufBuilder<Kernel>> {
self.kernel_builder.replace(builder)
}
pub fn with_tracepoint_events(
&mut self,
builder: RingBufBuilder<Tracepoint>) -> Self {
Self {
pages: self.pages,
target_pids: self.target_pids.take(),
target_cpus: self.target_cpus.take(),
wakeup_watermark: self.wakeup_watermark,
kernel_builder: self.kernel_builder.take(),
event_builder: Some(builder),
profiling_builder: self.profiling_builder.take(),
cswitch_builder: self.cswitch_builder.take(),
bpf_builder: self.bpf_builder.take(),
soft_page_faults_builder: self.soft_page_faults_builder.take(),
hard_page_faults_builder: self.hard_page_faults_builder.take(),
hooks: self.hooks.take(),
}
}
pub fn take_tracepoint_events(
&mut self) -> Option<RingBufBuilder<Tracepoint>> {
self.event_builder.take()
}
pub fn replace_tracepoint_events(
&mut self,
builder: RingBufBuilder<Tracepoint>) -> Option<RingBufBuilder<Tracepoint>> {
self.event_builder.replace(builder)
}
pub fn with_profiling_events(
&mut self,
builder: RingBufBuilder<Profiling>) -> Self {
Self {
pages: self.pages,
target_pids: self.target_pids.take(),
target_cpus: self.target_cpus.take(),
wakeup_watermark: self.wakeup_watermark,
kernel_builder: self.kernel_builder.take(),
event_builder: self.event_builder.take(),
profiling_builder: Some(builder),
cswitch_builder: self.cswitch_builder.take(),
bpf_builder: self.bpf_builder.take(),
soft_page_faults_builder: self.soft_page_faults_builder.take(),
hard_page_faults_builder: self.hard_page_faults_builder.take(),
hooks: self.hooks.take(),
}
}
pub fn take_profiling_events(
&mut self) -> Option<RingBufBuilder<Profiling>> {
self.profiling_builder.take()
}
pub fn replace_profiling_events(
&mut self,
builder: RingBufBuilder<Profiling>) -> Option<RingBufBuilder<Profiling>> {
self.profiling_builder.replace(builder)
}
pub fn with_cswitch_events(
&mut self,
builder: RingBufBuilder<ContextSwitches>) -> Self {
Self {
pages: self.pages,
target_pids: self.target_pids.take(),
target_cpus: self.target_cpus.take(),
wakeup_watermark: self.wakeup_watermark,
kernel_builder: self.kernel_builder.take(),
event_builder: self.event_builder.take(),
profiling_builder: self.profiling_builder.take(),
cswitch_builder: Some(builder),
bpf_builder: self.bpf_builder.take(),
soft_page_faults_builder: self.soft_page_faults_builder.take(),
hard_page_faults_builder: self.hard_page_faults_builder.take(),
hooks: self.hooks.take(),
}
}
pub fn take_cswitch_events(
&mut self) -> Option<RingBufBuilder<ContextSwitches>> {
self.cswitch_builder.take()
}
pub fn replace_cswitch_events(
&mut self,
builder: RingBufBuilder<ContextSwitches>) -> Option<RingBufBuilder<ContextSwitches>> {
self.cswitch_builder.replace(builder)
}
pub fn with_bpf_events(
&mut self,
builder: RingBufBuilder<Bpf>) -> Self {
Self {
pages: self.pages,
target_pids: self.target_pids.take(),
target_cpus: self.target_cpus.take(),
wakeup_watermark: self.wakeup_watermark,
kernel_builder: self.kernel_builder.take(),
event_builder: self.event_builder.take(),
profiling_builder: self.profiling_builder.take(),
cswitch_builder: self.cswitch_builder.take(),
bpf_builder: Some(builder),
soft_page_faults_builder: self.soft_page_faults_builder.take(),
hard_page_faults_builder: self.hard_page_faults_builder.take(),
hooks: self.hooks.take(),
}
}
pub fn take_bpf_events(
&mut self) -> Option<RingBufBuilder<Bpf>> {
self.bpf_builder.take()
}
pub fn replace_bpf_events(
&mut self,
builder: RingBufBuilder<Bpf>) -> Option<RingBufBuilder<Bpf>> {
self.bpf_builder.replace(builder)
}
pub fn with_soft_page_faults_events(
&mut self,
builder: RingBufBuilder<PageFaults>) -> Self {
Self {
pages: self.pages,
target_pids: self.target_pids.take(),
target_cpus: self.target_cpus.take(),
wakeup_watermark: self.wakeup_watermark,
kernel_builder: self.kernel_builder.take(),
event_builder: self.event_builder.take(),
profiling_builder: self.profiling_builder.take(),
cswitch_builder: self.cswitch_builder.take(),
bpf_builder: self.bpf_builder.take(),
soft_page_faults_builder: Some(builder),
hard_page_faults_builder: self.hard_page_faults_builder.take(),
hooks: self.hooks.take(),
}
}
pub fn take_soft_page_faults_events(
&mut self) -> Option<RingBufBuilder<PageFaults>> {
self.soft_page_faults_builder.take()
}
pub fn replace_soft_page_faults_events(
&mut self,
builder: RingBufBuilder<PageFaults>) -> Option<RingBufBuilder<PageFaults>> {
self.soft_page_faults_builder.replace(builder)
}
pub fn with_hard_page_faults_events(
&mut self,
builder: RingBufBuilder<PageFaults>) -> Self {
Self {
pages: self.pages,
target_pids: self.target_pids.take(),
target_cpus: self.target_cpus.take(),
wakeup_watermark: self.wakeup_watermark,
kernel_builder: self.kernel_builder.take(),
event_builder: self.event_builder.take(),
profiling_builder: self.profiling_builder.take(),
cswitch_builder: self.cswitch_builder.take(),
bpf_builder: self.bpf_builder.take(),
soft_page_faults_builder: self.soft_page_faults_builder.take(),
hard_page_faults_builder: Some(builder),
hooks: self.hooks.take(),
}
}
pub fn take_hard_page_faults_events(
&mut self) -> Option<RingBufBuilder<PageFaults>> {
self.hard_page_faults_builder.take()
}
pub fn replace_hard_page_faults_events(
&mut self,
builder: RingBufBuilder<PageFaults>) -> Option<RingBufBuilder<PageFaults>> {
self.hard_page_faults_builder.replace(builder)
}
pub fn with_hooks(
&mut self,
builder_hook: impl FnOnce(&mut RingBufSessionBuilder) + 'static,
session_hook: impl FnOnce(&mut PerfSession) + 'static) -> Self {
let mut hooks = self.hooks.take().unwrap_or_default();
hooks.push(
RingBufSessionHook::new(
Box::new(builder_hook),
Box::new(session_hook)));
Self {
pages: self.pages,
target_pids: self.target_pids.take(),
target_cpus: self.target_cpus.take(),
wakeup_watermark: self.wakeup_watermark,
kernel_builder: self.kernel_builder.take(),
event_builder: self.event_builder.take(),
profiling_builder: self.profiling_builder.take(),
cswitch_builder: self.cswitch_builder.take(),
bpf_builder: self.bpf_builder.take(),
soft_page_faults_builder: self.soft_page_faults_builder.take(),
hard_page_faults_builder: self.hard_page_faults_builder.take(),
hooks: Some(hooks),
}
}
pub fn build(&mut self) -> IOResult<PerfSession> {
debug!(
"RingBufSessionBuilder::build: pages={}, target_pids={:?}",
self.pages, self.target_pids
);
let mut hooks = self.hooks.take();
if let Some(hooks) = &mut hooks {
for hook in hooks {
if let Some(hook) = hook.builder_hook() {
(hook)(self);
}
}
}
if let Some(bytes) = self.wakeup_watermark.take() {
let ring_bytes = ring_data_bytes(self.pages);
if (bytes as usize) >= ring_bytes {
let msg = format!(
"wakeup_watermark ({} bytes) must be smaller than the \
per-CPU ring data area ({} bytes from {} pages); the \
kernel would never wake the perf fds",
bytes, ring_bytes, self.pages);
warn!("RingBufSessionBuilder::build: {}", msg);
return Err(io_error(&msg));
}
let mut kernel = self.kernel_builder
.take()
.unwrap_or_else(RingBufBuilder::for_kernel);
kernel.set_wakeup_watermark(bytes);
self.kernel_builder = Some(kernel);
debug!(
"RingBufSessionBuilder::build: wakeup_watermark applied, bytes={}, ring_bytes={}",
bytes, ring_bytes);
}
let mut source = RingBufDataSource::new(
self.pages,
self.target_pids.take(),
self.target_cpus.take(),
self.kernel_builder.take(),
self.event_builder.take(),
self.profiling_builder.take(),
self.cswitch_builder.take(),
self.bpf_builder.take(),
self.soft_page_faults_builder.take(),
self.hard_page_faults_builder.take());
source.build()?;
let mut session = PerfSession::new(Box::new(source));
if let Some(hooks) = &mut hooks {
for hook in hooks {
if let Some(hook) = hook.session_hook() {
(hook)(&mut session);
}
}
}
info!("PerfSession created successfully");
Ok(session)
}
}
pub struct RingBufDataSource {
readers: Vec<CpuRingReader>,
cursors: Vec<CpuRingCursor>,
temp: Vec<u8>,
leader_ids: HashMap<u32, u64>,
ring_bufs: HashMap<u64, CpuRingBuf>,
pages: usize,
enabled: bool,
target_pids: Option<Vec<i32>>,
target_cpus: Option<Vec<u16>>,
kernel_builder: Option<RingBufBuilder<Kernel>>,
event_builder: Option<RingBufBuilder<Tracepoint>>,
profiling_builder: Option<RingBufBuilder<Profiling>>,
cswitch_builder: Option<RingBufBuilder<ContextSwitches>>,
bpf_builder: Option<RingBufBuilder<Bpf>>,
soft_page_faults_builder: Option<RingBufBuilder<PageFaults>>,
hard_page_faults_builder: Option<RingBufBuilder<PageFaults>>,
next_time: Option<u64>,
oldest_cpu: Option<usize>,
cpu_to_reader_index: Vec<Option<usize>>,
in_process_ring_buf: Option<InProcessRingBuf>,
}
impl RingBufDataSource {
fn new(
pages: usize,
target_pids: Option<Vec<i32>>,
target_cpus: Option<Vec<u16>>,
kernel_builder: Option<RingBufBuilder<Kernel>>,
event_builder: Option<RingBufBuilder<Tracepoint>>,
profiling_builder: Option<RingBufBuilder<Profiling>>,
cswitch_builder: Option<RingBufBuilder<ContextSwitches>>,
bpf_builder: Option<RingBufBuilder<Bpf>>,
soft_page_faults_builder: Option<RingBufBuilder<PageFaults>>,
hard_page_faults_builder: Option<RingBufBuilder<PageFaults>>) -> Self {
Self {
readers: Vec::new(),
cursors: Vec::new(),
temp: Vec::new(),
leader_ids: HashMap::new(),
ring_bufs: HashMap::new(),
pages,
target_pids,
target_cpus,
kernel_builder,
event_builder,
profiling_builder,
cswitch_builder,
bpf_builder,
soft_page_faults_builder,
hard_page_faults_builder,
next_time: None,
oldest_cpu: None,
cpu_to_reader_index: Vec::new(),
enabled: false,
in_process_ring_buf: None,
}
}
fn add_cpu_bufs(
target_pid: Option<i32>,
target_cpus: &Option<Vec<u16>>,
leader_ids: &HashMap<u32, u64>,
ring_bufs: &mut HashMap<u64, CpuRingBuf>,
common_buf: &CommonRingBuf,
mut fds: Option<&mut Vec<PerfDataFile>>) -> IOResult<()> {
for i in 0..cpu_count() {
if let Some(target_cpus) = target_cpus {
if !target_cpus.contains(&(i as u16)) {
continue;
}
}
let leader_id = leader_ids[&i];
let leader = &ring_bufs[&leader_id];
let mut cpu_buf = common_buf.for_cpu(i);
cpu_buf.open(target_pid)?;
match cpu_buf.id() {
Some(id) => {
cpu_buf.redirect_to(leader)?;
if let Some(fds) = fds.as_mut() {
fds.push(
PerfDataFile::new(
id,
cpu_buf.fd.unwrap()));
}
debug!(
"add_cpu_bufs: cpu={}, id={}, leader_id={}, target_pid={:?}",
i, id, leader_id, target_pid
);
ring_bufs.insert(id, cpu_buf);
},
None => {
warn!(
"add_cpu_bufs failed: no buffer ID returned, cpu={}, target_pid={:?}",
i, target_pid
);
return Err(io_error(
"Internal error getting buffer ID."));
}
}
}
Ok(())
}
fn tasks_for_pids(pids: &mut Vec<i32>) {
let mut tasks = HashSet::new();
for pid in pids.drain(..) {
tasks.insert(pid);
procfs::iter_proc_tasks(
pid as u32,
|task| { tasks.insert(task as i32); });
}
for task in tasks.drain() {
pids.push(task);
}
}
fn build(&mut self) -> IOResult<()> {
debug!(
"RingBufDataSource::build: pages={}, target_pids={:?}",
self.pages, self.target_pids
);
let common = self.kernel_builder
.get_or_insert_with(RingBufBuilder::for_kernel)
.build();
let empty_pids = Vec::new();
let target_pids = &mut self.target_pids.as_mut();
let target_cpus = &self.target_cpus;
let pids = match target_pids {
Some(pids) => {
Self::tasks_for_pids(pids);
debug!("build: found {} task(s) for target PIDs", pids.len());
pids
},
None => { &empty_pids },
};
self.cpu_to_reader_index.resize(cpu_count() as usize, None);
for i in 0..cpu_count() {
let mut cpu_buf = common.for_cpu(i);
if pids.is_empty() {
cpu_buf.open(None)?;
} else {
cpu_buf.open(Some(pids[0]))?;
}
match cpu_buf.id() {
Some(id) => {
self.leader_ids.insert(i, id);
let reader = cpu_buf.create_reader(self.pages)?;
let reader_index = self.readers.len();
self.cpu_to_reader_index[i as usize] = Some(reader_index);
self.readers.push(reader);
self.cursors.push(CpuRingCursor::default());
debug!("build: leader ring buffer created, cpu={}, id={}", i, id);
self.ring_bufs.insert(id, cpu_buf);
},
None => {
warn!("build failed: no buffer ID returned for leader ring, cpu={}", i);
return Err(io_error(
"Internal error getting buffer ID."));
}
}
}
debug!(
"build: leader ring buffers created, cpu_count={}, ring_buf_count={}",
cpu_count(), self.ring_bufs.len()
);
if !pids.is_empty() {
for pid in &pids[1..] {
Self::add_cpu_bufs(
Some(*pid),
&None,
&self.leader_ids,
&mut self.ring_bufs,
&common,
None)?;
}
}
if let Some(profiling_builder) = self.profiling_builder.as_mut() {
debug!("build: adding profiling event buffers");
let common = profiling_builder.build();
if pids.is_empty() {
Self::add_cpu_bufs(
None,
target_cpus,
&self.leader_ids,
&mut self.ring_bufs,
&common,
None)?;
} else {
for pid in pids {
Self::add_cpu_bufs(
Some(*pid),
target_cpus,
&self.leader_ids,
&mut self.ring_bufs,
&common,
None)?;
}
}
}
if let Some(cswitch_builder) = self.cswitch_builder.as_mut() {
debug!("build: adding context switch event buffers");
let common = cswitch_builder.build();
if pids.is_empty() {
Self::add_cpu_bufs(
None,
target_cpus,
&self.leader_ids,
&mut self.ring_bufs,
&common,
None)?;
} else {
for pid in pids {
Self::add_cpu_bufs(
Some(*pid),
target_cpus,
&self.leader_ids,
&mut self.ring_bufs,
&common,
None)?;
}
}
}
if let Some(faults_builder) = self.soft_page_faults_builder.as_mut() {
debug!("build: adding soft page faults event buffers");
let common = faults_builder.build();
if pids.is_empty() {
Self::add_cpu_bufs(
None,
target_cpus,
&self.leader_ids,
&mut self.ring_bufs,
&common,
None)?;
} else {
for pid in pids {
Self::add_cpu_bufs(
Some(*pid),
target_cpus,
&self.leader_ids,
&mut self.ring_bufs,
&common,
None)?;
}
}
}
if let Some(faults_builder) = self.hard_page_faults_builder.as_mut() {
debug!("build: adding hard page faults event buffers");
let common = faults_builder.build();
if pids.is_empty() {
Self::add_cpu_bufs(
None,
target_cpus,
&self.leader_ids,
&mut self.ring_bufs,
&common,
None)?;
} else {
for pid in pids {
Self::add_cpu_bufs(
Some(*pid),
target_cpus,
&self.leader_ids,
&mut self.ring_bufs,
&common,
None)?;
}
}
}
let mut in_process = InProcessRingBuf::new(self.pages);
let reader = in_process.create_reader();
self.readers.push(reader);
self.cursors.push(CpuRingCursor::default());
self.in_process_ring_buf = Some(in_process);
let common_attrs = Rc::new(RingBufBuilder::common_attributes());
self.ring_bufs.insert(0, CpuRingBuf::new(0, common_attrs));
info!(
"RingBufDataSource built successfully: ring_buf_count={}, reader_count={}",
self.ring_bufs.len(), self.readers.len()
);
Ok(())
}
fn enable(&mut self) -> IOResult<()> {
debug!("RingBufDataSource::enable: enabling {} ring buffers", self.ring_bufs.len());
for rb in self.ring_bufs.values() {
if !rb.is_open() {
continue;
}
rb.enable()?;
}
self.enabled = true;
info!("RingBufDataSource enabled: ring_buf_count={}", self.ring_bufs.len());
Ok(())
}
fn disable(&mut self) -> IOResult<()> {
debug!("RingBufDataSource::disable: disabling {} ring buffers", self.ring_bufs.len());
for rb in self.ring_bufs.values() {
if !rb.is_open() {
continue;
}
rb.disable()?;
}
self.enabled = false;
info!("RingBufDataSource disabled: ring_buf_count={}", self.ring_bufs.len());
Ok(())
}
fn read_time<'a>(
reader: &'a CpuRingReader,
cursor: &'a mut CpuRingCursor,
ring_bufs: &'a HashMap<u64, CpuRingBuf>) -> Option<(u64, &'a CpuRingBuf)> {
let mut start = 0;
let slice = reader.data_slice();
if !cursor.more() {
return None;
}
match reader.peek_header(
cursor,
slice,
&mut start) {
Ok(header) => {
let min_size = abi::Header::data_offset() + 16;
if (header.size as usize) < min_size {
warn!(
"read_time: invalid header size={}, skipping record",
header.size
);
let skip = std::cmp::max(
header.size as usize,
abi::Header::data_offset(),
) as u16;
cursor.advance(skip);
return None;
}
let id_offset: u16;
let mut time_offset: Option<u16> = None;
if header.entry_type == abi::PERF_RECORD_SAMPLE {
id_offset = abi::Header::data_offset() as u16;
} else {
time_offset = Some(header.size - 16);
id_offset = header.size - 8;
}
let id = reader.peek_u64(
cursor,
id_offset as u64);
let Some(buf) = ring_bufs.get(&id) else {
warn!(
"read_time: no ring buffer found for id={}, skipping record",
id
);
cursor.advance(header.size);
return None;
};
if time_offset.is_none() {
time_offset = Some(buf.sample_time_offset());
}
let time = reader.peek_u64(
cursor,
time_offset.unwrap() as u64);
Some((time, buf))
},
Err(_) => None,
}
}
fn find_current_buffer(
&mut self) {
let mut oldest_time: Option<u64> = None;
let mut next_time: Option<u64> = None;
let mut oldest_cpu: Option<usize> = None;
for i in 0..self.readers.len() {
let reader = &mut self.readers[i];
let cursor = &mut self.cursors[i];
if let Some((time, _rb)) = Self::read_time(
reader,
cursor,
&self.ring_bufs) {
match oldest_time {
Some(prev_time) => {
if time < prev_time {
next_time = oldest_time;
oldest_time = Some(time);
oldest_cpu = Some(i);
} else {
match next_time {
Some(current_next_time) => {
if time < current_next_time {
next_time = Some(time);
}
},
None => {
next_time = Some(time);
}
}
}
},
None => {
oldest_time = Some(time);
oldest_cpu = Some(i);
},
}
}
}
self.oldest_cpu = oldest_cpu;
self.next_time = next_time;
trace!(
"find_current_buffer: oldest_cpu={:?}, oldest_time={:?}, next_time={:?}",
oldest_cpu, oldest_time, next_time
);
}
}
impl PerfDataSource for RingBufDataSource {
fn enable(&mut self) -> IOResult<()> {
self.enable()
}
fn disable(&mut self) -> IOResult<()> {
self.disable()
}
fn target_pids(&self) -> Option<&[i32]> {
match &self.target_pids {
Some(pids) => { Some(&pids) },
None => { None },
}
}
fn create_bpf_files(
&mut self,
event: Option<&Event>) -> IOResult<Vec<PerfDataFile>> {
let mut files = Vec::new();
if let Some(bpf_builder) = self.bpf_builder.as_mut() {
debug!("create_bpf_files: creating BPF event buffers");
let mut common = bpf_builder.build();
let mut target_cpus = &None;
if let Some(event) = &event {
debug!("create_bpf_files: event_name={}, event_id={}", event.name(), event.id());
if event.has_no_callstack_flag() {
debug!("create_bpf_files: event has no_callstack flag, disabling callstack");
common = common.without_callstack();
}
if !event.has_no_cpu_mask_flag() {
target_cpus = &self.target_cpus;
}
}
match &self.target_pids {
None => {
Self::add_cpu_bufs(
None,
target_cpus,
&self.leader_ids,
&mut self.ring_bufs,
&common,
Some(&mut files))?;
},
Some(pids) => {
for pid in pids {
Self::add_cpu_bufs(
Some(*pid),
target_cpus,
&self.leader_ids,
&mut self.ring_bufs,
&common,
Some(&mut files))?;
}
},
}
info!("BPF files created: file_count={}", files.len());
} else {
warn!("create_bpf_files: no BPF builder configured");
}
Ok(files)
}
fn add_event(
&mut self,
event: &Event) -> IOResult<()> {
if let Some(event_builder) = self.event_builder.as_mut() {
debug!("add_event: adding event_name={}, event_id={}", event.name(), event.id());
let mut common = event_builder.build(event.id() as u64);
if event.has_no_callstack_flag() {
debug!("add_event: event has no_callstack flag, disabling callstack");
common = common.without_callstack();
}
let target_cpus = match event.has_no_cpu_mask_flag() {
true => { &None },
false => { &self.target_cpus },
};
let before: std::collections::HashSet<u64> = self.ring_bufs.keys().copied().collect();
match &self.target_pids {
None => {
Self::add_cpu_bufs(
None,
target_cpus,
&self.leader_ids,
&mut self.ring_bufs,
&common,
None)?;
},
Some(pids) => {
for pid in pids {
Self::add_cpu_bufs(
Some(*pid),
target_cpus,
&self.leader_ids,
&mut self.ring_bufs,
&common,
None)?;
}
},
}
if let Some(filter_str) = event.extension().perf_filter() {
let filter_cstr = std::ffi::CString::new(filter_str)
.map_err(|_| io_error("Invalid perf filter string"))?;
for (&id, buf) in self.ring_bufs.iter() {
if !before.contains(&id) {
if let Err(err) = buf.set_filter(&filter_cstr) {
warn!(
"add_event: failed to set perf filter on cpu={}, event={}: {}",
buf.cpu,
event.name(),
err
);
}
}
}
debug!(
"add_event: perf filter applied, event_name={}, filter={:?}",
event.name(),
filter_str
);
}
info!("Event added: event_name={}, event_id={}", event.name(), event.id());
} else {
warn!("add_event: no event builder configured");
}
Ok(())
}
fn begin_reading(&mut self) {
trace!("begin_reading: starting read cycle for {} readers", self.readers.len());
for i in 0..self.readers.len() {
let reader = &mut self.readers[i];
let cursor = &mut self.cursors[i];
reader.begin_reading(cursor);
}
self.find_current_buffer();
trace!(
"begin_reading: oldest_cpu={:?}, next_time={:?}",
self.oldest_cpu, self.next_time
);
}
fn read(
&mut self,
timeout: Duration) -> Option<PerfData<'_>> {
if self.oldest_cpu.is_none() {
trace!("read: no data available, sleeping for {:?}", timeout);
std::thread::sleep(timeout);
return None;
}
let cpu = self.oldest_cpu.unwrap();
let reader = &self.readers[cpu];
let cursor = &mut self.cursors[cpu];
let ancillary: AncillaryData;
match Self::read_time(
reader,
cursor,
&self.ring_bufs) {
Some((time, rb)) => {
if let Some(next_time) = self.next_time {
if time > next_time {
trace!(
"read: time {} exceeds next_time {}, switching buffers",
time, next_time
);
return None;
}
}
ancillary = rb.ancillary();
},
None => {
trace!("read: no more data in buffer for cpu={}", cpu);
return None;
}
}
match reader.read(
cursor,
&mut self.temp) {
Ok(raw_data) => {
trace!("read: data read from cpu={}, size={}", cpu, raw_data.len());
let perf_data = PerfData {
ancillary,
raw_data,
};
Some(perf_data)
},
Err(e) => {
warn!("read failed: cpu={}, error={}", cpu, e);
None
},
}
}
fn end_reading(&mut self) {
if let Some(oldest_cpu) = self.oldest_cpu {
trace!("end_reading: completing read for cpu={}", oldest_cpu);
let reader = &mut self.readers[oldest_cpu];
let cursor = &mut self.cursors[oldest_cpu];
reader.end_reading(cursor);
}
}
fn more(&self) -> bool {
if self.oldest_cpu.is_some() {
return true;
}
self.enabled
}
fn begin_reading_cpu(
&mut self,
cpu: u32) -> Option<CpuReadContext> {
let Some(Some(index)) = self.cpu_to_reader_index.get(cpu as usize) else {
return None;
};
let index = *index;
let reader = &self.readers[index];
let cursor = &mut self.cursors[index];
reader.begin_reading(cursor);
trace!("begin_reading_cpu: primed cpu={}, index={}", cpu, index);
Some(CpuReadContext::new(cpu, index))
}
fn read_cpu(
&mut self,
ctx: &CpuReadContext) -> Option<PerfData<'_>> {
let index = ctx.reader_index();
let reader = &self.readers[index];
let cursor = &mut self.cursors[index];
let ancillary = match Self::read_time(
reader,
cursor,
&self.ring_bufs) {
Some((_time, rb)) => rb.ancillary(),
None => {
trace!("read_cpu: no data for cpu={}", ctx.cpu());
return None;
}
};
match reader.read(
cursor,
&mut self.temp) {
Ok(raw_data) => {
trace!("read_cpu: data read from cpu={}, size={}", ctx.cpu(), raw_data.len());
Some(PerfData {
ancillary,
raw_data,
})
},
Err(e) => {
warn!("read_cpu failed: cpu={}, error={}", ctx.cpu(), e);
None
},
}
}
fn end_reading_cpu(
&mut self,
ctx: &CpuReadContext) {
let index = ctx.reader_index();
trace!("end_reading_cpu: completing read for cpu={}", ctx.cpu());
let reader = &mut self.readers[index];
let cursor = &mut self.cursors[index];
reader.end_reading(cursor);
}
#[cfg(target_os = "linux")]
fn perf_fds(&self) -> Vec<(u32, BorrowedFd<'_>)> {
let mut cpus: Vec<u32> = self.leader_ids.keys().copied().collect();
cpus.sort_unstable();
let mut fds = Vec::with_capacity(cpus.len());
for cpu in cpus {
let leader_id = self.leader_ids[&cpu];
if let Some(buf) = self.ring_bufs.get(&leader_id) {
if let Some(fd) = buf.borrowed_fd() {
fds.push((cpu, fd));
}
}
}
fds
}
fn take_in_process_writer(
&mut self) -> Option<InProcessRingBufWriter> {
self.in_process_ring_buf.as_mut().map(|rb| rb.writer())
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
#[test]
fn config() {
let kernel = RingBufBuilder::for_kernel()
.with_executable_mmap_records()
.with_comm_records()
.with_task_records()
.with_cswitch_records();
let freq = 1000;
let profiling = RingBufBuilder::for_profiling(
freq)
.with_callchain_data();
let _builder = RingBufSessionBuilder::new()
.with_page_count(1)
.with_kernel_events(kernel)
.with_profiling_events(profiling);
}
#[test]
fn read_time_skips_unknown_ring_buffer_id() {
const TEST_RECORD_TYPE: u32 = 1024;
let mut ring_buf = InProcessRingBuf::new(1);
let mut writer = ring_buf.writer();
let reader = ring_buf.create_reader();
let mut cursor = CpuRingCursor::default();
let mut record = Vec::new();
let mut payload = Vec::new();
let unknown_time = 10u64;
let unknown_id = 99u64;
payload.extend_from_slice(&unknown_time.to_ne_bytes());
payload.extend_from_slice(&unknown_id.to_ne_bytes());
abi::Header::write(TEST_RECORD_TYPE, 0, &payload, &mut record);
writer.write(&record);
record.clear();
payload.clear();
let known_time = 20u64;
let known_id = 7u64;
payload.extend_from_slice(&known_time.to_ne_bytes());
payload.extend_from_slice(&known_id.to_ne_bytes());
abi::Header::write(TEST_RECORD_TYPE, 0, &payload, &mut record);
writer.write(&record);
reader.begin_reading(&mut cursor);
let mut ring_bufs = std::collections::HashMap::new();
let common_attrs = std::rc::Rc::new(RingBufBuilder::common_attributes());
ring_bufs.insert(known_id, CpuRingBuf::new(0, common_attrs));
let first = RingBufDataSource::read_time(
&reader,
&mut cursor,
&ring_bufs);
assert!(first.is_none());
assert_eq!(24, cursor.start());
let (time, _buf) = RingBufDataSource::read_time(
&reader,
&mut cursor,
&ring_bufs).unwrap();
assert_eq!(known_time, time);
assert_eq!(24, cursor.start());
}
#[test]
#[ignore]
fn profile() {
let freq = 1000;
let profiling = RingBufBuilder::for_profiling(
freq)
.with_callchain_data();
let mut session = RingBufSessionBuilder::new()
.with_page_count(8)
.with_profiling_events(profiling)
.build()
.unwrap();
session.set_read_timeout(Duration::from_millis(0));
let samples = Arc::new(AtomicUsize::new(0));
let callback_samples = samples.clone();
let time_data = session.time_data_ref();
let ancillary = session.ancillary_data();
let prof_event = session.cpu_profile_event();
let atomic_time = Arc::new(AtomicUsize::new(0));
prof_event.add_callback(move |data| {
let full_data = data.full_data();
let time = time_data.try_get_u64(full_data).unwrap() as usize;
let prev = atomic_time.load(Ordering::Relaxed);
let mut cpu: u32 = 0;
ancillary.read(|ancillary| {
cpu = ancillary.cpu();
});
assert!(time >= prev);
callback_samples.fetch_add(1, Ordering::Relaxed);
atomic_time.store(time, Ordering::Relaxed);
Ok(())
});
session.enable().unwrap();
let now = std::time::Instant::now();
while now.elapsed().as_millis() < 100 {
}
session.disable().unwrap();
let now = std::time::Instant::now();
session.parse_all().unwrap();
println!("Took {}us", now.elapsed().as_micros());
let count = samples.load(Ordering::Relaxed);
println!("Got {} samples", count);
assert!(count >= 100);
}
#[test]
fn wakeup_watermark_rejected_when_equal_to_ring_size() {
let pages = 1;
let ring_bytes = ring_data_bytes(pages);
let watermark = ring_bytes as u32;
let err = match RingBufSessionBuilder::new()
.with_page_count(pages)
.with_wakeup_watermark(watermark)
.build()
{
Ok(_) => panic!("build() must reject watermark == ring size"),
Err(err) => err,
};
let expected = format!(
"wakeup_watermark ({} bytes) must be smaller than the per-CPU \
ring data area ({} bytes from {} pages); the kernel would \
never wake the perf fds",
watermark, ring_bytes, pages);
assert_eq!(format!("{}", err), expected);
assert_eq!(err.kind(), std::io::ErrorKind::Other);
}
#[test]
fn wakeup_watermark_rejected_when_above_ring_size() {
let pages = 1;
let ring_bytes = ring_data_bytes(pages);
let watermark = (ring_bytes + 1) as u32;
let err = match RingBufSessionBuilder::new()
.with_page_count(pages)
.with_wakeup_watermark(watermark)
.build()
{
Ok(_) => panic!("build() must reject watermark > ring size"),
Err(err) => err,
};
let expected = format!(
"wakeup_watermark ({} bytes) must be smaller than the per-CPU \
ring data area ({} bytes from {} pages); the kernel would \
never wake the perf fds",
watermark, ring_bytes, pages);
assert_eq!(format!("{}", err), expected);
assert_eq!(err.kind(), std::io::ErrorKind::Other);
}
#[test]
fn wakeup_watermark_accepted_when_below_ring_size() {
let pages = 1;
let ring_bytes = ring_data_bytes(pages);
let watermark = (ring_bytes - 1) as u32;
if let Err(err) = RingBufSessionBuilder::new()
.with_page_count(pages)
.with_wakeup_watermark(watermark)
.build()
{
let msg = format!("{}", err);
assert!(
!msg.contains("wakeup_watermark"),
"watermark < ring size must not trigger validation error, \
got: {}",
msg);
}
}
}