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
22pub 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 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
175pub 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 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 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
309struct 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}