use std::future::Future;
use std::io::Result;
use std::pin::Pin;
use std::task::{Context, Poll};
use std::time::{Duration, Instant};
use orengine_macros::{poll_for_io_request, poll_for_time_bounded_io_request};
use crate as orengine;
use crate::io::io_request_data::IoRequestData;
use crate::io::sys::{AsRawFd, RawFd};
use crate::io::worker::{local_worker, IoWorker};
use crate::io::{Buffer, FixedBufferMut};
pub struct PeekBytes<'buf> {
fd: RawFd,
buf: &'buf mut [u8],
io_request_data: Option<IoRequestData>,
}
impl<'buf> PeekBytes<'buf> {
pub fn new(fd: RawFd, buf: &'buf mut [u8]) -> Self {
Self {
fd,
buf,
io_request_data: None,
}
}
}
impl Future for PeekBytes<'_> {
type Output = Result<usize>;
#[allow(
clippy::cast_possible_truncation,
reason = "It never peek more than u32::MAX"
)]
fn poll(self: Pin<&mut Self>, cx: &mut Context) -> Poll<Self::Output> {
let this = unsafe { self.get_unchecked_mut() };
let ret;
poll_for_io_request!((
local_worker().peek(
this.fd,
this.buf.as_mut_ptr(),
this.buf.len() as u32,
unsafe { this.io_request_data.as_mut().unwrap_unchecked() }
),
ret
));
}
}
unsafe impl Send for PeekBytes<'_> {}
pub struct PeekFixed<'buf> {
fd: RawFd,
ptr: *mut u8,
len: u32,
fixed_index: u16,
io_request_data: Option<IoRequestData>,
phantom_data: std::marker::PhantomData<&'buf Buffer>,
}
impl PeekFixed<'_> {
pub fn new(fd: RawFd, ptr: *mut u8, len: u32, fixed_index: u16) -> Self {
Self {
fd,
ptr,
len,
fixed_index,
io_request_data: None,
phantom_data: std::marker::PhantomData,
}
}
}
impl Future for PeekFixed<'_> {
type Output = Result<u32>;
#[allow(
clippy::cast_possible_truncation,
reason = "It never peek more than u32::MAX"
)]
fn poll(self: Pin<&mut Self>, cx: &mut Context) -> Poll<Self::Output> {
let this = unsafe { self.get_unchecked_mut() };
let ret;
poll_for_io_request!((
local_worker().peek_fixed(this.fd, this.ptr, this.len, this.fixed_index, unsafe {
this.io_request_data.as_mut().unwrap_unchecked()
}),
ret as u32
));
}
}
unsafe impl Send for PeekFixed<'_> {}
pub struct PeekBytesWithDeadline<'buf> {
fd: RawFd,
buf: &'buf mut [u8],
io_request_data: Option<IoRequestData>,
deadline: Instant,
}
impl<'buf> PeekBytesWithDeadline<'buf> {
pub fn new(fd: RawFd, buf: &'buf mut [u8], deadline: Instant) -> Self {
Self {
fd,
buf,
io_request_data: None,
deadline,
}
}
}
impl Future for PeekBytesWithDeadline<'_> {
type Output = Result<usize>;
#[allow(
clippy::cast_possible_truncation,
reason = "It never peek more than u32::MAX"
)]
fn poll(self: Pin<&mut Self>, cx: &mut Context) -> Poll<Self::Output> {
let this = unsafe { self.get_unchecked_mut() };
let worker = local_worker();
let ret;
poll_for_time_bounded_io_request!((
worker.peek_with_deadline(
this.fd,
this.buf.as_mut_ptr(),
this.buf.len() as u32,
unsafe { this.io_request_data.as_mut().unwrap_unchecked() },
&mut this.deadline
),
ret
));
}
}
unsafe impl Send for PeekBytesWithDeadline<'_> {}
pub struct PeekFixedWithDeadline<'buf> {
fd: RawFd,
ptr: *mut u8,
len: u32,
fixed_index: u16,
io_request_data: Option<IoRequestData>,
deadline: Instant,
phantom_data: std::marker::PhantomData<&'buf Buffer>,
}
impl PeekFixedWithDeadline<'_> {
pub fn new(fd: RawFd, ptr: *mut u8, len: u32, fixed_index: u16, deadline: Instant) -> Self {
Self {
fd,
ptr,
len,
fixed_index,
io_request_data: None,
deadline,
phantom_data: std::marker::PhantomData,
}
}
}
impl Future for PeekFixedWithDeadline<'_> {
type Output = Result<u32>;
#[allow(
clippy::cast_possible_truncation,
reason = "It never peek more than u32::MAX"
)]
fn poll(self: Pin<&mut Self>, cx: &mut Context) -> Poll<Self::Output> {
let this = unsafe { self.get_unchecked_mut() };
let worker = local_worker();
let ret;
poll_for_time_bounded_io_request!((
worker.peek_fixed_with_deadline(
this.fd,
this.ptr,
this.len,
this.fixed_index,
unsafe { this.io_request_data.as_mut().unwrap_unchecked() },
&mut this.deadline
),
ret as u32
));
}
}
unsafe impl Send for PeekFixedWithDeadline<'_> {}
pub trait AsyncPeek: AsRawFd {
#[inline(always)]
fn peek_bytes(&mut self, buf: &mut [u8]) -> impl Future<Output = Result<usize>> {
PeekBytes::new(self.as_raw_fd(), buf)
}
#[inline(always)]
async fn peek(&mut self, buf: &mut impl FixedBufferMut) -> Result<u32> {
if buf.is_fixed() {
PeekFixed::new(
self.as_raw_fd(),
buf.as_mut_ptr(),
buf.len_u32(),
buf.fixed_index(),
)
.await
} else {
#[allow(
clippy::cast_possible_truncation,
reason = "It never peek more than u32::MAX"
)]
PeekBytes::new(self.as_raw_fd(), buf.as_bytes_mut())
.await
.map(|r| r as u32)
}
}
#[inline(always)]
fn peek_bytes_with_deadline(
&mut self,
buf: &mut [u8],
deadline: Instant,
) -> impl Future<Output = Result<usize>> {
PeekBytesWithDeadline::new(self.as_raw_fd(), buf, deadline)
}
#[inline(always)]
async fn peek_with_deadline(
&mut self,
buf: &mut impl FixedBufferMut,
deadline: Instant,
) -> Result<u32> {
if buf.is_fixed() {
PeekFixedWithDeadline::new(
self.as_raw_fd(),
buf.as_mut_ptr(),
buf.len_u32(),
buf.fixed_index(),
deadline,
)
.await
} else {
#[allow(
clippy::cast_possible_truncation,
reason = "It never peek more than u32::MAX"
)]
PeekBytesWithDeadline::new(self.as_raw_fd(), buf.as_bytes_mut(), deadline)
.await
.map(|r| r as u32)
}
}
#[inline(always)]
fn peek_bytes_with_timeout(
&mut self,
buf: &mut [u8],
timeout: Duration,
) -> impl Future<Output = Result<usize>> {
self.peek_bytes_with_deadline(buf, Instant::now() + timeout)
}
#[inline(always)]
fn peek_with_timeout(
&mut self,
buf: &mut impl FixedBufferMut,
timeout: Duration,
) -> impl Future<Output = Result<u32>> {
self.peek_with_deadline(buf, Instant::now() + timeout)
}
#[inline(always)]
async fn peek_bytes_exact(&mut self, buf: &mut [u8]) -> Result<()> {
let mut peeked = 0;
while peeked < buf.len() {
peeked += self.peek_bytes(&mut buf[peeked..]).await?;
}
Ok(())
}
#[inline(always)]
async fn peek_exact(&mut self, buf: &mut impl FixedBufferMut) -> Result<()> {
if buf.is_fixed() {
let mut peeked = 0;
#[allow(
clippy::cast_possible_wrap,
reason = "We believe it never peek u32::MAX bytes"
)]
while peeked < buf.len_u32() {
peeked += PeekFixed::new(
self.as_raw_fd(),
unsafe { buf.as_mut_ptr().offset(peeked as isize) },
buf.len_u32() - peeked,
buf.fixed_index(),
)
.await?;
}
} else {
let mut peeked = 0;
let slice = buf.as_bytes_mut();
while peeked < slice.len() {
peeked += self.peek_bytes(&mut slice[peeked..]).await?;
}
}
Ok(())
}
#[inline(always)]
async fn peek_bytes_exact_with_deadline(
&mut self,
buf: &mut [u8],
deadline: Instant,
) -> Result<()> {
let mut peeked = 0;
while peeked < buf.len() {
peeked += self
.peek_bytes_with_deadline(&mut buf[peeked..], deadline)
.await?;
}
Ok(())
}
#[inline(always)]
async fn peek_exact_with_deadline(
&mut self,
buf: &mut impl FixedBufferMut,
deadline: Instant,
) -> Result<()> {
if buf.is_fixed() {
let mut peeked = 0;
#[allow(
clippy::cast_possible_wrap,
reason = "We believe it never peek u32::MAX bytes"
)]
while peeked < buf.len_u32() {
peeked += PeekFixedWithDeadline::new(
self.as_raw_fd(),
unsafe { buf.as_mut_ptr().offset(peeked as isize) },
buf.len_u32() - peeked,
buf.fixed_index(),
deadline,
)
.await?;
}
} else {
let mut peeked = 0;
let slice = buf.as_bytes_mut();
while peeked < slice.len() {
peeked += self
.peek_bytes_with_deadline(&mut slice[peeked..], deadline)
.await?;
}
}
Ok(())
}
#[inline(always)]
fn peek_bytes_exact_with_timeout(
&mut self,
buf: &mut [u8],
timeout: Duration,
) -> impl Future<Output = Result<()>> {
self.peek_bytes_exact_with_deadline(buf, Instant::now() + timeout)
}
#[inline(always)]
fn peek_exact_with_timeout(
&mut self,
buf: &mut impl FixedBufferMut,
timeout: Duration,
) -> impl Future<Output = Result<()>> {
self.peek_exact_with_deadline(buf, Instant::now() + timeout)
}
}