#![expect(
clippy::unwrap_used,
reason = "ratchet MOIRAI-UNWRAP-1: pre-existing debt"
)]
use super::channel::MpmcChannel;
use crate::channel::error::{ChannelError, Result};
use std::sync::atomic::Ordering;
impl<T: Send> MpmcChannel<T> {
pub(super) fn send_unbounded(&self, value: T) -> Result<()> {
let (mutex, not_full, not_empty) = &self.state;
let guard = mutex.lock().unwrap();
let mut guard = self.wait_while(guard, not_full, &self.sender_waiter_count, |state| {
!state.closed && state.capacity.is_some_and(|cap| state.queue.len() >= cap)
});
if guard.closed {
return Err(ChannelError::Closed);
}
guard.queue.push_back(value);
drop(guard);
if self.receiver_waiter_count.load(Ordering::Acquire) > 0 {
not_empty.notify_one();
}
Ok(())
}
pub(super) fn try_send_unbounded(&self, value: T) -> Result<()> {
let (mutex, _, not_empty) = &self.state;
let mut guard = mutex.lock().unwrap();
if guard.closed {
return Err(ChannelError::Closed);
}
if guard.capacity.is_some_and(|cap| guard.queue.len() >= cap) {
return Err(ChannelError::Full);
}
guard.queue.push_back(value);
drop(guard);
if self.receiver_waiter_count.load(Ordering::Acquire) > 0 {
not_empty.notify_one();
}
Ok(())
}
pub(super) fn recv_unbounded(&self) -> Result<T> {
let (mutex, not_full, not_empty) = &self.state;
let guard = mutex.lock().unwrap();
let mut guard = self.wait_while(guard, not_empty, &self.receiver_waiter_count, |state| {
state.queue.is_empty() && !state.closed
});
match guard.queue.pop_front() {
Some(value) => {
drop(guard);
if self.sender_waiter_count.load(Ordering::Acquire) > 0 {
not_full.notify_one();
}
Ok(value)
}
_ => Err(ChannelError::Closed),
}
}
pub(super) fn try_recv_unbounded(&self) -> Result<T> {
let (mutex, not_full, _) = &self.state;
let mut guard = mutex.lock().unwrap();
match guard.queue.pop_front() {
Some(value) => {
drop(guard);
if self.sender_waiter_count.load(Ordering::Acquire) > 0 {
not_full.notify_one();
}
Ok(value)
}
_ => {
if guard.closed {
Err(ChannelError::Closed)
} else {
Err(ChannelError::Empty)
}
}
}
}
pub(super) fn is_empty_unbounded(&self) -> bool {
let (mutex, _, _) = &self.state;
let guard = mutex.lock().unwrap();
guard.queue.is_empty()
}
pub(super) fn is_full_unbounded(&self) -> bool {
let (mutex, _, _) = &self.state;
let guard = mutex.lock().unwrap();
guard.capacity.is_some_and(|cap| guard.queue.len() >= cap)
}
pub(super) fn capacity_unbounded(&self) -> Option<usize> {
let (mutex, _, _) = &self.state;
let guard = mutex.lock().unwrap();
guard.capacity
}
}