1use std::io::{self, Write};
21use std::str::FromStr;
22use std::sync::{Arc, Mutex, RwLock};
23
24use serde::{Deserialize, Serialize};
25use tracing::field::{Field, Visit};
26use tracing::span;
27use tracing::{Event, Level, Subscriber};
28use tracing_subscriber::Layer;
29use tracing_subscriber::fmt::MakeWriter;
30use tracing_subscriber::layer::Context;
31use tracing_subscriber::registry::LookupSpan;
32use tracing_subscriber::util::TryInitError;
33
34#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
36#[serde(rename_all = "lowercase")]
37pub enum LogFormat {
38 Json,
39 Plain,
40 Yaml,
42}
43
44impl LogFormat {
45 pub const fn as_str(self) -> &'static str {
47 match self {
48 Self::Json => "json",
49 Self::Plain => "plain",
50 Self::Yaml => "yaml",
51 }
52 }
53}
54
55impl std::fmt::Display for LogFormat {
56 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
57 formatter.write_str(self.as_str())
58 }
59}
60
61impl FromStr for LogFormat {
62 type Err = String;
63
64 fn from_str(value: &str) -> Result<Self, Self::Err> {
65 match value {
66 "json" => Ok(Self::Json),
67 "plain" => Ok(Self::Plain),
68 "yaml" => Ok(Self::Yaml),
69 _ => Err("invalid log format: expected json, plain, or yaml".to_string()),
70 }
71 }
72}
73
74impl From<LogFormat> for crate::OutputFormat {
75 fn from(value: LogFormat) -> Self {
76 match value {
77 LogFormat::Json => Self::Json,
78 LogFormat::Plain => Self::Plain,
79 LogFormat::Yaml => Self::Yaml,
80 }
81 }
82}
83
84trait LogSink: Send + Sync {
85 fn write_line(&self, line: &str, metadata: Option<&tracing::Metadata<'_>>) -> io::Result<()>;
86}
87
88struct SharedSink {
95 inner: RwLock<Arc<dyn LogSink>>,
96}
97
98impl SharedSink {
99 fn new(sink: Arc<dyn LogSink>) -> Self {
100 Self {
101 inner: RwLock::new(sink),
102 }
103 }
104
105 fn replace(&self, sink: Arc<dyn LogSink>) {
106 match self.inner.write() {
107 Ok(mut current) => *current = sink,
108 Err(poisoned) => *poisoned.into_inner() = sink,
109 }
110 }
111}
112
113impl LogSink for SharedSink {
114 fn write_line(&self, line: &str, metadata: Option<&tracing::Metadata<'_>>) -> io::Result<()> {
115 let sink = self
116 .inner
117 .read()
118 .map_err(|_| io::Error::other("AFDATA log sink lock is poisoned"))?;
119 sink.write_line(line, metadata)
120 }
121}
122
123struct MakeWriterSink<W> {
124 writer: Mutex<W>,
125}
126
127impl<W> LogSink for MakeWriterSink<W>
128where
129 W: Send + for<'writer> MakeWriter<'writer>,
130 for<'writer> <W as MakeWriter<'writer>>::Writer: Write,
131{
132 fn write_line(&self, line: &str, metadata: Option<&tracing::Metadata<'_>>) -> io::Result<()> {
133 let factory = self
134 .writer
135 .lock()
136 .map_err(|_| io::Error::other("AFDATA log writer lock is poisoned"))?;
137 let mut writer = match metadata {
138 Some(metadata) => factory.make_writer_for(metadata),
139 None => factory.make_writer(),
140 };
141 let mut framed = String::with_capacity(line.len() + 1);
148 framed.push_str(line);
149 framed.push('\n');
150 writer.write_all(framed.as_bytes())?;
151 writer.flush()
152 }
153}
154
155pub struct AfdataLayer {
157 format: LogFormat,
158 redactor: crate::Redactor,
159 sink: Arc<SharedSink>,
160}
161
162impl std::fmt::Debug for AfdataLayer {
163 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
164 formatter
165 .debug_struct("AfdataLayer")
166 .field("format", &self.format)
167 .field("redactor", &self.redactor)
168 .finish_non_exhaustive()
169 }
170}
171
172#[derive(Clone)]
174pub struct StructuredLogHandle {
175 format: LogFormat,
176 redactor: crate::Redactor,
177 sink: Arc<SharedSink>,
178}
179
180impl std::fmt::Debug for StructuredLogHandle {
181 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
182 formatter
183 .debug_struct("StructuredLogHandle")
184 .field("format", &self.format)
185 .field("redactor", &self.redactor)
186 .finish_non_exhaustive()
187 }
188}
189
190pub fn try_init(
202 filter: tracing_subscriber::EnvFilter,
203 format: LogFormat,
204 redactor: crate::Redactor,
205) -> Result<(), TryInitError> {
206 use tracing_subscriber::layer::SubscriberExt;
207 use tracing_subscriber::util::SubscriberInitExt;
208
209 tracing_subscriber::registry()
210 .with(filter)
211 .with(AfdataLayer::new(format, redactor))
212 .try_init()
213}
214
215impl AfdataLayer {
216 #[allow(clippy::disallowed_methods)]
221 pub fn new(format: LogFormat, redactor: crate::Redactor) -> Self {
222 Self {
223 format,
224 redactor,
225 sink: make_shared_sink(io::stderr),
226 }
227 }
228
229 pub fn with_writer<W>(self, writer: W) -> Self
237 where
238 W: Send + for<'writer> MakeWriter<'writer> + 'static,
239 for<'writer> <W as MakeWriter<'writer>>::Writer: Write,
240 {
241 self.sink.replace(make_sink(writer));
242 self
243 }
244
245 pub fn structured_log_handle(&self) -> StructuredLogHandle {
247 StructuredLogHandle {
248 format: self.format,
249 redactor: self.redactor.clone(),
250 sink: Arc::clone(&self.sink),
251 }
252 }
253
254 fn output_options(&self) -> crate::OutputOptions {
255 crate::OutputOptions {
256 redaction: self.redactor.clone(),
257 style: crate::PlainStyle::Readable,
258 }
259 }
260
261 fn format_value(&self, value: &serde_json::Value) -> String {
262 crate::render(value, self.format.into(), &self.output_options())
263 }
264}
265
266fn make_sink<W>(writer: W) -> Arc<dyn LogSink>
267where
268 W: Send + for<'writer> MakeWriter<'writer> + 'static,
269 for<'writer> <W as MakeWriter<'writer>>::Writer: Write,
270{
271 Arc::new(MakeWriterSink {
272 writer: Mutex::new(writer),
273 })
274}
275
276fn make_shared_sink<W>(writer: W) -> Arc<SharedSink>
277where
278 W: Send + for<'writer> MakeWriter<'writer> + 'static,
279 for<'writer> <W as MakeWriter<'writer>>::Writer: Write,
280{
281 Arc::new(SharedSink::new(make_sink(writer)))
282}
283
284impl StructuredLogHandle {
285 pub fn emit(&self, payload: serde_json::Value) -> io::Result<()> {
300 let event = crate::json_log(stamp_log_metadata(payload)).build();
301 let options = crate::OutputOptions {
302 redaction: self.redactor.clone(),
303 style: crate::PlainStyle::Readable,
304 };
305 let line = crate::render(event.as_value(), self.format.into(), &options);
306 self.sink.write_line(&line, None)
307 }
308}
309
310struct SpanFields(Vec<(String, serde_json::Value)>);
312
313impl<S> Layer<S> for AfdataLayer
314where
315 S: Subscriber + for<'a> LookupSpan<'a>,
316{
317 fn on_new_span(&self, attrs: &span::Attributes<'_>, id: &span::Id, ctx: Context<'_, S>) {
318 let mut visitor = JsonVisitor::new();
319 attrs.record(&mut visitor);
320
321 if let Some(span) = ctx.span(id) {
322 span.extensions_mut().insert(SpanFields(visitor.fields));
323 }
324 }
325
326 fn on_record(&self, id: &span::Id, values: &span::Record<'_>, ctx: Context<'_, S>) {
327 if let Some(span) = ctx.span(id) {
328 let mut visitor = JsonVisitor::new();
329 values.record(&mut visitor);
330
331 let mut extensions = span.extensions_mut();
332 if let Some(existing) = extensions.get_mut::<SpanFields>() {
333 existing.0.extend(visitor.fields);
334 } else {
335 extensions.insert(SpanFields(visitor.fields));
336 }
337 }
338 }
339
340 fn on_event(&self, event: &Event<'_>, ctx: Context<'_, S>) {
341 let meta = event.metadata();
342
343 let mut visitor = JsonVisitor::new();
345 event.record(&mut visitor);
346
347 let mut map = serde_json::Map::with_capacity(4 + visitor.fields.len());
349
350 let level = match *meta.level() {
351 Level::TRACE | Level::DEBUG => crate::LogLevel::Debug,
354 Level::INFO => crate::LogLevel::Info,
355 Level::WARN => crate::LogLevel::Warn,
356 Level::ERROR => crate::LogLevel::Error,
357 };
358
359 let message = visitor
361 .message
362 .take()
363 .unwrap_or_else(|| "(no message)".to_string());
364
365 if let Some(scope) = ctx.event_scope(event) {
367 for span in scope.from_root() {
368 let extensions = span.extensions();
369 if let Some(fields) = extensions.get::<SpanFields>() {
370 for (k, v) in &fields.0 {
371 if !is_reserved_log_field(k) {
372 map.insert(k.clone(), v.clone());
373 }
374 }
375 }
376 }
377 }
378
379 for (k, v) in visitor.fields {
383 if !is_reserved_log_field(&k) {
384 map.insert(k, v);
385 }
386 }
387
388 map.insert(
390 "timestamp_epoch_ms".into(),
391 serde_json::Value::Number(chrono::Utc::now().timestamp_millis().into()),
392 );
393 map.insert(
394 "level".to_string(),
395 serde_json::Value::String(level.as_str().to_string()),
396 );
397 map.insert("message".to_string(), serde_json::Value::String(message));
398 let builder = crate::json_log(serde_json::Value::Object(map));
399 let value = builder.build();
400
401 let line = self.format_value(value.as_value());
403
404 let _ = self.sink.write_line(&line, Some(meta));
405 }
406}
407
408fn is_reserved_log_field(field: &str) -> bool {
409 matches!(field, "level" | "message" | "timestamp_epoch_ms")
410}
411
412fn stamp_log_metadata(payload: serde_json::Value) -> serde_json::Value {
418 let serde_json::Value::Object(mut map) = payload else {
419 return payload;
420 };
421 map.entry("level".to_string())
422 .or_insert_with(|| serde_json::Value::String("info".to_string()));
423 map.insert(
424 "timestamp_epoch_ms".to_string(),
425 serde_json::Value::Number(chrono::Utc::now().timestamp_millis().into()),
426 );
427 serde_json::Value::Object(map)
428}
429
430struct JsonVisitor {
432 message: Option<String>,
433 fields: Vec<(String, serde_json::Value)>,
434}
435
436impl JsonVisitor {
437 fn new() -> Self {
438 Self {
439 message: None,
440 fields: Vec::new(),
441 }
442 }
443}
444
445impl Visit for JsonVisitor {
446 fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
447 let val = format!("{:?}", value);
448 if field.name() == "message" {
449 self.message = Some(val);
450 } else {
451 self.fields
457 .push((field.name().to_string(), serde_json::Value::String(val)));
458 }
459 }
460
461 fn record_str(&mut self, field: &Field, value: &str) {
462 if field.name() == "message" {
463 self.message = Some(value.to_string());
464 } else {
465 self.fields.push((
466 field.name().to_string(),
467 serde_json::Value::String(value.to_string()),
468 ));
469 }
470 }
471
472 fn record_i64(&mut self, field: &Field, value: i64) {
473 self.fields.push((
474 field.name().to_string(),
475 serde_json::Value::Number(value.into()),
476 ));
477 }
478
479 fn record_u64(&mut self, field: &Field, value: u64) {
480 self.fields.push((
481 field.name().to_string(),
482 serde_json::Value::Number(value.into()),
483 ));
484 }
485
486 fn record_f64(&mut self, field: &Field, value: f64) {
487 if let Some(n) = serde_json::Number::from_f64(value) {
488 self.fields
489 .push((field.name().to_string(), serde_json::Value::Number(n)));
490 } else {
491 self.fields.push((
492 field.name().to_string(),
493 serde_json::Value::String(value.to_string()),
494 ));
495 }
496 }
497
498 fn record_bool(&mut self, field: &Field, value: bool) {
499 self.fields
500 .push((field.name().to_string(), serde_json::Value::Bool(value)));
501 }
502}
503
504#[cfg(test)]
505mod tests {
506 use super::*;
507 use serde_json::json;
508
509 #[test]
515 fn code_field_is_accepted_by_log_builder() {
516 let value = crate::json_log(json!({"code": "cache_miss"})).build();
517 assert_eq!(value.as_value()["log"]["code"], "cache_miss");
518 }
519
520 #[test]
521 fn secret_named_field_is_redacted_at_emit() {
522 let line = crate::render(
523 &json!({
524 "code": "info",
525 "api_key_secret": "sk-live-123",
526 }),
527 crate::OutputFormat::Json,
528 &crate::OutputOptions::default(),
529 );
530 assert!(line.contains("\"api_key_secret\":\"***\""), "{line}");
531 assert!(!line.contains("sk-live-123"), "{line}");
532 }
533
534 #[test]
535 fn non_secret_field_whose_value_mentions_secret_is_not_redacted() {
536 let line = crate::render(
539 &json!({
540 "code": "info",
541 "note": "see the api_key_secret field in docs",
542 }),
543 crate::OutputFormat::Json,
544 &crate::OutputOptions::default(),
545 );
546 assert!(
547 line.contains("see the api_key_secret field in docs"),
548 "{line}"
549 );
550 }
551
552 #[test]
553 fn secret_typed_field_is_redacted_regardless_of_record_path() {
554 let line = crate::render(
557 &json!({
558 "code": "warn",
559 "db_password_secret": 1234,
560 }),
561 crate::OutputFormat::Json,
562 &crate::OutputOptions::default(),
563 );
564 assert!(line.contains("\"db_password_secret\":\"***\""), "{line}");
565 }
566
567 #[test]
568 fn legacy_secret_names_are_redacted_when_layer_has_options() {
569 let value = crate::json_log(json!({
570 "level": "info",
571 "message": "authorization appears in message but is not name-redacted",
572 "timestamp_epoch_ms": 1,
573 "authorization": "Bearer legacy",
574 "request_url": "https://example.test/path?authorization=legacy&ok=1",
575 }))
576 .build();
577 let redactor = crate::Redactor::new().secret_names(vec!["authorization".to_string()]);
578
579 let formats = [LogFormat::Json, LogFormat::Plain, LogFormat::Yaml];
580
581 for format in formats {
582 let layer = AfdataLayer::new(format, redactor.clone());
583 let line = layer.format_value(value.as_value());
584 assert!(line.contains("***"), "{line}");
585 assert!(
586 !line.contains("Bearer legacy"),
587 "legacy field value should be redacted: {line}"
588 );
589 assert!(
590 !line.contains("authorization=legacy"),
591 "legacy URL query parameter should be redacted: {line}"
592 );
593 assert!(
594 line.contains("authorization appears in message"),
595 "message is free-form and should remain readable: {line}"
596 );
597 }
598 }
599
600 #[test]
601 fn legacy_secret_names_are_visible_without_layer_options() {
602 let value = crate::json_log(json!({
603 "level": "info",
604 "message": "ready",
605 "timestamp_epoch_ms": 1,
606 "authorization": "Bearer visible",
607 }))
608 .build();
609 let layer = AfdataLayer::new(LogFormat::Json, crate::Redactor::new());
610
611 let line = layer.format_value(value.as_value());
612 assert!(
613 line.contains("\"authorization\":\"Bearer visible\""),
614 "{line}"
615 );
616 }
617
618 #[test]
619 fn log_format_round_trips_text_and_serde() {
620 for format in [LogFormat::Json, LogFormat::Plain, LogFormat::Yaml] {
621 assert_eq!(format.to_string().parse(), Ok(format));
622 let encoded = serde_json::to_string(&format).unwrap_or_default();
623 assert_eq!(
624 serde_json::from_str::<LogFormat>(&encoded).ok(),
625 Some(format)
626 );
627 }
628
629 let canary = "canary-log-format-secret";
630 let error = canary.parse::<LogFormat>().unwrap_err();
631 assert!(!error.contains(canary));
632 assert!(error.contains("json"));
633 }
634
635 #[derive(Clone)]
636 struct MemoryMakeWriter {
637 bytes: Arc<Mutex<Vec<u8>>>,
638 }
639
640 struct MemoryWriter {
641 bytes: Arc<Mutex<Vec<u8>>>,
642 }
643
644 impl Write for MemoryWriter {
645 fn write(&mut self, buffer: &[u8]) -> io::Result<usize> {
646 let mut bytes = self
647 .bytes
648 .lock()
649 .map_err(|_| io::Error::other("test buffer lock poisoned"))?;
650 bytes.extend_from_slice(buffer);
651 Ok(buffer.len())
652 }
653
654 fn flush(&mut self) -> io::Result<()> {
655 Ok(())
656 }
657 }
658
659 impl<'writer> MakeWriter<'writer> for MemoryMakeWriter {
660 type Writer = MemoryWriter;
661
662 fn make_writer(&'writer self) -> Self::Writer {
663 MemoryWriter {
664 bytes: Arc::clone(&self.bytes),
665 }
666 }
667 }
668
669 #[test]
670 fn structured_handle_and_tracing_share_writer_and_keep_nested_json() {
671 use tracing_subscriber::layer::SubscriberExt;
672
673 let bytes = Arc::new(Mutex::new(Vec::new()));
674 let layer = AfdataLayer::new(LogFormat::Json, crate::Redactor::new()).with_writer(
675 MemoryMakeWriter {
676 bytes: Arc::clone(&bytes),
677 },
678 );
679 let structured = layer.structured_log_handle();
680 structured
681 .emit(json!({
682 "level": "info",
683 "message": "configuration loaded",
684 "configuration": {
685 "region": "test",
686 "credential_secret": "do-not-log"
687 }
688 }))
689 .unwrap_or_else(|error| panic!("{error}"));
690
691 let subscriber = tracing_subscriber::registry().with(layer);
692 tracing::subscriber::with_default(subscriber, || {
693 tracing::info!(request_id = 7_u64, "ordinary event");
694 });
695
696 let output = {
697 let bytes = bytes
698 .lock()
699 .unwrap_or_else(std::sync::PoisonError::into_inner);
700 String::from_utf8(bytes.clone()).unwrap_or_else(|error| panic!("{error}"))
701 };
702 let lines: Vec<&str> = output.lines().collect();
703 assert_eq!(lines.len(), 2, "{output}");
704 let structured_value: serde_json::Value =
705 serde_json::from_str(lines[0]).unwrap_or_else(|error| panic!("{error}"));
706 assert_eq!(
707 structured_value["log"]["configuration"]["region"],
708 serde_json::json!("test")
709 );
710 assert_eq!(
711 structured_value["log"]["configuration"]["credential_secret"],
712 serde_json::json!("***")
713 );
714 assert!(lines[1].contains("ordinary event"), "{output}");
715 }
716
717 #[test]
721 fn handle_taken_before_with_writer_follows_the_new_writer() {
722 let bytes = Arc::new(Mutex::new(Vec::new()));
723 let layer = AfdataLayer::new(LogFormat::Json, crate::Redactor::new());
724 let structured = layer.structured_log_handle();
726 let _layer = layer.with_writer(MemoryMakeWriter {
727 bytes: Arc::clone(&bytes),
728 });
729
730 structured
731 .emit(json!({"message": "after rewiring"}))
732 .unwrap_or_else(|error| panic!("{error}"));
733
734 let output = {
735 let bytes = bytes
736 .lock()
737 .unwrap_or_else(std::sync::PoisonError::into_inner);
738 String::from_utf8(bytes.clone()).unwrap_or_else(|error| panic!("{error}"))
739 };
740 assert!(
741 output.contains("after rewiring"),
742 "handle wrote somewhere else entirely: {output:?}"
743 );
744 }
745
746 #[test]
749 fn directly_emitted_events_carry_level_and_timestamp() {
750 let bytes = Arc::new(Mutex::new(Vec::new()));
751 let layer = AfdataLayer::new(LogFormat::Json, crate::Redactor::new()).with_writer(
752 MemoryMakeWriter {
753 bytes: Arc::clone(&bytes),
754 },
755 );
756 let structured = layer.structured_log_handle();
757 structured
758 .emit(json!({"message": "defaulted"}))
759 .unwrap_or_else(|error| panic!("{error}"));
760 structured
761 .emit(json!({"level": "warn", "message": "explicit"}))
762 .unwrap_or_else(|error| panic!("{error}"));
763
764 let output = {
765 let bytes = bytes
766 .lock()
767 .unwrap_or_else(std::sync::PoisonError::into_inner);
768 String::from_utf8(bytes.clone()).unwrap_or_else(|error| panic!("{error}"))
769 };
770 let lines: Vec<&str> = output.lines().collect();
771 assert_eq!(lines.len(), 2, "{output}");
772 for (line, expected_level) in lines.iter().zip(["info", "warn"]) {
773 let value: serde_json::Value =
774 serde_json::from_str(line).unwrap_or_else(|error| panic!("{error}"));
775 assert_eq!(value["log"]["level"], serde_json::json!(expected_level));
776 assert!(
777 value["log"]["timestamp_epoch_ms"].is_number(),
778 "missing timestamp: {line}"
779 );
780 }
781 }
782
783 #[test]
784 fn trace_maps_to_debug_and_fields_cannot_override_metadata() {
785 use tracing_subscriber::layer::SubscriberExt;
786
787 let bytes = Arc::new(Mutex::new(Vec::new()));
788 let layer = AfdataLayer::new(LogFormat::Json, crate::Redactor::new()).with_writer(
789 MemoryMakeWriter {
790 bytes: Arc::clone(&bytes),
791 },
792 );
793 let subscriber = tracing_subscriber::registry().with(layer);
794 tracing::subscriber::with_default(subscriber, || {
795 tracing::trace!(level = "error", timestamp_epoch_ms = 1_i64, "trace event");
796 });
797
798 let output = {
799 let bytes = bytes
800 .lock()
801 .unwrap_or_else(std::sync::PoisonError::into_inner);
802 String::from_utf8(bytes.clone()).unwrap_or_else(|error| panic!("{error}"))
803 };
804 let value: serde_json::Value =
805 serde_json::from_str(output.trim()).unwrap_or_else(|error| panic!("{error}"));
806 assert_eq!(value["log"]["level"], "debug");
807 assert_eq!(value["log"]["message"], "trace event");
808 assert_ne!(value["log"]["timestamp_epoch_ms"], 1);
809 }
810}