Skip to main content

ahri_tre_observability/
lib.rs

1//! Safe, bounded operational events. Libraries never install a subscriber.
2
3mod 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
13/// Maximum serialized JSON size, excluding its terminating newline.
14pub 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    /// Only these fixed authentication routes may enter the event contract.
78    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    /// Resolve only release-owned kinds; arbitrary input never enters an event.
92    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/// Already-public protocol codes retained without error messages or details.
214#[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/// Optional numeric evidence has a closed vocabulary and contains no content.
233#[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/// Non-authoritative correlation; parentage is meaningful only within a process.
242#[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    /// Join an existing opaque Operation. This grants no workflow authority.
261    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
279/// Explicitly pass `context()` across execution boundaries; no ambient span is
280/// entered across an async suspension. Completing a span emits one event.
281pub 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    /// Timestamp evidence when the owner obtains it, before asynchronous delivery.
293    /// A later status projection must retain this instant, never replace it.
294    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    /// Number an existing transport attempt; this never authorizes a retry.
300    pub fn with_attempt(mut self, attempt: u32) -> Self {
301        self.attempt = Some(attempt);
302        self
303    }
304
305    /// Emit one bounded dispatch entry before work begins.
306    pub fn record_start(&self) {
307        self.emit(None, None, &[]);
308    }
309
310    /// Attach only a closed, already-public code; no message or detail is accepted.
311    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        // Bound allocation even if a caller supplies excessive optional evidence.
348        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/// Immutable provenance copied from the owning executable's resolved state.
400/// Validation prevents free text from entering the operational schema.
401#[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        // Stable releases and numbered release candidates only: arbitrary
449        // prerelease/build labels are not evidence.
450        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/// Release-owned process signals; no provider errors or listener details enter logs.
474#[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}