use std::{
cell::{Cell, RefCell},
collections::VecDeque,
future::poll_fn,
rc::{Rc, Weak},
task::{Poll, Waker},
time::Duration,
};
use super::{ModuleLifecyclePhase, RuntimeFailure};
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[repr(u8)]
pub enum DiagnosticSource {
Lifecycle = 0,
Invocation = 1,
Admission = 2,
Supervision = 3,
Shutdown = 4,
RuntimeFailure = 5,
}
impl DiagnosticSource {
const COUNT: u8 = 6;
const fn bit(self) -> u8 {
1 << (self as u8)
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct DiagnosticFilter {
mask: u8,
}
impl DiagnosticFilter {
pub const fn none() -> Self {
Self { mask: 0 }
}
pub const fn all() -> Self {
Self {
mask: (1 << DiagnosticSource::COUNT) - 1,
}
}
pub const fn only(source: DiagnosticSource) -> Self {
Self { mask: source.bit() }
}
#[must_use]
pub const fn with_source(self, source: DiagnosticSource) -> Self {
Self {
mask: self.mask | source.bit(),
}
}
pub const fn includes(self, source: DiagnosticSource) -> bool {
self.mask & source.bit() != 0
}
}
impl Default for DiagnosticFilter {
fn default() -> Self {
Self::all()
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum RuntimeFailureKind {
Unavailable,
UnknownOperation,
AmbiguousBinding,
ProtocolViolation,
MissingModuleFactory,
UnavailableExecutionClass,
InvalidResolvedPlan,
AdmissionClosed,
ResourceExhausted,
DeadlineExceeded,
Cancelled,
Internal,
ModuleFailure,
ModuleRestartExhausted,
}
impl From<&RuntimeFailure> for RuntimeFailureKind {
fn from(error: &RuntimeFailure) -> Self {
match error {
RuntimeFailure::Unavailable { .. } => Self::Unavailable,
RuntimeFailure::UnknownOperation { .. } => Self::UnknownOperation,
RuntimeFailure::AmbiguousBinding { .. } => Self::AmbiguousBinding,
RuntimeFailure::ProtocolViolation { .. } => Self::ProtocolViolation,
RuntimeFailure::MissingModuleFactory { .. } => Self::MissingModuleFactory,
RuntimeFailure::UnavailableExecutionClass { .. } => Self::UnavailableExecutionClass,
RuntimeFailure::InvalidResolvedPlan { .. } => Self::InvalidResolvedPlan,
RuntimeFailure::AdmissionClosed => Self::AdmissionClosed,
RuntimeFailure::ResourceExhausted { .. } => Self::ResourceExhausted,
RuntimeFailure::DeadlineExceeded { .. } => Self::DeadlineExceeded,
RuntimeFailure::Cancelled { .. } => Self::Cancelled,
RuntimeFailure::Internal { .. } => Self::Internal,
RuntimeFailure::ModuleFailure { .. } => Self::ModuleFailure,
RuntimeFailure::ModuleRestartExhausted { .. } => Self::ModuleRestartExhausted,
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum DiagnosticOutcome {
Succeeded,
DomainError,
RuntimeFailure(RuntimeFailureKind),
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum DiagnosticAdmission {
Accepted,
Unavailable,
Exhausted,
Closed,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum DiagnosticShutdownOutcome {
Clean,
RuntimeFailure,
Timeout,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum DiagnosticEvent {
AppStarted { module_count: usize },
AppReady,
LifecycleStarted {
instance: String,
generation: u64,
phase: ModuleLifecyclePhase,
},
LifecycleCompleted {
instance: String,
generation: u64,
phase: ModuleLifecyclePhase,
outcome: DiagnosticOutcome,
elapsed: Duration,
},
InvocationStarted {
request_id: u64,
caller_instance: Option<String>,
provider_instance: Option<String>,
capability: &'static str,
operation: Option<&'static str>,
},
InvocationCompleted {
request_id: u64,
caller_instance: Option<String>,
provider_instance: Option<String>,
capability: &'static str,
operation: Option<&'static str>,
outcome: DiagnosticOutcome,
elapsed: Duration,
},
AdmissionRejected {
request_id: u64,
caller_instance: Option<String>,
provider_instance: Option<String>,
capability: &'static str,
operation: Option<&'static str>,
outcome: DiagnosticAdmission,
},
EventAdmission {
request_id: u64,
publisher_instance: String,
subscriber_instance: String,
capability: &'static str,
operation: Option<&'static str>,
outcome: DiagnosticAdmission,
},
GenerationUnavailable { instance: String, generation: u64 },
GenerationReady { instance: String, generation: u64 },
RestartScheduled {
instance: String,
attempt: usize,
delay: Duration,
},
RestartExhausted {
instance: String,
attempts: usize,
terminal: bool,
},
RuntimeFailure {
instance: Option<String>,
kind: RuntimeFailureKind,
},
ShutdownAdmissionClosed,
ShutdownCleanupStarted { timeout: Duration },
ShutdownCompleted {
outcome: DiagnosticShutdownOutcome,
elapsed: Duration,
},
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct DiagnosticRecord {
pub sequence: u64,
pub timestamp: Duration,
pub source: DiagnosticSource,
pub event: DiagnosticEvent,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum DiagnosticSubscribeError {
ZeroCapacity,
}
#[derive(Debug, Default)]
struct RuntimeDiagnosticsState {
observers: RefCell<Vec<Weak<DiagnosticObserverState>>>,
next_sequence: Cell<u64>,
}
impl Drop for RuntimeDiagnosticsState {
fn drop(&mut self) {
for observer in self
.observers
.get_mut()
.drain(..)
.filter_map(|observer| observer.upgrade())
{
observer.connected.set(false);
observer.wake_receiver();
}
}
}
#[derive(Clone, Debug)]
pub struct RuntimeDiagnostics {
state: Rc<RuntimeDiagnosticsState>,
}
impl RuntimeDiagnostics {
pub fn new() -> Self {
Self {
state: Rc::new(RuntimeDiagnosticsState::default()),
}
}
pub fn subscribe(
&self,
filter: DiagnosticFilter,
capacity: usize,
) -> Result<DiagnosticObserver, DiagnosticSubscribeError> {
if capacity == 0 {
return Err(DiagnosticSubscribeError::ZeroCapacity);
}
let observer = Rc::new(DiagnosticObserverState {
filter,
capacity,
queue: RefCell::new(VecDeque::with_capacity(capacity)),
dropped: Cell::new(0),
connected: Cell::new(true),
receiver_waker: RefCell::new(None),
});
let mut observers = self.state.observers.borrow_mut();
observers.retain(|observer| observer.upgrade().is_some());
observers.push(Rc::downgrade(&observer));
Ok(DiagnosticObserver { state: observer })
}
pub fn subscribe_all(
&self,
capacity: usize,
) -> Result<DiagnosticObserver, DiagnosticSubscribeError> {
self.subscribe(DiagnosticFilter::all(), capacity)
}
pub fn observer_count(&self) -> usize {
let mut observers = self.state.observers.borrow_mut();
observers.retain(|observer| observer.upgrade().is_some());
observers.len()
}
pub(crate) fn emit<F>(&self, source: DiagnosticSource, timestamp: Duration, build: F)
where
F: FnOnce(u64) -> DiagnosticEvent,
{
let interested = self
.state
.observers
.borrow()
.iter()
.filter_map(Weak::upgrade)
.any(|observer| observer.filter.includes(source));
if !interested {
return;
}
let sequence = self.state.next_sequence.get();
self.state.next_sequence.set(sequence.saturating_add(1));
let record = DiagnosticRecord {
sequence,
timestamp,
source,
event: build(sequence),
};
self.state.observers.borrow_mut().retain(|observer| {
let Some(observer) = observer.upgrade() else {
return false;
};
if observer.filter.includes(source) {
observer.enqueue(record.clone());
}
true
});
}
pub(crate) fn emit_runtime_failure(
&self,
timestamp: Duration,
instance: Option<&str>,
error: &RuntimeFailure,
) {
let kind = RuntimeFailureKind::from(error);
self.emit(DiagnosticSource::RuntimeFailure, timestamp, |_| {
DiagnosticEvent::RuntimeFailure {
instance: instance.map(str::to_owned),
kind,
}
});
}
}
impl Default for RuntimeDiagnostics {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug)]
struct DiagnosticObserverState {
filter: DiagnosticFilter,
capacity: usize,
queue: RefCell<VecDeque<DiagnosticRecord>>,
dropped: Cell<u64>,
connected: Cell<bool>,
receiver_waker: RefCell<Option<Waker>>,
}
impl DiagnosticObserverState {
fn enqueue(&self, record: DiagnosticRecord) {
let mut queue = self.queue.borrow_mut();
if queue.len() >= self.capacity {
self.dropped.set(self.dropped.get().saturating_add(1));
return;
}
queue.push_back(record);
drop(queue);
self.wake_receiver();
}
fn wake_receiver(&self) {
if let Some(waker) = self.receiver_waker.borrow_mut().take() {
waker.wake();
}
}
}
#[derive(Debug)]
pub struct DiagnosticObserver {
state: Rc<DiagnosticObserverState>,
}
impl DiagnosticObserver {
pub async fn recv(&mut self) -> Option<DiagnosticRecord> {
poll_fn(|context| {
if let Some(record) = self.try_recv() {
return Poll::Ready(Some(record));
}
if !self.state.connected.get() {
return Poll::Ready(None);
}
self.state
.receiver_waker
.replace(Some(context.waker().clone()));
if let Some(record) = self.try_recv() {
self.state.receiver_waker.borrow_mut().take();
return Poll::Ready(Some(record));
}
Poll::Pending
})
.await
}
pub fn try_recv(&self) -> Option<DiagnosticRecord> {
self.state.queue.borrow_mut().pop_front()
}
pub fn try_next(&self) -> Option<DiagnosticRecord> {
self.try_recv()
}
pub fn dropped_count(&self) -> u64 {
self.state.dropped.get()
}
pub fn pending_count(&self) -> usize {
self.state.queue.borrow().len()
}
pub fn capacity(&self) -> usize {
self.state.capacity
}
pub fn filter(&self) -> DiagnosticFilter {
self.state.filter
}
}
pub(crate) fn diagnostic_operation(
operations: &'static [&'static str],
operation: &str,
) -> Option<&'static str> {
operations
.iter()
.copied()
.find(|candidate| *candidate == operation)
}
#[cfg(test)]
mod tests {
use super::{DiagnosticEvent, DiagnosticFilter, DiagnosticSource, RuntimeDiagnostics};
use std::time::Duration;
#[test]
fn does_not_build_a_record_without_an_interested_observer() {
let diagnostics = RuntimeDiagnostics::new();
let built = std::cell::Cell::new(false);
diagnostics.emit(DiagnosticSource::Lifecycle, Duration::ZERO, |_| {
built.set(true);
DiagnosticEvent::AppReady
});
assert!(!built.get());
}
#[test]
fn filters_sources_before_building_a_record() {
let diagnostics = RuntimeDiagnostics::new();
let observer = diagnostics
.subscribe(DiagnosticFilter::only(DiagnosticSource::Invocation), 1)
.expect("observer capacity is positive");
let built = std::cell::Cell::new(false);
diagnostics.emit(DiagnosticSource::Lifecycle, Duration::ZERO, |_| {
built.set(true);
DiagnosticEvent::AppReady
});
assert!(!built.get());
assert!(observer.try_recv().is_none());
}
}