mod prepared;
#[cfg(test)]
mod tests;
use super::model::{CompletionDelivery, DeliveryScope};
use crate::editor::events::{AppEvent, EventSender};
use parking_lot::{Condvar, Mutex};
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::Arc;
use strop_core::worker::{Outcome, WorkerId};
const PROVIDER_SLOTS: usize = 5;
const MAX_RETAINED_DELIVERIES: usize = 16;
const MAX_RETAINED_BYTES: usize = 24 * 1024 * 1024;
struct Retained {
delivery: Arc<CompletionDelivery>,
bytes: usize,
}
struct Slot {
retained: Retained,
observed: bool,
}
struct State {
slots: [Option<Slot>; PROVIDER_SLOTS],
retired: Vec<Retained>,
prepared: Option<prepared::Ready>,
bytes: usize,
high_water_bytes: usize,
destroyed: u64,
capacity_changed: bool,
stopped: Option<Outcome<()>>,
wake: Option<EventSender>,
}
impl Default for State {
fn default() -> Self {
Self {
slots: std::array::from_fn(|_| None),
retired: Vec::with_capacity(MAX_RETAINED_DELIVERIES),
prepared: None,
bytes: 0,
high_water_bytes: 0,
destroyed: 0,
capacity_changed: false,
stopped: None,
wake: None,
}
}
}
#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
pub(super) struct CompletionDrain {
pub deliveries: Vec<Arc<CompletionDelivery>>,
pub capacity_changed: bool,
pub stopped: Option<Outcome<()>>,
pub usage: PublicationUsage,
}
#[derive(Debug, Clone, Copy, Default, serde::Serialize, serde::Deserialize)]
pub(super) struct PublicationUsage {
pub retained_deliveries: usize,
pub charged_bytes: usize,
pub high_water_bytes: usize,
pub destroyed: u64,
}
#[derive(Default)]
pub(super) struct CompletionMailbox {
query: AtomicU64,
selection: AtomicU64,
preparation: AtomicU64,
closed: AtomicBool,
state: Mutex<State>,
available: Condvar,
}
impl CompletionMailbox {
pub fn connect(&self, sender: EventSender) {
self.state.lock().wake = Some(sender);
}
pub fn query(&self, request: Option<WorkerId>) {
self.query
.store(request.map_or(0, WorkerId::get), Ordering::Release);
self.selection.store(0, Ordering::Release);
self.preparation.store(0, Ordering::Release);
self.released();
}
pub fn select(&self, selection: Option<WorkerId>) {
self.selection
.store(selection.map_or(0, WorkerId::get), Ordering::Release);
self.preparation.store(0, Ordering::Release);
self.released();
}
pub fn prepare(&self, work: WorkerId) {
self.preparation.store(work.get(), Ordering::Release);
self.released();
}
pub fn owns(&self, query: WorkerId, scope: DeliveryScope) -> bool {
!self.closed.load(Ordering::Acquire)
&& self.query.load(Ordering::Acquire) == query.get()
&& match scope {
DeliveryScope::Query => true,
DeliveryScope::Selection(selection) => {
self.selection.load(Ordering::Acquire) == selection.get()
}
DeliveryScope::Preparation { selection, work } => {
self.selection.load(Ordering::Acquire) == selection.get()
&& self.preparation.load(Ordering::Acquire) == work.get()
}
}
}
pub fn publish(&self, delivery: CompletionDelivery) -> bool {
if !self.owns(delivery.query.request, delivery.scope) {
return false;
}
let bytes = delivery.payload.charged_bytes();
debug_assert!(bytes <= MAX_RETAINED_BYTES);
let slot = delivery.payload.slot();
let retained = Retained {
delivery: Arc::new(delivery),
bytes,
};
loop {
self.retire_native();
let mut state = self.state.lock();
if !self.owns(retained.delivery.query.request, retained.delivery.scope) {
return false;
}
if let Some(previous) = state.slots[slot].take() {
state.retired.push(previous.retained);
drop(state);
continue;
}
let count = state.retired.len() + state.slots.iter().flatten().count();
if count < MAX_RETAINED_DELIVERIES && state.bytes + bytes <= MAX_RETAINED_BYTES {
state.bytes += bytes;
state.high_water_bytes = state.high_water_bytes.max(state.bytes);
state.slots[slot] = Some(Slot {
retained,
observed: false,
});
let wake = state.wake.clone();
drop(state);
wake_owner(wake);
return true;
}
if state
.retired
.iter()
.any(|entry| Arc::strong_count(&entry.delivery) == 1)
{
drop(state);
continue;
}
self.available.wait(&mut state);
}
}
pub fn capacity(&self) {
let mut state = self.state.lock();
state.capacity_changed = true;
let wake = state.wake.clone();
drop(state);
wake_owner(wake);
}
pub fn take(&self) -> CompletionDrain {
let mut state = self.state.lock();
let mut deliveries = Vec::new();
for slot in state.slots.iter_mut().flatten() {
if !slot.observed {
slot.observed = true;
deliveries.push(slot.retained.delivery.clone());
}
}
CompletionDrain {
deliveries,
capacity_changed: std::mem::take(&mut state.capacity_changed),
stopped: state.stopped.take(),
usage: PublicationUsage {
retained_deliveries: state.retired.len()
+ state.slots.iter().flatten().count()
+ usize::from(state.prepared.is_some()),
charged_bytes: state.bytes,
high_water_bytes: state.high_water_bytes,
destroyed: state.destroyed,
},
}
}
pub fn released(&self) {
let _state = self.state.lock();
self.available.notify_all();
}
pub fn retire_native(&self) {
loop {
let mut state = self.state.lock();
if state
.prepared
.as_ref()
.is_some_and(|ready| !self.owns(ready.query, ready.scope))
{
let prepared = state.prepared.take();
state.bytes -= prepared::PLAN_BYTES;
drop(state);
drop(prepared);
self.available.notify_all();
continue;
}
for index in 0..PROVIDER_SLOTS {
let stale = state.slots[index].as_ref().is_some_and(|slot| {
!self.owns(
slot.retained.delivery.query.request,
slot.retained.delivery.scope,
)
});
if stale {
if let Some(slot) = state.slots[index].take() {
state.retired.push(slot.retained);
}
}
}
let obsolete = state
.retired
.iter()
.position(|entry| Arc::strong_count(&entry.delivery) == 1)
.map(|index| state.retired.swap_remove(index));
let Some(obsolete) = obsolete else { return };
state.bytes -= obsolete.bytes;
state.destroyed = state.destroyed.saturating_add(1);
drop(state);
drop(obsolete);
self.available.notify_all();
}
}
pub fn finish_native(&self, outcome: Outcome<()>) {
self.closed.store(true, Ordering::Release);
let mut state = self.state.lock();
state.stopped = Some(outcome);
self.available.notify_all();
let wake = state.wake.clone();
drop(state);
wake_owner(wake);
loop {
self.retire_native();
let mut state = self.state.lock();
if state.bytes == 0 {
break;
}
if state
.retired
.iter()
.any(|entry| Arc::strong_count(&entry.delivery) == 1)
{
drop(state);
continue;
}
self.available.wait(&mut state);
}
}
}
fn wake_owner(sender: Option<EventSender>) {
if let Some(sender) = sender {
let _ = sender.send(AppEvent::Completion);
}
}