#![cfg(feature = "io-uring")]
use crate::io_uring_backend::provided_buffer_ring::ProvidedBufferRing;
use crate::ZmqError;
use bytes::Bytes;
use io_uring::IoUring;
use std::fmt;
pub struct BufferRingManager {
ring: ProvidedBufferRing,
}
impl fmt::Debug for BufferRingManager {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("BufferRingManager")
.field("ring", &self.ring)
.finish_non_exhaustive()
}
}
impl BufferRingManager {
pub fn new(
ring: &IoUring,
ring_entries: u16,
bgid: u16,
buffer_capacity: usize,
) -> Result<Self, ZmqError> {
tracing::info!(
"Initializing BufferRingManager with bgid: {}, requested_entries: {}, capacity_per_buffer: {}",
bgid,
ring_entries,
buffer_capacity
);
let ring = ProvidedBufferRing::new(ring, ring_entries, bgid, buffer_capacity)?;
tracing::info!(
"BufferRingManager: provided-buffer ring registered successfully for bgid: {}",
bgid
);
Ok(Self { ring })
}
pub fn group_id(&self) -> u16 {
self.ring.group_id()
}
pub fn unregister(&self, ring: &IoUring) {
self.ring.unregister(ring);
}
pub fn reprovide_buffer(&self, buffer_id: u16) -> Result<(), ZmqError> {
self.ring.reprovide(buffer_id)
}
pub fn take_and_replenish_buffer(
&self,
buffer_id: u16,
available_len: usize,
) -> Result<Bytes, ZmqError> {
self.ring.take(buffer_id, available_len)
}
}