use std::sync::mpsc::{Receiver, Sender};
use std::sync::Arc;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
mod cmdline;
pub use cmdline::{read_process_argv, read_process_cmdline};
mod file_handles;
pub use file_handles::read_process_file_handles;
mod process_watch;
pub(crate) use process_watch::ProcessWatchEmitter;
pub use process_watch::{
CaptureSource, DumpResult, ObservationGrade, ObservationPolicy, ProcessEvent, ProcessEventKind,
ProcessIdentity, ProcessObservation, ProcessObservationCapabilities, ProcessObservationError,
ProcessWatch, ProcessWatchConfigurationError, ProcessWatchCursor, ProcessWatchGap,
ProcessWatchLoss, ProcessWatchMatch, ProcessWatchRead, ProcessWatchSubscriber, StackCapture,
StackDump,
};
pub(crate) type DescendantPumpStop =
running_process_platform_internal::platform::process::DescendantMonitorStop;
#[non_exhaustive]
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum TraceScope {
LaunchedProcessTree,
SystemWide,
}
impl TraceScope {
pub const ALL: [TraceScope; 2] = [TraceScope::LaunchedProcessTree, TraceScope::SystemWide];
pub fn as_str(self) -> &'static str {
match self {
TraceScope::LaunchedProcessTree => "launched-process-tree",
TraceScope::SystemWide => "system-wide",
}
}
}
#[non_exhaustive]
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum EventCategory {
Lifecycle,
File,
Network,
Process,
}
impl EventCategory {
pub const ALL: [EventCategory; 4] = [
EventCategory::Lifecycle,
EventCategory::File,
EventCategory::Network,
EventCategory::Process,
];
pub fn as_str(self) -> &'static str {
match self {
EventCategory::Lifecycle => "lifecycle",
EventCategory::File => "file",
EventCategory::Network => "network",
EventCategory::Process => "process",
}
}
}
#[non_exhaustive]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CapabilitySupport {
Supported,
Partial,
Unavailable,
}
impl CapabilitySupport {
pub fn as_str(self) -> &'static str {
match self {
CapabilitySupport::Supported => "supported",
CapabilitySupport::Partial => "partial",
CapabilitySupport::Unavailable => "unavailable",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CategoryCapability {
pub category: EventCategory,
pub support: CapabilitySupport,
pub backend: &'static str,
pub reason: &'static str,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ObserverCapabilities {
scope: TraceScope,
categories: Vec<CategoryCapability>,
}
fn detect_file_backend(scope: TraceScope) -> (CapabilitySupport, &'static str, &'static str) {
detect_backend(
scope,
running_process_platform_internal::platform::process::ObserverCategory::File,
)
}
fn detect_network_backend(scope: TraceScope) -> (CapabilitySupport, &'static str, &'static str) {
detect_backend(
scope,
running_process_platform_internal::platform::process::ObserverCategory::Network,
)
}
fn detect_process_backend(scope: TraceScope) -> (CapabilitySupport, &'static str, &'static str) {
detect_backend(
scope,
running_process_platform_internal::platform::process::ObserverCategory::Process,
)
}
fn detect_backend(
scope: TraceScope,
category: running_process_platform_internal::platform::process::ObserverCategory,
) -> (CapabilitySupport, &'static str, &'static str) {
use running_process_platform_internal::platform::process::{
observer_backend, ObserverScope, ObserverSupport,
};
let scope = match scope {
TraceScope::SystemWide => ObserverScope::SystemWide,
TraceScope::LaunchedProcessTree => ObserverScope::LaunchedProcessTree,
};
let backend = observer_backend(scope, category);
let support = match backend.support {
ObserverSupport::Supported => CapabilitySupport::Supported,
ObserverSupport::Partial => CapabilitySupport::Partial,
ObserverSupport::Unavailable => CapabilitySupport::Unavailable,
};
(support, backend.backend, backend.reason)
}
impl ObserverCapabilities {
pub fn negotiate() -> Self {
Self::negotiate_for_scope(TraceScope::SystemWide)
}
pub fn negotiate_for_scope(scope: TraceScope) -> Self {
let categories = EventCategory::ALL
.iter()
.map(|&category| match category {
EventCategory::Lifecycle => CategoryCapability {
category,
support: CapabilitySupport::Supported,
backend: "portable-lifecycle",
reason: "started/exited emitted from the crate spawn and reap path",
},
EventCategory::File => {
let (support, backend, reason) = detect_file_backend(scope);
CategoryCapability {
category,
support,
backend,
reason,
}
}
EventCategory::Network => {
let (support, backend, reason) = detect_network_backend(scope);
CategoryCapability {
category,
support,
backend,
reason,
}
}
EventCategory::Process => {
let (support, backend, reason) = detect_process_backend(scope);
CategoryCapability {
category,
support,
backend,
reason,
}
}
})
.collect();
Self { scope, categories }
}
pub fn scope(&self) -> TraceScope {
self.scope
}
pub fn categories(&self) -> &[CategoryCapability] {
&self.categories
}
pub fn category(&self, category: EventCategory) -> &CategoryCapability {
self.categories
.iter()
.find(|entry| entry.category == category)
.expect("ObserverCapabilities always contains every EventCategory")
}
pub fn support(&self, category: EventCategory) -> CapabilitySupport {
self.category(category).support
}
pub fn is_supported(&self, category: EventCategory) -> bool {
self.support(category) == CapabilitySupport::Supported
}
pub fn to_table_rows(&self) -> Vec<[String; 4]> {
self.categories
.iter()
.map(|entry| {
[
entry.category.as_str().to_string(),
entry.support.as_str().to_string(),
entry.backend.to_string(),
entry.reason.to_string(),
]
})
.collect()
}
pub fn render_summary(&self) -> String {
let rows = self.to_table_rows();
let mut widths = [0usize; 3];
for row in &rows {
for (i, cell) in row[..3].iter().enumerate() {
widths[i] = widths[i].max(cell.len());
}
}
let mut out = format!("observer capabilities (scope={}):\n", self.scope.as_str());
for row in &rows {
out.push_str(&format!(
" {cat:<cw$} {sup:<sw$} {bk:<bw$} {reason}\n",
cat = row[0],
sup = row[1],
bk = row[2],
reason = row[3],
cw = widths[0],
sw = widths[1],
bw = widths[2],
));
}
out
}
}
#[non_exhaustive]
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ObserverEventKind {
Started,
Exited {
exit_code: i32,
},
DescendantStarted,
DescendantExited,
FileOpen {
path: std::path::PathBuf,
flags: u32,
},
FileWrite {
path: std::path::PathBuf,
byte_count: u64,
},
FileClose {
path: std::path::PathBuf,
},
FileUnlink {
path: std::path::PathBuf,
},
FileRename {
from: std::path::PathBuf,
to: std::path::PathBuf,
},
}
impl ObserverEventKind {
pub fn as_str(&self) -> &'static str {
match self {
ObserverEventKind::Started => "started",
ObserverEventKind::Exited { .. } => "exited",
ObserverEventKind::DescendantStarted => "descendant-started",
ObserverEventKind::DescendantExited => "descendant-exited",
ObserverEventKind::FileOpen { .. } => "file-open",
ObserverEventKind::FileWrite { .. } => "file-write",
ObserverEventKind::FileClose { .. } => "file-close",
ObserverEventKind::FileUnlink { .. } => "file-unlink",
ObserverEventKind::FileRename { .. } => "file-rename",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ObserverEvent {
pub category: EventCategory,
pub kind: ObserverEventKind,
pub pid: u32,
pub ppid: Option<u32>,
pub timestamp_ms: u128,
}
impl ObserverEvent {
fn now(category: EventCategory, kind: ObserverEventKind, pid: u32) -> Self {
Self::now_with_parent(category, kind, pid, None)
}
fn now_with_parent(
category: EventCategory,
kind: ObserverEventKind,
pid: u32,
ppid: Option<u32>,
) -> Self {
let timestamp_ms = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_millis())
.unwrap_or(0);
Self {
category,
kind,
pid,
ppid,
timestamp_ms,
}
}
pub fn new_now(category: EventCategory, kind: ObserverEventKind, pid: u32) -> Self {
Self::now(category, kind, pid)
}
pub fn new_now_with_parent(
category: EventCategory,
kind: ObserverEventKind,
pid: u32,
ppid: Option<u32>,
) -> Self {
Self::now_with_parent(category, kind, pid, ppid)
}
}
#[derive(Debug, Clone)]
pub struct ObserverConfig {
categories: Vec<EventCategory>,
}
impl ObserverConfig {
pub fn lifecycle() -> Self {
Self {
categories: vec![EventCategory::Lifecycle],
}
}
pub fn with_categories(categories: impl IntoIterator<Item = EventCategory>) -> Self {
Self {
categories: categories.into_iter().collect(),
}
}
pub fn observes(&self, category: EventCategory) -> bool {
self.categories.contains(&category)
}
pub fn categories(&self) -> &[EventCategory] {
&self.categories
}
}
pub fn observe_launched_tree(root_pid: u32, config: ObserverConfig) -> ObserverSubscriber {
let (emitter, subscriber) = ObserverEmitter::new(config);
crate::descendant_monitor::start_attached(root_pid, Some(&emitter));
subscriber
}
pub struct ObserverSubscriber {
rx: Receiver<ObserverEvent>,
descendant_stop: Arc<DescendantPumpStop>,
}
impl ObserverSubscriber {
#[cfg(feature = "client")]
pub(crate) fn from_receiver(rx: Receiver<ObserverEvent>) -> Self {
Self {
rx,
descendant_stop: Arc::new(DescendantPumpStop::new()),
}
}
pub fn recv(&self) -> Option<ObserverEvent> {
self.rx.recv().ok()
}
pub fn recv_timeout(
&self,
timeout: Duration,
) -> Result<ObserverEvent, std::sync::mpsc::RecvTimeoutError> {
self.rx.recv_timeout(timeout)
}
pub fn try_recv(&self) -> Option<ObserverEvent> {
self.rx.try_recv().ok()
}
pub fn drain(&self) -> Vec<ObserverEvent> {
let mut events = Vec::new();
while let Ok(event) = self.rx.try_recv() {
events.push(event);
}
events
}
pub fn receiver(&self) -> &Receiver<ObserverEvent> {
&self.rx
}
pub fn stop(&self) {
self.descendant_stop.stop();
}
}
impl Drop for ObserverSubscriber {
fn drop(&mut self) {
self.stop();
}
}
pub(crate) struct ObserverEmitter {
config: ObserverConfig,
tx: Sender<ObserverEvent>,
#[allow(dead_code)]
descendant_stop: Arc<DescendantPumpStop>,
}
impl ObserverEmitter {
pub(crate) fn new(config: ObserverConfig) -> (Self, ObserverSubscriber) {
let (tx, rx) = std::sync::mpsc::channel();
let descendant_stop = Arc::new(DescendantPumpStop::new());
(
Self {
config,
tx,
descendant_stop: Arc::clone(&descendant_stop),
},
ObserverSubscriber {
rx,
descendant_stop,
},
)
}
pub(crate) fn emit_started(&self, pid: u32) {
if !self.config.observes(EventCategory::Lifecycle) {
return;
}
let _ = self.tx.send(ObserverEvent::now(
EventCategory::Lifecycle,
ObserverEventKind::Started,
pid,
));
}
pub(crate) fn emit_exited(&self, pid: u32, exit_code: i32) {
if !self.config.observes(EventCategory::Lifecycle) {
return;
}
let _ = self.tx.send(ObserverEvent::now(
EventCategory::Lifecycle,
ObserverEventKind::Exited { exit_code },
pid,
));
}
#[allow(dead_code)]
pub(crate) fn descendant_sink(&self) -> Option<Sender<ObserverEvent>> {
if self.config.observes(EventCategory::Process) {
Some(self.tx.clone())
} else {
None
}
}
#[allow(dead_code)]
pub(crate) fn descendant_pump(
&self,
) -> Option<(Sender<ObserverEvent>, Arc<DescendantPumpStop>)> {
self.descendant_sink()
.map(|sink| (sink, Arc::clone(&self.descendant_stop)))
}
}
#[cfg(test)]
mod tests;