use std::collections::VecDeque;
use std::io;
use std::os::windows::io::{AsRawHandle, BorrowedHandle, FromRawHandle, OwnedHandle};
use std::sync::{Arc, Condvar, Mutex, MutexGuard, OnceLock};
use std::time::{Duration, Instant};
use windows_sys::Win32::Foundation::{DUPLICATE_SAME_ACCESS, DuplicateHandle, FALSE, HANDLE};
use windows_sys::Win32::System::Threading::{GetCurrentProcess, ResetEvent, SetEvent};
use windows_threadpool_sys::wait::WaitableHandle;
use crate::completion::{Completion, EnumerationId, TerminalOutcome};
pub(crate) const MINIMUM_COMPLETION_CAPACITY: usize = 2;
pub(crate) struct CompletionRing {
state: Mutex<RingState>,
arrived: Condvar,
doorbell: OnceLock<WaitableHandle>,
}
struct RingState {
queue: VecDeque<Completion>,
capacity: usize,
reserved: usize,
sessions: usize,
active: usize,
}
impl RingState {
fn data_room(&self) -> usize {
self.capacity - self.queue.len() - self.reserved
}
fn can_reserve(&self) -> bool {
self.data_room() > 0 && self.reserved + 1 < self.capacity
}
fn closed(&self) -> bool {
self.sessions == 0 && self.active == 0
}
fn pending(&self) -> bool {
!self.queue.is_empty() || self.closed()
}
}
impl CompletionRing {
pub(crate) fn new(capacity: usize) -> Self {
assert!(
capacity >= MINIMUM_COMPLETION_CAPACITY,
"a completion ring must be able to hold one terminal and one entry"
);
Self {
state: Mutex::new(RingState {
queue: VecDeque::new(),
capacity,
reserved: 0,
sessions: 1,
active: 0,
}),
arrived: Condvar::new(),
doorbell: OnceLock::new(),
}
}
fn lock(&self) -> MutexGuard<'_, RingState> {
self.state
.lock()
.unwrap_or_else(|poison| poison.into_inner())
}
pub(crate) fn capacity(&self) -> usize {
self.lock().capacity
}
pub(crate) fn len(&self) -> usize {
self.lock().queue.len()
}
pub(crate) fn has_data_room(&self) -> bool {
self.lock().data_room() > 0
}
pub(crate) fn is_closed(&self) -> bool {
self.lock().closed()
}
#[cfg(test)]
pub(crate) fn reserved(&self) -> usize {
self.lock().reserved
}
#[cfg(test)]
pub(crate) fn is_pending(&self) -> bool {
self.lock().pending()
}
pub(crate) fn add_session(&self) {
self.lock().sessions += 1;
}
pub(crate) fn remove_session(&self) {
{
let mut state = self.lock();
state.sessions -= 1;
self.refresh_doorbell(&state);
}
self.arrived.notify_all();
}
pub(crate) fn reserve_terminal(
self: &Arc<Self>,
enumeration: EnumerationId,
) -> Option<TerminalSlot> {
let mut state = self.lock();
if !state.can_reserve() {
return None;
}
state.reserved += 1;
state.active += 1;
Some(TerminalSlot {
ring: Arc::clone(self),
enumeration,
})
}
#[allow(
clippy::result_large_err,
reason = "handing the record back by value is the point: it is returned \
intact rather than boxed, reallocated, or dropped"
)]
pub(crate) fn try_send_entry(&self, record: Completion) -> Result<(), Completion> {
{
let mut state = self.lock();
if state.data_room() == 0 {
return Err(record);
}
state.queue.push_back(record);
self.refresh_doorbell(&state);
}
self.arrived.notify_all();
Ok(())
}
pub(crate) fn try_take(&self) -> Option<Completion> {
let mut state = self.lock();
let record = state.queue.pop_front();
if record.is_some() {
self.refresh_doorbell(&state);
}
record
}
pub(crate) fn take_blocking(&self, timeout: Option<Duration>) -> Option<Completion> {
let deadline = timeout.map(|timeout| Instant::now() + timeout);
let mut state = self.lock();
loop {
if let Some(record) = state.queue.pop_front() {
self.refresh_doorbell(&state);
return Some(record);
}
if state.closed() {
return None;
}
state = match deadline {
None => self
.arrived
.wait(state)
.unwrap_or_else(|poison| poison.into_inner()),
Some(deadline) => {
let remaining = deadline.checked_duration_since(Instant::now())?;
self.arrived
.wait_timeout(state, remaining)
.unwrap_or_else(|poison| poison.into_inner())
.0
}
};
}
}
pub(crate) fn doorbell(&self) -> io::Result<BorrowedHandle<'_>> {
if self.doorbell.get().is_none() {
let state = self.lock();
if self.doorbell.get().is_none() {
let event = WaitableHandle::event(true, state.pending())?;
let _ = self.doorbell.set(event);
}
}
Ok(self
.doorbell
.get()
.expect("the doorbell was just created")
.handle())
}
pub(crate) fn doorbell_owned(&self) -> io::Result<OwnedHandle> {
duplicate(self.doorbell()?)
}
fn refresh_doorbell(&self, state: &RingState) {
let Some(event) = self.doorbell.get() else {
return;
};
let handle = event.handle().as_raw_handle() as HANDLE;
unsafe {
if state.pending() {
SetEvent(handle);
} else {
ResetEvent(handle);
}
}
}
fn release_reservation(&self) {
{
let mut state = self.lock();
state.reserved -= 1;
state.active -= 1;
self.refresh_doorbell(&state);
}
self.arrived.notify_all();
}
fn fill_reservation(&self, record: Completion) {
{
let mut state = self.lock();
state.reserved -= 1;
state.active -= 1;
state.queue.push_back(record);
self.refresh_doorbell(&state);
}
self.arrived.notify_all();
}
}
pub(crate) struct TerminalSlot {
ring: Arc<CompletionRing>,
enumeration: EnumerationId,
}
impl TerminalSlot {
pub(crate) fn send(self, outcome: TerminalOutcome) {
let record = Completion::Terminal {
enumeration: self.enumeration,
outcome,
};
self.ring.fill_reservation(record);
std::mem::forget(self);
}
}
impl Drop for TerminalSlot {
fn drop(&mut self) {
self.ring.release_reservation();
}
}
fn duplicate(handle: BorrowedHandle<'_>) -> io::Result<OwnedHandle> {
let mut copy: HANDLE = std::ptr::null_mut();
let ok = unsafe {
DuplicateHandle(
GetCurrentProcess(),
handle.as_raw_handle() as HANDLE,
GetCurrentProcess(),
&raw mut copy,
0,
FALSE,
DUPLICATE_SAME_ACCESS,
)
};
if ok == 0 {
return Err(io::Error::last_os_error());
}
Ok(unsafe { OwnedHandle::from_raw_handle(copy as _) })
}
#[cfg(test)]
mod tests;