1use 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#[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 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 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 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 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}