Skip to content

Commit 757b571

Browse files
committed
perf(memtrack): bound ring reads and resume paused pids at low fill
A regular poll tick drained the whole ring in one `poll(ZERO)`, and pressure-stopped processes were resumed only after that tick ended with the ring completely empty. On a large memory benchmark suite with stack capture, single ticks of the stacks ring ran 0.4-1.7 s, every pressure stop landed inside one, and stopped processes waited 238-513 ms (median) to resume, 5-7% of the run. Ticks now read in chunks of 1024 records with `consume_raw_n` and check the fill between chunks. Paused processes resume once the ring is below a quarter full instead of empty; BPF stops them at three quarters, so the two thresholds leave room between stop and resume. `drain()` still reads the ring fully, and shutdown still releases unconditionally. Closes COD-3659
1 parent 9cc4405 commit 757b571

3 files changed

Lines changed: 31 additions & 16 deletions

File tree

‎crates/memtrack/src/ebpf/memtrack/maps.rs‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -93,8 +93,8 @@ impl MemtrackBpf {
9393
))
9494
}
9595

96-
/// Callback that resumes every pressure-stopped process.
97-
pub(super) fn on_ring_drained(&self) -> Box<dyn Fn() + Send> {
96+
/// Callback that resumes pressure-stopped processes at the low watermark.
97+
pub(super) fn on_ring_low_fill(&self) -> Box<dyn Fn() + Send> {
9898
let stopped = self.stopped.clone();
9999
Box::new(move || {
100100
if let Err(error) = stopped.release_pressure() {

‎crates/memtrack/src/ebpf/memtrack/mod.rs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -282,7 +282,7 @@ impl MemtrackBpf {
282282
resolve,
283283
tx,
284284
poll_interval_ms,
285-
Some(self.on_ring_drained()),
285+
Some(self.on_ring_low_fill()),
286286
))
287287
}
288288

‎crates/memtrack/src/ebpf/poller.rs‎

Lines changed: 28 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,9 @@ fn consume_all(ringbuf: &RingBuffer, ring: *mut libbpf_sys::ring) {
4141
}
4242
}
4343

44+
/// Records read per regular-tick chunk. Bounds the work between fill checks.
45+
const CONSUME_CHUNK_RECORDS: usize = BATCH_ITEMS;
46+
4447
fn poll_iteration<T>(
4548
control: std::result::Result<Sender<()>, RecvTimeoutError>,
4649
consume: impl FnOnce(),
@@ -72,8 +75,9 @@ fn poll_iteration<T>(
7275
/// Polls a BPF ring buffer in a background thread, parsing raw entries with a
7376
/// user-supplied closure and forwarding them to an mpsc channel in batches.
7477
///
75-
/// The poll thread runs until the poller is dropped, doing a final full
76-
/// `consume()` on shutdown so no buffered entries are lost.
78+
/// Regular ticks consume bounded chunks and check the low-fill release
79+
/// watermark between chunks. Shutdown performs a final full `consume()` so no
80+
/// buffered entries are lost.
7781
pub struct RingBufferPoller {
7882
ctl: Option<Sender<Sender<()>>>,
7983
poll_thread: Option<JoinHandle<()>>,
@@ -85,7 +89,7 @@ impl RingBufferPoller {
8589
parse: F,
8690
tx: Sender<Vec<T>>,
8791
poll_interval_ms: u64,
88-
on_drained: Option<Box<dyn Fn() + Send>>,
92+
on_low_fill: Option<Box<dyn Fn() + Send>>,
8993
) -> Result<Self>
9094
where
9195
M: MapCore,
@@ -123,23 +127,34 @@ impl RingBufferPoller {
123127
// SAFETY: the built `RingBuffer` holds exactly the one ring added above.
124128
let ring =
125129
unsafe { libbpf_sys::ring_buffer__ring(ringbuf.as_libbpf_object().as_ptr(), 0) };
130+
// Resume below 1/4 fill; BPF stops at 3/4, which leaves hysteresis.
131+
let release_if_low = || {
132+
if let Some(on_low_fill) = &on_low_fill
133+
&& unsafe { libbpf_sys::ring__avail_data_size(ring) }
134+
< unsafe { libbpf_sys::ring__size(ring) } / 4
135+
{
136+
on_low_fill();
137+
}
138+
};
126139
while poll_iteration(
127140
ctl_rx.recv_timeout(Duration::from_millis(poll_interval_ms)),
128141
|| consume_all(&ringbuf, ring),
129142
|| {
130-
let _ = ringbuf.poll(Duration::ZERO);
143+
// A short or failed chunk ends the tick; the loop body below
144+
// checks the fill after it.
145+
while ringbuf.consume_raw_n(CONSUME_CHUNK_RECORDS)
146+
== CONSUME_CHUNK_RECORDS as i32
147+
{
148+
release_if_low();
149+
}
131150
},
132151
&batch,
133152
&tx,
134153
) {
135-
if let Some(on_drained) = &on_drained
136-
&& unsafe { libbpf_sys::ring__avail_data_size(ring) } == 0
137-
{
138-
on_drained();
139-
}
154+
release_if_low();
140155
}
141-
if let Some(on_drained) = &on_drained {
142-
on_drained();
156+
if let Some(on_low_fill) = &on_low_fill {
157+
on_low_fill();
143158
}
144159
});
145160

@@ -192,7 +207,7 @@ impl ThreadedRingBufferPoller {
192207
resolve: R,
193208
tx: Sender<Vec<U>>,
194209
poll_interval_ms: u64,
195-
on_drained: Option<Box<dyn Fn() + Send>>,
210+
on_low_fill: Option<Box<dyn Fn() + Send>>,
196211
) -> Result<Self>
197212
where
198213
M: MapCore,
@@ -202,7 +217,7 @@ impl ThreadedRingBufferPoller {
202217
R: Fn(T) -> U + Send + 'static,
203218
{
204219
let (parsed_tx, parsed_rx) = mpsc::channel::<Vec<T>>();
205-
let ring = RingBufferPoller::new(rb_map, parse, parsed_tx, poll_interval_ms, on_drained)?;
220+
let ring = RingBufferPoller::new(rb_map, parse, parsed_tx, poll_interval_ms, on_low_fill)?;
206221
let resolver = std::thread::spawn(move || {
207222
for batch in parsed_rx {
208223
let resolved = batch.into_iter().map(&resolve).collect();

0 commit comments

Comments
 (0)