use crate::{
RuntimeError,
futures::net::socket::configure,
modules::{fd::Fd, int_check::IntCheck},
};
use std::{
collections::VecDeque,
sync::{Arc, Mutex, MutexGuard, PoisonError},
};
pub(crate) struct Core<T> {
state: Mutex<State<T>>,
items_ring: Fd,
items_bell: Fd,
room_ring: Fd,
room_bell: Fd,
}
pub(crate) struct State<T> {
queue: VecDeque<T>,
capacity: Option<usize>,
senders: usize,
receivers: usize,
}
pub(crate) enum Refused<T> {
Closed,
Full(T),
}
impl<T> Core<T> {
pub(crate) fn open(capacity: Option<usize>) -> Result<Arc<Self>, RuntimeError> {
let mut pair = [0; 2];
unsafe { libc::socketpair(libc::AF_UNIX, libc::SOCK_STREAM, 0, pair.as_mut_ptr()) }
.check()?;
let (first, second) = (Fd::new(pair[0]), Fd::new(pair[1]));
configure(first.raw())?;
configure(second.raw())?;
let items_ring = Fd::new(unsafe { libc::dup(first.raw()) }.check()?);
let room_ring = Fd::new(unsafe { libc::dup(second.raw()) }.check()?);
let core = Self {
state: Mutex::new(State {
queue: VecDeque::new(),
capacity,
senders: 1,
receivers: 1,
}),
items_ring,
items_bell: second,
room_ring,
room_bell: first,
};
if capacity.is_some() {
ring(&core.room_ring);
}
Ok(Arc::new(core))
}
#[inline(always)]
pub(crate) fn items_bell(&self) -> libc::c_int {
self.items_bell.raw()
}
#[inline(always)]
pub(crate) fn room_bell(&self) -> libc::c_int {
self.room_bell.raw()
}
fn lock(&self) -> MutexGuard<'_, State<T>> {
self.state.lock().unwrap_or_else(PoisonError::into_inner)
}
pub(crate) fn push(&self, value: T) -> Result<(), Refused<T>> {
let mut state = self.lock();
if state.receivers == 0 {
return Err(Refused::Closed);
}
if state
.capacity
.is_some_and(|capacity| state.queue.len() >= capacity)
{
return Err(Refused::Full(value));
}
let was_empty = state.queue.is_empty();
state.queue.push_back(value);
if was_empty && state.senders > 0 {
ring(&self.items_ring);
}
if state.capacity == Some(state.queue.len()) {
drain(&self.room_bell);
}
Ok(())
}
pub(crate) fn pop(&self) -> Result<T, RuntimeError> {
let mut state = self.lock();
let was_full = state.capacity == Some(state.queue.len());
let Some(value) = state.queue.pop_front() else {
return match state.senders {
0 => Err(RuntimeError::Closed),
_ => Err(RuntimeError::NotReady),
};
};
if state.queue.is_empty() && state.senders > 0 {
drain(&self.items_bell);
}
if was_full {
ring(&self.room_ring);
}
Ok(value)
}
pub(crate) fn add_sender(&self) {
self.lock().senders += 1;
}
pub(crate) fn drop_sender(&self) {
let mut state = self.lock();
state.senders -= 1;
if state.senders == 0 && state.queue.is_empty() {
ring(&self.items_ring);
}
}
pub(crate) fn add_receiver(&self) {
self.lock().receivers += 1;
}
pub(crate) fn drop_receiver(&self) {
let mut state = self.lock();
state.receivers -= 1;
if state.receivers == 0 && state.capacity == Some(state.queue.len()) {
ring(&self.room_ring);
}
}
pub(crate) fn len(&self) -> usize {
self.lock().queue.len()
}
}
fn ring(fd: &Fd) {
let byte = 1u8;
let _ = unsafe { libc::write(fd.raw(), (&byte as *const u8).cast(), 1) };
}
fn drain(fd: &Fd) {
let mut bytes = [0u8; 16];
while unsafe { libc::read(fd.raw(), bytes.as_mut_ptr().cast(), bytes.len()) } > 0 {}
}