use std::collections::VecDeque;
use std::num::NonZeroUsize;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Condvar, Mutex, Weak};
use std::time::Duration;
use serde::{Deserialize, Serialize};
use thiserror::Error;
use crate::capability::CapabilityDenied;
use crate::component::ComponentId;
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub enum LifecycleState {
Registered,
Starting,
Running,
Stopping,
Stopped,
Failed,
Removed,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct LifecycleEvent {
pub sequence: u64,
pub component_id: ComponentId,
pub kind: LifecycleEventKind,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub enum LifecycleEventKind {
Transition {
from: Option<LifecycleState>,
to: LifecycleState,
},
CapabilityDenied(CapabilityDenied),
}
pub struct LifecycleSubscription {
queue: Arc<SubscriberQueue>,
}
impl LifecycleSubscription {
pub fn recv_timeout(&self, timeout: Duration) -> Result<LifecycleEvent, EventReceiveError> {
let guard = self
.queue
.state
.lock()
.map_err(|_| EventReceiveError::Poisoned)?;
let (mut guard, _wait) = self
.queue
.ready
.wait_timeout_while(guard, timeout, |state| {
state.events.is_empty() && !state.closed
})
.map_err(|_| EventReceiveError::Poisoned)?;
if let Some(event) = guard.events.pop_front() {
return Ok(event);
}
if guard.closed {
Err(EventReceiveError::Closed)
} else {
Err(EventReceiveError::Timeout)
}
}
pub fn try_recv(&self) -> Result<LifecycleEvent, EventTryReceiveError> {
let mut state = self
.queue
.state
.lock()
.map_err(|_| EventTryReceiveError::Poisoned)?;
if let Some(event) = state.events.pop_front() {
Ok(event)
} else if state.closed {
Err(EventTryReceiveError::Closed)
} else {
Err(EventTryReceiveError::Empty)
}
}
#[must_use]
pub fn lagged_events(&self) -> usize {
self.queue.lagged.load(Ordering::Acquire)
}
}
impl Drop for LifecycleSubscription {
fn drop(&mut self) {
if let Ok(mut state) = self.queue.state.lock() {
state.closed = true;
self.queue.ready.notify_all();
}
}
}
#[derive(Clone, Copy, Debug, Error, Eq, PartialEq)]
pub enum EventReceiveError {
#[error("lifecycle event receive timed out")]
Timeout,
#[error("lifecycle event subscription is closed")]
Closed,
#[error("lifecycle event subscription synchronization is poisoned")]
Poisoned,
}
#[derive(Clone, Copy, Debug, Error, Eq, PartialEq)]
pub enum EventTryReceiveError {
#[error("lifecycle event subscription is empty")]
Empty,
#[error("lifecycle event subscription is closed")]
Closed,
#[error("lifecycle event subscription synchronization is poisoned")]
Poisoned,
}
#[derive(Default)]
pub(crate) struct EventHub {
next_sequence: AtomicU64,
subscribers: Mutex<Vec<Weak<SubscriberQueue>>>,
}
impl EventHub {
pub(crate) fn subscribe(
&self,
capacity: NonZeroUsize,
) -> Result<LifecycleSubscription, EventPublishError> {
let queue = Arc::new(SubscriberQueue {
capacity: capacity.get(),
state: Mutex::new(QueueState::default()),
ready: Condvar::new(),
lagged: AtomicUsize::new(0),
});
self.subscribers
.lock()
.map_err(|_| EventPublishError::Poisoned)?
.push(Arc::downgrade(&queue));
Ok(LifecycleSubscription { queue })
}
pub(crate) fn publish_transition(
&self,
component_id: ComponentId,
from: Option<LifecycleState>,
to: LifecycleState,
) -> Result<(), EventPublishError> {
self.publish(component_id, LifecycleEventKind::Transition { from, to })
}
pub(crate) fn publish_denial(&self, denial: CapabilityDenied) -> Result<(), EventPublishError> {
self.publish(
denial.component_id,
LifecycleEventKind::CapabilityDenied(denial),
)
}
fn publish(
&self,
component_id: ComponentId,
kind: LifecycleEventKind,
) -> Result<(), EventPublishError> {
let mut subscribers = self
.subscribers
.lock()
.map_err(|_| EventPublishError::Poisoned)?;
subscribers.retain(|subscriber| subscriber.strong_count() > 0);
if subscribers.is_empty() {
return Ok(());
}
let event = LifecycleEvent {
sequence: self.next_sequence.fetch_add(1, Ordering::AcqRel),
component_id,
kind,
};
for subscriber in subscribers.iter().filter_map(Weak::upgrade) {
let mut state = subscriber
.state
.lock()
.map_err(|_| EventPublishError::Poisoned)?;
if state.events.len() == subscriber.capacity {
let _discarded = state.events.pop_front();
subscriber.lagged.fetch_add(1, Ordering::AcqRel);
}
state.events.push_back(event.clone());
subscriber.ready.notify_one();
}
Ok(())
}
}
#[derive(Debug, Error)]
pub(crate) enum EventPublishError {
#[error("lifecycle event stream synchronization is poisoned")]
Poisoned,
}
struct SubscriberQueue {
capacity: usize,
state: Mutex<QueueState>,
ready: Condvar,
lagged: AtomicUsize,
}
#[derive(Default)]
struct QueueState {
events: VecDeque<LifecycleEvent>,
closed: bool,
}