mod attribution;
pub(crate) mod clock;
mod events;
pub(crate) mod polls;
pub(crate) mod spans;
mod types;
use events::*;
use rustc_hash::FxHashMap;
#[cfg(test)]
use spans::span_builder;
use spans::{interval_pairing, legacy};
pub(crate) use types::{
DecodeResult, DecodeStats, EnclosingSpanSummary, ResolvedPoll, ResolvedSample, ResolvedSpan,
SchedulingDelayKind,
};
const SOURCE_CPU_PROFILE: u8 = 0;
fn parse_source_key(key: &str) -> (String, String, String) {
let path = if let Some(rest) = key.strip_prefix("s3://") {
rest.split_once('/').map_or(rest, |(_, p)| p)
} else {
key
};
let parts: Vec<&str> = path.split('/').collect();
if let Some(anchor) = parts.iter().position(|p| is_date(p)) {
let date = parts.get(anchor).copied().unwrap_or("").to_string();
let service = parts.get(anchor + 2).copied().unwrap_or("").to_string();
let host = parts.get(anchor + 3).copied().unwrap_or("").to_string();
(date, service, host)
} else {
let date = parts.first().copied().unwrap_or("").to_string();
let service = parts.get(2).copied().unwrap_or("").to_string();
let host = parts.get(3).copied().unwrap_or("").to_string();
(date, service, host)
}
}
fn is_date(s: &str) -> bool {
let b = s.as_bytes();
b.len() == 10
&& b[4] == b'-'
&& b[7] == b'-'
&& b[..4].iter().all(u8::is_ascii_digit)
&& b[5..7].iter().all(u8::is_ascii_digit)
&& b[8..].iter().all(u8::is_ascii_digit)
}
pub(crate) fn decode_samples(data: &[u8], source_key: &str) -> anyhow::Result<DecodeResult> {
decode_samples_with_stats(data, source_key).map(|(result, _stats)| result)
}
pub(crate) fn decode_samples_with_stats(
data: &[u8],
source_key: &str,
) -> anyhow::Result<(DecodeResult, DecodeStats)> {
use std::time::Instant;
let mut stats = DecodeStats::default();
let t_wire = Instant::now();
let events::DecodedTrace {
interner,
mut addr_to_keys,
mut events,
mut clock_offset,
segment_metadata_boot_id,
legacy_enters,
legacy_exits,
legacy_closes,
} = events::decode_trace(data, source_key)?;
stats.wire_decode = t_wire.elapsed();
stats.span_events_decoded =
(legacy_enters.len() + legacy_exits.len() + legacy_closes.len()) as u64;
stats.events_decoded = events.len() as u64 + stats.span_events_decoded;
let timestamp_bounds = events
.iter()
.map(TraceEvent::timestamp_ns)
.chain(legacy_enters.iter().map(|(_, event)| event.timestamp_ns))
.chain(legacy_exits.iter().map(|(_, event)| event.timestamp_ns))
.chain(legacy_closes.iter().map(|event| event.timestamp_ns))
.fold(None, |bounds: Option<(u64, u64)>, timestamp| {
Some(match bounds {
Some((min, max)) => (min.min(timestamp), max.max(timestamp)),
None => (timestamp, timestamp),
})
});
if let (Some(offset), Some((min, max))) = (clock_offset, timestamp_bounds)
&& !offset.is_valid_for(clock::MonoNs(min), clock::MonoNs(max))
{
use dial9_core::rate_limited;
rate_limited!(std::time::Duration::from_secs(60), {
tracing::warn!(
source_key,
min_mono_ns = min,
max_mono_ns = max,
"ignoring clock sync that cannot convert the trace timestamp range"
);
});
clock_offset = None;
}
tracing::info!("sorting {} events", events.len());
let t_sort = Instant::now();
events.sort_unstable_by_key(|e| e.timestamp_ns());
for entries in addr_to_keys.values_mut() {
entries.sort_unstable_by_key(|(d, _)| *d);
}
stats.sort_events = t_sort.elapsed();
let t_polls = Instant::now();
let mut poll_timeline = polls::PollTimeline::reconstruct(&events);
stats.poll_reconstruct = t_polls.elapsed();
let t_samples = Instant::now();
let mut stacks_dict: FxHashMap<[u8; 16], Vec<String>> = FxHashMap::default();
let mut stack_cache: FxHashMap<Vec<u64>, [u8; 16]> = FxHashMap::default();
let mut samples = Vec::new();
let (parsed_date, parsed_service, parsed_host) = parse_source_key(source_key);
for event in &events {
match event {
TraceEvent::WorkerPark(_)
| TraceEvent::WorkerUnpark(_)
| TraceEvent::PollStart(_)
| TraceEvent::PollEnd(_)
| TraceEvent::TaskSpawn(_)
| TraceEvent::TaskTerminate(_)
| TraceEvent::Wake(_) => {}
TraceEvent::CpuSample(s) => {
let (worker_id, poll_duration_ns, spawn_location) = poll_timeline.attribute_sample(
s.tid,
clock::MonoNs(s.timestamp_ns),
s.source as u8,
);
let stack_id = if let Some(&cached) = stack_cache.get(&s.callchain) {
cached
} else {
let mut hasher = blake3::Hasher::new();
let mut first = true;
let mut frame_strings: Vec<String> = Vec::new();
for &addr in &s.callchain {
if let Some(entries) = addr_to_keys.get(&addr) {
for (_, key) in entries {
let name = interner.resolve(key);
if !first {
hasher.update(b"\x00");
}
hasher.update(name.as_bytes());
frame_strings.push(name.to_string());
first = false;
}
} else {
let hex = format!("0x{addr:x}");
if !first {
hasher.update(b"\x00");
}
hasher.update(hex.as_bytes());
frame_strings.push(hex);
first = false;
}
}
if frame_strings.is_empty() {
continue;
}
let hash = hasher.finalize();
let mut id = [0u8; 16];
id.copy_from_slice(&hash.as_bytes()[..16]);
stacks_dict.entry(id).or_insert(frame_strings);
stack_cache.insert(s.callchain.clone(), id);
id
};
let wall_ns = clock::MonoNs(s.timestamp_ns)
.to_wall_or_raw(clock_offset)
.raw();
samples.push(ResolvedSample {
timestamp_ns: wall_ns,
stack_id,
worker_id,
source: s.source as u8,
source_key: source_key.to_string(),
host: parsed_host.clone(),
service: parsed_service.clone(),
date: parsed_date.clone(),
poll_duration_ns,
spawn_location,
enclosing_spans: Vec::new(),
});
}
}
}
let resolved_polls =
poll_timeline.resolved(clock_offset, &parsed_host, &parsed_service, &parsed_date);
stats.sample_resolve = t_samples.elapsed();
let boot_id: String = match segment_metadata_boot_id {
Some(meta_bid) => meta_bid,
None => extract_boot_id_from_path_qualified(source_key)
.0
.to_string(),
};
let mut resolved_spans: Vec<ResolvedSpan> = Vec::new();
let mut legacy_intervals: FxHashMap<u64, Vec<interval_pairing::MonoInterval>> =
FxHashMap::default();
let t_spans = Instant::now();
if !legacy_enters.is_empty() || !legacy_closes.is_empty() {
let legacy_resolution = resolve_legacy_spans(
&legacy_enters,
&legacy_exits,
&legacy_closes,
poll_timeline.records(),
source_key,
&boot_id,
clock_offset,
&parsed_host,
&parsed_service,
&parsed_date,
);
legacy_intervals = legacy_resolution.instance_intervals;
resolved_spans.extend(legacy_resolution.spans);
}
stats.span_resolve = t_spans.elapsed();
let t_attr = Instant::now();
attribution::attribute_samples_to_spans(
&mut samples,
&mut resolved_spans,
&legacy_intervals,
&boot_id,
clock_offset,
);
stats.sample_attribution = t_attr.elapsed();
Ok((
(
samples,
stacks_dict.into_iter().collect(),
resolved_polls,
resolved_spans,
),
stats,
))
}
#[cfg(test)]
fn compute_span_uid(boot_id: &str, span_instance_id: u64) -> [u8; 16] {
span_builder::compute_span_uid(boot_id, span_instance_id)
}
#[cfg(test)]
fn extract_boot_id_from_path(source_key: &str) -> &str {
extract_boot_id_from_path_qualified(source_key).0
}
fn extract_boot_id_from_path_qualified(source_key: &str) -> (&str, bool) {
let path = if let Some(rest) = source_key.strip_prefix("s3://") {
rest.split_once('/').map_or(rest, |(_, p)| p)
} else {
source_key
};
let parts: Vec<&str> = path.rsplitn(3, '/').collect();
if parts.len() >= 2 && !parts[1].is_empty() {
let candidate = parts[1];
let is_namespaced = is_boot_id_format(candidate);
(candidate, is_namespaced)
} else {
let fallback = path.rsplit_once('/').map_or(path, |(dir, _)| dir);
(fallback, false)
}
}
fn is_boot_id_format(s: &str) -> bool {
let Some((alpha, digits)) = s.split_once('-') else {
return false;
};
alpha.len() == 4
&& alpha.bytes().all(|b| b.is_ascii_lowercase())
&& !digits.is_empty()
&& digits.bytes().all(|b| b.is_ascii_digit())
}
#[cfg(test)]
fn compute_span_type_uid(
kind: &str,
target: &str,
name: &str,
file: Option<&str>,
line: Option<u32>,
) -> [u8; 16] {
span_builder::compute_span_type_uid(kind, target, name, file, line)
}
struct SpanResolution {
spans: Vec<ResolvedSpan>,
instance_intervals: FxHashMap<u64, Vec<interval_pairing::MonoInterval>>,
}
#[allow(clippy::too_many_arguments)]
fn resolve_legacy_spans(
legacy_enters: &[(String, LegacySpanEnterEvent)],
legacy_exits: &[(String, LegacySpanExitEvent)],
legacy_closes: &[LegacySpanCloseEvent],
polls: &[polls::PollRecord],
source_key: &str,
boot_id: &str,
clock_offset: Option<clock::ClockOffset>,
host: &str,
service: &str,
date: &str,
) -> SpanResolution {
let result = legacy::resolve_legacy_spans(
legacy_enters,
legacy_exits,
legacy_closes,
polls,
source_key,
boot_id,
clock_offset,
host,
service,
date,
);
SpanResolution {
spans: result.spans,
instance_intervals: result.instance_intervals,
}
}
#[cfg(test)]
fn resolve_span_task(worker_polls: &[(u64, u64, u64)], enter_ts: u64) -> Option<u64> {
let worker_polls: Vec<_> = worker_polls
.iter()
.map(|&(start, end, task_id)| (clock::MonoNs(start), clock::MonoNs(end), task_id))
.collect();
legacy::resolve_span_task(&worker_polls, clock::MonoNs(enter_ts))
}
#[cfg(test)]
fn attribute_legacy_span_from_polls(
entered: &[(u64, u64)],
task_polls: &[(u64, u64)],
) -> (u64, u64) {
let entered: Vec<_> = entered
.iter()
.map(|&(start, end)| (clock::MonoNs(start), clock::MonoNs(end)))
.collect();
let task_polls: Vec<_> = task_polls
.iter()
.map(|&(start, end)| (clock::MonoNs(start), clock::MonoNs(end)))
.collect();
legacy::attribute_legacy_span_from_polls(&entered, &task_polls)
}
#[cfg(test)]
fn union_intervals(intervals: &[(u64, u64)]) -> u64 {
let intervals: Vec<_> = intervals
.iter()
.map(|&(start, end)| (clock::MonoNs(start), clock::MonoNs(end)))
.collect();
interval_pairing::union_interval_duration(&intervals).raw()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parse_source_key_is_date_anchored() {
assert_eq!(
parse_source_key("2026-06-19/1300/svc/host-a/boot/0-0.bin.gz"),
(
"2026-06-19".to_string(),
"svc".to_string(),
"host-a".to_string()
)
);
assert_eq!(
parse_source_key("traces/2026-06-19/1300/svc/host-a/boot/0-0.bin.gz"),
(
"2026-06-19".to_string(),
"svc".to_string(),
"host-a".to_string()
)
);
assert_eq!(
parse_source_key("s3://bucket/traces/2026-06-19/1300/svc/host-a/boot/0-0.bin.gz"),
(
"2026-06-19".to_string(),
"svc".to_string(),
"host-a".to_string()
)
);
assert_eq!(
parse_source_key("a/b/c/d"),
("a".to_string(), "c".to_string(), "d".to_string())
);
}
fn load_demo_trace() -> Vec<u8> {
let data = std::fs::read(concat!(
env!("CARGO_MANIFEST_DIR"),
"/ui/public/demo-trace.bin"
))
.unwrap();
let mut dec = flate2::read::GzDecoder::new(data.as_slice());
let mut buf = Vec::new();
std::io::Read::read_to_end(&mut dec, &mut buf).unwrap();
buf
}
#[test]
fn test_stack_id_deterministic() {
let decompressed = load_demo_trace();
let (s1, d1, _, _) = decode_samples(&decompressed, "test").unwrap();
let (s2, d2, _, _) = decode_samples(&decompressed, "test").unwrap();
assert_eq!(s1.len(), s2.len());
assert_eq!(d1.len(), d2.len());
for (a, b) in s1.iter().zip(s2.iter()) {
assert_eq!(a.stack_id, b.stack_id);
assert_eq!(a.timestamp_ns, b.timestamp_ns);
assert_eq!(a.worker_id, b.worker_id);
assert_eq!(a.source, b.source);
}
}
#[test]
fn test_decode_demo_trace() {
let decompressed = load_demo_trace();
let (samples, stacks, polls, _spans) =
decode_samples(&decompressed, "demo-trace.bin").unwrap();
assert!(!samples.is_empty(), "expected CPU samples in demo trace");
assert!(!stacks.is_empty(), "expected stacks in dictionary");
for sample in &samples {
assert!(stacks.contains_key(&sample.stack_id));
}
let min_ts = samples.iter().map(|s| s.timestamp_ns).min().unwrap();
assert!(
min_ts > 1_500_000_000_000_000_000,
"timestamps should be wall-clock epoch ns, got {min_ts}"
);
assert!(!polls.is_empty(), "expected poll spans in demo trace");
let attributed = samples
.iter()
.filter(|s| s.poll_duration_ns.is_some())
.count();
assert!(
attributed > 0,
"expected some samples attributed to a poll, got 0"
);
eprintln!(
"decoded {} samples ({} poll-attributed), {} unique stacks, {} polls",
samples.len(),
attributed,
stacks.len(),
polls.len(),
);
}
#[test]
fn invalid_clock_sync_is_ignored_instead_of_clamped() {
use dial9_trace_format::TraceEvent;
use dial9_trace_format::encoder::Encoder;
#[derive(TraceEvent)]
struct ClockSyncEvent {
#[traceevent(timestamp)]
timestamp_ns: u64,
realtime_ns: u64,
}
#[derive(TraceEvent)]
struct PollStartEvent {
#[traceevent(timestamp)]
timestamp_ns: u64,
worker_id: u64,
task_id: u64,
}
#[derive(TraceEvent)]
struct PollEndEvent {
#[traceevent(timestamp)]
timestamp_ns: u64,
worker_id: u64,
}
let mut encoder = Encoder::new();
encoder
.write(&ClockSyncEvent {
timestamp_ns: 100,
realtime_ns: 1,
})
.unwrap();
encoder
.write(&PollStartEvent {
timestamp_ns: 10,
worker_id: 1,
task_id: 7,
})
.unwrap();
encoder
.write(&PollEndEvent {
timestamp_ns: 20,
worker_id: 1,
})
.unwrap();
let (_, _, polls, _) = decode_samples(&encoder.into_inner(), "test").unwrap();
assert_eq!(polls.len(), 1);
assert_eq!(polls[0].start_ns, 10);
assert_eq!(polls[0].end_ns, 20);
}
#[test]
fn test_worker_id_inferred_from_park_unpark() {
let decompressed = load_demo_trace();
let (samples, _, _, _) = decode_samples(&decompressed, "test").unwrap();
let worker_samples = samples.iter().filter(|s| s.worker_id.is_some()).count();
assert!(
worker_samples > 0,
"expected some samples attributed to a worker via tid correlation"
);
eprintln!(
"{} of {} samples attributed to a worker (worker_id = Some)",
worker_samples,
samples.len()
);
}
#[test]
fn test_decode_real_trace() {
let path = "/tmp/dial9-ingest-test/2026-06-19/1459/shale/ip-10-2-123-116.us-west-2.compute.internal/kxgw-1/1781881195-9725.bin.gz";
if !std::path::Path::new(path).exists() {
eprintln!("skipping: real trace not available");
return;
}
let compressed = std::fs::read(path).unwrap();
let decompressed = {
use std::io::Read;
let mut dec = flate2::read::GzDecoder::new(compressed.as_slice());
let mut buf = Vec::new();
dec.read_to_end(&mut buf).unwrap();
buf
};
let (samples, stacks, _polls, _spans) = decode_samples(&decompressed, path).unwrap();
eprintln!(
"decoded {} samples, {} unique stacks",
samples.len(),
stacks.len()
);
assert!(!samples.is_empty(), "expected CPU samples in real trace");
assert!(!stacks.is_empty(), "expected stacks in dictionary");
}
#[test]
fn test_span_resolution_produces_valid_uids() {
let uid1 = compute_span_uid("boot-abc", 42);
let uid2 = compute_span_uid("boot-abc", 42);
assert_eq!(uid1, uid2, "span_uid must be deterministic");
let uid3 = compute_span_uid("boot-abc", 43);
assert_ne!(
uid1, uid3,
"different instance_ids must produce different uids"
);
let type_uid1 = compute_span_type_uid(
"tracing",
"my_crate",
"handle_request",
Some("src/main.rs"),
Some(10),
);
let type_uid2 = compute_span_type_uid(
"tracing",
"my_crate",
"handle_request",
Some("src/main.rs"),
Some(10),
);
assert_eq!(type_uid1, type_uid2, "span_type_uid must be deterministic");
let type_uid3 = compute_span_type_uid(
"tracing",
"my_crate",
"other_fn",
Some("src/main.rs"),
Some(20),
);
assert_ne!(
type_uid1, type_uid3,
"different names must produce different type_uids"
);
}
#[test]
fn test_cross_source_identity_same_boot_id() {
let uid1 = compute_span_uid("boot-abc", 42);
let uid2 = compute_span_uid("boot-abc", 42);
assert_eq!(
uid1, uid2,
"same boot-id + instance_id must produce same span_uid across segments"
);
}
#[test]
fn test_cross_source_identity_different_boot_id() {
let uid1 = compute_span_uid("boot-abc", 42);
let uid2 = compute_span_uid("boot-xyz", 42);
assert_ne!(
uid1, uid2,
"different boot-ids must produce different span_uids"
);
}
#[test]
fn test_recycled_raw_ids_do_not_collide() {
let uid1 = compute_span_uid("boot-1", 100);
let uid2 = compute_span_uid("boot-1", 200);
assert_ne!(uid1, uid2, "different instance_ids must never collide");
}
#[test]
fn test_samples_last_ordering() {
let source_key = "2026-06-19/1300/svc/host/boot/0.bin";
let empty_trace = {
use dial9_trace_format::encoder::Encoder;
let enc = Encoder::new();
enc.into_inner()
};
let result = decode_samples(&empty_trace, source_key);
assert!(result.is_ok());
let (samples, _stacks, _polls, spans) = result.unwrap();
assert!(samples.is_empty());
assert!(spans.is_empty());
}
#[test]
fn test_union_intervals_helper() {
assert_eq!(union_intervals(&[]), 0);
assert_eq!(union_intervals(&[(10, 20)]), 10);
assert_eq!(union_intervals(&[(10, 20), (30, 40)]), 20);
assert_eq!(union_intervals(&[(10, 30), (20, 40)]), 30);
assert_eq!(union_intervals(&[(10, 40), (15, 25)]), 30);
assert_eq!(union_intervals(&[(10, 20), (20, 30)]), 20);
assert_eq!(union_intervals(&[(10, 20), (15, 25), (22, 35)]), 25);
}
#[test]
fn test_resolve_span_task_binary_search() {
let polls = [(0, 100, 11), (200, 300, 22), (400, 500, 33)];
assert_eq!(resolve_span_task(&polls, 250), Some(22));
assert_eq!(resolve_span_task(&polls, 400), Some(33));
assert_eq!(resolve_span_task(&polls, 100), None);
assert_eq!(resolve_span_task(&polls, 200), Some(22));
assert_eq!(resolve_span_task(&polls, 150), None);
assert_eq!(resolve_span_task(&polls, 600), None);
assert_eq!(resolve_span_task(&[], 10), None);
}
#[test]
fn test_attribute_legacy_span_from_polls() {
let (on_cpu, wait) =
attribute_legacy_span_from_polls(&[(0, 1000)], &[(100, 200), (500, 600)]);
assert_eq!(on_cpu, 200);
assert_eq!(wait, 800);
let (on_cpu, wait) =
attribute_legacy_span_from_polls(&[(300, 700)], &[(100, 400), (600, 900)]);
assert_eq!(on_cpu, 100 + 100); assert_eq!(wait, 400 - 200);
let (on_cpu, wait) = attribute_legacy_span_from_polls(&[(0, 500)], &[]);
assert_eq!(on_cpu, 0);
assert_eq!(wait, 500);
let (on_cpu, wait) = attribute_legacy_span_from_polls(&[(0, 500)], &[(0, 500)]);
assert_eq!(on_cpu, 500);
assert_eq!(wait, 0);
let (on_cpu, wait) =
attribute_legacy_span_from_polls(&[(0, 300), (200, 500)], &[(100, 400)]);
assert_eq!(on_cpu, 300);
assert_eq!(wait, 200);
assert_eq!(on_cpu + wait, 500);
}
#[test]
fn test_extract_boot_id_from_path() {
assert_eq!(
extract_boot_id_from_path("2026-06-19/1300/svc/host/boot-abc/0-0.bin.gz"),
"boot-abc"
);
assert_eq!(
extract_boot_id_from_path(
"s3://bucket/traces/2026-06-19/1300/svc/host/my-boot/file.bin"
),
"my-boot"
);
assert_eq!(extract_boot_id_from_path("file.bin"), "file.bin");
}
#[test]
fn test_extract_boot_id_from_path_qualified() {
let (bid, namespaced) =
extract_boot_id_from_path_qualified("2026-06-19/1300/svc/host/abcd-12345/0-0.bin.gz");
assert_eq!(bid, "abcd-12345");
assert!(namespaced, "4-alpha-digits should be namespaced");
let (bid, namespaced) = extract_boot_id_from_path_qualified(
"s3://bucket/traces/2026-06-19/1300/svc/host/qmxz-481/file.bin",
);
assert_eq!(bid, "qmxz-481");
assert!(namespaced, "S3 URI with valid boot_id should be namespaced");
let (bid, namespaced) =
extract_boot_id_from_path_qualified("2026-06-19/1300/svc/host/my-boot/file.bin");
assert_eq!(bid, "my-boot");
assert!(!namespaced, "my-boot is not a valid boot_id format");
let (bid, namespaced) = extract_boot_id_from_path_qualified("file.bin");
assert_eq!(bid, "file.bin");
assert!(!namespaced, "flat path is not namespaced");
let (bid, namespaced) = extract_boot_id_from_path_qualified("some/dir/abcdef/file.bin");
assert_eq!(bid, "abcdef");
assert!(!namespaced, "no dash means not boot_id format");
}
#[test]
fn test_is_boot_id_format() {
assert!(is_boot_id_format("abcd-123"));
assert!(is_boot_id_format("qmxz-481"));
assert!(is_boot_id_format("zzzz-99999"));
assert!(!is_boot_id_format("abc-123")); assert!(!is_boot_id_format("abcde-123")); assert!(!is_boot_id_format("ABCD-123")); assert!(!is_boot_id_format("abcd-")); assert!(!is_boot_id_format("abcd")); assert!(!is_boot_id_format("my-boot")); assert!(!is_boot_id_format("")); assert!(!is_boot_id_format("1234-5678")); }
#[test]
fn test_decode_samples_metadata_identity_quality() {
use dial9_trace_format::TraceEvent;
use dial9_trace_format::encoder::Encoder;
#[derive(TraceEvent)]
struct ClockSyncEvent {
#[traceevent(timestamp)]
timestamp_ns: u64,
realtime_ns: u64,
}
#[derive(TraceEvent)]
struct SegmentMetadataEvent {
#[traceevent(timestamp)]
timestamp_ns: u64,
entries: Vec<(String, String)>,
}
let build_trace_with_metadata = |boot_id: Option<&str>| -> Vec<u8> {
let mut enc = Encoder::new();
enc.write(&ClockSyncEvent {
timestamp_ns: 100,
realtime_ns: 1_700_000_000_000_000_000 + 100,
})
.unwrap();
if let Some(bid) = boot_id {
enc.write(&SegmentMetadataEvent {
timestamp_ns: 101,
entries: vec![("boot_id".to_string(), bid.to_string())],
})
.unwrap();
}
enc.into_inner()
};
let data_with = build_trace_with_metadata(Some("test-boot-abc123"));
let (_, _, _, spans_with) = decode_samples(
&data_with,
"2026-06-19/1300/svc/host/some-path-boot/0.bin.gz",
)
.unwrap();
assert!(spans_with.is_empty());
let data_without = build_trace_with_metadata(None);
let result = decode_samples(
&data_without,
"2026-06-19/1300/svc/host/some-path-boot/0.bin.gz",
);
assert!(result.is_ok());
}
#[test]
fn test_parse_legacy_span_schema_name() {
let info = parse_legacy_span_schema_name(
"SpanEnter:metrics_service::routes::record_metric:examples/metrics-service/src/routes.rs:26",
).unwrap();
assert_eq!(info.target, "metrics_service::routes");
assert_eq!(info.name, "record_metric");
assert_eq!(
info.file.as_deref(),
Some("examples/metrics-service/src/routes.rs")
);
assert_eq!(info.line, Some(26));
let info = parse_legacy_span_schema_name(
"SpanExit:metrics_service::ddb::query_metric:examples/metrics-service/src/ddb.rs:122",
)
.unwrap();
assert_eq!(info.target, "metrics_service::ddb");
assert_eq!(info.name, "query_metric");
assert_eq!(
info.file.as_deref(),
Some("examples/metrics-service/src/ddb.rs")
);
assert_eq!(info.line, Some(122));
let info =
parse_legacy_span_schema_name("SpanEnter:a::b::c::d::my_span:src/lib.rs:99").unwrap();
assert_eq!(info.target, "a::b::c::d");
assert_eq!(info.name, "my_span");
assert_eq!(info.file.as_deref(), Some("src/lib.rs"));
assert_eq!(info.line, Some(99));
let info = parse_legacy_span_schema_name("SpanEnter__ShaleOperation").unwrap();
assert_eq!(info.target, "");
assert_eq!(info.name, "ShaleOperation");
assert_eq!(info.file, None);
assert_eq!(info.line, None);
assert!(parse_legacy_span_schema_name("SpanEnter").is_none());
assert!(parse_legacy_span_schema_name("SpanEnter:a::b:file").is_none());
}
#[test]
fn test_find_first_single_colon() {
assert_eq!(find_first_single_colon("a:b:c"), Some(1));
assert_eq!(
find_first_single_colon(
"metrics_service::routes::record_metric:examples/src/routes.rs"
),
Some(38) );
assert_eq!(find_first_single_colon("a::b::c"), None);
assert_eq!(find_first_single_colon("abc:"), Some(3));
assert_eq!(find_first_single_colon(""), None);
}
#[test]
fn parse_legacy_span_schema_name_preserves_windows_path() {
let info = parse_legacy_span_schema_name(r"SpanEnter:svc::op:C:\src\lib.rs:42").unwrap();
assert_eq!(info.target, "svc");
assert_eq!(info.name, "op");
assert_eq!(info.file.as_deref(), Some(r"C:\src\lib.rs"));
assert_eq!(info.line, Some(42));
}
#[test]
fn test_decode_demo_trace_legacy_spans() {
let decompressed = load_demo_trace();
let (samples, _stacks, _polls, spans) = decode_samples(
&decompressed,
"2026-01-01/1300/svc/host/demo-boot/demo-trace.bin",
)
.unwrap();
assert!(
!spans.is_empty(),
"expected legacy span rows from demo trace, got 0"
);
for span in &spans {
assert_eq!(
span.identity_quality, "legacy",
"demo trace spans must have identity_quality='legacy'"
);
assert!(
!span.details_complete,
"legacy spans must have details_complete=false"
);
}
let record_metric_spans: Vec<_> = spans
.iter()
.filter(|s| s.name == "record_metric" && s.target.contains("routes"))
.collect();
assert!(
!record_metric_spans.is_empty(),
"expected record_metric spans from demo trace"
);
let query_metric_spans: Vec<_> =
spans.iter().filter(|s| s.name == "query_metric").collect();
assert!(
!query_metric_spans.is_empty(),
"expected query_metric spans from demo trace"
);
let sample_span = &record_metric_spans[0];
assert!(
sample_span.target.contains("metrics_service"),
"target should contain metrics_service, got: {}",
sample_span.target
);
assert!(
sample_span.callsite_file.is_some(),
"callsite_file should be parsed from schema name"
);
assert!(
sample_span.callsite_line.is_some(),
"callsite_line should be parsed from schema name"
);
let spans_with_elapsed: Vec<_> = spans.iter().filter(|s| s.elapsed_ns > 0).collect();
assert!(
!spans_with_elapsed.is_empty(),
"expected some spans with non-zero elapsed_ns"
);
let spans_with_active: Vec<_> = spans
.iter()
.filter(|s| s.observed_active_wall_ns > 0)
.collect();
assert!(
!spans_with_active.is_empty(),
"expected some spans with observed active wall time"
);
let samples_with_spans = samples
.iter()
.filter(|s| !s.enclosing_spans.is_empty())
.count();
assert!(
samples_with_spans > 0,
"expected some samples attributed to legacy spans, got 0"
);
eprintln!(
"decoded {} legacy spans ({} record_metric, {} query_metric), {} samples with span attribution",
spans.len(),
record_metric_spans.len(),
query_metric_spans.len(),
samples_with_spans,
);
}
#[test]
fn test_legacy_span_reconstruction_synthetic() {
use dial9_trace_format::TraceEvent;
use dial9_trace_format::encoder::{Encoder, Schema};
use dial9_trace_format::schema::FieldDef;
use dial9_trace_format::types::{FieldType, FieldValue};
#[derive(TraceEvent)]
struct ClockSyncEvent {
#[traceevent(timestamp)]
timestamp_ns: u64,
realtime_ns: u64,
}
let enter_schema = Schema::new(
"SpanEnter:my_crate::handler::do_work:src/handler.rs:42",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("parent_span_id", FieldType::OptionalVarint),
FieldDef::new("span_name", FieldType::String),
],
);
let exit_schema = Schema::new(
"SpanExit:my_crate::handler::do_work:src/handler.rs:42",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("span_name", FieldType::String),
],
);
let close_schema = Schema::new(
"SpanCloseEvent",
vec![FieldDef::new("span_id", FieldType::Varint)],
);
let mut enc = Encoder::new();
enc.write(&ClockSyncEvent {
timestamp_ns: 10,
realtime_ns: 1_700_000_000_000_000_010,
})
.unwrap();
enc.write_event(
&enter_schema,
&[
FieldValue::Varint(100), FieldValue::Varint(0), FieldValue::Varint(1), FieldValue::None, FieldValue::String("do_work".to_string()), ],
)
.unwrap();
enc.write_event(
&exit_schema,
&[
FieldValue::Varint(200), FieldValue::Varint(0), FieldValue::Varint(1), FieldValue::String("do_work".to_string()), ],
)
.unwrap();
enc.write_event(
&close_schema,
&[
FieldValue::Varint(250), FieldValue::Varint(1), ],
)
.unwrap();
let data = enc.into_inner();
let source_key = "2026-06-19/1300/svc/host/test-boot/0.bin";
let (_, _, _, spans) = decode_samples(&data, source_key).unwrap();
assert_eq!(spans.len(), 1, "should produce exactly one legacy span");
let span = &spans[0];
assert_eq!(span.name, "do_work");
assert_eq!(span.target, "my_crate::handler");
assert_eq!(span.callsite_file.as_deref(), Some("src/handler.rs"));
assert_eq!(span.callsite_line, Some(42));
assert_eq!(span.identity_quality, "legacy");
assert!(!span.details_complete);
assert_eq!(span.kind, "tracing");
assert!(span.elapsed_ns > 0, "elapsed_ns should be > 0");
assert_eq!(span.observed_active_wall_ns, 100);
}
#[test]
fn test_legacy_span_captures_user_attributes() {
use dial9_trace_format::TraceEvent;
use dial9_trace_format::encoder::{Encoder, Schema};
use dial9_trace_format::schema::FieldDef;
use dial9_trace_format::types::{FieldType, FieldValue};
#[derive(TraceEvent)]
struct ClockSyncEvent {
#[traceevent(timestamp)]
timestamp_ns: u64,
realtime_ns: u64,
}
let enter_schema = Schema::new(
"SpanEnter:svc::routes::handle:src/routes.rs:7",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("parent_span_id", FieldType::OptionalVarint),
FieldDef::new("span_name", FieldType::String),
FieldDef::new("request_id", FieldType::String),
FieldDef::new("status_code", FieldType::Varint),
],
);
let exit_schema = Schema::new(
"SpanExit:svc::routes::handle:src/routes.rs:7",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("span_name", FieldType::String),
FieldDef::new("request_id", FieldType::String),
FieldDef::new("status_code", FieldType::Varint),
],
);
let close_schema = Schema::new(
"SpanCloseEvent",
vec![FieldDef::new("span_id", FieldType::Varint)],
);
let mut enc = Encoder::new();
enc.write(&ClockSyncEvent {
timestamp_ns: 10,
realtime_ns: 1_700_000_000_000_000_010,
})
.unwrap();
enc.write_event(
&enter_schema,
&[
FieldValue::Varint(100),
FieldValue::Varint(0),
FieldValue::Varint(1),
FieldValue::None,
FieldValue::String("handle".to_string()),
FieldValue::String("5d051ec2-999b-4a25-93b6-0f9cf83fa8b2".to_string()),
FieldValue::Varint(0), ],
)
.unwrap();
enc.write_event(
&exit_schema,
&[
FieldValue::Varint(200),
FieldValue::Varint(0),
FieldValue::Varint(1),
FieldValue::String("handle".to_string()),
FieldValue::String("5d051ec2-999b-4a25-93b6-0f9cf83fa8b2".to_string()),
FieldValue::Varint(500), ],
)
.unwrap();
enc.write_event(
&close_schema,
&[FieldValue::Varint(250), FieldValue::Varint(1)],
)
.unwrap();
let (_, _, _, spans) =
decode_samples(&enc.into_inner(), "2026-07-17/1746/svc/host/boot/0.bin").unwrap();
assert_eq!(spans.len(), 1);
let attrs: std::collections::HashMap<&str, &str> = spans[0]
.attributes
.iter()
.map(|(k, v)| (k.as_str(), v.as_str()))
.collect();
assert_eq!(
attrs.get("request_id"),
Some(&"5d051ec2-999b-4a25-93b6-0f9cf83fa8b2")
);
assert_eq!(attrs.get("status_code"), Some(&"500"));
assert!(!attrs.contains_key("worker_id"));
assert!(!attrs.contains_key("span_id"));
assert!(!attrs.contains_key("span_name"));
}
#[test]
fn test_legacy_equal_timestamp_exit_before_enter_stays_unbalanced() {
use dial9_trace_format::TraceEvent;
use dial9_trace_format::encoder::{Encoder, Schema};
use dial9_trace_format::schema::FieldDef;
use dial9_trace_format::types::{FieldType, FieldValue};
#[derive(TraceEvent)]
struct ClockSyncEvent {
#[traceevent(timestamp)]
timestamp_ns: u64,
realtime_ns: u64,
}
let enter_schema = Schema::new(
"SpanEnter__EqualTimestamp",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("parent_span_id", FieldType::OptionalVarint),
FieldDef::new("span_name", FieldType::String),
],
);
let exit_schema = Schema::new(
"SpanExit__EqualTimestamp",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("span_name", FieldType::String),
],
);
let close_schema = Schema::new(
"SpanCloseEvent",
vec![FieldDef::new("span_id", FieldType::Varint)],
);
let mut enc = Encoder::new();
enc.write(&ClockSyncEvent {
timestamp_ns: 10,
realtime_ns: 1_700_000_000_000_000_010,
})
.unwrap();
enc.write_event(
&exit_schema,
&[
FieldValue::Varint(100),
FieldValue::Varint(0),
FieldValue::Varint(7),
FieldValue::String("equal_timestamp".to_string()),
],
)
.unwrap();
enc.write_event(
&enter_schema,
&[
FieldValue::Varint(100),
FieldValue::Varint(0),
FieldValue::Varint(7),
FieldValue::None,
FieldValue::String("equal_timestamp".to_string()),
],
)
.unwrap();
enc.write_event(
&close_schema,
&[FieldValue::Varint(101), FieldValue::Varint(7)],
)
.unwrap();
let (_, _, _, spans) = decode_samples(
&enc.into_inner(),
"2026-07-15/1714/svc/host/test-boot/0.bin",
)
.unwrap();
assert_eq!(spans.len(), 1);
let span = &spans[0];
assert_eq!(span.unbalanced_exits, 1);
assert_eq!(span.unbalanced_enters, 1);
assert_eq!(span.observed_active_wall_ns, 0);
assert!(!span.details_complete);
}
#[test]
fn test_struct_derived_span_name_convention() {
use dial9_trace_format::TraceEvent;
use dial9_trace_format::encoder::{Encoder, Schema};
use dial9_trace_format::schema::FieldDef;
use dial9_trace_format::types::{FieldType, FieldValue};
#[derive(TraceEvent)]
struct ClockSyncEvent {
#[traceevent(timestamp)]
timestamp_ns: u64,
realtime_ns: u64,
}
let enter_schema = Schema::new(
"SpanEnter__ShaleOperation",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("parent_span_id", FieldType::OptionalVarint),
FieldDef::new("span_name", FieldType::String),
FieldDef::new("request_id", FieldType::String),
],
);
let exit_schema = Schema::new(
"SpanExit__ShaleOperation",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("span_name", FieldType::String),
],
);
let mut enc = Encoder::new();
enc.write(&ClockSyncEvent {
timestamp_ns: 10,
realtime_ns: 1_700_000_000_000_000_010,
})
.unwrap();
enc.write_event(
&enter_schema,
&[
FieldValue::Varint(1000),
FieldValue::Varint(3),
FieldValue::Varint(42),
FieldValue::None,
FieldValue::String("/jobs/next".to_string()),
FieldValue::String("req-abc".to_string()),
],
)
.unwrap();
enc.write_event(
&exit_schema,
&[
FieldValue::Varint(5000),
FieldValue::Varint(3),
FieldValue::Varint(42),
FieldValue::String("/jobs/next".to_string()),
],
)
.unwrap();
let data = enc.into_inner();
let source_key = "2026-07-15/1714/shale/host/test-boot/0.bin";
let (_, _, _, spans) = decode_samples(&data, source_key).unwrap();
assert_eq!(
spans.len(),
1,
"struct-derived (__) span events must be picked up, not dropped"
);
let span = &spans[0];
assert_eq!(span.name, "/jobs/next");
assert_eq!(span.kind, "tracing");
assert_eq!(span.identity_quality, "legacy");
assert_eq!(span.observed_active_wall_ns, 4000);
assert!(span.elapsed_ns > 0, "elapsed_ns should be > 0");
}
#[test]
fn test_struct_derived_schema_suffix_disambiguates_shared_runtime_name() {
use dial9_trace_format::encoder::{Encoder, Schema};
use dial9_trace_format::schema::FieldDef;
use dial9_trace_format::types::{FieldType, FieldValue};
fn enter_schema(name: &'static str) -> Schema {
Schema::new(
name,
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("parent_span_id", FieldType::OptionalVarint),
FieldDef::new("span_name", FieldType::String),
],
)
}
fn exit_schema(name: &'static str) -> Schema {
Schema::new(
name,
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("span_name", FieldType::String),
],
)
}
let first_enter = enter_schema("SpanEnter__FirstOperation");
let first_exit = exit_schema("SpanExit__FirstOperation");
let second_enter = enter_schema("SpanEnter__SecondOperation");
let second_exit = exit_schema("SpanExit__SecondOperation");
let mut enc = Encoder::new();
for (schema, timestamp, span_id) in [
(&first_enter, 100, 1),
(&first_exit, 200, 1),
(&second_enter, 300, 2),
(&second_exit, 400, 2),
] {
let values = if schema.name().starts_with("SpanEnter__") {
vec![
FieldValue::Varint(timestamp),
FieldValue::Varint(0),
FieldValue::Varint(span_id),
FieldValue::None,
FieldValue::String("shared-runtime-name".to_string()),
]
} else {
vec![
FieldValue::Varint(timestamp),
FieldValue::Varint(0),
FieldValue::Varint(span_id),
FieldValue::String("shared-runtime-name".to_string()),
]
};
enc.write_event(schema, &values).unwrap();
}
let (_, _, _, spans) = decode_samples(
&enc.into_inner(),
"2026-07-15/1714/svc/host/test-boot/0.bin",
)
.unwrap();
assert_eq!(spans.len(), 2);
assert!(spans.iter().all(|span| span.name == "shared-runtime-name"));
assert_ne!(
spans[0].span_type_uid, spans[1].span_type_uid,
"struct schema suffixes must remain distinct type identities"
);
}
#[test]
fn test_struct_derived_schema_preserves_distinct_runtime_names() {
use dial9_trace_format::encoder::{Encoder, Schema};
use dial9_trace_format::schema::FieldDef;
use dial9_trace_format::types::{FieldType, FieldValue};
let enter = Schema::new(
"SpanEnter__SharedOperation",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("parent_span_id", FieldType::OptionalVarint),
FieldDef::new("span_name", FieldType::String),
],
);
let exit = Schema::new(
"SpanExit__SharedOperation",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("span_name", FieldType::String),
],
);
let mut enc = Encoder::new();
for (timestamp, span_id, runtime_name, schema) in [
(100, 1, "/jobs", &enter),
(200, 1, "/jobs", &exit),
(300, 2, "/jobs/next", &enter),
(400, 2, "/jobs/next", &exit),
] {
let values = if schema.name().starts_with("SpanEnter__") {
vec![
FieldValue::Varint(timestamp),
FieldValue::Varint(0),
FieldValue::Varint(span_id),
FieldValue::None,
FieldValue::String(runtime_name.to_string()),
]
} else {
vec![
FieldValue::Varint(timestamp),
FieldValue::Varint(0),
FieldValue::Varint(span_id),
FieldValue::String(runtime_name.to_string()),
]
};
enc.write_event(schema, &values).unwrap();
}
let (_, _, _, spans) = decode_samples(
&enc.into_inner(),
"2026-07-15/1714/svc/host/test-boot/0.bin",
)
.unwrap();
assert_eq!(spans.len(), 2);
assert_ne!(spans[0].name, spans[1].name);
assert_ne!(
spans[0].span_type_uid, spans[1].span_type_uid,
"runtime names sharing one struct schema must remain distinct type identities"
);
}
#[test]
fn test_legacy_span_survives_worker_migration() {
use dial9_trace_format::TraceEvent;
use dial9_trace_format::encoder::{Encoder, Schema};
use dial9_trace_format::schema::FieldDef;
use dial9_trace_format::types::{FieldType, FieldValue};
#[derive(TraceEvent)]
struct ClockSyncEvent {
#[traceevent(timestamp)]
timestamp_ns: u64,
realtime_ns: u64,
}
let enter_schema = Schema::new(
"SpanEnter__ShaleOperation",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("parent_span_id", FieldType::OptionalVarint),
FieldDef::new("span_name", FieldType::String),
],
);
let exit_schema = Schema::new(
"SpanExit__ShaleOperation",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("span_name", FieldType::String),
],
);
let mut enc = Encoder::new();
enc.write(&ClockSyncEvent {
timestamp_ns: 10,
realtime_ns: 1_700_000_000_000_000_010,
})
.unwrap();
enc.write_event(
&enter_schema,
&[
FieldValue::Varint(1000),
FieldValue::Varint(3),
FieldValue::Varint(42),
FieldValue::None,
FieldValue::String("/jobs/next".to_string()),
],
)
.unwrap();
enc.write_event(
&exit_schema,
&[
FieldValue::Varint(5000),
FieldValue::Varint(7),
FieldValue::Varint(42),
FieldValue::String("/jobs/next".to_string()),
],
)
.unwrap();
let data = enc.into_inner();
let source_key = "2026-07-15/1714/shale/host/test-boot/0.bin";
let (_, _, _, spans) = decode_samples(&data, source_key).unwrap();
assert_eq!(
spans.len(),
1,
"span whose task migrated workers must still be paired, not dropped"
);
assert_eq!(spans[0].observed_active_wall_ns, 4000);
}
#[test]
fn test_legacy_span_cpu_wait_attribution_from_polls() {
use dial9_trace_format::TraceEvent;
use dial9_trace_format::encoder::{Encoder, Schema};
use dial9_trace_format::schema::FieldDef;
use dial9_trace_format::types::{FieldType, FieldValue};
#[derive(TraceEvent)]
struct ClockSyncEvent {
#[traceevent(timestamp)]
timestamp_ns: u64,
realtime_ns: u64,
}
let unpark_schema = Schema::new(
"WorkerUnparkEvent",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("local_queue", FieldType::Varint),
FieldDef::new("cpu_time_ns", FieldType::Varint),
FieldDef::new("sched_wait_ns", FieldType::OptionalVarint),
FieldDef::new("tid", FieldType::Varint),
],
);
let poll_start_schema = Schema::new(
"PollStartEvent",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("local_queue", FieldType::Varint),
FieldDef::new("task_id", FieldType::Varint),
FieldDef::new("spawn_loc", FieldType::String),
],
);
let poll_end_schema = Schema::new(
"PollEndEvent",
vec![FieldDef::new("worker_id", FieldType::Varint)],
);
let enter_schema = Schema::new(
"SpanEnter__ShaleOperation",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("parent_span_id", FieldType::OptionalVarint),
FieldDef::new("span_name", FieldType::String),
],
);
let exit_schema = Schema::new(
"SpanExit__ShaleOperation",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("span_name", FieldType::String),
],
);
let mut enc = Encoder::new();
enc.write(&ClockSyncEvent {
timestamp_ns: 10,
realtime_ns: 1_700_000_000_000_000_010,
})
.unwrap();
enc.write_event(
&unpark_schema,
&[
FieldValue::Varint(500), FieldValue::Varint(3), FieldValue::Varint(0), FieldValue::Varint(0), FieldValue::None, FieldValue::Varint(500), ],
)
.unwrap();
enc.write_event(
&poll_start_schema,
&[
FieldValue::Varint(900), FieldValue::Varint(3), FieldValue::Varint(0), FieldValue::Varint(77), FieldValue::String("app::handler".to_string()),
],
)
.unwrap();
enc.write_event(
&poll_end_schema,
&[FieldValue::Varint(1100), FieldValue::Varint(3)],
)
.unwrap();
enc.write_event(
&enter_schema,
&[
FieldValue::Varint(1000),
FieldValue::Varint(3),
FieldValue::Varint(42),
FieldValue::None,
FieldValue::String("/jobs/next".to_string()),
],
)
.unwrap();
enc.write_event(
&poll_start_schema,
&[
FieldValue::Varint(5000),
FieldValue::Varint(3),
FieldValue::Varint(0),
FieldValue::Varint(77),
FieldValue::String("app::handler".to_string()),
],
)
.unwrap();
enc.write_event(
&poll_end_schema,
&[FieldValue::Varint(5100), FieldValue::Varint(3)],
)
.unwrap();
enc.write_event(
&exit_schema,
&[
FieldValue::Varint(9000),
FieldValue::Varint(3),
FieldValue::Varint(42),
FieldValue::String("/jobs/next".to_string()),
],
)
.unwrap();
let data = enc.into_inner();
let source_key = "2026-07-15/1714/shale/host/test-boot/0.bin";
let (_, _, _, spans) = decode_samples(&data, source_key).unwrap();
assert_eq!(spans.len(), 1);
let s = &spans[0];
assert_eq!(s.observed_active_wall_ns, 8000);
assert_eq!(s.on_cpu_ns_est, Some(200));
assert_eq!(s.async_wait_ns, Some(7800));
let sum = s.on_cpu_ns_est.unwrap_or(0)
+ s.blocked_ns_est.unwrap_or(0)
+ s.async_wait_ns.unwrap_or(0)
+ s.scheduler_delay_ns.unwrap_or(0)
+ s.unknown_ns;
assert_eq!(
sum, s.elapsed_ns,
"five-way attribution must sum to elapsed"
);
assert_eq!(s.attribution_flags & 0b0100, 0, "task-resolved bit cleared");
}
#[test]
fn test_legacy_recycled_span_ids() {
use dial9_trace_format::TraceEvent;
use dial9_trace_format::encoder::{Encoder, Schema};
use dial9_trace_format::schema::FieldDef;
use dial9_trace_format::types::{FieldType, FieldValue};
#[derive(TraceEvent)]
struct ClockSyncEvent {
#[traceevent(timestamp)]
timestamp_ns: u64,
realtime_ns: u64,
}
let enter_schema = Schema::new(
"SpanEnter:app::op:src/lib.rs:10",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("parent_span_id", FieldType::OptionalVarint),
FieldDef::new("span_name", FieldType::String),
],
);
let exit_schema = Schema::new(
"SpanExit:app::op:src/lib.rs:10",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("span_name", FieldType::String),
],
);
let close_schema = Schema::new(
"SpanCloseEvent",
vec![FieldDef::new("span_id", FieldType::Varint)],
);
let mut enc = Encoder::new();
enc.write(&ClockSyncEvent {
timestamp_ns: 10,
realtime_ns: 1_700_000_000_000_000_010,
})
.unwrap();
enc.write_event(
&enter_schema,
&[
FieldValue::Varint(100),
FieldValue::Varint(0),
FieldValue::Varint(1),
FieldValue::None,
FieldValue::String("op".to_string()),
],
)
.unwrap();
enc.write_event(
&exit_schema,
&[
FieldValue::Varint(200),
FieldValue::Varint(0),
FieldValue::Varint(1),
FieldValue::String("op".to_string()),
],
)
.unwrap();
enc.write_event(
&close_schema,
&[FieldValue::Varint(250), FieldValue::Varint(1)],
)
.unwrap();
let data = enc.into_inner();
let (_, _, _, spans) = decode_samples(&data, "test/path/boot/0.bin").unwrap();
assert_eq!(spans.len(), 1);
assert_eq!(spans[0].identity_quality, "legacy");
}
#[test]
fn test_legacy_recycled_span_ids_are_close_delimited() {
use dial9_trace_format::TraceEvent;
use dial9_trace_format::encoder::{Encoder, Schema};
use dial9_trace_format::schema::FieldDef;
use dial9_trace_format::types::{FieldType, FieldValue};
#[derive(TraceEvent)]
struct ClockSyncEvent {
#[traceevent(timestamp)]
timestamp_ns: u64,
realtime_ns: u64,
}
let enter_schema = Schema::new(
"SpanEnter:app::request::handle:src/lib.rs:10",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("parent_span_id", FieldType::OptionalVarint),
FieldDef::new("span_name", FieldType::String),
],
);
let exit_schema = Schema::new(
"SpanExit:app::request::handle:src/lib.rs:10",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("span_name", FieldType::String),
],
);
let close_schema = Schema::new(
"SpanCloseEvent",
vec![FieldDef::new("span_id", FieldType::Varint)],
);
let mut enc = Encoder::new();
enc.write(&ClockSyncEvent {
timestamp_ns: 10,
realtime_ns: 1_700_000_000_000_000_010,
})
.unwrap();
let enter = |enc: &mut Encoder, ts: u64| {
enc.write_event(
&enter_schema,
&[
FieldValue::Varint(ts),
FieldValue::Varint(0),
FieldValue::Varint(1),
FieldValue::None,
FieldValue::String("handle".to_string()),
],
)
.unwrap();
};
let exit = |enc: &mut Encoder, ts: u64| {
enc.write_event(
&exit_schema,
&[
FieldValue::Varint(ts),
FieldValue::Varint(0),
FieldValue::Varint(1),
FieldValue::String("handle".to_string()),
],
)
.unwrap();
};
let close = |enc: &mut Encoder, ts: u64| {
enc.write_event(
&close_schema,
&[FieldValue::Varint(ts), FieldValue::Varint(1)],
)
.unwrap();
};
enter(&mut enc, 1_000);
exit(&mut enc, 1_100);
close(&mut enc, 1_100);
enter(&mut enc, 10_000_001_000);
exit(&mut enc, 10_000_001_100);
close(&mut enc, 10_000_001_100);
let data = enc.into_inner();
let (_, _, _, mut spans) = decode_samples(&data, "test/path/boot/0.bin").unwrap();
assert_eq!(
spans.len(),
2,
"each close-delimited reuse of a recycled span_id must be its own span"
);
spans.sort_by_key(|s| s.start_ns);
for span in &spans {
assert_eq!(
span.elapsed_ns, 100,
"each instance must keep its own short lifecycle, not the 10s span between reuses"
);
assert_eq!(span.observed_active_wall_ns, 100);
assert_eq!(span.name, "handle");
}
assert_ne!(spans[0].span_uid, spans[1].span_uid);
}
}
#[cfg(test)]
mod decode_test;
#[cfg(test)]
mod parser_parity_test;