mod span;
mod span_v1;
use crate::span::v04::Span;
use crate::span::v1::TracerPayload;
use crate::span::TraceData;
use span::LogSpan;
use span_v1::{ChunkContextV1, LogSpanV1};
use std::io::Write;
const TRACE_PREFIX: &[u8] = b"{\"traces\":[[";
const TRACE_SUFFIX: &[u8] = b"]]}\n";
const TRACE_FORMAT_OVERHEAD: usize = TRACE_PREFIX.len() + TRACE_SUFFIX.len();
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct EncodeStats {
pub spans_written: usize,
pub spans_dropped: usize,
}
pub fn encode_traces<T: TraceData>(
traces: &[Vec<Span<T>>],
out: &mut impl Write,
max_line_size: usize,
) -> std::io::Result<EncodeStats> {
let mut stats = EncodeStats::default();
let mut span_buf: Vec<u8> = Vec::with_capacity(512);
let mut line: Vec<u8> = Vec::with_capacity(max_line_size.min(64 * 1024));
line.extend_from_slice(TRACE_PREFIX);
let mut line_span_count: usize = 0;
for trace in traces {
for span in trace {
span_buf.clear();
serde_json::to_writer(&mut span_buf, &LogSpan(span)).map_err(std::io::Error::other)?;
let span_len = span_buf.len();
if span_len + TRACE_FORMAT_OVERHEAD > max_line_size {
stats.spans_dropped += 1;
tracing::debug!(
span_len,
max_line_size,
"Span too large to send to logs, dropping"
);
continue;
}
let comma = usize::from(line_span_count > 0);
if line_span_count > 0
&& line.len() + comma + span_len + TRACE_SUFFIX.len() > max_line_size
{
flush_line(out, &mut line)?;
line_span_count = 0;
}
if line_span_count > 0 {
line.push(b',');
}
line.extend_from_slice(&span_buf);
line_span_count += 1;
stats.spans_written += 1;
}
}
if line_span_count > 0 {
flush_line(out, &mut line)?;
}
Ok(stats)
}
pub fn encode_traces_v1<T: TraceData>(
payload: &TracerPayload<T>,
out: &mut impl Write,
max_line_size: usize,
) -> std::io::Result<EncodeStats> {
let mut stats = EncodeStats::default();
let mut span_buf: Vec<u8> = Vec::with_capacity(512);
let mut line: Vec<u8> = Vec::with_capacity(max_line_size.min(64 * 1024));
line.extend_from_slice(TRACE_PREFIX);
let mut line_span_count: usize = 0;
for chunk in &payload.chunks {
let ctx = ChunkContextV1::new(
chunk,
&payload.env,
&payload.app_version,
&payload.attributes,
);
for span in &chunk.spans {
span_buf.clear();
serde_json::to_writer(&mut span_buf, &LogSpanV1(span, &ctx))
.map_err(std::io::Error::other)?;
let span_len = span_buf.len();
if span_len + TRACE_FORMAT_OVERHEAD > max_line_size {
stats.spans_dropped += 1;
tracing::debug!(
span_len,
max_line_size,
"Span too large to send to logs, dropping"
);
continue;
}
let comma = usize::from(line_span_count > 0);
if line_span_count > 0
&& line.len() + comma + span_len + TRACE_SUFFIX.len() > max_line_size
{
flush_line(out, &mut line)?;
line_span_count = 0;
}
if line_span_count > 0 {
line.push(b',');
}
line.extend_from_slice(&span_buf);
line_span_count += 1;
stats.spans_written += 1;
}
}
if line_span_count > 0 {
flush_line(out, &mut line)?;
}
Ok(stats)
}
fn flush_line(out: &mut impl Write, line: &mut Vec<u8>) -> std::io::Result<()> {
line.extend_from_slice(TRACE_SUFFIX);
out.write_all(line)?;
line.clear();
line.extend_from_slice(TRACE_PREFIX);
Ok(())
}
#[cfg(test)]
#[allow(clippy::useless_conversion)]
mod tests {
use super::*;
use crate::span::v04::SpanSlice;
use serde_json::Value;
const MAX: usize = 64 * 1024;
fn lines(out: &[u8]) -> Vec<String> {
String::from_utf8(out.to_vec())
.unwrap()
.lines()
.map(|s| s.to_string())
.collect()
}
#[test]
fn golden_known_span() {
let span = SpanSlice {
service: "my-fn".into(),
name: "aws.lambda".into(),
resource: "my-fn".into(),
r#type: "serverless".into(),
trace_id: 1,
span_id: 2,
parent_id: 0,
start: 1717200000000000000,
duration: 1500000,
error: 0,
meta: [("env".into(), "prod".into())].into_iter().collect(),
metrics: [("_sampling_priority_v1".into(), 1.0)]
.into_iter()
.collect(),
..Default::default()
};
let mut out = Vec::new();
let stats = encode_traces(&[vec![span]], &mut out, MAX).unwrap();
assert_eq!(stats.spans_written, 1);
assert_eq!(stats.spans_dropped, 0);
let expected = concat!(
"{\"traces\":[[",
"{\"trace_id\":\"0000000000000001\",",
"\"span_id\":\"0000000000000002\",",
"\"parent_id\":\"0000000000000000\",",
"\"service\":\"my-fn\",",
"\"name\":\"aws.lambda\",",
"\"resource\":\"my-fn\",",
"\"type\":\"serverless\",",
"\"error\":0,",
"\"start\":1717200000000000000,",
"\"duration\":1500000,",
"\"meta\":{\"env\":\"prod\"},",
"\"metrics\":{\"_sampling_priority_v1\":1.0}",
"}]]}\n",
);
assert_eq!(String::from_utf8(out).unwrap(), expected);
}
#[test]
fn meta_struct_is_omitted() {
let span = SpanSlice {
trace_id: 1,
span_id: 2,
meta_struct: [("_dd.appsec.json".into(), [0x81u8, 0xa4].as_slice())]
.into_iter()
.collect(),
..Default::default()
};
let mut out = Vec::new();
encode_traces(&[vec![span]], &mut out, MAX).unwrap();
let text = String::from_utf8(out).unwrap();
assert!(
!text.contains("meta_struct"),
"meta_struct must not be emitted: {text}"
);
let v: Value = serde_json::from_str(text.trim_end()).unwrap();
assert!(v["traces"][0][0]["trace_id"].is_string());
}
#[test]
fn hex_high_bits_and_root_parent() {
let trace_id: u128 = (0xABCDu128 << 64) | 0x1;
let span = SpanSlice {
trace_id,
span_id: 0xFF,
parent_id: 0,
..Default::default()
};
let mut out = Vec::new();
encode_traces(&[vec![span]], &mut out, MAX).unwrap();
let text = String::from_utf8(out).unwrap();
assert!(
text.contains("\"trace_id\":\"000000000000abcd0000000000000001\""),
"got: {text}"
);
assert!(text.contains("\"span_id\":\"00000000000000ff\""));
assert!(text.contains("\"parent_id\":\"0000000000000000\""));
}
#[test]
fn error_is_integer() {
let span = SpanSlice {
error: 1,
..Default::default()
};
let mut out = Vec::new();
encode_traces(&[vec![span]], &mut out, MAX).unwrap();
let text = String::from_utf8(out).unwrap();
assert!(text.contains("\"error\":1,"), "got: {text}");
}
#[test]
fn empty_maps_and_type_omitted() {
let span = SpanSlice {
name: "op".into(),
..Default::default()
};
let mut out = Vec::new();
encode_traces(&[vec![span]], &mut out, MAX).unwrap();
let text = String::from_utf8(out).unwrap();
assert!(!text.contains("\"meta\""), "got: {text}");
assert!(!text.contains("\"metrics\""));
assert!(!text.contains("\"meta_struct\""));
assert!(!text.contains("\"span_links\""));
assert!(!text.contains("\"span_events\""));
assert!(!text.contains("\"type\""));
}
#[test]
fn string_escaping() {
let span = SpanSlice {
name: "say \"hi\"\n".into(),
..Default::default()
};
let mut out = Vec::new();
encode_traces(&[vec![span]], &mut out, MAX).unwrap();
for line in lines(&out) {
let parsed: Value = serde_json::from_str(&line).unwrap();
assert_eq!(
parsed["traces"][0][0]["name"].as_str().unwrap(),
"say \"hi\"\n"
);
}
}
#[test]
fn size_cap_batches_into_multiple_lines() {
let make = |id: u64| SpanSlice {
name: "x".into(),
span_id: id,
..Default::default()
};
let trace: Vec<SpanSlice> = (1..=6).map(make).collect();
let mut one = Vec::new();
encode_traces(&[vec![make(1)]], &mut one, MAX).unwrap();
let span_len = one.len() - TRACE_FORMAT_OVERHEAD;
let cap = TRACE_FORMAT_OVERHEAD + span_len * 2 + 1;
let mut out = Vec::new();
let stats = encode_traces(&[trace], &mut out, cap).unwrap();
assert_eq!(stats.spans_written, 6);
assert_eq!(stats.spans_dropped, 0);
let emitted = lines(&out);
assert_eq!(emitted.len(), 3, "expected 3 lines, got {emitted:?}");
for line in &emitted {
assert!(line.len() <= cap, "line over cap: {} > {cap}", line.len());
let parsed: Value = serde_json::from_str(line).unwrap();
assert_eq!(parsed["traces"][0].as_array().unwrap().len(), 2);
}
}
#[test]
fn oversize_single_span_dropped() {
let big = "a".repeat(10_000);
let span = SpanSlice {
name: big.as_str().into(),
span_id: 1,
..Default::default()
};
let small = SpanSlice {
name: "ok".into(),
span_id: 2,
..Default::default()
};
let mut out = Vec::new();
let stats = encode_traces(&[vec![span, small]], &mut out, 1024).unwrap();
assert_eq!(stats.spans_dropped, 1);
assert_eq!(stats.spans_written, 1);
let emitted = lines(&out);
assert_eq!(emitted.len(), 1);
let parsed: Value = serde_json::from_str(&emitted[0]).unwrap();
assert_eq!(parsed["traces"][0][0]["name"].as_str().unwrap(), "ok");
}
#[test]
fn metric_non_finite_serializes_null() {
let span = SpanSlice {
span_id: 1,
metrics: [("nan".into(), f64::NAN), ("inf".into(), f64::INFINITY)]
.into_iter()
.collect(),
..Default::default()
};
let mut out = Vec::new();
encode_traces(&[vec![span]], &mut out, MAX).unwrap();
let emitted = lines(&out);
assert_eq!(emitted.len(), 1);
let parsed: Value = serde_json::from_str(&emitted[0]).unwrap();
let metrics = &parsed["traces"][0][0]["metrics"];
assert!(metrics["nan"].is_null(), "got: {metrics}");
assert!(metrics["inf"].is_null(), "got: {metrics}");
}
#[test]
fn multi_trace_flattened_into_one_line() {
let span_a = SpanSlice {
name: "a".into(),
span_id: 1,
..Default::default()
};
let span_b = SpanSlice {
name: "b".into(),
span_id: 2,
..Default::default()
};
let mut out = Vec::new();
let stats = encode_traces(&[vec![span_a], vec![span_b]], &mut out, MAX).unwrap();
assert_eq!(stats.spans_written, 2);
let emitted = lines(&out);
assert_eq!(emitted.len(), 1, "expected one line, got {emitted:?}");
let parsed: Value = serde_json::from_str(&emitted[0]).unwrap();
let inner = parsed["traces"][0].as_array().unwrap();
assert_eq!(inner.len(), 2);
assert_eq!(parsed["traces"].as_array().unwrap().len(), 1);
}
#[test]
fn span_links_and_events_emitted_when_present() {
use crate::span::v04::{SpanEvent, SpanLink};
let span = SpanSlice {
span_id: 1,
span_links: vec![SpanLink {
trace_id: 7,
span_id: 8,
..Default::default()
}],
span_events: vec![SpanEvent {
time_unix_nano: 123,
name: "evt".into(),
..Default::default()
}],
..Default::default()
};
let mut out = Vec::new();
encode_traces(&[vec![span]], &mut out, MAX).unwrap();
let emitted = lines(&out);
assert_eq!(emitted.len(), 1);
let parsed: Value = serde_json::from_str(&emitted[0]).unwrap();
let span_json = &parsed["traces"][0][0];
assert_eq!(span_json["span_links"][0]["trace_id"], 7);
assert_eq!(span_json["span_links"][0]["span_id"], 8);
assert_eq!(span_json["span_events"][0]["name"], "evt");
assert_eq!(span_json["span_events"][0]["time_unix_nano"], 123);
}
#[test]
fn span_link_flags_sentinel_bit_masked() {
use crate::span::v04::SpanLink;
fn encoded_flags(flags: u32) -> Value {
let span = SpanSlice {
span_id: 1,
span_links: vec![SpanLink {
trace_id: 7,
span_id: 8,
flags,
..Default::default()
}],
..Default::default()
};
let mut out = Vec::new();
encode_traces(&[vec![span]], &mut out, MAX).unwrap();
let emitted = lines(&out);
let parsed: Value = serde_json::from_str(&emitted[0]).unwrap();
parsed["traces"][0][0]["span_links"][0]["flags"].clone()
}
assert_eq!(
encoded_flags(0x8000_0001),
1,
"the sentinel bit must not leak into the JSON log"
);
assert_eq!(
encoded_flags(0x8000_0000),
0,
"an explicit drop decision must still emit flags: 0, not omit the field"
);
}
#[test]
fn empty_inner_trace_writes_nothing() {
let traces: Vec<Vec<SpanSlice>> = vec![vec![]];
let mut out = Vec::new();
let stats = encode_traces(&traces, &mut out, MAX).unwrap();
assert!(out.is_empty());
assert_eq!(stats.spans_written, 0);
assert_eq!(stats.spans_dropped, 0);
}
#[test]
fn empty_input_writes_nothing() {
let traces: Vec<Vec<SpanSlice>> = vec![];
let mut out = Vec::new();
let stats = encode_traces(&traces, &mut out, MAX).unwrap();
assert_eq!(stats, EncodeStats::default());
assert!(out.is_empty());
}
#[test]
fn forwarder_is_trace_contract() {
let trace: Vec<SpanSlice> = (1u64..=5)
.map(|i| SpanSlice {
service: "svc".into(),
name: "op".into(),
trace_id: u128::from(i),
span_id: i + 100,
..Default::default()
})
.collect();
let mut out = Vec::new();
encode_traces(&[trace], &mut out, 200).unwrap();
let text = String::from_utf8(out).unwrap();
assert!(!text.is_empty());
for chunk in text.split_inclusive('\n') {
assert!(chunk.ends_with('\n'), "line not newline-terminated");
let line = chunk.trim_end_matches('\n');
let parsed: Value = serde_json::from_str(line).unwrap();
let traces = parsed.get("traces").and_then(Value::as_array).unwrap();
assert!(!traces.is_empty());
let first = traces[0].as_array().unwrap();
assert!(!first.is_empty());
let trace_id = first[0].get("trace_id").unwrap();
assert!(trace_id.is_string());
assert!(!trace_id.as_str().unwrap().is_empty());
}
}
}