Skip to main content

ahri_tre_observability/
lifecycle.rs

1use crate::{ConfigurationProvenance, FailureCategory, Outcome, TARGET};
2use serde::{Deserialize, Serialize};
3use std::time::Instant;
4use uuid::Uuid;
5
6/// Existing process and recovery boundaries; no caller-controlled phase names.
7#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
8#[serde(rename_all = "snake_case")]
9pub enum Phase {
10    Configuration,
11    InjectedSecrets,
12    ManagedSecrets,
13    Dependencies,
14    Recovery,
15    RecoveryCleanup,
16    Listener,
17    Serving,
18    Shutdown,
19    SessionDrain,
20    Maintenance,
21}
22
23/// Observes existing work without changing its result, cancellation or authority.
24/// There is deliberately no Drop completion: a lost process cannot claim an outcome.
25pub struct LifecycleSpan {
26    record: Lifecycle,
27    started: Instant,
28}
29
30impl LifecycleSpan {
31    pub fn start(phase: Phase, provenance: Option<ConfigurationProvenance>) -> Self {
32        Self::start_operation(phase, provenance, None)
33    }
34
35    /// Recovery joins only the existing Operation, never an invented request ID.
36    pub fn start_operation(
37        phase: Phase,
38        provenance: Option<ConfigurationProvenance>,
39        operation_id: Option<Uuid>,
40    ) -> Self {
41        let span = Self {
42            record: Lifecycle {
43                phase,
44                span_id: Uuid::new_v4(),
45                operation_id,
46                outcome: None,
47                failure_category: None,
48                duration_ms: None,
49                provenance,
50            },
51            started: Instant::now(),
52        };
53        span.emit();
54        span
55    }
56
57    pub fn run<T, E>(
58        phase: Phase,
59        provenance: Option<ConfigurationProvenance>,
60        category: FailureCategory,
61        work: impl FnOnce() -> Result<T, E>,
62    ) -> Result<T, E> {
63        let span = Self::start(phase, provenance);
64        let result = work();
65        span.complete(&result, category);
66        result
67    }
68
69    pub fn complete<T, E>(self, result: &Result<T, E>, category: FailureCategory) {
70        if result.is_ok() {
71            self.finish(Outcome::Success, None);
72        } else {
73            self.finish(Outcome::Unavailable, Some(category));
74        }
75    }
76
77    pub fn finish(mut self, outcome: Outcome, category: Option<FailureCategory>) {
78        self.record.outcome = Some(outcome);
79        self.record.failure_category = category;
80        self.record.duration_ms =
81            Some(self.started.elapsed().as_millis().min(u64::MAX as u128) as u64);
82        self.emit();
83    }
84
85    fn emit(&self) {
86        if let Ok(record) = serde_json::to_string(&self.record) {
87            tracing::info!(target: TARGET, lifecycle = record.as_str());
88        }
89    }
90}
91
92#[derive(Serialize, Deserialize)]
93#[serde(deny_unknown_fields)]
94pub(crate) struct Lifecycle {
95    pub phase: Phase,
96    span_id: Uuid,
97    #[serde(skip_serializing_if = "Option::is_none")]
98    operation_id: Option<Uuid>,
99    #[serde(skip_serializing_if = "Option::is_none")]
100    pub outcome: Option<Outcome>,
101    #[serde(skip_serializing_if = "Option::is_none")]
102    failure_category: Option<FailureCategory>,
103    #[serde(skip_serializing_if = "Option::is_none")]
104    duration_ms: Option<u64>,
105    #[serde(skip_serializing_if = "Option::is_none")]
106    pub provenance: Option<ConfigurationProvenance>,
107}