Skip to main content

ahri_tre_observability/
delivery.rs

1use crate::lifecycle::Lifecycle;
2use crate::{
3    Completion, ConfigurationProvenance, MAX_EVENT_BYTES, Outcome, ProcessEvent, ServiceRole,
4    TARGET,
5};
6use chrono::{SecondsFormat, Utc};
7use serde::Serialize;
8use std::io::{self, Write};
9use std::sync::{
10    Arc, OnceLock,
11    atomic::{AtomicU64, Ordering},
12    mpsc::{self, SyncSender},
13};
14use std::time::Duration;
15use tracing::{
16    Event, Subscriber,
17    field::{Field, Visit},
18};
19use tracing_subscriber::{Layer, layer::Context, prelude::*};
20use uuid::Uuid;
21
22/// Drop newest when all 256 event slots are occupied. The worker can hold one
23/// additional event. Failed writes are discarded, never retried indefinitely.
24pub const QUEUE_CAPACITY: usize = 256;
25
26struct Process {
27    role: ServiceRole,
28    id: Uuid,
29    provenance: OnceLock<ConfigurationProvenance>,
30    dropped: AtomicU64,
31}
32
33impl Process {
34    fn drop_events(&self, count: u64) {
35        let _ = self
36            .dropped
37            .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |previous| {
38                Some(previous.saturating_add(count))
39            });
40    }
41
42    fn encode(
43        &self,
44        mut completion: Option<Completion>,
45        dropped_events: Option<u64>,
46        signal: Option<ProcessEvent>,
47        mut lifecycle: Option<Lifecycle>,
48    ) -> Vec<u8> {
49        let provenance = lifecycle.as_mut().and_then(|event| event.provenance.take());
50        let severity = match completion
51            .as_ref()
52            .and_then(|event| event.outcome)
53            .or_else(|| lifecycle.as_ref().and_then(|event| event.outcome))
54        {
55            Some(Outcome::Success | Outcome::Accepted | Outcome::Skipped) => "info",
56            None if lifecycle.is_some()
57                || completion.is_some()
58                || matches!(signal, Some(ProcessEvent::Ready)) =>
59            {
60                "info"
61            }
62            Some(Outcome::InternalFailure) => "error",
63            _ => "warn",
64        };
65        let summary = ahri_tre_security::sanitize_free_text(if dropped_events.is_some() {
66            "Operational event delivery recovered."
67        } else if let Some(lifecycle) = &lifecycle {
68            if lifecycle.outcome.is_some() {
69                "Process phase completed."
70            } else {
71                "Process phase started."
72            }
73        } else {
74            match signal {
75                Some(ProcessEvent::BootstrapFailed) => "Service bootstrap failed.",
76                Some(_) => "Service state observed.",
77                None if completion
78                    .as_ref()
79                    .is_some_and(|event| event.outcome.is_none()) =>
80                {
81                    "Operation stage started."
82                }
83                None => "Operation stage completed.",
84            }
85        });
86        let encode = || {
87            serde_json::to_vec(&Envelope {
88                schema_version: 1,
89                timestamp: Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true),
90                severity,
91                event: if let Some(signal) = signal {
92                    match signal {
93                        ProcessEvent::Ready => "process_ready",
94                        ProcessEvent::BootstrapFailed => "bootstrap_failed",
95                        ProcessEvent::RecoveryUnavailable => "recovery_unavailable",
96                        ProcessEvent::ConnectionRejected => "connection_rejected",
97                        ProcessEvent::AdministrationRejected => "administration_rejected",
98                    }
99                } else if let Some(lifecycle) = &lifecycle {
100                    if lifecycle.outcome.is_some() {
101                        "process_phase_completed"
102                    } else {
103                        "process_phase_started"
104                    }
105                } else if dropped_events.is_some() {
106                    "events_dropped"
107                } else if completion
108                    .as_ref()
109                    .is_some_and(|event| event.outcome.is_none())
110                {
111                    "operation_started"
112                } else {
113                    "operation_completed"
114                },
115                service_role: self.role,
116                package_version: env!("CARGO_PKG_VERSION"),
117                process_instance_id: self.id,
118                provenance: provenance.as_ref().or_else(|| {
119                    if lifecycle
120                        .as_ref()
121                        .is_some_and(|event| matches!(event.phase, crate::Phase::Configuration))
122                    {
123                        None
124                    } else {
125                        self.provenance.get()
126                    }
127                }),
128                lifecycle: lifecycle.as_ref(),
129                completion: completion.as_ref(),
130                dropped_events,
131                summary: &summary,
132            })
133            .expect("closed operational fields serialize")
134        };
135        let mut bytes = encode();
136        if bytes.len() > MAX_EVENT_BYTES {
137            if let Some(completion) = completion.as_mut() {
138                completion.evidence.clear();
139                completion.truncated = true;
140            }
141            bytes = self.encode(completion, dropped_events, signal, lifecycle);
142            // Recursive output already includes the delimiter.
143            return bytes;
144        }
145        bytes.push(b'\n');
146        bytes
147    }
148}
149
150#[derive(Serialize)]
151struct Envelope<'a> {
152    schema_version: u64,
153    timestamp: String,
154    severity: &'static str,
155    event: &'static str,
156    service_role: ServiceRole,
157    package_version: &'static str,
158    process_instance_id: Uuid,
159    #[serde(flatten)]
160    provenance: Option<&'a ConfigurationProvenance>,
161    #[serde(flatten)]
162    completion: Option<&'a Completion>,
163    #[serde(flatten)]
164    lifecycle: Option<&'a Lifecycle>,
165    #[serde(skip_serializing_if = "Option::is_none")]
166    dropped_events: Option<u64>,
167    summary: &'a str,
168}
169
170enum Delivery {
171    Event(Vec<u8>),
172    Flush(SyncSender<bool>),
173}
174
175/// Owns bounded asynchronous delivery. Dropping this handle never joins a
176/// potentially blocked writer. `flush` is an explicit, bounded best effort.
177pub struct Logging {
178    process: Arc<Process>,
179    sender: SyncSender<Delivery>,
180}
181
182impl Logging {
183    pub fn with_writer(role: ServiceRole, writer: impl Write + Send + 'static) -> io::Result<Self> {
184        let process = Arc::new(Process {
185            role,
186            id: Uuid::new_v4(),
187            provenance: OnceLock::new(),
188            dropped: AtomicU64::new(0),
189        });
190        let (sender, receiver) = mpsc::sync_channel(QUEUE_CAPACITY);
191        let worker = Arc::clone(&process);
192        std::thread::Builder::new()
193            .name("operational-events".into())
194            .spawn(move || {
195                let mut writer = FramedWriter {
196                    inner: writer,
197                    incomplete: false,
198                };
199                while let Ok(delivery) = receiver.recv() {
200                    match delivery {
201                        Delivery::Event(bytes) => {
202                            if writer.write_event(&bytes).is_err() {
203                                worker.drop_events(1);
204                                continue;
205                            }
206                            let dropped = worker.dropped.swap(0, Ordering::Relaxed);
207                            if dropped > 0
208                                && writer
209                                    .write_event(&worker.encode(None, Some(dropped), None, None))
210                                    .is_err()
211                            {
212                                worker.drop_events(dropped);
213                            }
214                        }
215                        Delivery::Flush(acknowledge) => {
216                            let _ = acknowledge.try_send(writer.inner.flush().is_ok());
217                        }
218                    }
219                }
220            })?;
221        Ok(Self { process, sender })
222    }
223
224    /// Executables alone install this fixed subscriber. Ambient filter settings
225    /// and arbitrary library log fields are never interpreted as logging policy.
226    pub fn stderr(role: ServiceRole) -> io::Result<Self> {
227        let logging = Self::with_writer(role, io::stderr())?;
228        tracing_subscriber::registry()
229            .with(logging.layer())
230            .try_init()
231            .map_err(io::Error::other)?;
232        Ok(logging)
233    }
234
235    pub fn layer(&self) -> OperationalLayer {
236        OperationalLayer {
237            process: Arc::clone(&self.process),
238            sender: self.sender.clone(),
239        }
240    }
241
242    pub fn set_provenance(&self, provenance: ConfigurationProvenance) {
243        let _ = self.process.provenance.set(provenance);
244    }
245
246    pub fn flush(&self, timeout: Duration) -> bool {
247        let (sender, receiver) = mpsc::sync_channel(1);
248        self.sender.try_send(Delivery::Flush(sender)).is_ok()
249            && receiver.recv_timeout(timeout).unwrap_or(false)
250    }
251}
252
253pub struct OperationalLayer {
254    process: Arc<Process>,
255    sender: SyncSender<Delivery>,
256}
257
258impl<S: Subscriber> Layer<S> for OperationalLayer {
259    fn on_event(&self, event: &Event<'_>, _context: Context<'_, S>) {
260        if event.metadata().target() != TARGET {
261            return;
262        }
263        let mut visitor = CompletionVisitor {
264            completion: None,
265            signal: None,
266            lifecycle: None,
267        };
268        event.record(&mut visitor);
269        if visitor.completion.is_some() || visitor.signal.is_some() || visitor.lifecycle.is_some() {
270            let bytes =
271                self.process
272                    .encode(visitor.completion, None, visitor.signal, visitor.lifecycle);
273            if self.sender.try_send(Delivery::Event(bytes)).is_err() {
274                self.process.drop_events(1);
275            }
276        }
277    }
278}
279
280struct CompletionVisitor {
281    completion: Option<Completion>,
282    signal: Option<ProcessEvent>,
283    lifecycle: Option<Lifecycle>,
284}
285
286impl Visit for CompletionVisitor {
287    fn record_str(&mut self, field: &Field, value: &str) {
288        // Revalidate even the private tracing bridge before output. This rejects
289        // arbitrary DTOs, unknown fields, names, and untyped identifiers.
290        if field.name() == "record" && value.len() <= 2 * MAX_EVENT_BYTES {
291            self.completion = serde_json::from_str(value).ok();
292        } else if field.name() == "lifecycle" && value.len() <= MAX_EVENT_BYTES {
293            self.lifecycle = serde_json::from_str::<Lifecycle>(value)
294                .ok()
295                .map(|mut event| {
296                    event.provenance = event
297                        .provenance
298                        .and_then(ConfigurationProvenance::validated);
299                    event
300                });
301        } else if field.name() == "signal" && value.len() < 64 {
302            self.signal = serde_json::from_str(value).ok();
303        }
304    }
305
306    fn record_debug(&mut self, _field: &Field, _value: &dyn std::fmt::Debug) {}
307}
308
309// A failed sink may have consumed an irreversible prefix. Delimit that damaged
310// record before the next attempt so recovery never corrupts a later JSON event.
311struct FramedWriter<W> {
312    inner: W,
313    incomplete: bool,
314}
315
316impl<W: Write> FramedWriter<W> {
317    fn write_event(&mut self, bytes: &[u8]) -> io::Result<()> {
318        if self.incomplete {
319            self.inner.write_all(b"\n")?;
320            self.incomplete = false;
321        }
322        let mut remaining = bytes;
323        while !remaining.is_empty() {
324            match self.inner.write(remaining) {
325                Ok(0) => return Err(io::ErrorKind::WriteZero.into()),
326                Ok(written) => {
327                    self.incomplete = remaining[written - 1] != b'\n';
328                    remaining = &remaining[written..];
329                }
330                Err(error) if error.kind() == io::ErrorKind::Interrupted => continue,
331                Err(error) => return Err(error),
332            }
333        }
334        self.inner.flush()
335    }
336}