use dial9_trace_format::decoder::{Decoder, RawEvent};
use dial9_trace_format::types::FieldValueRef;
use lasso::{Rodeo, Spur};
use rustc_hash::FxHashMap;
use serde::Deserialize;
use super::clock::ClockOffset;
const BASE_SPAN_FIELDS: &[&str] = &[
"timestamp_ns",
"worker_id",
"span_id",
"span_instance_id",
"tid",
"parent_span_id",
"span_name",
];
fn extract_span_attributes(ev: &RawEvent<'_, '_>) -> Vec<(String, String)> {
let mut attrs = Vec::new();
for (name, value) in ev.field_names().zip(ev.fields.iter()) {
if BASE_SPAN_FIELDS.contains(&name) {
continue;
}
let rendered = match value {
FieldValueRef::String(s) => Some((*s).to_string()),
FieldValueRef::PooledString(id) => ev.string_pool.get(*id).map(|s| s.to_string()),
FieldValueRef::I64(v) => Some(v.to_string()),
FieldValueRef::Varint(v) => Some(v.to_string()),
FieldValueRef::F64(v) => Some(v.to_string()),
FieldValueRef::Bool(v) => Some(v.to_string()),
FieldValueRef::Bytes(_)
| FieldValueRef::StackFrames(_)
| FieldValueRef::PooledStackFrames(_)
| FieldValueRef::StringMap(_)
| FieldValueRef::List(_)
| FieldValueRef::Map(_)
| FieldValueRef::None => None,
_ => None,
};
if let Some(rendered) = rendered {
attrs.push((name.to_string(), rendered));
}
}
attrs
}
#[derive(Debug, Deserialize)]
pub(crate) struct CpuSample {
pub(crate) timestamp_ns: u64,
pub(crate) tid: u32,
pub(crate) source: u64,
pub(crate) callchain: Vec<u64>,
}
#[derive(Debug, Deserialize)]
pub(crate) struct WorkerPark {
pub(crate) timestamp_ns: u64,
pub(crate) worker_id: u64,
pub(crate) tid: u32,
}
#[derive(Debug, Deserialize)]
pub(crate) struct WorkerUnpark {
pub(crate) timestamp_ns: u64,
pub(crate) worker_id: u64,
pub(crate) tid: u32,
}
#[derive(Debug, Deserialize)]
pub(crate) struct PollStart {
pub(crate) timestamp_ns: u64,
pub(crate) worker_id: u64,
pub(crate) task_id: u64,
#[serde(default)]
pub(crate) spawn_loc: Option<String>,
}
#[derive(Debug, Deserialize)]
pub(crate) struct PollEnd {
pub(crate) timestamp_ns: u64,
pub(crate) worker_id: u64,
}
#[derive(Debug, Deserialize)]
pub(crate) struct TaskSpawn {
pub(crate) timestamp_ns: u64,
pub(crate) task_id: u64,
#[serde(default)]
pub(crate) instrumented: Option<bool>,
}
#[derive(Debug, Deserialize)]
pub(crate) struct TaskTerminate {
pub(crate) timestamp_ns: u64,
pub(crate) task_id: u64,
}
#[derive(Debug, Deserialize)]
pub(crate) struct WakeEvent {
pub(crate) timestamp_ns: u64,
pub(crate) waker_task_id: u64,
pub(crate) woken_task_id: u64,
}
#[derive(Debug, Deserialize)]
pub(crate) struct ClockSync {
pub(crate) timestamp_ns: u64,
pub(crate) realtime_ns: u64,
}
#[derive(Debug, Deserialize)]
pub(crate) struct SymbolEntry {
pub(crate) addr: u64,
pub(crate) inline_depth: u64,
pub(crate) symbol_name: String,
}
#[derive(Debug, Deserialize)]
#[allow(dead_code)]
pub(crate) struct LegacySpanEnterEvent {
pub(crate) timestamp_ns: u64,
#[serde(default)]
pub(crate) worker_id: u64,
#[serde(default)]
pub(crate) span_id: u64,
#[serde(default)]
pub(crate) parent_span_id: Option<u64>,
#[serde(default)]
pub(crate) span_name: Option<String>,
#[serde(skip)]
pub(crate) decode_sequence: u64,
#[serde(skip)]
pub(crate) attributes: Vec<(String, String)>,
}
#[derive(Debug, Deserialize)]
#[allow(dead_code)]
pub(crate) struct LegacySpanExitEvent {
pub(crate) timestamp_ns: u64,
#[serde(default)]
pub(crate) worker_id: u64,
#[serde(default)]
pub(crate) span_id: u64,
#[serde(default)]
pub(crate) span_name: Option<String>,
#[serde(skip)]
pub(crate) decode_sequence: u64,
#[serde(skip)]
pub(crate) attributes: Vec<(String, String)>,
}
#[derive(Debug, Deserialize)]
pub(crate) struct LegacySpanCloseEvent {
pub(crate) timestamp_ns: u64,
#[serde(default)]
pub(crate) span_id: u64,
#[serde(skip)]
pub(crate) decode_sequence: u64,
}
#[derive(Debug, Clone)]
pub(crate) struct LegacySpanSchemaInfo {
pub(crate) target: String,
pub(crate) name: String,
pub(crate) file: Option<String>,
pub(crate) line: Option<u32>,
}
pub(crate) fn parse_legacy_span_schema_name(schema_name: &str) -> Option<LegacySpanSchemaInfo> {
if let Some(name) = schema_name
.strip_prefix("SpanEnter__")
.or_else(|| schema_name.strip_prefix("SpanExit__"))
.filter(|name| !name.is_empty())
{
return Some(LegacySpanSchemaInfo {
target: String::new(),
name: name.to_string(),
file: None,
line: None,
});
}
let rest = schema_name
.strip_prefix("SpanEnter:")
.or_else(|| schema_name.strip_prefix("SpanExit:"))?;
let line_colon_pos = rest.rfind(':')?;
let line_str = &rest[line_colon_pos + 1..];
let line: u32 = line_str.parse().ok()?;
let before_line = &rest[..line_colon_pos];
let file_colon_pos = find_first_single_colon(before_line)?;
let target_name = &before_line[..file_colon_pos];
let file = &before_line[file_colon_pos + 1..];
let (target, name) = if let Some(last_dcolon) = target_name.rfind("::") {
(&target_name[..last_dcolon], &target_name[last_dcolon + 2..])
} else {
("", target_name)
};
Some(LegacySpanSchemaInfo {
target: target.to_string(),
name: name.to_string(),
file: if file.is_empty() {
None
} else {
Some(file.to_string())
},
line: Some(line),
})
}
pub(crate) fn find_first_single_colon(s: &str) -> Option<usize> {
let bytes = s.as_bytes();
let mut i = 0;
while i < bytes.len() {
if bytes[i] == b':' {
let preceded_by_colon = i > 0 && bytes[i - 1] == b':';
let followed_by_colon = i + 1 < bytes.len() && bytes[i + 1] == b':';
if !preceded_by_colon && !followed_by_colon {
return Some(i);
}
if followed_by_colon {
i += 2;
continue;
}
}
i += 1;
}
None
}
pub(crate) enum TraceEvent {
CpuSample(CpuSample),
WorkerPark(WorkerPark),
WorkerUnpark(WorkerUnpark),
PollStart(PollStart),
PollEnd(PollEnd),
TaskSpawn(TaskSpawn),
TaskTerminate(TaskTerminate),
Wake(WakeEvent),
}
impl TraceEvent {
pub(crate) fn timestamp_ns(&self) -> u64 {
match self {
Self::CpuSample(e) => e.timestamp_ns,
Self::WorkerPark(e) => e.timestamp_ns,
Self::WorkerUnpark(e) => e.timestamp_ns,
Self::PollStart(e) => e.timestamp_ns,
Self::PollEnd(e) => e.timestamp_ns,
Self::TaskSpawn(e) => e.timestamp_ns,
Self::TaskTerminate(e) => e.timestamp_ns,
Self::Wake(e) => e.timestamp_ns,
}
}
}
pub(crate) struct DecodedTrace {
pub(crate) interner: Rodeo,
pub(crate) addr_to_keys: FxHashMap<u64, Vec<(u64, Spur)>>,
pub(crate) events: Vec<TraceEvent>,
pub(crate) clock_offset: Option<ClockOffset>,
pub(crate) segment_metadata_boot_id: Option<String>,
pub(crate) legacy_enters: Vec<(String, LegacySpanEnterEvent)>,
pub(crate) legacy_exits: Vec<(String, LegacySpanExitEvent)>,
pub(crate) legacy_closes: Vec<LegacySpanCloseEvent>,
}
pub(crate) fn decode_trace(data: &[u8], source_key: &str) -> anyhow::Result<DecodedTrace> {
let mut decoder = Decoder::new(data).ok_or_else(|| anyhow::anyhow!("invalid trace header"))?;
let mut interner = Rodeo::default();
let mut addr_to_keys: FxHashMap<u64, Vec<(u64, Spur)>> = FxHashMap::default();
let mut events: Vec<TraceEvent> = Vec::new();
let mut clock_offset: Option<ClockOffset> = None;
let mut segment_metadata_boot_id: Option<String> = None;
let mut span_decode_sequence: u64 = 0;
let mut legacy_enters: Vec<(String, LegacySpanEnterEvent)> = Vec::new(); let mut legacy_exits: Vec<(String, LegacySpanExitEvent)> = Vec::new();
let mut legacy_closes: Vec<LegacySpanCloseEvent> = Vec::new();
let mut legacy_enter_decode_errors: u64 = 0;
let mut legacy_exit_decode_errors: u64 = 0;
let mut legacy_close_decode_errors: u64 = 0;
decoder
.for_each_event(|ev| match ev.name {
"ClockSyncEvent" => {
if let Ok(cs) = ev.deserialize::<ClockSync>()
&& cs.realtime_ns > 0
&& cs.timestamp_ns > 0
&& clock_offset.is_none()
{
clock_offset = Some(ClockOffset::from_clock_sync(
cs.realtime_ns,
cs.timestamp_ns,
));
}
}
"SegmentMetadataEvent" => {
#[derive(serde::Deserialize)]
struct SegmentMeta {
#[serde(default)]
entries: std::collections::HashMap<String, String>,
}
if let Ok(meta) = ev.deserialize::<SegmentMeta>()
&& segment_metadata_boot_id.is_none()
&& let Some(bid) = meta.entries.get("boot_id")
&& !bid.is_empty()
{
segment_metadata_boot_id = Some(bid.clone());
}
}
"CpuSampleEvent" | "CpuSample" => {
if let Ok(s) = ev.deserialize::<CpuSample>()
&& !s.callchain.is_empty()
{
events.push(TraceEvent::CpuSample(s));
}
}
"WorkerParkEvent" => {
if let Ok(p) = ev.deserialize::<WorkerPark>() {
events.push(TraceEvent::WorkerPark(p));
}
}
"WorkerUnparkEvent" => {
if let Ok(u) = ev.deserialize::<WorkerUnpark>() {
events.push(TraceEvent::WorkerUnpark(u));
}
}
"PollStartEvent" => {
if let Ok(p) = ev.deserialize::<PollStart>() {
events.push(TraceEvent::PollStart(p));
}
}
"PollEndEvent" => {
if let Ok(p) = ev.deserialize::<PollEnd>() {
events.push(TraceEvent::PollEnd(p));
}
}
"TaskSpawnEvent" => {
if let Ok(event) = ev.deserialize::<TaskSpawn>() {
events.push(TraceEvent::TaskSpawn(event));
}
}
"TaskTerminateEvent" => {
if let Ok(event) = ev.deserialize::<TaskTerminate>() {
events.push(TraceEvent::TaskTerminate(event));
}
}
"WakeEventEvent" => {
if let Ok(event) = ev.deserialize::<WakeEvent>() {
events.push(TraceEvent::Wake(event));
}
}
"SymbolTableEntry" => {
if let Ok(sym) = ev.deserialize::<SymbolEntry>() {
let key = interner.get_or_intern(&sym.symbol_name);
addr_to_keys
.entry(sym.addr)
.or_default()
.push((sym.inline_depth, key));
}
}
name if name == "SpanCloseEvent" || name.starts_with("SpanClose__") => {
match ev.deserialize::<LegacySpanCloseEvent>() {
Ok(mut lc) if lc.span_id > 0 => {
lc.decode_sequence = span_decode_sequence;
span_decode_sequence += 1;
legacy_closes.push(lc);
}
_ => {
legacy_close_decode_errors += 1;
}
}
}
name if name.starts_with("SpanEnter:") || name.starts_with("SpanEnter__") => {
match ev.deserialize::<LegacySpanEnterEvent>() {
Ok(mut le) if le.span_id > 0 => {
le.decode_sequence = span_decode_sequence;
span_decode_sequence += 1;
le.attributes = extract_span_attributes(&ev);
legacy_enters.push((name.to_string(), le));
}
_ => {
legacy_enter_decode_errors += 1;
}
}
}
name if name.starts_with("SpanExit:") || name.starts_with("SpanExit__") => {
match ev.deserialize::<LegacySpanExitEvent>() {
Ok(mut le) if le.span_id > 0 => {
le.decode_sequence = span_decode_sequence;
span_decode_sequence += 1;
le.attributes = extract_span_attributes(&ev);
legacy_exits.push((name.to_string(), le));
}
_ => {
legacy_exit_decode_errors += 1;
}
}
}
_ => {}
})
.map_err(|e| anyhow::anyhow!("decode error: {e}"))?;
if legacy_close_decode_errors > 0
|| legacy_enter_decode_errors > 0
|| legacy_exit_decode_errors > 0
{
use dial9_core::rate_limited;
rate_limited!(std::time::Duration::from_secs(60), {
tracing::warn!(
source_key,
legacy_enter_errors = legacy_enter_decode_errors,
legacy_exit_errors = legacy_exit_decode_errors,
legacy_close_errors = legacy_close_decode_errors,
"skipped malformed span event(s) during decode"
);
});
}
Ok(DecodedTrace {
interner,
addr_to_keys,
events,
clock_offset,
segment_metadata_boot_id,
legacy_enters,
legacy_exits,
legacy_closes,
})
}