use anyhow::Result;
use parking_lot::Mutex;
use std::collections::btree_map::Entry;
use std::collections::{BTreeMap, BTreeSet};
use std::sync::Arc;
use std::sync::atomic::{AtomicU32, Ordering};
use crate::events::{EventAwaiter, EventHandle, EventManager, EventPoison, EventStatus};
#[derive(Copy, Clone, Debug, PartialEq, Eq, Hash)]
pub(crate) struct RemoteEventKey {
pub system_id: u64,
pub local_index: u32,
}
impl RemoteEventKey {
pub fn from_handle(handle: EventHandle) -> Self {
Self {
system_id: handle.system_id(),
local_index: handle.local_index(),
}
}
}
#[derive(Clone, Debug)]
pub(crate) enum CompletionKind {
Triggered,
Poisoned(Arc<EventPoison>),
}
pub(crate) struct CompletedEventInfo {
pub highest_generation: u32,
pub poisoned_generations: BTreeSet<u32>,
}
pub(crate) struct RemoteEvent {
proxy_manager: EventManager,
known_triggered: AtomicU32,
proxy_handles: Mutex<BTreeMap<u32, EventHandle>>,
completions: Mutex<BTreeMap<u32, Arc<CompletionKind>>>,
pending: Mutex<BTreeSet<u32>>,
}
pub(crate) enum WaitRegistration {
Ready,
Pending(EventAwaiter),
Poisoned(Arc<EventPoison>),
}
impl RemoteEvent {
pub fn new(proxy_manager: EventManager) -> Self {
Self {
proxy_manager,
known_triggered: AtomicU32::new(0),
proxy_handles: Mutex::new(BTreeMap::new()),
completions: Mutex::new(BTreeMap::new()),
pending: Mutex::new(BTreeSet::new()),
}
}
pub fn known_generation(&self) -> u32 {
self.known_triggered.load(Ordering::Acquire)
}
fn update_known_generation(&self, generation: u32) {
let _ = self
.known_triggered
.fetch_update(Ordering::AcqRel, Ordering::Acquire, |current| {
(generation > current).then_some(generation)
});
}
pub fn register_waiter(&self, generation: u32) -> Result<WaitRegistration> {
if generation <= self.known_generation() {
let completions = self.completions.lock();
return Ok(match completions.get(&generation) {
Some(completion) => match completion.as_ref() {
CompletionKind::Triggered => WaitRegistration::Ready,
CompletionKind::Poisoned(p) => WaitRegistration::Poisoned(p.clone()),
},
None => WaitRegistration::Ready,
});
}
let proxy_handle = {
let mut proxy_handles = self.proxy_handles.lock();
match proxy_handles.entry(generation) {
Entry::Occupied(e) => *e.get(),
Entry::Vacant(e) => {
let event = self.proxy_manager.new_event()?;
*e.insert(event.into_handle())
}
}
};
Ok(match self.proxy_manager.awaiter(proxy_handle) {
Ok(waiter) => WaitRegistration::Pending(waiter),
Err(_) => WaitRegistration::Ready,
})
}
pub fn add_pending(&self, generation: u32) -> bool {
self.pending.lock().insert(generation)
}
pub fn complete_generation(&self, generation: u32, completion: CompletionKind) {
{
let mut completions = self.completions.lock();
completions.insert(generation, Arc::new(completion.clone()));
const MAX_COMPLETION_HISTORY: u32 = 100;
let known_gen = self.known_generation();
if known_gen > MAX_COMPLETION_HISTORY {
completions.retain(|&g, _| g > known_gen - MAX_COMPLETION_HISTORY);
}
}
self.update_known_generation(generation);
let handles_to_complete = {
let mut proxy_handles = self.proxy_handles.lock();
let gens_to_wake: Vec<u32> = proxy_handles
.range(..=generation)
.map(|(g, _)| *g)
.collect();
let mut handles = Vec::new();
for generation in gens_to_wake {
if let Some(handle) = proxy_handles.remove(&generation) {
handles.push(handle);
}
}
handles
};
for handle in handles_to_complete {
match &completion {
CompletionKind::Triggered => {
let _ = self.proxy_manager.trigger(handle);
}
CompletionKind::Poisoned(p) => {
let _ = self.proxy_manager.poison(handle, p.reason());
}
}
}
let mut pending = self.pending.lock();
let to_remove: Vec<u32> = pending
.iter()
.copied()
.take_while(|g| *g <= generation)
.collect();
for g in to_remove {
pending.remove(&g);
}
}
pub fn status_for(&self, generation: u32) -> EventStatus {
if generation <= self.known_generation() {
let completions = self.completions.lock();
match completions.get(&generation) {
Some(completion) => match completion.as_ref() {
CompletionKind::Triggered => EventStatus::Ready,
CompletionKind::Poisoned(_) => EventStatus::Poisoned,
},
None => EventStatus::Ready,
}
} else {
EventStatus::Pending
}
}
pub fn poisoned_generations(&self) -> BTreeSet<u32> {
self.completions
.lock()
.iter()
.filter_map(|(generation, kind)| match kind.as_ref() {
CompletionKind::Poisoned(_) => Some(*generation),
_ => None,
})
.collect()
}
pub fn is_cacheable(&self) -> bool {
let pending = self.pending.lock();
let proxy_handles = self.proxy_handles.lock();
pending.is_empty() && proxy_handles.is_empty()
}
}