Skip to content

Commit 5fc8625

Browse files
harrison001claude
andcommitted
Simplify ring buffer initialization by removing unnecessary spawn_blocking
- Changed ring_buffer_processor to take &Map reference instead of owned Map - Removed spawn_blocking wrapper around ring buffer creation since Map reference is sufficient - Simplified error handling by building ring buffer synchronously in the async context 🤖 Generated with [Claude Code](https://claude.ai/code) Co-Authored-By: Claude <noreply@anthropic.com>
1 parent 533adb8 commit 5fc8625

1 file changed

Lines changed: 19 additions & 27 deletions

File tree

kernel-agent/src/lib.rs

Lines changed: 19 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -228,9 +228,9 @@ impl EbpfLoader {
228228
let shutdown = Arc::clone(&self.shutdown_signal);
229229
let config = self.config.clone();
230230

231-
// Spawn the processor with the map (transfer ownership)
231+
// Spawn the processor with the map reference
232232
tokio::spawn(Self::ring_buffer_processor(
233-
rb_map,
233+
&rb_map,
234234
sender,
235235
metrics,
236236
rate_limiter,
@@ -243,7 +243,7 @@ impl EbpfLoader {
243243

244244
#[cfg(target_os = "linux")]
245245
async fn ring_buffer_processor(
246-
rb_map: libbpf_rs::Map,
246+
rb_map: &libbpf_rs::Map,
247247
sender: mpsc::Sender<RawEvent>,
248248
metrics: Arc<RwLock<EbpfMetrics>>,
249249
rate_limiter: Arc<Semaphore>,
@@ -255,33 +255,25 @@ impl EbpfLoader {
255255
// Create a channel for events from the ring buffer callback
256256
let (event_tx, mut event_rx) = mpsc::unbounded_channel::<Vec<u8>>();
257257

258-
// Build ring buffer in a blocking context to avoid Send issues
259-
let rb_result = tokio::task::spawn_blocking(move || -> anyhow::Result<_> {
260-
let mut rb_builder = RingBufferBuilder::new();
261-
262-
rb_builder.add(&rb_map, {
263-
let event_tx = event_tx.clone();
264-
move |data: &[u8]| {
265-
// Copy data and send through channel for async processing
266-
if let Err(_) = event_tx.send(data.to_vec()) {
267-
// Channel closed, stop processing
268-
return -1;
269-
}
270-
0
258+
// Build ring buffer synchronously since Map doesn't implement Send
259+
let mut rb_builder = RingBufferBuilder::new();
260+
261+
rb_builder.add(rb_map, {
262+
let event_tx = event_tx.clone();
263+
move |data: &[u8]| {
264+
// Copy data and send through channel for async processing
265+
if let Err(_) = event_tx.send(data.to_vec()) {
266+
// Channel closed, stop processing
267+
return -1;
271268
}
272-
});
273-
274-
rb_builder.build().map_err(|e| anyhow::anyhow!("Ring buffer build error: {}", e))
275-
}).await;
276-
277-
let rb = match rb_result {
278-
Ok(Ok(rb)) => rb,
279-
Ok(Err(e)) => {
280-
error!("Cannot create ring buffer: {}", e);
281-
return;
269+
0
282270
}
271+
});
272+
273+
let rb = match rb_builder.build() {
274+
Ok(rb) => rb,
283275
Err(e) => {
284-
error!("Ring buffer task error: {}", e);
276+
error!("Cannot create ring buffer: {}", e);
285277
return;
286278
}
287279
};

0 commit comments

Comments
 (0)