use crate::io::config::IoWorkerConfig;
use crate::io::io_request_data::IoRequestData;
use crate::io::sys::{MessageRecvHeader, OpenHow, OsMessageHeader, RawFd, WorkerSys};
use crate::BUG_MESSAGE;
use nix::libc;
use nix::libc::sockaddr;
use std::cell::UnsafeCell;
use std::net::Shutdown;
use std::time::{Duration, Instant};
thread_local! {
pub(crate) static LOCAL_WORKER: UnsafeCell<Option<WorkerSys>> = const {
UnsafeCell::new(None)
};
}
pub(crate) fn get_local_worker_ref() -> &'static mut Option<WorkerSys> {
LOCAL_WORKER.with(|local_worker| unsafe { &mut *local_worker.get() })
}
pub(crate) unsafe fn init_local_worker(config: IoWorkerConfig) {
assert!(!get_local_worker_ref().is_some(), "{BUG_MESSAGE}");
*get_local_worker_ref() = Some(WorkerSys::new(config));
}
#[inline(always)]
pub(crate) fn local_worker() -> &'static mut WorkerSys {
#[cfg(debug_assertions)]
{
get_local_worker_ref().as_mut().expect(
"An attempt to call io-operation has failed, \
because an Executor has no io-worker. Look at the config of the Executor.",
)
}
#[cfg(not(debug_assertions))]
unsafe {
get_local_worker_ref().as_mut().unwrap_unchecked()
}
}
pub(crate) trait IoWorker {
fn new(config: IoWorkerConfig) -> Self;
fn register_time_bounded_io_task(
&mut self,
io_request_data: &IoRequestData,
deadline: &mut Instant,
);
fn deregister_time_bounded_io_task(&mut self, deadline: &Instant);
fn has_work(&self) -> bool;
fn must_poll(&mut self, timeout_option: Option<Duration>);
fn socket(
&mut self,
domain: socket2::Domain,
sock_type: socket2::Type,
protocol: socket2::Protocol,
request_ptr: *mut IoRequestData,
);
fn accept(
&mut self,
listen_fd: RawFd,
addr: *mut sockaddr,
addrlen: *mut libc::socklen_t,
request_ptr: *mut IoRequestData,
);
fn accept_with_deadline(
&mut self,
listen_fd: RawFd,
addr: *mut sockaddr,
addrlen: *mut libc::socklen_t,
request_ptr: *mut IoRequestData,
deadline: &mut Instant,
) {
self.register_time_bounded_io_task(unsafe { &*request_ptr } as _, deadline);
self.accept(listen_fd, addr, addrlen, request_ptr);
}
fn connect(
&mut self,
socket_fd: RawFd,
addr_ptr: *const sockaddr,
addr_len: libc::socklen_t,
request_ptr: *mut IoRequestData,
);
fn connect_with_deadline(
&mut self,
socket_fd: RawFd,
addr_ptr: *const sockaddr,
addr_len: libc::socklen_t,
request_ptr: *mut IoRequestData,
deadline: &mut Instant,
) {
self.register_time_bounded_io_task(unsafe { &*request_ptr } as _, deadline);
self.connect(socket_fd, addr_ptr, addr_len, request_ptr);
}
fn poll_fd_read(&mut self, fd: RawFd, request_ptr: *mut IoRequestData);
fn poll_fd_read_with_deadline(
&mut self,
fd: RawFd,
request_ptr: *mut IoRequestData,
deadline: &mut Instant,
) {
self.register_time_bounded_io_task(unsafe { &*request_ptr } as _, deadline);
self.poll_fd_read(fd, request_ptr);
}
fn poll_fd_write(&mut self, fd: RawFd, request_ptr: *mut IoRequestData);
fn poll_fd_write_with_deadline(
&mut self,
fd: RawFd,
request_ptr: *mut IoRequestData,
deadline: &mut Instant,
) {
self.register_time_bounded_io_task(unsafe { &*request_ptr } as _, deadline);
self.poll_fd_write(fd, request_ptr);
}
fn recv(&mut self, fd: RawFd, ptr: *mut u8, len: u32, request_ptr: *mut IoRequestData);
fn recv_fixed(
&mut self,
fd: RawFd,
ptr: *mut u8,
len: u32,
buf_index: u16,
request_ptr: *mut IoRequestData,
);
fn recv_with_deadline(
&mut self,
fd: RawFd,
ptr: *mut u8,
len: u32,
request_ptr: *mut IoRequestData,
deadline: &mut Instant,
) {
self.register_time_bounded_io_task(unsafe { &*request_ptr } as _, deadline);
self.recv(fd, ptr, len, request_ptr);
}
fn recv_fixed_with_deadline(
&mut self,
fd: RawFd,
ptr: *mut u8,
len: u32,
buf_index: u16,
request_ptr: *mut IoRequestData,
deadline: &mut Instant,
) {
self.register_time_bounded_io_task(unsafe { &*request_ptr } as _, deadline);
self.recv_fixed(fd, ptr, len, buf_index, request_ptr);
}
fn recv_from(
&mut self,
fd: RawFd,
msg_header: &mut MessageRecvHeader,
request_ptr: *mut IoRequestData,
);
fn recv_from_with_deadline(
&mut self,
fd: RawFd,
msg_header: &mut MessageRecvHeader,
request_ptr: *mut IoRequestData,
deadline: &mut Instant,
) {
self.register_time_bounded_io_task(unsafe { &*request_ptr } as _, deadline);
self.recv_from(fd, msg_header, request_ptr);
}
fn send(&mut self, fd: RawFd, ptr: *const u8, len: u32, request_ptr: *mut IoRequestData);
fn send_fixed(
&mut self,
fd: RawFd,
ptr: *const u8,
len: u32,
buf_index: u16,
request_ptr: *mut IoRequestData,
);
fn send_with_deadline(
&mut self,
fd: RawFd,
ptr: *const u8,
len: u32,
request_ptr: *mut IoRequestData,
deadline: &mut Instant,
) {
self.register_time_bounded_io_task(unsafe { &*request_ptr } as _, deadline);
self.send(fd, ptr, len, request_ptr);
}
fn send_fixed_with_deadline(
&mut self,
fd: RawFd,
ptr: *const u8,
len: u32,
buf_index: u16,
request_ptr: *mut IoRequestData,
deadline: &mut Instant,
) {
self.register_time_bounded_io_task(unsafe { &*request_ptr } as _, deadline);
self.send_fixed(fd, ptr, len, buf_index, request_ptr);
}
fn send_to(
&mut self,
fd: RawFd,
msg_header: *const OsMessageHeader,
request_ptr: *mut IoRequestData,
);
fn send_to_with_deadline(
&mut self,
fd: RawFd,
msg_header: *const OsMessageHeader,
request_ptr: *mut IoRequestData,
deadline: &mut Instant,
) {
self.register_time_bounded_io_task(unsafe { &*request_ptr } as _, deadline);
self.send_to(fd, msg_header, request_ptr);
}
fn peek(&mut self, fd: RawFd, ptr: *mut u8, len: u32, request_ptr: *mut IoRequestData);
fn peek_fixed(
&mut self,
fd: RawFd,
ptr: *mut u8,
len: u32,
buf_index: u16,
request_ptr: *mut IoRequestData,
);
fn peek_with_deadline(
&mut self,
fd: RawFd,
ptr: *mut u8,
len: u32,
request_ptr: *mut IoRequestData,
deadline: &mut Instant,
) {
self.register_time_bounded_io_task(unsafe { &*request_ptr } as _, deadline);
self.peek(fd, ptr, len, request_ptr);
}
fn peek_fixed_with_deadline(
&mut self,
fd: RawFd,
ptr: *mut u8,
len: u32,
buf_index: u16,
request_ptr: *mut IoRequestData,
deadline: &mut Instant,
) {
self.register_time_bounded_io_task(unsafe { &*request_ptr } as _, deadline);
self.peek_fixed(fd, ptr, len, buf_index, request_ptr);
}
fn peek_from(
&mut self,
fd: RawFd,
msg: &mut MessageRecvHeader,
request_ptr: *mut IoRequestData,
);
fn peek_from_with_deadline(
&mut self,
fd: RawFd,
msg: &mut MessageRecvHeader,
request_ptr: *mut IoRequestData,
deadline: &mut Instant,
) {
self.register_time_bounded_io_task(unsafe { &*request_ptr } as _, deadline);
self.peek_from(fd, msg, request_ptr);
}
fn shutdown(&mut self, fd: RawFd, how: Shutdown, request_ptr: *mut IoRequestData);
fn open(
&mut self,
path: *const libc::c_char,
open_how: *const OpenHow,
request_ptr: *mut IoRequestData,
);
fn fallocate(
&mut self,
fd: RawFd,
offset: u64,
len: u64,
flags: i32,
request_ptr: *mut IoRequestData,
);
fn sync_all(&mut self, fd: RawFd, request_ptr: *mut IoRequestData);
fn sync_data(&mut self, fd: RawFd, request_ptr: *mut IoRequestData);
fn read(&mut self, fd: RawFd, ptr: *mut u8, len: u32, request_ptr: *mut IoRequestData);
fn read_fixed(
&mut self,
fd: RawFd,
ptr: *mut u8,
len: u32,
buf_index: u16,
request_ptr: *mut IoRequestData,
);
fn pread(
&mut self,
fd: RawFd,
ptr: *mut u8,
len: u32,
offset: usize,
request_ptr: *mut IoRequestData,
);
fn pread_fixed(
&mut self,
fd: RawFd,
ptr: *mut u8,
len: u32,
buf_index: u16,
offset: usize,
request_ptr: *mut IoRequestData,
);
fn write(&mut self, fd: RawFd, ptr: *const u8, len: u32, request_ptr: *mut IoRequestData);
fn write_fixed(
&mut self,
fd: RawFd,
ptr: *const u8,
len: u32,
buf_index: u16,
request_ptr: *mut IoRequestData,
);
fn pwrite(
&mut self,
fd: RawFd,
ptr: *const u8,
len: u32,
offset: usize,
request_ptr: *mut IoRequestData,
);
fn pwrite_fixed(
&mut self,
fd: RawFd,
ptr: *const u8,
len: u32,
buf_index: u16,
offset: usize,
request_ptr: *mut IoRequestData,
);
fn close(&mut self, fd: RawFd, request_ptr: *mut IoRequestData);
fn rename(
&mut self,
old_path: *const libc::c_char,
new_path: *const libc::c_char,
request_ptr: *mut IoRequestData,
);
fn create_dir(&mut self, path: *const libc::c_char, mode: u32, request_ptr: *mut IoRequestData);
fn remove_file(&mut self, path: *const libc::c_char, request_ptr: *mut IoRequestData);
fn remove_dir(&mut self, path: *const libc::c_char, request_ptr: *mut IoRequestData);
}