use std::io::{self, Write};
use std::str::FromStr;
use std::sync::{Arc, Mutex, RwLock};
use serde::{Deserialize, Serialize};
use tracing::field::{Field, Visit};
use tracing::span;
use tracing::{Event, Level, Subscriber};
use tracing_subscriber::Layer;
use tracing_subscriber::fmt::MakeWriter;
use tracing_subscriber::layer::Context;
use tracing_subscriber::registry::LookupSpan;
use tracing_subscriber::util::TryInitError;
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum LogFormat {
Json,
Plain,
Yaml,
}
impl LogFormat {
pub const fn as_str(self) -> &'static str {
match self {
Self::Json => "json",
Self::Plain => "plain",
Self::Yaml => "yaml",
}
}
}
impl std::fmt::Display for LogFormat {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str(self.as_str())
}
}
impl FromStr for LogFormat {
type Err = String;
fn from_str(value: &str) -> Result<Self, Self::Err> {
match value {
"json" => Ok(Self::Json),
"plain" => Ok(Self::Plain),
"yaml" => Ok(Self::Yaml),
_ => Err("invalid log format: expected json, plain, or yaml".to_string()),
}
}
}
impl From<LogFormat> for crate::OutputFormat {
fn from(value: LogFormat) -> Self {
match value {
LogFormat::Json => Self::Json,
LogFormat::Plain => Self::Plain,
LogFormat::Yaml => Self::Yaml,
}
}
}
trait LogSink: Send + Sync {
fn write_line(&self, line: &str, metadata: Option<&tracing::Metadata<'_>>) -> io::Result<()>;
}
struct SharedSink {
inner: RwLock<Arc<dyn LogSink>>,
}
impl SharedSink {
fn new(sink: Arc<dyn LogSink>) -> Self {
Self {
inner: RwLock::new(sink),
}
}
fn replace(&self, sink: Arc<dyn LogSink>) {
match self.inner.write() {
Ok(mut current) => *current = sink,
Err(poisoned) => *poisoned.into_inner() = sink,
}
}
}
impl LogSink for SharedSink {
fn write_line(&self, line: &str, metadata: Option<&tracing::Metadata<'_>>) -> io::Result<()> {
let sink = self
.inner
.read()
.map_err(|_| io::Error::other("AFDATA log sink lock is poisoned"))?;
sink.write_line(line, metadata)
}
}
struct MakeWriterSink<W> {
writer: Mutex<W>,
}
impl<W> LogSink for MakeWriterSink<W>
where
W: Send + for<'writer> MakeWriter<'writer>,
for<'writer> <W as MakeWriter<'writer>>::Writer: Write,
{
fn write_line(&self, line: &str, metadata: Option<&tracing::Metadata<'_>>) -> io::Result<()> {
let factory = self
.writer
.lock()
.map_err(|_| io::Error::other("AFDATA log writer lock is poisoned"))?;
let mut writer = match metadata {
Some(metadata) => factory.make_writer_for(metadata),
None => factory.make_writer(),
};
let mut framed = String::with_capacity(line.len() + 1);
framed.push_str(line);
framed.push('\n');
writer.write_all(framed.as_bytes())?;
writer.flush()
}
}
pub struct AfdataLayer {
format: LogFormat,
redactor: crate::Redactor,
sink: Arc<SharedSink>,
}
impl std::fmt::Debug for AfdataLayer {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("AfdataLayer")
.field("format", &self.format)
.field("redactor", &self.redactor)
.finish_non_exhaustive()
}
}
#[derive(Clone)]
pub struct StructuredLogHandle {
format: LogFormat,
redactor: crate::Redactor,
sink: Arc<SharedSink>,
}
impl std::fmt::Debug for StructuredLogHandle {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("StructuredLogHandle")
.field("format", &self.format)
.field("redactor", &self.redactor)
.finish_non_exhaustive()
}
}
pub fn try_init(
filter: tracing_subscriber::EnvFilter,
format: LogFormat,
redactor: crate::Redactor,
) -> Result<(), TryInitError> {
use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::util::SubscriberInitExt;
tracing_subscriber::registry()
.with(filter)
.with(AfdataLayer::new(format, redactor))
.try_init()
}
impl AfdataLayer {
#[allow(clippy::disallowed_methods)]
pub fn new(format: LogFormat, redactor: crate::Redactor) -> Self {
Self {
format,
redactor,
sink: make_shared_sink(io::stderr),
}
}
pub fn with_writer<W>(self, writer: W) -> Self
where
W: Send + for<'writer> MakeWriter<'writer> + 'static,
for<'writer> <W as MakeWriter<'writer>>::Writer: Write,
{
self.sink.replace(make_sink(writer));
self
}
pub fn structured_log_handle(&self) -> StructuredLogHandle {
StructuredLogHandle {
format: self.format,
redactor: self.redactor.clone(),
sink: Arc::clone(&self.sink),
}
}
fn output_options(&self) -> crate::OutputOptions {
crate::OutputOptions {
redaction: self.redactor.clone(),
style: crate::PlainStyle::Readable,
}
}
fn format_value(&self, value: &serde_json::Value) -> String {
crate::render(value, self.format.into(), &self.output_options())
}
}
fn make_sink<W>(writer: W) -> Arc<dyn LogSink>
where
W: Send + for<'writer> MakeWriter<'writer> + 'static,
for<'writer> <W as MakeWriter<'writer>>::Writer: Write,
{
Arc::new(MakeWriterSink {
writer: Mutex::new(writer),
})
}
fn make_shared_sink<W>(writer: W) -> Arc<SharedSink>
where
W: Send + for<'writer> MakeWriter<'writer> + 'static,
for<'writer> <W as MakeWriter<'writer>>::Writer: Write,
{
Arc::new(SharedSink::new(make_sink(writer)))
}
impl StructuredLogHandle {
pub fn emit(&self, payload: serde_json::Value) -> io::Result<()> {
let event = crate::json_log(stamp_log_metadata(payload)).build();
let options = crate::OutputOptions {
redaction: self.redactor.clone(),
style: crate::PlainStyle::Readable,
};
let line = crate::render(event.as_value(), self.format.into(), &options);
self.sink.write_line(&line, None)
}
}
struct SpanFields(Vec<(String, serde_json::Value)>);
impl<S> Layer<S> for AfdataLayer
where
S: Subscriber + for<'a> LookupSpan<'a>,
{
fn on_new_span(&self, attrs: &span::Attributes<'_>, id: &span::Id, ctx: Context<'_, S>) {
let mut visitor = JsonVisitor::new();
attrs.record(&mut visitor);
if let Some(span) = ctx.span(id) {
span.extensions_mut().insert(SpanFields(visitor.fields));
}
}
fn on_record(&self, id: &span::Id, values: &span::Record<'_>, ctx: Context<'_, S>) {
if let Some(span) = ctx.span(id) {
let mut visitor = JsonVisitor::new();
values.record(&mut visitor);
let mut extensions = span.extensions_mut();
if let Some(existing) = extensions.get_mut::<SpanFields>() {
existing.0.extend(visitor.fields);
} else {
extensions.insert(SpanFields(visitor.fields));
}
}
}
fn on_event(&self, event: &Event<'_>, ctx: Context<'_, S>) {
let meta = event.metadata();
let mut visitor = JsonVisitor::new();
event.record(&mut visitor);
let mut map = serde_json::Map::with_capacity(4 + visitor.fields.len());
let level = match *meta.level() {
Level::TRACE | Level::DEBUG => crate::LogLevel::Debug,
Level::INFO => crate::LogLevel::Info,
Level::WARN => crate::LogLevel::Warn,
Level::ERROR => crate::LogLevel::Error,
};
let message = visitor
.message
.take()
.unwrap_or_else(|| "(no message)".to_string());
if let Some(scope) = ctx.event_scope(event) {
for span in scope.from_root() {
let extensions = span.extensions();
if let Some(fields) = extensions.get::<SpanFields>() {
for (k, v) in &fields.0 {
if !is_reserved_log_field(k) {
map.insert(k.clone(), v.clone());
}
}
}
}
}
for (k, v) in visitor.fields {
if !is_reserved_log_field(&k) {
map.insert(k, v);
}
}
map.insert(
"timestamp_epoch_ms".into(),
serde_json::Value::Number(chrono::Utc::now().timestamp_millis().into()),
);
map.insert(
"level".to_string(),
serde_json::Value::String(level.as_str().to_string()),
);
map.insert("message".to_string(), serde_json::Value::String(message));
let builder = crate::json_log(serde_json::Value::Object(map));
let value = builder.build();
let line = self.format_value(value.as_value());
let _ = self.sink.write_line(&line, Some(meta));
}
}
fn is_reserved_log_field(field: &str) -> bool {
matches!(field, "level" | "message" | "timestamp_epoch_ms")
}
fn stamp_log_metadata(payload: serde_json::Value) -> serde_json::Value {
let serde_json::Value::Object(mut map) = payload else {
return payload;
};
map.entry("level".to_string())
.or_insert_with(|| serde_json::Value::String("info".to_string()));
map.insert(
"timestamp_epoch_ms".to_string(),
serde_json::Value::Number(chrono::Utc::now().timestamp_millis().into()),
);
serde_json::Value::Object(map)
}
struct JsonVisitor {
message: Option<String>,
fields: Vec<(String, serde_json::Value)>,
}
impl JsonVisitor {
fn new() -> Self {
Self {
message: None,
fields: Vec::new(),
}
}
}
impl Visit for JsonVisitor {
fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
let val = format!("{:?}", value);
if field.name() == "message" {
self.message = Some(val);
} else {
self.fields
.push((field.name().to_string(), serde_json::Value::String(val)));
}
}
fn record_str(&mut self, field: &Field, value: &str) {
if field.name() == "message" {
self.message = Some(value.to_string());
} else {
self.fields.push((
field.name().to_string(),
serde_json::Value::String(value.to_string()),
));
}
}
fn record_i64(&mut self, field: &Field, value: i64) {
self.fields.push((
field.name().to_string(),
serde_json::Value::Number(value.into()),
));
}
fn record_u64(&mut self, field: &Field, value: u64) {
self.fields.push((
field.name().to_string(),
serde_json::Value::Number(value.into()),
));
}
fn record_f64(&mut self, field: &Field, value: f64) {
if let Some(n) = serde_json::Number::from_f64(value) {
self.fields
.push((field.name().to_string(), serde_json::Value::Number(n)));
} else {
self.fields.push((
field.name().to_string(),
serde_json::Value::String(value.to_string()),
));
}
}
fn record_bool(&mut self, field: &Field, value: bool) {
self.fields
.push((field.name().to_string(), serde_json::Value::Bool(value)));
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn code_field_is_accepted_by_log_builder() {
let value = crate::json_log(json!({"code": "cache_miss"})).build();
assert_eq!(value.as_value()["log"]["code"], "cache_miss");
}
#[test]
fn secret_named_field_is_redacted_at_emit() {
let line = crate::render(
&json!({
"code": "info",
"api_key_secret": "sk-live-123",
}),
crate::OutputFormat::Json,
&crate::OutputOptions::default(),
);
assert!(line.contains("\"api_key_secret\":\"***\""), "{line}");
assert!(!line.contains("sk-live-123"), "{line}");
}
#[test]
fn non_secret_field_whose_value_mentions_secret_is_not_redacted() {
let line = crate::render(
&json!({
"code": "info",
"note": "see the api_key_secret field in docs",
}),
crate::OutputFormat::Json,
&crate::OutputOptions::default(),
);
assert!(
line.contains("see the api_key_secret field in docs"),
"{line}"
);
}
#[test]
fn secret_typed_field_is_redacted_regardless_of_record_path() {
let line = crate::render(
&json!({
"code": "warn",
"db_password_secret": 1234,
}),
crate::OutputFormat::Json,
&crate::OutputOptions::default(),
);
assert!(line.contains("\"db_password_secret\":\"***\""), "{line}");
}
#[test]
fn legacy_secret_names_are_redacted_when_layer_has_options() {
let value = crate::json_log(json!({
"level": "info",
"message": "authorization appears in message but is not name-redacted",
"timestamp_epoch_ms": 1,
"authorization": "Bearer legacy",
"request_url": "https://example.test/path?authorization=legacy&ok=1",
}))
.build();
let redactor = crate::Redactor::new().secret_names(vec!["authorization".to_string()]);
let formats = [LogFormat::Json, LogFormat::Plain, LogFormat::Yaml];
for format in formats {
let layer = AfdataLayer::new(format, redactor.clone());
let line = layer.format_value(value.as_value());
assert!(line.contains("***"), "{line}");
assert!(
!line.contains("Bearer legacy"),
"legacy field value should be redacted: {line}"
);
assert!(
!line.contains("authorization=legacy"),
"legacy URL query parameter should be redacted: {line}"
);
assert!(
line.contains("authorization appears in message"),
"message is free-form and should remain readable: {line}"
);
}
}
#[test]
fn legacy_secret_names_are_visible_without_layer_options() {
let value = crate::json_log(json!({
"level": "info",
"message": "ready",
"timestamp_epoch_ms": 1,
"authorization": "Bearer visible",
}))
.build();
let layer = AfdataLayer::new(LogFormat::Json, crate::Redactor::new());
let line = layer.format_value(value.as_value());
assert!(
line.contains("\"authorization\":\"Bearer visible\""),
"{line}"
);
}
#[test]
fn log_format_round_trips_text_and_serde() {
for format in [LogFormat::Json, LogFormat::Plain, LogFormat::Yaml] {
assert_eq!(format.to_string().parse(), Ok(format));
let encoded = serde_json::to_string(&format).unwrap_or_default();
assert_eq!(
serde_json::from_str::<LogFormat>(&encoded).ok(),
Some(format)
);
}
let canary = "canary-log-format-secret";
let error = canary.parse::<LogFormat>().unwrap_err();
assert!(!error.contains(canary));
assert!(error.contains("json"));
}
#[derive(Clone)]
struct MemoryMakeWriter {
bytes: Arc<Mutex<Vec<u8>>>,
}
struct MemoryWriter {
bytes: Arc<Mutex<Vec<u8>>>,
}
impl Write for MemoryWriter {
fn write(&mut self, buffer: &[u8]) -> io::Result<usize> {
let mut bytes = self
.bytes
.lock()
.map_err(|_| io::Error::other("test buffer lock poisoned"))?;
bytes.extend_from_slice(buffer);
Ok(buffer.len())
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
impl<'writer> MakeWriter<'writer> for MemoryMakeWriter {
type Writer = MemoryWriter;
fn make_writer(&'writer self) -> Self::Writer {
MemoryWriter {
bytes: Arc::clone(&self.bytes),
}
}
}
#[test]
fn structured_handle_and_tracing_share_writer_and_keep_nested_json() {
use tracing_subscriber::layer::SubscriberExt;
let bytes = Arc::new(Mutex::new(Vec::new()));
let layer = AfdataLayer::new(LogFormat::Json, crate::Redactor::new()).with_writer(
MemoryMakeWriter {
bytes: Arc::clone(&bytes),
},
);
let structured = layer.structured_log_handle();
structured
.emit(json!({
"level": "info",
"message": "configuration loaded",
"configuration": {
"region": "test",
"credential_secret": "do-not-log"
}
}))
.unwrap_or_else(|error| panic!("{error}"));
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, || {
tracing::info!(request_id = 7_u64, "ordinary event");
});
let output = {
let bytes = bytes
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
String::from_utf8(bytes.clone()).unwrap_or_else(|error| panic!("{error}"))
};
let lines: Vec<&str> = output.lines().collect();
assert_eq!(lines.len(), 2, "{output}");
let structured_value: serde_json::Value =
serde_json::from_str(lines[0]).unwrap_or_else(|error| panic!("{error}"));
assert_eq!(
structured_value["log"]["configuration"]["region"],
serde_json::json!("test")
);
assert_eq!(
structured_value["log"]["configuration"]["credential_secret"],
serde_json::json!("***")
);
assert!(lines[1].contains("ordinary event"), "{output}");
}
#[test]
fn handle_taken_before_with_writer_follows_the_new_writer() {
let bytes = Arc::new(Mutex::new(Vec::new()));
let layer = AfdataLayer::new(LogFormat::Json, crate::Redactor::new());
let structured = layer.structured_log_handle();
let _layer = layer.with_writer(MemoryMakeWriter {
bytes: Arc::clone(&bytes),
});
structured
.emit(json!({"message": "after rewiring"}))
.unwrap_or_else(|error| panic!("{error}"));
let output = {
let bytes = bytes
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
String::from_utf8(bytes.clone()).unwrap_or_else(|error| panic!("{error}"))
};
assert!(
output.contains("after rewiring"),
"handle wrote somewhere else entirely: {output:?}"
);
}
#[test]
fn directly_emitted_events_carry_level_and_timestamp() {
let bytes = Arc::new(Mutex::new(Vec::new()));
let layer = AfdataLayer::new(LogFormat::Json, crate::Redactor::new()).with_writer(
MemoryMakeWriter {
bytes: Arc::clone(&bytes),
},
);
let structured = layer.structured_log_handle();
structured
.emit(json!({"message": "defaulted"}))
.unwrap_or_else(|error| panic!("{error}"));
structured
.emit(json!({"level": "warn", "message": "explicit"}))
.unwrap_or_else(|error| panic!("{error}"));
let output = {
let bytes = bytes
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
String::from_utf8(bytes.clone()).unwrap_or_else(|error| panic!("{error}"))
};
let lines: Vec<&str> = output.lines().collect();
assert_eq!(lines.len(), 2, "{output}");
for (line, expected_level) in lines.iter().zip(["info", "warn"]) {
let value: serde_json::Value =
serde_json::from_str(line).unwrap_or_else(|error| panic!("{error}"));
assert_eq!(value["log"]["level"], serde_json::json!(expected_level));
assert!(
value["log"]["timestamp_epoch_ms"].is_number(),
"missing timestamp: {line}"
);
}
}
#[test]
fn trace_maps_to_debug_and_fields_cannot_override_metadata() {
use tracing_subscriber::layer::SubscriberExt;
let bytes = Arc::new(Mutex::new(Vec::new()));
let layer = AfdataLayer::new(LogFormat::Json, crate::Redactor::new()).with_writer(
MemoryMakeWriter {
bytes: Arc::clone(&bytes),
},
);
let subscriber = tracing_subscriber::registry().with(layer);
tracing::subscriber::with_default(subscriber, || {
tracing::trace!(level = "error", timestamp_epoch_ms = 1_i64, "trace event");
});
let output = {
let bytes = bytes
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
String::from_utf8(bytes.clone()).unwrap_or_else(|error| panic!("{error}"))
};
let value: serde_json::Value =
serde_json::from_str(output.trim()).unwrap_or_else(|error| panic!("{error}"));
assert_eq!(value["log"]["level"], "debug");
assert_eq!(value["log"]["message"], "trace event");
assert_ne!(value["log"]["timestamp_epoch_ms"], 1);
}
}