frame-core 0.2.0

Component model, lifecycle, process isolation, and WASM module host
Documentation
//! Typed, bounded lifecycle event delivery.

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;

/// An externally observable component lifecycle state.
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub enum LifecycleState {
    /// Definition is registered but has not started.
    Registered,
    /// Module loading and child liveness checks are in progress.
    Starting,
    /// Every declared child completed a mailbox round-trip.
    Running,
    /// Ordered child and supervisor drain is in progress.
    Stopping,
    /// All component processes have normal tombstones.
    Stopped,
    /// Start or supervision failed; the associated status carries the reason.
    Failed,
    /// Modules and registration were removed without residue.
    Removed,
}

/// One ordered item on the component lifecycle and security-news stream.
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct LifecycleEvent {
    /// Registry-wide monotonic sequence number.
    pub sequence: u64,
    /// Component associated with this news item.
    pub component_id: ComponentId,
    /// Typed transition or capability denial payload.
    pub kind: LifecycleEventKind,
}

/// The typed payload carried by one lifecycle stream item.
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub enum LifecycleEventKind {
    /// One component lifecycle transition.
    Transition {
        /// State before this transition, absent for initial registration.
        from: Option<LifecycleState>,
        /// State after this transition.
        to: LifecycleState,
    },
    /// One consuming act was denied by the fresh capability check.
    CapabilityDenied(CapabilityDenied),
}

/// A bounded subscription that reports every event it had to discard.
pub struct LifecycleSubscription {
    queue: Arc<SubscriberQueue>,
}

impl LifecycleSubscription {
    /// Waits up to `timeout` for the next retained event.
    ///
    /// # Errors
    ///
    /// Returns a typed timeout, closure, or synchronization failure.
    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)
        }
    }

    /// Returns the next retained event without waiting.
    ///
    /// # Errors
    ///
    /// Returns closure, emptiness, or synchronization failures distinctly.
    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)
        }
    }

    /// Returns the cumulative number of events dropped from this subscriber.
    #[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();
        }
    }
}

/// Failure while waiting for an event.
#[derive(Clone, Copy, Debug, Error, Eq, PartialEq)]
pub enum EventReceiveError {
    /// No event arrived before the caller's deadline.
    #[error("lifecycle event receive timed out")]
    Timeout,
    /// The subscription was closed.
    #[error("lifecycle event subscription is closed")]
    Closed,
    /// Subscriber synchronization was poisoned by a panic.
    #[error("lifecycle event subscription synchronization is poisoned")]
    Poisoned,
}

/// Failure while polling for an event.
#[derive(Clone, Copy, Debug, Error, Eq, PartialEq)]
pub enum EventTryReceiveError {
    /// No retained event is currently available.
    #[error("lifecycle event subscription is empty")]
    Empty,
    /// The subscription was closed.
    #[error("lifecycle event subscription is closed")]
    Closed,
    /// Subscriber synchronization was poisoned by a panic.
    #[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,
}