use std::collections::VecDeque;
use std::sync::{Mutex, MutexGuard};
use crate::completion::EnumerationId;
use crate::completion_ring::TerminalSlot;
use crate::engine::EngineState;
pub(crate) struct BeginMessage {
pub(crate) enumeration: EnumerationId,
pub(crate) engine: EngineState,
pub(crate) terminal: TerminalSlot,
pub(crate) retire: RetireSlot,
}
impl std::fmt::Debug for BeginMessage {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("BeginMessage")
.field("enumeration", &self.enumeration)
.field("engine", &self.engine)
.finish_non_exhaustive()
}
}
#[derive(Debug)]
pub(crate) enum ControlMessage {
Begin(Box<BeginMessage>),
Cancel(EnumerationId),
Retire(EnumerationId),
Abandon,
}
pub(crate) struct SubmissionRing {
state: Mutex<RingState>,
}
struct RingState {
queue: VecDeque<ControlMessage>,
capacity: usize,
reserved: usize,
draining: bool,
abandoned: bool,
}
impl RingState {
fn free(&self) -> usize {
self.capacity - self.queue.len() - self.reserved
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
pub(crate) enum SubmitRejection {
Full,
Abandoned,
}
impl SubmissionRing {
pub(crate) fn new(capacity: usize) -> Self {
assert!(
capacity > 0,
"a submission ring must hold at least one message"
);
Self {
state: Mutex::new(RingState {
queue: VecDeque::new(),
capacity,
reserved: 0,
draining: false,
abandoned: false,
}),
}
}
fn lock(&self) -> MutexGuard<'_, RingState> {
self.state
.lock()
.unwrap_or_else(|poison| poison.into_inner())
}
pub(crate) fn capacity(&self) -> usize {
self.lock().capacity
}
#[cfg(test)]
pub(crate) fn len(&self) -> usize {
self.lock().queue.len()
}
pub(crate) fn is_abandoned(&self) -> bool {
self.lock().abandoned
}
#[cfg(test)]
pub(crate) fn reserved(&self) -> usize {
self.lock().reserved
}
pub(crate) fn reserve_cancel(&self) -> Option<CancelSlot> {
let mut state = self.lock();
if state.free() == 0 {
return None;
}
state.reserved += 1;
Some(CancelSlot { _private: () })
}
pub(crate) fn reserve_retire(&self) -> Option<RetireSlot> {
let mut state = self.lock();
if state.free() == 0 {
return None;
}
state.reserved += 1;
Some(RetireSlot { _private: () })
}
pub(crate) fn reserve_abandon(&self) -> Option<AbandonSlot> {
let mut state = self.lock();
if state.free() == 0 {
return None;
}
state.reserved += 1;
Some(AbandonSlot { _private: () })
}
pub(crate) fn try_push(
&self,
message: ControlMessage,
) -> Result<PushOutcome, (ControlMessage, SubmitRejection)> {
let mut state = self.lock();
if state.abandoned {
return Err((message, SubmitRejection::Abandoned));
}
if state.free() == 0 {
return Err((message, SubmitRejection::Full));
}
state.queue.push_back(message);
Ok(claim_drain(&mut state))
}
pub(crate) fn push_cancel(&self, _slot: CancelSlot, enumeration: EnumerationId) -> PushOutcome {
let mut state = self.lock();
state.reserved -= 1;
state.queue.push_back(ControlMessage::Cancel(enumeration));
claim_drain(&mut state)
}
pub(crate) fn push_retire(&self, _slot: RetireSlot, enumeration: EnumerationId) -> PushOutcome {
let mut state = self.lock();
state.reserved -= 1;
state.queue.push_back(ControlMessage::Retire(enumeration));
claim_drain(&mut state)
}
pub(crate) fn push_abandon(&self, _slot: AbandonSlot) -> PushOutcome {
let mut state = self.lock();
state.reserved -= 1;
state.abandoned = true;
state.queue.push_back(ControlMessage::Abandon);
claim_drain(&mut state)
}
pub(crate) fn take_for_service(&self) -> Option<ControlMessage> {
let mut state = self.lock();
match state.queue.pop_front() {
Some(message) => Some(message),
None => {
state.draining = false;
None
}
}
}
fn release_reservation(&self) {
self.lock().reserved -= 1;
}
}
fn claim_drain(state: &mut RingState) -> PushOutcome {
if state.draining {
PushOutcome::DrainAlreadyScheduled
} else {
state.draining = true;
PushOutcome::RingDoorbell
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[must_use = "a queued message is not serviced until the doorbell is rung"]
pub(crate) enum PushOutcome {
RingDoorbell,
DrainAlreadyScheduled,
}
pub(crate) struct CancelSlot {
_private: (),
}
pub(crate) struct RetireSlot {
_private: (),
}
pub(crate) struct AbandonSlot {
_private: (),
}
pub(crate) fn release_cancel_slot(ring: &SubmissionRing, _slot: CancelSlot) {
ring.release_reservation();
}
pub(crate) fn release_retire_slot(ring: &SubmissionRing, _slot: RetireSlot) {
ring.release_reservation();
}
#[cfg(test)]
mod tests;