saluki_common/resource_tracking/
stats.rs

1use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering::Relaxed};
2
3/// Statistics for an resource group.
4pub struct ResourceStats {
5    allocated_bytes: AtomicUsize,
6    allocated_objects: AtomicUsize,
7    deallocated_bytes: AtomicUsize,
8    deallocated_objects: AtomicUsize,
9    cpu_time_nanos: AtomicU64,
10}
11
12impl ResourceStats {
13    pub(super) const fn new() -> Self {
14        Self {
15            allocated_bytes: AtomicUsize::new(0),
16            allocated_objects: AtomicUsize::new(0),
17            deallocated_bytes: AtomicUsize::new(0),
18            deallocated_objects: AtomicUsize::new(0),
19            cpu_time_nanos: AtomicU64::new(0),
20        }
21    }
22
23    /// Returns `true` if the given group has allocated any memory at all.
24    pub fn has_allocated(&self) -> bool {
25        self.allocated_bytes.load(Relaxed) > 0
26    }
27
28    #[inline]
29    pub(super) fn track_allocation(&self, size: usize) {
30        self.allocated_bytes.fetch_add(size, Relaxed);
31        self.allocated_objects.fetch_add(1, Relaxed);
32    }
33
34    #[inline]
35    pub(super) fn track_deallocation(&self, size: usize) {
36        self.deallocated_bytes.fetch_add(size, Relaxed);
37        self.deallocated_objects.fetch_add(1, Relaxed);
38    }
39
40    #[inline]
41    pub(super) fn track_cpu_time(&self, nanos: u64) {
42        self.cpu_time_nanos.fetch_add(nanos, Relaxed);
43    }
44
45    /// Captures a snapshot of the current statistics based on the delta from a previous snapshot.
46    ///
47    /// This can be used to keep a single local snapshot of the last delta, and then both track the delta since that
48    /// snapshot, as well as update the snapshot to the current statistics.
49    ///
50    /// Callers should generally create their snapshot via [`ResourceStatsSnapshot::empty`] and then use this method
51    /// to get their snapshot delta, utilize those delta values in whatever way is necessary, and then merge the
52    /// snapshot delta into the primary snapshot via [`ResourceStatsSnapshot::merge`] to make it current.
53    pub fn snapshot_delta(&self, previous: &ResourceStatsSnapshot) -> ResourceStatsSnapshot {
54        ResourceStatsSnapshot {
55            allocated_bytes: self.allocated_bytes.load(Relaxed) - previous.allocated_bytes,
56            allocated_objects: self.allocated_objects.load(Relaxed) - previous.allocated_objects,
57            deallocated_bytes: self.deallocated_bytes.load(Relaxed) - previous.deallocated_bytes,
58            deallocated_objects: self.deallocated_objects.load(Relaxed) - previous.deallocated_objects,
59            cpu_time_nanos: self.cpu_time_nanos.load(Relaxed) - previous.cpu_time_nanos,
60        }
61    }
62}
63
64/// Snapshot of allocation statistics for a group.
65pub struct ResourceStatsSnapshot {
66    /// Number of allocated bytes since the last snapshot.
67    pub allocated_bytes: usize,
68
69    /// Number of allocated objects since the last snapshot.
70    pub allocated_objects: usize,
71
72    /// Number of deallocated bytes since the last snapshot.
73    pub deallocated_bytes: usize,
74
75    /// Number of deallocated objects since the last snapshot.
76    pub deallocated_objects: usize,
77
78    /// Cumulative CPU time in nanoseconds consumed by this group since the last snapshot.
79    pub cpu_time_nanos: u64,
80}
81
82impl ResourceStatsSnapshot {
83    /// Creates an empty `ResourceStatsSnapshot`.
84    pub const fn empty() -> Self {
85        Self {
86            allocated_bytes: 0,
87            allocated_objects: 0,
88            deallocated_bytes: 0,
89            deallocated_objects: 0,
90            cpu_time_nanos: 0,
91        }
92    }
93
94    /// Returns the number of live allocated bytes.
95    ///
96    /// Saturates at zero rather than underflowing. The allocation and deallocation counters this is derived from are
97    /// read independently, so a deallocation landing between the two reads can momentarily make the difference
98    /// negative.
99    pub fn live_bytes(&self) -> usize {
100        self.allocated_bytes.saturating_sub(self.deallocated_bytes)
101    }
102
103    /// Returns the number of live allocated objects.
104    ///
105    /// Saturates at zero, for the same reason as [`live_bytes`][Self::live_bytes].
106    pub fn live_objects(&self) -> usize {
107        self.allocated_objects.saturating_sub(self.deallocated_objects)
108    }
109
110    /// Merges `other` into `self`.
111    ///
112    /// This can be used to accumulate the total number of (de)allocated bytes and objects when handling the deltas
113    /// generated from `ResourceStats::consume`.
114    pub fn merge(&mut self, other: &Self) {
115        self.allocated_bytes += other.allocated_bytes;
116        self.allocated_objects += other.allocated_objects;
117        self.deallocated_bytes += other.deallocated_bytes;
118        self.deallocated_objects += other.deallocated_objects;
119        self.cpu_time_nanos += other.cpu_time_nanos;
120    }
121}
122
123/// Returns the current thread's CPU time in nanoseconds, or `None` if unavailable.
124#[cfg(target_os = "linux")]
125#[inline]
126pub(crate) fn thread_cpu_time_nanos() -> Option<u64> {
127    // NOTE: `CLOCK_THREAD_CPUTIME_ID` is not vDSO accelerated and degrades to a full syscall.
128    //
129    // Practically speaking, this ends up being roughly ~40-50x slower than a vDSO call at around 850ns or so. This is
130    // acceptable for our use case, because we're only tracking thread CPU time on task enter/exit, and only doing so on
131    // root task futures which are running for appreciable amounts of time so the rate at which we're calling this is fairly
132    // low.
133
134    let mut ts = libc::timespec { tv_sec: 0, tv_nsec: 0 };
135    // SAFETY: We pass a valid reference to the `timespec` struct, and `CLOCK_THREAD_CPUTIME_ID` has been available
136    // since Linux 2.6.12.
137    let ret = unsafe { libc::clock_gettime(libc::CLOCK_THREAD_CPUTIME_ID, &mut ts) };
138    if ret == 0 {
139        Some(ts.tv_sec as u64 * 1_000_000_000 + ts.tv_nsec as u64)
140    } else {
141        None
142    }
143}
144
145/// Returns the current thread's CPU time in nanoseconds, or `None` if unavailable.
146#[cfg(not(target_os = "linux"))]
147#[inline]
148pub(crate) fn thread_cpu_time_nanos() -> Option<u64> {
149    None
150}
151
152#[cfg(test)]
153mod tests {
154    use super::*;
155
156    #[cfg(target_os = "linux")]
157    #[test]
158    fn thread_cpu_time_returns_some() {
159        let t = thread_cpu_time_nanos();
160        assert!(t.is_some());
161        assert!(t.unwrap() > 0);
162    }
163
164    #[cfg(target_os = "linux")]
165    #[test]
166    fn thread_cpu_time_is_monotonic() {
167        let t1 = thread_cpu_time_nanos().unwrap();
168        // Do some work to consume CPU.
169        let mut sum = 0u64;
170        for i in 0..100_000 {
171            sum = sum.wrapping_add(i);
172        }
173        std::hint::black_box(sum);
174        let t2 = thread_cpu_time_nanos().unwrap();
175        assert!(t2 > t1);
176    }
177
178    #[cfg(not(target_os = "linux"))]
179    #[test]
180    fn thread_cpu_time_returns_none_on_non_linux() {
181        assert!(thread_cpu_time_nanos().is_none());
182    }
183}