Skip to main content

ahri_tre_app/
session_observations.rs

1//! Bounded, non-persistent evidence owned by one live Session's workflows.
2
3use ahri_tre_protocol::diagnostics::{
4    DiagnosticComponent, DiagnosticEvidenceBasis, DiagnosticFinding, DiagnosticFindingSeverity,
5    DiagnosticFindingStatus, DiagnosticNextAction, DiagnosticObservation,
6    DiagnosticObservationOrigin, DiagnosticScope,
7};
8use chrono::{DateTime, Utc};
9use std::sync::{Arc, Mutex};
10
11const CATALOG_AUTHENTICATION: usize = 3;
12const CATALOG_READ: usize = 4;
13
14/// Holds only the latest authentication, metadata and Lake observations.
15/// Clones belong to the same live Session, including its admitted background work.
16#[derive(Clone)]
17pub struct SessionObservations(Arc<Mutex<Observations>>);
18
19struct Observations {
20    created_at: DateTime<Utc>,
21    sequence: u64,
22    findings: [Option<(u64, DiagnosticFinding)>; 5],
23}
24
25impl Default for SessionObservations {
26    fn default() -> Self {
27        Self(Arc::new(Mutex::new(Observations {
28            created_at: Utc::now(),
29            sequence: 0,
30            findings: [None, None, None, None, None],
31        })))
32    }
33}
34
35impl SessionObservations {
36    /// Capture under the Session's workflow lock before recording this request.
37    pub fn cursor(&self) -> u64 {
38        self.0
39            .lock()
40            .unwrap_or_else(|error| error.into_inner())
41            .sequence
42    }
43
44    pub fn record(
45        &self,
46        component: DiagnosticComponent,
47        origin: DiagnosticObservationOrigin,
48        succeeded: bool,
49    ) {
50        self.record_at(component, origin, succeeded, Utc::now());
51    }
52
53    pub fn record_at(
54        &self,
55        component: DiagnosticComponent,
56        origin: DiagnosticObservationOrigin,
57        succeeded: bool,
58        observed_at: DateTime<Utc>,
59    ) {
60        let Some(index) = component_index(component) else {
61            return;
62        };
63        let mut finding = finding(component, observed_at);
64        finding.status = if succeeded {
65            DiagnosticFindingStatus::Passed
66        } else {
67            DiagnosticFindingStatus::Failed
68        };
69        finding.severity = if succeeded {
70            DiagnosticFindingSeverity::Info
71        } else {
72            DiagnosticFindingSeverity::Error
73        };
74        finding.summary = match (component, succeeded) {
75            (DiagnosticComponent::Authentication, true) => {
76                "Retained Session authentication was accepted; provider was not contacted"
77            }
78            (DiagnosticComponent::Authentication, false) => {
79                "Retained Session authentication was rejected"
80            }
81            (DiagnosticComponent::Metadata, true) => "Metadata operation succeeded",
82            (DiagnosticComponent::Metadata, false) => "Metadata operation failed",
83            (_, true) => "Lake operation succeeded",
84            (_, false) => "Lake operation failed",
85        }
86        .into();
87        let observation = finding
88            .observation
89            .as_mut()
90            .expect("owned finding has evidence");
91        observation.origin = Some(origin);
92        observation.basis = DiagnosticEvidenceBasis::HistoricalObservation;
93        self.retain(index, finding);
94    }
95
96    /// Consume only stage evidence from the Lake operation that already ran.
97    pub fn record_catalog(
98        &self,
99        catalog: &ahri_tre_lake::LakeCatalogObservations,
100        origin: DiagnosticObservationOrigin,
101    ) {
102        use ahri_tre_lake::CatalogObservationResult as Result;
103        use ahri_tre_protocol::diagnostics::DiagnosticFailureCategory as Category;
104        for (index, evidence) in [
105            (CATALOG_AUTHENTICATION, catalog.authentication()),
106            (CATALOG_READ, catalog.read()),
107        ] {
108            let Some(evidence) = evidence else { continue };
109            let mut finding = initial_finding(index, evidence.observed_at_utc);
110            let (status, summary, category) = match evidence.result {
111                Result::Passed => (
112                    DiagnosticFindingStatus::Passed,
113                    if index == CATALOG_AUTHENTICATION {
114                        "Lake catalog authentication succeeded"
115                    } else {
116                        "Lake catalog read succeeded"
117                    },
118                    None,
119                ),
120                Result::AuthenticationRejected => (
121                    DiagnosticFindingStatus::Failed,
122                    "Lake catalog authentication was rejected",
123                    Some(Category::Authentication),
124                ),
125                Result::AuthenticationNotEstablished => (
126                    DiagnosticFindingStatus::Skipped,
127                    "Lake catalog authentication was not established; attach failed without confirmed authentication rejection",
128                    Some(Category::Unavailable),
129                ),
130                Result::ReadFailed => (
131                    DiagnosticFindingStatus::Failed,
132                    "Lake catalog read failed after authentication was established",
133                    Some(Category::Query),
134                ),
135                Result::ReadNotAttempted => (
136                    DiagnosticFindingStatus::Skipped,
137                    "Lake catalog read was not checked because authentication was not established",
138                    Some(Category::AuthorityUnavailable),
139                ),
140            };
141            finding.status = status;
142            finding.summary = summary.into();
143            finding.severity = match status {
144                DiagnosticFindingStatus::Failed => DiagnosticFindingSeverity::Error,
145                DiagnosticFindingStatus::Skipped => DiagnosticFindingSeverity::Warning,
146                _ => DiagnosticFindingSeverity::Info,
147            };
148            let observation = finding
149                .observation
150                .as_mut()
151                .expect("owned finding has evidence");
152            observation.origin = Some(origin);
153            observation.failure_category = category;
154            observation.basis = if status == DiagnosticFindingStatus::Skipped {
155                DiagnosticEvidenceBasis::NotChecked
156            } else {
157                DiagnosticEvidenceBasis::HistoricalObservation
158            };
159            self.retain(index, finding);
160        }
161    }
162
163    fn retain(&self, index: usize, finding: DiagnosticFinding) {
164        let observed_at = finding
165            .observation
166            .as_ref()
167            .expect("owned finding has evidence")
168            .observed_at_utc;
169        let mut state = self.0.lock().unwrap_or_else(|error| error.into_inner());
170        // Authentication may finish before waiting for the workflow lock.
171        // A delayed request must not replace evidence acquired more recently.
172        if state.findings[index]
173            .as_ref()
174            .and_then(|(_, finding)| finding.observation.as_ref())
175            .is_some_and(|retained| retained.observed_at_utc > observed_at)
176        {
177            return;
178        }
179        state.sequence = state.sequence.saturating_add(1);
180        state.findings[index] = Some((state.sequence, finding));
181    }
182
183    /// Status/close pass None. A workflow may identify only evidence obtained
184    /// since its cursor as current; opening evidence always remains historical.
185    pub fn snapshot(&self, after: Option<u64>) -> Vec<DiagnosticFinding> {
186        let state = self.0.lock().unwrap_or_else(|error| error.into_inner());
187        (0..state.findings.len())
188            .map(|index| {
189                let Some((sequence, retained)) = &state.findings[index] else {
190                    return initial_finding(index, state.created_at);
191                };
192                let mut finding = retained.clone();
193                let observation = finding
194                    .observation
195                    .as_mut()
196                    .expect("owned finding has evidence");
197                if after.is_some_and(|cursor| *sequence > cursor)
198                    && observation.origin == Some(DiagnosticObservationOrigin::Workflow)
199                    && observation.basis == DiagnosticEvidenceBasis::HistoricalObservation
200                {
201                    observation.basis = DiagnosticEvidenceBasis::AuthenticatedAdapterObservation;
202                }
203                finding
204            })
205            .collect()
206    }
207}
208
209fn component_index(component: DiagnosticComponent) -> Option<usize> {
210    match component {
211        DiagnosticComponent::Authentication => Some(0),
212        DiagnosticComponent::Metadata => Some(1),
213        DiagnosticComponent::Lake => Some(2),
214        _ => None,
215    }
216}
217
218fn finding(component: DiagnosticComponent, observed_at_utc: DateTime<Utc>) -> DiagnosticFinding {
219    let (id, action) = match component {
220        DiagnosticComponent::Authentication => {
221            ("session.authentication", DiagnosticNextAction::AuthStatus)
222        }
223        DiagnosticComponent::Metadata => (
224            "session.metadata",
225            DiagnosticNextAction::ConfigPreflightHelp,
226        ),
227        _ => ("session.lake", DiagnosticNextAction::ConfigPreflightHelp),
228    };
229    DiagnosticFinding {
230        id: id.into(),
231        component: None,
232        status: DiagnosticFindingStatus::Skipped,
233        severity: DiagnosticFindingSeverity::Info,
234        summary: "No adapter observation is available".into(),
235        detail: None,
236        evidence: Vec::new(),
237        observation: Some(DiagnosticObservation {
238            component,
239            scope: DiagnosticScope::OwnerSession,
240            observed_at_utc,
241            basis: DiagnosticEvidenceBasis::NotChecked,
242            origin: None,
243            failure_category: None,
244            next_action: Some(action),
245        }),
246    }
247}
248
249fn initial_finding(index: usize, observed_at: DateTime<Utc>) -> DiagnosticFinding {
250    let component = match index {
251        0 => DiagnosticComponent::Authentication,
252        1 => DiagnosticComponent::Metadata,
253        _ => DiagnosticComponent::Lake,
254    };
255    let mut finding = finding(component, observed_at);
256    if index >= CATALOG_AUTHENTICATION {
257        finding.id = if index == CATALOG_AUTHENTICATION {
258            "session.lake.catalog_authentication"
259        } else {
260            "session.lake.catalog_read"
261        }
262        .into();
263        finding.summary = if index == CATALOG_AUTHENTICATION {
264            "Lake opening did not reach catalog authentication"
265        } else {
266            "Lake opening did not reach the catalog read"
267        }
268        .into();
269        finding
270            .observation
271            .as_mut()
272            .expect("owned finding has evidence")
273            .origin = Some(DiagnosticObservationOrigin::SessionOpen);
274    }
275    finding
276}