pub mod arrow;
pub mod convert;
pub mod decode;
pub mod error;
pub mod output;
pub mod schemas;
pub mod transform;
pub mod value;
#[cfg(all(feature = "wasm", target_arch = "wasm32"))]
pub mod wasm;
#[cfg(feature = "ffi")]
pub mod ffi;
use ::arrow::record_batch::RecordBatch;
pub use arrow::{
exp_histogram_schema, extract_min_timestamp_micros, extract_service_name, gauge_schema,
group_batch_by_service, histogram_schema, logs_schema, sum_schema, traces_schema,
values_to_arrow, PartitionedBatch, PartitionedMetrics, ServiceGroupedBatches,
};
pub use decode::{
count_skipped_metric_data_points, decode_logs, decode_metrics, decode_traces,
normalise_json_value, normalize_json_bytes, DecodeMetricsResult, InputFormat, MetricSkipCounts,
SkippedMetrics,
};
pub use error::{Error, Result};
#[cfg(feature = "parquet")]
pub use output::to_parquet;
pub use output::{to_ipc, to_json};
pub use schemas::{schema_def, schema_defs, FieldType, SchemaDef, SchemaField};
pub use value::{KeyString, ObjectMap, Value};
#[derive(Debug)]
pub struct MetricBatches {
pub gauge: Option<RecordBatch>,
pub sum: Option<RecordBatch>,
pub histogram: Option<RecordBatch>,
pub exp_histogram: Option<RecordBatch>,
pub skipped: SkippedMetrics,
}
#[derive(Debug)]
pub struct JsonMetricBatches {
pub gauge: Vec<serde_json::Value>,
pub sum: Vec<serde_json::Value>,
pub histogram: Vec<serde_json::Value>,
pub exp_histogram: Vec<serde_json::Value>,
pub skipped: SkippedMetrics,
}
#[derive(Debug, Default)]
pub struct MetricValues {
pub gauge: Vec<Value>,
pub sum: Vec<Value>,
pub histogram: Vec<Value>,
pub exp_histogram: Vec<Value>,
}
pub fn transform_logs(bytes: &[u8], format: InputFormat) -> Result<RecordBatch> {
let values = decode_logs(bytes, format)?;
let transformed = apply_log_transform(values);
let batch = values_to_arrow(&transformed, &logs_schema())?;
Ok(batch)
}
pub fn transform_logs_json(bytes: &[u8], format: InputFormat) -> Result<Vec<serde_json::Value>> {
let values = decode_logs(bytes, format)?;
let transformed = apply_log_transform(values);
values_to_json(transformed, "log")
}
pub fn transform_traces(bytes: &[u8], format: InputFormat) -> Result<RecordBatch> {
let values = decode_traces(bytes, format)?;
let transformed = apply_trace_transform(values);
let batch = values_to_arrow(&transformed, &traces_schema())?;
Ok(batch)
}
pub fn transform_traces_json(bytes: &[u8], format: InputFormat) -> Result<Vec<serde_json::Value>> {
let values = decode_traces(bytes, format)?;
let transformed = apply_trace_transform(values);
values_to_json(transformed, "span")
}
pub fn transform_metrics(bytes: &[u8], format: InputFormat) -> Result<MetricBatches> {
let decode_result = decode_metrics(bytes, format)?;
let metric_values = apply_metric_transform(decode_result.values);
let gauge = if metric_values.gauge.is_empty() {
None
} else {
Some(values_to_arrow(&metric_values.gauge, &gauge_schema())?)
};
let sum = if metric_values.sum.is_empty() {
None
} else {
Some(values_to_arrow(&metric_values.sum, &sum_schema())?)
};
let histogram = if metric_values.histogram.is_empty() {
None
} else {
Some(values_to_arrow(
&metric_values.histogram,
&histogram_schema(),
)?)
};
let exp_histogram = if metric_values.exp_histogram.is_empty() {
None
} else {
Some(values_to_arrow(
&metric_values.exp_histogram,
&exp_histogram_schema(),
)?)
};
Ok(MetricBatches {
gauge,
sum,
histogram,
exp_histogram,
skipped: decode_result.skipped,
})
}
pub fn transform_metrics_json(bytes: &[u8], format: InputFormat) -> Result<JsonMetricBatches> {
let decode_result = decode_metrics(bytes, format)?;
let metric_values = apply_metric_transform(decode_result.values);
Ok(JsonMetricBatches {
gauge: values_to_json(metric_values.gauge, "gauge metric")?,
sum: values_to_json(metric_values.sum, "sum metric")?,
histogram: values_to_json(metric_values.histogram, "histogram metric")?,
exp_histogram: values_to_json(metric_values.exp_histogram, "exp_histogram metric")?,
skipped: decode_result.skipped,
})
}
pub fn transform_logs_partitioned(
bytes: &[u8],
format: InputFormat,
) -> Result<ServiceGroupedBatches> {
let batch = transform_logs(bytes, format)?;
Ok(group_batch_by_service(batch))
}
pub fn transform_traces_partitioned(
bytes: &[u8],
format: InputFormat,
) -> Result<ServiceGroupedBatches> {
let batch = transform_traces(bytes, format)?;
Ok(group_batch_by_service(batch))
}
pub fn transform_metrics_partitioned(
bytes: &[u8],
format: InputFormat,
) -> Result<PartitionedMetrics> {
let batches = transform_metrics(bytes, format)?;
let gauge = match batches.gauge {
Some(batch) => group_batch_by_service(batch),
None => ServiceGroupedBatches::default(),
};
let sum = match batches.sum {
Some(batch) => group_batch_by_service(batch),
None => ServiceGroupedBatches::default(),
};
let histogram = match batches.histogram {
Some(batch) => group_batch_by_service(batch),
None => ServiceGroupedBatches::default(),
};
let exp_histogram = match batches.exp_histogram {
Some(batch) => group_batch_by_service(batch),
None => ServiceGroupedBatches::default(),
};
Ok(PartitionedMetrics {
gauge,
sum,
histogram,
exp_histogram,
skipped: batches.skipped,
})
}
pub fn apply_log_transform(values: Vec<Value>) -> Vec<Value> {
values.into_iter().map(transform::transform_log).collect()
}
pub fn apply_trace_transform(values: Vec<Value>) -> Vec<Value> {
values.into_iter().map(transform::transform_trace).collect()
}
pub fn apply_metric_transform(values: Vec<Value>) -> MetricValues {
let mut gauge_count = 0;
let mut sum_count = 0;
let mut histogram_count = 0;
let mut exp_histogram_count = 0;
for value in &values {
match extract_metric_type(value) {
"gauge" => gauge_count += 1,
"sum" => sum_count += 1,
"histogram" => histogram_count += 1,
"exp_histogram" => exp_histogram_count += 1,
_ => {}
}
}
let mut result = MetricValues {
gauge: Vec::with_capacity(gauge_count),
sum: Vec::with_capacity(sum_count),
histogram: Vec::with_capacity(histogram_count),
exp_histogram: Vec::with_capacity(exp_histogram_count),
};
for value in values {
match extract_metric_type(&value) {
"gauge" => {
result.gauge.push(transform::transform_gauge(value));
}
"sum" => {
result.sum.push(transform::transform_sum(value));
}
"histogram" => {
result.histogram.push(transform::transform_histogram(value));
}
"exp_histogram" => {
result
.exp_histogram
.push(transform::transform_exp_histogram(value));
}
_ => {
}
}
}
result
}
fn values_to_json(values: Vec<Value>, label: &str) -> Result<Vec<serde_json::Value>> {
let mut out = Vec::with_capacity(values.len());
for (idx, value) in values.into_iter().enumerate() {
let json = crate::convert::value_to_json(&value).ok_or_else(|| {
Error::InvalidInput(format!(
"{label} record {idx} contains unrepresentable JSON value"
))
})?;
out.push(json);
}
Ok(out)
}
fn extract_metric_type(value: &Value) -> &'static str {
if let Value::Object(map) = value {
if let Some(Value::Bytes(bytes)) = map.get("_metric_type") {
return match bytes.as_ref() {
b"gauge" => "gauge",
b"sum" => "sum",
b"histogram" => "histogram",
b"exp_histogram" => "exp_histogram",
_ => "",
};
}
}
""
}
#[cfg(test)]
mod tests {
use super::*;
use opentelemetry_proto::tonic::{
collector::logs::v1::ExportLogsServiceRequest,
collector::metrics::v1::ExportMetricsServiceRequest,
collector::trace::v1::ExportTraceServiceRequest,
common::v1::{any_value, AnyValue, InstrumentationScope, KeyValue},
logs::v1::{LogRecord, ResourceLogs, ScopeLogs},
metrics::v1::{
metric::Data, Gauge, Metric, NumberDataPoint, ResourceMetrics, ScopeMetrics, Sum,
},
resource::v1::Resource,
trace::v1::{ResourceSpans, ScopeSpans, Span},
};
use prost::Message;
fn create_test_log_request() -> ExportLogsServiceRequest {
ExportLogsServiceRequest {
resource_logs: vec![ResourceLogs {
resource: Some(Resource {
attributes: vec![KeyValue {
key: "service.name".to_string(),
value: Some(AnyValue {
value: Some(any_value::Value::StringValue("test-service".to_string())),
}),
}],
..Default::default()
}),
scope_logs: vec![ScopeLogs {
scope: Some(InstrumentationScope {
name: "test-lib".to_string(),
version: "1.0.0".to_string(),
..Default::default()
}),
log_records: vec![LogRecord {
time_unix_nano: 1_700_000_000_000_000_000,
observed_time_unix_nano: 1_700_000_000_100_000_000,
severity_number: 9,
severity_text: "INFO".to_string(),
body: Some(AnyValue {
value: Some(any_value::Value::StringValue(
"Test log message".to_string(),
)),
}),
attributes: vec![KeyValue {
key: "log.key".to_string(),
value: Some(AnyValue {
value: Some(any_value::Value::StringValue("log-value".to_string())),
}),
}],
trace_id: vec![0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15],
span_id: vec![0, 1, 2, 3, 4, 5, 6, 7],
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
}
}
fn create_test_trace_request() -> ExportTraceServiceRequest {
ExportTraceServiceRequest {
resource_spans: vec![ResourceSpans {
resource: Some(Resource {
attributes: vec![KeyValue {
key: "service.name".to_string(),
value: Some(AnyValue {
value: Some(any_value::Value::StringValue("test-service".to_string())),
}),
}],
..Default::default()
}),
scope_spans: vec![ScopeSpans {
scope: Some(InstrumentationScope {
name: "test-lib".to_string(),
version: "1.0.0".to_string(),
..Default::default()
}),
spans: vec![Span {
trace_id: vec![0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15],
span_id: vec![0, 1, 2, 3, 4, 5, 6, 7],
parent_span_id: vec![],
name: "test-span".to_string(),
kind: 1, start_time_unix_nano: 1_700_000_000_000_000_000,
end_time_unix_nano: 1_700_000_000_100_000_000,
attributes: vec![KeyValue {
key: "span.key".to_string(),
value: Some(AnyValue {
value: Some(any_value::Value::StringValue(
"span-value".to_string(),
)),
}),
}],
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
}
}
fn create_test_metrics_request() -> ExportMetricsServiceRequest {
ExportMetricsServiceRequest {
resource_metrics: vec![ResourceMetrics {
resource: Some(Resource {
attributes: vec![KeyValue {
key: "service.name".to_string(),
value: Some(AnyValue {
value: Some(any_value::Value::StringValue("test-service".to_string())),
}),
}],
..Default::default()
}),
scope_metrics: vec![ScopeMetrics {
scope: Some(InstrumentationScope {
name: "test-lib".to_string(),
version: "1.0.0".to_string(),
..Default::default()
}),
metrics: vec![
Metric {
name: "test.gauge".to_string(),
description: "A test gauge".to_string(),
unit: "1".to_string(),
data: Some(Data::Gauge(Gauge {
data_points: vec![NumberDataPoint {
time_unix_nano: 1_700_000_000_000_000_000,
start_time_unix_nano: 1_699_999_000_000_000_000,
value: Some(
opentelemetry_proto::tonic::metrics::v1::number_data_point::Value::AsDouble(42.5),
),
attributes: vec![KeyValue {
key: "metric.key".to_string(),
value: Some(AnyValue {
value: Some(any_value::Value::StringValue(
"metric-value".to_string(),
)),
}),
}],
..Default::default()
}],
})),
..Default::default()
},
Metric {
name: "test.sum".to_string(),
description: "A test sum".to_string(),
unit: "bytes".to_string(),
data: Some(Data::Sum(Sum {
data_points: vec![NumberDataPoint {
time_unix_nano: 1_700_000_000_000_000_000,
start_time_unix_nano: 1_699_999_000_000_000_000,
value: Some(
opentelemetry_proto::tonic::metrics::v1::number_data_point::Value::AsDouble(100.0),
),
attributes: vec![],
..Default::default()
}],
aggregation_temporality: 2, is_monotonic: true,
})),
..Default::default()
},
],
..Default::default()
}],
..Default::default()
}],
}
}
fn create_gauge_only_metrics_request() -> ExportMetricsServiceRequest {
ExportMetricsServiceRequest {
resource_metrics: vec![ResourceMetrics {
resource: Some(Resource {
attributes: vec![KeyValue {
key: "service.name".to_string(),
value: Some(AnyValue {
value: Some(any_value::Value::StringValue("test-service".to_string())),
}),
}],
..Default::default()
}),
scope_metrics: vec![ScopeMetrics {
scope: Some(InstrumentationScope::default()),
metrics: vec![Metric {
name: "test.gauge".to_string(),
data: Some(Data::Gauge(Gauge {
data_points: vec![NumberDataPoint {
time_unix_nano: 1_700_000_000_000_000_000,
value: Some(
opentelemetry_proto::tonic::metrics::v1::number_data_point::Value::AsDouble(1.0),
),
..Default::default()
}],
})),
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
}
}
fn create_sum_only_metrics_request() -> ExportMetricsServiceRequest {
ExportMetricsServiceRequest {
resource_metrics: vec![ResourceMetrics {
resource: Some(Resource {
attributes: vec![KeyValue {
key: "service.name".to_string(),
value: Some(AnyValue {
value: Some(any_value::Value::StringValue("test-service".to_string())),
}),
}],
..Default::default()
}),
scope_metrics: vec![ScopeMetrics {
scope: Some(InstrumentationScope::default()),
metrics: vec![Metric {
name: "test.sum".to_string(),
data: Some(Data::Sum(Sum {
data_points: vec![NumberDataPoint {
time_unix_nano: 1_700_000_000_000_000_000,
value: Some(
opentelemetry_proto::tonic::metrics::v1::number_data_point::Value::AsDouble(2.0),
),
..Default::default()
}],
aggregation_temporality: 1,
is_monotonic: false,
})),
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
}
}
#[test]
fn test_transform_logs_protobuf() {
let request = create_test_log_request();
let bytes = request.encode_to_vec();
let batch = transform_logs(&bytes, InputFormat::Protobuf).unwrap();
assert_eq!(batch.num_rows(), 1);
assert!(batch.num_columns() > 0);
let schema = batch.schema();
assert!(schema.field_with_name("timestamp").is_ok());
assert!(schema.field_with_name("service_name").is_ok());
assert!(schema.field_with_name("severity_number").is_ok());
}
#[test]
fn test_transform_logs_json() {
let json = r#"{
"resourceLogs": [{
"resource": { "attributes": [{ "key": "service.name", "value": { "stringValue": "json-svc" } }]},
"scopeLogs": [{
"scope": { "name": "lib", "version": "1" },
"logRecords": [{
"timeUnixNano": "1700000000000000000",
"observedTimeUnixNano": "1700000000100000000",
"severityNumber": 9,
"severityText": "INFO",
"body": { "stringValue": "JSON log" }
}]
}]
}]
}"#;
let batch = transform_logs(json.as_bytes(), InputFormat::Json).unwrap();
assert_eq!(batch.num_rows(), 1);
}
#[test]
fn test_transform_logs_empty() {
let request = ExportLogsServiceRequest {
resource_logs: vec![],
};
let bytes = request.encode_to_vec();
let batch = transform_logs(&bytes, InputFormat::Protobuf).unwrap();
assert_eq!(batch.num_rows(), 0);
}
#[test]
fn test_transform_traces_protobuf() {
let request = create_test_trace_request();
let bytes = request.encode_to_vec();
let batch = transform_traces(&bytes, InputFormat::Protobuf).unwrap();
assert_eq!(batch.num_rows(), 1);
let schema = batch.schema();
assert!(schema.field_with_name("timestamp").is_ok());
assert!(schema.field_with_name("trace_id").is_ok());
assert!(schema.field_with_name("span_id").is_ok());
assert!(schema.field_with_name("span_name").is_ok());
}
#[test]
fn test_transform_traces_json() {
let json = r#"{
"resourceSpans": [{
"resource": { "attributes": [{ "key": "service.name", "value": { "stringValue": "json-svc" } }]},
"scopeSpans": [{
"scope": { "name": "lib" },
"spans": [{
"traceId": "00010203040506070809101112131415",
"spanId": "0001020304050607",
"name": "json-span",
"kind": 1,
"startTimeUnixNano": "1700000000000000000",
"endTimeUnixNano": "1700000000100000000"
}]
}]
}]
}"#;
let batch = transform_traces(json.as_bytes(), InputFormat::Json).unwrap();
assert_eq!(batch.num_rows(), 1);
}
#[test]
fn test_transform_traces_empty() {
let request = ExportTraceServiceRequest {
resource_spans: vec![],
};
let bytes = request.encode_to_vec();
let batch = transform_traces(&bytes, InputFormat::Protobuf).unwrap();
assert_eq!(batch.num_rows(), 0);
}
#[test]
fn test_transform_metrics_protobuf() {
let request = create_test_metrics_request();
let bytes = request.encode_to_vec();
let batches = transform_metrics(&bytes, InputFormat::Protobuf).unwrap();
assert!(batches.gauge.is_some());
assert!(batches.sum.is_some());
let gauge = batches.gauge.unwrap();
let sum = batches.sum.unwrap();
assert_eq!(gauge.num_rows(), 1);
assert_eq!(sum.num_rows(), 1);
let gauge_schema = gauge.schema();
assert!(gauge_schema.field_with_name("metric_name").is_ok());
assert!(gauge_schema.field_with_name("value").is_ok());
let sum_schema = sum.schema();
assert!(sum_schema
.field_with_name("aggregation_temporality")
.is_ok());
assert!(sum_schema.field_with_name("is_monotonic").is_ok());
}
#[test]
fn test_transform_metrics_gauge_only() {
let request = create_gauge_only_metrics_request();
let bytes = request.encode_to_vec();
let batches = transform_metrics(&bytes, InputFormat::Protobuf).unwrap();
assert!(batches.gauge.is_some());
assert!(batches.sum.is_none());
}
#[test]
fn test_transform_metrics_sum_only() {
let request = create_sum_only_metrics_request();
let bytes = request.encode_to_vec();
let batches = transform_metrics(&bytes, InputFormat::Protobuf).unwrap();
assert!(batches.gauge.is_none());
assert!(batches.sum.is_some());
}
#[test]
fn test_transform_metrics_empty() {
let request = ExportMetricsServiceRequest {
resource_metrics: vec![],
};
let bytes = request.encode_to_vec();
let batches = transform_metrics(&bytes, InputFormat::Protobuf).unwrap();
assert!(batches.gauge.is_none());
assert!(batches.sum.is_none());
}
#[test]
fn test_apply_log_transform() {
let request = create_test_log_request();
let bytes = request.encode_to_vec();
let decoded = decode_logs(&bytes, InputFormat::Protobuf).unwrap();
let transformed = apply_log_transform(decoded);
assert_eq!(transformed.len(), 1);
if let Value::Object(map) = &transformed[0] {
let ts_key: KeyString = "timestamp".into();
let svc_key: KeyString = "service_name".into();
assert!(map.get(&ts_key).is_some());
assert!(map.get(&svc_key).is_some());
} else {
panic!("Expected object value");
}
}
#[test]
fn test_apply_trace_transform() {
let request = create_test_trace_request();
let bytes = request.encode_to_vec();
let decoded = decode_traces(&bytes, InputFormat::Protobuf).unwrap();
let transformed = apply_trace_transform(decoded);
assert_eq!(transformed.len(), 1);
if let Value::Object(map) = &transformed[0] {
let ts_key: KeyString = "timestamp".into();
let span_key: KeyString = "span_name".into();
assert!(map.get(&ts_key).is_some());
assert!(map.get(&span_key).is_some());
} else {
panic!("Expected object value");
}
}
#[test]
fn test_apply_metric_transform() {
let request = create_test_metrics_request();
let bytes = request.encode_to_vec();
let decode_result = decode_metrics(&bytes, InputFormat::Protobuf).unwrap();
let transformed = apply_metric_transform(decode_result.values);
assert_eq!(transformed.gauge.len(), 1);
assert_eq!(transformed.sum.len(), 1);
}
#[test]
fn test_apply_log_transform_empty() {
let transformed = apply_log_transform(vec![]);
assert!(transformed.is_empty());
}
#[test]
fn test_apply_trace_transform_empty() {
let transformed = apply_trace_transform(vec![]);
assert!(transformed.is_empty());
}
#[test]
fn test_apply_metric_transform_empty() {
let transformed = apply_metric_transform(vec![]);
assert!(transformed.gauge.is_empty());
assert!(transformed.sum.is_empty());
}
#[test]
fn test_transform_logs_invalid_protobuf() {
let result = transform_logs(b"not valid protobuf", InputFormat::Protobuf);
assert!(result.is_err());
}
#[test]
fn test_transform_logs_invalid_json() {
let result = transform_logs(b"not valid json", InputFormat::Json);
assert!(result.is_err());
}
#[test]
fn test_transform_traces_invalid_protobuf() {
let result = transform_traces(b"not valid protobuf", InputFormat::Protobuf);
assert!(result.is_err());
}
#[test]
fn test_transform_metrics_invalid_protobuf() {
let result = transform_metrics(b"not valid protobuf", InputFormat::Protobuf);
assert!(result.is_err());
}
#[test]
fn test_metric_batches_debug() {
let batches = MetricBatches {
gauge: None,
sum: None,
histogram: None,
exp_histogram: None,
skipped: SkippedMetrics::default(),
};
let debug_str = format!("{batches:?}");
assert!(debug_str.contains("MetricBatches"));
}
#[test]
fn test_metric_values_default() {
let values = MetricValues::default();
assert!(values.gauge.is_empty());
assert!(values.sum.is_empty());
assert!(values.histogram.is_empty());
assert!(values.exp_histogram.is_empty());
}
#[test]
fn test_metric_values_debug() {
let values = MetricValues::default();
let debug_str = format!("{values:?}");
assert!(debug_str.contains("MetricValues"));
}
#[test]
fn test_full_pipeline_logs_to_json_output() {
let request = create_test_log_request();
let bytes = request.encode_to_vec();
let batch = transform_logs(&bytes, InputFormat::Protobuf).unwrap();
let json_output = to_json(&batch).unwrap();
assert!(!json_output.is_empty());
let json_str = String::from_utf8(json_output).unwrap();
assert!(json_str.contains('\n') || !json_str.is_empty());
}
#[test]
fn test_full_pipeline_traces_to_ipc_output() {
let request = create_test_trace_request();
let bytes = request.encode_to_vec();
let batch = transform_traces(&bytes, InputFormat::Protobuf).unwrap();
let ipc_output = to_ipc(&batch).unwrap();
assert!(!ipc_output.is_empty());
}
#[test]
fn test_full_pipeline_metrics_round_trip() {
let request = create_test_metrics_request();
let bytes = request.encode_to_vec();
let batches = transform_metrics(&bytes, InputFormat::Protobuf).unwrap();
if let Some(gauge) = &batches.gauge {
let json = to_json(gauge).unwrap();
assert!(!json.is_empty());
}
if let Some(sum) = &batches.sum {
let json = to_json(sum).unwrap();
assert!(!json.is_empty());
}
}
#[test]
fn test_timestamp_not_epoch_traces() {
let request = ExportTraceServiceRequest {
resource_spans: vec![ResourceSpans {
resource: Some(Resource {
attributes: vec![KeyValue {
key: "service.name".to_string(),
value: Some(AnyValue {
value: Some(any_value::Value::StringValue("test-service".to_string())),
}),
}],
..Default::default()
}),
scope_spans: vec![ScopeSpans {
scope: Some(InstrumentationScope::default()),
spans: vec![Span {
trace_id: vec![0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15],
span_id: vec![0, 1, 2, 3, 4, 5, 6, 7],
name: "test-span".to_string(),
start_time_unix_nano: 1_703_265_600_000_000_000, end_time_unix_nano: 1_703_265_600_100_000_000,
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
};
let bytes = request.encode_to_vec();
let batch = transform_traces(&bytes, InputFormat::Protobuf).unwrap();
assert_eq!(batch.num_rows(), 1);
let schema = batch.schema();
let ts_idx = schema.index_of("timestamp").unwrap();
let ts_column = batch
.column(ts_idx)
.as_any()
.downcast_ref::<::arrow::array::TimestampMicrosecondArray>()
.expect("timestamp should be TimestampMicrosecondArray");
let ts_value = ts_column.value(0);
let expected_micros: i64 = 1_703_265_600_000_000;
assert_eq!(
ts_value, expected_micros,
"Timestamp should be Dec 22, 2023, not epoch (1970)"
);
assert!(
ts_value > 1_600_000_000_000_000,
"Timestamp {ts_value} appears to be too small, possibly 1970 date"
);
}
#[test]
fn test_timestamp_not_epoch_logs() {
let request = ExportLogsServiceRequest {
resource_logs: vec![ResourceLogs {
resource: Some(Resource {
attributes: vec![KeyValue {
key: "service.name".to_string(),
value: Some(AnyValue {
value: Some(any_value::Value::StringValue("test-service".to_string())),
}),
}],
..Default::default()
}),
scope_logs: vec![ScopeLogs {
scope: Some(InstrumentationScope::default()),
log_records: vec![LogRecord {
time_unix_nano: 1_703_265_600_000_000_000, observed_time_unix_nano: 1_703_265_600_100_000_000,
severity_number: 9,
severity_text: "INFO".to_string(),
body: Some(AnyValue {
value: Some(any_value::Value::StringValue("Test log".to_string())),
}),
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
};
let bytes = request.encode_to_vec();
let batch = transform_logs(&bytes, InputFormat::Protobuf).unwrap();
assert_eq!(batch.num_rows(), 1);
let schema = batch.schema();
let ts_idx = schema.index_of("timestamp").unwrap();
let ts_column = batch
.column(ts_idx)
.as_any()
.downcast_ref::<::arrow::array::TimestampMicrosecondArray>()
.expect("timestamp should be TimestampMicrosecondArray");
let ts_value = ts_column.value(0);
let expected_micros: i64 = 1_703_265_600_000_000;
assert_eq!(
ts_value, expected_micros,
"Timestamp should be Dec 22, 2023, not epoch (1970)"
);
assert!(
ts_value > 1_600_000_000_000_000,
"Timestamp {ts_value} appears to be too small, possibly 1970 date"
);
}
#[test]
fn test_timestamp_not_epoch_metrics() {
let request = ExportMetricsServiceRequest {
resource_metrics: vec![ResourceMetrics {
resource: Some(Resource {
attributes: vec![KeyValue {
key: "service.name".to_string(),
value: Some(AnyValue {
value: Some(any_value::Value::StringValue("test-service".to_string())),
}),
}],
..Default::default()
}),
scope_metrics: vec![ScopeMetrics {
scope: Some(InstrumentationScope::default()),
metrics: vec![Metric {
name: "test.gauge".to_string(),
data: Some(Data::Gauge(Gauge {
data_points: vec![NumberDataPoint {
time_unix_nano: 1_703_265_600_000_000_000, value: Some(
opentelemetry_proto::tonic::metrics::v1::number_data_point::Value::AsDouble(42.0),
),
..Default::default()
}],
})),
..Default::default()
}],
..Default::default()
}],
..Default::default()
}],
};
let bytes = request.encode_to_vec();
let batches = transform_metrics(&bytes, InputFormat::Protobuf).unwrap();
let gauge = batches.gauge.expect("should have gauge metrics");
assert_eq!(gauge.num_rows(), 1);
let schema = gauge.schema();
let ts_idx = schema.index_of("timestamp").unwrap();
let ts_column = gauge
.column(ts_idx)
.as_any()
.downcast_ref::<::arrow::array::TimestampMicrosecondArray>()
.expect("timestamp should be TimestampMicrosecondArray");
let ts_value = ts_column.value(0);
let expected_micros: i64 = 1_703_265_600_000_000;
assert_eq!(
ts_value, expected_micros,
"Timestamp should be Dec 22, 2023, not epoch (1970)"
);
assert!(
ts_value > 1_600_000_000_000_000,
"Timestamp {ts_value} appears to be too small, possibly 1970 date"
);
}
}