1use std::{
2 cell::{Cell, RefCell},
3 collections::VecDeque,
4 future::poll_fn,
5 rc::{Rc, Weak},
6 task::{Poll, Waker},
7 time::Duration,
8};
9
10use super::{ModuleLifecyclePhase, RuntimeFailure};
11
12#[derive(Clone, Copy, Debug, Eq, PartialEq)]
14#[repr(u8)]
15pub enum DiagnosticSource {
16 Lifecycle = 0,
18 Invocation = 1,
20 Admission = 2,
22 Supervision = 3,
24 Shutdown = 4,
26 RuntimeFailure = 5,
28}
29
30impl DiagnosticSource {
31 const COUNT: u8 = 6;
32
33 const fn bit(self) -> u8 {
34 1 << (self as u8)
35 }
36}
37
38#[derive(Clone, Copy, Debug, Eq, PartialEq)]
40pub struct DiagnosticFilter {
41 mask: u8,
42}
43
44impl DiagnosticFilter {
45 pub const fn none() -> Self {
47 Self { mask: 0 }
48 }
49
50 pub const fn all() -> Self {
52 Self {
53 mask: (1 << DiagnosticSource::COUNT) - 1,
54 }
55 }
56
57 pub const fn only(source: DiagnosticSource) -> Self {
59 Self { mask: source.bit() }
60 }
61
62 #[must_use]
64 pub const fn with_source(self, source: DiagnosticSource) -> Self {
65 Self {
66 mask: self.mask | source.bit(),
67 }
68 }
69
70 pub const fn includes(self, source: DiagnosticSource) -> bool {
72 self.mask & source.bit() != 0
73 }
74}
75
76impl Default for DiagnosticFilter {
77 fn default() -> Self {
78 Self::all()
79 }
80}
81
82#[derive(Clone, Copy, Debug, Eq, PartialEq)]
88pub enum RuntimeFailureKind {
89 Unavailable,
91 UnknownOperation,
93 AmbiguousBinding,
95 ProtocolViolation,
97 MissingModuleFactory,
99 UnavailableExecutionClass,
101 InvalidResolvedPlan,
103 AdmissionClosed,
105 ResourceExhausted,
107 DeadlineExceeded,
109 Cancelled,
111 Internal,
113 ModuleFailure,
115 ModuleRestartExhausted,
117}
118
119impl From<&RuntimeFailure> for RuntimeFailureKind {
120 fn from(error: &RuntimeFailure) -> Self {
121 match error {
122 RuntimeFailure::Unavailable { .. } => Self::Unavailable,
123 RuntimeFailure::UnknownOperation { .. } => Self::UnknownOperation,
124 RuntimeFailure::AmbiguousBinding { .. } => Self::AmbiguousBinding,
125 RuntimeFailure::ProtocolViolation { .. } => Self::ProtocolViolation,
126 RuntimeFailure::MissingModuleFactory { .. } => Self::MissingModuleFactory,
127 RuntimeFailure::UnavailableExecutionClass { .. } => Self::UnavailableExecutionClass,
128 RuntimeFailure::InvalidResolvedPlan { .. } => Self::InvalidResolvedPlan,
129 RuntimeFailure::AdmissionClosed => Self::AdmissionClosed,
130 RuntimeFailure::ResourceExhausted { .. } => Self::ResourceExhausted,
131 RuntimeFailure::DeadlineExceeded { .. } => Self::DeadlineExceeded,
132 RuntimeFailure::Cancelled { .. } => Self::Cancelled,
133 RuntimeFailure::Internal { .. } => Self::Internal,
134 RuntimeFailure::ModuleFailure { .. } => Self::ModuleFailure,
135 RuntimeFailure::ModuleRestartExhausted { .. } => Self::ModuleRestartExhausted,
136 }
137 }
138}
139
140#[derive(Clone, Copy, Debug, Eq, PartialEq)]
142pub enum DiagnosticOutcome {
143 Succeeded,
145 DomainError,
147 RuntimeFailure(RuntimeFailureKind),
149}
150
151#[derive(Clone, Copy, Debug, Eq, PartialEq)]
153pub enum DiagnosticAdmission {
154 Accepted,
156 Unavailable,
158 Exhausted,
160 Closed,
162}
163
164#[derive(Clone, Copy, Debug, Eq, PartialEq)]
166pub enum DiagnosticShutdownOutcome {
167 Clean,
169 RuntimeFailure,
171 Timeout,
173}
174
175#[derive(Clone, Debug, Eq, PartialEq)]
183pub enum DiagnosticEvent {
184 AppStarted { module_count: usize },
186 AppReady,
188 LifecycleStarted {
190 instance: String,
191 generation: u64,
192 phase: ModuleLifecyclePhase,
193 },
194 LifecycleCompleted {
196 instance: String,
197 generation: u64,
198 phase: ModuleLifecyclePhase,
199 outcome: DiagnosticOutcome,
200 elapsed: Duration,
201 },
202 InvocationStarted {
204 request_id: u64,
205 caller_instance: Option<String>,
206 provider_instance: Option<String>,
207 capability: &'static str,
208 operation: Option<&'static str>,
209 },
210 InvocationCompleted {
212 request_id: u64,
213 caller_instance: Option<String>,
214 provider_instance: Option<String>,
215 capability: &'static str,
216 operation: Option<&'static str>,
217 outcome: DiagnosticOutcome,
218 elapsed: Duration,
219 },
220 AdmissionRejected {
222 request_id: u64,
223 caller_instance: Option<String>,
224 provider_instance: Option<String>,
225 capability: &'static str,
226 operation: Option<&'static str>,
227 outcome: DiagnosticAdmission,
228 },
229 EventAdmission {
231 request_id: u64,
232 publisher_instance: String,
233 subscriber_instance: String,
234 capability: &'static str,
235 operation: Option<&'static str>,
236 outcome: DiagnosticAdmission,
237 },
238 GenerationUnavailable { instance: String, generation: u64 },
240 GenerationReady { instance: String, generation: u64 },
242 RestartScheduled {
244 instance: String,
245 attempt: usize,
246 delay: Duration,
247 },
248 RestartExhausted {
250 instance: String,
251 attempts: usize,
252 terminal: bool,
253 },
254 RuntimeFailure {
256 instance: Option<String>,
257 kind: RuntimeFailureKind,
258 },
259 ShutdownAdmissionClosed,
261 ShutdownCleanupStarted { timeout: Duration },
263 ShutdownCompleted {
265 outcome: DiagnosticShutdownOutcome,
266 elapsed: Duration,
267 },
268}
269
270#[derive(Clone, Debug, Eq, PartialEq)]
272pub struct DiagnosticRecord {
273 pub sequence: u64,
275 pub timestamp: Duration,
277 pub source: DiagnosticSource,
279 pub event: DiagnosticEvent,
281}
282
283#[derive(Clone, Copy, Debug, Eq, PartialEq)]
285pub enum DiagnosticSubscribeError {
286 ZeroCapacity,
288}
289
290#[derive(Debug, Default)]
291struct RuntimeDiagnosticsState {
292 observers: RefCell<Vec<Weak<DiagnosticObserverState>>>,
293 next_sequence: Cell<u64>,
294}
295
296impl Drop for RuntimeDiagnosticsState {
297 fn drop(&mut self) {
298 for observer in self
299 .observers
300 .get_mut()
301 .drain(..)
302 .filter_map(|observer| observer.upgrade())
303 {
304 observer.connected.set(false);
305 observer.wake_receiver();
306 }
307 }
308}
309
310#[derive(Clone, Debug)]
316pub struct RuntimeDiagnostics {
317 state: Rc<RuntimeDiagnosticsState>,
318}
319
320impl RuntimeDiagnostics {
321 pub fn new() -> Self {
323 Self {
324 state: Rc::new(RuntimeDiagnosticsState::default()),
325 }
326 }
327
328 pub fn subscribe(
330 &self,
331 filter: DiagnosticFilter,
332 capacity: usize,
333 ) -> Result<DiagnosticObserver, DiagnosticSubscribeError> {
334 if capacity == 0 {
335 return Err(DiagnosticSubscribeError::ZeroCapacity);
336 }
337 let observer = Rc::new(DiagnosticObserverState {
338 filter,
339 capacity,
340 queue: RefCell::new(VecDeque::with_capacity(capacity)),
341 dropped: Cell::new(0),
342 connected: Cell::new(true),
343 receiver_waker: RefCell::new(None),
344 });
345 let mut observers = self.state.observers.borrow_mut();
346 observers.retain(|observer| observer.upgrade().is_some());
347 observers.push(Rc::downgrade(&observer));
348 Ok(DiagnosticObserver { state: observer })
349 }
350
351 pub fn subscribe_all(
353 &self,
354 capacity: usize,
355 ) -> Result<DiagnosticObserver, DiagnosticSubscribeError> {
356 self.subscribe(DiagnosticFilter::all(), capacity)
357 }
358
359 pub fn observer_count(&self) -> usize {
361 let mut observers = self.state.observers.borrow_mut();
362 observers.retain(|observer| observer.upgrade().is_some());
363 observers.len()
364 }
365
366 pub(crate) fn emit<F>(&self, source: DiagnosticSource, timestamp: Duration, build: F)
367 where
368 F: FnOnce(u64) -> DiagnosticEvent,
369 {
370 let interested = self
371 .state
372 .observers
373 .borrow()
374 .iter()
375 .filter_map(Weak::upgrade)
376 .any(|observer| observer.filter.includes(source));
377 if !interested {
378 return;
379 }
380
381 let sequence = self.state.next_sequence.get();
382 self.state.next_sequence.set(sequence.saturating_add(1));
383 let record = DiagnosticRecord {
384 sequence,
385 timestamp,
386 source,
387 event: build(sequence),
388 };
389 self.state.observers.borrow_mut().retain(|observer| {
390 let Some(observer) = observer.upgrade() else {
391 return false;
392 };
393 if observer.filter.includes(source) {
394 observer.enqueue(record.clone());
395 }
396 true
397 });
398 }
399
400 pub(crate) fn emit_runtime_failure(
401 &self,
402 timestamp: Duration,
403 instance: Option<&str>,
404 error: &RuntimeFailure,
405 ) {
406 let kind = RuntimeFailureKind::from(error);
407 self.emit(DiagnosticSource::RuntimeFailure, timestamp, |_| {
408 DiagnosticEvent::RuntimeFailure {
409 instance: instance.map(str::to_owned),
410 kind,
411 }
412 });
413 }
414}
415
416impl Default for RuntimeDiagnostics {
417 fn default() -> Self {
418 Self::new()
419 }
420}
421
422#[derive(Debug)]
423struct DiagnosticObserverState {
424 filter: DiagnosticFilter,
425 capacity: usize,
426 queue: RefCell<VecDeque<DiagnosticRecord>>,
427 dropped: Cell<u64>,
428 connected: Cell<bool>,
429 receiver_waker: RefCell<Option<Waker>>,
430}
431
432impl DiagnosticObserverState {
433 fn enqueue(&self, record: DiagnosticRecord) {
434 let mut queue = self.queue.borrow_mut();
435 if queue.len() >= self.capacity {
436 self.dropped.set(self.dropped.get().saturating_add(1));
437 return;
438 }
439 queue.push_back(record);
440 drop(queue);
441 self.wake_receiver();
442 }
443
444 fn wake_receiver(&self) {
445 if let Some(waker) = self.receiver_waker.borrow_mut().take() {
446 waker.wake();
447 }
448 }
449}
450
451#[derive(Debug)]
453pub struct DiagnosticObserver {
454 state: Rc<DiagnosticObserverState>,
455}
456
457impl DiagnosticObserver {
458 pub async fn recv(&mut self) -> Option<DiagnosticRecord> {
463 poll_fn(|context| {
464 if let Some(record) = self.try_recv() {
465 return Poll::Ready(Some(record));
466 }
467 if !self.state.connected.get() {
468 return Poll::Ready(None);
469 }
470 self.state
471 .receiver_waker
472 .replace(Some(context.waker().clone()));
473 if let Some(record) = self.try_recv() {
474 self.state.receiver_waker.borrow_mut().take();
475 return Poll::Ready(Some(record));
476 }
477 Poll::Pending
478 })
479 .await
480 }
481
482 pub fn try_recv(&self) -> Option<DiagnosticRecord> {
484 self.state.queue.borrow_mut().pop_front()
485 }
486
487 pub fn try_next(&self) -> Option<DiagnosticRecord> {
489 self.try_recv()
490 }
491
492 pub fn dropped_count(&self) -> u64 {
494 self.state.dropped.get()
495 }
496
497 pub fn pending_count(&self) -> usize {
499 self.state.queue.borrow().len()
500 }
501
502 pub fn capacity(&self) -> usize {
504 self.state.capacity
505 }
506
507 pub fn filter(&self) -> DiagnosticFilter {
509 self.state.filter
510 }
511}
512
513pub(crate) fn diagnostic_operation(
514 operations: &'static [&'static str],
515 operation: &str,
516) -> Option<&'static str> {
517 operations
518 .iter()
519 .copied()
520 .find(|candidate| *candidate == operation)
521}
522
523#[cfg(test)]
524mod tests {
525 use super::{DiagnosticEvent, DiagnosticFilter, DiagnosticSource, RuntimeDiagnostics};
526 use std::time::Duration;
527
528 #[test]
529 fn does_not_build_a_record_without_an_interested_observer() {
530 let diagnostics = RuntimeDiagnostics::new();
531 let built = std::cell::Cell::new(false);
532
533 diagnostics.emit(DiagnosticSource::Lifecycle, Duration::ZERO, |_| {
534 built.set(true);
535 DiagnosticEvent::AppReady
536 });
537
538 assert!(!built.get());
539 }
540
541 #[test]
542 fn filters_sources_before_building_a_record() {
543 let diagnostics = RuntimeDiagnostics::new();
544 let observer = diagnostics
545 .subscribe(DiagnosticFilter::only(DiagnosticSource::Invocation), 1)
546 .expect("observer capacity is positive");
547 let built = std::cell::Cell::new(false);
548
549 diagnostics.emit(DiagnosticSource::Lifecycle, Duration::ZERO, |_| {
550 built.set(true);
551 DiagnosticEvent::AppReady
552 });
553
554 assert!(!built.get());
555 assert!(observer.try_recv().is_none());
556 }
557}