diff --git a/core/common/src/buffer_type.rs b/core/common/src/buffer_type.rs index 7124e55..1735168 100644 --- a/core/common/src/buffer_type.rs +++ b/core/common/src/buffer_type.rs @@ -1,20 +1,25 @@ -#[cfg(feature = "monitoring-structs")] -use crate::otel_metrics::Metrics; +//! eBPF data structures and buffer size definitions. +//! +//! This module contains: +//! - C-compatible structs emitted by the eBPF programs (`PacketLossMetrics`, `SchedStatWait`, etc.). +//! - [`BufferSize`] for pre-allocating per-CPU byte buffers. +//! - [`IpProtocols`] for L4 protocol reconstruction. +//! +//! The consumer logic has been moved to [`crate::consumer`]. + +use aya::maps::perf::PerfEventArrayBuffer; #[cfg(feature = "buffer-reader")] use aya::maps::{MapData, PerfEventArray}; -use aya::{maps::perf::PerfEventArrayBuffer, util::online_cpus}; +use aya::util::online_cpus; use bytemuck_derive::Zeroable; use bytes::BytesMut; use std::net::Ipv4Addr; -#[cfg(feature = "buffer-reader")] -#[cfg(feature = "monitoring-structs")] -use std::sync::Arc; use tracing::{error, info, warn}; -// -// IpProtocols enum to reconstruct the packet protocol based on the -// IPV4 Header Protocol code -// +/// +/// IpProtocols enum to reconstruct the packet protocol based on the +/// IPV4 Header Protocol code +/// #[derive(Debug)] #[repr(u8)] @@ -24,11 +29,11 @@ pub enum IpProtocols { UDP = 17, } -// -// TryFrom Trait implementation for IpProtocols enum -// This is used to reconstruct the packet protocol based on the -// IPV4 Header Protocol code -// +/// +/// TryFrom Trait implementation for IpProtocols enum +/// This is used to reconstruct the packet protocol based on the +/// IPV4 Header Protocol code +/// impl TryFrom for IpProtocols { type Error = (); @@ -42,10 +47,10 @@ impl TryFrom for IpProtocols { } } -// -// Structure PacketLog -//This structure is used to store the packet information -// +/// +/// Structure PacketLog +/// This structure is used to store the packet information +/// #[cfg(feature = "network-structs")] #[repr(C)] #[derive(Clone, Copy, Zeroable)] @@ -91,11 +96,11 @@ pub struct TcpPacketRegistry { unsafe impl aya::Pod for TcpPacketRegistry {} #[cfg(feature = "monitoring-structs")] -pub const TASK_COMM_LEN: usize = 16; // linux/sched.h +pub const TASK_COMM_LEN: usize = 16; #[cfg(feature = "monitoring-structs")] #[repr(C, packed)] #[derive(Clone, Copy, Zeroable)] -pub struct NetworkMetrics { +pub struct PacketLossMetrics { pub tgid: u32, pub comm: [u8; TASK_COMM_LEN], pub ts_us: u64, @@ -108,7 +113,7 @@ pub struct NetworkMetrics { pub sk_drops: i32, // Offset 136 } #[cfg(feature = "monitoring-structs")] -unsafe impl aya::Pod for NetworkMetrics {} +unsafe impl aya::Pod for PacketLossMetrics {} #[cfg(feature = "monitoring-structs")] #[repr(C, packed)] @@ -132,12 +137,11 @@ unsafe impl aya::Pod for TimeStampMetrics {} #[repr(C, packed)] #[derive(Clone, Copy, Zeroable)] pub struct CpuFrequency { - //pub cpu_id: u32, - //pub cpu_freq: u32, pub bytes_alloc: u32, pub pid: u32, pub command: [u8; 16], } +#[cfg(feature = "monitoring-structs")] unsafe impl aya::Pod for CpuFrequency {} #[cfg(feature = "monitoring-structs")] @@ -184,727 +188,21 @@ pub struct CpuIdle { #[cfg(feature = "monitoring-structs")] unsafe impl aya::Pod for CpuIdle {} -// docs: -// This function perform a byte swap from little-endian to big-endian -// It's used to reconstruct the correct IPv4 address from the u32 representation -// -// Takes a u32 address in big-endian format and returns a Ipv4Addr with reversed octets -// +/// Perform a byte swap from little-endian to big-endian. +/// +/// Used to reconstruct the correct IPv4 address from the u32 representation. +/// Takes a `u32` address in big-endian format and returns an [`Ipv4Addr`] with reversed octets. #[inline(always)] pub fn reverse_be_addr(addr: u32) -> Ipv4Addr { let octects = addr.to_be_bytes(); let [a, b, c, d] = [octects[3], octects[2], octects[1], octects[0]]; - let reversed_ip = Ipv4Addr::new(a, b, c, d); - reversed_ip -} - -// enum BuffersType -#[cfg(feature = "buffer-reader")] -pub enum BufferType { - #[cfg(feature = "network-structs")] - PacketLog, - #[cfg(feature = "network-structs")] - TcpPacketRegistry, - #[cfg(feature = "network-structs")] - VethLog, - #[cfg(feature = "monitoring-structs")] - NetworkMetrics, - #[cfg(feature = "monitoring-structs")] - TimeStampMetrics, - #[cfg(feature = "monitoring-structs")] - CpuFrequency, - #[cfg(feature = "monitoring-structs")] - MemAlloc, - #[cfg(feature = "monitoring-structs")] - SchedStatWait, - #[cfg(feature = "monitoring-structs")] - SchedStatRuntime, - #[cfg(feature = "monitoring-structs")] - CpuIdle, -} - -#[cfg(feature = "buffer-reader")] -impl BufferType { - #[cfg(feature = "network-structs")] - pub async fn read_packet_log(buffers: &mut [BytesMut], tot_events: i32, offset: i32) { - for i in offset..tot_events { - let vec_bytes = &buffers[i as usize]; - if vec_bytes.len() < std::mem::size_of::() { - error!( - "Corrupted Packet log data. Raw data: {}. Readed {} bytes expected {} bytes", - vec_bytes - .iter() - .map(|b| format!("{:02x}", b)) - .collect::>() - .join(" "), - vec_bytes.len(), - std::mem::size_of::() - ); - continue; - } - if vec_bytes.len() >= std::mem::size_of::() { - let pl: PacketLog = - unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; // reading raw bytes - - // extracting struct info from bytes - let src_ip = reverse_be_addr(pl.src_ip); - let dst_ip = reverse_be_addr(pl.dst_ip); - let src_port = u16::from_be(pl.src_port); - let dst_port = u16::from_be(pl.dst_port); - let event_id = pl.pid; - let protocol = pl.proto; - - // protocol extraction - match IpProtocols::try_from(protocol) { - Ok(proto) => { - info!( - "Event Id: {} Protocol: {:?} SRC: {}:{} -> DST: {}:{}", - event_id, proto, src_ip, src_port, dst_ip, dst_port - ); - } - Err(e) => { - error!("Unknown protocol. Data maybe corrupted. Reason:{:?}", e); - } - } - } - } - } - #[cfg(feature = "network-structs")] - pub async fn read_tcp_registry_log(buffers: &mut [BytesMut], tot_events: i32, offset: i32) { - for i in offset..tot_events { - let vec_bytes = &buffers[i as usize]; - if vec_bytes.len() < std::mem::size_of::() { - error!( - "Corrupted data Tcp Registry data. Raw data: {}. Readed {} bytes expected {} bytes", - vec_bytes - .iter() - .map(|b| format!("{:02x}", b)) - .collect::>() - .join(" "), - vec_bytes.len(), - std::mem::size_of::() - ); - continue; - } - if vec_bytes.len() >= std::mem::size_of::() { - let pl: TcpPacketRegistry = - unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; // reading raw bytes - - // extracting struct info from bytes - let src = reverse_be_addr(pl.src_ip); - let dst = reverse_be_addr(pl.dst_ip); - let src_port = u16::from_be(pl.src_port); - let dst_port = u16::from_be(pl.dst_port); - let event_id = pl.pid; - let command = pl.command.to_vec(); - let end = command - .iter() - .position(|&x| x == 0) - .unwrap_or(command.len()); - let command_str = String::from_utf8_lossy(&command[..end]).to_string(); - let cgroup_id = pl.cgroup_id; - let protocol = pl.proto; - - // protocol extraction - match IpProtocols::try_from(protocol) { - Ok(proto) => { - info!( - "Event Id: {} Protocol: {:?} SRC: {}:{} -> DST: {}:{} Command: {} Cgroup_id: {}", - event_id, - proto, - src, - src_port, - dst, - dst_port, - command_str, - cgroup_id //proc_content - ); - } - Err(e) => { - error!("Unknown protocol. Data maybe corrupted. Reason:{:?}", e); - } - } - } - } - } - #[cfg(feature = "network-structs")] - pub async fn read_and_handle_veth_log(buffers: &mut [BytesMut], tot_events: i32, offset: i32) { - for i in offset..tot_events { - let vec_bytes = &buffers[i as usize]; - if vec_bytes.len() < std::mem::size_of::() { - error!( - "Corrupted data VethLog data. Raw data: {}. Readed {} bytes expected {} bytes", - vec_bytes - .iter() - .map(|b| format!("{:02x}", b)) - .collect::>() - .join(" "), - vec_bytes.len(), - std::mem::size_of::() - ); - continue; - } - if vec_bytes.len() >= std::mem::size_of::() { - let vthl: VethLog = - unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; // reading raw bytes - - // extracting struct info from bytes - let name_bytes = vthl.name; - let dev_addr_bytes = vthl.dev_addr; - let name = std::str::from_utf8(&name_bytes); - let state = vthl.state; - - let dev_addr = dev_addr_bytes; - let netns = vthl.netns; - let mut event_type = String::new(); - - // event_type extraction - match vthl.event_type { - 1 => { - event_type = "creation".to_string(); - match name { - Ok(veth_name) => { - info!( - "[{}] Veth Event: Type: {} Name: {} Dev_addr: {:x?} State: {}", - netns, - event_type, - veth_name.trim_end_matches("\0"), - dev_addr, - state - ); - } - Err(e) => { - error!( - "Failed to extract veth name during event_type = creation (1).Reason:{}", - e - ); - } - } - } - 2 => { - event_type = "deletion".to_string(); - match name { - Ok(veth_name) => { - info!( - "[{}] Veth Event: Type: {} Name: {} Dev_addr: {:x?} State: {}", - netns, - event_type, - veth_name.trim_end_matches("\0"), - dev_addr, - state - ); - } - Err(e) => { - error!( - "Failed to extract veth name during event_type = deletion (2).Reason:{}", - e - ); - } - } - } - _ => { - warn!("Unknown event type") - } - } - } - } - } - #[cfg(feature = "monitoring-structs")] - /// Continuously read [`NetworkMetrics`] events and record OpenTelemetry - /// observations. - /// - /// This helper mirrors the core behaviour of - /// [`cortexbrain_common::buffer_type::read_perf_buffer`] but adds the OTel - /// instrumentation layer. - /// - /// # Loop - /// - /// 1. For every CPU buffer call `read_events`. - /// 2. Parse each raw [`BytesMut`] into [`NetworkMetrics`] using an - /// unaligned read (the struct is `#[repr(C, packed)]` and `Pod`). - /// 3. Call [`Metrics::record_network_metrics`]. - /// 4. Retain the legacy `tracing::info!` log for human-readable local output. - /// 5. Sleep 100 ms between polls. - /// - /// # Safety - /// - /// `std::ptr::read_unaligned` is safe here because the eBPF program writes - /// exactly the `NetworkMetrics` layout into the ring buffer and the struct - /// implements [`aya::Pod`]. - /// Continuously read [`TimeStampMetrics`] events and record OpenTelemetry - /// observations. - /// - /// Counterpart to [`read_network_buffer`] for the `time_stamp_events` map. - - pub async fn read_network_metrics( - buffers: &mut [BytesMut], - tot_events: i32, - offset: i32, - exporter: &str, - metrics: Arc, - ) { - for i in offset..tot_events { - let vec_bytes = &buffers[i as usize]; - if vec_bytes.len() < std::mem::size_of::() { - error!( - "Corrupted Network Metrics data. Raw data: {}. Readed {} bytes expected {} bytes", - vec_bytes - .iter() - .map(|b| format!("{:02x}", b)) - .collect::>() - .join(" "), - vec_bytes.len(), - std::mem::size_of::() - ); - continue; - } - if vec_bytes.len() >= std::mem::size_of::() { - let net_metrics: NetworkMetrics = - unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; - - match exporter { - "otlp" => metrics.record_network_metrics(&net_metrics), - _ => continue, // skip - } - let tgid = net_metrics.tgid; - let comm = String::from_utf8_lossy(&net_metrics.comm); - let ts_us = net_metrics.ts_us; - let sk_drop_count = net_metrics.sk_drops; - let sk_err = net_metrics.sk_err; - let sk_err_soft = net_metrics.sk_err_soft; - let sk_backlog_len = net_metrics.sk_backlog_len; - let sk_write_memory_queued = net_metrics.sk_write_memory_queued; - let sk_ack_backlog = net_metrics.sk_ack_backlog; - let sk_receive_buffer_size = net_metrics.sk_receive_buffer_size; - - info!( - "tgid: {}, comm: {}, ts_us: {}, sk_drops: {}, sk_err: {}, sk_err_soft: {}, sk_backlog_len: {}, sk_write_memory_queued: {}, sk_ack_backlog: {}, sk_receive_buffer_size: {}", - tgid, - comm, - ts_us, - sk_drop_count, - sk_err, - sk_err_soft, - sk_backlog_len, - sk_write_memory_queued, - sk_ack_backlog, - sk_receive_buffer_size - ); - } - } - } - #[cfg(feature = "monitoring-structs")] - pub async fn read_timestamp_metrics( - buffers: &mut [BytesMut], - tot_events: i32, - offset: i32, - exporter: &str, - metrics: Arc, - ) { - for i in offset..tot_events { - let vec_bytes = &buffers[i as usize]; - if vec_bytes.len() < std::mem::size_of::() { - error!( - "Corrupted Network Metrics data. Raw data: {}. Readed {} bytes expected {} bytes", - vec_bytes - .iter() - .map(|b| format!("{:02x}", b)) - .collect::>() - .join(" "), - vec_bytes.len(), - std::mem::size_of::() - ); - continue; - } - if vec_bytes.len() >= std::mem::size_of::() { - let time_stamp_event: TimeStampMetrics = - unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; - - match exporter { - "otlp" => metrics.record_timestamp_metrics(&time_stamp_event), - _ => continue, - } - - let delta_us = time_stamp_event.delta_us; - let ts_us = time_stamp_event.ts_us; - let tgid = time_stamp_event.tgid; - let comm = String::from_utf8_lossy(&time_stamp_event.comm); - let lport = time_stamp_event.lport; - let dport_be = time_stamp_event.dport_be; - let af = time_stamp_event.af; - info!( - "TimeStampEvent - delta_us: {}, ts_us: {}, tgid: {}, comm: {}, lport: {}, dport_be: {}, af: {}", - delta_us, ts_us, tgid, comm, lport, dport_be, af - ); - } - } - } - - #[cfg(feature = "monitoring-structs")] - pub async fn read_cpu_frequency( - buffers: &mut [BytesMut], - tot_events: i32, - offset: i32, - exporter: &str, - metrics: Arc, - ) { - for i in offset..tot_events { - let vec_bytes = &buffers[i as usize]; - if vec_bytes.len() < std::mem::size_of::() { - error!( - "Corrupted Cpu Frequency Metrics data. Raw data: {}. Readed {} bytes expected {} bytes", - vec_bytes - .iter() - .map(|b| format!("{:02x}", b)) - .collect::>() - .join(" "), - vec_bytes.len(), - std::mem::size_of::() - ); - continue; - } - if vec_bytes.len() >= std::mem::size_of::() { - let cpu_freq_metrics: CpuFrequency = - unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; - - match exporter { - "otlp" => metrics.record_cpu_bytes_alloc(&cpu_freq_metrics), - _ => continue, - } - - //let cpu_id = cpu_freq_metrics.cpu_id; - //let cpu_freq = cpu_freq_metrics.cpu_freq; - let bytes_alloc = cpu_freq_metrics.bytes_alloc; - //info!( - // "Cpu id: {} Cpu frequency: {} Bytes alloc: {}", - // cpu_id, cpu_freq, bytes_alloc - //); - let pid = cpu_freq_metrics.pid; - let command = cpu_freq_metrics.command; - info!( - "Cpu Bytes alloc: {} pid : {} command: {:?}", - bytes_alloc, pid, command - ); - } - } - } - - #[cfg(feature = "monitoring-structs")] - pub async fn read_mem_alloc( - buffers: &mut [BytesMut], - tot_events: i32, - offset: i32, - exporter: &str, - metrics: Arc, - ) { - for i in offset..tot_events { - let vec_bytes = &buffers[i as usize]; - if vec_bytes.len() < std::mem::size_of::() { - error!( - "Corrupted MemAlloc data. Raw data: {}. Readed {} bytes expected {} bytes", - vec_bytes - .iter() - .map(|b| format!("{:02x}", b)) - .collect::>() - .join(" "), - vec_bytes.len(), - std::mem::size_of::() - ); - continue; - } - if vec_bytes.len() >= std::mem::size_of::() { - let mem_alloc: MemAlloc = - unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; - - match exporter { - "otlp" => metrics.record_enter_mem_alloc(&mem_alloc), - _ => continue, - } - - let tgid = mem_alloc.tgid; - let command = String::from_utf8_lossy(&mem_alloc.command); - let addr = mem_alloc.addr; - let length = mem_alloc.length; - - info!( - "MemAlloc - tgid: {}, command: {}, addr: {}, length: {}", - tgid, command, addr, length - ); - } - } - } - - #[cfg(feature = "monitoring-structs")] - pub async fn read_sched_stat_wait( - buffers: &mut [BytesMut], - tot_events: i32, - offset: i32, - exporter: &str, - metrics: Arc, - ) { - for i in offset..tot_events { - let vec_bytes = &buffers[i as usize]; - if vec_bytes.len() < std::mem::size_of::() { - error!( - "Corrupted SchedStatWait data. Raw data: {}. Readed {} bytes expected {} bytes", - vec_bytes - .iter() - .map(|b| format!("{:02x}", b)) - .collect::>() - .join(" "), - vec_bytes.len(), - std::mem::size_of::() - ); - continue; - } - if vec_bytes.len() >= std::mem::size_of::() { - let sched_stat_wait: SchedStatWait = - unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; - - match exporter { - "otlp" => metrics.record_sched_stat_wait(&sched_stat_wait), - _ => continue, - } - - let tgid = sched_stat_wait.tgid; - let command = String::from_utf8_lossy(&sched_stat_wait.command); - let delay = sched_stat_wait.delay; - - info!( - "SchedStatWait - tgid: {}, command: {}, delay: {}", - tgid, command, delay - ); - } - } - } - - #[cfg(feature = "monitoring-structs")] - pub async fn read_sched_stat_runtime( - buffers: &mut [BytesMut], - tot_events: i32, - offset: i32, - exporter: &str, - metrics: Arc, - ) { - for i in offset..tot_events { - let vec_bytes = &buffers[i as usize]; - if vec_bytes.len() < std::mem::size_of::() { - error!( - "Corrupted SchedStatRuntime data. Raw data: {}. Readed {} bytes expected {} bytes", - vec_bytes - .iter() - .map(|b| format!("{:02x}", b)) - .collect::>() - .join(" "), - vec_bytes.len(), - std::mem::size_of::() - ); - continue; - } - if vec_bytes.len() >= std::mem::size_of::() { - let sched_stat_runtime: SchedStatRuntime = - unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; - - match exporter { - "otlp" => metrics.record_sched_stat_runtime(&sched_stat_runtime), - _ => continue, - } - - let tgid = sched_stat_runtime.tgid; - let command = String::from_utf8_lossy(&sched_stat_runtime.command); - let runtime = sched_stat_runtime.runtime; - - info!( - "SchedStatRuntime - tgid: {}, command: {}, runtime: {}", - tgid, command, runtime - ); - } - } - } - - #[cfg(feature = "monitoring-structs")] - pub async fn read_cpu_idle( - buffers: &mut [BytesMut], - tot_events: i32, - offset: i32, - exporter: &str, - metrics: Arc, - ) { - for i in offset..tot_events { - let vec_bytes = &buffers[i as usize]; - if vec_bytes.len() < std::mem::size_of::() { - error!( - "Corrupted CpuIdle data. Raw data: {}. Readed {} bytes expected {} bytes", - vec_bytes - .iter() - .map(|b| format!("{:02x}", b)) - .collect::>() - .join(" "), - vec_bytes.len(), - std::mem::size_of::() - ); - continue; - } - if vec_bytes.len() >= std::mem::size_of::() { - let cpu_idle: CpuIdle = - unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; - - match exporter { - "otlp" => metrics.record_cpu_idle(&cpu_idle), - _ => continue, - } - - let cpu_id = cpu_idle.cpu_id; - let state = cpu_idle.state; - - info!( - "CpuIdle state changed - cpu_id: {}, state: {}", - cpu_id, state - ); - } - } - } -} - -// docs: read buffer function: -// template function that take a mut perf_event_array_buffer of type T and a mutable buffer of Vec -#[cfg(feature = "buffer-reader")] -pub async fn read_perf_buffer>( - mut array_buffers: Vec>, - mut buffers: Vec, - buffer_type: BufferType, - #[cfg(feature = "monitoring-structs")] metrics: Option>, -) { - // loop over the buffers - loop { - for buf in array_buffers.iter_mut() { - match buf.read_events(&mut buffers) { - Ok(events) => { - // triggered if some events are lost - if events.lost > 0 { - tracing::debug!("Lost events: {} ", events.lost); - } - // triggered if some events are readed - if events.read > 0 { - tracing::debug!("Readed events: {}", events.read); - let offset = 0; - let tot_events = events.read as i32; - - //read the events in the buffer - match buffer_type { - #[cfg(feature = "network-structs")] - BufferType::PacketLog => { - BufferType::read_packet_log(&mut buffers, tot_events, offset).await - } - #[cfg(feature = "network-structs")] - BufferType::TcpPacketRegistry => { - BufferType::read_tcp_registry_log(&mut buffers, tot_events, offset) - .await - } - #[cfg(feature = "network-structs")] - BufferType::VethLog => { - BufferType::read_and_handle_veth_log( - &mut buffers, - tot_events, - offset, - ) - .await - } - #[cfg(feature = "monitoring-structs")] - BufferType::NetworkMetrics => { - BufferType::read_network_metrics( - &mut buffers, - tot_events, - offset, - "otlp", - metrics - .clone() - .expect("Metrics required for NetworkMetrics"), - ) - .await - } - #[cfg(feature = "monitoring-structs")] - BufferType::TimeStampMetrics => { - BufferType::read_timestamp_metrics( - &mut buffers, - tot_events, - offset, - "otlp", - metrics - .clone() - .expect("Metric required for TimeStampMetrics"), - ) - .await - } - #[cfg(feature = "monitoring-structs")] - BufferType::CpuFrequency => { - BufferType::read_cpu_frequency( - &mut buffers, - tot_events, - offset, - "otlp", - metrics.clone().expect("Metric required for CpuFrequency"), - ) - .await - } - #[cfg(feature = "monitoring-structs")] - BufferType::MemAlloc => { - BufferType::read_mem_alloc( - &mut buffers, - tot_events, - offset, - "otlp", - metrics.clone().expect("Metric required for MemAlloc"), - ) - .await - } - #[cfg(feature = "monitoring-structs")] - BufferType::SchedStatWait => { - BufferType::read_sched_stat_wait( - &mut buffers, - tot_events, - offset, - "otlp", - metrics.clone().expect("Metric required for SchedStatWait"), - ) - .await - } - #[cfg(feature = "monitoring-structs")] - BufferType::SchedStatRuntime => { - BufferType::read_sched_stat_runtime( - &mut buffers, - tot_events, - offset, - "otlp", - metrics - .clone() - .expect("Metric required for SchedStatRuntime"), - ) - .await - } - #[cfg(feature = "monitoring-structs")] - BufferType::CpuIdle => { - BufferType::read_cpu_idle( - &mut buffers, - tot_events, - offset, - "otlp", - metrics.clone().expect("Metric required for CpuIdle"), - ) - .await - } - } - } - } - Err(e) => { - error!("Cannot read events from buffer. Reason: {} ", e); - } - } - } - tokio::time::sleep(std::time::Duration::from_millis(100)).await; // small sleep - } + Ipv4Addr::new(a, b, c, d) } +/// Buffer size presets for per-CPU perf-buffer allocation. +/// +/// Each variant carries a multiplier that determines how many struct-sized +/// slots are pre-allocated per CPU in [`BufferSize::set_buffer`]. #[cfg(feature = "buffer-reader")] pub enum BufferSize { #[cfg(feature = "network-structs")] @@ -928,8 +226,10 @@ pub enum BufferSize { #[cfg(feature = "monitoring-structs")] CpuIdle, } + #[cfg(feature = "buffer-reader")] impl BufferSize { + /// Return the size in bytes of the struct associated with this variant. pub fn get_size(&self) -> usize { match self { #[cfg(feature = "network-structs")] @@ -939,7 +239,7 @@ impl BufferSize { #[cfg(feature = "network-structs")] BufferSize::TcpEvents => std::mem::size_of::(), #[cfg(feature = "monitoring-structs")] - BufferSize::NetworkMetricsEvents => std::mem::size_of::(), + BufferSize::NetworkMetricsEvents => std::mem::size_of::(), #[cfg(feature = "monitoring-structs")] BufferSize::TimeMetricsEvents => std::mem::size_of::(), #[cfg(feature = "monitoring-structs")] @@ -954,21 +254,14 @@ impl BufferSize { BufferSize::CpuIdle => std::mem::size_of::(), } } - pub fn set_buffer(&self) -> Vec { - // iter returns and iterator of cpu ids, - // we need only the total number of cpus to set the buffer size so we use .len() to get - // the count of total cpus and then we allocate a buffer for each cpu with a capacity - // based on the structure size * a factor to have a bigger buffer to avoid overflows and lost events - // Old buffers where 1024 bytes long. Now we set different buffer size based on - // the frequence of the events. - // ClassifierNetEvents are triggered by the TC classifier program, events has high frequency - // VethEvents are triggered by the creation and deletion of veth interfaces, events has small frequency compared to classifier events - // TcpEvents are triggered by TCP events and connections. Events has similar frequency to ClassifierNetEvents. + /// Allocate one `BytesMut` per CPU with capacity tuned to the event type. + pub fn set_buffer(&self) -> Vec { + use aya::util::online_cpus; - let tot_cpu = online_cpus().iter().len(); // total number of cpus + let tot_cpu = online_cpus().iter().len(); - // TODO: finish to do all the calculations for the buffer sizes + // TODO: finish buffer size calculations match self { #[cfg(feature = "network-structs")] BufferSize::ClassifierNetEvents => { @@ -977,7 +270,7 @@ impl BufferSize { } #[cfg(feature = "network-structs")] BufferSize::VethEvents => { - let capacity = self.get_size() * 100; // Allocates 4Kb of memory for the buffers + let capacity = self.get_size() * 100; return vec![BytesMut::with_capacity(capacity); tot_cpu]; } #[cfg(feature = "network-structs")] @@ -1024,11 +317,10 @@ impl BufferSize { } } +/// Open a [`PerfEventArrayBuffer`] for every online CPU and append them to `vec_of_buffers`. #[cfg(feature = "buffer-reader")] pub fn fill_buffers( - //buf: PerfEventArrayBuffer, mut vec_of_buffers: Vec>, - //buffers: Vec, mut events_array: PerfEventArray, ) -> Vec> { for cpu_id in online_cpus() diff --git a/core/common/src/consumer.rs b/core/common/src/consumer.rs new file mode 100644 index 0000000..bd7d267 --- /dev/null +++ b/core/common/src/consumer.rs @@ -0,0 +1,779 @@ +//! Perf-buffer consumers for eBPF events. +//! +//! This module provides the [`Consumer`] enum and its associated `read` methods +//! that parse raw bytes from [`aya::maps::perf::PerfEventArrayBuffer`] into +//! strongly-typed eBPF structs and forward them to the OpenTelemetry metrics +//! pipeline via [`crate::otel_metrics::Metrics`]. +//! +//! Each consumer method: +//! 1. Validates the raw byte buffer length against the expected struct size. +//! 2. Performs an unaligned read into the `#[repr(C, packed)]` struct. +//! 3. Builds [`crate::metadata::Metadata`] (with optional Docker/K8s enrichment). +//! 4. Records the observation through [`Metrics::record_*`]. + +#[cfg(feature = "monitoring-structs")] +use crate::buffer_type::{ + CpuFrequency, CpuIdle, MemAlloc, PacketLossMetrics, SchedStatRuntime, SchedStatWait, + TimeStampMetrics, +}; +#[cfg(feature = "network-structs")] +use crate::buffer_type::{PacketLog, TcpPacketRegistry, VethLog}; +#[cfg(feature = "monitoring-structs")] +use crate::metadata::Metadata; +#[cfg(feature = "monitoring-structs")] +use crate::otel_metrics::Metrics; +use bytes::BytesMut; +#[cfg(feature = "monitoring-structs")] +use std::sync::Arc; +use tracing::{error, info, warn}; + +/// Discriminator for perf-buffer event types consumed by the collector. +/// +/// Each variant maps to an eBPF program output struct and a dedicated +/// `read_*` method that knows how to parse it. +#[cfg(feature = "buffer-reader")] +pub enum Consumer { + #[cfg(feature = "network-structs")] + PacketLog, + #[cfg(feature = "network-structs")] + TcpPacketRegistry, + #[cfg(feature = "network-structs")] + VethLog, + #[cfg(feature = "monitoring-structs")] + PacketLossMetrics, + #[cfg(feature = "monitoring-structs")] + TimeStampMetrics, + #[cfg(feature = "monitoring-structs")] + CpuFrequency, + #[cfg(feature = "monitoring-structs")] + MemAlloc, + #[cfg(feature = "monitoring-structs")] + SchedStatWait, + #[cfg(feature = "monitoring-structs")] + SchedStatRuntime, + #[cfg(feature = "monitoring-structs")] + CpuIdle, +} + +#[cfg(feature = "buffer-reader")] +impl Consumer { + /// Read and log [`PacketLog`] events from the perf buffer. + /// + /// Parses IPv4 addresses, ports and L4 protocol from raw eBPF bytes and + /// emits human-readable `tracing::info!` lines. + #[cfg(feature = "network-structs")] + pub async fn read_packet_log(buffers: &mut [BytesMut], tot_events: i32, offset: i32) { + use crate::buffer_type::{IpProtocols, reverse_be_addr}; + + for i in offset..tot_events { + let vec_bytes = &buffers[i as usize]; + if vec_bytes.len() < std::mem::size_of::() { + error!( + "Corrupted Packet log data. Raw data: {}. Readed {} bytes expected {} bytes", + vec_bytes + .iter() + .map(|b| format!("{:02x}", b)) + .collect::>() + .join(" "), + vec_bytes.len(), + std::mem::size_of::() + ); + continue; + } + if vec_bytes.len() >= std::mem::size_of::() { + let pl: PacketLog = + unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; + + let src_ip = reverse_be_addr(pl.src_ip); + let dst_ip = reverse_be_addr(pl.dst_ip); + let src_port = u16::from_be(pl.src_port); + let dst_port = u16::from_be(pl.dst_port); + let event_id = pl.pid; + let protocol = pl.proto; + + match IpProtocols::try_from(protocol) { + Ok(proto) => { + info!( + "Event Id: {} Protocol: {:?} SRC: {}:{} -> DST: {}:{}", + event_id, proto, src_ip, src_port, dst_ip, dst_port + ); + } + Err(e) => { + error!("Unknown protocol. Data maybe corrupted. Reason:{:?}", e); + } + } + } + } + } + + /// Read and log [`TcpPacketRegistry`] events from the perf buffer. + /// + /// Similar to [`read_packet_log`] but additionally prints the command name + /// and cgroup ID extracted from the eBPF struct. + #[cfg(feature = "network-structs")] + pub async fn read_tcp_registry_log(buffers: &mut [BytesMut], tot_events: i32, offset: i32) { + use crate::buffer_type::{IpProtocols, reverse_be_addr}; + + for i in offset..tot_events { + let vec_bytes = &buffers[i as usize]; + if vec_bytes.len() < std::mem::size_of::() { + error!( + "Corrupted data Tcp Registry data. Raw data: {}. Readed {} bytes expected {} bytes", + vec_bytes + .iter() + .map(|b| format!("{:02x}", b)) + .collect::>() + .join(" "), + vec_bytes.len(), + std::mem::size_of::() + ); + continue; + } + if vec_bytes.len() >= std::mem::size_of::() { + let pl: TcpPacketRegistry = + unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; + + let src = reverse_be_addr(pl.src_ip); + let dst = reverse_be_addr(pl.dst_ip); + let src_port = u16::from_be(pl.src_port); + let dst_port = u16::from_be(pl.dst_port); + let event_id = pl.pid; + let command = pl.command.to_vec(); + let end = command + .iter() + .position(|&x| x == 0) + .unwrap_or(command.len()); + let command_str = String::from_utf8_lossy(&command[..end]).to_string(); + let cgroup_id = pl.cgroup_id; + let protocol = pl.proto; + + match IpProtocols::try_from(protocol) { + Ok(proto) => { + info!( + "Event Id: {} Protocol: {:?} SRC: {}:{} -> DST: {}:{} Command: {} Cgroup_id: {}", + event_id, proto, src, src_port, dst, dst_port, command_str, cgroup_id + ); + } + Err(e) => { + error!("Unknown protocol. Data maybe corrupted. Reason:{:?}", e); + } + } + } + } + } + + /// Read and log [`VethLog`] events from the perf buffer. + /// + /// Distinguishes between veth interface creation (event_type == 1) and + /// deletion (event_type == 2) and logs the interface name, MAC address and + /// netns. + #[cfg(feature = "network-structs")] + pub async fn read_and_handle_veth_log(buffers: &mut [BytesMut], tot_events: i32, offset: i32) { + for i in offset..tot_events { + let vec_bytes = &buffers[i as usize]; + if vec_bytes.len() < std::mem::size_of::() { + error!( + "Corrupted data VethLog data. Raw data: {}. Readed {} bytes expected {} bytes", + vec_bytes + .iter() + .map(|b| format!("{:02x}", b)) + .collect::>() + .join(" "), + vec_bytes.len(), + std::mem::size_of::() + ); + continue; + } + if vec_bytes.len() >= std::mem::size_of::() { + let vthl: VethLog = + unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; + + let name_bytes = vthl.name; + let dev_addr_bytes = vthl.dev_addr; + let name = std::str::from_utf8(&name_bytes); + let state = vthl.state; + let dev_addr = dev_addr_bytes; + let netns = vthl.netns; + let mut event_type = String::new(); + + match vthl.event_type { + 1 => { + event_type = "creation".to_string(); + match name { + Ok(veth_name) => { + info!( + "[{}] Veth Event: Type: {} Name: {} Dev_addr: {:x?} State: {}", + netns, + event_type, + veth_name.trim_end_matches("\0"), + dev_addr, + state + ); + } + Err(e) => { + error!( + "Failed to extract veth name during event_type = creation (1).Reason:{}", + e + ); + } + } + } + 2 => { + event_type = "deletion".to_string(); + match name { + Ok(veth_name) => { + info!( + "[{}] Veth Event: Type: {} Name: {} Dev_addr: {:x?} State: {}", + netns, + event_type, + veth_name.trim_end_matches("\0"), + dev_addr, + state + ); + } + Err(e) => { + error!( + "Failed to extract veth name during event_type = deletion (2).Reason:{}", + e + ); + } + } + } + _ => { + warn!("Unknown event type") + } + } + } + } + } + + /// Read [`PacketLossMetrics`] events and record OpenTelemetry observations. + /// + /// # Arguments + /// - `buffers` — raw byte buffers populated by `PerfEventArrayBuffer::read_events`. + /// - `tot_events` — number of events to process. + /// - `offset` — start index in `buffers`. + /// - `exporter` — `"otlp"` forwards to [`Metrics`], any other value is skipped. + /// - `metrics` — shared [`Metrics`] handle. + /// + /// # Safety + /// Uses `std::ptr::read_unaligned` on `#[repr(C, packed)]` structs that implement [`aya::Pod`]. + /// + /// # Metadata enrichment + /// If `exporter == "otlp"`, constructs [`Metadata`] and calls [`Metadata::enrich()`] to resolve + /// Docker container name from `/proc//cgroup` when available. + #[cfg(feature = "monitoring-structs")] + pub async fn read_packet_loss_metrics( + buffers: &mut [BytesMut], + tot_events: i32, + offset: i32, + exporter: &str, + metrics: Arc, + ) { + for i in offset..tot_events { + let vec_bytes = &buffers[i as usize]; + if vec_bytes.len() < std::mem::size_of::() { + error!( + "Corrupted Network Metrics data. Raw data: {}. Readed {} bytes expected {} bytes", + vec_bytes + .iter() + .map(|b| format!("{:02x}", b)) + .collect::>() + .join(" "), + vec_bytes.len(), + std::mem::size_of::() + ); + continue; + } + if vec_bytes.len() >= std::mem::size_of::() { + let packet_loss: PacketLossMetrics = + unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; + + match exporter { + "otlp" => { + let mut metadata = + Metadata::from_ebpf(Some(packet_loss.tgid), &packet_loss.comm); + metadata.enrich(); + metrics.record_packet_loss_metrics(&packet_loss, &metadata); + } + _ => continue, + } + + let tgid = packet_loss.tgid; + let comm = String::from_utf8_lossy(&packet_loss.comm); + let ts_us = packet_loss.ts_us; + let sk_drop_count = packet_loss.sk_drops; + let sk_err = packet_loss.sk_err; + let sk_err_soft = packet_loss.sk_err_soft; + let sk_backlog_len = packet_loss.sk_backlog_len; + let sk_write_memory_queued = packet_loss.sk_write_memory_queued; + let sk_ack_backlog = packet_loss.sk_ack_backlog; + let sk_receive_buffer_size = packet_loss.sk_receive_buffer_size; + + info!( + "tgid: {}, comm: {}, ts_us: {}, sk_drops: {}, sk_err: {}, sk_err_soft: {}, sk_backlog_len: {}, sk_write_memory_queued: {}, sk_ack_backlog: {}, sk_receive_buffer_size: {}", + tgid, + comm, + ts_us, + sk_drop_count, + sk_err, + sk_err_soft, + sk_backlog_len, + sk_write_memory_queued, + sk_ack_backlog, + sk_receive_buffer_size + ); + } + } + } + + /// Read [`TimeStampMetrics`] events and record OpenTelemetry observations. + /// + /// Counterpart to [`read_packet_loss_metrics`] for the `time_stamp_events` map. + #[cfg(feature = "monitoring-structs")] + pub async fn read_timestamp_metrics( + buffers: &mut [BytesMut], + tot_events: i32, + offset: i32, + exporter: &str, + metrics: Arc, + ) { + for i in offset..tot_events { + let vec_bytes = &buffers[i as usize]; + if vec_bytes.len() < std::mem::size_of::() { + error!( + "Corrupted Network Metrics data. Raw data: {}. Readed {} bytes expected {} bytes", + vec_bytes + .iter() + .map(|b| format!("{:02x}", b)) + .collect::>() + .join(" "), + vec_bytes.len(), + std::mem::size_of::() + ); + continue; + } + if vec_bytes.len() >= std::mem::size_of::() { + let time_stamp_event: TimeStampMetrics = + unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; + + match exporter { + "otlp" => { + let mut metadata = Metadata::from_ebpf( + Some(time_stamp_event.tgid), + &time_stamp_event.comm, + ); + metadata.enrich(); + metrics.record_timestamp_metrics(&time_stamp_event, &metadata); + } + _ => continue, + } + + let delta_us = time_stamp_event.delta_us; + let ts_us = time_stamp_event.ts_us; + let tgid = time_stamp_event.tgid; + let comm = String::from_utf8_lossy(&time_stamp_event.comm); + let lport = time_stamp_event.lport; + let dport_be = time_stamp_event.dport_be; + let af = time_stamp_event.af; + info!( + "TimeStampEvent - delta_us: {}, ts_us: {}, tgid: {}, comm: {}, lport: {}, dport_be: {}, af: {}", + delta_us, ts_us, tgid, comm, lport, dport_be, af + ); + } + } + } + + /// Read [`CpuFrequency`] events and record OpenTelemetry observations. + #[cfg(feature = "monitoring-structs")] + pub async fn read_cpu_frequency( + buffers: &mut [BytesMut], + tot_events: i32, + offset: i32, + exporter: &str, + metrics: Arc, + ) { + for i in offset..tot_events { + let vec_bytes = &buffers[i as usize]; + if vec_bytes.len() < std::mem::size_of::() { + error!( + "Corrupted Cpu Frequency Metrics data. Raw data: {}. Readed {} bytes expected {} bytes", + vec_bytes + .iter() + .map(|b| format!("{:02x}", b)) + .collect::>() + .join(" "), + vec_bytes.len(), + std::mem::size_of::() + ); + continue; + } + if vec_bytes.len() >= std::mem::size_of::() { + let cpu_freq_metrics: CpuFrequency = + unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; + + match exporter { + "otlp" => { + let mut metadata = Metadata::from_ebpf( + Some(cpu_freq_metrics.pid), + &cpu_freq_metrics.command, + ); + metadata.enrich(); + metrics.record_cpu_bytes_alloc(&cpu_freq_metrics, &metadata); + } + _ => continue, + } + + let bytes_alloc = cpu_freq_metrics.bytes_alloc; + let pid = cpu_freq_metrics.pid; + let command = cpu_freq_metrics.command; + info!( + "Cpu Bytes alloc: {} pid : {} command: {:?}", + bytes_alloc, pid, command + ); + } + } + } + + /// Read [`MemAlloc`] events and record OpenTelemetry observations. + #[cfg(feature = "monitoring-structs")] + pub async fn read_mem_alloc( + buffers: &mut [BytesMut], + tot_events: i32, + offset: i32, + exporter: &str, + metrics: Arc, + ) { + for i in offset..tot_events { + let vec_bytes = &buffers[i as usize]; + if vec_bytes.len() < std::mem::size_of::() { + error!( + "Corrupted MemAlloc data. Raw data: {}. Readed {} bytes expected {} bytes", + vec_bytes + .iter() + .map(|b| format!("{:02x}", b)) + .collect::>() + .join(" "), + vec_bytes.len(), + std::mem::size_of::() + ); + continue; + } + if vec_bytes.len() >= std::mem::size_of::() { + let mem_alloc: MemAlloc = + unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; + + match exporter { + "otlp" => { + let mut metadata = + Metadata::from_ebpf(Some(mem_alloc.tgid), &mem_alloc.command); + metadata.enrich(); + metrics.record_enter_mem_alloc(&mem_alloc, &metadata); + } + _ => continue, + } + + let tgid = mem_alloc.tgid; + let command = String::from_utf8_lossy(&mem_alloc.command); + let addr = mem_alloc.addr; + let length = mem_alloc.length; + + info!( + "MemAlloc - tgid: {}, command: {}, addr: {}, length: {}", + tgid, command, addr, length + ); + } + } + } + + /// Read [`SchedStatWait`] events and record OpenTelemetry observations. + #[cfg(feature = "monitoring-structs")] + pub async fn read_sched_stat_wait( + buffers: &mut [BytesMut], + tot_events: i32, + offset: i32, + exporter: &str, + metrics: Arc, + ) { + for i in offset..tot_events { + let vec_bytes = &buffers[i as usize]; + if vec_bytes.len() < std::mem::size_of::() { + error!( + "Corrupted SchedStatWait data. Raw data: {}. Readed {} bytes expected {} bytes", + vec_bytes + .iter() + .map(|b| format!("{:02x}", b)) + .collect::>() + .join(" "), + vec_bytes.len(), + std::mem::size_of::() + ); + continue; + } + if vec_bytes.len() >= std::mem::size_of::() { + let sched_stat_wait: SchedStatWait = + unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; + + match exporter { + "otlp" => { + let mut metadata = Metadata::from_ebpf( + Some(sched_stat_wait.tgid), + &sched_stat_wait.command, + ); + metadata.enrich(); + metrics.record_sched_stat_wait(&sched_stat_wait, &metadata); + } + _ => continue, + } + + let tgid = sched_stat_wait.tgid; + let command = String::from_utf8_lossy(&sched_stat_wait.command); + let delay = sched_stat_wait.delay; + + info!( + "SchedStatWait - tgid: {}, command: {}, delay: {}", + tgid, command, delay + ); + } + } + } + + /// Read [`SchedStatRuntime`] events and record OpenTelemetry observations. + #[cfg(feature = "monitoring-structs")] + pub async fn read_sched_stat_runtime( + buffers: &mut [BytesMut], + tot_events: i32, + offset: i32, + exporter: &str, + metrics: Arc, + ) { + for i in offset..tot_events { + let vec_bytes = &buffers[i as usize]; + if vec_bytes.len() < std::mem::size_of::() { + error!( + "Corrupted SchedStatRuntime data. Raw data: {}. Readed {} bytes expected {} bytes", + vec_bytes + .iter() + .map(|b| format!("{:02x}", b)) + .collect::>() + .join(" "), + vec_bytes.len(), + std::mem::size_of::() + ); + continue; + } + if vec_bytes.len() >= std::mem::size_of::() { + let sched_stat_runtime: SchedStatRuntime = + unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; + + match exporter { + "otlp" => { + let mut metadata = Metadata::from_ebpf( + Some(sched_stat_runtime.tgid), + &sched_stat_runtime.command, + ); + metadata.enrich(); + metrics.record_sched_stat_runtime(&sched_stat_runtime, &metadata); + } + _ => continue, + } + + let tgid = sched_stat_runtime.tgid; + let command = String::from_utf8_lossy(&sched_stat_runtime.command); + let runtime = sched_stat_runtime.runtime; + + info!( + "SchedStatRuntime - tgid: {}, command: {}, runtime: {}", + tgid, command, runtime + ); + } + } + } + + /// Read [`CpuIdle`] events and record OpenTelemetry observations. + #[cfg(feature = "monitoring-structs")] + pub async fn read_cpu_idle( + buffers: &mut [BytesMut], + tot_events: i32, + offset: i32, + exporter: &str, + metrics: Arc, + ) { + for i in offset..tot_events { + let vec_bytes = &buffers[i as usize]; + if vec_bytes.len() < std::mem::size_of::() { + error!( + "Corrupted CpuIdle data. Raw data: {}. Readed {} bytes expected {} bytes", + vec_bytes + .iter() + .map(|b| format!("{:02x}", b)) + .collect::>() + .join(" "), + vec_bytes.len(), + std::mem::size_of::() + ); + continue; + } + if vec_bytes.len() >= std::mem::size_of::() { + let cpu_idle: CpuIdle = + unsafe { std::ptr::read_unaligned(vec_bytes.as_ptr() as *const _) }; + + match exporter { + "otlp" => { + let metadata = Metadata::from_ebpf(None, &[]); + metrics.record_cpu_idle(&cpu_idle, &metadata); + } + _ => continue, + } + + let cpu_id = cpu_idle.cpu_id; + let state = cpu_idle.state; + + info!( + "CpuIdle state changed - cpu_id: {}, state: {}", + cpu_id, state + ); + } + } + } +} + +/// Read perf-buffer events in a loop and dispatch to the appropriate [`Consumer`] handler. +/// +/// This function runs indefinitely (or until the process receives `SIGINT`). +/// It polls every CPU buffer every 100 ms, reads available events, and routes +/// them to the matching `Consumer::read_*` method. +/// +/// # Arguments +/// - `array_buffers` — per-CPU `PerfEventArrayBuffer` handles opened by [`fill_buffers`]. +/// - `buffers` — pre-allocated `BytesMut` scratch space sized by [`BufferSize::set_buffer`]. +/// - `consumer` — discriminator that selects which `read_*` method to invoke. +/// - `metrics` — optional [`Metrics`] handle; required when `consumer` is a monitoring variant. +#[cfg(feature = "buffer-reader")] +pub async fn read_perf_buffer>( + mut array_buffers: Vec>, + mut buffers: Vec, + consumer: Consumer, + #[cfg(feature = "monitoring-structs")] metrics: Option>, +) { + loop { + for buf in array_buffers.iter_mut() { + match buf.read_events(&mut buffers) { + Ok(events) => { + if events.lost > 0 { + tracing::debug!("Lost events: {} ", events.lost); + } + if events.read > 0 { + tracing::debug!("Readed events: {}", events.read); + let offset = 0; + let tot_events = events.read as i32; + + match consumer { + #[cfg(feature = "network-structs")] + Consumer::PacketLog => { + Consumer::read_packet_log(&mut buffers, tot_events, offset).await + } + #[cfg(feature = "network-structs")] + Consumer::TcpPacketRegistry => { + Consumer::read_tcp_registry_log(&mut buffers, tot_events, offset) + .await + } + #[cfg(feature = "network-structs")] + Consumer::VethLog => { + Consumer::read_and_handle_veth_log(&mut buffers, tot_events, offset) + .await + } + #[cfg(feature = "monitoring-structs")] + Consumer::PacketLossMetrics => { + Consumer::read_packet_loss_metrics( + &mut buffers, + tot_events, + offset, + "otlp", + metrics + .clone() + .expect("Metrics required for PacketLossMetrics"), + ) + .await + } + #[cfg(feature = "monitoring-structs")] + Consumer::TimeStampMetrics => { + Consumer::read_timestamp_metrics( + &mut buffers, + tot_events, + offset, + "otlp", + metrics + .clone() + .expect("Metric required for TimeStampMetrics"), + ) + .await + } + #[cfg(feature = "monitoring-structs")] + Consumer::CpuFrequency => { + Consumer::read_cpu_frequency( + &mut buffers, + tot_events, + offset, + "otlp", + metrics.clone().expect("Metric required for CpuFrequency"), + ) + .await + } + #[cfg(feature = "monitoring-structs")] + Consumer::MemAlloc => { + Consumer::read_mem_alloc( + &mut buffers, + tot_events, + offset, + "otlp", + metrics.clone().expect("Metric required for MemAlloc"), + ) + .await + } + #[cfg(feature = "monitoring-structs")] + Consumer::SchedStatWait => { + Consumer::read_sched_stat_wait( + &mut buffers, + tot_events, + offset, + "otlp", + metrics.clone().expect("Metric required for SchedStatWait"), + ) + .await + } + #[cfg(feature = "monitoring-structs")] + Consumer::SchedStatRuntime => { + Consumer::read_sched_stat_runtime( + &mut buffers, + tot_events, + offset, + "otlp", + metrics + .clone() + .expect("Metric required for SchedStatRuntime"), + ) + .await + } + #[cfg(feature = "monitoring-structs")] + Consumer::CpuIdle => { + Consumer::read_cpu_idle( + &mut buffers, + tot_events, + offset, + "otlp", + metrics.clone().expect("Metric required for CpuIdle"), + ) + .await + } + } + } + } + Err(e) => { + error!("Cannot read events from buffer. Reason: {} ", e); + } + } + } + tokio::time::sleep(std::time::Duration::from_millis(100)).await; + } +} diff --git a/core/common/src/lib.rs b/core/common/src/lib.rs index 15c4ad7..a0b2877 100644 --- a/core/common/src/lib.rs +++ b/core/common/src/lib.rs @@ -13,3 +13,4 @@ pub mod map_handlers; pub mod otel_metrics; #[cfg(feature = "program-handlers")] pub mod program_handlers; +mod semantic; diff --git a/core/common/src/otel_metrics.rs b/core/common/src/otel_metrics.rs index 79d08b8..e3c1fcc 100644 --- a/core/common/src/otel_metrics.rs +++ b/core/common/src/otel_metrics.rs @@ -12,9 +12,10 @@ //! telemetry by process. use crate::buffer_type::{ - CpuFrequency, CpuIdle, MemAlloc, NetworkMetrics, SchedStatRuntime, SchedStatWait, + CpuFrequency, CpuIdle, MemAlloc, PacketLossMetrics, SchedStatRuntime, SchedStatWait, TimeStampMetrics, }; +use crate::semantic::Semantic; use opentelemetry::KeyValue; use opentelemetry::metrics::{Counter, Gauge, Histogram, Meter}; pub struct Metrics { @@ -23,7 +24,7 @@ pub struct Metrics { /// Total number of network-related events produced by the `net_metrics` /// eBPF map. - pub packets_total: Counter, + pub socket_events_total: Counter, /// Observed socket drop count (`sk_drops`) from the kernel sock struct. pub sk_drops: Gauge, @@ -33,11 +34,7 @@ pub struct Metrics { /// Histogram of `delta_us` values supplied by the `time_stamp_events` /// perf buffer. - pub delta_us: Histogram, - - /// Histogram of `ts_us` values seen in both `net_metrics` and - /// `time_stamp_events`. - pub ts_us: Histogram, + pub tcp_latency_us: Histogram, /// Cpu bytes alloc total events pub cpu_bytes_alloc_events_total: Counter, @@ -54,154 +51,179 @@ pub struct Metrics { /// Observed scheduler wait time in nanoseconds (sched_stat_wait). pub sched_stat_wait: Gauge, + /// Distribution of scheduler wait times in nanoseconds (sched_stat_wait). + pub sched_stat_wait_distribution: Histogram, + /// Observed scheduler runtime in nanoseconds (sched_stat_runtime). pub sched_stat_runtime: Gauge, + /// Distribution of scheduler runtimes in nanoseconds (sched_stat_runtime). + pub sched_stat_runtime_distribution: Histogram, + /// Current CPU idle C-state per cpu_id, updated only on state change. pub cpu_idle_state: Gauge, } +// TODO: add identity metrics with TC classifier packet counts +// TODO: introduce a metric called total_tcp_packets total_udp_packets impl Metrics { /// Initialise all instruments backed by the supplied [`Meter`]. pub fn new(meter: &Meter) -> Self { // total events let events_total = meter - .u64_counter("events_total") - .with_description("Total number of eBPF events processed") + .u64_counter(Semantic::TotalEvents.title()) + .with_description(Semantic::TotalEvents.description()) + .with_unit("1") .build(); - // total packets - let packets_total = meter - .u64_counter("packets_total") - .with_description("Total number of network events processed") + // total socket events + let socket_events_total = meter + .u64_counter(Semantic::SocketTotalEvents.title()) + .with_description(Semantic::SocketTotalEvents.description()) + .with_unit("1") .build(); // socket drops let sk_drops = meter - .i64_gauge("sk_drops") - .with_description("Socket drop count per event") + .i64_gauge(Semantic::SocketDrops.title()) + .with_description(Semantic::SocketDrops.description()) + .with_unit("1") .build(); // socket errors let sk_err = meter - .i64_gauge("sk_err") - .with_description("Socket error count per event") + .i64_gauge(Semantic::SocketErrorsCount.title()) + .with_description(Semantic::SocketErrorsCount.description()) + .with_unit("1") .build(); - // delta microseconds - let delta_us = meter - .u64_histogram("delta_us") - .with_description("Distribution of delta_us values from timestamp events") - .build(); - - // timestamp microseconds grouped - let ts_us = meter - .u64_histogram("ts_us") - .with_description("Distribution of timestamp values from eBPF events") + // tcp latency microseconds + let tcp_latency_us = meter + .u64_histogram(Semantic::Latency.title()) + .with_description(Semantic::Latency.description()) + .with_unit("us") .build(); // cpu bytes alloc total events let cpu_bytes_alloc_events_total = meter - .u64_counter("bytes_alloc_events_total") - .with_description("Total bytes_alloc events occuring in the CPU") + .u64_counter(Semantic::PerCpuTotalEvents.title()) + .with_description(Semantic::PerCpuTotalEvents.description()) + .with_unit("1") .build(); // cpu bytes allocation let cpu_bytes_alloc = meter - .i64_gauge("cpu_bytes_alloc") - .with_description("Cpu bytes allocation per event") + .i64_gauge(Semantic::PerCpuBytesAllocated.title()) + .with_description(Semantic::PerCpuBytesAllocated.description()) + .with_unit("bytes") .build(); // memory allocation (mmap) events total let mem_alloc_events_total = meter - .u64_counter("mem_alloc_events_total") - .with_description("Total number of memory allocation (mmap) events processed") + .u64_counter(Semantic::TotalMemoryAllocationEvents.title()) + .with_description(Semantic::TotalMemoryAllocationEvents.description()) + .with_unit("1") .build(); // bytes requested via mmap syscalls let enter_mem_alloc = meter - .i64_gauge("enter_mem_alloc") - .with_description("Bytes requested via mmap syscalls") + .i64_gauge(Semantic::RequestedMemoryBytes.title()) + .with_description(Semantic::RequestedMemoryBytes.description()) + .with_unit("bytes") .build(); // scheduler wait time in nanoseconds let sched_stat_wait = meter - .i64_gauge("sched_stat_wait") - .with_description("Scheduler wait time in nanoseconds from sched_stat_wait") + .i64_gauge(Semantic::SchedulerWaitTime.title()) + .with_description(Semantic::SchedulerWaitTime.description()) + .with_unit("ns") + .build(); + + // distribution of scheduler wait times + let sched_stat_wait_distribution = meter + .u64_histogram(Semantic::SchedulerWaitTimeDistribution.title()) + .with_description(Semantic::SchedulerWaitTimeDistribution.description()) + .with_unit("ns") .build(); // scheduler runtime in nanoseconds let sched_stat_runtime = meter - .i64_gauge("sched_stat_runtime") - .with_description("Scheduler runtime in nanoseconds from sched_stat_runtime") + .i64_gauge(Semantic::SchedulerRuntime.title()) + .with_description(Semantic::SchedulerRuntime.description()) + .with_unit("ns") + .build(); + + // distribution of scheduler runtimes + let sched_stat_runtime_distribution = meter + .u64_histogram(Semantic::SchedulerRuntimeDistribution.title()) + .with_description(Semantic::SchedulerRuntimeDistribution.description()) + .with_unit("ns") .build(); // current CPU idle C-state per cpu_id let cpu_idle_state = meter - .i64_gauge("cpu_idle_state") - .with_description("Current CPU idle C-state per cpu_id, updated only on state change") + .i64_gauge(Semantic::CpuIdleState.title()) + .with_description(Semantic::CpuIdleState.description()) .build(); - Self { events_total, - packets_total, + socket_events_total, sk_drops, sk_err, - delta_us, - ts_us, + tcp_latency_us, cpu_bytes_alloc, cpu_bytes_alloc_events_total, mem_alloc_events_total, enter_mem_alloc, sched_stat_wait, + sched_stat_wait_distribution, sched_stat_runtime, + sched_stat_runtime_distribution, cpu_idle_state, } } - /// Record a single [`NetworkMetrics`] event. + /// Record a single [`PacketLossMetrics`] event. /// - /// Increments `events_total` and `packets_total`, records `sk_drops` and - /// `sk_err` as gauges, and observes `ts_us` in the timestamp histogram. + /// Increments `events_total` and `socket_events_total`, records `sk_drops` + /// and `sk_err` as gauges. /// /// Every observation carries: /// - /// -`tgid` – task group ID. - /// - `comm` – command name (null-terminated bytes converted to a UTF-8 + /// - `tgid` – task group ID. + /// - `command` – command name (null-terminated bytes converted to a UTF-8 /// string and trimmed). - pub fn record_network_metrics(&self, m: &NetworkMetrics) { + pub fn record_packet_loss_metrics(&self, m: &PacketLossMetrics) { let comm = String::from_utf8_lossy(&m.comm); let comm_trimmed = comm.trim_end_matches('\0').to_string(); let attrs = &[ KeyValue::new("tgid", m.tgid as i64), - KeyValue::new("comm", comm_trimmed), + KeyValue::new("command", comm_trimmed), ]; self.events_total.add(1, attrs); - self.packets_total.add(1, attrs); + self.socket_events_total.add(1, attrs); self.sk_drops.record(m.sk_drops as i64, attrs); self.sk_err.record(m.sk_err as i64, attrs); - self.ts_us.record(m.ts_us, attrs); } /// Record a single [`TimeStampMetrics`] event. /// - /// Increments `events_total`, and records `delta_us` and `ts_us` in their - /// respective histograms. + /// Increments `events_total`, and records `delta_us` in the latency + /// histogram. /// - /// Every observation carries `tgid` and `comm` (see - /// [`record_network_metrics`]). + /// Every observation carries `tgid` and `command` (see + /// [`record_packet_loss_metrics`]). pub fn record_timestamp_metrics(&self, m: &TimeStampMetrics) { let comm = String::from_utf8_lossy(&m.comm); let comm_trimmed = comm.trim_end_matches('\0').to_string(); let attrs = &[ KeyValue::new("tgid", m.tgid as i64), - KeyValue::new("comm", comm_trimmed), + KeyValue::new("command", comm_trimmed), ]; self.events_total.add(1, attrs); - self.delta_us.record(m.delta_us, attrs); - self.ts_us.record(m.ts_us, attrs); + self.tcp_latency_us.record(m.delta_us, attrs); } pub fn record_cpu_bytes_alloc(&self, m: &CpuFrequency) { @@ -231,14 +253,16 @@ impl Metrics { KeyValue::new("command", command), ]; + self.events_total.add(1, attrs); self.mem_alloc_events_total.add(1, attrs); self.enter_mem_alloc.record(m.length as i64, attrs); } /// Record a single [`SchedStatWait`] event. /// - /// Records `delay` in the `sched_stat_wait` gauge. No shared or dedicated - /// counter is incremented, as requested. + /// Increments `events_total`, records `delay` in the `sched_stat_wait` + /// gauge, and observes `delay` in the `sched_stat_wait_distribution` + /// histogram. pub fn record_sched_stat_wait(&self, m: &SchedStatWait) { let comm = String::from_utf8_lossy(&m.command); let command = comm.trim_end_matches('\0').to_string(); @@ -247,13 +271,16 @@ impl Metrics { KeyValue::new("command", command), ]; + self.events_total.add(1, attrs); self.sched_stat_wait.record(m.delay as i64, attrs); + self.sched_stat_wait_distribution.record(m.delay, attrs); } /// Record a single [`SchedStatRuntime`] event. /// - /// Records `runtime` in the `sched_stat_runtime` gauge. No shared or - /// dedicated counter is incremented, as requested. + /// Increments `events_total`, records `runtime` in the `sched_stat_runtime` + /// gauge, and observes `runtime` in the `sched_stat_runtime_distribution` + /// histogram. pub fn record_sched_stat_runtime(&self, m: &SchedStatRuntime) { let comm = String::from_utf8_lossy(&m.command); let command = comm.trim_end_matches('\0').to_string(); @@ -262,7 +289,9 @@ impl Metrics { KeyValue::new("command", command), ]; + self.events_total.add(1, attrs); self.sched_stat_runtime.record(m.runtime as i64, attrs); + self.sched_stat_runtime_distribution.record(m.runtime, attrs); } /// Record a single [`CpuIdle`] event. @@ -272,6 +301,7 @@ impl Metrics { pub fn record_cpu_idle(&self, m: &CpuIdle) { let attrs = &[KeyValue::new("cpu_id", m.cpu_id as i64)]; + self.events_total.add(1, attrs); self.cpu_idle_state.record(m.state as i64, attrs); } } diff --git a/core/common/src/semantic.rs b/core/common/src/semantic.rs new file mode 100644 index 0000000..a0381d8 --- /dev/null +++ b/core/common/src/semantic.rs @@ -0,0 +1,71 @@ +/// semantic conventions + +pub enum Semantic { + TotalEvents, + SocketTotalEvents, + SocketDrops, + SocketErrorsCount, + Latency, + PerCpuTotalEvents, + PerCpuBytesAllocated, + SchedulerRuntime, + SchedulerRuntimeDistribution, + SchedulerWaitTime, + SchedulerWaitTimeDistribution, + TotalMemoryAllocationEvents, + RequestedMemoryBytes, + CpuIdleState, +} + +impl Semantic { + pub fn title(&self) -> &'static str { + match self { + Semantic::TotalEvents => "events_total", + Semantic::SocketTotalEvents => "socket_events_total", + Semantic::SocketDrops => "sk_drops", + Semantic::SocketErrorsCount => "sk_err", + Semantic::Latency => "latency_us", + Semantic::PerCpuTotalEvents => "bytes_alloc_events_total", + Semantic::PerCpuBytesAllocated => "cpu_bytes_alloc", + Semantic::SchedulerRuntime => "sched_stat_runtime", + Semantic::SchedulerRuntimeDistribution => "sched_stat_runtime_distribution", + Semantic::SchedulerWaitTime => "sched_stat_wait", + Semantic::SchedulerWaitTimeDistribution => "sched_stat_wait_distribution", + Semantic::TotalMemoryAllocationEvents => "mem_alloc_events_total", + Semantic::RequestedMemoryBytes => "enter_mem_alloc", + Semantic::CpuIdleState => "cpu_idle_state", + } + } + pub fn description(&self) -> &'static str { + match self { + Semantic::TotalEvents => { + "Total number of eBPF events processed across all perf buffers" + } + Semantic::SocketTotalEvents => "Total number of socket state events processed", + Semantic::SocketDrops => "Socket drop count per event", + Semantic::SocketErrorsCount => "Socket error count per event", + Semantic::Latency => "Distribution of latency values from timestamp events", + Semantic::PerCpuTotalEvents => "Total bytes_alloc events occurring in the CPU", + Semantic::PerCpuBytesAllocated => "CPU bytes allocation per event", + Semantic::SchedulerRuntime => { + "Scheduler runtime in nanoseconds from sched_stat_runtime" + } + Semantic::SchedulerRuntimeDistribution => { + "Distribution of scheduler runtimes in nanoseconds from sched_stat_runtime" + } + Semantic::SchedulerWaitTime => { + "Scheduler wait time in nanoseconds from sched_stat_wait" + } + Semantic::SchedulerWaitTimeDistribution => { + "Distribution of scheduler wait times in nanoseconds from sched_stat_wait" + } + Semantic::TotalMemoryAllocationEvents => { + "Total number of memory allocation (mmap) events processed" + } + Semantic::RequestedMemoryBytes => "Bytes requested via mmap syscalls", + Semantic::CpuIdleState => { + "Current CPU idle C-state per cpu_id, updated only on state change" + } + } + } +} diff --git a/core/common/src/service_discovery.rs b/core/common/src/service_discovery.rs new file mode 100644 index 0000000..af2efd1 --- /dev/null +++ b/core/common/src/service_discovery.rs @@ -0,0 +1,166 @@ +#[cfg(feature = "experimental")] +use anyhow::Error; +#[cfg(feature = "experimental")] +use std::fs; + +/// Supported runtime prefixes for extracting the container ID from a cgroup path. +const RUNTIME_PREFIXES: &[&str] = &["docker-", "cri-containerd-", "crio-"]; + +/// Extract the container ID (e.g. Docker runtime ID) from a cgroup filesystem path. +/// +/// Supports multiple prefixes: docker-, cri-containerd-, crio-. +#[cfg(feature = "experimental")] +pub fn extract_container_id(cgroup_path: &str) -> Result { + let splits: Vec<&str> = cgroup_path.split('/').collect(); + + for prefix in RUNTIME_PREFIXES { + if let Ok(index) = extract_target_from_splits(&splits, prefix) { + let id = splits[index] + .trim_start_matches(prefix) + .trim_end_matches(".scope"); + return Ok(id.to_string()); + } + } + + Err(Error::msg(format!( + "No known runtime prefix found in cgroup path: {}", + cgroup_path + ))) +} + +/// Extract the Pod UID from a Kubernetes cgroup filesystem path. +/// +/// Example inputs: +/// - `/sys/fs/cgroup/kubepods.slice/kubepods-besteffort.slice/kubepods-besteffort-pod231bd2d7_0f09_4781_a4e1_e4ea026342dd.slice` +/// - `/sys/fs/cgroup/kubelet.slice/kubelet-kubepods.slice/kubelet-kubepods-besteffort.slice/kubelet-kubepods-besteffort-pod231bd2d7_0f09_4781_a4e1_e4ea026342dd.slice` +#[cfg(feature = "experimental")] +pub fn extract_pod_uid(cgroup_path: &str) -> Result { + let splits: Vec<&str> = cgroup_path.split('/').collect(); + + let index = extract_target_from_splits(&splits, "-pod")?; + + let pod_split = splits[index] + .trim_start_matches("kubelet-kubepods-besteffort-") + .trim_start_matches("kubelet-kubepods-burstable-") + .trim_start_matches("kubepods-besteffort-") + .trim_start_matches("kubepods-burstable-"); + + let uid_ = pod_split + .trim_start_matches("pod") + .trim_end_matches(".slice"); + + let uid = uid_.replace('_', "-"); + Ok(uid) +} + +/// Scan a given cgroup directory and return subdirectory paths. +/// +/// If `path` does not exist or is empty, falls back to the default K8s kubepods.slice path. +#[cfg(feature = "experimental")] +pub fn scan_cgroup_paths(path: &str) -> Result, Error> { + let mut cgroup_paths: Vec = Vec::new(); + let default_path = "/sys/fs/cgroup/kubepods.slice"; + + let target_path = if path.is_empty() || fs::metadata(path).is_err() { + default_path + } else { + path + }; + + let entries = match fs::read_dir(target_path) { + Ok(entries) => entries, + Err(e) => { + tracing::error!("Error reading cgroup directory {:?}: {}", target_path, e); + return Ok(cgroup_paths); + } + }; + + for entry in entries { + if let Ok(entry) = entry { + let path = entry.path(); + if path.is_dir() { + if let Some(path_str) = path.to_str() { + cgroup_paths.push(path_str.to_string()); + } + } + } + } + + Ok(cgroup_paths) +} + +#[cfg(feature = "experimental")] +fn extract_target_from_splits(splits: &[&str], target: &str) -> Result { + for (index, split) in splits.iter().enumerate() { + if split.contains(target) { + return Ok(index); + } + } + Err(Error::msg(format!("'{target}' word not found in split"))) +} + +#[cfg(feature = "experimental")] +mod tests { + use super::*; + + #[test] + fn extract_uid_from_string() { + let cgroup_paths = vec![ + "/sys/fs/cgroup/kubepods.slice/kubepods-besteffort.slice/kubepods-besteffort-pod231bd2d7_0f09_4781_a4e1_e4ea026342dd.slice".to_string(), + "/sys/fs/cgroup/kubelet.slice/kubelet-kubepods.slice/kubelet-kubepods-besteffort.slice/kubelet-kubepods-besteffort-pod231bd2d7_0f09_4781_a4e1_e4ea026342dd.slice".to_string(), + ]; + + let mut uid_vec = Vec::::new(); + + for cgroup_path in cgroup_paths { + let uid = extract_pod_uid(&cgroup_path) + .map_err(|e| format!("An error occurred {}", e)) + .unwrap(); + uid_vec.push(uid); + } + + let check = vec![ + "231bd2d7-0f09-4781-a4e1-e4ea026342dd".to_string(), + "231bd2d7-0f09-4781-a4e1-e4ea026342dd".to_string(), + ]; + + assert_eq!(uid_vec, check); + } + + #[test] + fn test_extract_target_index() { + let cgroup_paths = vec![ + "/sys/fs/cgroup/kubepods.slice/kubepods-besteffort.slice/kubepods-besteffort-pod231bd2d7_0f09_4781_a4e1_e4ea026342dd.slice".to_string(), + "/sys/fs/cgroup/kubelet.slice/kubelet-kubepods.slice/kubelet-kubepods-besteffort.slice/kubelet-kubepods-besteffort-pod231bd2d7_0f09_4781_a4e1_e4ea026342dd.slice".to_string(), + ]; + + let mut index_vec = Vec::::new(); + for cgroup_path in cgroup_paths { + let splits: Vec<&str> = cgroup_path.split('/').collect(); + + let target_index = extract_target_from_splits(&splits, "-pod").unwrap(); + index_vec.push(target_index); + } + let index_check = vec![6, 7]; + assert_eq!(index_vec, index_check); + } + + #[test] + fn extract_docker_id() { + let cgroup_paths = vec![ + "/sys/fs/cgroup/kubepods.slice/kubepods-besteffort.slice/kubepods-besteffort-pod17fd3f7c_37e4_4009_8c38_e58b30691af3.slice/docker-13abd64c0ba349975a762476c9703b642d18077eabeb3aa1d941132048afc861.scope".to_string(), + "/sys/fs/cgroup/kubelet.slice/kubelet-kubepods.slice/kubelet-kubepods-besteffort.slice/kubelet-kubepods-besteffort-pod17fd3f7c_37e4_4009_8c38_e58b30691af3.slice/docker-13abd64c0ba349975a762476c9703b642d18077eabeb3aa1d941132048afc861.scope".to_string(), + ]; + + let mut id_vec = Vec::::new(); + for cgroup_path in cgroup_paths { + let id = extract_container_id(&cgroup_path).unwrap(); + id_vec.push(id); + } + let id_check = vec![ + "13abd64c0ba349975a762476c9703b642d18077eabeb3aa1d941132048afc861".to_string(), + "13abd64c0ba349975a762476c9703b642d18077eabeb3aa1d941132048afc861".to_string(), + ]; + assert_eq!(id_vec, id_check); + } +} diff --git a/core/src/components/metrics/build-local-metrics.sh b/core/src/components/metrics/build-local-metrics.sh index b861e6e..d0ed752 100755 --- a/core/src/components/metrics/build-local-metrics.sh +++ b/core/src/components/metrics/build-local-metrics.sh @@ -3,7 +3,7 @@ # Building identity files echo "Building the metrics-tracer files" pushd ../metrics_tracer -./build-metrics-tracer.sh +./build-local-metrics-tracer.sh popd echo "Copying metrics_tracer binaries" diff --git a/core/src/components/metrics/src/helpers.rs b/core/src/components/metrics/src/helpers.rs index 141cad6..6c526f0 100644 --- a/core/src/components/metrics/src/helpers.rs +++ b/core/src/components/metrics/src/helpers.rs @@ -7,7 +7,7 @@ use std::sync::Arc; use tokio::signal; use tracing::{error, info}; -use cortexbrain_common::buffer_type::{BufferType, read_perf_buffer}; +use cortexbrain_common::consumer::{Consumer, read_perf_buffer}; use cortexbrain_common::otel_metrics::Metrics; /// Listen for eBPF perf-buffer events and record OpenTelemetry metrics. @@ -103,7 +103,7 @@ pub async fn event_listener(bpf_maps: BpfMapsData, meter: Meter) -> Result<(), a read_perf_buffer( array_buffers, buffers, - BufferType::NetworkMetrics, + Consumer::PacketLossMetrics, Some(metrics), ) .await; @@ -118,7 +118,7 @@ pub async fn event_listener(bpf_maps: BpfMapsData, meter: Meter) -> Result<(), a read_perf_buffer( array_buffers, buffers, - BufferType::TimeStampMetrics, + Consumer::TimeStampMetrics, Some(metrics), ) .await; @@ -133,7 +133,7 @@ pub async fn event_listener(bpf_maps: BpfMapsData, meter: Meter) -> Result<(), a read_perf_buffer( array_buffers, buffers, - BufferType::CpuFrequency, + Consumer::CpuFrequency, Some(metrics), ) .await; @@ -145,7 +145,7 @@ pub async fn event_listener(bpf_maps: BpfMapsData, meter: Meter) -> Result<(), a let mut array_buffers = cpu_idle_perf_buffer; let mut buffers = cpu_idle_buffers; tokio::spawn(async move { - read_perf_buffer(array_buffers, buffers, BufferType::CpuIdle, Some(metrics)).await; + read_perf_buffer(array_buffers, buffers, Consumer::CpuIdle, Some(metrics)).await; }) }; @@ -154,7 +154,7 @@ pub async fn event_listener(bpf_maps: BpfMapsData, meter: Meter) -> Result<(), a let mut array_buffers = mem_alloc_perf_buffer; let mut buffers = mem_alloc_buffers; tokio::spawn(async move { - read_perf_buffer(array_buffers, buffers, BufferType::MemAlloc, Some(metrics)).await; + read_perf_buffer(array_buffers, buffers, Consumer::MemAlloc, Some(metrics)).await; }) }; @@ -166,7 +166,7 @@ pub async fn event_listener(bpf_maps: BpfMapsData, meter: Meter) -> Result<(), a read_perf_buffer( array_buffers, buffers, - BufferType::SchedStatWait, + Consumer::SchedStatWait, Some(metrics), ) .await; @@ -181,7 +181,7 @@ pub async fn event_listener(bpf_maps: BpfMapsData, meter: Meter) -> Result<(), a read_perf_buffer( array_buffers, buffers, - BufferType::SchedStatRuntime, + Consumer::SchedStatRuntime, Some(metrics), ) .await; diff --git a/core/src/components/metrics/src/main.rs b/core/src/components/metrics/src/main.rs index bcf7de7..3a7bcbb 100644 --- a/core/src/components/metrics/src/main.rs +++ b/core/src/components/metrics/src/main.rs @@ -79,10 +79,14 @@ async fn main() -> Result<(), anyhow::Error> { info!("BPF maps pinned successfully to {}", bpf_map_save_path); { - load_program(bpf.clone(), "metrics_tracer", "tcp_identify_packet_loss") - .context( - "An error occurred during the execution of load_program function", - )?; + load_program( + bpf.clone(), + "packet_loss_tracer", + "tcp_identify_packet_loss", + ) + .context( + "An error occurred during the execution of load_program function", + )?; load_program(tcp_bpf, "tcp_v4_connect", "tcp_v4_connect") .context("An error occurred during the execution of load_and_attach_tcp_programs function")?; @@ -90,12 +94,8 @@ async fn main() -> Result<(), anyhow::Error> { load_program(tcp_v6_bpf, "tcp_v6_connect", "tcp_v6_connect") .context("An error occurred during the execution of load_and_attach_tcp_programs function")?; - load_program( - tcp_rev_bpf, - "tcp_rcv_state_process", - "tcp_rcv_state_process", - ) - .context( + load_program(tcp_rev_bpf, "tcp_latency_monitor", "tcp_rcv_state_process") + .context( "An error occurred during the execution of load_program function", )?; load_tracepoint_program( diff --git a/core/src/components/metrics_tracer/src/cpu.rs b/core/src/components/metrics_tracer/src/cpu.rs index 374ea36..dbab7d4 100644 --- a/core/src/components/metrics_tracer/src/cpu.rs +++ b/core/src/components/metrics_tracer/src/cpu.rs @@ -2,7 +2,7 @@ //tracepoint:power:cpu_frequency_limits //tracepoint:power:cpu_idle //tracepoint:power:cpu_idle_miss -use aya_ebpf::{EbpfContext, programs::TracePointContext}; +use aya_ebpf::{EbpfContext, helpers::bpf_get_current_pid_tgid, programs::TracePointContext}; use aya_log_ebpf::info; use crate::data_structures::{CPU_FREQUENCY, CPU_IDLE, CPU_IDLE_LAST_STATE, CpuFrequency, CpuIdle}; @@ -37,7 +37,9 @@ pub fn per_cpu_bytes_alloc(ctx: &TracePointContext) -> Result<((u32, u32, [u8; 1 let bytes_alloc_offset = 64; let pid_offset = 4; let bytes_alloc = unsafe { ctx.read_at(bytes_alloc_offset) }?; - let pid = unsafe { ctx.read_at(pid_offset) }?; + //let tgid: u32 = unsafe { ctx.read_at(tgid_offset) }?; + let pid_tgid: u64 = bpf_get_current_pid_tgid(); + let tgid: u32 = (pid_tgid >> 32) as u32; let command = ctx.command()?; //let cpu_freq_data = CpuFrequency { @@ -47,29 +49,33 @@ pub fn per_cpu_bytes_alloc(ctx: &TracePointContext) -> Result<((u32, u32, [u8; 1 //CPU_FREQUENCY.output(&ctx, &cpu_freq_data, 0); - Ok((bytes_alloc, pid, command)) + Ok((bytes_alloc, tgid, command)) } pub fn sched_stat_wait(ctx: &TracePointContext) -> Result<((u32, u64, [u8; 16])), i64> { let pid_offset = 4; let delay_offset = 16; - let pid = unsafe { ctx.read_at(pid_offset) }?; + //let tgid: u32 = unsafe { ctx.read_at(tgid_offset) }?; + let pid_tgid: u64 = bpf_get_current_pid_tgid(); + let tgid: u32 = (pid_tgid >> 32) as u32; let delay = unsafe { ctx.read_at(delay_offset) }?; let command = ctx.command()?; - Ok((pid, delay, command)) + Ok((tgid, delay, command)) } pub fn sched_stat_runtime(ctx: &TracePointContext) -> Result<((u32, u64, [u8; 16])), i64> { let pid_offset = 4; let runtime_offset = 16; - let pid = unsafe { ctx.read_at(pid_offset) }?; + //let tgid: u32 = unsafe { ctx.read_at(tgid_offset) }?; + let pid_tgid: u64 = bpf_get_current_pid_tgid(); + let tgid: u32 = (pid_tgid >> 32) as u32; let runtime = unsafe { ctx.read_at(runtime_offset) }?; let command = ctx.command()?; - Ok((pid, runtime, command)) + Ok((tgid, runtime, command)) } diff --git a/core/src/components/metrics_tracer/src/data_structures.rs b/core/src/components/metrics_tracer/src/data_structures.rs index d74ef35..cfc80b6 100644 --- a/core/src/components/metrics_tracer/src/data_structures.rs +++ b/core/src/components/metrics_tracer/src/data_structures.rs @@ -6,7 +6,7 @@ use aya_ebpf::{ pub const TASK_COMM_LEN: usize = 16; #[repr(C, packed)] -pub struct NetworkMetrics { +pub struct PacketLossMetrics { pub tgid: u32, pub comm: [u8; TASK_COMM_LEN], pub ts_us: u64, @@ -27,7 +27,8 @@ pub struct TimeStampStartInfo { pub tgid: u32, } -// Event we send to userspace when latency is computed +/// Event we send to userspace when latency is computed +/// used to compute tcp_delta_us, tcp_ts_us metrics #[repr(C, packed)] #[derive(Copy, Clone)] pub struct TimeStampEvent { @@ -50,7 +51,7 @@ pub struct CpuFrequency { //pub(crate) cpu_id: u32, //pub(crate) cpu_freq: u32, pub(crate) bytes_alloc: u32, - pub(crate) pid: u32, + pub(crate) tgid: u32, pub(crate) command: [u8; 16], } @@ -86,6 +87,25 @@ pub struct CpuIdle { pub(crate) state: u32, } +#[repr(C, packed)] +#[derive(Copy, Clone, Debug, bytemuck::Pod, bytemuck::Zeroable)] +pub struct SslEvent { + pub tgid: u32, + pub comm: [u8; TASK_COMM_LEN], + pub ts_us: u64, + pub direction: u8, // 0 = read, 1 = write + pub size: i32, // return value (bytes transferred or <0 on error) + pub requested: i32, // num argument passed to SSL_read/SSL_write +} + +#[repr(C, packed)] +#[derive(Copy, Clone, Debug, bytemuck::Pod, bytemuck::Zeroable)] +pub struct SslCtxInfo { + pub ssl: u64, // SSL* pointer + pub requested: i32, + pub ts_ns: u64, +} + // Map: connect-start timestamp by socket pointer #[map(name = "time_stamp_start")] pub static mut TIME_STAMP_START: HashMap<*mut core::ffi::c_void, TimeStampStartInfo> = @@ -97,7 +117,7 @@ pub static mut TIME_STAMP_EVENTS: PerfEventArray = PerfEventArray::::new(0); #[map(name = "net_metrics")] -pub static NET_METRICS: PerfEventArray = PerfEventArray::new(0); +pub static NET_METRICS: PerfEventArray = PerfEventArray::new(0); #[map(name = "cpu_frequency")] pub static CPU_FREQUENCY: PerfEventArray = PerfEventArray::new(0); @@ -117,3 +137,10 @@ pub static CPU_IDLE: PerfEventArray = PerfEventArray::new(0); #[map(name = "cpu_idle_last_state")] pub static mut CPU_IDLE_LAST_STATE: HashMap = HashMap::::with_max_entries(256, 0); + +#[map(name = "ssl_ctx_map")] +pub static mut SSL_CTX_MAP: HashMap = + HashMap::::with_max_entries(4096, 0); + +#[map(name = "ssl_events")] +pub static SSL_EVENTS: PerfEventArray = PerfEventArray::new(0); diff --git a/core/src/components/metrics_tracer/src/main.rs b/core/src/components/metrics_tracer/src/main.rs index 55faa45..1da3451 100644 --- a/core/src/components/metrics_tracer/src/main.rs +++ b/core/src/components/metrics_tracer/src/main.rs @@ -6,6 +6,7 @@ mod bindings; mod cpu; mod data_structures; mod memory; +mod network; use crate::bindings::net_device; use crate::cpu::{cpu_idle, per_cpu_bytes_alloc, sched_stat_runtime, sched_stat_wait}; @@ -13,12 +14,13 @@ use crate::data_structures::CpuFrequency; use crate::data_structures::NET_METRICS; use crate::data_structures::{CPU_FREQUENCY, SchedStatWait}; use crate::data_structures::{ - CPU_IDLE, NetworkMetrics, TASK_COMM_LEN, TIME_STAMP_EVENTS, TIME_STAMP_START, TimeStampEvent, - TimeStampStartInfo, + CPU_IDLE, PacketLossMetrics, TASK_COMM_LEN, TIME_STAMP_EVENTS, TIME_STAMP_START, + TimeStampEvent, TimeStampStartInfo, }; use crate::data_structures::{MEM_ALLOC, SCHED_STAT_RUNTIME, SCHED_STAT_WAIT}; use crate::data_structures::{MemAlloc, SchedStatRuntime}; use crate::memory::enter_mmap; +use crate::network::{detect_packet_loss, on_connect, on_rcv_state_process}; use aya_ebpf::EbpfContext; use aya_ebpf::helpers::bpf_get_current_pid_tgid; use aya_ebpf::helpers::generated::{bpf_ktime_get_ns, bpf_perf_event_output}; @@ -35,7 +37,7 @@ const AF_INET6: u16 = 10; const TCP_SYN_SENT: u8 = 2; #[kprobe] -fn metrics_tracer(ctx: ProbeContext) -> u32 { +fn packet_loss_tracer(ctx: ProbeContext) -> u32 { match try_metrics_tracer(ctx) { Ok(ret) => ret, Err(ret) => ret.try_into().unwrap_or(1), @@ -43,64 +45,7 @@ fn metrics_tracer(ctx: ProbeContext) -> u32 { } fn try_metrics_tracer(ctx: ProbeContext) -> Result { - let sk_pointer = ctx.arg::<*const u8>(0).ok_or(1i64)?; - - if sk_pointer.is_null() { - return Err(1); - } - - let tgid = (unsafe { bpf_get_current_pid_tgid() } >> 32) as u32; - let comm = unsafe { bpf_get_current_comm() }.map_err(|_| 1i64)?; - let ts_us: u64 = unsafe { bpf_ktime_get_ns() } / 1_000; - let sk_err_offset = 284; - let sk_err_soft_offset = 600; - let sk_backlog_len_offset = 196; - let sk_write_memory_queued_offset = 376; - let sk_receive_buffer_size_offset = 244; - let sk_ack_backlog_offset = 604; - let sk_drops_offset = 136; - - let sk_err = unsafe { - bpf_probe_read_kernel::(sk_pointer.add(sk_err_offset) as *const i32).map_err(|_| 1)? - }; - let sk_err_soft = unsafe { - bpf_probe_read_kernel::(sk_pointer.add(sk_err_soft_offset) as *const i32) - .map_err(|_| 1)? - }; - let sk_backlog_len = unsafe { - bpf_probe_read_kernel::(sk_pointer.add(sk_backlog_len_offset) as *const i32) - .map_err(|_| 1)? - }; - let sk_write_memory_queued = unsafe { - bpf_probe_read_kernel::(sk_pointer.add(sk_write_memory_queued_offset) as *const i32) - .map_err(|_| 1)? - }; - let sk_receive_buffer_size = unsafe { - bpf_probe_read_kernel::(sk_pointer.add(sk_receive_buffer_size_offset) as *const i32) - .map_err(|_| 1)? - }; - let sk_ack_backlog = unsafe { - bpf_probe_read_kernel::(sk_pointer.add(sk_ack_backlog_offset) as *const u32) - .map_err(|_| 1)? - }; - let sk_drops = unsafe { - bpf_probe_read_kernel::(sk_pointer.add(sk_drops_offset) as *const i32) - .map_err(|_| 1)? - }; - - let net_metrics = NetworkMetrics { - tgid: tgid, - comm: comm, - ts_us: ts_us, - sk_err: sk_err, - sk_err_soft: sk_err_soft, - sk_backlog_len: sk_backlog_len, - sk_write_memory_queued: sk_write_memory_queued, - sk_receive_buffer_size: sk_receive_buffer_size, - sk_ack_backlog: sk_ack_backlog, - sk_drops: sk_drops, - }; - + let net_metrics = detect_packet_loss(&ctx)?; unsafe { NET_METRICS.output(&ctx, &net_metrics, 0); } @@ -108,7 +53,7 @@ fn try_metrics_tracer(ctx: ProbeContext) -> Result { Ok(0) } -// Monitor on tcp_sendmsg, tcp_v4_connect +/// Monitor on tcp_sendmsg, tcp_v6_connect #[kprobe] fn tcp_v6_connect(ctx: ProbeContext) -> u32 { match on_connect(ctx) { @@ -117,7 +62,7 @@ fn tcp_v6_connect(ctx: ProbeContext) -> u32 { } } -// Monitor on tcp_sendmsg, tcp_v4_connect +/// Monitor on tcp_sendmsg, tcp_v4_connect #[kprobe] fn tcp_v4_connect(ctx: ProbeContext) -> u32 { match on_connect(ctx) { @@ -126,149 +71,14 @@ fn tcp_v4_connect(ctx: ProbeContext) -> u32 { } } -fn on_connect(ctx: ProbeContext) -> Result<(), i64> { - let sk = ctx.arg::<*mut bindings::sock>(0).ok_or(1i64)?; - if sk.is_null() { - return Err(1); - } - - let tgid = (unsafe { bpf_get_current_pid_tgid() } >> 32) as u32; - let mut start = TimeStampStartInfo { - comm: [0; TASK_COMM_LEN], - ts_ns: unsafe { bpf_ktime_get_ns() }, - tgid, - }; - unsafe { - let comm_result = bpf_get_current_comm(); - if let Ok(comm) = comm_result { - start.comm.copy_from_slice(&comm); - } - let map_ptr = &raw mut TIME_STAMP_START; - (*map_ptr) - .insert(&(sk as *mut core::ffi::c_void), &start, 0) - .map_err(|_| 1)?; - } - Ok(()) -} - #[kprobe] -fn tcp_rcv_state_process(ctx: ProbeContext) -> u32 { +fn tcp_latency_monitor(ctx: ProbeContext) -> u32 { match on_rcv_state_process(ctx) { Ok(_) => 0, Err(e) => e as u32, } } -fn on_rcv_state_process(ctx: ProbeContext) -> Result<(), i64> { - // On some kernels, kprobe wrapper puts `sk` at arg0; on others arg1. - let sk = ctx.arg::<*mut bindings::sock>(0).unwrap_or(ptr::null_mut()); - let sk = if sk.is_null() { - ctx.arg::<*mut bindings::sock>(1).ok_or(1i64)? - } else { - sk - }; - - if sk.is_null() { - return Err(1); - } - - let skc_daddr_off = 0; - let skc_rcv_saddr_off = 4; - let skc_dport_off = 12; - let skc_num_off = 14; - let skc_family_off = 16; - let skc_state_off = 18; - let skc_v6_daddr_off = 56; - let skc_v6_rcv_saddr_off = 72; - - let state = unsafe { bpf_probe_read_kernel::((sk as usize + skc_state_off) as *const u8) } - .map_err(|_| 1)?; - - if state != TCP_SYN_SENT { - return Ok(()); - } - - let start = unsafe { - let map_ptr = &raw const TIME_STAMP_START; - (*map_ptr).get(&((sk as usize) as *mut core::ffi::c_void)) - } - .ok_or(1i64)?; - let now = unsafe { bpf_ktime_get_ns() }; - let delta = now as i64 - start.ts_ns as i64; - if delta <= 0 { - unsafe { - let map_ptr = &raw mut TIME_STAMP_START; - let _ = (*map_ptr).remove(&((sk as usize) as *mut core::ffi::c_void)); - } - return Ok(()); - } - - let mut ev = TimeStampEvent { - delta_us: (delta as u64) / 1_000, - ts_us: now / 1_000, - tgid: start.tgid, - comm: start.comm, - lport: 0, - dport_be: 0, - af: 0, - saddr_v4: 0, - daddr_v4: 0, - saddr_v6: [0; 4], - daddr_v6: [0; 4], - }; - - // family, ports - ev.af = unsafe { - bpf_probe_read_kernel::((sk as usize + skc_family_off) as *const u16).map_err(|_| 1)? - }; - ev.lport = unsafe { - bpf_probe_read_kernel::((sk as usize + skc_num_off) as *const u16).map_err(|_| 1)? - }; - ev.dport_be = unsafe { - bpf_probe_read_kernel::((sk as usize + skc_dport_off) as *const u16).map_err(|_| 1)? - }; - - if ev.af == AF_INET { - ev.saddr_v4 = unsafe { - bpf_probe_read_kernel::((sk as usize + skc_rcv_saddr_off) as *const u32) - .map_err(|_| 1)? - }; - ev.daddr_v4 = unsafe { - bpf_probe_read_kernel::((sk as usize + skc_daddr_off) as *const u32) - .map_err(|_| 1)? - }; - } else { - // read 16 bytes as four u32 words - for i in 0..4 { - ev.saddr_v6[i] = unsafe { - bpf_probe_read_kernel::( - (sk as usize + skc_v6_rcv_saddr_off + i * 4) as *const u32, - ) - .map_err(|_| 1)? - }; - ev.daddr_v6[i] = unsafe { - bpf_probe_read_kernel::((sk as usize + skc_v6_daddr_off + i * 4) as *const u32) - .map_err(|_| 1)? - }; - } - } - - // emit + cleanup - unsafe { - bpf_perf_event_output( - ctx.as_ptr(), - &raw const TIME_STAMP_EVENTS as *const _ as *mut _, - 0, // BPF_F_CURRENT_CPU - &ev as *const _ as *mut _, - (mem::size_of::() as u32).into(), - ); - let map_ptr = &raw mut TIME_STAMP_START; - let _ = (*map_ptr).remove(&((sk as usize) as *mut core::ffi::c_void)); - } - - Ok(()) -} - #[tracepoint] fn trace_cpu_frequency(ctx: TracePointContext) -> u32 { match trace_cpu_metrics(&ctx) { @@ -286,13 +96,13 @@ fn trace_cpu_idle(ctx: TracePointContext) -> u32 { } fn trace_cpu_metrics(ctx: &TracePointContext) -> Result<(), i64> { - let (bytes_alloc, pid, command) = per_cpu_bytes_alloc(ctx)?; + let (bytes_alloc, tgid, command) = per_cpu_bytes_alloc(ctx)?; //let (cpu_id, cpu_freq) = cpu_frequency(&ctx)?; let cpu_metrics = CpuFrequency { // cpu_id, // cpu_freq, bytes_alloc, - pid, + tgid, command, }; diff --git a/core/src/components/metrics_tracer/src/memory.rs b/core/src/components/metrics_tracer/src/memory.rs index de14c7a..a7ee4d4 100644 --- a/core/src/components/metrics_tracer/src/memory.rs +++ b/core/src/components/metrics_tracer/src/memory.rs @@ -1,4 +1,4 @@ -use aya_ebpf::{EbpfContext, programs::TracePointContext}; +use aya_ebpf::{EbpfContext, helpers::bpf_get_current_pid_tgid, programs::TracePointContext}; /// Read the fields of the `syscalls:sys_enter_mmap` tracepoint. pub fn enter_mmap(ctx: &TracePointContext) -> Result<((u32, u64, u64, [u8; 16])), i64> { @@ -7,7 +7,9 @@ pub fn enter_mmap(ctx: &TracePointContext) -> Result<((u32, u64, u64, [u8; 16])) let addr_offset = 16; let len_offset = 24; - let tgid: u32 = unsafe { ctx.read_at(tgid_offset) }?; + //let tgid: u32 = unsafe { ctx.read_at(tgid_offset) }?; + let pid_tgid: u64 = bpf_get_current_pid_tgid(); + let tgid: u32 = (pid_tgid >> 32) as u32; let addr: u64 = unsafe { ctx.read_at(addr_offset) }?; let len: u64 = unsafe { ctx.read_at(len_offset) }?; let command = ctx.command()?; diff --git a/core/src/components/metrics_tracer/src/mod.rs b/core/src/components/metrics_tracer/src/mod.rs index 56f80cc..e62fb54 100644 --- a/core/src/components/metrics_tracer/src/mod.rs +++ b/core/src/components/metrics_tracer/src/mod.rs @@ -2,3 +2,4 @@ mod bindings; mod cpu; mod data_structures; mod memory; +mod network; diff --git a/core/src/components/metrics_tracer/src/network.rs b/core/src/components/metrics_tracer/src/network.rs new file mode 100644 index 0000000..4cc50bb --- /dev/null +++ b/core/src/components/metrics_tracer/src/network.rs @@ -0,0 +1,217 @@ +use crate::bindings::{self, net_device}; +use crate::data_structures::{NET_METRICS, PacketLossMetrics}; +use crate::data_structures::{ + TASK_COMM_LEN, TIME_STAMP_EVENTS, TIME_STAMP_START, TimeStampEvent, TimeStampStartInfo, +}; +use aya_ebpf::EbpfContext; +use aya_ebpf::helpers::bpf_get_current_pid_tgid; +use aya_ebpf::helpers::generated::{bpf_ktime_get_ns, bpf_perf_event_output}; +use aya_ebpf::helpers::{ + bpf_get_current_comm, bpf_probe_read_kernel, bpf_probe_read_kernel_str_bytes, +}; +use aya_ebpf::macros::{kprobe, map, tracepoint}; +use aya_ebpf::maps::{HashMap, PerfEventArray}; +use aya_ebpf::programs::ProbeContext; +use core::{mem, ptr}; + +const AF_INET: u16 = 2; +const AF_INET6: u16 = 10; +const TCP_SYN_SENT: u8 = 2; + +/// packet loss tracer +pub fn detect_packet_loss(ctx: &ProbeContext) -> Result { + let sk_pointer = ctx.arg::<*const u8>(0).ok_or(1i64)?; + + if sk_pointer.is_null() { + return Err(1); + } + + let tgid = (unsafe { bpf_get_current_pid_tgid() } >> 32) as u32; + let comm = unsafe { bpf_get_current_comm() }.map_err(|_| 1i64)?; + let ts_us: u64 = unsafe { bpf_ktime_get_ns() } / 1_000; + let sk_err_offset = 284; + let sk_err_soft_offset = 600; + let sk_backlog_len_offset = 196; + let sk_write_memory_queued_offset = 376; + let sk_receive_buffer_size_offset = 244; + let sk_ack_backlog_offset = 604; + let sk_drops_offset = 136; + + let sk_err = unsafe { + bpf_probe_read_kernel::(sk_pointer.add(sk_err_offset) as *const i32).map_err(|_| 1)? + }; + let sk_err_soft = unsafe { + bpf_probe_read_kernel::(sk_pointer.add(sk_err_soft_offset) as *const i32) + .map_err(|_| 1)? + }; + let sk_backlog_len = unsafe { + bpf_probe_read_kernel::(sk_pointer.add(sk_backlog_len_offset) as *const i32) + .map_err(|_| 1)? + }; + let sk_write_memory_queued = unsafe { + bpf_probe_read_kernel::(sk_pointer.add(sk_write_memory_queued_offset) as *const i32) + .map_err(|_| 1)? + }; + let sk_receive_buffer_size = unsafe { + bpf_probe_read_kernel::(sk_pointer.add(sk_receive_buffer_size_offset) as *const i32) + .map_err(|_| 1)? + }; + let sk_ack_backlog = unsafe { + bpf_probe_read_kernel::(sk_pointer.add(sk_ack_backlog_offset) as *const u32) + .map_err(|_| 1)? + }; + let sk_drops = unsafe { + bpf_probe_read_kernel::(sk_pointer.add(sk_drops_offset) as *const i32) + .map_err(|_| 1)? + }; + + let packet_loss_metrics = PacketLossMetrics { + tgid: tgid, + comm: comm, + ts_us: ts_us, + sk_err: sk_err, + sk_err_soft: sk_err_soft, + sk_backlog_len: sk_backlog_len, + sk_write_memory_queued: sk_write_memory_queued, + sk_receive_buffer_size: sk_receive_buffer_size, + sk_ack_backlog: sk_ack_backlog, + sk_drops: sk_drops, + }; + + Ok(packet_loss_metrics) +} + +pub fn on_connect(ctx: ProbeContext) -> Result<(), i64> { + let sk = ctx.arg::<*mut bindings::sock>(0).ok_or(1i64)?; + if sk.is_null() { + return Err(1); + } + + let tgid = (unsafe { bpf_get_current_pid_tgid() } >> 32) as u32; + let mut start = TimeStampStartInfo { + comm: [0; TASK_COMM_LEN], + ts_ns: unsafe { bpf_ktime_get_ns() }, + tgid, + }; + unsafe { + let comm_result = bpf_get_current_comm(); + if let Ok(comm) = comm_result { + start.comm.copy_from_slice(&comm); + } + let map_ptr = &raw mut TIME_STAMP_START; + (*map_ptr) + .insert(&(sk as *mut core::ffi::c_void), &start, 0) + .map_err(|_| 1)?; + } + Ok(()) +} + +pub fn on_rcv_state_process(ctx: ProbeContext) -> Result<(), i64> { + // On some kernels, kprobe wrapper puts `sk` at arg0; on others arg1. + let sk = ctx.arg::<*mut bindings::sock>(0).unwrap_or(ptr::null_mut()); + let sk = if sk.is_null() { + ctx.arg::<*mut bindings::sock>(1).ok_or(1i64)? + } else { + sk + }; + + if sk.is_null() { + return Err(1); + } + + let skc_daddr_off = 0; + let skc_rcv_saddr_off = 4; + let skc_dport_off = 12; + let skc_num_off = 14; + let skc_family_off = 16; + let skc_state_off = 18; + let skc_v6_daddr_off = 56; + let skc_v6_rcv_saddr_off = 72; + + let state = unsafe { bpf_probe_read_kernel::((sk as usize + skc_state_off) as *const u8) } + .map_err(|_| 1)?; + + if state != TCP_SYN_SENT { + return Ok(()); + } + + let start = unsafe { + let map_ptr = &raw const TIME_STAMP_START; + (*map_ptr).get(&((sk as usize) as *mut core::ffi::c_void)) + } + .ok_or(1i64)?; + let now = unsafe { bpf_ktime_get_ns() }; + let delta = now as i64 - start.ts_ns as i64; + if delta <= 0 { + unsafe { + let map_ptr = &raw mut TIME_STAMP_START; + let _ = (*map_ptr).remove(&((sk as usize) as *mut core::ffi::c_void)); + } + return Ok(()); + } + + let mut ev = TimeStampEvent { + delta_us: (delta as u64) / 1_000, + ts_us: now / 1_000, + tgid: start.tgid, + comm: start.comm, + lport: 0, + dport_be: 0, + af: 0, + saddr_v4: 0, + daddr_v4: 0, + saddr_v6: [0; 4], + daddr_v6: [0; 4], + }; + + // family, ports + ev.af = unsafe { + bpf_probe_read_kernel::((sk as usize + skc_family_off) as *const u16).map_err(|_| 1)? + }; + ev.lport = unsafe { + bpf_probe_read_kernel::((sk as usize + skc_num_off) as *const u16).map_err(|_| 1)? + }; + ev.dport_be = unsafe { + bpf_probe_read_kernel::((sk as usize + skc_dport_off) as *const u16).map_err(|_| 1)? + }; + + if ev.af == AF_INET { + ev.saddr_v4 = unsafe { + bpf_probe_read_kernel::((sk as usize + skc_rcv_saddr_off) as *const u32) + .map_err(|_| 1)? + }; + ev.daddr_v4 = unsafe { + bpf_probe_read_kernel::((sk as usize + skc_daddr_off) as *const u32) + .map_err(|_| 1)? + }; + } else { + // read 16 bytes as four u32 words + for i in 0..4 { + ev.saddr_v6[i] = unsafe { + bpf_probe_read_kernel::( + (sk as usize + skc_v6_rcv_saddr_off + i * 4) as *const u32, + ) + .map_err(|_| 1)? + }; + ev.daddr_v6[i] = unsafe { + bpf_probe_read_kernel::((sk as usize + skc_v6_daddr_off + i * 4) as *const u32) + .map_err(|_| 1)? + }; + } + } + + // emit + cleanup + unsafe { + bpf_perf_event_output( + ctx.as_ptr(), + &raw const TIME_STAMP_EVENTS as *const _ as *mut _, + 0, // BPF_F_CURRENT_CPU + &ev as *const _ as *mut _, + (mem::size_of::() as u32).into(), + ); + let map_ptr = &raw mut TIME_STAMP_START; + let _ = (*map_ptr).remove(&((sk as usize) as *mut core::ffi::c_void)); + } + + Ok(()) +}