use std::collections::{BTreeMap, HashMap, HashSet};
use std::io::{BufWriter, Read, Write};
use std::path::Path;
use anyhow::{Context, bail, ensure};
#[cfg(test)]
use dial9_trace_format::decoder::DecodedFrame;
use dial9_trace_format::decoder::Decoder;
use dial9_trace_format::encoder::{Encoder, Schema};
use dial9_trace_format::schema::{FieldAnnotation, FieldDef, SchemaEntry};
#[cfg(test)]
use dial9_trace_format::types::FieldValueRef;
use dial9_trace_format::types::{FieldType, FieldValue};
use serde::{Deserialize, Serialize};
const SHAPE_VERSION: u32 = 2;
const QUANTUM_NS: u64 = 10_000;
const SYNTHETIC_EPOCH_NS: u64 = 1_577_836_800_000_000_000;
const GZIP_MAGIC: [u8; 2] = [0x1f, 0x8b];
const SAFE_ANNOTATION_KEYS: &[&str] = &["unit", "metrique.unit", "kind"];
const SAFE_UNITS: &[&str] = &["ns", "us", "ms", "s", "bytes", "count"];
const SAFE_FIELD_KINDS: &[&str] = &["gauge", "counter", "updown-counter"];
fn is_safe_annotation_value(key: &str, value: &str) -> bool {
match key {
"unit" | "metrique.unit" => SAFE_UNITS.contains(&value),
"kind" => SAFE_FIELD_KINDS.contains(&value),
_ => false,
}
}
const MAX_DYNAMIC_DEPTH: usize = 8;
const MAX_SHAPE_JSON_BYTES: u64 = 32 * 1024 * 1024;
const MAX_RAW_TRACE_BYTES: u64 = 2 * 1024 * 1024 * 1024;
const MAX_BYTES_LENGTH: u32 = 64 * 1024 * 1024;
const MAX_CONTAINER_ELEMENTS: usize = 1_000_000;
const MAX_GENERATED_EVENTS: u64 = 100_000_000;
const MAX_GENERATED_OUTPUT_BYTES: u64 = 16 * 1024 * 1024 * 1024;
const MAX_RETAINED_GENERATION_BYTES: u64 = 512 * 1024 * 1024;
const MAX_WIRE_SCHEMAS: usize = (u16::MAX - dial9_trace_format::STATIC_WIRE_ID_LIMIT) as usize;
const MAX_SCHEMA_ANNOTATIONS: usize = u16::MAX as usize;
const BUILTIN_SCHEMAS: &[&str] = &[
"PollStartEvent",
"PollEndEvent",
"WorkerParkEvent",
"WorkerUnparkEvent",
"QueueSampleEvent",
"TaskSpawnEvent",
"TaskTerminateEvent",
"WakeEventEvent",
"CpuSampleEvent",
"TaskDumpEvent",
"AllocEvent",
"FreeEvent",
"MemoryProfileOverflowEvent",
"ClockSyncEvent",
"SegmentMetadataEvent",
"SymbolTableEntry",
"ProcessResourceUsageEvent",
"TcpAcceptQueueEvent",
"SpanCloseEvent",
];
type BuiltinFieldSignature = &'static [(&'static str, FieldType)];
fn builtin_signatures(schema_name: &str) -> Option<&'static [BuiltinFieldSignature]> {
use FieldType::{
Bool as B, OptionalPooledString as OPS, OptionalVarint as OV, PooledStackFrames as PStack,
PooledString as PS, StackFrames as Stack, String as S, StringMap as SM, U8, U16, U32,
Varint as V,
};
match schema_name {
"PollStartEvent" => Some(&[
&[
("worker_id", V),
("local_queue", U8),
("task_id", V),
("spawn_loc", PS),
],
&[
("worker_id", V),
("local_queue", U8),
("task_id", U32),
("spawn_loc", PS),
],
&[
("worker_id", V),
("local_queue", V),
("task_id", V),
("spawn_loc", PS),
],
]),
"PollEndEvent" => Some(&[&[("worker_id", V)]]),
"WorkerParkEvent" => Some(&[
&[
("worker_id", V),
("local_queue", U8),
("cpu_time_ns", V),
("tid", U32),
],
&[("worker_id", V), ("local_queue", U8), ("cpu_time_ns", V)],
&[
("worker_id", V),
("local_queue", V),
("cpu_time_ns", V),
("tid", V),
],
&[("worker_id", V), ("local_queue", V), ("cpu_time_ns", V)],
]),
"WorkerUnparkEvent" => Some(&[
&[
("worker_id", V),
("local_queue", U8),
("cpu_time_ns", V),
("sched_wait_ns", OV),
("tid", U32),
],
&[
("worker_id", V),
("local_queue", U8),
("cpu_time_ns", V),
("sched_wait_ns", V),
("tid", U32),
],
&[
("worker_id", V),
("local_queue", U8),
("cpu_time_ns", V),
("sched_wait_ns", V),
],
&[
("worker_id", V),
("local_queue", V),
("cpu_time_ns", V),
("sched_wait_ns", OV),
("tid", V),
],
&[
("worker_id", V),
("local_queue", V),
("cpu_time_ns", V),
("sched_wait_ns", V),
("tid", V),
],
&[
("worker_id", V),
("local_queue", V),
("cpu_time_ns", V),
("sched_wait_ns", V),
],
]),
"QueueSampleEvent" => Some(&[
&[("global_queue", U8), ("active_tasks", V)],
&[("global_queue", V), ("active_tasks", V)],
&[("global_queue", U8)],
&[("global_queue", V)],
]),
"TaskSpawnEvent" => Some(&[
&[("task_id", V), ("spawn_loc", PS), ("instrumented", B)],
&[("task_id", V), ("spawn_loc", PS)],
&[("task_id", U32), ("spawn_loc", PS)],
]),
"TaskTerminateEvent" => Some(&[&[("task_id", V)], &[("task_id", U32)]]),
"WakeEventEvent" => Some(&[
&[
("waker_task_id", V),
("woken_task_id", V),
("target_worker", U8),
],
&[
("waker_task_id", U32),
("woken_task_id", U32),
("target_worker", U8),
],
&[
("waker_task_id", V),
("woken_task_id", V),
("target_worker", V),
],
]),
"CpuSampleEvent" => Some(&[
&[
("worker_id", V),
("tid", U32),
("source", U8),
("thread_name", OPS),
("callchain", PStack),
("cpu", OV),
],
&[
("worker_id", V),
("tid", U32),
("source", U8),
("thread_name", OPS),
("callchain", PStack),
],
&[
("worker_id", V),
("tid", U32),
("source", U8),
("thread_name", OPS),
("callchain", Stack),
],
&[
("worker_id", V),
("tid", U32),
("source", U8),
("thread_name", PS),
("callchain", Stack),
],
&[
("worker_id", V),
("tid", V),
("source", V),
("thread_name", OPS),
("callchain", PStack),
("cpu", OV),
],
&[
("worker_id", V),
("tid", V),
("source", V),
("thread_name", OPS),
("callchain", PStack),
],
&[
("worker_id", V),
("tid", V),
("source", V),
("thread_name", OPS),
("callchain", Stack),
],
&[
("worker_id", V),
("tid", V),
("source", V),
("thread_name", PS),
("callchain", Stack),
],
]),
"TaskDumpEvent" => Some(&[&[("task_id", V), ("callchain", PStack)]]),
"AllocEvent" => Some(&[
&[
("tid", U32),
("size", V),
("addr", V),
("callchain", PStack),
],
&[("tid", V), ("size", V), ("addr", V), ("callchain", PStack)],
]),
"FreeEvent" => Some(&[
&[
("tid", U32),
("addr", V),
("size", V),
("alloc_timestamp_ns", V),
],
&[
("tid", V),
("addr", V),
("size", V),
("alloc_timestamp_ns", V),
],
]),
"MemoryProfileOverflowEvent" => Some(&[&[("dropped_allocs", V), ("dropped_frees", V)]]),
"ClockSyncEvent" => Some(&[&[("realtime_ns", V)]]),
"SegmentMetadataEvent" => Some(&[&[("entries", SM)]]),
"SymbolTableEntry" => Some(&[
&[
("addr", V),
("size", V),
("symbol_name", PS),
("inline_depth", V),
("source_file", PS),
("source_line", V),
],
&[
("addr", V),
("size", V),
("symbol_name", PS),
("inline_depth", V),
],
&[("base_addr", V), ("size", V), ("symbol_name", PS)],
]),
"ProcessResourceUsageEvent" => Some(&[&[
("user_cpu_ns", V),
("system_cpu_ns", V),
("max_rss_bytes", V),
("minor_faults", V),
("major_faults", V),
("block_input_ops", V),
("block_output_ops", V),
("voluntary_context_switches", V),
("involuntary_context_switches", V),
]]),
"TcpAcceptQueueEvent" => Some(&[
&[
("socket_cookie", V),
("socket_inode", V),
("ip_version", U8),
("local_addr", S),
("local_port", U16),
("pending_connections", U32),
("backlog_limit", U32),
],
&[
("socket_cookie", V),
("socket_inode", V),
("ip_version", V),
("local_addr", S),
("local_port", V),
("pending_connections", V),
("backlog_limit", V),
],
]),
"SpanCloseEvent" => Some(&[&[("span_id", V)]]),
_ => None,
}
}
fn validate_builtin_schema_signature(schema_name: &str, fields: &[FieldDef]) -> anyhow::Result<()> {
let signatures = builtin_signatures(schema_name).ok_or_else(|| {
anyhow::anyhow!("builtin schema '{schema_name}' has no canonical signature")
})?;
let matches = signatures.iter().any(|signature| {
signature.len() == fields.len()
&& signature
.iter()
.zip(fields)
.all(|((name, field_type), field)| {
*name == field.name() && *field_type == field.field_type()
})
});
ensure!(
matches,
"builtin schema '{schema_name}' signature {:?} does not exactly match a known complete signature {:?}",
fields
.iter()
.map(|field| (field.name(), field.field_type()))
.collect::<Vec<_>>(),
signatures
);
Ok(())
}
fn is_known_builtin_field(schema_name: &str, field_name: &str) -> bool {
builtin_signatures(schema_name).is_some_and(|signatures| {
signatures
.iter()
.any(|signature| signature.iter().any(|(name, _)| *name == field_name))
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub(crate) enum NamespaceId {
Task,
Span,
Tid,
Addr,
SocketCookie,
SocketInode,
Anon(u16),
}
impl serde::Serialize for NamespaceId {
fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
match self {
NamespaceId::Task => serializer.serialize_u16(0),
NamespaceId::Span => serializer.serialize_u16(1),
NamespaceId::Tid => serializer.serialize_u16(2),
NamespaceId::Addr => serializer.serialize_u16(3),
NamespaceId::SocketCookie => serializer.serialize_u16(4),
NamespaceId::SocketInode => serializer.serialize_u16(5),
NamespaceId::Anon(n) => serializer.serialize_u16(100 + n),
}
}
}
impl<'de> serde::Deserialize<'de> for NamespaceId {
fn deserialize<D: serde::Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
let n = u16::deserialize(deserializer)?;
Ok(match n {
0 => NamespaceId::Task,
1 => NamespaceId::Span,
2 => NamespaceId::Tid,
3 => NamespaceId::Addr,
4 => NamespaceId::SocketCookie,
5 => NamespaceId::SocketInode,
other if other >= 100 => NamespaceId::Anon(other - 100),
_ => {
return Err(serde::de::Error::custom(format!(
"namespace id {n} is in reserved range (6..99)"
)));
}
})
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub(crate) enum FieldSemantics {
Structural,
Identity,
TimestampRef,
RealtimeOffset,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct FieldRepeatMeta {
pub semantics: FieldSemantics,
#[serde(skip_serializing_if = "Option::is_none")]
pub namespace: Option<NamespaceId>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct TraceShape {
pub version: u32,
pub summary: ShapeSummary,
pub schemas: Vec<ShapeSchema>,
pub events: Vec<ShapeEvent>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct ShapeSummary {
pub event_count: u64,
pub duration_ns: u64,
pub event_type_counts: BTreeMap<String, u64>,
pub worker_cardinality: u32,
pub task_cardinality: u32,
#[serde(skip_serializing_if = "Option::is_none")]
pub poll_duration_quantiles: Option<Quantiles>,
#[serde(skip_serializing_if = "Option::is_none")]
pub span_duration_quantiles: Option<Quantiles>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct Quantiles {
pub count: u64,
pub min_ns: u64,
pub p50_ns: u64,
pub p90_ns: u64,
pub p99_ns: u64,
pub max_ns: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct ShapeSchema {
pub name: String,
pub has_timestamp: bool,
pub fields: Vec<ShapeField>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub annotations: Vec<ShapeAnnotation>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct ShapeField {
pub name: String,
pub field_type: u8,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub repeat_meta: Option<FieldRepeatMeta>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct ShapeAnnotation {
pub field_index: u16,
pub key: String,
pub value: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct ShapeEvent {
pub schema_index: u32,
#[serde(skip_serializing_if = "Option::is_none")]
pub timestamp_offset_ns: Option<u64>,
pub values: Vec<ShapeValue>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "t", content = "v")]
pub(crate) enum ShapeValue {
U(u64),
I(i64),
F(f64),
B(bool),
S(String),
PS(String),
Bytes(u32),
Stack(Vec<u64>),
PStack(Vec<u64>),
StringMap(Vec<(String, String)>),
List(Vec<ShapeValue>),
Map(Vec<(ShapeValue, ShapeValue)>),
None,
}
pub(crate) fn extract(trace_path: &Path, output_path: &Path) -> anyhow::Result<()> {
let data = read_trace_file(trace_path)?;
let shape = extract_shape(&data)?;
let json = serde_json::to_string_pretty(&shape).context("serialize shape JSON")?;
std::fs::write(output_path, json.as_bytes())
.with_context(|| format!("write {}", output_path.display()))?;
Ok(())
}
pub(crate) fn generate(shape_path: &Path, output_path: &Path, repeat: u32) -> anyhow::Result<()> {
ensure!(repeat >= 1, "--repeat must be >= 1, got {repeat}");
let metadata =
std::fs::metadata(shape_path).with_context(|| format!("stat {}", shape_path.display()))?;
ensure!(
metadata.len() <= MAX_SHAPE_JSON_BYTES,
"shape JSON exceeds maximum size ({} bytes > {MAX_SHAPE_JSON_BYTES} byte limit)",
metadata.len()
);
let file = std::fs::File::open(shape_path)
.with_context(|| format!("open {}", shape_path.display()))?;
let take_limit = MAX_SHAPE_JSON_BYTES.saturating_add(1);
let mut bounded = std::io::BufReader::new(file).take(take_limit);
let mut json = String::new();
bounded
.read_to_string(&mut json)
.with_context(|| format!("read {}", shape_path.display()))?;
ensure!(
(json.len() as u64) <= MAX_SHAPE_JSON_BYTES,
"shape JSON actual read size ({} bytes) exceeds limit ({MAX_SHAPE_JSON_BYTES} byte limit)",
json.len()
);
let shape: TraceShape = serde_json::from_str(&json).context("parse shape JSON")?;
validate_shape(&shape)?;
write_generated_trace(&shape, output_path, repeat)
}
pub(crate) fn synthesize(trace_path: &Path, output_path: &Path, repeat: u32) -> anyhow::Result<()> {
ensure!(repeat >= 1, "--repeat must be >= 1, got {repeat}");
let data = read_trace_file(trace_path)?;
let shape = extract_shape(&data)?;
drop(data);
write_generated_trace(&shape, output_path, repeat)
}
fn write_generated_trace(
shape: &TraceShape,
output_path: &Path,
repeat: u32,
) -> anyhow::Result<()> {
validate_repeat_preflight(shape, repeat)?;
let file = std::fs::File::create(output_path)
.with_context(|| format!("create {}", output_path.display()))?;
let writer = BufWriter::new(file);
let limited = LimitedWriter::new(writer, MAX_GENERATED_OUTPUT_BYTES);
generate_to_writer(shape, repeat, limited)
.with_context(|| format!("write {}", output_path.display()))?;
Ok(())
}
struct LimitedWriter<W> {
inner: W,
written: u64,
limit: u64,
}
impl<W> LimitedWriter<W> {
fn new(inner: W, limit: u64) -> Self {
Self {
inner,
written: 0,
limit,
}
}
}
impl<W: Write> Write for LimitedWriter<W> {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
let new_total = self.written.saturating_add(buf.len() as u64);
if new_total > self.limit {
return Err(std::io::Error::other(format!(
"generated output exceeds hard limit ({} bytes written + {} = {new_total} > {} limit)",
self.written,
buf.len(),
self.limit
)));
}
let n = self.inner.write(buf)?;
self.written += n as u64;
Ok(n)
}
fn flush(&mut self) -> std::io::Result<()> {
self.inner.flush()
}
}
fn read_trace_file(path: &Path) -> anyhow::Result<Vec<u8>> {
let metadata = std::fs::metadata(path).with_context(|| format!("stat {}", path.display()))?;
let file_size = metadata.len();
ensure!(
file_size <= MAX_RAW_TRACE_BYTES,
"trace file exceeds maximum size ({file_size} bytes > {MAX_RAW_TRACE_BYTES} byte limit)"
);
let mut file = std::fs::File::open(path).with_context(|| format!("open {}", path.display()))?;
let mut magic = [0u8; 2];
let magic_len = std::io::Read::read(&mut file, &mut magic)
.with_context(|| format!("read magic from {}", path.display()))?;
if magic_len >= 2 && magic == GZIP_MAGIC {
use std::io::Seek;
file.seek(std::io::SeekFrom::Start(0))
.with_context(|| format!("seek {}", path.display()))?;
read_gzip_bounded(file, MAX_RAW_TRACE_BYTES, MAX_RAW_TRACE_BYTES)
} else {
use std::io::Seek;
file.seek(std::io::SeekFrom::Start(0))
.with_context(|| format!("seek {}", path.display()))?;
let take_limit = MAX_RAW_TRACE_BYTES.saturating_add(1);
let mut bounded = file.take(take_limit);
let mut raw = Vec::with_capacity(file_size as usize);
bounded
.read_to_end(&mut raw)
.with_context(|| format!("read {}", path.display()))?;
ensure!(
(raw.len() as u64) <= MAX_RAW_TRACE_BYTES,
"trace file actual read size ({} bytes) exceeds limit ({MAX_RAW_TRACE_BYTES} byte limit)",
raw.len()
);
Ok(raw)
}
}
fn read_gzip_bounded<R: Read>(
reader: R,
compressed_limit: u64,
decompressed_limit: u64,
) -> anyhow::Result<Vec<u8>> {
let compressed_take_limit = compressed_limit.saturating_add(1);
let decompressed_take_limit = decompressed_limit.saturating_add(1);
let compressed = reader.take(compressed_take_limit);
let decoder = flate2::read::GzDecoder::new(compressed);
let mut bounded_output = decoder.take(decompressed_take_limit);
let mut out = Vec::new();
bounded_output
.read_to_end(&mut out)
.context("decompress gzip trace")?;
let decoder = bounded_output.into_inner();
let compressed = decoder.into_inner();
let compressed_consumed = compressed_take_limit.saturating_sub(compressed.limit());
ensure!(
compressed_consumed <= compressed_limit,
"compressed trace exceeds maximum size ({compressed_consumed} bytes consumed > {compressed_limit} byte limit)"
);
ensure!(
(out.len() as u64) <= decompressed_limit,
"decompressed trace exceeds maximum size ({} bytes > {} byte limit)",
out.len(),
decompressed_limit
);
Ok(out)
}
fn normalize_field_type(ft: FieldType) -> FieldType {
match ft {
FieldType::U8 | FieldType::U16 | FieldType::U32 => FieldType::Varint,
FieldType::OptionalU8 | FieldType::OptionalU16 | FieldType::OptionalU32 => {
FieldType::OptionalVarint
}
other => other,
}
}
fn quantize_ns(ns: u64) -> u64 {
(ns / QUANTUM_NS) * QUANTUM_NS
}
fn quantize_numeric(val: u64) -> u64 {
if val == 0 {
return 0;
}
let bits = 64u32.saturating_sub(val.leading_zeros());
if bits <= 10 {
return val; }
let shift = bits - 10; let top = val >> shift;
let rounded = if shift >= 54 {
top
} else {
top.saturating_add(1)
};
rounded.checked_shl(shift).unwrap_or(u64::MAX)
}
fn quantize_i64(val: i64) -> i64 {
if val == i64::MIN {
let mag = quantize_numeric(val.unsigned_abs());
-(mag.min(i64::MAX as u64) as i64)
} else {
let sign = val.signum();
let mag = quantize_numeric(val.unsigned_abs());
sign * (mag.min(i64::MAX as u64) as i64)
}
}
fn quantize_f64(val: f64) -> anyhow::Result<f64> {
ensure!(
val.is_finite(),
"non-finite f64 value cannot be stored in shape JSON"
);
if val == 0.0 {
return Ok(0.0);
}
let magnitude = val.abs().log10().floor() as i32;
let rounded = if magnitude < -306 {
format!("{val:.2e}")
.parse::<f64>()
.context("round very small finite f64")?
} else {
let factor = 10f64.powi(3 - 1 - magnitude);
(val * factor).round() / factor
};
ensure!(
rounded.is_finite(),
"finite f64 quantization produced a non-finite result"
);
Ok(rounded)
}
fn privacy_bucket_u64(val: u64) -> u64 {
if val == 0 {
return 0;
}
let bits = 64u32.saturating_sub(val.leading_zeros()); let bucket = 1u64.checked_shl(bits.saturating_sub(1)).unwrap_or(u64::MAX);
if val <= 2 {
return val + 2;
}
let offset = (bucket >> 1).saturating_add(bucket >> 2); let result = bucket.saturating_add(offset).saturating_add(1);
if result == val {
result.saturating_add(bucket >> 2)
} else {
result
}
}
fn privacy_bucket_i64(val: i64) -> i64 {
if val == 0 {
return 0;
}
let sign = val.signum();
let mag = privacy_bucket_u64(val.unsigned_abs());
sign * (mag.min(i64::MAX as u64) as i64)
}
fn is_builtin_schema(name: &str) -> bool {
BUILTIN_SCHEMAS.contains(&name)
}
fn span_prefix(name: &str) -> Option<(&str, &str)> {
if let Some(suffix) = name.strip_prefix("SpanEnter:") {
Some(("SpanEnter:", suffix))
} else if let Some(suffix) = name.strip_prefix("SpanExit:") {
Some(("SpanExit:", suffix))
} else {
None
}
}
fn classify_field(schema_name: &str, field_name: &str, ft: FieldType) -> FieldRepeatMeta {
let inner = normalize_field_type(ft).inner();
if field_name == "worker_id" || field_name == "target_worker" {
return FieldRepeatMeta {
semantics: FieldSemantics::Structural,
namespace: None,
};
}
if field_name == "alloc_timestamp_ns" {
return FieldRepeatMeta {
semantics: FieldSemantics::TimestampRef,
namespace: None,
};
}
if schema_name == "ClockSyncEvent" && field_name == "realtime_ns" {
return FieldRepeatMeta {
semantics: FieldSemantics::RealtimeOffset,
namespace: None,
};
}
let is_varint = matches!(inner, FieldType::Varint);
if is_varint {
if matches!(field_name, "task_id" | "waker_task_id" | "woken_task_id") {
return FieldRepeatMeta {
semantics: FieldSemantics::Identity,
namespace: Some(NamespaceId::Task),
};
}
if matches!(field_name, "span_id" | "parent_span_id") {
return FieldRepeatMeta {
semantics: FieldSemantics::Identity,
namespace: Some(NamespaceId::Span),
};
}
if field_name == "thread_id" || field_name == "tid" {
return FieldRepeatMeta {
semantics: FieldSemantics::Identity,
namespace: Some(NamespaceId::Tid),
};
}
if matches!(
field_name,
"address" | "alloc_address" | "symbol_address" | "addr" | "base_addr"
) {
return FieldRepeatMeta {
semantics: FieldSemantics::Identity,
namespace: Some(NamespaceId::Addr),
};
}
if field_name == "socket_cookie" {
return FieldRepeatMeta {
semantics: FieldSemantics::Identity,
namespace: Some(NamespaceId::SocketCookie),
};
}
if field_name == "socket_inode" {
return FieldRepeatMeta {
semantics: FieldSemantics::Identity,
namespace: Some(NamespaceId::SocketInode),
};
}
if field_name.ends_with("_id") && field_name != "worker_id" {
return FieldRepeatMeta {
semantics: FieldSemantics::Identity,
namespace: None, };
}
}
if matches!(inner, FieldType::StackFrames | FieldType::PooledStackFrames) {
return FieldRepeatMeta {
semantics: FieldSemantics::Identity,
namespace: Some(NamespaceId::Addr),
};
}
FieldRepeatMeta {
semantics: FieldSemantics::Structural,
namespace: None,
}
}
struct ExtractContext {
callsite_names: HashMap<String, String>,
callsite_counter: u32,
schema_names: HashMap<String, String>,
custom_schema_counter: u32,
field_names: HashMap<(String, String), String>,
custom_field_counter: u32,
anon_namespace_map: HashMap<String, u16>,
next_anon_namespace: u16,
task_ids: HashMap<u64, u64>,
next_task_id: u64,
span_ids: HashMap<u64, u64>,
next_span_id: u64,
thread_ids: HashMap<u64, u64>,
next_thread_id: u64,
addresses: HashMap<u64, u64>,
next_address: u64,
generic_ids: HashMap<(u16, u64), u64>,
next_generic_id: HashMap<u16, u64>,
socket_cookies: HashMap<u64, u64>,
next_socket_cookie: u64,
socket_inodes: HashMap<u64, u64>,
next_socket_inode: u64,
string_remap: HashMap<String, String>,
next_string_id: u64,
}
impl ExtractContext {
fn new() -> Self {
Self {
callsite_names: HashMap::new(),
callsite_counter: 0,
schema_names: HashMap::new(),
custom_schema_counter: 0,
field_names: HashMap::new(),
custom_field_counter: 0,
anon_namespace_map: HashMap::new(),
next_anon_namespace: 0,
task_ids: HashMap::new(),
next_task_id: 1,
span_ids: HashMap::new(),
next_span_id: 1,
thread_ids: HashMap::new(),
next_thread_id: 1,
addresses: HashMap::new(),
next_address: 0x1000,
generic_ids: HashMap::new(),
next_generic_id: HashMap::new(),
socket_cookies: HashMap::new(),
next_socket_cookie: 1,
socket_inodes: HashMap::new(),
next_socket_inode: 1,
string_remap: HashMap::new(),
next_string_id: 0,
}
}
fn sanitize_schema_name(&mut self, name: &str) -> String {
if is_builtin_schema(name) {
return name.to_string();
}
if let Some((prefix, suffix)) = span_prefix(name) {
let callsite = self
.callsite_names
.entry(suffix.to_string())
.or_insert_with(|| {
let id = self.callsite_counter;
self.callsite_counter += 1;
format!("callsite_{id:04}")
});
return format!("{prefix}{callsite}");
}
self.schema_names
.entry(name.to_string())
.or_insert_with(|| {
let id = self.custom_schema_counter;
self.custom_schema_counter += 1;
format!("CustomEvent_{id:04}")
})
.clone()
}
fn sanitize_field_name(
&mut self,
schema_name: &str,
original_schema_name: &str,
field_name: &str,
) -> String {
if is_builtin_schema(original_schema_name) {
if builtin_signatures(original_schema_name).is_some() {
if is_known_builtin_field(original_schema_name, field_name) {
return field_name.to_string();
}
} else {
return field_name.to_string();
}
} else if span_prefix(original_schema_name).is_some() {
if matches!(
field_name,
"worker_id" | "span_id" | "parent_span_id" | "span_name" | "task_id"
) {
return field_name.to_string();
}
}
let key = (schema_name.to_string(), field_name.to_string());
self.field_names
.entry(key)
.or_insert_with(|| {
let id = self.custom_field_counter;
self.custom_field_counter += 1;
format!("field_{id:04}")
})
.clone()
}
fn get_anon_namespace(&mut self, field_name: &str) -> anyhow::Result<NamespaceId> {
if let Some(&ns_num) = self.anon_namespace_map.get(field_name) {
return Ok(NamespaceId::Anon(ns_num));
}
const MAX_ANON_NAMESPACES: u16 = u16::MAX - 100;
ensure!(
self.next_anon_namespace < MAX_ANON_NAMESPACES,
"too many anonymous namespaces ({} >= {MAX_ANON_NAMESPACES}); \
trace has too many distinct *_id field names",
self.next_anon_namespace
);
let n = self.next_anon_namespace;
self.next_anon_namespace = n
.checked_add(1)
.ok_or_else(|| anyhow::anyhow!("anonymous namespace counter overflow"))?;
self.anon_namespace_map.insert(field_name.to_string(), n);
Ok(NamespaceId::Anon(n))
}
fn remap_task_id(&mut self, id: u64) -> u64 {
if id == 0 {
return 0;
}
let next = &mut self.next_task_id;
*self.task_ids.entry(id).or_insert_with(|| {
let v = *next;
*next += 1;
v
})
}
fn remap_span_id(&mut self, id: u64) -> u64 {
if id == 0 {
return 0;
}
let next = &mut self.next_span_id;
*self.span_ids.entry(id).or_insert_with(|| {
let v = *next;
*next += 1;
v
})
}
fn remap_thread_id(&mut self, id: u64) -> u64 {
if id == 0 {
return 0;
}
let next = &mut self.next_thread_id;
*self.thread_ids.entry(id).or_insert_with(|| {
let v = *next;
*next += 1;
v
})
}
fn remap_address(&mut self, addr: u64) -> u64 {
if addr == 0 {
return 0;
}
let next = &mut self.next_address;
*self.addresses.entry(addr).or_insert_with(|| {
let v = *next;
*next += 0x10;
v
})
}
fn remap_socket_cookie(&mut self, id: u64) -> u64 {
if id == 0 {
return 0;
}
let next = &mut self.next_socket_cookie;
*self.socket_cookies.entry(id).or_insert_with(|| {
let v = *next;
*next += 1;
v
})
}
fn remap_socket_inode(&mut self, id: u64) -> u64 {
if id == 0 {
return 0;
}
let next = &mut self.next_socket_inode;
*self.socket_inodes.entry(id).or_insert_with(|| {
let v = *next;
*next += 1;
v
})
}
fn remap_generic_id(&mut self, ns: NamespaceId, id: u64) -> u64 {
if id == 0 {
return 0;
}
let ns_num = match ns {
NamespaceId::Anon(n) => n,
_ => return id, };
let key = (ns_num, id);
let counter = self.next_generic_id.entry(ns_num).or_insert(1);
*self.generic_ids.entry(key).or_insert_with(|| {
let v = *counter;
*counter += 1;
v
})
}
fn sanitize_string(&mut self, s: &str) -> String {
self.string_remap
.entry(s.to_string())
.or_insert_with(|| {
let id = self.next_string_id;
self.next_string_id += 1;
format!("s_{id:04}")
})
.clone()
}
#[allow(clippy::too_many_arguments)]
#[allow(clippy::collapsible_if)]
fn remap_value_ref(
&mut self,
schema_name: &str,
field_name: &str,
meta: &FieldRepeatMeta,
value: &dial9_trace_format::types::FieldValueRef<'_>,
ft: FieldType,
string_pool: &dial9_trace_format::decoder::StringPool,
stack_pool: &dial9_trace_format::decoder::StackPool,
source_base_ts: u64,
event_offset_ns: u64,
is_builtin: bool,
) -> anyhow::Result<ShapeValue> {
use dial9_trace_format::types::FieldValueRef;
if matches!(value, FieldValueRef::None) {
ensure!(
ft.is_optional(),
"None value for non-optional field '{field_name}'"
);
return Ok(ShapeValue::None);
}
let inner_ft = normalize_field_type(ft).inner();
match meta.semantics {
FieldSemantics::Structural => {
if field_name == "worker_id" || field_name == "target_worker" {
if let FieldValueRef::Varint(v) = value {
return Ok(ShapeValue::U(*v));
}
}
self.convert_structural_ref(value, inner_ft, string_pool, stack_pool, is_builtin, 0)
}
FieldSemantics::Identity => {
let ns = meta.namespace.unwrap_or(NamespaceId::Addr);
self.convert_identity_ref(schema_name, field_name, ns, value, inner_ft, stack_pool)
}
FieldSemantics::TimestampRef => {
if let FieldValueRef::Varint(v) = value {
let offset = v.saturating_sub(source_base_ts);
Ok(ShapeValue::U(quantize_ns(offset)))
} else {
bail!("TimestampRef field '{field_name}' must be Varint");
}
}
FieldSemantics::RealtimeOffset => {
if let FieldValueRef::Varint(_) = value {
Ok(ShapeValue::U(quantize_ns(event_offset_ns)))
} else {
bail!("RealtimeOffset field '{field_name}' must be Varint");
}
}
}
}
fn convert_structural_ref(
&mut self,
value: &dial9_trace_format::types::FieldValueRef<'_>,
_inner_ft: FieldType,
string_pool: &dial9_trace_format::decoder::StringPool,
stack_pool: &dial9_trace_format::decoder::StackPool,
is_builtin: bool,
depth: usize,
) -> anyhow::Result<ShapeValue> {
use dial9_trace_format::types::FieldValueRef;
ensure!(
depth <= MAX_DYNAMIC_DEPTH,
"excessive nesting depth ({depth}) in dynamic value"
);
match value {
FieldValueRef::Varint(v) => {
if is_builtin {
Ok(ShapeValue::U(quantize_numeric(*v)))
} else {
Ok(ShapeValue::U(privacy_bucket_u64(*v)))
}
}
FieldValueRef::I64(v) => {
if is_builtin {
Ok(ShapeValue::I(quantize_i64(*v)))
} else {
Ok(ShapeValue::I(privacy_bucket_i64(*v)))
}
}
FieldValueRef::F64(v) => Ok(ShapeValue::F(quantize_f64(*v)?)),
FieldValueRef::Bool(v) => Ok(ShapeValue::B(*v)),
FieldValueRef::String(s) => Ok(ShapeValue::S(self.sanitize_string(s))),
FieldValueRef::Bytes(b) => {
ensure!(
b.len() <= MAX_BYTES_LENGTH as usize,
"Bytes field length {} exceeds limit {MAX_BYTES_LENGTH}",
b.len()
);
Ok(ShapeValue::Bytes(b.len() as u32))
}
FieldValueRef::PooledString(id) => {
let resolved = string_pool
.get(*id)
.ok_or_else(|| anyhow::anyhow!("missing pooled string id {:?}", id.raw_id()))?;
Ok(ShapeValue::PS(self.sanitize_string(resolved)))
}
FieldValueRef::StackFrames(frames_ref) => {
ensure!(
(frames_ref.count() as usize) <= MAX_CONTAINER_ELEMENTS,
"StackFrames has {} elements, exceeds limit {MAX_CONTAINER_ELEMENTS}",
frames_ref.count()
);
let remapped: Vec<u64> = frames_ref.iter().map(|a| self.remap_address(a)).collect();
Ok(ShapeValue::Stack(remapped))
}
FieldValueRef::PooledStackFrames(id) => {
let frames = stack_pool
.get(*id)
.ok_or_else(|| anyhow::anyhow!("missing pooled stack id {:?}", id.raw_id()))?;
ensure!(
frames.len() <= MAX_CONTAINER_ELEMENTS,
"PooledStackFrames has {} elements, exceeds limit {MAX_CONTAINER_ELEMENTS}",
frames.len()
);
let remapped: Vec<u64> = frames.iter().map(|a| self.remap_address(*a)).collect();
Ok(ShapeValue::PStack(remapped))
}
FieldValueRef::StringMap(map_ref) => {
let count = map_ref.iter().count();
ensure!(
count <= MAX_CONTAINER_ELEMENTS,
"StringMap has {count} entries, exceeds limit {MAX_CONTAINER_ELEMENTS}"
);
let sanitized: Vec<(String, String)> = map_ref
.iter()
.map(|(k, v)| (self.sanitize_string(k), self.sanitize_string(v)))
.collect();
Ok(ShapeValue::StringMap(sanitized))
}
FieldValueRef::List(list_ref) => {
let count = list_ref.iter().count();
ensure!(
count <= MAX_CONTAINER_ELEMENTS,
"List has {count} elements, exceeds limit {MAX_CONTAINER_ELEMENTS}"
);
let mapped: Vec<ShapeValue> = list_ref
.iter()
.map(|item| {
self.convert_structural_ref(
item,
FieldType::Varint,
string_pool,
stack_pool,
is_builtin,
depth + 1,
)
})
.collect::<anyhow::Result<_>>()?;
Ok(ShapeValue::List(mapped))
}
FieldValueRef::Map(map_ref) => {
let count = map_ref.iter().count();
ensure!(
count <= MAX_CONTAINER_ELEMENTS,
"Map has {count} entries, exceeds limit {MAX_CONTAINER_ELEMENTS}"
);
let mapped: Vec<(ShapeValue, ShapeValue)> = map_ref
.iter()
.map(|(k, v)| {
let kv = self.convert_structural_ref(
k,
FieldType::Varint,
string_pool,
stack_pool,
is_builtin,
depth + 1,
)?;
let vv = self.convert_structural_ref(
v,
FieldType::Varint,
string_pool,
stack_pool,
is_builtin,
depth + 1,
)?;
Ok((kv, vv))
})
.collect::<anyhow::Result<_>>()?;
Ok(ShapeValue::Map(mapped))
}
FieldValueRef::None => Ok(ShapeValue::None),
_ => bail!("unsupported FieldValueRef variant for structural field"),
}
}
fn convert_identity_ref(
&mut self,
_schema_name: &str,
field_name: &str,
ns: NamespaceId,
value: &dial9_trace_format::types::FieldValueRef<'_>,
_inner_ft: FieldType,
stack_pool: &dial9_trace_format::decoder::StackPool,
) -> anyhow::Result<ShapeValue> {
use dial9_trace_format::types::FieldValueRef;
match value {
FieldValueRef::Varint(v) => {
let remapped = match ns {
NamespaceId::Task => self.remap_task_id(*v),
NamespaceId::Span => self.remap_span_id(*v),
NamespaceId::Tid => self.remap_thread_id(*v),
NamespaceId::Addr => self.remap_address(*v),
NamespaceId::SocketCookie => self.remap_socket_cookie(*v),
NamespaceId::SocketInode => self.remap_socket_inode(*v),
NamespaceId::Anon(_) => self.remap_generic_id(ns, *v),
};
Ok(ShapeValue::U(remapped))
}
FieldValueRef::StackFrames(frames_ref) => {
ensure!(
(frames_ref.count() as usize) <= MAX_CONTAINER_ELEMENTS,
"identity StackFrames for field '{field_name}' has {} elements, exceeds limit {MAX_CONTAINER_ELEMENTS}",
frames_ref.count()
);
let remapped: Vec<u64> = frames_ref.iter().map(|a| self.remap_address(a)).collect();
Ok(ShapeValue::Stack(remapped))
}
FieldValueRef::PooledStackFrames(id) => {
let frames = stack_pool.get(*id).ok_or_else(|| {
anyhow::anyhow!(
"missing pooled stack id {:?} for field '{field_name}'",
id.raw_id()
)
})?;
ensure!(
frames.len() <= MAX_CONTAINER_ELEMENTS,
"identity PooledStackFrames for field '{field_name}' has {} elements, exceeds limit {MAX_CONTAINER_ELEMENTS}",
frames.len()
);
let remapped: Vec<u64> = frames.iter().map(|a| self.remap_address(*a)).collect();
Ok(ShapeValue::PStack(remapped))
}
_ => bail!(
"identity field '{field_name}' has unexpected value type {:?}",
value
),
}
}
}
fn compute_quantiles(values: &mut [u64]) -> Option<Quantiles> {
if values.is_empty() {
return None;
}
values.sort_unstable();
let n = values.len();
Some(Quantiles {
count: n as u64,
min_ns: values[0],
p50_ns: values[n / 2],
p90_ns: values[n * 90 / 100],
p99_ns: values[n * 99 / 100],
max_ns: values[n - 1],
})
}
#[allow(clippy::collapsible_if)] pub(crate) fn extract_shape(data: &[u8]) -> anyhow::Result<TraceShape> {
let mut decoder1 = Decoder::new(data).context("invalid trace file header")?;
let mut global_min_ts: Option<u64> = None;
let mut global_max_ts: Option<u64> = None;
decoder1
.try_for_each_event(|ev| {
let ts = ev.timestamp_ns.ok_or_else(|| {
anyhow::anyhow!("event for schema '{}' has no timestamp", ev.name)
})?;
global_min_ts = Some(global_min_ts.map_or(ts, |m| m.min(ts)));
global_max_ts = Some(global_max_ts.map_or(ts, |m| m.max(ts)));
Ok::<(), anyhow::Error>(())
})
.map_err(|e| match e {
dial9_trace_format::decoder::TryForEachError::Decode(d) => {
anyhow::anyhow!("pass 1 decode error at byte {}: {}", d.pos, d.message)
}
dial9_trace_format::decoder::TryForEachError::User(u) => u.context("pass 1"),
})?;
ensure!(
decoder1.position() == decoder1.data_len(),
"pass 1: trailing {} bytes after last frame",
decoder1.data_len() - decoder1.position()
);
let source_base_ts = global_min_ts.unwrap_or(0);
let mut decoder2 = Decoder::new(data).context("invalid trace file header (pass 2)")?;
let mut ctx = ExtractContext::new();
let mut schema_index_map: HashMap<String, u32> = HashMap::new(); let mut schemas: Vec<ShapeSchema> = Vec::new();
let mut sanitized_defs: HashMap<String, ShapeSchema> = HashMap::new();
let mut original_fields: HashMap<String, Vec<(String, FieldType)>> = HashMap::new();
let mut events: Vec<ShapeEvent> = Vec::new();
let mut event_type_counts: BTreeMap<String, u64> = BTreeMap::new();
decoder2
.try_for_each_event(|ev| {
let original_name = ev.name;
{
let sanitized_name = ctx.sanitize_schema_name(original_name);
if is_builtin_schema(original_name) {
validate_builtin_schema_signature(original_name, ev.schema.fields())?;
}
let orig_flds: Vec<(String, FieldType)> = ev
.schema
.fields()
.iter()
.map(|f| (f.name().to_string(), f.field_type()))
.collect();
let fields: Vec<ShapeField> = ev
.schema
.fields()
.iter()
.map(|f| {
let fname =
ctx.sanitize_field_name(&sanitized_name, original_name, f.name());
let norm_ft = normalize_field_type(f.field_type());
let mut meta = classify_field(original_name, f.name(), f.field_type());
if meta.semantics == FieldSemantics::Identity && meta.namespace.is_none() {
meta.namespace = Some(ctx.get_anon_namespace(f.name())?);
}
Ok(ShapeField {
name: fname,
field_type: norm_ft as u8,
repeat_meta: if meta.semantics == FieldSemantics::Structural
&& meta.namespace.is_none()
{
None
} else {
Some(meta)
},
})
})
.collect::<anyhow::Result<Vec<_>>>()?;
let annotations: Vec<ShapeAnnotation> = ev
.schema
.annotations()
.iter()
.filter_map(|a| {
if SAFE_ANNOTATION_KEYS.contains(&a.key())
&& is_safe_annotation_value(a.key(), a.value())
{
if (a.field_index() as usize) >= fields.len() {
return None;
}
Some(ShapeAnnotation {
field_index: a.field_index(),
key: a.key().to_string(),
value: a.value().to_string(),
})
} else {
None
}
})
.collect();
let shape_schema = ShapeSchema {
name: sanitized_name.clone(),
has_timestamp: ev.schema.has_timestamp(),
fields,
annotations,
};
if let Some(existing) = sanitized_defs.get(&sanitized_name) {
let existing_json = serde_json::to_string(existing).unwrap_or_default();
let new_json = serde_json::to_string(&shape_schema).unwrap_or_default();
if existing_json != new_json {
return Err(anyhow::anyhow!(
"schema '{}' (sanitized '{}') re-registered after reset with \
different definition (field count/type/name/annotations/repeat \
metadata differ)",
original_name,
sanitized_name
));
}
if !schema_index_map.contains_key(original_name) {
let existing_idx = schemas
.iter()
.position(|s| s.name == sanitized_name)
.unwrap();
schema_index_map.insert(original_name.to_string(), existing_idx as u32);
}
} else {
let idx = schemas.len() as u32;
schema_index_map.insert(original_name.to_string(), idx);
sanitized_defs.insert(sanitized_name, shape_schema.clone());
schemas.push(shape_schema);
}
original_fields.insert(original_name.to_string(), orig_flds);
}
let schema_idx = schema_index_map[original_name];
let ts = ev.timestamp_ns.ok_or_else(|| {
anyhow::anyhow!("event for schema '{}' has no timestamp", original_name)
})?;
let offset = quantize_ns(ts.saturating_sub(source_base_ts));
let orig_flds = original_fields.get(original_name).ok_or_else(|| {
anyhow::anyhow!("missing original fields for '{}'", original_name)
})?;
let is_builtin = is_builtin_schema(original_name);
let mut remapped_values: Vec<ShapeValue> = Vec::with_capacity(ev.fields.len());
for (i, fvr) in ev.fields.iter().enumerate() {
let (fname, ftype) = orig_flds.get(i).ok_or_else(|| {
anyhow::anyhow!(
"field index {i} out of range for schema '{}'",
original_name
)
})?;
let mut meta = classify_field(original_name, fname, *ftype);
if meta.semantics == FieldSemantics::Identity && meta.namespace.is_none() {
meta.namespace = Some(ctx.get_anon_namespace(fname)?);
}
let sv = ctx.remap_value_ref(
original_name,
fname,
&meta,
fvr,
*ftype,
ev.string_pool,
ev.stack_pool,
source_base_ts,
offset,
is_builtin,
)?;
remapped_values.push(sv);
}
let sanitized_name = &schemas[schema_idx as usize].name;
*event_type_counts.entry(sanitized_name.clone()).or_default() += 1;
events.push(ShapeEvent {
schema_index: schema_idx,
timestamp_offset_ns: Some(offset),
values: remapped_values,
});
Ok::<(), anyhow::Error>(())
})
.map_err(|e| match e {
dial9_trace_format::decoder::TryForEachError::Decode(d) => {
anyhow::anyhow!("pass 2 decode error at byte {}: {}", d.pos, d.message)
}
dial9_trace_format::decoder::TryForEachError::User(u) => u.context("pass 2"),
})?;
ensure!(
decoder2.position() == decoder2.data_len(),
"pass 2: trailing {} bytes after last frame",
decoder2.data_len() - decoder2.position()
);
let duration_ns = match (global_min_ts, global_max_ts) {
(Some(mn), Some(mx)) => quantize_ns(mx.saturating_sub(mn)),
_ => 0,
};
let mut shape = TraceShape {
version: SHAPE_VERSION,
summary: ShapeSummary {
event_count: events.len() as u64,
duration_ns,
event_type_counts,
worker_cardinality: 0,
task_cardinality: 0,
poll_duration_quantiles: None,
span_duration_quantiles: None,
},
schemas,
events,
};
let recomputed = recompute_summary(&shape);
shape.summary.worker_cardinality = recomputed.worker_cardinality;
shape.summary.task_cardinality = recomputed.task_cardinality;
shape.summary.poll_duration_quantiles = recomputed.poll_duration_quantiles;
shape.summary.span_duration_quantiles = recomputed.span_duration_quantiles;
validate_shape(&shape).context("extracted shape failed validation")?;
Ok(shape)
}
pub(crate) fn retain_events(
shape: &mut TraceShape,
mut keep: impl FnMut(&ShapeSchema, &ShapeEvent) -> bool,
) -> anyhow::Result<()> {
let schemas = &shape.schemas;
shape
.events
.retain(|event| keep(&schemas[event.schema_index as usize], event));
ensure!(
!shape.events.is_empty(),
"event filter removed every event from the trace shape"
);
rebuild_summary(shape);
validate_shape(shape).context("filtered shape failed validation")
}
pub(crate) fn rebase_timeline_from_first_nonzero_event(
shape: &mut TraceShape,
) -> anyhow::Result<usize> {
let original_len = shape.events.len();
shape.events.retain(|event| {
let timestamp = event.timestamp_offset_ns;
timestamp.is_some_and(|timestamp| timestamp > 0)
|| (timestamp == Some(0)
&& shape.schemas[event.schema_index as usize].name == "SymbolTableEntry")
});
ensure!(
!shape.events.is_empty(),
"trace shape has no nonzero-timestamp events"
);
let mut base = shape
.events
.iter()
.filter_map(|event| event.timestamp_offset_ns)
.filter(|timestamp| *timestamp > 0)
.min()
.context("trace shape has no positive timestamped events")?;
for event in &shape.events {
let schema = &shape.schemas[event.schema_index as usize];
for (field, value) in schema.fields.iter().zip(event.values.iter()) {
if field_shifts_with_timeline(field)
&& let ShapeValue::U(value) = value
&& *value > 0
{
base = base.min(*value);
}
}
}
for event in &mut shape.events {
let timestamp = event
.timestamp_offset_ns
.context("trace shape event lost its timestamp")?;
if timestamp > 0 {
event.timestamp_offset_ns = Some(
timestamp
.checked_sub(base)
.context("trace shape timestamp precedes selected timeline base")?,
);
}
let schema = &shape.schemas[event.schema_index as usize];
for (field, value) in schema.fields.iter().zip(event.values.iter_mut()) {
if field_shifts_with_timeline(field)
&& let ShapeValue::U(value) = value
&& *value > 0
{
*value = value
.checked_sub(base)
.context("trace shape timestamp reference precedes timeline base")?;
}
}
}
rebuild_summary(shape);
validate_shape(shape).context("rebased shape failed validation")?;
Ok(original_len - shape.events.len())
}
fn field_shifts_with_timeline(field: &ShapeField) -> bool {
field.repeat_meta.as_ref().is_some_and(|metadata| {
matches!(
metadata.semantics,
FieldSemantics::TimestampRef | FieldSemantics::RealtimeOffset
)
})
}
fn rebuild_summary(shape: &mut TraceShape) {
let mut event_type_counts = BTreeMap::new();
for event in &shape.events {
let name = &shape.schemas[event.schema_index as usize].name;
*event_type_counts.entry(name.clone()).or_default() += 1;
}
let recomputed = recompute_summary(shape);
shape.summary = ShapeSummary {
event_count: recomputed.event_count,
duration_ns: recomputed.duration_ns,
event_type_counts,
worker_cardinality: recomputed.worker_cardinality,
task_cardinality: recomputed.task_cardinality,
poll_duration_quantiles: recomputed.poll_duration_quantiles,
span_duration_quantiles: recomputed.span_duration_quantiles,
};
}
const REJECTED_FIXED_TAGS: &[u8] = &[
FieldType::U8 as u8,
FieldType::U16 as u8,
FieldType::U32 as u8,
FieldType::OptionalU8 as u8,
FieldType::OptionalU16 as u8,
FieldType::OptionalU32 as u8,
];
#[allow(clippy::collapsible_if)]
fn validate_shape(shape: &TraceShape) -> anyhow::Result<()> {
ensure!(
shape.version == SHAPE_VERSION,
"unsupported shape version {} (expected {SHAPE_VERSION})",
shape.version
);
ensure!(!shape.schemas.is_empty(), "shape has no schemas");
ensure!(!shape.events.is_empty(), "shape has no events");
ensure!(
shape.schemas.len() <= MAX_WIRE_SCHEMAS,
"shape has {} schemas, exceeds wire type-ID capacity ({MAX_WIRE_SCHEMAS})",
shape.schemas.len()
);
let mut seen_names: HashSet<&str> = HashSet::new();
for (si, schema) in shape.schemas.iter().enumerate() {
ensure!(
!schema.name.is_empty(),
"schema at index {si} has empty name"
);
ensure!(
seen_names.insert(&schema.name),
"duplicate schema name '{}' at index {si}",
schema.name
);
ensure!(
schema.name.len() <= u16::MAX as usize,
"schema '{}' name is {} bytes, exceeds u16 wire limit (65535)",
schema.name,
schema.name.len()
);
ensure!(
schema.fields.len() <= u16::MAX as usize,
"schema '{}' has {} fields, exceeds u16 wire limit (65535)",
schema.name,
schema.fields.len()
);
let mut seen_field_names: HashSet<&str> = HashSet::new();
for (fi, field) in schema.fields.iter().enumerate() {
ensure!(
!field.name.is_empty(),
"schema '{}' field {fi} has empty name",
schema.name
);
ensure!(
field.name.len() <= u16::MAX as usize,
"schema '{}' field '{}' name is {} bytes, exceeds u16 wire limit (65535)",
schema.name,
field.name,
field.name.len()
);
ensure!(
seen_field_names.insert(&field.name),
"schema '{}' has duplicate field name '{}' at index {fi}",
schema.name,
field.name
);
ensure!(
FieldType::from_tag(field.field_type).is_some(),
"schema '{}' field '{}' has invalid type tag {:#x}",
schema.name,
field.name,
field.field_type
);
ensure!(
!REJECTED_FIXED_TAGS.contains(&field.field_type),
"schema '{}' field '{}' uses rejected fixed-width tag {:#x}; \
use Varint/OptionalVarint instead (see fixed-width normalization docs)",
schema.name,
field.name,
field.field_type
);
if let Some(meta) = &field.repeat_meta {
let ft = FieldType::from_tag(field.field_type).unwrap();
let inner = ft.inner();
match meta.semantics {
FieldSemantics::Identity => {
ensure!(
meta.namespace.is_some(),
"schema '{}' field '{}': Identity semantics requires a namespace",
schema.name,
field.name
);
let ok = matches!(
inner,
FieldType::Varint
| FieldType::StackFrames
| FieldType::PooledStackFrames
);
ensure!(
ok,
"schema '{}' field '{}': Identity semantics incompatible with type {:?}",
schema.name,
field.name,
inner
);
if matches!(inner, FieldType::StackFrames | FieldType::PooledStackFrames) {
ensure!(
meta.namespace == Some(NamespaceId::Addr),
"schema '{}' field '{}': stack Identity must use Addr namespace, got {:?}",
schema.name,
field.name,
meta.namespace
);
}
}
FieldSemantics::TimestampRef => {
ensure!(
inner == FieldType::Varint,
"schema '{}' field '{}': TimestampRef requires Varint type, got {:?}",
schema.name,
field.name,
inner
);
ensure!(
meta.namespace.is_none(),
"schema '{}' field '{}': TimestampRef forbids namespace",
schema.name,
field.name
);
}
FieldSemantics::RealtimeOffset => {
ensure!(
inner == FieldType::Varint,
"schema '{}' field '{}': RealtimeOffset requires Varint type, got {:?}",
schema.name,
field.name,
inner
);
ensure!(
meta.namespace.is_none(),
"schema '{}' field '{}': RealtimeOffset forbids namespace",
schema.name,
field.name
);
}
FieldSemantics::Structural => {
ensure!(
meta.namespace.is_none(),
"schema '{}' field '{}': Structural semantics forbids namespace",
schema.name,
field.name
);
}
}
}
}
ensure!(
schema.annotations.len() <= MAX_SCHEMA_ANNOTATIONS,
"schema '{}' has {} annotations, exceeds u16 wire limit ({MAX_SCHEMA_ANNOTATIONS})",
schema.name,
schema.annotations.len()
);
let mut seen_annotations: HashSet<(u16, &str)> = HashSet::new();
for ann in &schema.annotations {
ensure!(
(ann.field_index as usize) < schema.fields.len(),
"schema '{}' annotation field_index {} >= field count {}",
schema.name,
ann.field_index,
schema.fields.len()
);
ensure!(
SAFE_ANNOTATION_KEYS.contains(&ann.key.as_str()),
"schema '{}' annotation has unsafe key '{}'",
schema.name,
ann.key
);
ensure!(
is_safe_annotation_value(&ann.key, &ann.value),
"schema '{}' annotation '{}' has unsafe value '{}'",
schema.name,
ann.key,
ann.value
);
ensure!(
seen_annotations.insert((ann.field_index, &ann.key)),
"schema '{}' has duplicate annotation (field_index={}, key='{}')",
schema.name,
ann.field_index,
ann.key
);
}
if !schema.has_timestamp {
let has_events = shape.events.iter().any(|e| e.schema_index == si as u32);
ensure!(
!has_events,
"schema '{}' has has_timestamp=false but has events; Encoder cannot emit these",
schema.name
);
}
}
let actual_max_offset = shape
.events
.iter()
.filter_map(|e| e.timestamp_offset_ns)
.max()
.unwrap_or(0);
let mut actual_type_counts: BTreeMap<String, u64> = BTreeMap::new();
for (i, event) in shape.events.iter().enumerate() {
let si = event.schema_index as usize;
ensure!(
si < shape.schemas.len(),
"event {i} references schema index {si} but only {} schemas exist",
shape.schemas.len()
);
let schema = &shape.schemas[si];
ensure!(
event.values.len() == schema.fields.len(),
"event {i} has {} values but schema '{}' has {} fields",
event.values.len(),
schema.name,
schema.fields.len()
);
if schema.has_timestamp {
ensure!(
event.timestamp_offset_ns.is_some(),
"event {i} for timestamped schema '{}' missing timestamp_offset_ns",
schema.name
);
let offset = event.timestamp_offset_ns.unwrap();
ensure!(
offset <= actual_max_offset,
"event {i} offset {offset} > derived max offset {actual_max_offset}"
);
}
for (fi, (sv, field)) in event.values.iter().zip(schema.fields.iter()).enumerate() {
validate_value_compat(sv, field.field_type, &schema.name, &field.name, i, fi)?;
}
*actual_type_counts.entry(schema.name.clone()).or_default() += 1;
}
ensure!(
shape.summary.event_count == shape.events.len() as u64,
"summary event_count {} != actual {}",
shape.summary.event_count,
shape.events.len()
);
for (name, expected) in &shape.summary.event_type_counts {
ensure!(
*expected > 0,
"summary event_type_counts['{name}'] = 0; zero-count keys must not be present"
);
let actual = actual_type_counts.get(name).copied().unwrap_or(0);
ensure!(
actual == *expected,
"summary event_type_counts['{name}'] = {expected} but actual = {actual}"
);
}
for (name, actual) in &actual_type_counts {
let expected = shape
.summary
.event_type_counts
.get(name)
.copied()
.unwrap_or(0);
ensure!(
*actual == expected,
"actual events have {actual} '{name}' events but summary does not list it (expected {expected})"
);
}
ensure!(
shape.summary.duration_ns == actual_max_offset,
"summary duration_ns {} != actual max offset {}",
shape.summary.duration_ns,
actual_max_offset
);
let recomputed = recompute_summary(shape);
ensure!(
shape.summary.worker_cardinality == recomputed.worker_cardinality,
"summary worker_cardinality {} != recomputed {}",
shape.summary.worker_cardinality,
recomputed.worker_cardinality
);
ensure!(
shape.summary.task_cardinality == recomputed.task_cardinality,
"summary task_cardinality {} != recomputed {}",
shape.summary.task_cardinality,
recomputed.task_cardinality
);
validate_quantiles_match(
"poll_duration",
&shape.summary.poll_duration_quantiles,
&recomputed.poll_duration_quantiles,
)?;
validate_quantiles_match(
"span_duration",
&shape.summary.span_duration_quantiles,
&recomputed.span_duration_quantiles,
)?;
Ok(())
}
fn validate_quantiles_match(
label: &str,
declared: &Option<Quantiles>,
recomputed: &Option<Quantiles>,
) -> anyhow::Result<()> {
match (declared, recomputed) {
(None, None) => Ok(()),
(Some(_), None) => {
bail!("summary {label}_quantiles is present but recomputation found no data")
}
(None, Some(_)) => {
bail!("summary {label}_quantiles is absent but recomputation found data")
}
(Some(d), Some(r)) => {
ensure!(
d.count == r.count,
"summary {label}_quantiles.count {} != recomputed {}",
d.count,
r.count
);
ensure!(
d.min_ns == r.min_ns,
"summary {label}_quantiles.min_ns {} != recomputed {}",
d.min_ns,
r.min_ns
);
ensure!(
d.p50_ns == r.p50_ns,
"summary {label}_quantiles.p50_ns {} != recomputed {}",
d.p50_ns,
r.p50_ns
);
ensure!(
d.p90_ns == r.p90_ns,
"summary {label}_quantiles.p90_ns {} != recomputed {}",
d.p90_ns,
r.p90_ns
);
ensure!(
d.p99_ns == r.p99_ns,
"summary {label}_quantiles.p99_ns {} != recomputed {}",
d.p99_ns,
r.p99_ns
);
ensure!(
d.max_ns == r.max_ns,
"summary {label}_quantiles.max_ns {} != recomputed {}",
d.max_ns,
r.max_ns
);
Ok(())
}
}
}
#[allow(clippy::collapsible_if)] fn recompute_summary(shape: &TraceShape) -> ShapeSummary {
let mut worker_ids: HashSet<u64> = HashSet::new();
let mut task_ids: HashSet<u64> = HashSet::new();
let mut poll_starts: HashMap<u64, Vec<u64>> = HashMap::new(); let mut poll_durations: Vec<u64> = Vec::new();
let mut span_enters: HashMap<u64, Vec<u64>> = HashMap::new(); let mut span_durations: Vec<u64> = Vec::new();
for event in &shape.events {
let schema = &shape.schemas[event.schema_index as usize];
let ts = event.timestamp_offset_ns.unwrap_or(0);
for (fi, field) in schema.fields.iter().enumerate() {
if field.name == "worker_id" {
if let Some(ShapeValue::U(wid)) = event.values.get(fi) {
if *wid != 254 && *wid != 255 {
worker_ids.insert(*wid);
}
}
}
if let Some(meta) = &field.repeat_meta {
if meta.semantics == FieldSemantics::Identity
&& meta.namespace == Some(NamespaceId::Task)
{
if let Some(ShapeValue::U(tid)) = event.values.get(fi) {
if *tid != 0 {
task_ids.insert(*tid);
}
}
}
}
}
if schema.name == "PollStartEvent" {
for (fi, field) in schema.fields.iter().enumerate() {
if field.name == "worker_id" {
if let Some(ShapeValue::U(wid)) = event.values.get(fi) {
poll_starts.entry(*wid).or_default().push(ts);
}
break;
}
}
}
if schema.name == "PollEndEvent" {
for (fi, field) in schema.fields.iter().enumerate() {
if field.name == "worker_id" {
if let Some(ShapeValue::U(wid)) = event.values.get(fi) {
if let Some(starts) = poll_starts.get_mut(wid) {
if let Some(start_ts) = starts.pop() {
if ts >= start_ts {
poll_durations.push(ts - start_ts);
}
}
}
}
break;
}
}
}
if schema.name.starts_with("SpanEnter:") {
for (fi, field) in schema.fields.iter().enumerate() {
if field.name == "span_id" {
if let Some(ShapeValue::U(sid)) = event.values.get(fi) {
span_enters.entry(*sid).or_default().push(ts);
}
break;
}
}
}
if schema.name.starts_with("SpanExit:") {
for (fi, field) in schema.fields.iter().enumerate() {
if field.name == "span_id" {
if let Some(ShapeValue::U(sid)) = event.values.get(fi) {
if let Some(enters) = span_enters.get_mut(sid) {
if let Some(enter_ts) = enters.pop() {
if ts >= enter_ts {
span_durations.push(ts - enter_ts);
}
}
}
}
break;
}
}
}
}
let poll_quantiles = compute_quantiles(&mut poll_durations);
let span_quantiles = compute_quantiles(&mut span_durations);
ShapeSummary {
event_count: shape.events.len() as u64,
duration_ns: shape
.events
.iter()
.filter_map(|e| e.timestamp_offset_ns)
.max()
.unwrap_or(0),
event_type_counts: BTreeMap::new(), worker_cardinality: worker_ids.len() as u32,
task_cardinality: task_ids.len() as u32,
poll_duration_quantiles: poll_quantiles,
span_duration_quantiles: span_quantiles,
}
}
fn validate_value_compat(
sv: &ShapeValue,
field_type_tag: u8,
schema_name: &str,
field_name: &str,
event_idx: usize,
field_idx: usize,
) -> anyhow::Result<()> {
let ft = FieldType::from_tag(field_type_tag).unwrap(); let is_optional = ft.is_optional();
let inner = ft.inner();
if matches!(sv, ShapeValue::None) {
ensure!(
is_optional,
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
None value but field is not optional"
);
return Ok(());
}
match (sv, inner) {
(ShapeValue::U(_), FieldType::Varint) => Ok(()),
(ShapeValue::I(_), FieldType::I64) => Ok(()),
(ShapeValue::F(_), FieldType::F64) => Ok(()),
(ShapeValue::B(_), FieldType::Bool) => Ok(()),
(ShapeValue::S(_), FieldType::String) => Ok(()),
(ShapeValue::PS(_), FieldType::PooledString) => Ok(()),
(ShapeValue::Bytes(len), FieldType::Bytes) => {
ensure!(
*len <= MAX_BYTES_LENGTH,
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
Bytes length {len} exceeds limit {MAX_BYTES_LENGTH}"
);
Ok(())
}
(ShapeValue::Stack(addrs), FieldType::StackFrames) => {
ensure!(
addrs.len() <= MAX_CONTAINER_ELEMENTS,
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
Stack has {} elements, exceeds limit {MAX_CONTAINER_ELEMENTS}",
addrs.len()
);
Ok(())
}
(ShapeValue::PStack(addrs), FieldType::PooledStackFrames) => {
ensure!(
addrs.len() <= MAX_CONTAINER_ELEMENTS,
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
PStack has {} elements, exceeds limit {MAX_CONTAINER_ELEMENTS}",
addrs.len()
);
Ok(())
}
(ShapeValue::StringMap(pairs), FieldType::StringMap) => {
ensure!(
pairs.len() <= MAX_CONTAINER_ELEMENTS,
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
StringMap has {} entries, exceeds limit {MAX_CONTAINER_ELEMENTS}",
pairs.len()
);
Ok(())
}
(ShapeValue::List(items), FieldType::DynamicList) => {
ensure!(
items.len() <= MAX_CONTAINER_ELEMENTS,
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
List has {} elements, exceeds limit {MAX_CONTAINER_ELEMENTS}",
items.len()
);
validate_nested_values(items, schema_name, field_name, event_idx, field_idx, 1)
}
(ShapeValue::Map(pairs), FieldType::DynamicMap) => {
ensure!(
pairs.len() <= MAX_CONTAINER_ELEMENTS,
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
Map has {} entries, exceeds limit {MAX_CONTAINER_ELEMENTS}",
pairs.len()
);
for (k, v) in pairs {
ensure!(
!matches!(k, ShapeValue::None),
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
nested None in dynamic map key"
);
ensure!(
!matches!(v, ShapeValue::None),
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
nested None in dynamic map value"
);
validate_nested_value(k, schema_name, field_name, event_idx, field_idx, 1)?;
validate_nested_value(v, schema_name, field_name, event_idx, field_idx, 1)?;
}
Ok(())
}
_ => bail!(
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
value type {:?} incompatible with field type {:?}",
std::mem::discriminant(sv),
inner
),
}
}
fn validate_nested_values(
items: &[ShapeValue],
schema_name: &str,
field_name: &str,
event_idx: usize,
field_idx: usize,
depth: usize,
) -> anyhow::Result<()> {
ensure!(
depth <= MAX_DYNAMIC_DEPTH,
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
excessive nesting depth ({depth})"
);
for item in items {
ensure!(
!matches!(item, ShapeValue::None),
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
nested None in dynamic container is not allowed"
);
validate_nested_value(item, schema_name, field_name, event_idx, field_idx, depth)?;
}
Ok(())
}
fn validate_nested_value(
item: &ShapeValue,
schema_name: &str,
field_name: &str,
event_idx: usize,
field_idx: usize,
depth: usize,
) -> anyhow::Result<()> {
match item {
ShapeValue::Bytes(len) => {
ensure!(
*len <= MAX_BYTES_LENGTH,
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
nested Bytes length {len} exceeds limit {MAX_BYTES_LENGTH}"
);
}
ShapeValue::Stack(addrs) | ShapeValue::PStack(addrs) => {
ensure!(
addrs.len() <= MAX_CONTAINER_ELEMENTS,
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
nested stack has {} elements, exceeds limit {MAX_CONTAINER_ELEMENTS}",
addrs.len()
);
}
ShapeValue::List(inner_items) => {
ensure!(
inner_items.len() <= MAX_CONTAINER_ELEMENTS,
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
nested List has {} elements at depth {depth}",
inner_items.len()
);
validate_nested_values(
inner_items,
schema_name,
field_name,
event_idx,
field_idx,
depth + 1,
)?;
}
ShapeValue::Map(inner_pairs) => {
ensure!(
inner_pairs.len() <= MAX_CONTAINER_ELEMENTS,
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
nested Map has {} entries at depth {depth}",
inner_pairs.len()
);
for (k, v) in inner_pairs {
ensure!(
!matches!(k, ShapeValue::None),
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
nested None in dynamic map key"
);
ensure!(
!matches!(v, ShapeValue::None),
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
nested None in dynamic map value"
);
validate_nested_value(k, schema_name, field_name, event_idx, field_idx, depth + 1)?;
validate_nested_value(v, schema_name, field_name, event_idx, field_idx, depth + 1)?;
}
}
ShapeValue::None => {
bail!(
"event {event_idx} field {field_idx} ('{field_name}' in '{schema_name}'): \
nested None in dynamic container"
);
}
_ => {} }
Ok(())
}
fn validate_repeat_preflight(shape: &TraceShape, repeat: u32) -> anyhow::Result<()> {
let total_events = (shape.events.len() as u64)
.checked_mul(repeat as u64)
.ok_or_else(|| anyhow::anyhow!("events.len() * repeat overflows u64"))?;
ensure!(
total_events <= MAX_GENERATED_EVENTS,
"total event count {total_events} exceeds limit {MAX_GENERATED_EVENTS}"
);
let estimated_size = estimate_encoded_size(shape, repeat)?;
ensure!(
estimated_size <= MAX_GENERATED_OUTPUT_BYTES,
"estimated output size {estimated_size} bytes exceeds limit {MAX_GENERATED_OUTPUT_BYTES}"
);
let actual_duration = shape
.events
.iter()
.filter_map(|e| e.timestamp_offset_ns)
.max()
.unwrap_or(0);
let template_duration = actual_duration.max(QUANTUM_NS);
let gap = QUANTUM_NS;
let duration_plus_gap = template_duration
.checked_add(gap)
.ok_or_else(|| anyhow::anyhow!("preflight: template_duration + gap overflow"))?;
let last_rep = (repeat as u64).saturating_sub(1);
let max_time_shift = last_rep
.checked_mul(duration_plus_gap)
.ok_or_else(|| anyhow::anyhow!("preflight: rep * (duration+gap) overflow at last rep"))?;
for (i, event) in shape.events.iter().enumerate() {
let ts_offset = event.timestamp_offset_ns.unwrap_or(0);
ts_offset.checked_add(max_time_shift).ok_or_else(|| {
anyhow::anyhow!(
"preflight: event {i} timestamp offset {ts_offset} + time_shift {max_time_shift} overflow"
)
})?;
}
let strides = compute_namespace_strides(shape)?;
for (ns, stride) in &strides {
let max_id_shift = last_rep.checked_mul(*stride).ok_or_else(|| {
anyhow::anyhow!("preflight: namespace {ns:?} stride {stride} * last_rep overflow")
})?;
let max_val = find_namespace_max_value(shape, *ns);
if max_val > 0 {
max_val.checked_add(max_id_shift).ok_or_else(|| {
anyhow::anyhow!(
"preflight: namespace {ns:?} max_val {max_val} + shift {max_id_shift} overflow"
)
})?;
}
}
for (i, event) in shape.events.iter().enumerate() {
let schema = &shape.schemas[event.schema_index as usize];
for (fi, field) in schema.fields.iter().enumerate() {
if let Some(meta) = &field.repeat_meta {
match meta.semantics {
FieldSemantics::TimestampRef => {
if let Some(ShapeValue::U(v)) = event.values.get(fi) {
v.checked_add(max_time_shift).ok_or_else(|| {
anyhow::anyhow!(
"preflight: event {i} field {fi} TimestampRef {v} + time_shift overflow"
)
})?;
}
}
FieldSemantics::RealtimeOffset => {
if let Some(ShapeValue::U(v)) = event.values.get(fi) {
SYNTHETIC_EPOCH_NS
.checked_add(*v)
.and_then(|x| x.checked_add(max_time_shift))
.ok_or_else(|| {
anyhow::anyhow!(
"preflight: event {i} field {fi} RealtimeOffset overflow"
)
})?;
}
}
_ => {}
}
}
}
}
const MAX_EVENT_WORKING_SET: u64 = MAX_BYTES_LENGTH as u64; for (i, event) in shape.events.iter().enumerate() {
let working_set = event.values.iter().try_fold(0u64, |total, value| {
total
.checked_add(estimate_value_working_set(value))
.ok_or_else(|| anyhow::anyhow!("preflight: event working-set overflow"))
})?;
ensure!(
working_set <= MAX_EVENT_WORKING_SET,
"preflight: event {i} working set {working_set} exceeds \
per-event limit {MAX_EVENT_WORKING_SET}"
);
}
let schema_bytes = estimate_registered_schema_memory(shape);
let pooled_string_bytes = shape.events.iter().fold(0u64, |total, event| {
event.values.iter().fold(total, |event_total, value| {
event_total.saturating_add(sum_pooled_string_memory(value))
})
});
let pooled_stack_bytes_per_template = shape.events.iter().fold(0u64, |total, event| {
event.values.iter().fold(total, |event_total, value| {
event_total.saturating_add(sum_pooled_stack_memory(value))
})
});
let retained_pooled_stack_bytes = pooled_stack_bytes_per_template
.checked_mul(repeat as u64)
.ok_or_else(|| anyhow::anyhow!("preflight: retained pooled stack bytes overflow"))?;
let retained_generation_bytes = schema_bytes
.checked_add(pooled_string_bytes)
.and_then(|bytes| bytes.checked_add(retained_pooled_stack_bytes))
.ok_or_else(|| anyhow::anyhow!("preflight: retained generation memory overflow"))?;
ensure!(
retained_generation_bytes <= MAX_RETAINED_GENERATION_BYTES,
"preflight: retained generation memory {retained_generation_bytes} exceeds \
limit {MAX_RETAINED_GENERATION_BYTES} (schemas={schema_bytes}, \
pooled_strings={pooled_string_bytes}, pooled_stacks={retained_pooled_stack_bytes})"
);
Ok(())
}
fn estimate_value_working_set(sv: &ShapeValue) -> u64 {
const VALUE_OVERHEAD: u64 = 64;
match sv {
ShapeValue::U(_)
| ShapeValue::I(_)
| ShapeValue::F(_)
| ShapeValue::B(_)
| ShapeValue::None => VALUE_OVERHEAD,
ShapeValue::S(value) => VALUE_OVERHEAD.saturating_add(value.len() as u64),
ShapeValue::PS(value) => {
VALUE_OVERHEAD.saturating_add((value.len() as u64).saturating_mul(2))
}
ShapeValue::Bytes(len) => VALUE_OVERHEAD.saturating_add(*len as u64),
ShapeValue::Stack(frames) => {
VALUE_OVERHEAD.saturating_add((frames.len() as u64).saturating_mul(8))
}
ShapeValue::PStack(frames) => {
VALUE_OVERHEAD.saturating_add((frames.len() as u64).saturating_mul(24))
}
ShapeValue::StringMap(pairs) => pairs.iter().fold(VALUE_OVERHEAD, |total, (key, value)| {
total
.saturating_add(64)
.saturating_add(key.len() as u64)
.saturating_add(value.len() as u64)
}),
ShapeValue::List(items) => items.iter().fold(VALUE_OVERHEAD, |total, item| {
total
.saturating_add(VALUE_OVERHEAD)
.saturating_add(estimate_value_working_set(item))
}),
ShapeValue::Map(pairs) => pairs.iter().fold(VALUE_OVERHEAD, |total, (key, value)| {
total
.saturating_add(VALUE_OVERHEAD * 2)
.saturating_add(estimate_value_working_set(key))
.saturating_add(estimate_value_working_set(value))
}),
}
}
fn sum_pooled_string_memory(sv: &ShapeValue) -> u64 {
const POOL_ENTRY_OVERHEAD: u64 = 128;
match sv {
ShapeValue::PS(value) => POOL_ENTRY_OVERHEAD.saturating_add(value.len() as u64),
ShapeValue::List(items) => items
.iter()
.map(sum_pooled_string_memory)
.fold(0u64, u64::saturating_add),
ShapeValue::Map(pairs) => pairs
.iter()
.flat_map(|(key, value)| [key, value])
.map(sum_pooled_string_memory)
.fold(0u64, u64::saturating_add),
_ => 0,
}
}
fn sum_pooled_stack_memory(sv: &ShapeValue) -> u64 {
const POOL_ENTRY_OVERHEAD: u64 = 128;
match sv {
ShapeValue::PStack(frames) => {
POOL_ENTRY_OVERHEAD.saturating_add((frames.len() as u64).saturating_mul(8))
}
ShapeValue::List(items) => items
.iter()
.map(sum_pooled_stack_memory)
.fold(0u64, u64::saturating_add),
ShapeValue::Map(pairs) => pairs
.iter()
.flat_map(|(key, value)| [key, value])
.map(sum_pooled_stack_memory)
.fold(0u64, u64::saturating_add),
_ => 0,
}
}
fn estimate_registered_schema_memory(shape: &TraceShape) -> u64 {
shape.schemas.iter().fold(0u64, |total, schema| {
let fields = schema.fields.iter().fold(0u64, |field_total, field| {
field_total
.saturating_add(192)
.saturating_add((field.name.len() as u64).saturating_mul(3))
});
let annotations = schema.annotations.iter().fold(0u64, |ann_total, ann| {
ann_total
.saturating_add(192)
.saturating_add(((ann.key.len() + ann.value.len()) as u64).saturating_mul(3))
});
total
.saturating_add(256)
.saturating_add((schema.name.len() as u64).saturating_mul(3))
.saturating_add(fields)
.saturating_add(annotations)
})
}
fn find_namespace_max_value(shape: &TraceShape, target_ns: NamespaceId) -> u64 {
let mut max_val: u64 = 0;
for event in &shape.events {
let schema = &shape.schemas[event.schema_index as usize];
for (fi, field) in schema.fields.iter().enumerate() {
let Some(sv) = event.values.get(fi) else {
continue;
};
if target_ns == NamespaceId::Addr {
max_val = max_val.max(max_nested_stack_address(sv));
}
if let Some(meta) = &field.repeat_meta
&& meta.semantics == FieldSemantics::Identity
&& meta.namespace == Some(target_ns)
{
max_val = max_val.max(max_identity_value(sv));
}
}
}
max_val
}
fn max_identity_value(sv: &ShapeValue) -> u64 {
match sv {
ShapeValue::U(v) => *v,
ShapeValue::Stack(addrs) | ShapeValue::PStack(addrs) => {
addrs.iter().copied().max().unwrap_or(0)
}
_ => 0,
}
}
fn max_nested_stack_address(sv: &ShapeValue) -> u64 {
match sv {
ShapeValue::Stack(addrs) | ShapeValue::PStack(addrs) => {
addrs.iter().copied().max().unwrap_or(0)
}
ShapeValue::List(items) => items
.iter()
.map(max_nested_stack_address)
.max()
.unwrap_or(0),
ShapeValue::Map(pairs) => pairs
.iter()
.flat_map(|(k, v)| [max_nested_stack_address(k), max_nested_stack_address(v)])
.max()
.unwrap_or(0),
_ => 0,
}
}
fn estimate_encoded_size(shape: &TraceShape, repeat: u32) -> anyhow::Result<u64> {
let header_overhead: u64 = 64; let schema_overhead: u64 = shape
.schemas
.iter()
.map(|s| {
let base = 1u64 + 2 + 2 + s.name.len() as u64 + 1 + 2;
let fields_cost: u64 = s
.fields
.iter()
.map(|f| 2u64 + f.name.len() as u64 + 1)
.sum();
let ann_cost: u64 = if s.annotations.is_empty() {
0
} else {
6 + s
.annotations
.iter()
.map(|a| 8u64 + a.key.len() as u64 + a.value.len() as u64)
.sum::<u64>()
};
base.saturating_add(fields_cost).saturating_add(ann_cost)
})
.fold(0u64, u64::saturating_add);
let mut per_template_bytes: u64 = 0;
for event in &shape.events {
let mut event_size: u64 = 20;
let schema = &shape.schemas[event.schema_index as usize];
for (fi, sv) in event.values.iter().enumerate() {
let field = &schema.fields[fi];
event_size = event_size.saturating_add(estimate_value_size(sv, field.field_type));
}
per_template_bytes = per_template_bytes.saturating_add(event_size);
}
let repeated = per_template_bytes
.checked_mul(repeat as u64)
.ok_or_else(|| anyhow::anyhow!("size estimate: per_template * repeat overflows u64"))?;
let total = header_overhead
.saturating_add(schema_overhead)
.saturating_add(repeated);
Ok(total)
}
fn estimate_value_size(sv: &ShapeValue, field_type_tag: u8) -> u64 {
let optional_prefix =
FieldType::from_tag(field_type_tag).is_some_and(FieldType::is_optional) as u64;
let value_size = match sv {
ShapeValue::U(_) | ShapeValue::I(_) => 10, ShapeValue::F(_) => 8,
ShapeValue::B(_) => 1,
ShapeValue::S(s) => 5 + s.len() as u64, ShapeValue::PS(s) => {
let pool_frame = 13u64 + s.len() as u64;
pool_frame + 4
}
ShapeValue::Bytes(len) => 5 + *len as u64,
ShapeValue::Stack(addrs) => 5 + addrs.len() as u64 * 10, ShapeValue::PStack(addrs) => {
let pool_frame = 13u64 + addrs.len() as u64 * 8;
pool_frame + 4
}
ShapeValue::StringMap(pairs) => {
5 + pairs
.iter()
.map(|(k, v)| 10 + k.len() as u64 + v.len() as u64)
.sum::<u64>()
}
ShapeValue::List(items) => {
4 + items
.iter()
.map(|i| 1u64.saturating_add(estimate_value_size(i, 0)))
.sum::<u64>()
}
ShapeValue::Map(pairs) => {
4 + pairs
.iter()
.map(|(k, v)| {
2u64.saturating_add(estimate_value_size(k, 0))
.saturating_add(estimate_value_size(v, 0))
})
.sum::<u64>()
}
ShapeValue::None => 0,
};
optional_prefix.saturating_add(value_size)
}
#[allow(clippy::collapsible_if)] fn compute_namespace_strides(shape: &TraceShape) -> anyhow::Result<HashMap<NamespaceId, u64>> {
let mut max_values: HashMap<NamespaceId, u64> = HashMap::new();
for event in &shape.events {
let schema = &shape.schemas[event.schema_index as usize];
for (fi, sv) in event.values.iter().enumerate() {
let field = &schema.fields[fi];
let nested_addr = max_nested_stack_address(sv);
if nested_addr != 0 {
let entry = max_values.entry(NamespaceId::Addr).or_insert(0);
*entry = (*entry).max(nested_addr);
}
if let Some(m) = field.repeat_meta.as_ref()
&& m.semantics == FieldSemantics::Identity
&& let Some(ns) = m.namespace
{
let entry = max_values.entry(ns).or_insert(0);
*entry = (*entry).max(max_identity_value(sv));
}
}
}
let mut strides: HashMap<NamespaceId, u64> = HashMap::new();
for (ns, max_val) in max_values {
let stride = max_val
.checked_add(1)
.ok_or_else(|| anyhow::anyhow!("namespace '{ns:?}' max value overflow for stride"))?;
strides.insert(ns, stride);
}
Ok(strides)
}
#[cfg(test)]
fn generate_trace(shape: &TraceShape, repeat: u32) -> anyhow::Result<Vec<u8>> {
generate_trace_with_bases(shape, repeat, 0, SYNTHETIC_EPOCH_NS)
}
pub(crate) fn generate_trace_with_bases(
shape: &TraceShape,
repeat: u32,
timestamp_base_ns: u64,
realtime_base_ns: u64,
) -> anyhow::Result<Vec<u8>> {
let mut buf = Vec::new();
validate_repeat_preflight(shape, repeat)?;
generate_to_writer_with_bases(shape, repeat, timestamp_base_ns, realtime_base_ns, &mut buf)?;
Ok(buf)
}
fn generate_to_writer<W: Write>(shape: &TraceShape, repeat: u32, writer: W) -> anyhow::Result<()> {
generate_to_writer_with_bases(shape, repeat, 0, SYNTHETIC_EPOCH_NS, writer)
}
fn generate_to_writer_with_bases<W: Write>(
shape: &TraceShape,
repeat: u32,
timestamp_base_ns: u64,
realtime_base_ns: u64,
writer: W,
) -> anyhow::Result<()> {
let actual_duration = shape
.events
.iter()
.filter_map(|e| e.timestamp_offset_ns)
.max()
.unwrap_or(0);
let template_duration = actual_duration.max(QUANTUM_NS);
let gap = QUANTUM_NS;
let strides = compute_namespace_strides(shape)?;
let mut encoder = Encoder::new_to(writer).context("write trace header")?;
let schemas: Vec<Schema> = shape
.schemas
.iter()
.map(|ss| {
let fields: Vec<FieldDef> = ss
.fields
.iter()
.map(|f| {
let ft = FieldType::from_tag(f.field_type).expect("validated");
FieldDef::new(&f.name, ft)
})
.collect();
let annotations: Vec<FieldAnnotation> = ss
.annotations
.iter()
.map(|ann| {
FieldAnnotation::new(ann.field_index, ann.key.clone(), ann.value.clone())
})
.collect();
let entry =
SchemaEntry::with_annotations(&ss.name, ss.has_timestamp, fields, annotations);
Schema::from_entry(entry)
})
.collect();
for schema in &schemas {
encoder
.register_existing(schema)
.with_context(|| format!("register schema '{}'", schema.name()))?;
}
for rep in 0..repeat {
let time_shift = (rep as u64)
.checked_mul(
template_duration
.checked_add(gap)
.ok_or_else(|| anyhow::anyhow!("template_duration + gap overflow"))?,
)
.ok_or_else(|| anyhow::anyhow!("rep * (duration+gap) overflow at rep={rep}"))?;
let timeline = GenerationTimeline {
time_shift_ns: time_shift,
realtime_base_ns,
};
for event in &shape.events {
let schema = &schemas[event.schema_index as usize];
let shape_schema = &shape.schemas[event.schema_index as usize];
let ts_offset = event.timestamp_offset_ns.unwrap_or(0);
let ts_ns = timestamp_base_ns
.checked_add(ts_offset)
.and_then(|timestamp| timestamp.checked_add(time_shift))
.ok_or_else(|| anyhow::anyhow!("timestamp overflow at rep={rep}"))?;
let mut values = Vec::with_capacity(shape_schema.fields.len() + 1);
values.push(FieldValue::Varint(ts_ns));
for (fi, sv) in event.values.iter().enumerate() {
let field = &shape_schema.fields[fi];
let ft = FieldType::from_tag(field.field_type).expect("validated");
let meta = field.repeat_meta.as_ref();
let fv = shape_value_to_field_value(
sv,
ft,
meta,
rep,
timeline,
&strides,
&mut encoder,
)?;
values.push(fv);
}
encoder
.write_event(schema, &values)
.with_context(|| format!("write event for schema '{}'", schema.name()))?;
}
}
encoder.flush().context("flush encoder")?;
let mut inner = encoder.into_inner();
inner.flush().context("flush writer")?;
Ok(())
}
#[derive(Clone, Copy)]
struct GenerationTimeline {
time_shift_ns: u64,
realtime_base_ns: u64,
}
fn shape_value_to_field_value<W: Write>(
sv: &ShapeValue,
ft: FieldType,
meta: Option<&FieldRepeatMeta>,
rep: u32,
timeline: GenerationTimeline,
strides: &HashMap<NamespaceId, u64>,
encoder: &mut Encoder<W>,
) -> anyhow::Result<FieldValue> {
let is_optional = ft.is_optional();
if matches!(sv, ShapeValue::None) {
ensure!(is_optional, "None value for non-optional field");
return Ok(FieldValue::None);
}
let semantics = meta
.map(|m| m.semantics)
.unwrap_or(FieldSemantics::Structural);
let namespace = meta.and_then(|m| m.namespace);
match sv {
ShapeValue::U(v) => {
let shifted = match semantics {
FieldSemantics::Identity => {
if *v == 0 {
0
} else {
let ns = namespace.unwrap_or(NamespaceId::Addr);
let stride = strides.get(&ns).copied().unwrap_or(1);
let id_shift = (rep as u64).checked_mul(stride).ok_or_else(|| {
anyhow::anyhow!("identity shift overflow ns={ns:?} rep={rep}")
})?;
v.checked_add(id_shift)
.ok_or_else(|| anyhow::anyhow!("identity value overflow ns={ns:?}"))?
}
}
FieldSemantics::TimestampRef => {
v.checked_add(timeline.time_shift_ns)
.ok_or_else(|| anyhow::anyhow!("alloc_timestamp_ns overflow"))?
}
FieldSemantics::RealtimeOffset => {
timeline
.realtime_base_ns
.checked_add(*v)
.and_then(|x| x.checked_add(timeline.time_shift_ns))
.ok_or_else(|| anyhow::anyhow!("realtime_ns overflow"))?
}
FieldSemantics::Structural => *v,
};
Ok(FieldValue::Varint(shifted))
}
ShapeValue::I(v) => Ok(FieldValue::I64(*v)),
ShapeValue::F(v) => Ok(FieldValue::F64(*v)),
ShapeValue::B(v) => Ok(FieldValue::Bool(*v)),
ShapeValue::S(s) => Ok(FieldValue::String(s.clone())),
ShapeValue::PS(s) => {
let id = encoder.intern_string(s)?;
Ok(FieldValue::PooledString(id))
}
ShapeValue::Bytes(len) => Ok(FieldValue::Bytes(vec![0u8; *len as usize])),
ShapeValue::Stack(addrs) => {
let shifted: Vec<u64> = if semantics == FieldSemantics::Identity && rep > 0 {
let ns = namespace.unwrap_or(NamespaceId::Addr);
let stride = strides.get(&ns).copied().unwrap_or(1);
let id_shift = (rep as u64)
.checked_mul(stride)
.ok_or_else(|| anyhow::anyhow!("stack addr shift overflow"))?;
addrs
.iter()
.map(|a| {
if *a == 0 {
Ok(0u64)
} else {
a.checked_add(id_shift)
.ok_or_else(|| anyhow::anyhow!("stack addr overflow"))
}
})
.collect::<anyhow::Result<_>>()?
} else {
addrs.clone()
};
Ok(FieldValue::StackFrames(shifted.into()))
}
ShapeValue::PStack(addrs) => {
let shifted: Vec<u64> = if semantics == FieldSemantics::Identity && rep > 0 {
let ns = namespace.unwrap_or(NamespaceId::Addr);
let stride = strides.get(&ns).copied().unwrap_or(1);
let id_shift = (rep as u64)
.checked_mul(stride)
.ok_or_else(|| anyhow::anyhow!("pooled stack addr shift overflow"))?;
addrs
.iter()
.map(|a| {
if *a == 0 {
Ok(0u64)
} else {
a.checked_add(id_shift)
.ok_or_else(|| anyhow::anyhow!("pooled stack addr overflow"))
}
})
.collect::<anyhow::Result<_>>()?
} else {
addrs.clone()
};
let id = encoder.intern_stack_frames(&shifted)?;
Ok(FieldValue::PooledStackFrames(id))
}
ShapeValue::StringMap(pairs) => {
let mapped: Vec<(Vec<u8>, Vec<u8>)> = pairs
.iter()
.map(|(k, v)| (k.as_bytes().to_vec(), v.as_bytes().to_vec()))
.collect();
Ok(FieldValue::StringMap(mapped))
}
ShapeValue::List(items) => {
let mapped: Vec<FieldValue> = items
.iter()
.map(|item| {
convert_nested_shape_value(item, rep, timeline.time_shift_ns, strides, encoder)
})
.collect::<anyhow::Result<_>>()?;
Ok(FieldValue::List(mapped))
}
ShapeValue::Map(pairs) => {
let mapped: Vec<(FieldValue, FieldValue)> = pairs
.iter()
.map(|(k, v)| {
let kv = convert_nested_shape_value(
k,
rep,
timeline.time_shift_ns,
strides,
encoder,
)?;
let vv = convert_nested_shape_value(
v,
rep,
timeline.time_shift_ns,
strides,
encoder,
)?;
Ok((kv, vv))
})
.collect::<anyhow::Result<_>>()?;
Ok(FieldValue::Map(mapped))
}
ShapeValue::None => Ok(FieldValue::None),
}
}
#[allow(clippy::only_used_in_recursion)]
fn convert_nested_shape_value<W: Write>(
sv: &ShapeValue,
rep: u32,
time_shift: u64,
strides: &HashMap<NamespaceId, u64>,
encoder: &mut Encoder<W>,
) -> anyhow::Result<FieldValue> {
match sv {
ShapeValue::U(v) => Ok(FieldValue::Varint(*v)),
ShapeValue::I(v) => Ok(FieldValue::I64(*v)),
ShapeValue::F(v) => Ok(FieldValue::F64(*v)),
ShapeValue::B(v) => Ok(FieldValue::Bool(*v)),
ShapeValue::S(s) => Ok(FieldValue::String(s.clone())),
ShapeValue::PS(s) => {
let id = encoder.intern_string(s)?;
Ok(FieldValue::PooledString(id))
}
ShapeValue::Bytes(len) => Ok(FieldValue::Bytes(vec![0u8; *len as usize])),
ShapeValue::Stack(addrs) => {
let shifted = shift_nested_addrs(addrs, rep, strides)?;
Ok(FieldValue::StackFrames(shifted.into()))
}
ShapeValue::PStack(addrs) => {
let shifted = shift_nested_addrs(addrs, rep, strides)?;
let id = encoder.intern_stack_frames(&shifted)?;
Ok(FieldValue::PooledStackFrames(id))
}
ShapeValue::StringMap(pairs) => {
let mapped: Vec<(Vec<u8>, Vec<u8>)> = pairs
.iter()
.map(|(k, v)| (k.as_bytes().to_vec(), v.as_bytes().to_vec()))
.collect();
Ok(FieldValue::StringMap(mapped))
}
ShapeValue::List(items) => {
let mapped: Vec<FieldValue> = items
.iter()
.map(|item| convert_nested_shape_value(item, rep, time_shift, strides, encoder))
.collect::<anyhow::Result<_>>()?;
Ok(FieldValue::List(mapped))
}
ShapeValue::Map(pairs) => {
let mapped: Vec<(FieldValue, FieldValue)> = pairs
.iter()
.map(|(k, v)| {
let kv = convert_nested_shape_value(k, rep, time_shift, strides, encoder)?;
let vv = convert_nested_shape_value(v, rep, time_shift, strides, encoder)?;
Ok((kv, vv))
})
.collect::<anyhow::Result<_>>()?;
Ok(FieldValue::Map(mapped))
}
ShapeValue::None => Ok(FieldValue::None),
}
}
fn shift_nested_addrs(
addrs: &[u64],
rep: u32,
strides: &HashMap<NamespaceId, u64>,
) -> anyhow::Result<Vec<u64>> {
if rep == 0 {
return Ok(addrs.to_vec());
}
let stride = strides.get(&NamespaceId::Addr).copied().unwrap_or(1);
let id_shift = (rep as u64)
.checked_mul(stride)
.ok_or_else(|| anyhow::anyhow!("nested stack addr shift overflow"))?;
addrs
.iter()
.map(|a| {
if *a == 0 {
Ok(0u64)
} else {
a.checked_add(id_shift)
.ok_or_else(|| anyhow::anyhow!("nested stack addr value overflow"))
}
})
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
use dial9_trace_format::encoder::{Encoder, Schema};
use dial9_trace_format::schema::FieldDef;
use dial9_trace_format::types::{FieldType, FieldValue};
const BASE_TS: u64 = 1_700_000_000_000_000_000;
fn register_poll_start_schema(enc: &mut Encoder) -> Schema {
enc.register_schema(
"PollStartEvent",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("local_queue", FieldType::U8),
FieldDef::new("task_id", FieldType::Varint),
FieldDef::new("spawn_loc", FieldType::PooledString),
],
)
.unwrap()
}
fn write_poll_start(enc: &mut Encoder, schema: &Schema, ts: u64, worker: u64, task: u64) {
let spawn_loc = enc.intern_string("test/source.rs").unwrap();
enc.write_event(
schema,
&[
FieldValue::Varint(ts),
FieldValue::Varint(worker),
FieldValue::Varint(0),
FieldValue::Varint(task),
FieldValue::PooledString(spawn_loc),
],
)
.unwrap();
}
fn register_task_spawn_schema(enc: &mut Encoder) -> Schema {
enc.register_schema(
"TaskSpawnEvent",
vec![
FieldDef::new("task_id", FieldType::Varint),
FieldDef::new("spawn_loc", FieldType::PooledString),
FieldDef::new("instrumented", FieldType::Bool),
],
)
.unwrap()
}
fn write_task_spawn(enc: &mut Encoder, schema: &Schema, ts: u64, task: u64) {
let spawn_loc = enc.intern_string("test/source.rs").unwrap();
enc.write_event(
schema,
&[
FieldValue::Varint(ts),
FieldValue::Varint(task),
FieldValue::PooledString(spawn_loc),
FieldValue::Bool(true),
],
)
.unwrap();
}
fn make_test_trace() -> Vec<u8> {
let mut enc = Encoder::new();
let poll_start = enc
.register_schema(
"PollStartEvent",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("local_queue", FieldType::U8),
FieldDef::new("task_id", FieldType::Varint),
FieldDef::new("spawn_loc", FieldType::PooledString),
],
)
.unwrap();
let poll_end = enc
.register_schema(
"PollEndEvent",
vec![FieldDef::new("worker_id", FieldType::Varint)],
)
.unwrap();
let custom = enc
.register_schema(
"MySecretService",
vec![
FieldDef::new("secret_name", FieldType::String),
FieldDef::new("task_id", FieldType::Varint),
FieldDef::new("payload", FieldType::Bytes),
],
)
.unwrap();
let spawn_loc = enc.intern_string("private/source.rs").unwrap();
enc.write_event(
&poll_start,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(0),
FieldValue::Varint(1),
FieldValue::Varint(10),
FieldValue::PooledString(spawn_loc),
],
)
.unwrap();
enc.write_event(
&poll_end,
&[FieldValue::Varint(BASE_TS + 50_000), FieldValue::Varint(0)],
)
.unwrap();
enc.write_event(
&custom,
&[
FieldValue::Varint(BASE_TS + 100_000),
FieldValue::String("password-123".into()),
FieldValue::Varint(42),
FieldValue::Bytes(vec![0xDE, 0xAD, 0xBE, 0xEF]),
],
)
.unwrap();
enc.write_event(
&poll_start,
&[
FieldValue::Varint(BASE_TS + 200_000),
FieldValue::Varint(1),
FieldValue::Varint(2),
FieldValue::Varint(11),
FieldValue::PooledString(spawn_loc),
],
)
.unwrap();
enc.finish()
}
fn make_span_trace() -> Vec<u8> {
let mut enc = Encoder::new();
let enter = enc
.register_schema(
"SpanEnter:my_secret_handler",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
],
)
.unwrap();
let exit = enc
.register_schema(
"SpanExit:my_secret_handler",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
],
)
.unwrap();
enc.write_event(
&enter,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(0),
FieldValue::Varint(999),
],
)
.unwrap();
enc.write_event(
&exit,
&[
FieldValue::Varint(BASE_TS + 100_000),
FieldValue::Varint(0),
FieldValue::Varint(999),
],
)
.unwrap();
enc.finish()
}
#[test]
fn extract_removes_secrets_and_absolute_timestamps() {
let trace = make_test_trace();
let shape = extract_shape(&trace).unwrap();
let json = serde_json::to_string_pretty(&shape).unwrap();
assert!(!json.contains("password-123"), "secret string leaked");
assert!(
!json.contains("MySecretService"),
"custom schema name leaked"
);
assert!(!json.contains("secret_name"), "custom field name leaked");
assert!(!json.contains("1700000000"), "absolute timestamp leaked");
assert!(json.contains("PollStartEvent"));
assert!(json.contains("PollEndEvent"));
assert!(json.contains("worker_id"));
}
#[test]
fn extract_summary_counts() {
let trace = make_test_trace();
let shape = extract_shape(&trace).unwrap();
assert_eq!(shape.summary.event_count, 4);
assert_eq!(shape.version, SHAPE_VERSION);
assert!(shape.summary.duration_ns > 0);
assert_eq!(
*shape
.summary
.event_type_counts
.get("PollStartEvent")
.unwrap(),
2
);
assert_eq!(
*shape.summary.event_type_counts.get("PollEndEvent").unwrap(),
1
);
}
#[test]
fn extract_poll_quantiles() {
let trace = make_test_trace();
let shape = extract_shape(&trace).unwrap();
let q = shape.summary.poll_duration_quantiles.as_ref().unwrap();
assert_eq!(q.count, 1);
assert!(q.min_ns <= q.max_ns);
}
#[test]
fn span_quantiles_from_enter_exit() {
let trace = make_span_trace();
let shape = extract_shape(&trace).unwrap();
let q = shape.summary.span_duration_quantiles.as_ref().unwrap();
assert_eq!(q.count, 1);
assert!(q.min_ns > 0);
}
#[test]
fn span_callsite_correlated_and_sanitized() {
let trace = make_span_trace();
let shape = extract_shape(&trace).unwrap();
let json = serde_json::to_string(&shape).unwrap();
assert!(!json.contains("my_secret_handler"), "span callsite leaked");
assert!(json.contains("SpanEnter:"));
assert!(json.contains("SpanExit:"));
let enter = shape
.schemas
.iter()
.find(|s| s.name.starts_with("SpanEnter:"))
.unwrap();
let exit = shape
.schemas
.iter()
.find(|s| s.name.starts_with("SpanExit:"))
.unwrap();
assert_eq!(
enter.name.strip_prefix("SpanEnter:"),
exit.name.strip_prefix("SpanExit:"),
"enter/exit callsites not correlated"
);
}
#[test]
fn gzip_input() {
let trace = make_test_trace();
let mut gz = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::default());
gz.write_all(&trace).unwrap();
let compressed = gz.finish().unwrap();
let tmp = tempfile::NamedTempFile::new().unwrap();
std::fs::write(tmp.path(), &compressed).unwrap();
let data = read_trace_file(tmp.path()).unwrap();
let shape = extract_shape(&data).unwrap();
assert_eq!(shape.summary.event_count, 4);
}
#[test]
fn generate_roundtrip_decodes() {
let trace = make_test_trace();
let shape = extract_shape(&trace).unwrap();
let generated = generate_trace(&shape, 1).unwrap();
let mut dec = Decoder::new(&generated).unwrap();
let mut event_count = 0u64;
loop {
match dec.next_frame().unwrap() {
Some(DecodedFrame::Event { .. }) => event_count += 1,
Some(_) => {}
None => break,
}
}
assert_eq!(event_count, shape.summary.event_count);
}
#[test]
fn generate_repeat_scales_events() {
let trace = make_test_trace();
let shape = extract_shape(&trace).unwrap();
let generated = generate_trace(&shape, 3).unwrap();
let mut dec = Decoder::new(&generated).unwrap();
let mut count = 0u64;
loop {
match dec.next_frame().unwrap() {
Some(DecodedFrame::Event { .. }) => count += 1,
Some(_) => {}
None => break,
}
}
assert_eq!(count, shape.summary.event_count * 3);
}
#[test]
fn repeat_preserves_worker_id() {
let trace = make_test_trace();
let shape = extract_shape(&trace).unwrap();
let generated = generate_trace(&shape, 2).unwrap();
let mut dec = Decoder::new(&generated).unwrap();
let mut worker_ids: HashSet<u64> = HashSet::new();
loop {
match dec.next_frame().unwrap() {
Some(DecodedFrame::Event {
type_id, values, ..
}) => {
let entry = dec.registry().get(type_id).unwrap();
if (entry.name() == "PollStartEvent" || entry.name() == "PollEndEvent")
&& let Some(FieldValue::Varint(wid)) = values.first()
{
worker_ids.insert(*wid);
}
}
Some(_) => {}
None => break,
}
}
assert!(worker_ids.contains(&0));
assert!(worker_ids.contains(&1));
assert_eq!(worker_ids.len(), 2);
}
#[test]
fn repeat_makes_task_namespaces_disjoint() {
let mut enc = Encoder::new();
let spawn = enc
.register_schema(
"TestTaskSpawnEvent",
vec![FieldDef::new("task_id", FieldType::Varint)],
)
.unwrap();
enc.write_event(
&spawn,
&[FieldValue::Varint(BASE_TS), FieldValue::Varint(100)],
)
.unwrap();
enc.write_event(
&spawn,
&[
FieldValue::Varint(BASE_TS + 10_000),
FieldValue::Varint(200),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let generated = generate_trace(&shape, 3).unwrap();
let mut dec = Decoder::new(&generated).unwrap();
let mut all_task_ids: Vec<u64> = Vec::new();
loop {
match dec.next_frame().unwrap() {
Some(DecodedFrame::Event { values, .. }) => {
if let Some(FieldValue::Varint(tid)) = values.first() {
all_task_ids.push(*tid);
}
}
Some(_) => {}
None => break,
}
}
assert_eq!(all_task_ids.len(), 6); let rep0: HashSet<u64> = all_task_ids[..2].iter().copied().collect();
let rep1: HashSet<u64> = all_task_ids[2..4].iter().copied().collect();
let rep2: HashSet<u64> = all_task_ids[4..].iter().copied().collect();
assert!(rep0.is_disjoint(&rep1), "rep0 and rep1 task IDs overlap");
assert!(rep1.is_disjoint(&rep2), "rep1 and rep2 task IDs overlap");
}
#[test]
fn no_source_addresses_in_shape() {
let mut enc = Encoder::new();
let cpu = enc
.register_schema(
"TestCpuSampleEvent",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("tid", FieldType::Varint),
FieldDef::new("callchain", FieldType::StackFrames),
],
)
.unwrap();
let real_addrs = vec![0x7fff_dead_beef_u64, 0x5555_cafe_babe_u64];
enc.write_event(
&cpu,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(0),
FieldValue::Varint(12345),
FieldValue::StackFrames(real_addrs.into()),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let json = serde_json::to_string(&shape).unwrap();
assert!(!json.contains("deadbeef"), "real address leaked");
assert!(!json.contains("cafebabe"), "real address leaked");
}
#[test]
fn task_id_correlation_across_schemas() {
let mut enc = Encoder::new();
let spawn = enc
.register_schema(
"TestTaskSpawnEvent",
vec![FieldDef::new("task_id", FieldType::Varint)],
)
.unwrap();
let wake = enc
.register_schema(
"TestWakeEvent",
vec![
FieldDef::new("waker_task_id", FieldType::Varint),
FieldDef::new("woken_task_id", FieldType::Varint),
],
)
.unwrap();
let real_id = 0xABCD_1234_u64;
enc.write_event(
&spawn,
&[FieldValue::Varint(BASE_TS), FieldValue::Varint(real_id)],
)
.unwrap();
enc.write_event(
&wake,
&[
FieldValue::Varint(BASE_TS + 10_000),
FieldValue::Varint(real_id),
FieldValue::Varint(real_id),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let spawn_ev = &shape.events[0];
let wake_ev = &shape.events[1];
let spawn_tid = match &spawn_ev.values[0] {
ShapeValue::U(v) => *v,
_ => panic!(),
};
let waker_tid = match &wake_ev.values[0] {
ShapeValue::U(v) => *v,
_ => panic!(),
};
let woken_tid = match &wake_ev.values[1] {
ShapeValue::U(v) => *v,
_ => panic!(),
};
assert_eq!(spawn_tid, waker_tid, "task_id/waker not correlated");
assert_eq!(spawn_tid, woken_tid, "task_id/woken not correlated");
assert_ne!(spawn_tid, real_id, "original ID leaked");
}
#[test]
fn pooled_string_missing_id_errors() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"TestEvent",
vec![FieldDef::new("name", FieldType::PooledString)],
)
.unwrap();
let _id = enc.intern_string("hello").unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::PooledString(dial9_trace_format::types::InternedString::from_raw(0)),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace);
assert!(shape.is_ok());
let mut enc2 = Encoder::new();
let schema2 = enc2
.register_schema(
"TestEvent",
vec![FieldDef::new("name", FieldType::PooledString)],
)
.unwrap();
enc2.write_event(
&schema2,
&[
FieldValue::Varint(BASE_TS),
FieldValue::PooledString(dial9_trace_format::types::InternedString::from_raw(999)),
],
)
.unwrap();
let trace2 = enc2.finish();
let result = extract_shape(&trace2);
assert!(
result.is_err(),
"expected error for missing pooled string ID"
);
let err_msg = result.unwrap_err().to_string();
assert!(
err_msg.contains("missing pooled string") || err_msg.contains("pass 2"),
"unexpected error: {err_msg}"
);
}
#[test]
fn out_of_order_timestamps_global_base() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"TestTimestampEvent",
vec![FieldDef::new("worker_id", FieldType::Varint)],
)
.unwrap();
enc.write_event(
&schema,
&[FieldValue::Varint(BASE_TS + 200_000), FieldValue::Varint(0)],
)
.unwrap();
enc.write_event(
&schema,
&[FieldValue::Varint(BASE_TS), FieldValue::Varint(0)],
)
.unwrap();
enc.write_event(
&schema,
&[FieldValue::Varint(BASE_TS + 100_000), FieldValue::Varint(0)],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let offsets: Vec<u64> = shape
.events
.iter()
.map(|e| e.timestamp_offset_ns.unwrap())
.collect();
assert_eq!(offsets[1], 0, "minimum timestamp should produce offset 0");
assert!(
offsets[0] > offsets[1],
"first event should have larger offset"
);
}
#[test]
fn fixed_width_normalization() {
assert_eq!(normalize_field_type(FieldType::U8), FieldType::Varint);
assert_eq!(normalize_field_type(FieldType::U16), FieldType::Varint);
assert_eq!(normalize_field_type(FieldType::U32), FieldType::Varint);
assert_eq!(
normalize_field_type(FieldType::OptionalU8),
FieldType::OptionalVarint
);
assert_eq!(
normalize_field_type(FieldType::OptionalU16),
FieldType::OptionalVarint
);
assert_eq!(
normalize_field_type(FieldType::OptionalU32),
FieldType::OptionalVarint
);
assert_eq!(normalize_field_type(FieldType::Varint), FieldType::Varint);
assert_eq!(normalize_field_type(FieldType::I64), FieldType::I64);
}
#[test]
fn validate_rejects_bad_version() {
let shape = TraceShape {
version: 999,
summary: ShapeSummary {
event_count: 0,
duration_ns: 0,
event_type_counts: BTreeMap::new(),
worker_cardinality: 0,
task_cardinality: 0,
poll_duration_quantiles: None,
span_duration_quantiles: None,
},
schemas: vec![ShapeSchema {
name: "T".into(),
has_timestamp: true,
fields: vec![],
annotations: vec![],
}],
events: vec![ShapeEvent {
schema_index: 0,
timestamp_offset_ns: Some(0),
values: vec![],
}],
};
assert!(validate_shape(&shape).is_err());
}
#[test]
fn validate_rejects_empty_events() {
let shape = TraceShape {
version: SHAPE_VERSION,
summary: ShapeSummary {
event_count: 0,
duration_ns: 0,
event_type_counts: BTreeMap::new(),
worker_cardinality: 0,
task_cardinality: 0,
poll_duration_quantiles: None,
span_duration_quantiles: None,
},
schemas: vec![ShapeSchema {
name: "T".into(),
has_timestamp: true,
fields: vec![],
annotations: vec![],
}],
events: vec![],
};
assert!(validate_shape(&shape).is_err());
}
#[test]
fn validate_rejects_invalid_schema_ref() {
let shape = TraceShape {
version: SHAPE_VERSION,
summary: ShapeSummary {
event_count: 1,
duration_ns: 0,
event_type_counts: [("T".into(), 1)].into_iter().collect(),
worker_cardinality: 0,
task_cardinality: 0,
poll_duration_quantiles: None,
span_duration_quantiles: None,
},
schemas: vec![ShapeSchema {
name: "T".into(),
has_timestamp: true,
fields: vec![],
annotations: vec![],
}],
events: vec![ShapeEvent {
schema_index: 99,
timestamp_offset_ns: Some(0),
values: vec![],
}],
};
assert!(validate_shape(&shape).is_err());
}
#[test]
fn validate_rejects_fixed_width_tags() {
let shape = TraceShape {
version: SHAPE_VERSION,
summary: ShapeSummary {
event_count: 1,
duration_ns: 0,
event_type_counts: [("T".into(), 1)].into_iter().collect(),
worker_cardinality: 0,
task_cardinality: 0,
poll_duration_quantiles: None,
span_duration_quantiles: None,
},
schemas: vec![ShapeSchema {
name: "T".into(),
has_timestamp: true,
fields: vec![ShapeField {
name: "x".into(),
field_type: FieldType::U8 as u8,
repeat_meta: None,
}],
annotations: vec![],
}],
events: vec![ShapeEvent {
schema_index: 0,
timestamp_offset_ns: Some(0),
values: vec![ShapeValue::U(1)],
}],
};
let err = validate_shape(&shape).unwrap_err();
assert!(err.to_string().contains("rejected fixed-width"), "{err}");
}
#[test]
fn validate_rejects_untimestamped_events() {
let shape = TraceShape {
version: SHAPE_VERSION,
summary: ShapeSummary {
event_count: 1,
duration_ns: 0,
event_type_counts: [("T".into(), 1)].into_iter().collect(),
worker_cardinality: 0,
task_cardinality: 0,
poll_duration_quantiles: None,
span_duration_quantiles: None,
},
schemas: vec![ShapeSchema {
name: "T".into(),
has_timestamp: false,
fields: vec![],
annotations: vec![],
}],
events: vec![ShapeEvent {
schema_index: 0,
timestamp_offset_ns: None,
values: vec![],
}],
};
let err = validate_shape(&shape).unwrap_err();
assert!(err.to_string().contains("has_timestamp=false"), "{err}");
}
#[test]
fn validate_rejects_type_mismatch() {
let shape = TraceShape {
version: SHAPE_VERSION,
summary: ShapeSummary {
event_count: 1,
duration_ns: 0,
event_type_counts: [("T".into(), 1)].into_iter().collect(),
worker_cardinality: 0,
task_cardinality: 0,
poll_duration_quantiles: None,
span_duration_quantiles: None,
},
schemas: vec![ShapeSchema {
name: "T".into(),
has_timestamp: true,
fields: vec![ShapeField {
name: "x".into(),
field_type: FieldType::Varint as u8,
repeat_meta: None,
}],
annotations: vec![],
}],
events: vec![ShapeEvent {
schema_index: 0,
timestamp_offset_ns: Some(0),
values: vec![ShapeValue::S("wrong".into())], }],
};
assert!(validate_shape(&shape).is_err());
}
#[test]
fn validate_rejects_none_for_non_optional() {
let shape = TraceShape {
version: SHAPE_VERSION,
summary: ShapeSummary {
event_count: 1,
duration_ns: 0,
event_type_counts: [("T".into(), 1)].into_iter().collect(),
worker_cardinality: 0,
task_cardinality: 0,
poll_duration_quantiles: None,
span_duration_quantiles: None,
},
schemas: vec![ShapeSchema {
name: "T".into(),
has_timestamp: true,
fields: vec![ShapeField {
name: "x".into(),
field_type: FieldType::Varint as u8,
repeat_meta: None,
}],
annotations: vec![],
}],
events: vec![ShapeEvent {
schema_index: 0,
timestamp_offset_ns: Some(0),
values: vec![ShapeValue::None],
}],
};
assert!(validate_shape(&shape).is_err());
}
#[test]
fn overflow_repeat_rejected() {
let trace = make_test_trace();
let mut shape = extract_shape(&trace).unwrap();
shape.events[0].timestamp_offset_ns = Some(u64::MAX - 1);
shape.summary.duration_ns = u64::MAX - 1;
let result = generate_trace(&shape, 2);
assert!(result.is_err(), "should fail for overflow");
}
#[test]
fn quantize_numeric_no_overflow() {
assert_eq!(quantize_numeric(0), 0);
assert_eq!(quantize_numeric(1), 1);
assert_eq!(quantize_numeric(1023), 1023);
let result = quantize_numeric(u64::MAX);
assert!(result > 0);
assert!(result >= (1u64 << 54));
}
#[test]
fn quantize_i64_extremes() {
assert_eq!(quantize_i64(0), 0);
let min_result = quantize_i64(i64::MIN);
assert!(min_result < 0);
let max_result = quantize_i64(i64::MAX);
assert!(max_result > 0);
}
#[test]
fn quantize_f64_rejects_non_finite() {
assert!(quantize_f64(f64::NAN).is_err());
assert!(quantize_f64(f64::INFINITY).is_err());
assert!(quantize_f64(f64::NEG_INFINITY).is_err());
assert!(quantize_f64(1.23456).is_ok());
}
#[test]
fn malformed_input_rejected() {
assert!(extract_shape(&[]).is_err());
assert!(extract_shape(&[0x01, 0x02]).is_err());
}
#[test]
fn optional_field_roundtrip() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"TestEvent",
vec![
FieldDef::new("required", FieldType::Varint),
FieldDef::new("optional", FieldType::OptionalVarint),
],
)
.unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(42),
FieldValue::None,
],
)
.unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS + 10_000),
FieldValue::Varint(7),
FieldValue::Varint(99),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
assert_eq!(shape.events.len(), 2);
assert!(matches!(&shape.events[0].values[1], ShapeValue::None));
assert!(matches!(&shape.events[1].values[1], ShapeValue::U(_)));
let generated = generate_trace(&shape, 1).unwrap();
let mut dec = Decoder::new(&generated).unwrap();
let mut events = Vec::new();
loop {
match dec.next_frame().unwrap() {
Some(DecodedFrame::Event { values, .. }) => events.push(values),
Some(_) => {}
None => break,
}
}
assert_eq!(events.len(), 2);
assert!(matches!(&events[0][1], FieldValue::None));
assert!(matches!(&events[1][1], FieldValue::Varint(_)));
}
#[test]
fn string_map_roundtrip() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"TestEvent",
vec![FieldDef::new("tags", FieldType::StringMap)],
)
.unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::StringMap(vec![
(b"key1".to_vec(), b"val1".to_vec()),
(b"secret_key".to_vec(), b"secret_val".to_vec()),
]),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let json = serde_json::to_string(&shape).unwrap();
assert!(!json.contains("secret_key"));
assert!(!json.contains("secret_val"));
let generated = generate_trace(&shape, 1).unwrap();
let mut dec = Decoder::new(&generated).unwrap();
loop {
match dec.next_frame().unwrap() {
Some(DecodedFrame::Event { values, .. }) => {
assert!(matches!(&values[0], FieldValue::StringMap(_)));
break;
}
Some(_) => {}
None => panic!("no events"),
}
}
}
#[test]
fn clock_sync_uses_synthetic_epoch() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"ClockSyncEvent",
vec![FieldDef::new("realtime_ns", FieldType::Varint)],
)
.unwrap();
let real_epoch = 1_700_000_000_000_000_000u64;
enc.write_event(
&schema,
&[FieldValue::Varint(BASE_TS), FieldValue::Varint(real_epoch)],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let json = serde_json::to_string(&shape).unwrap();
assert!(
!json.contains("1700000000000000000"),
"source realtime leaked"
);
assert!(
!json.contains(&SYNTHETIC_EPOCH_NS.to_string()),
"synthetic epoch appeared in JSON"
);
let clock_event = &shape.events[0];
let realtime_val = match &clock_event.values[0] {
ShapeValue::U(v) => *v,
other => panic!("expected U for realtime_ns, got {other:?}"),
};
assert_eq!(
realtime_val, 0,
"realtime should be event's monotonic offset (0 for first event)"
);
let generated = generate_trace(&shape, 1).unwrap();
let mut dec = Decoder::new(&generated).unwrap();
loop {
match dec.next_frame().unwrap() {
Some(DecodedFrame::Event { values, .. }) => {
if let Some(FieldValue::Varint(rt)) = values.first() {
assert!(
*rt >= SYNTHETIC_EPOCH_NS,
"generated realtime should be >= synthetic epoch"
);
assert_ne!(
*rt, real_epoch,
"generated realtime must differ from source"
);
}
break;
}
Some(_) => {}
None => panic!("no events"),
}
}
}
#[test]
fn alloc_timestamp_shifts_with_event() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"FreeEvent",
vec![
FieldDef::new("tid", FieldType::Varint),
FieldDef::new("addr", FieldType::Varint),
FieldDef::new("size", FieldType::Varint),
FieldDef::new("alloc_timestamp_ns", FieldType::Varint),
],
)
.unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS + 100_000),
FieldValue::Varint(1001),
FieldValue::Varint(0xCAFE),
FieldValue::Varint(1024),
FieldValue::Varint(BASE_TS + 50_000),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let generated = generate_trace(&shape, 2).unwrap();
let mut dec = Decoder::new(&generated).unwrap();
let mut alloc_ts_values: Vec<u64> = Vec::new();
let mut event_ts_values: Vec<u64> = Vec::new();
loop {
match dec.next_frame().unwrap() {
Some(DecodedFrame::Event {
values,
timestamp_ns,
..
}) => {
if let Some(ts) = timestamp_ns {
event_ts_values.push(ts);
}
if let Some(FieldValue::Varint(v)) = values.get(3) {
alloc_ts_values.push(*v);
}
}
Some(_) => {}
None => break,
}
}
assert_eq!(alloc_ts_values.len(), 2);
assert!(
alloc_ts_values[1] > alloc_ts_values[0],
"alloc_timestamp_ns should shift with repetition"
);
}
#[test]
fn file_roundtrip() {
let trace = make_test_trace();
let tmp = tempfile::tempdir().unwrap();
let trace_path = tmp.path().join("input.bin");
let shape_path = tmp.path().join("shape.json");
let output_path = tmp.path().join("output.bin");
std::fs::write(&trace_path, &trace).unwrap();
extract(&trace_path, &shape_path).unwrap();
generate(&shape_path, &output_path, 1).unwrap();
let output = std::fs::read(&output_path).unwrap();
let mut dec = Decoder::new(&output).unwrap();
let mut event_count = 0u64;
loop {
match dec.next_frame().unwrap() {
Some(DecodedFrame::Event { .. }) => event_count += 1,
Some(_) => {}
None => break,
}
}
assert_eq!(
dec.position(),
dec.data_len(),
"decoder did not consume all bytes: position {} != data_len {}",
dec.position(),
dec.data_len()
);
assert!(event_count > 0, "no events in generated trace");
let shape_json = std::fs::read_to_string(&shape_path).unwrap();
let shape: TraceShape = serde_json::from_str(&shape_json).unwrap();
assert_eq!(
event_count, shape.summary.event_count,
"generated event count doesn't match shape summary"
);
}
#[test]
fn synthesize_matches_in_memory_generation_and_removes_source_secrets() {
let trace = make_test_trace();
let shape = extract_shape(&trace).unwrap();
let expected = generate_trace(&shape, 2).unwrap();
let tmp = tempfile::tempdir().unwrap();
let source_path = tmp.path().join("source.bin");
let output_path = tmp.path().join("synthetic.bin");
std::fs::write(&source_path, &trace).unwrap();
synthesize(&source_path, &output_path, 2).unwrap();
let generated = std::fs::read(&output_path).unwrap();
assert_eq!(generated, expected);
for secret in ["password-123", "MySecretService", "secret_name"] {
assert!(
!generated
.windows(secret.len())
.any(|window| window == secret.as_bytes()),
"source secret leaked into synthetic trace: {secret}"
);
}
let mut decoder = Decoder::new(&generated).unwrap();
let mut event_count = 0u64;
while let Some(frame) = decoder.next_frame().unwrap() {
if matches!(frame, DecodedFrame::Event { .. }) {
event_count += 1;
}
}
assert_eq!(event_count, shape.summary.event_count * 2);
assert_eq!(decoder.position(), decoder.data_len());
}
#[test]
fn synthesize_rejects_zero_repeat_before_creating_output() {
let trace = make_test_trace();
let tmp = tempfile::tempdir().unwrap();
let source_path = tmp.path().join("source.bin");
let output_path = tmp.path().join("synthetic.bin");
std::fs::write(&source_path, trace).unwrap();
let error = synthesize(&source_path, &output_path, 0).unwrap_err();
assert!(error.to_string().contains("--repeat must be >= 1"));
assert!(!output_path.exists());
}
#[test]
fn tagged_serde_roundtrip() {
let values = vec![
ShapeValue::U(42),
ShapeValue::I(-7),
ShapeValue::F(1.234),
ShapeValue::B(true),
ShapeValue::S("hello".into()),
ShapeValue::PS("pooled".into()),
ShapeValue::Bytes(100),
ShapeValue::Stack(vec![0x1000, 0x2000]),
ShapeValue::PStack(vec![0x3000, 0x4000]),
ShapeValue::StringMap(vec![("k".into(), "v".into())]),
ShapeValue::List(vec![ShapeValue::U(1), ShapeValue::S("x".into())]),
ShapeValue::Map(vec![(ShapeValue::U(1), ShapeValue::B(false))]),
ShapeValue::None,
];
for v in &values {
let json = serde_json::to_string(v).unwrap();
let parsed: ShapeValue = serde_json::from_str(&json).unwrap();
let json2 = serde_json::to_string(&parsed).unwrap();
assert_eq!(json, json2, "value {v:?} didn't roundtrip");
}
}
#[test]
fn dynamic_list_roundtrip() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"TestEvent",
vec![FieldDef::new("items", FieldType::DynamicList)],
)
.unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::List(vec![
FieldValue::Varint(1),
FieldValue::String("hello".into()),
FieldValue::Bool(true),
]),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
match &shape.events[0].values[0] {
ShapeValue::List(items) => assert_eq!(items.len(), 3),
other => panic!("expected List, got {other:?}"),
}
let generated = generate_trace(&shape, 1).unwrap();
let mut dec = Decoder::new(&generated).unwrap();
loop {
match dec.next_frame().unwrap() {
Some(DecodedFrame::Event { values, .. }) => {
match &values[0] {
FieldValue::List(items) => assert_eq!(items.len(), 3),
other => panic!("expected List, got {other:?}"),
}
break;
}
Some(_) => {}
None => panic!("no events"),
}
}
}
#[test]
fn pooled_stack_roundtrip() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"TestCpuSampleEvent",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("callchain", FieldType::PooledStackFrames),
],
)
.unwrap();
let frames = vec![0x1000u64, 0x2000, 0x3000];
let pooled = enc.intern_stack_frames(&frames).unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(0),
FieldValue::PooledStackFrames(pooled),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
match &shape.events[0].values[1] {
ShapeValue::PStack(addrs) => assert_eq!(addrs.len(), 3),
other => panic!("expected PStack, got {other:?}"),
}
let generated = generate_trace(&shape, 1).unwrap();
let mut dec = Decoder::new(&generated).unwrap();
loop {
match dec.next_frame().unwrap() {
Some(DecodedFrame::Event { values, .. }) => {
match &values[1] {
FieldValue::PooledStackFrames(_) => {}
other => panic!("expected PooledStackFrames, got {other:?}"),
}
break;
}
Some(_) => {}
None => panic!("no events"),
}
}
}
#[test]
fn addr_namespace_disjoint_on_repeat() {
let mut enc = Encoder::new();
let sym = enc
.register_schema(
"TestSymbolTableEntry",
vec![
FieldDef::new("addr", FieldType::Varint),
FieldDef::new("size", FieldType::Varint),
],
)
.unwrap();
enc.write_event(
&sym,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(0x1000),
FieldValue::Varint(0x100),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let generated = generate_trace(&shape, 2).unwrap();
let mut dec = Decoder::new(&generated).unwrap();
let mut all_addrs: Vec<(u64, u64)> = Vec::new();
loop {
match dec.next_frame().unwrap() {
Some(DecodedFrame::Event { values, .. }) => {
if let (Some(FieldValue::Varint(a)), Some(FieldValue::Varint(b))) =
(values.first(), values.get(1))
{
all_addrs.push((*a, *b));
}
}
Some(_) => {}
None => break,
}
}
assert_eq!(all_addrs.len(), 2);
assert_ne!(
all_addrs[0], all_addrs[1],
"addresses should differ across reps"
);
}
#[test]
fn extract_demo_trace_succeeds() {
let demo_path =
std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("ui/public/demo-trace.bin");
if !demo_path.exists() {
eprintln!(
"skipping: demo-trace.bin not found at {}",
demo_path.display()
);
return;
}
let data = read_trace_file(&demo_path).unwrap();
let shape = extract_shape(&data).expect("demo trace extraction must succeed");
assert!(
shape.summary.event_count > 1000,
"demo trace should have many events"
);
assert!(
shape.schemas.len() > 5,
"demo trace should have many schemas"
);
let json = serde_json::to_string(&shape).unwrap();
assert!(!json.contains("1700000000"), "source timestamps leaked");
let generated = generate_trace(&shape, 1).expect("generate from demo shape must succeed");
let mut dec = Decoder::new(&generated).unwrap();
let mut count = 0u64;
dec.for_each_event(|_| count += 1)
.expect("decode generated trace");
assert_eq!(count, shape.summary.event_count);
}
#[test]
fn clock_privacy_source_epoch_not_in_json() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"ClockSyncEvent",
vec![FieldDef::new("realtime_ns", FieldType::Varint)],
)
.unwrap();
let source_realtime = 1_700_000_000_123_456_789u64;
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(source_realtime),
],
)
.unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS + 100_000),
FieldValue::Varint(source_realtime + 100_000),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let json = serde_json::to_string(&shape).unwrap();
assert!(!json.contains("1700000000"), "source epoch leaked");
let ev1 = &shape.events[1];
let rt_val = match &ev1.values[0] {
ShapeValue::U(v) => *v,
other => panic!("expected U, got {other:?}"),
};
assert_eq!(rt_val, quantize_ns(100_000));
let generated = generate_trace(&shape, 1).unwrap();
let mut dec = Decoder::new(&generated).unwrap();
let mut gen_realtimes = Vec::new();
dec.for_each_event(|ev| {
if ev.name == "ClockSyncEvent"
&& let Some(dial9_trace_format::types::FieldValueRef::Varint(rt)) =
ev.fields.first()
{
gen_realtimes.push(*rt);
}
})
.unwrap();
assert_eq!(gen_realtimes.len(), 2);
for rt in &gen_realtimes {
assert_ne!(
*rt, source_realtime,
"generated realtime must differ from source"
);
assert!(
*rt >= SYNTHETIC_EPOCH_NS,
"should be based on synthetic epoch"
);
}
}
#[test]
fn alloc_timestamp_remains_relative_monotonic() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"AllocEvent",
vec![
FieldDef::new("tid", FieldType::Varint),
FieldDef::new("size", FieldType::Varint),
FieldDef::new("addr", FieldType::Varint),
FieldDef::new("callchain", FieldType::PooledStackFrames),
],
)
.unwrap();
let free_schema = enc
.register_schema(
"FreeEvent",
vec![
FieldDef::new("tid", FieldType::Varint),
FieldDef::new("addr", FieldType::Varint),
FieldDef::new("size", FieldType::Varint),
FieldDef::new("alloc_timestamp_ns", FieldType::Varint),
],
)
.unwrap();
let frames = enc.intern_stack_frames(&[0x1000]).unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(1001),
FieldValue::Varint(1024),
FieldValue::Varint(0xCAFE),
FieldValue::PooledStackFrames(frames),
],
)
.unwrap();
enc.write_event(
&free_schema,
&[
FieldValue::Varint(BASE_TS + 50_000),
FieldValue::Varint(1001),
FieldValue::Varint(0xCAFE),
FieldValue::Varint(512),
FieldValue::Varint(BASE_TS), ],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let free_ev = &shape.events[1];
let alloc_ts = match &free_ev.values[3] {
ShapeValue::U(v) => *v,
other => panic!("expected U for alloc_timestamp_ns, got {other:?}"),
};
assert_eq!(alloc_ts, 0, "alloc_timestamp_ns should be relative offset");
let free_schema_shape = &shape.schemas[1];
let alloc_ts_field = &free_schema_shape.fields[3];
assert_eq!(
alloc_ts_field.repeat_meta.as_ref().unwrap().semantics,
FieldSemantics::TimestampRef
);
}
#[test]
fn namespace_json_has_no_customer_field_text() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"TestIdentityEvent",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("task_id", FieldType::Varint),
],
)
.unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(0),
FieldValue::Varint(42),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let json = serde_json::to_string(&shape).unwrap();
assert!(!json.contains("\"task\""), "string namespace 'task' leaked");
assert!(!json.contains("\"span\""), "string namespace 'span' leaked");
assert!(!json.contains("\"tid\""), "string namespace 'tid' leaked");
assert!(!json.contains("\"addr\""), "string namespace 'addr' leaked");
assert!(
!json.contains("\"socket_cookie\""),
"string namespace leaked"
);
assert!(
json.contains("\"namespace\":0"),
"task namespace should be serialized as 0"
);
}
#[test]
fn malicious_builtin_schema_is_rejected() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"PollStartEvent",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("customer_name", FieldType::String),
FieldDef::new("customer_account_id", FieldType::Varint),
],
)
.unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(0),
FieldValue::String("John Doe".into()),
FieldValue::Varint(12345),
],
)
.unwrap();
let trace = enc.finish();
let result = extract_shape(&trace);
assert!(
result.is_err(),
"builtin schema with non-canonical fields should be rejected"
);
let err_msg = format!("{:#}", result.unwrap_err());
assert!(
err_msg.contains("does not exactly match a known complete signature"),
"unexpected error: {err_msg}"
);
}
#[test]
fn base_addr_remapped_as_addr_identity() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"TestEvent",
vec![FieldDef::new("base_addr", FieldType::Varint)],
)
.unwrap();
let real_addr = 0x7FFF_DEAD_BEEF_u64;
enc.write_event(
&schema,
&[FieldValue::Varint(BASE_TS), FieldValue::Varint(real_addr)],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let json = serde_json::to_string(&shape).unwrap();
assert!(!json.contains("deadbeef"), "real base_addr leaked");
let ev = &shape.events[0];
let val = match &ev.values[0] {
ShapeValue::U(v) => *v,
other => panic!("expected U, got {other:?}"),
};
assert_ne!(val, real_addr, "base_addr should be remapped");
assert_ne!(val, 0);
}
#[test]
fn custom_small_integer_not_retained() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"MyCustomEvent",
vec![
FieldDef::new("count", FieldType::Varint),
FieldDef::new("score", FieldType::Varint),
],
)
.unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(3), FieldValue::Varint(7), ],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let ev = &shape.events[0];
let count_val = match &ev.values[0] {
ShapeValue::U(v) => *v,
other => panic!("expected U, got {other:?}"),
};
let score_val = match &ev.values[1] {
ShapeValue::U(v) => *v,
other => panic!("expected U, got {other:?}"),
};
assert_ne!(
count_val, 3,
"source small integer 3 should not be retained"
);
assert_ne!(
score_val, 7,
"source small integer 7 should not be retained"
);
assert_eq!(count_val, privacy_bucket_u64(3));
assert_eq!(score_val, privacy_bucket_u64(7));
assert!(count_val > 0 && count_val <= 12);
assert!(score_val > 0 && score_val <= 28);
}
#[test]
fn generic_id_namespace_serialized_as_number() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"TestEvent",
vec![
FieldDef::new("correlation_id", FieldType::Varint),
FieldDef::new("request_id", FieldType::Varint),
],
)
.unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(999),
FieldValue::Varint(888),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let json = serde_json::to_string(&shape).unwrap();
assert!(
!json.contains("correlation_id"),
"original field name leaked in namespace"
);
assert!(
!json.contains("request_id"),
"original field name leaked in namespace"
);
assert!(
json.contains("\"namespace\":100"),
"first anon namespace should be 100"
);
assert!(
json.contains("\"namespace\":101"),
"second anon namespace should be 101"
);
}
#[test]
fn span_preserves_only_allowed_fields() {
let mut enc = Encoder::new();
let enter = enc
.register_schema(
"SpanEnter:my_handler",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
FieldDef::new("parent_span_id", FieldType::Varint),
FieldDef::new("customer_data", FieldType::String),
],
)
.unwrap();
enc.write_event(
&enter,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(0),
FieldValue::Varint(42),
FieldValue::Varint(0),
FieldValue::String("secret".into()),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let schema = &shape.schemas[0];
assert_eq!(schema.fields[0].name, "worker_id");
assert_eq!(schema.fields[1].name, "span_id");
assert_eq!(schema.fields[2].name, "parent_span_id");
assert!(schema.fields[3].name.starts_with("field_"));
assert!(!schema.fields[3].name.contains("customer"));
}
#[test]
fn dynamic_list_with_pooled_string_and_stack() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"TestEvent",
vec![FieldDef::new("items", FieldType::DynamicList)],
)
.unwrap();
let ps_id = enc.intern_string("hello_pooled").unwrap();
let stack_id = enc.intern_stack_frames(&[0x1000, 0x2000, 0x3000]).unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::List(vec![
FieldValue::Varint(42),
FieldValue::PooledString(ps_id),
FieldValue::PooledStackFrames(stack_id),
FieldValue::String("inline_str".into()),
]),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let items = match &shape.events[0].values[0] {
ShapeValue::List(items) => items,
other => panic!("expected List, got {other:?}"),
};
assert_eq!(items.len(), 4);
assert!(matches!(&items[0], ShapeValue::U(_)));
assert!(matches!(&items[1], ShapeValue::PS(_)));
assert!(matches!(&items[2], ShapeValue::PStack(addrs) if addrs.len() == 3));
assert!(matches!(&items[3], ShapeValue::S(_)));
let generated = generate_trace(&shape, 1).unwrap();
let mut dec = Decoder::new(&generated).unwrap();
dec.for_each_event(|ev| {
if ev.name == "TestEvent" {
let list = match &ev.fields[0] {
dial9_trace_format::types::FieldValueRef::List(l) => l,
other => panic!("expected List, got {other:?}"),
};
let items: Vec<_> = list.iter().collect();
assert_eq!(items.len(), 4);
assert!(matches!(
items[0],
dial9_trace_format::types::FieldValueRef::Varint(_)
));
assert!(matches!(
items[1],
dial9_trace_format::types::FieldValueRef::PooledString(_)
));
if let dial9_trace_format::types::FieldValueRef::PooledString(id) = items[1] {
assert!(
ev.string_pool.get(*id).is_some(),
"pooled string should resolve from pool"
);
}
assert!(matches!(
items[2],
dial9_trace_format::types::FieldValueRef::PooledStackFrames(_)
));
if let dial9_trace_format::types::FieldValueRef::PooledStackFrames(id) = items[2] {
let frames = ev.stack_pool.get(*id);
assert!(frames.is_some(), "pooled stack should resolve from pool");
assert_eq!(frames.unwrap().len(), 3);
}
assert!(matches!(
items[3],
dial9_trace_format::types::FieldValueRef::String(_)
));
}
})
.unwrap();
}
#[test]
fn dynamic_map_with_pooled_values() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"TestEvent",
vec![FieldDef::new("metadata", FieldType::DynamicMap)],
)
.unwrap();
let ps_key = enc.intern_string("map_key").unwrap();
let ps_val = enc.intern_string("map_value").unwrap();
let stack_id = enc.intern_stack_frames(&[0xABC, 0xDEF]).unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Map(vec![
(
FieldValue::PooledString(ps_key),
FieldValue::PooledStackFrames(stack_id),
),
(
FieldValue::String("inline_key".into()),
FieldValue::PooledString(ps_val),
),
]),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let pairs = match &shape.events[0].values[0] {
ShapeValue::Map(pairs) => pairs,
other => panic!("expected Map, got {other:?}"),
};
assert_eq!(pairs.len(), 2);
assert!(matches!(&pairs[0].0, ShapeValue::PS(_)));
assert!(matches!(&pairs[0].1, ShapeValue::PStack(addrs) if addrs.len() == 2));
assert!(matches!(&pairs[1].0, ShapeValue::S(_)));
assert!(matches!(&pairs[1].1, ShapeValue::PS(_)));
let generated = generate_trace(&shape, 1).unwrap();
let mut dec = Decoder::new(&generated).unwrap();
dec.for_each_event(|ev| {
if ev.name == "TestEvent" {
let map = match &ev.fields[0] {
dial9_trace_format::types::FieldValueRef::Map(m) => m,
other => panic!("expected Map, got {other:?}"),
};
let entries: Vec<_> = map.iter().collect();
assert_eq!(entries.len(), 2);
assert!(matches!(
entries[0].0,
dial9_trace_format::types::FieldValueRef::PooledString(_)
));
assert!(matches!(
entries[0].1,
dial9_trace_format::types::FieldValueRef::PooledStackFrames(_)
));
if let dial9_trace_format::types::FieldValueRef::PooledString(id) = entries[0].0 {
assert!(ev.string_pool.get(*id).is_some());
}
if let dial9_trace_format::types::FieldValueRef::PooledStackFrames(id) =
entries[0].1
{
let frames = ev.stack_pool.get(*id).unwrap();
assert_eq!(frames.len(), 2);
}
}
})
.unwrap();
}
#[test]
fn validation_rejects_nested_none() {
let shape = TraceShape {
version: SHAPE_VERSION,
summary: ShapeSummary {
event_count: 1,
duration_ns: 0,
event_type_counts: [("T".into(), 1)].into_iter().collect(),
worker_cardinality: 0,
task_cardinality: 0,
poll_duration_quantiles: None,
span_duration_quantiles: None,
},
schemas: vec![ShapeSchema {
name: "T".into(),
has_timestamp: true,
fields: vec![ShapeField {
name: "items".into(),
field_type: FieldType::DynamicList as u8,
repeat_meta: None,
}],
annotations: vec![],
}],
events: vec![ShapeEvent {
schema_index: 0,
timestamp_offset_ns: Some(0),
values: vec![ShapeValue::List(vec![ShapeValue::U(1), ShapeValue::None])],
}],
};
let err = validate_shape(&shape).unwrap_err();
assert!(
err.to_string().contains("nested None"),
"unexpected error: {err}"
);
}
#[test]
fn privacy_bucket_transforms_all_nonzero() {
assert_eq!(privacy_bucket_u64(0), 0);
assert_ne!(privacy_bucket_u64(1), 1, "1 must be transformed");
assert_ne!(privacy_bucket_u64(2), 2, "2 must be transformed");
assert_ne!(privacy_bucket_u64(3), 3, "3 must be transformed");
assert_ne!(privacy_bucket_u64(4), 4, "4 must be transformed");
assert_ne!(privacy_bucket_u64(5), 5);
assert_ne!(privacy_bucket_u64(7), 7);
assert_ne!(privacy_bucket_u64(8), 8);
assert_ne!(privacy_bucket_u64(16), 16);
assert_ne!(privacy_bucket_u64(100), 100);
assert_ne!(privacy_bucket_u64(1024), 1024);
let large = privacy_bucket_u64(u64::MAX);
assert!(large > 0);
assert_ne!(large, u64::MAX);
for shift in 0..63 {
let v = 1u64 << shift;
assert_ne!(
privacy_bucket_u64(v),
v,
"power of 2 (2^{shift} = {v}) must be transformed"
);
}
for v in [3u64, 7, 15, 31, 100, 1000, 10000, 1_000_000] {
let bucketed = privacy_bucket_u64(v);
assert!(
bucketed > 0 && bucketed <= v * 4,
"privacy_bucket_u64({v}) = {bucketed}, expected within 4x"
);
}
assert_eq!(privacy_bucket_i64(0), 0);
assert_ne!(privacy_bucket_i64(1), 1);
assert_ne!(privacy_bucket_i64(-1), -1);
let neg = privacy_bucket_i64(-100);
assert!(neg < 0, "negative should stay negative");
assert_ne!(neg, -100);
}
fn minimal_shape(fields: Vec<ShapeField>, values: Vec<ShapeValue>) -> TraceShape {
TraceShape {
version: SHAPE_VERSION,
summary: ShapeSummary {
event_count: 1,
duration_ns: 0,
event_type_counts: [("T".into(), 1)].into_iter().collect(),
worker_cardinality: 0,
task_cardinality: 0,
poll_duration_quantiles: None,
span_duration_quantiles: None,
},
schemas: vec![ShapeSchema {
name: "T".into(),
has_timestamp: true,
fields,
annotations: vec![],
}],
events: vec![ShapeEvent {
schema_index: 0,
timestamp_offset_ns: Some(0),
values,
}],
}
}
#[test]
fn validate_rejects_empty_schema_name() {
let shape = TraceShape {
version: SHAPE_VERSION,
summary: ShapeSummary {
event_count: 1,
duration_ns: 0,
event_type_counts: [("".into(), 1)].into_iter().collect(),
worker_cardinality: 0,
task_cardinality: 0,
poll_duration_quantiles: None,
span_duration_quantiles: None,
},
schemas: vec![ShapeSchema {
name: "".into(),
has_timestamp: true,
fields: vec![],
annotations: vec![],
}],
events: vec![ShapeEvent {
schema_index: 0,
timestamp_offset_ns: Some(0),
values: vec![],
}],
};
let err = validate_shape(&shape).unwrap_err();
assert!(err.to_string().contains("empty name"), "{err}");
}
#[test]
fn validate_rejects_duplicate_schema_names() {
let shape = TraceShape {
version: SHAPE_VERSION,
summary: ShapeSummary {
event_count: 2,
duration_ns: 0,
event_type_counts: [("Dup".into(), 2)].into_iter().collect(),
worker_cardinality: 0,
task_cardinality: 0,
poll_duration_quantiles: None,
span_duration_quantiles: None,
},
schemas: vec![
ShapeSchema {
name: "Dup".into(),
has_timestamp: true,
fields: vec![],
annotations: vec![],
},
ShapeSchema {
name: "Dup".into(),
has_timestamp: true,
fields: vec![],
annotations: vec![],
},
],
events: vec![
ShapeEvent {
schema_index: 0,
timestamp_offset_ns: Some(0),
values: vec![],
},
ShapeEvent {
schema_index: 1,
timestamp_offset_ns: Some(0),
values: vec![],
},
],
};
let err = validate_shape(&shape).unwrap_err();
assert!(err.to_string().contains("duplicate schema"), "{err}");
}
#[test]
fn validate_rejects_duplicate_field_names() {
let shape = minimal_shape(
vec![
ShapeField {
name: "x".into(),
field_type: FieldType::Varint as u8,
repeat_meta: None,
},
ShapeField {
name: "x".into(),
field_type: FieldType::Varint as u8,
repeat_meta: None,
},
],
vec![ShapeValue::U(1), ShapeValue::U(2)],
);
let err = validate_shape(&shape).unwrap_err();
assert!(err.to_string().contains("duplicate field name"), "{err}");
}
#[test]
fn validate_identity_requires_namespace() {
let shape = minimal_shape(
vec![ShapeField {
name: "id".into(),
field_type: FieldType::Varint as u8,
repeat_meta: Some(FieldRepeatMeta {
semantics: FieldSemantics::Identity,
namespace: None, }),
}],
vec![ShapeValue::U(1)],
);
let err = validate_shape(&shape).unwrap_err();
assert!(err.to_string().contains("requires a namespace"), "{err}");
}
#[test]
fn validate_identity_stack_requires_addr_namespace() {
let shape = minimal_shape(
vec![ShapeField {
name: "stack".into(),
field_type: FieldType::StackFrames as u8,
repeat_meta: Some(FieldRepeatMeta {
semantics: FieldSemantics::Identity,
namespace: Some(NamespaceId::Task), }),
}],
vec![ShapeValue::Stack(vec![0x1000])],
);
let err = validate_shape(&shape).unwrap_err();
assert!(err.to_string().contains("Addr namespace"), "{err}");
}
#[test]
fn validate_timestamp_ref_forbids_namespace() {
let shape = minimal_shape(
vec![ShapeField {
name: "ts".into(),
field_type: FieldType::Varint as u8,
repeat_meta: Some(FieldRepeatMeta {
semantics: FieldSemantics::TimestampRef,
namespace: Some(NamespaceId::Addr),
}),
}],
vec![ShapeValue::U(100)],
);
let err = validate_shape(&shape).unwrap_err();
assert!(err.to_string().contains("forbids namespace"), "{err}");
}
#[test]
fn validate_timestamp_ref_requires_varint() {
let shape = minimal_shape(
vec![ShapeField {
name: "ts".into(),
field_type: FieldType::String as u8,
repeat_meta: Some(FieldRepeatMeta {
semantics: FieldSemantics::TimestampRef,
namespace: None,
}),
}],
vec![ShapeValue::S("hello".into())],
);
let err = validate_shape(&shape).unwrap_err();
assert!(err.to_string().contains("requires Varint"), "{err}");
}
#[test]
fn validate_realtime_offset_forbids_namespace() {
let shape = minimal_shape(
vec![ShapeField {
name: "rt".into(),
field_type: FieldType::Varint as u8,
repeat_meta: Some(FieldRepeatMeta {
semantics: FieldSemantics::RealtimeOffset,
namespace: Some(NamespaceId::Task),
}),
}],
vec![ShapeValue::U(100)],
);
let err = validate_shape(&shape).unwrap_err();
assert!(err.to_string().contains("forbids namespace"), "{err}");
}
#[test]
fn validate_structural_forbids_namespace() {
let shape = minimal_shape(
vec![ShapeField {
name: "count".into(),
field_type: FieldType::Varint as u8,
repeat_meta: Some(FieldRepeatMeta {
semantics: FieldSemantics::Structural,
namespace: Some(NamespaceId::Addr),
}),
}],
vec![ShapeValue::U(42)],
);
let err = validate_shape(&shape).unwrap_err();
assert!(
err.to_string()
.contains("Structural semantics forbids namespace"),
"{err}"
);
}
#[test]
fn validate_rejects_duplicate_annotations() {
let mut shape = minimal_shape(
vec![ShapeField {
name: "x".into(),
field_type: FieldType::Varint as u8,
repeat_meta: None,
}],
vec![ShapeValue::U(1)],
);
shape.schemas[0].annotations = vec![
ShapeAnnotation {
field_index: 0,
key: "metrique.unit".into(),
value: "ns".into(),
},
ShapeAnnotation {
field_index: 0,
key: "metrique.unit".into(),
value: "us".into(),
},
];
let err = validate_shape(&shape).unwrap_err();
assert!(err.to_string().contains("duplicate annotation"), "{err}");
}
#[test]
fn validate_rejects_annotation_value_for_different_key() {
let mut shape = minimal_shape(
vec![ShapeField {
name: "x".into(),
field_type: FieldType::Varint as u8,
repeat_meta: None,
}],
vec![ShapeValue::U(1)],
);
shape.schemas[0].annotations = vec![ShapeAnnotation {
field_index: 0,
key: "kind".into(),
value: "ns".into(),
}];
let err = validate_shape(&shape).unwrap_err();
assert!(
err.to_string()
.contains("annotation 'kind' has unsafe value 'ns'"),
"{err}"
);
}
#[test]
fn validate_summary_type_counts_both_directions() {
let mut shape = minimal_shape(
vec![ShapeField {
name: "x".into(),
field_type: FieldType::Varint as u8,
repeat_meta: None,
}],
vec![ShapeValue::U(1)],
);
shape.summary.event_type_counts.insert("Phantom".into(), 5);
let err = validate_shape(&shape).unwrap_err();
assert!(err.to_string().contains("Phantom"), "{err}");
let mut shape2 = minimal_shape(
vec![ShapeField {
name: "x".into(),
field_type: FieldType::Varint as u8,
repeat_meta: None,
}],
vec![ShapeValue::U(1)],
);
shape2.summary.event_type_counts.clear(); let err2 = validate_shape(&shape2).unwrap_err();
assert!(err2.to_string().contains("does not list"), "{err2}");
}
#[test]
fn validate_summary_duration_mismatch() {
let mut shape = minimal_shape(
vec![ShapeField {
name: "x".into(),
field_type: FieldType::Varint as u8,
repeat_meta: None,
}],
vec![ShapeValue::U(1)],
);
shape.summary.duration_ns = 99999; let err = validate_shape(&shape).unwrap_err();
assert!(err.to_string().contains("duration_ns"), "{err}");
}
#[test]
fn validate_rejects_oversized_bytes() {
let shape = minimal_shape(
vec![ShapeField {
name: "data".into(),
field_type: FieldType::Bytes as u8,
repeat_meta: None,
}],
vec![ShapeValue::Bytes(MAX_BYTES_LENGTH + 1)],
);
let err = validate_shape(&shape).unwrap_err();
assert!(err.to_string().contains("exceeds limit"), "{err}");
}
#[test]
fn validate_rejects_oversized_stack() {
let ok_stack = ShapeValue::Stack(vec![0x1000, 0x2000]);
assert!(
validate_value_compat(&ok_stack, FieldType::StackFrames as u8, "T", "s", 0, 0).is_ok()
);
let small_over = vec![0u64; 1025]; assert!(
validate_value_compat(
&ShapeValue::Stack(small_over),
FieldType::StackFrames as u8,
"T",
"s",
0,
0
)
.is_ok()
);
let mut nested = ShapeValue::U(1);
for _ in 0..9 {
nested = ShapeValue::List(vec![nested]);
}
let err = validate_value_compat(&nested, FieldType::DynamicList as u8, "T", "items", 0, 0)
.unwrap_err();
assert!(
err.to_string().contains("nesting depth"),
"expected depth error: {err}"
);
}
#[test]
fn validate_rejects_too_many_events_for_repeat() {
let trace = make_test_trace();
let shape = extract_shape(&trace).unwrap();
let max_repeat = (MAX_GENERATED_EVENTS / shape.events.len() as u64) + 1;
let result = validate_repeat_preflight(&shape, max_repeat as u32);
assert!(result.is_err());
}
#[test]
fn no_output_file_on_validation_failure() {
let tmp = tempfile::tempdir().unwrap();
let shape_path = tmp.path().join("bad_shape.json");
let output_path = tmp.path().join("output.bin");
let bad_shape = TraceShape {
version: 999,
summary: ShapeSummary {
event_count: 1,
duration_ns: 0,
event_type_counts: [("T".into(), 1)].into_iter().collect(),
worker_cardinality: 0,
task_cardinality: 0,
poll_duration_quantiles: None,
span_duration_quantiles: None,
},
schemas: vec![ShapeSchema {
name: "T".into(),
has_timestamp: true,
fields: vec![],
annotations: vec![],
}],
events: vec![ShapeEvent {
schema_index: 0,
timestamp_offset_ns: Some(0),
values: vec![],
}],
};
let json = serde_json::to_string_pretty(&bad_shape).unwrap();
std::fs::write(&shape_path, &json).unwrap();
let result = generate(&shape_path, &output_path, 1);
assert!(result.is_err());
assert!(
!output_path.exists(),
"output file created despite validation failure"
);
}
#[test]
fn huge_bytes_rejected_without_allocation() {
let shape = minimal_shape(
vec![ShapeField {
name: "data".into(),
field_type: FieldType::Bytes as u8,
repeat_meta: None,
}],
vec![ShapeValue::Bytes(MAX_BYTES_LENGTH + 1)],
);
let err = validate_shape(&shape).unwrap_err();
assert!(err.to_string().contains("exceeds limit"), "{err}");
}
#[test]
fn streaming_output_decodes() {
let trace = make_test_trace();
let shape = extract_shape(&trace).unwrap();
let mut buf = Vec::new();
generate_to_writer(&shape, 1, &mut buf).unwrap();
let mut dec = Decoder::new(&buf).unwrap();
let mut count = 0u64;
loop {
match dec.next_frame().unwrap() {
Some(DecodedFrame::Event { .. }) => count += 1,
Some(_) => {}
None => break,
}
}
assert_eq!(count, shape.summary.event_count);
}
#[test]
fn worker_sentinel_excluded_from_cardinality() {
let mut enc = Encoder::new();
let schema = register_poll_start_schema(&mut enc);
for &wid in &[0u64, 1, 254, 255] {
write_poll_start(&mut enc, &schema, BASE_TS + wid * 10_000, wid, wid + 1);
}
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
assert_eq!(shape.summary.worker_cardinality, 2);
}
#[test]
fn task_id_zero_excluded_from_cardinality() {
let mut enc = Encoder::new();
let schema = register_task_spawn_schema(&mut enc);
for &tid in &[0u64, 1, 2] {
write_task_spawn(&mut enc, &schema, BASE_TS + tid * 10_000, tid);
}
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
assert_eq!(shape.summary.task_cardinality, 2);
}
#[test]
fn repeated_poll_start_end_pairing() {
let mut enc = Encoder::new();
let start = register_poll_start_schema(&mut enc);
let end = enc
.register_schema(
"PollEndEvent",
vec![FieldDef::new("worker_id", FieldType::Varint)],
)
.unwrap();
write_poll_start(&mut enc, &start, BASE_TS, 0, 1);
enc.write_event(
&end,
&[FieldValue::Varint(BASE_TS + 100_000), FieldValue::Varint(0)],
)
.unwrap();
write_poll_start(&mut enc, &start, BASE_TS + 200_000, 0, 1);
enc.write_event(
&end,
&[FieldValue::Varint(BASE_TS + 400_000), FieldValue::Varint(0)],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let q = shape.summary.poll_duration_quantiles.as_ref().unwrap();
assert_eq!(q.count, 2);
assert!(q.min_ns <= q.max_ns);
}
#[test]
fn repeated_span_enter_exit_pairing() {
let mut enc = Encoder::new();
let enter = enc
.register_schema(
"SpanEnter:handler",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
],
)
.unwrap();
let exit = enc
.register_schema(
"SpanExit:handler",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
],
)
.unwrap();
enc.write_event(
&enter,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(0),
FieldValue::Varint(42),
],
)
.unwrap();
enc.write_event(
&exit,
&[
FieldValue::Varint(BASE_TS + 100_000),
FieldValue::Varint(0),
FieldValue::Varint(42),
],
)
.unwrap();
enc.write_event(
&enter,
&[
FieldValue::Varint(BASE_TS + 200_000),
FieldValue::Varint(0),
FieldValue::Varint(42),
],
)
.unwrap();
enc.write_event(
&exit,
&[
FieldValue::Varint(BASE_TS + 500_000),
FieldValue::Varint(0),
FieldValue::Varint(42),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let q = shape.summary.span_duration_quantiles.as_ref().unwrap();
assert_eq!(q.count, 2);
}
#[test]
fn summary_cardinality_recomputed_correctly() {
let trace = make_test_trace();
let mut shape = extract_shape(&trace).unwrap();
let original_worker_card = shape.summary.worker_cardinality;
shape.summary.worker_cardinality = original_worker_card + 99;
let err = validate_shape(&shape).unwrap_err();
assert!(err.to_string().contains("worker_cardinality"), "{err}");
}
#[test]
fn summary_quantiles_mismatch_detected() {
let trace = make_test_trace();
let mut shape = extract_shape(&trace).unwrap();
if let Some(q) = shape.summary.poll_duration_quantiles.as_mut() {
q.count += 1;
}
let err = validate_shape(&shape).unwrap_err();
assert!(err.to_string().contains("poll_duration"), "{err}");
}
#[test]
fn target_worker_not_counted_in_cardinality() {
let mut enc = Encoder::new();
let wake = enc
.register_schema(
"WakeEventEvent",
vec![
FieldDef::new("waker_task_id", FieldType::Varint),
FieldDef::new("woken_task_id", FieldType::Varint),
FieldDef::new("target_worker", FieldType::Varint),
],
)
.unwrap();
enc.write_event(
&wake,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(1),
FieldValue::Varint(2),
FieldValue::Varint(5),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
assert_eq!(
shape.summary.worker_cardinality, 0,
"target_worker should not be counted"
);
}
#[test]
fn reset_same_name_different_schema_errors() {
let mut enc1 = Encoder::new();
let schema1 = enc1
.register_schema(
"SameCustom",
vec![FieldDef::new("field_a", FieldType::String)],
)
.unwrap();
enc1.write_event(
&schema1,
&[
FieldValue::Varint(BASE_TS),
FieldValue::String("hello".into()),
],
)
.unwrap();
let part1 = enc1.finish();
let mut enc2 = Encoder::new();
let schema2 = enc2
.register_schema(
"SameCustom",
vec![FieldDef::new("field_a", FieldType::Varint)],
)
.unwrap();
enc2.write_event(
&schema2,
&[
FieldValue::Varint(BASE_TS + 100_000),
FieldValue::Varint(42),
],
)
.unwrap();
let part2 = enc2.finish();
let mut combined = part1;
combined.extend_from_slice(&part2);
let result = extract_shape(&combined);
assert!(
result.is_err(),
"extraction should error on conflicting schema after reset"
);
let err_msg = format!("{:#}", result.unwrap_err());
assert!(
err_msg.contains("different definition")
|| err_msg.contains("re-registered after reset")
|| err_msg.contains("does not exactly match a known complete signature"),
"unexpected error: {err_msg}"
);
}
#[test]
fn reset_identical_schema_works() {
let mut enc1 = Encoder::new();
let schema1 = enc1
.register_schema(
"SameCustom",
vec![FieldDef::new("name", FieldType::PooledString)],
)
.unwrap();
let ps1 = enc1.intern_string("pooled_val").unwrap();
enc1.write_event(
&schema1,
&[FieldValue::Varint(BASE_TS), FieldValue::PooledString(ps1)],
)
.unwrap();
let part1 = enc1.finish();
let mut enc2 = Encoder::new();
let schema2 = enc2
.register_schema(
"SameCustom",
vec![FieldDef::new("name", FieldType::PooledString)],
)
.unwrap();
let ps2 = enc2.intern_string("another_val").unwrap();
enc2.write_event(
&schema2,
&[
FieldValue::Varint(BASE_TS + 100_000),
FieldValue::PooledString(ps2),
],
)
.unwrap();
let part2 = enc2.finish();
let mut combined = part1;
combined.extend_from_slice(&part2);
let shape = extract_shape(&combined).expect("identical schema after reset should work");
assert_eq!(shape.events.len(), 2);
}
#[test]
fn custom_schema_original_worker_id_roundtrip() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"MyCustomEvent",
vec![FieldDef::new("worker_id", FieldType::Varint)],
)
.unwrap();
enc.write_event(
&schema,
&[FieldValue::Varint(BASE_TS), FieldValue::Varint(3)],
)
.unwrap();
enc.write_event(
&schema,
&[FieldValue::Varint(BASE_TS + 10_000), FieldValue::Varint(5)],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let schema_fields: Vec<&str> = shape.schemas[0]
.fields
.iter()
.map(|f| f.name.as_str())
.collect();
assert!(
schema_fields[0].starts_with("field_"),
"custom schema worker_id should be anonymized"
);
assert_eq!(shape.summary.worker_cardinality, 0);
let generated = generate_trace(&shape, 1).unwrap();
let re_shape = extract_shape(&generated).unwrap();
assert_eq!(
re_shape.summary.worker_cardinality,
shape.summary.worker_cardinality
);
}
#[test]
fn same_timestamp_poll_start_end() {
let mut enc = Encoder::new();
let start = register_poll_start_schema(&mut enc);
let end = enc
.register_schema(
"PollEndEvent",
vec![FieldDef::new("worker_id", FieldType::Varint)],
)
.unwrap();
write_poll_start(&mut enc, &start, BASE_TS, 0, 1);
enc.write_event(&end, &[FieldValue::Varint(BASE_TS), FieldValue::Varint(0)])
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let q = shape.summary.poll_duration_quantiles.as_ref().unwrap();
assert_eq!(q.count, 1);
assert_eq!(q.min_ns, 0);
assert_eq!(q.max_ns, 0);
}
#[test]
fn same_timestamp_span_enter_exit() {
let mut enc = Encoder::new();
let enter = enc
.register_schema(
"SpanEnter:my_handler",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
],
)
.unwrap();
let exit = enc
.register_schema(
"SpanExit:my_handler",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("span_id", FieldType::Varint),
],
)
.unwrap();
enc.write_event(
&enter,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(0),
FieldValue::Varint(42),
],
)
.unwrap();
enc.write_event(
&exit,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(0),
FieldValue::Varint(42),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let q = shape.summary.span_duration_quantiles.as_ref().unwrap();
assert_eq!(q.count, 1);
assert_eq!(q.min_ns, 0);
}
#[test]
fn extra_zero_count_key_rejected() {
let mut shape = minimal_shape(
vec![ShapeField {
name: "x".into(),
field_type: FieldType::Varint as u8,
repeat_meta: None,
}],
vec![ShapeValue::U(1)],
);
shape.summary.event_type_counts.insert("Phantom".into(), 0);
let err = validate_shape(&shape).unwrap_err();
assert!(
err.to_string().contains("zero-count"),
"expected zero-count rejection: {err}"
);
}
#[test]
fn clock_sync_extra_monotonic_ns_rejected() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"ClockSyncEvent",
vec![
FieldDef::new("realtime_ns", FieldType::Varint),
FieldDef::new("monotonic_ns", FieldType::Varint),
],
)
.unwrap();
let source_realtime = 1_700_000_000_000_000_000u64;
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(source_realtime),
FieldValue::Varint(BASE_TS),
],
)
.unwrap();
let trace = enc.finish();
let result = extract_shape(&trace);
assert!(
result.is_err(),
"ClockSyncEvent with extra monotonic_ns should be rejected"
);
let err_msg = format!("{:#}", result.unwrap_err());
assert!(
err_msg.contains("does not exactly match a known complete signature")
&& err_msg.contains("monotonic_ns"),
"unexpected error: {err_msg}"
);
}
#[test]
fn clock_sync_realtime_only_works() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"ClockSyncEvent",
vec![FieldDef::new("realtime_ns", FieldType::Varint)],
)
.unwrap();
let source_realtime = 1_700_000_000_000_000_000u64;
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(source_realtime),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let json = serde_json::to_string(&shape).unwrap();
assert!(
!json.contains("1700000000000000000"),
"source realtime leaked"
);
}
#[test]
fn generate_flush_error_propagated() {
struct FailFlush {
buf: Vec<u8>,
}
impl Write for FailFlush {
fn write(&mut self, data: &[u8]) -> std::io::Result<usize> {
self.buf.extend_from_slice(data);
Ok(data.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Err(std::io::Error::other("flush failed intentionally"))
}
}
let trace = make_test_trace();
let shape = extract_shape(&trace).unwrap();
let writer = FailFlush { buf: Vec::new() };
let result = generate_to_writer(&shape, 1, writer);
assert!(result.is_err(), "should propagate flush error");
let err_msg = result.unwrap_err().to_string();
assert!(
err_msg.contains("flush"),
"error should mention flush: {err_msg}"
);
}
#[cfg(target_os = "linux")]
#[test]
fn generate_dev_full_returns_error() {
use std::fs::OpenOptions;
let trace = make_test_trace();
let shape = extract_shape(&trace).unwrap();
let file = OpenOptions::new()
.write(true)
.open("/dev/full")
.expect("/dev/full should be openable");
let writer = BufWriter::new(file);
let result = generate_to_writer(&shape, 1, writer);
assert!(result.is_err(), "/dev/full should cause write/flush error");
}
#[test]
fn nested_address_repeat_disjoint() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"TestCpuSampleEvent",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("callchain", FieldType::StackFrames),
],
)
.unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(0),
FieldValue::StackFrames(vec![0x1000u64, 0x2000].into()),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let generated = generate_trace(&shape, 2).unwrap();
let mut dec = Decoder::new(&generated).unwrap();
let mut all_addrs: Vec<Vec<u64>> = Vec::new();
loop {
match dec.next_frame().unwrap() {
Some(DecodedFrame::Event { values, .. }) => {
if let Some(FieldValue::StackFrames(frames)) = values.get(1) {
all_addrs.push(frames.iter().copied().collect());
}
}
Some(_) => {}
None => break,
}
}
assert_eq!(all_addrs.len(), 2);
let set0: HashSet<u64> = all_addrs[0].iter().copied().collect();
let set1: HashSet<u64> = all_addrs[1].iter().copied().collect();
assert!(
set0.is_disjoint(&set1),
"stack addresses should be disjoint across reps: {:?} vs {:?}",
set0,
set1
);
}
#[test]
fn nested_dynamic_stack_addresses_are_disjoint_across_repeats() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"NestedStacks",
vec![
FieldDef::new("list", FieldType::DynamicList),
FieldDef::new("map", FieldType::DynamicMap),
],
)
.unwrap();
let pooled = enc.intern_stack_frames(&[0x1010, 0x1020]).unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::List(vec![FieldValue::StackFrames(vec![0x1000, 0x1010].into())]),
FieldValue::Map(vec![(
FieldValue::String("stack".into()),
FieldValue::PooledStackFrames(pooled),
)]),
],
)
.unwrap();
let shape = extract_shape(&enc.finish()).unwrap();
let generated = generate_trace(&shape, 2).unwrap();
let mut dec = Decoder::new(&generated).unwrap();
let mut addresses_by_event = Vec::new();
loop {
match dec.next_frame().unwrap() {
Some(DecodedFrame::Event { values, .. }) => {
let mut addresses = HashSet::new();
if let Some(FieldValue::List(items)) = values.first() {
for item in items {
if let FieldValue::StackFrames(frames) = item {
addresses.extend(frames.iter().copied());
}
}
}
if let Some(FieldValue::Map(entries)) = values.get(1) {
for (_, value) in entries {
if let FieldValue::PooledStackFrames(id) = value {
addresses
.extend(dec.stack_pool().get(*id).unwrap().iter().copied());
}
}
}
addresses_by_event.push(addresses);
}
Some(_) => {}
None => break,
}
}
assert_eq!(addresses_by_event.len(), 2);
assert!(
addresses_by_event[0].is_disjoint(&addresses_by_event[1]),
"nested stack identities overlap: {:?} vs {:?}",
addresses_by_event[0],
addresses_by_event[1]
);
}
#[test]
fn overflow_shape_no_output_file_created() {
let tmp = tempfile::tempdir().unwrap();
let shape_path = tmp.path().join("overflow_shape.json");
let output_path = tmp.path().join("output.bin");
let mut shape = minimal_shape(
vec![ShapeField {
name: "ts_ref".into(),
field_type: FieldType::Varint as u8,
repeat_meta: Some(FieldRepeatMeta {
semantics: FieldSemantics::TimestampRef,
namespace: None,
}),
}],
vec![ShapeValue::U(u64::MAX - 100)],
);
shape.events[0].timestamp_offset_ns = Some(u64::MAX - 100);
shape.summary.duration_ns = u64::MAX - 100;
let json = serde_json::to_string_pretty(&shape).unwrap();
std::fs::write(&shape_path, &json).unwrap();
let result = generate(&shape_path, &output_path, 2);
assert!(result.is_err(), "should fail due to overflow in preflight");
assert!(
!output_path.exists(),
"output file must not be created on preflight overflow"
);
}
#[test]
fn max_shape_json_bytes_is_bounded() {
assert_eq!(MAX_SHAPE_JSON_BYTES, 32 * 1024 * 1024);
}
#[test]
fn extraction_bounds_nested_stringmap() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"TestEvent",
vec![FieldDef::new("data", FieldType::DynamicList)],
)
.unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::List(vec![FieldValue::StringMap(vec![
(b"key1".to_vec(), b"val1".to_vec()),
(b"key2".to_vec(), b"val2".to_vec()),
])]),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
match &shape.events[0].values[0] {
ShapeValue::List(items) => match &items[0] {
ShapeValue::StringMap(pairs) => assert_eq!(pairs.len(), 2),
other => panic!("expected StringMap, got {other:?}"),
},
other => panic!("expected List, got {other:?}"),
}
}
#[test]
fn anon_namespace_checked_allocator() {
let mut ctx = ExtractContext::new();
let ns0 = ctx.get_anon_namespace("field_a_id").unwrap();
let ns1 = ctx.get_anon_namespace("field_b_id").unwrap();
assert_eq!(ns0, NamespaceId::Anon(0));
assert_eq!(ns1, NamespaceId::Anon(1));
let ns0_again = ctx.get_anon_namespace("field_a_id").unwrap();
assert_eq!(ns0_again, NamespaceId::Anon(0));
}
#[test]
fn anon_namespace_near_limit() {
let mut ctx = ExtractContext::new();
ctx.next_anon_namespace = u16::MAX - 100 - 1;
let ns = ctx.get_anon_namespace("almost_full_id");
assert!(ns.is_ok());
ctx.next_anon_namespace = u16::MAX - 100;
let ns_err = ctx.get_anon_namespace("overflow_id");
assert!(ns_err.is_err(), "should error at max anonymous namespaces");
assert!(
ns_err
.unwrap_err()
.to_string()
.contains("too many anonymous namespaces")
);
}
#[test]
fn namespace_id_reserved_range_rejected() {
for n in [6u16, 50, 99] {
let json = format!("{n}");
let result: Result<NamespaceId, _> = serde_json::from_str(&json);
assert!(
result.is_err(),
"namespace id {n} in reserved range should be rejected"
);
}
for n in [0u16, 1, 2, 3, 4, 5, 100, 200, 65535] {
let json = format!("{n}");
let result: Result<NamespaceId, _> = serde_json::from_str(&json);
assert!(result.is_ok(), "namespace id {n} should be valid");
}
}
#[test]
fn extract_generate_re_extract_demo() {
let demo_path =
std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("ui/public/demo-trace.bin");
if !demo_path.exists() {
eprintln!("skipping: demo-trace.bin not found");
return;
}
let data = read_trace_file(&demo_path).unwrap();
let shape = extract_shape(&data).expect("demo extract must succeed");
let generated = generate_trace(&shape, 1).expect("generate from demo shape");
let re_shape = extract_shape(&generated).expect("re-extract from generated trace");
assert_eq!(re_shape.summary.event_count, shape.summary.event_count);
assert_eq!(re_shape.schemas.len(), shape.schemas.len());
}
#[test]
fn builtin_schema_wrong_field_type_rejected() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"PollStartEvent",
vec![
FieldDef::new("worker_id", FieldType::String),
FieldDef::new("local_queue", FieldType::U8),
FieldDef::new("task_id", FieldType::Varint),
FieldDef::new("spawn_loc", FieldType::PooledString),
],
)
.unwrap();
let spawn_loc = enc.intern_string("source.rs").unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::String("not-a-worker".into()),
FieldValue::Varint(1),
FieldValue::Varint(2),
FieldValue::PooledString(spawn_loc),
],
)
.unwrap();
let error = extract_shape(&enc.finish()).unwrap_err();
assert!(format!("{error:#}").contains("signature"), "{error:#}");
}
#[test]
fn builtin_schema_mixed_raw_type_signature_rejected() {
let fields = [
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("local_queue", FieldType::Varint),
FieldDef::new("task_id", FieldType::U32),
FieldDef::new("spawn_loc", FieldType::PooledString),
];
let error = validate_builtin_schema_signature("PollStartEvent", &fields).unwrap_err();
assert!(error.to_string().contains("signature"), "{error:#}");
}
#[test]
fn builtin_schema_known_historical_variant_accepted() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"WorkerUnparkEvent",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("local_queue", FieldType::U8),
FieldDef::new("cpu_time_ns", FieldType::Varint),
FieldDef::new("sched_wait_ns", FieldType::Varint),
],
)
.unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(0),
FieldValue::Varint(100),
FieldValue::Varint(500_000),
FieldValue::Varint(10_000),
],
)
.unwrap();
assert!(extract_shape(&enc.finish()).is_ok());
}
#[test]
fn builtin_historical_signatures_are_complete_and_accepted() {
let cases = [
(
"WorkerParkEvent",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("local_queue", FieldType::U8),
FieldDef::new("cpu_time_ns", FieldType::Varint),
],
),
(
"WorkerUnparkEvent",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("local_queue", FieldType::U8),
FieldDef::new("cpu_time_ns", FieldType::Varint),
FieldDef::new("sched_wait_ns", FieldType::Varint),
],
),
(
"TaskSpawnEvent",
vec![
FieldDef::new("task_id", FieldType::U32),
FieldDef::new("spawn_loc", FieldType::PooledString),
],
),
(
"CpuSampleEvent",
vec![
FieldDef::new("worker_id", FieldType::Varint),
FieldDef::new("tid", FieldType::U32),
FieldDef::new("source", FieldType::U8),
FieldDef::new("thread_name", FieldType::PooledString),
FieldDef::new("callchain", FieldType::StackFrames),
],
),
(
"SymbolTableEntry",
vec![
FieldDef::new("base_addr", FieldType::Varint),
FieldDef::new("size", FieldType::Varint),
FieldDef::new("symbol_name", FieldType::PooledString),
],
),
];
for (name, fields) in cases {
validate_builtin_schema_signature(name, &fields).unwrap();
}
}
#[test]
fn gzip_compressed_input_limit_is_enforced() {
let trace = make_test_trace();
let mut encoder = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::default());
encoder.write_all(&trace).unwrap();
let compressed = encoder.finish().unwrap();
let limit = compressed.len() as u64 - 1;
let error = read_gzip_bounded(compressed.as_slice(), limit, 1024 * 1024).unwrap_err();
assert!(
format!("{error:#}").contains("compressed trace exceeds maximum size"),
"{error:#}"
);
}
#[test]
fn gzip_decompressed_output_limit_is_enforced() {
let trace = make_test_trace();
let mut encoder = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::default());
encoder.write_all(&trace).unwrap();
let compressed = encoder.finish().unwrap();
let error = read_gzip_bounded(
compressed.as_slice(),
compressed.len() as u64,
trace.len() as u64 - 1,
)
.unwrap_err();
assert!(
format!("{error:#}").contains("decompressed trace exceeds maximum size"),
"{error:#}"
);
}
#[test]
fn deeply_nested_dynamic_list_returns_none() {
let mut data = Vec::new();
let depth = 40; for _ in 0..depth {
data.extend_from_slice(&1u32.to_le_bytes()); data.push(FieldType::DynamicList as u8); }
data.extend_from_slice(&0u32.to_le_bytes());
let result =
dial9_trace_format::types::FieldValueRef::decode(FieldType::DynamicList, &data, 0);
assert!(
result.is_none(),
"deeply nested dynamic list should return None (recursion budget exceeded)"
);
}
#[test]
fn deeply_nested_dynamic_map_returns_none() {
let mut data = Vec::new();
let depth = 40;
for _ in 0..depth {
data.extend_from_slice(&1u32.to_le_bytes()); data.push(FieldType::Varint as u8); data.push(0x01); data.push(FieldType::DynamicMap as u8); }
data.extend_from_slice(&0u32.to_le_bytes());
let result =
dial9_trace_format::types::FieldValueRef::decode(FieldType::DynamicMap, &data, 0);
assert!(
result.is_none(),
"deeply nested dynamic map should return None"
);
}
#[test]
fn stringmap_malformed_klen_overflow_returns_none() {
let mut data = Vec::new();
data.extend_from_slice(&1u32.to_le_bytes()); data.extend_from_slice(&u32::MAX.to_le_bytes());
let result =
dial9_trace_format::types::FieldValueRef::decode(FieldType::StringMap, &data, 0);
assert!(
result.is_none(),
"StringMap with overflow klen should return None"
);
}
#[test]
fn estimator_accounts_for_dynamic_type_tags() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"TestEvent",
vec![FieldDef::new("items", FieldType::DynamicList)],
)
.unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::List(vec![
FieldValue::Varint(1),
FieldValue::Varint(2),
FieldValue::String("hello".into()),
]),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let estimated = estimate_encoded_size(&shape, 1).unwrap();
let actual = generate_trace(&shape, 1).unwrap().len() as u64;
assert!(
estimated >= actual,
"estimator ({estimated}) should be >= actual ({actual})"
);
}
#[test]
fn estimator_accounts_for_optional_prefixes_and_annotations() {
let fields: Vec<ShapeField> = (0..100)
.map(|i| ShapeField {
name: format!("optional_bool_{i}"),
field_type: FieldType::OptionalBool as u8,
repeat_meta: None,
})
.collect();
let mut shape = minimal_shape(fields, vec![ShapeValue::B(true); 100]);
shape.schemas[0].annotations = (0..100)
.map(|i| ShapeAnnotation {
field_index: i,
key: "unit".into(),
value: "ns".into(),
})
.collect();
validate_shape(&shape).unwrap();
let estimated = estimate_encoded_size(&shape, 1).unwrap();
let actual = generate_trace(&shape, 1).unwrap().len() as u64;
assert!(
estimated >= actual,
"estimator ({estimated}) should cover optional prefixes and annotation framing ({actual})"
);
}
#[test]
fn repeated_pooled_stack_retention_rejected_before_output_creation() {
let shape = minimal_shape(
vec![ShapeField {
name: "callchain".into(),
field_type: FieldType::PooledStackFrames as u8,
repeat_meta: Some(FieldRepeatMeta {
semantics: FieldSemantics::Identity,
namespace: Some(NamespaceId::Addr),
}),
}],
vec![ShapeValue::PStack(vec![1; 1024])],
);
let tmp = tempfile::tempdir().unwrap();
let output = tmp.path().join("output.bin");
let error = write_generated_trace(&shape, &output, 70_000).unwrap_err();
assert!(
error.to_string().contains("retained generation memory"),
"unexpected error: {error:#}"
);
assert!(
!output.exists(),
"preflight rejection must not create output"
);
}
#[test]
fn quantize_f64_subnormal_stays_finite_and_json_roundtrips() {
for value in [f64::from_bits(1), -f64::from_bits(1), f64::MIN_POSITIVE] {
let rounded = quantize_f64(value).unwrap();
assert!(rounded.is_finite());
let encoded = serde_json::to_string(&ShapeValue::F(rounded)).unwrap();
let decoded: ShapeValue = serde_json::from_str(&encoded).unwrap();
assert!(matches!(decoded, ShapeValue::F(v) if v.is_finite()));
}
}
#[test]
fn per_event_non_bytes_working_set_rejected_before_output_creation() {
let shape = minimal_shape(
vec![ShapeField {
name: "stacks".into(),
field_type: FieldType::DynamicList as u8,
repeat_meta: None,
}],
vec![ShapeValue::List(
(0..9)
.map(|_| ShapeValue::Stack(vec![1; 1_000_000]))
.collect(),
)],
);
let tmp = tempfile::tempdir().unwrap();
let output = tmp.path().join("output.bin");
let error = write_generated_trace(&shape, &output, 1).unwrap_err();
assert!(error.to_string().contains("working set"), "{error:#}");
assert!(!output.exists());
}
#[test]
fn per_event_bytes_working_set_bounded() {
let shape = minimal_shape(
vec![
ShapeField {
name: "data1".into(),
field_type: FieldType::Bytes as u8,
repeat_meta: None,
},
ShapeField {
name: "data2".into(),
field_type: FieldType::Bytes as u8,
repeat_meta: None,
},
],
vec![ShapeValue::Bytes(1024), ShapeValue::Bytes(2048)],
);
assert!(validate_repeat_preflight(&shape, 1).is_ok());
}
#[test]
fn per_event_bytes_working_set_rejected_when_excessive() {
let excessive_bytes = MAX_BYTES_LENGTH; let shape = minimal_shape(
vec![
ShapeField {
name: "data1".into(),
field_type: FieldType::Bytes as u8,
repeat_meta: None,
},
ShapeField {
name: "data2".into(),
field_type: FieldType::Bytes as u8,
repeat_meta: None,
},
],
vec![
ShapeValue::Bytes(excessive_bytes),
ShapeValue::Bytes(excessive_bytes),
],
);
let err = validate_repeat_preflight(&shape, 1).unwrap_err();
let msg = err.to_string();
assert!(
msg.contains("working set") || msg.contains("per-event limit"),
"rejection message should mention working set limit: {msg}"
);
}
#[test]
fn safe_annotation_keys_and_values_preserved() {
use dial9_trace_format::schema::{FieldAnnotation, SchemaEntry};
let mut enc = Encoder::new();
let fields = vec![
FieldDef::new("latency_ns", FieldType::Varint),
FieldDef::new("duration_ms", FieldType::Varint),
];
let annotations = vec![
FieldAnnotation::new(0, "unit", "ns"),
FieldAnnotation::new(1, "metrique.unit", "ms"),
FieldAnnotation::new(0, "kind", "counter"),
FieldAnnotation::new(1, "kind", "histogram"),
];
let entry = SchemaEntry::with_annotations("TestEvent", true, fields, annotations);
let schema = dial9_trace_format::encoder::Schema::from_entry(entry);
enc.register_existing(&schema).unwrap();
enc.write_event(
&schema,
&[
FieldValue::Varint(BASE_TS),
FieldValue::Varint(1000),
FieldValue::Varint(5),
],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let schema_shape = &shape.schemas[0];
assert_eq!(schema_shape.annotations.len(), 3);
let ann_keys: Vec<&str> = schema_shape
.annotations
.iter()
.map(|a| a.key.as_str())
.collect();
assert!(
ann_keys.contains(&"unit"),
"derive-style 'unit' annotation not preserved"
);
assert!(
ann_keys.contains(&"metrique.unit"),
"legacy 'metrique.unit' annotation not preserved"
);
assert!(
ann_keys.contains(&"kind"),
"derive-style 'kind' annotation not preserved"
);
}
#[test]
fn annotation_unit_s_preserved() {
use dial9_trace_format::schema::{FieldAnnotation, SchemaEntry};
let mut enc = Encoder::new();
let fields = vec![FieldDef::new("timeout_s", FieldType::Varint)];
let annotations = vec![FieldAnnotation::new(0, "unit", "s")];
let entry = SchemaEntry::with_annotations("TestEvent", true, fields, annotations);
let schema = dial9_trace_format::encoder::Schema::from_entry(entry);
enc.register_existing(&schema).unwrap();
enc.write_event(
&schema,
&[FieldValue::Varint(BASE_TS), FieldValue::Varint(30)],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
assert_eq!(shape.schemas[0].annotations.len(), 1);
assert_eq!(shape.schemas[0].annotations[0].value, "s");
}
#[test]
fn privacy_bucket_u64_max_no_panic() {
let result = privacy_bucket_u64(u64::MAX);
assert_ne!(result, u64::MAX);
assert!(result > 0);
}
#[test]
fn privacy_bucket_deterministic() {
for v in [1u64, 2, 3, 7, 42, 100, 1024, u64::MAX / 2] {
assert_eq!(privacy_bucket_u64(v), privacy_bucket_u64(v));
}
}
#[test]
fn synthesize_on_nonexistent_source_errors_clearly() {
let tmp = tempfile::tempdir().unwrap();
let source = tmp.path().join("nonexistent.bin");
let output = tmp.path().join("output.bin");
let err = synthesize(&source, &output, 1).unwrap_err();
let msg = err.to_string();
assert!(
msg.contains("nonexistent") || msg.contains("stat") || msg.contains("No such file"),
"error should reference the missing file: {msg}"
);
assert!(!output.exists());
}
#[test]
fn generate_on_nonexistent_shape_errors_clearly() {
let tmp = tempfile::tempdir().unwrap();
let shape_path = tmp.path().join("nonexistent.json");
let output = tmp.path().join("output.bin");
let err = generate(&shape_path, &output, 1).unwrap_err();
let msg = err.to_string();
assert!(
msg.contains("nonexistent") || msg.contains("stat") || msg.contains("No such file"),
"error should reference the missing file: {msg}"
);
assert!(!output.exists());
}
#[test]
fn gzip_input_direct_synthesize() {
let trace = make_test_trace();
let mut gz = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::default());
gz.write_all(&trace).unwrap();
let compressed = gz.finish().unwrap();
let tmp = tempfile::tempdir().unwrap();
let source = tmp.path().join("trace.bin.gz");
let output = tmp.path().join("output.bin");
std::fs::write(&source, &compressed).unwrap();
synthesize(&source, &output, 1).unwrap();
assert!(output.exists());
let generated = std::fs::read(&output).unwrap();
let mut dec = Decoder::new(&generated).unwrap();
let mut count = 0u64;
loop {
match dec.next_frame().unwrap() {
Some(DecodedFrame::Event { .. }) => count += 1,
Some(_) => {}
None => break,
}
}
assert_eq!(count, 4);
}
#[test]
fn validate_rejects_schema_name_exceeding_u16() {
let long_name = "a".repeat(65536);
let shape = TraceShape {
version: SHAPE_VERSION,
schemas: vec![ShapeSchema {
name: long_name.clone(),
has_timestamp: true,
fields: vec![ShapeField {
name: "v".into(),
field_type: FieldType::Varint as u8,
repeat_meta: None,
}],
annotations: vec![],
}],
events: vec![ShapeEvent {
schema_index: 0,
timestamp_offset_ns: Some(1000),
values: vec![ShapeValue::U(42)],
}],
summary: ShapeSummary {
event_count: 1,
duration_ns: 0,
event_type_counts: [(long_name, 1)].into_iter().collect(),
worker_cardinality: 0,
task_cardinality: 0,
poll_duration_quantiles: None,
span_duration_quantiles: None,
},
};
let err = validate_shape(&shape).unwrap_err();
let msg = err.to_string();
assert!(
msg.contains("wire limit") || msg.contains("65535"),
"expected wire limit error, got: {msg}"
);
}
#[test]
fn validate_rejects_field_count_exceeding_u16() {
let fields: Vec<ShapeField> = (0..65536)
.map(|i| ShapeField {
name: format!("f{i}"),
field_type: FieldType::Varint as u8,
repeat_meta: None,
})
.collect();
let values: Vec<ShapeValue> = (0..65536).map(|i| ShapeValue::U(i as u64)).collect();
let shape = TraceShape {
version: SHAPE_VERSION,
schemas: vec![ShapeSchema {
name: "BigSchema".into(),
has_timestamp: true,
fields,
annotations: vec![],
}],
events: vec![ShapeEvent {
schema_index: 0,
timestamp_offset_ns: Some(1000),
values,
}],
summary: ShapeSummary {
event_count: 1,
duration_ns: 0,
event_type_counts: [("BigSchema".into(), 1)].into_iter().collect(),
worker_cardinality: 0,
task_cardinality: 0,
poll_duration_quantiles: None,
span_duration_quantiles: None,
},
};
let err = validate_shape(&shape).unwrap_err();
let msg = err.to_string();
assert!(
msg.contains("wire limit") || msg.contains("65535"),
"expected field count wire limit error, got: {msg}"
);
}
#[test]
fn validate_rejects_field_name_exceeding_u16() {
let long_field_name = "f".repeat(65536);
let shape = minimal_shape(
vec![ShapeField {
name: long_field_name,
field_type: FieldType::Varint as u8,
repeat_meta: None,
}],
vec![ShapeValue::U(1)],
);
let err = validate_shape(&shape).unwrap_err();
let msg = err.to_string();
assert!(
msg.contains("wire limit") || msg.contains("65535"),
"expected field name wire limit error, got: {msg}"
);
}
#[test]
fn generate_rejects_annotation_count_exceeding_u16_without_output() {
let mut shape = minimal_shape(
vec![ShapeField {
name: "v".into(),
field_type: FieldType::Varint as u8,
repeat_meta: None,
}],
vec![ShapeValue::U(1)],
);
shape.schemas[0].annotations = (0..=u16::MAX)
.map(|_| ShapeAnnotation {
field_index: 0,
key: "unit".into(),
value: "ns".into(),
})
.collect();
let tmp = tempfile::tempdir().unwrap();
let shape_path = tmp.path().join("shape.json");
let output = tmp.path().join("output.bin");
std::fs::write(&shape_path, serde_json::to_vec(&shape).unwrap()).unwrap();
let error = generate(&shape_path, &output, 1).unwrap_err();
assert!(
error.to_string().contains("annotations"),
"unexpected error: {error:#}"
);
assert!(
!output.exists(),
"validation rejection must not create output"
);
}
#[test]
fn generate_rejects_schema_count_exceeding_registry_capacity_without_output() {
let mut shape = minimal_shape(vec![], vec![]);
shape
.schemas
.extend((1..=MAX_WIRE_SCHEMAS).map(|i| ShapeSchema {
name: format!("T{i}"),
has_timestamp: true,
fields: vec![],
annotations: vec![],
}));
let tmp = tempfile::tempdir().unwrap();
let shape_path = tmp.path().join("shape.json");
let output = tmp.path().join("output.bin");
std::fs::write(&shape_path, serde_json::to_vec(&shape).unwrap()).unwrap();
let error = generate(&shape_path, &output, 1).unwrap_err();
assert!(
error.to_string().contains("type-ID capacity"),
"unexpected error: {error:#}"
);
assert!(
!output.exists(),
"validation rejection must not create output"
);
}
#[test]
fn field_value_ref_decode_offset_beyond_data_returns_none() {
let data = [0u8; 4];
let result = FieldValueRef::decode(FieldType::Varint, &data, 100);
assert!(
result.is_none(),
"offset beyond data.len() must return None"
);
}
#[test]
fn field_value_ref_decode_offset_at_boundary() {
let data = [42u8; 8];
let result = FieldValueRef::decode(FieldType::Varint, &data, 8);
assert!(result.is_none(), "offset == data.len() must return None");
}
#[test]
fn estimator_conservative_with_pooled_string() {
let mut enc = Encoder::new();
let schema = enc
.register_schema(
"PsEvent",
vec![FieldDef::new("label", FieldType::PooledString)],
)
.unwrap();
let id = enc.intern_string("test_pool_value").unwrap();
enc.write_event(
&schema,
&[FieldValue::Varint(BASE_TS), FieldValue::PooledString(id)],
)
.unwrap();
let trace = enc.finish();
let shape = extract_shape(&trace).unwrap();
let estimated = estimate_encoded_size(&shape, 1).unwrap();
let actual = generate_trace(&shape, 1).unwrap().len() as u64;
assert!(
estimated >= actual,
"estimator ({estimated}) must be >= actual ({actual}) for pooled strings"
);
}
#[test]
fn estimator_conservative_with_long_schema_name() {
let long_name = "VeryLongSchemaNameForTesting".repeat(10);
let shape = minimal_shape_with_name(
&long_name,
vec![ShapeField {
name: "value".into(),
field_type: FieldType::Varint as u8,
repeat_meta: None,
}],
vec![ShapeValue::U(42)],
);
let estimated = estimate_encoded_size(&shape, 1).unwrap();
let actual = generate_trace(&shape, 1).unwrap().len() as u64;
assert!(
estimated >= actual,
"estimator ({estimated}) must be >= actual ({actual}) with long schema name"
);
}
fn minimal_shape_with_name(
name: &str,
fields: Vec<ShapeField>,
values: Vec<ShapeValue>,
) -> TraceShape {
TraceShape {
version: SHAPE_VERSION,
schemas: vec![ShapeSchema {
name: name.to_string(),
has_timestamp: true,
fields,
annotations: vec![],
}],
events: vec![ShapeEvent {
schema_index: 0,
timestamp_offset_ns: Some(1000),
values,
}],
summary: ShapeSummary {
event_count: 1,
duration_ns: 0,
event_type_counts: [(name.to_string(), 1)].into_iter().collect(),
worker_cardinality: 0,
task_cardinality: 0,
poll_duration_quantiles: None,
span_duration_quantiles: None,
},
}
}
}