Skip to main content

agent_first_data/
afdata_tracing.rs

1//! AFDATA-compliant tracing layer.
2//!
3//! Outputs log events using agent-first-data's `render` function:
4//! - JSON: single-line JSONL (secrets redacted, original keys)
5//! - Plain: single-line logfmt (keys stripped, values formatted)
6//! - YAML: multi-line, structure-preserving
7//!   (original keys and values kept, secrets redacted)
8//!
9//! Span fields are flattened into every event line (e.g. `request_id`).
10//! All other tracing features (macros, spans, EnvFilter) work unchanged.
11//!
12//! # Usage
13//! ```ignore
14//! use agent_first_data::{Redactor, afdata_tracing::{self, LogFormat}};
15//! use tracing_subscriber::EnvFilter;
16//!
17//! afdata_tracing::try_init(EnvFilter::new("info"), LogFormat::Json, Redactor::new())?;
18//! ```
19
20use 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/// Output format for the AFDATA tracing layer.
35#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
36#[serde(rename_all = "lowercase")]
37pub enum LogFormat {
38    Json,
39    Plain,
40    /// Structure-preserving YAML.
41    Yaml,
42}
43
44impl LogFormat {
45    /// Return the canonical config spelling.
46    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
88/// The one destination a layer and every handle it issued write through.
89///
90/// `with_writer` swaps what is inside the cell rather than handing out a new
91/// one, so a `StructuredLogHandle` taken before configuration cannot keep
92/// writing to the old stream — a split that nothing would report at compile or
93/// run time.
94struct 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        // One write, not two. The inner `Mutex` only serializes this layer's
142        // own writes; a shared destination like `io::stderr` locks per call, so
143        // splitting the line from its newline lets another writer on the same
144        // stream — `CliEmitter`'s error events, under default split routing —
145        // land between them. That reader then loses both events to one
146        // malformed line, including the terminal one.
147        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
155/// A tracing Layer that outputs AFDATA-compliant log lines.
156pub 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/// Direct nested-JSON logger sharing an [`AfdataLayer`]'s writer and ordering.
173#[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
190/// Try to initialize tracing with AFDATA output.
191///
192/// Returns `Err` if a global tracing subscriber is already initialized. This is
193/// the convenience entry point for a global subscriber. Use
194/// [`AfdataLayer::new`] directly when composing a subscriber, injecting a
195/// writer, or retaining a [`StructuredLogHandle`].
196///
197/// # Arguments
198/// * `filter` - tracing_subscriber::EnvFilter controlling which events are recorded
199/// * `format` - LogFormat::Json, LogFormat::Plain, or LogFormat::Yaml
200/// * `redactor` - Redactor with optional custom secret field names and policy
201pub 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    /// Build a composable layer that writes to stderr.
217    ///
218    /// Diagnostic log events use stderr by default; [`AfdataLayer::with_writer`]
219    /// replaces that sanctioned process sink.
220    #[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    /// Replace the writer factory used by this layer.
230    ///
231    /// The factory is serialized behind a shared lock, so ordinary tracing and
232    /// [`StructuredLogHandle`] writes use the same sink in call order. Handles
233    /// taken before this call follow the new writer too: they share one cell
234    /// rather than a snapshot, so the two cannot end up addressing different
235    /// streams depending on the order the layer was configured in.
236    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    /// Get a direct nested-JSON logger sharing this layer's sink.
246    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    /// Emit a nested JSON log payload.
286    ///
287    /// Unlike tracing's scalar field visitor, objects and arrays stay
288    /// structured. Redaction and formatting are identical to the layer.
289    ///
290    /// The line carries `level` and `timestamp_epoch_ms` like every line the
291    /// layer writes, so one stream does not mix records that a level filter or
292    /// a timestamp correlation can read with records it silently drops. Supply
293    /// `level` in `payload` to override the default; the timestamp is always
294    /// stamped here, as it is for a tracing event.
295    ///
296    /// Note that a writer built with a level filter (for example
297    /// `io::stderr.with_max_level(..)`) cannot route these events: they have no
298    /// tracing metadata to route on, so they always take the default writer.
299    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
310/// Stored in span extensions to carry structured fields.
311struct 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        // Collect fields from the event
344        let mut visitor = JsonVisitor::new();
345        event.record(&mut visitor);
346
347        // Build output object with AFDATA field names.
348        let mut map = serde_json::Map::with_capacity(4 + visitor.fields.len());
349
350        let level = match *meta.level() {
351            // AFDATA's log level set starts at debug. TRACE remains
352            // distinguishable from info without inventing a fifth level.
353            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        // "message" field from the tracing macro's format string
360        let message = visitor
361            .message
362            .take()
363            .unwrap_or_else(|| "(no message)".to_string());
364
365        // Flatten span fields from root to leaf (child overrides parent on collision)
366        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        // Append event-level structured fields. Logs no longer use top-level
380        // protocol code; code may be a tool-defined field inside the log
381        // payload.
382        for (k, v) in visitor.fields {
383            if !is_reserved_log_field(&k) {
384                map.insert(k, v);
385            }
386        }
387
388        // Adapter metadata wins over same-named span/event fields.
389        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        // Format using the library's own output functions.
402        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
412/// Give a directly-emitted payload the same envelope fields the layer stamps on
413/// every tracing event, so both kinds of record are readable the same way.
414///
415/// A non-object payload is left alone: it has nowhere to put the fields, and
416/// wrapping it would change the shape the caller asked for.
417fn 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
430/// Visitor that collects tracing event fields into a JSON map.
431struct 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            // Push the raw value under its field name. Redaction happens at emit
452            // time in `on_event` via `render`, which redacts by field name
453            // (`_secret` suffix, `_url` scrubbing) —
454            // exactly like every other AFDATA surface. The visitor never scans
455            // rendered values for secret markers.
456            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    // The tracing layer redacts log fields the same way every AFDATA surface
510    // does: by FIELD NAME, applied by `output_*` at emit time — never by
511    // scanning a rendered value for the substring "_secret". These tests pin
512    // that contract (the visitor records raw values; emit redacts by name).
513
514    #[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        // A real secret value never contains the literal "_secret"; the old
537        // substring scan only ever produced false positives like this one.
538        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        // record_str / record_i64 etc. push raw values too; emit-time redaction
555        // covers every record_* path, not just record_debug.
556        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    /// Taking the handle before `with_writer` must not leave it addressing the
718    /// old stream. Nothing would report that split, so the ordering the docs
719    /// promise is pinned here rather than assumed.
720    #[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        // Handle first, writer second — the inverse of the usual order.
725        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    /// Every line on the stream carries the envelope a reader filters on, so a
747    /// level filter cannot silently drop the directly-emitted ones.
748    #[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}