#![cfg(feature = "io-uring")]
use crate::io_uring_backend::buffer_manager::BufferRingManager;
use crate::ZmqError;
use crate::io_uring_backend::connection_handler::{
HandlerIoOps, HandlerSqeBlueprint, UringConnectionHandler, UringWorkerInterface,
};
use crate::io_uring_backend::ops::UserData;
use crate::io_uring_backend::worker::internal_op_tracker::InternalOpTracker;
use io_uring::cqueue;
use io_uring::cqueue::Entry as CqeResult;
use std::os::unix::io::RawFd;
pub const IOURING_CQE_F_MORE: u32 = 1 << 1;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum MultishotFlowState {
Reading,
Pausing,
Paused,
}
#[derive(Debug)]
pub(crate) struct MultishotReader {
fd: RawFd,
buffer_group_id: u16,
active_op_user_data: Option<UserData>,
is_active: bool,
cancel_op_user_data: Option<UserData>,
flow_state: MultishotFlowState,
}
impl MultishotReader {
pub fn new(fd: RawFd, buffer_group_id: u16) -> Self {
Self {
fd,
buffer_group_id,
active_op_user_data: None,
is_active: false,
cancel_op_user_data: None,
flow_state: MultishotFlowState::Paused,
}
}
pub fn is_reading(&self) -> bool {
matches!(self.flow_state, MultishotFlowState::Reading) && self.is_active
}
pub fn is_pausing(&self) -> bool {
matches!(self.flow_state, MultishotFlowState::Pausing)
}
pub fn is_paused(&self) -> bool {
matches!(self.flow_state, MultishotFlowState::Paused)
}
fn acknowledge_pause(&mut self) {
tracing::debug!(
"[MultishotReader FD={}] Kernel-side read confirmed stopped — flow state: Paused.",
self.fd
);
self.flow_state = MultishotFlowState::Paused;
self.is_active = false;
self.active_op_user_data = None;
self.cancel_op_user_data = None;
}
pub fn buffer_group_id(&self) -> u16 {
self.buffer_group_id
}
pub fn prepare_recv_multi_intent(&self) -> Option<HandlerSqeBlueprint> {
if self.is_active || self.cancel_op_user_data.is_some() {
return None;
}
Some(HandlerSqeBlueprint::RequestNewRingReadMultishot {
fd: self.fd,
bgid: self.buffer_group_id,
})
}
pub fn mark_operation_submitted(&mut self, op_user_data: UserData) {
self.active_op_user_data = Some(op_user_data);
self.is_active = true;
self.cancel_op_user_data = None;
self.flow_state = MultishotFlowState::Reading;
tracing::debug!(
"[MultishotReader FD={}] Marked as active in kernel with UserData {} — flow state: Reading.",
self.fd,
op_user_data
);
}
pub fn prepare_cancel_intent(&mut self) -> Option<HandlerSqeBlueprint> {
if !self.is_reading() || self.active_op_user_data.is_none() {
return None;
}
self.flow_state = MultishotFlowState::Pausing;
tracing::debug!(
"[MultishotReader FD={}] Transitioning to Pausing — submitting ASYNC_CANCEL for UserData {:?}.",
self.fd,
self.active_op_user_data
);
Some(HandlerSqeBlueprint::RequestNewAsyncCancel {
fd: self.fd,
target_user_data: self.active_op_user_data.unwrap(),
})
}
pub fn mark_cancellation_submitted(
&mut self,
cancel_sqe_user_data: UserData,
_target_op_user_data: UserData,
) {
self.cancel_op_user_data = Some(cancel_sqe_user_data);
tracing::debug!(
"[MultishotReader FD={}] Cancellation submitted with UserData {}.",
self.fd,
cancel_sqe_user_data
);
}
pub fn process_cqe(
&mut self,
cqe: &CqeResult,
buffer_manager: &BufferRingManager,
owner_handler: &mut dyn UringConnectionHandler,
worker_interface: &UringWorkerInterface<'_>,
_internal_op_tracker_ref: &mut InternalOpTracker,
) -> Result<(HandlerIoOps, bool), ZmqError> {
let cqe_ud = cqe.user_data();
let cqe_res = cqe.result();
let cqe_flags = cqe.flags();
let mut ops_to_return = HandlerIoOps::new();
if Some(cqe_ud) == self.active_op_user_data {
if !self.is_active {
tracing::warn!(
"[MultishotReader FD={}] CQE (ud {}) for active_op_user_data, but reader not marked active_in_kernel. State inconsistency?",
self.fd,
cqe_ud
);
self.is_active = true;
}
if cqe_res < 0 {
let errno = -cqe_res;
if errno == libc::ECANCELED as i32 {
tracing::debug!(
"[MultishotReader FD={}] Original multishot (ud {}) received -ECANCELED; confirming Paused.",
self.fd,
cqe_ud
);
self.acknowledge_pause();
return Ok((ops_to_return, true));
}
if errno == libc::ENOBUFS {
tracing::debug!(
"[MultishotReader FD={}] Buffer ring exhausted (ENOBUFS). Notifying handler.",
self.fd
);
self.active_op_user_data = None;
self.is_active = false;
owner_handler.on_buffer_ring_exhausted();
return Ok((ops_to_return, true));
}
tracing::error!(
"[MultishotReader FD={}] Error on active multishot read (ud {}): errno {}. Terminating multishot.",
self.fd,
cqe_ud,
errno
);
self.active_op_user_data = None;
self.is_active = false;
return Ok((ops_to_return.set_error_close(), true));
}
if cqe_res == 0 {
tracing::debug!(
"[MultishotReader FD={}] Clean EOF on multishot read (ud {}). Notifying handler.",
self.fd,
cqe_ud
);
self.active_op_user_data = None;
self.is_active = false;
ops_to_return =
owner_handler.process_ring_read_bytes(bytes::Bytes::new(), worker_interface);
return Ok((ops_to_return, true));
}
let buffer_id_opt = cqueue::buffer_select(cqe_flags);
if buffer_id_opt.is_none() {
tracing::error!(
"[MultishotReader FD={}] Multishot CQE (ud {}) missing F_BUFFER flag or invalid BID! Flags: {:x}",
self.fd,
cqe_ud,
cqe_flags
);
self.active_op_user_data = None;
self.is_active = false;
return Ok((ops_to_return.set_error_close(), true));
}
let buffer_id = buffer_id_opt.unwrap(); let bytes_read = cqe_res as usize;
if bytes_read > 0 {
match buffer_manager.take_and_replenish_buffer(buffer_id, bytes_read) {
Ok(owned_bytes) => {
ops_to_return = owner_handler.process_ring_read_bytes(owned_bytes, worker_interface);
}
Err(e) => {
tracing::error!(
"[MultishotReader FD={}] Failed to copy buffer ID {} ({} bytes): {:?}. Terminating multishot.",
self.fd,
buffer_id,
bytes_read,
e
);
self.active_op_user_data = None;
self.is_active = false;
return Ok((ops_to_return.set_error_close(), true));
}
}
} else {
tracing::info!(
"[MultishotReader FD={}] EOF on multishot read (ud {}). Terminating multishot.",
self.fd,
cqe_ud
);
if let Err(e) = buffer_manager.reprovide_buffer(buffer_id) {
tracing::warn!(
"[MultishotReader FD={}] reprovide_buffer({}) on EOF failed: {:?}",
self.fd,
buffer_id,
e
);
}
ops_to_return =
owner_handler.process_ring_read_bytes(bytes::Bytes::new(), worker_interface);
}
if (cqe_flags & IOURING_CQE_F_MORE) == 0 || bytes_read == 0 {
tracing::debug!(
"[MultishotReader FD={}] Multishot read (ud {}) finished (no MORE flag or EOF). Bytes read: {}",
self.fd,
cqe_ud,
bytes_read
);
self.active_op_user_data = None;
self.is_active = false;
return Ok((ops_to_return, true));
} else {
tracing::trace!(
"[MultishotReader FD={}] Multishot read (ud {}) has MORE flag. Op remains active.",
self.fd,
cqe_ud
);
return Ok((ops_to_return, false));
}
} else if Some(cqe_ud) == self.cancel_op_user_data {
let is_cancel_success = cqe_res == 0 || cqe_res == -libc::ENOENT;
if !is_cancel_success {
tracing::warn!(
"[MultishotReader FD={}] ASYNC_CANCEL (ud {}) returned unexpected result {}; proceeding to Paused anyway.",
self.fd,
cqe_ud,
cqe_res
);
} else {
tracing::debug!(
"[MultishotReader FD={}] ASYNC_CANCEL (ud {}) confirmed (res {}). Transitioning to Paused.",
self.fd,
cqe_ud,
cqe_res
);
}
self.acknowledge_pause();
return Ok((ops_to_return, true));
} else {
tracing::error!(
"[MultishotReader FD={}] process_cqe called with non-matching UserData (ud {}). This indicates a logic error in cqe_processor's delegation.",
self.fd,
cqe_ud
);
return Err(ZmqError::Internal(
"MultishotReader::process_cqe called with non-matching UserData".into(),
));
}
}
pub fn is_active(&self) -> bool {
self.is_active && self.cancel_op_user_data.is_none()
}
pub(crate) fn set_active(&mut self, user_data: UserData) {
if self.active_op_user_data == Some(user_data) {
self.is_active = true;
tracing::debug!(
"[MultishotReader FD={}] Marked as active with UserData {}.",
self.fd,
user_data
);
} else {
tracing::error!(
fd = self.fd,
expected = ?self.active_op_user_data,
got = user_data,
"MultishotReader set_active out-of-order: tracking ID mismatch"
);
}
}
pub(crate) fn matches_cqe_user_data(&self, cqe_user_data: UserData) -> bool {
self.active_op_user_data == Some(cqe_user_data)
|| self.cancel_op_user_data == Some(cqe_user_data)
}
pub(crate) fn set_inactive_due_to_close(&mut self) {
tracing::debug!(
"[MultishotReader FD={}] Marked as inactive due to FD closure.",
self.fd
);
self.is_active = false;
self.active_op_user_data = None;
self.cancel_op_user_data = None;
self.flow_state = MultishotFlowState::Paused;
}
}