1mod delivery;
4mod lifecycle;
5
6pub use lifecycle::{LifecycleSpan, Phase};
7
8pub use delivery::{Logging, OperationalLayer, QUEUE_CAPACITY};
9use serde::{Deserialize, Serialize};
10use std::time::Instant;
11use uuid::Uuid;
12
13pub const MAX_EVENT_BYTES: usize = 8 * 1024;
15const TARGET: &str = "ahri_tre_observability";
16
17#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
18#[serde(rename_all = "snake_case")]
19pub enum ServiceRole {
20 Web,
21 TrustedRuntime,
22 ManagedRuntime,
23 Cli,
24}
25
26#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
27#[serde(rename_all = "snake_case")]
28pub enum Operation {
29 DatasetFromDatafile,
30 OperationList,
31 OperationGet,
32 OperationCancel,
33 OperationResultGet,
34 RuntimeLoginStart,
35 RuntimeLoginCallback,
36 RuntimeLoginPoll,
37 RuntimeLoginAcknowledge,
38 RuntimeCredentialStatus,
39 RuntimeCredentialRotate,
40 RuntimeCredentialLogout,
41 SessionOpen,
42 SessionStatus,
43 SessionClose,
44 DomainList,
45 BrowserDatastoreList,
46 BrowserStudyList,
47 BrowserStudyGet,
48 BrowserStudySearch,
49 BrowserResourceList,
50 BrowserDatasetSearch,
51 BrowserDataFileSearch,
52 BrowserDictionaryGet,
53 BrowserVariableSearch,
54 BrowserDomainList,
55 BrowserDomainGet,
56 BrowserDomainStudyList,
57 BrowserDomainVariableList,
58 BrowserDomainVariableGet,
59 BrowserVocabularyGet,
60 BrowserVocabularyItemList,
61 BrowserVocabularyMappingList,
62 BrowserSemanticCatalogGet,
63 BrowserModelEntitySearch,
64 BrowserModelRelationSearch,
65 BrowserAccessRequestSubmit,
66 BrowserAccessRequestList,
67 BrowserAccessRequestGet,
68 BrowserAccessRequestWithdraw,
69 BrowserCustodianAccessRequestList,
70 BrowserCustodianAccessRequestApprove,
71 BrowserCustodianAccessRequestReject,
72 OtherRequest,
73 RuntimeCapabilities,
74}
75
76impl Operation {
77 pub fn from_runtime_login_path(path: &str) -> Option<Self> {
79 Some(match path {
80 "/v1/runtime-login/start" => Self::RuntimeLoginStart,
81 "/v1/runtime-login/callback" => Self::RuntimeLoginCallback,
82 "/v1/runtime-login/poll" => Self::RuntimeLoginPoll,
83 "/v1/runtime-login/acknowledge" => Self::RuntimeLoginAcknowledge,
84 "/v1/runtime-credential/status" => Self::RuntimeCredentialStatus,
85 "/v1/runtime-credential/rotate" => Self::RuntimeCredentialRotate,
86 "/v1/runtime-credential/logout" => Self::RuntimeCredentialLogout,
87 _ => return None,
88 })
89 }
90
91 pub fn from_request_kind(kind: &str) -> Self {
93 match kind {
94 "ingest.dataset.from_datafile" => Self::DatasetFromDatafile,
95 "operation.list" => Self::OperationList,
96 "operation.get" => Self::OperationGet,
97 "operation.cancel" => Self::OperationCancel,
98 "operation.result.get" => Self::OperationResultGet,
99 "session.open" => Self::SessionOpen,
100 "session.status" => Self::SessionStatus,
101 "session.close" => Self::SessionClose,
102 "domain.list" => Self::DomainList,
103 "browser.datastore.list" => Self::BrowserDatastoreList,
104 "browser.study.list" => Self::BrowserStudyList,
105 "browser.study.get" => Self::BrowserStudyGet,
106 "browser.study.search" => Self::BrowserStudySearch,
107 "browser.resource.list" => Self::BrowserResourceList,
108 "browser.dataset.search" => Self::BrowserDatasetSearch,
109 "browser.datafile.search" => Self::BrowserDataFileSearch,
110 "browser.dictionary.get" => Self::BrowserDictionaryGet,
111 "browser.dictionary.variable.search" => Self::BrowserVariableSearch,
112 "browser.domain.list" => Self::BrowserDomainList,
113 "browser.domain.get" => Self::BrowserDomainGet,
114 "browser.domain.study.list" => Self::BrowserDomainStudyList,
115 "browser.domain.variable.list" => Self::BrowserDomainVariableList,
116 "browser.domain.variable.get" => Self::BrowserDomainVariableGet,
117 "browser.vocabulary.get" => Self::BrowserVocabularyGet,
118 "browser.vocabulary.item.list" => Self::BrowserVocabularyItemList,
119 "browser.vocabulary.mapping.list" => Self::BrowserVocabularyMappingList,
120 "browser.semantic_catalog.get" => Self::BrowserSemanticCatalogGet,
121 "browser.model.entity.search" => Self::BrowserModelEntitySearch,
122 "browser.model.relation.search" => Self::BrowserModelRelationSearch,
123 "browser.access_request.submit" => Self::BrowserAccessRequestSubmit,
124 "browser.access_request.list" => Self::BrowserAccessRequestList,
125 "browser.access_request.get" => Self::BrowserAccessRequestGet,
126 "browser.access_request.withdraw" => Self::BrowserAccessRequestWithdraw,
127 "browser.custodian.access_request.list" => Self::BrowserCustodianAccessRequestList,
128 "browser.custodian.access_request.approve" => {
129 Self::BrowserCustodianAccessRequestApprove
130 }
131 "browser.custodian.access_request.reject" => Self::BrowserCustodianAccessRequestReject,
132 "runtime.capabilities" => Self::RuntimeCapabilities,
133 _ => Self::OtherRequest,
134 }
135 }
136}
137
138#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
139#[serde(rename_all = "snake_case")]
140pub enum Stage {
141 Preparation,
142 MetadataCommit,
143 Compensation,
144 LakeStorage,
145 LakeRead,
146 LakeCatalogAuthentication,
147 LakeCatalogRead,
148 ProviderDiscovery,
149 TokenExchange,
150 CredentialValidation,
151 Admission,
152 Session,
153 Lake,
154 Cli,
155 ManagedRuntime,
156 Transport,
157 Authentication,
158 Authorization,
159 SecretCapability,
160 WorkflowEntry,
161 MetadataPreflight,
162 MetadataDecision,
163 MetadataConnection,
164 MetadataQuery,
165 Web,
166 TrustedRuntime,
167 Application,
168 Metadata,
169}
170
171#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
172#[serde(rename_all = "snake_case")]
173pub enum Outcome {
174 Interrupted,
175 Skipped,
176 Accepted,
177 Success,
178 Rejected,
179 Unavailable,
180 Timeout,
181 Cancelled,
182 InternalFailure,
183}
184
185#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
186#[serde(rename_all = "snake_case")]
187pub enum FailureCategory {
188 LegacyEnvironment,
189 Configuration,
190 Secret,
191 Listener,
192 Interrupted,
193 Metadata,
194 Storage,
195 Catalog,
196 Cleanup,
197 CommitUnknown,
198 Provider,
199 CredentialInvalid,
200 CredentialExpired,
201 Admission,
202 Adapter,
203 Authentication,
204 Authorization,
205 Connection,
206 Tls,
207 Query,
208 Compatibility,
209 Dependency,
210 Internal,
211}
212
213#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
215#[serde(rename_all = "snake_case")]
216pub enum PublicErrorCode {
217 ValidationFailed,
218 NotFound,
219 Conflict,
220 Unauthorized,
221 Forbidden,
222 PreconditionFailed,
223 UnsupportedProtocolVersion,
224 UnsupportedOperation,
225 OperationFailed,
226 OperationNotReady,
227 OperationResultUnavailable,
228 RateLimited,
229 Internal,
230}
231
232#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
234#[serde(rename_all = "snake_case")]
235pub enum Measurement {
236 ReturnedRecords(u64),
237 Rows(u64),
238 Columns(u64),
239}
240
241#[derive(Debug, Clone, Copy)]
243pub struct CorrelationContext {
244 request_id: Uuid,
245 operation: Operation,
246 parent_span_id: Option<Uuid>,
247 operation_id: Option<Uuid>,
248}
249
250impl CorrelationContext {
251 pub fn new(request_id: Uuid, operation: Operation) -> Self {
252 Self {
253 request_id,
254 operation,
255 parent_span_id: None,
256 operation_id: None,
257 }
258 }
259
260 pub fn with_operation_id(mut self, operation_id: Uuid) -> Self {
262 self.operation_id = Some(operation_id);
263 self
264 }
265
266 pub fn span(self, stage: Stage) -> OperationSpan {
267 OperationSpan {
268 context: self,
269 span_id: Uuid::new_v4(),
270 stage,
271 error_code: None,
272 attempt: None,
273 observed_at: None,
274 started: Instant::now(),
275 }
276 }
277}
278
279pub struct OperationSpan {
282 context: CorrelationContext,
283 span_id: Uuid,
284 stage: Stage,
285 error_code: Option<PublicErrorCode>,
286 attempt: Option<u32>,
287 observed_at: Option<chrono::DateTime<chrono::Utc>>,
288 started: Instant,
289}
290
291impl OperationSpan {
292 pub fn finish_observed(mut self, outcome: Outcome, category: Option<FailureCategory>) {
295 self.observed_at = Some(chrono::Utc::now());
296 self.finish(outcome, category);
297 }
298
299 pub fn with_attempt(mut self, attempt: u32) -> Self {
301 self.attempt = Some(attempt);
302 self
303 }
304
305 pub fn record_start(&self) {
307 self.emit(None, None, &[]);
308 }
309
310 pub fn with_error_code(mut self, code: PublicErrorCode) -> Self {
312 self.error_code = Some(code);
313 self
314 }
315
316 pub fn with_operation_id(mut self, operation_id: Uuid) -> Self {
317 self.context = self.context.with_operation_id(operation_id);
318 self
319 }
320
321 pub fn context(&self) -> CorrelationContext {
322 CorrelationContext {
323 parent_span_id: Some(self.span_id),
324 ..self.context
325 }
326 }
327
328 pub fn finish(self, outcome: Outcome, failure_category: Option<FailureCategory>) {
329 self.finish_with_evidence(outcome, failure_category, &[]);
330 }
331
332 pub fn finish_with_evidence(
333 self,
334 outcome: Outcome,
335 failure_category: Option<FailureCategory>,
336 evidence: &[Measurement],
337 ) {
338 self.emit(Some(outcome), failure_category, evidence);
339 }
340
341 fn emit(
342 &self,
343 outcome: Option<Outcome>,
344 failure_category: Option<FailureCategory>,
345 evidence: &[Measurement],
346 ) {
347 let completion = Completion {
349 request_id: self.context.request_id,
350 operation_id: self.context.operation_id,
351 operation: self.context.operation,
352 span_id: self.span_id,
353 parent_span_id: self.context.parent_span_id,
354 stage: self.stage,
355 outcome,
356 failure_category,
357 error_code: self.error_code,
358 attempt: self.attempt,
359 observed_at: self.observed_at,
360 duration_ms: outcome
361 .map(|_| self.started.elapsed().as_millis().min(u64::MAX as u128) as u64),
362 evidence: evidence.iter().take(256).copied().collect(),
363 truncated: evidence.len() > 256,
364 };
365 if let Ok(record) = serde_json::to_string(&completion) {
366 tracing::info!(target: TARGET, record = record.as_str());
367 }
368 }
369}
370
371#[derive(Serialize, Deserialize)]
372#[serde(deny_unknown_fields)]
373struct Completion {
374 request_id: Uuid,
375 #[serde(skip_serializing_if = "Option::is_none")]
376 operation_id: Option<Uuid>,
377 operation: Operation,
378 span_id: Uuid,
379 #[serde(skip_serializing_if = "Option::is_none")]
380 parent_span_id: Option<Uuid>,
381 stage: Stage,
382 #[serde(skip_serializing_if = "Option::is_none")]
383 outcome: Option<Outcome>,
384 #[serde(skip_serializing_if = "Option::is_none")]
385 failure_category: Option<FailureCategory>,
386 #[serde(skip_serializing_if = "Option::is_none")]
387 error_code: Option<PublicErrorCode>,
388 #[serde(skip_serializing_if = "Option::is_none")]
389 attempt: Option<u32>,
390 #[serde(skip_serializing_if = "Option::is_none")]
391 observed_at: Option<chrono::DateTime<chrono::Utc>>,
392 #[serde(skip_serializing_if = "Option::is_none")]
393 duration_ms: Option<u64>,
394 #[serde(default, skip_serializing_if = "Vec::is_empty")]
395 evidence: Vec<Measurement>,
396 truncated: bool,
397}
398
399#[derive(Debug, Clone, Serialize)]
402pub struct ConfigurationProvenance {
403 deployment_id: Uuid,
404 configuration_fingerprint: String,
405 configuration_schema_version: u64,
406 application_version: String,
407}
408
409impl<'de> Deserialize<'de> for ConfigurationProvenance {
410 fn deserialize<D: serde::Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
411 #[derive(Deserialize)]
412 #[serde(deny_unknown_fields)]
413 struct Fields {
414 deployment_id: Uuid,
415 configuration_fingerprint: String,
416 configuration_schema_version: u64,
417 application_version: String,
418 }
419 let fields = Fields::deserialize(deserializer)?;
420 Self::new(
421 fields.deployment_id,
422 &fields.configuration_fingerprint,
423 fields.configuration_schema_version,
424 &fields.application_version,
425 )
426 .ok_or_else(|| serde::de::Error::custom("invalid configuration provenance"))
427 }
428}
429
430impl ConfigurationProvenance {
431 pub(crate) fn validated(self) -> Option<Self> {
432 Self::new(
433 self.deployment_id,
434 &self.configuration_fingerprint,
435 self.configuration_schema_version,
436 &self.application_version,
437 )
438 }
439 pub fn new(
440 deployment_id: Uuid,
441 fingerprint: &str,
442 schema_version: u64,
443 application_version: &str,
444 ) -> Option<Self> {
445 if fingerprint.len() != 64 || !fingerprint.bytes().all(|byte| byte.is_ascii_hexdigit()) {
446 return None;
447 }
448 let version = semver::Version::parse(application_version).ok()?;
451 let release_candidate = version
452 .pre
453 .as_str()
454 .strip_prefix("rc.")
455 .is_some_and(|number| {
456 !number.is_empty() && number.bytes().all(|byte| byte.is_ascii_digit())
457 });
458 if (!version.pre.is_empty() && !release_candidate)
459 || !version.build.is_empty()
460 || application_version.len() > 64
461 {
462 return None;
463 }
464 Some(Self {
465 deployment_id,
466 configuration_fingerprint: fingerprint.to_owned(),
467 configuration_schema_version: schema_version,
468 application_version: application_version.to_owned(),
469 })
470 }
471}
472
473#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
475#[serde(rename_all = "snake_case")]
476pub enum ProcessEvent {
477 Ready,
478 BootstrapFailed,
479 RecoveryUnavailable,
480 ConnectionRejected,
481 AdministrationRejected,
482}
483
484impl ProcessEvent {
485 pub fn emit(self) {
486 if let Ok(signal) = serde_json::to_string(&self) {
487 tracing::info!(target: TARGET, signal = signal.as_str());
488 }
489 }
490}