Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
515 changes: 389 additions & 126 deletions rust/otap-dataflow/crates/admin/src/telemetry.rs

Large diffs are not rendered by default.

9 changes: 7 additions & 2 deletions rust/otap-dataflow/crates/otap/src/parquet_exporter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -444,6 +444,7 @@ mod test {
use otap_df_pdata::proto::opentelemetry::arrow::v1::ArrowPayloadType;
use otap_df_pdata::proto::opentelemetry::common::v1::{AnyValue, KeyValue, any_value::Value};
use otap_df_pdata::schema::consts;
use otap_df_telemetry::metrics::SnapshotValue;
use parquet::arrow::async_reader::ParquetRecordBatchStreamBuilder;
use tokio::fs::File;
use tokio::time::sleep;
Expand Down Expand Up @@ -1524,8 +1525,12 @@ mod test {
if desc.name == "exporter.pdata" {
saw_exporter_pdata = true;
for (_field, value) in iter {
if value.to_f64() > 0.0 {
any_positive = true;
if let SnapshotValue::Scalar(v) = value {
if v.to_f64() > 0.0 {
any_positive = true;
}
} else {
unreachable!("exporter.pdata metrics should be scalar values");
}
}
}
Expand Down
154 changes: 85 additions & 69 deletions rust/otap-dataflow/crates/telemetry-macros/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -130,96 +130,112 @@ pub fn derive_metric_set_handler(input: TokenStream) -> TokenStream {
let seg_opt = tp.path.segments.last();
if let Some(seg) = seg_opt {
let ident_ty = seg.ident.to_string();
// Expect generic arguments <u64> or <f64>
let value_type_variant = match &seg.arguments {
syn::PathArguments::AngleBracketed(ab) => {
if ab.args.len() != 1 {
return syn::Error::new(

// Handle Mmsc separately — it has no generic type parameter.
if ident_ty == "Mmsc" {
(
quote!(otap_df_telemetry::descriptor::Instrument::Mmsc),
quote!(Some(otap_df_telemetry::descriptor::Temporality::Delta)),
quote!(otap_df_telemetry::descriptor::MetricValueType::F64),
ident_ty,
)
} else {
// Expect generic arguments <u64> or <f64>
let value_type_variant = match &seg.arguments {
syn::PathArguments::AngleBracketed(ab) => {
if ab.args.len() != 1 {
return syn::Error::new(
seg.ident.span(),
"Metric field type must be one of Counter<u64|f64>, ObserveCounter<u64|f64>, UpDownCounter<u64|f64>, ObserveUpDownCounter<u64|f64>, Gauge<u64|f64>",
)
.to_compile_error()
.into();
}
match ab.args.first() {
Some(syn::GenericArgument::Type(syn::Type::Path(p)))
if p.path.is_ident("u64") =>
{
quote!(
}
match ab.args.first() {
Some(syn::GenericArgument::Type(syn::Type::Path(
p,
))) if p.path.is_ident("u64") => {
quote!(
otap_df_telemetry::descriptor::MetricValueType::U64
)
}
Some(syn::GenericArgument::Type(syn::Type::Path(p)))
if p.path.is_ident("f64") =>
{
quote!(
}
Some(syn::GenericArgument::Type(syn::Type::Path(
p,
))) if p.path.is_ident("f64") => {
quote!(
otap_df_telemetry::descriptor::MetricValueType::F64
)
}
_ => {
return syn::Error::new(
}
_ => {
return syn::Error::new(
seg.ident.span(),
"Metric field type must be one of Counter<u64|f64>, ObserveCounter<u64|f64>, UpDownCounter<u64|f64>, ObserveUpDownCounter<u64|f64>, Gauge<u64|f64>",
)
.to_compile_error()
.into();
}
}
}
}
_ => {
return syn::Error::new(
_ => {
return syn::Error::new(
seg.ident.span(),
"Metric field type must be one of Counter<u64|f64>, ObserveCounter<u64|f64>, UpDownCounter<u64|f64>, ObserveUpDownCounter<u64|f64>, Gauge<u64|f64>",
)
.to_compile_error()
.into();
}
};
let (instrument_variant, temporality_variant) = match ident_ty.as_str()
{
"Counter" => (
quote!(otap_df_telemetry::descriptor::Instrument::Counter),
quote!(Some(otap_df_telemetry::descriptor::Temporality::Delta)),
),
"ObserveCounter" => (
quote!(otap_df_telemetry::descriptor::Instrument::Counter),
quote!(Some(
otap_df_telemetry::descriptor::Temporality::Cumulative
)),
),
"UpDownCounter" => (
quote!(
}
};
let (instrument_variant, temporality_variant) = match ident_ty
.as_str()
{
"Counter" => (
quote!(otap_df_telemetry::descriptor::Instrument::Counter),
quote!(Some(
otap_df_telemetry::descriptor::Temporality::Delta
)),
),
"ObserveCounter" => (
quote!(otap_df_telemetry::descriptor::Instrument::Counter),
quote!(Some(
otap_df_telemetry::descriptor::Temporality::Cumulative
)),
),
"UpDownCounter" => (
quote!(
otap_df_telemetry::descriptor::Instrument::UpDownCounter
),
quote!(Some(otap_df_telemetry::descriptor::Temporality::Delta)),
),
"ObserveUpDownCounter" => (
quote!(
quote!(Some(
otap_df_telemetry::descriptor::Temporality::Delta
)),
),
"ObserveUpDownCounter" => (
quote!(
otap_df_telemetry::descriptor::Instrument::UpDownCounter
),
quote!(Some(
otap_df_telemetry::descriptor::Temporality::Cumulative
)),
),
"Gauge" => (
quote!(otap_df_telemetry::descriptor::Instrument::Gauge),
quote!(None),
),
other => {
return syn::Error::new(
seg.ident.span(),
format!("Unsupported metric instrument type: {other}"),
)
.to_compile_error()
.into();
}
};
(
instrument_variant,
temporality_variant,
value_type_variant,
ident_ty,
)
quote!(Some(
otap_df_telemetry::descriptor::Temporality::Cumulative
)),
),
"Gauge" => (
quote!(otap_df_telemetry::descriptor::Instrument::Gauge),
quote!(None),
),
other => {
return syn::Error::new(
seg.ident.span(),
format!("Unsupported metric instrument type: {other}"),
)
.to_compile_error()
.into();
}
};
(
instrument_variant,
temporality_variant,
value_type_variant,
ident_ty,
)
} // end of else (non-Mmsc types)
} else {
return syn::Error::new(
field.ty.span(),
Expand All @@ -246,10 +262,10 @@ pub fn derive_metric_set_handler(input: TokenStream) -> TokenStream {
metric_field_value_types.push(value_type_variant);

match instrument_ty_name.as_str() {
"Counter" | "UpDownCounter" => {
"Counter" | "UpDownCounter" | "Mmsc" => {
metric_field_clear_stmts.push(quote!( self.#field_ident.reset(); ));
metric_field_needs_flush_checks.push(quote!(
if !otap_df_telemetry::metrics::MetricValue::from(self.#field_ident.get()).is_zero() {
if !otap_df_telemetry::metrics::SnapshotValue::from(self.#field_ident.get()).is_zero() {
return true;
}
));
Expand Down Expand Up @@ -282,9 +298,9 @@ pub fn derive_metric_set_handler(input: TokenStream) -> TokenStream {
};
&#desc_ident
}
fn snapshot_values(&self) -> ::std::vec::Vec<otap_df_telemetry::metrics::MetricValue> {
fn snapshot_values(&self) -> ::std::vec::Vec<otap_df_telemetry::metrics::SnapshotValue> {
let mut out = ::std::vec::Vec::with_capacity(self.descriptor().metrics.len());
#( out.push(otap_df_telemetry::metrics::MetricValue::from(self.#metric_field_idents.get())); )*
#( out.push(otap_df_telemetry::metrics::SnapshotValue::from(self.#metric_field_idents.get())); )*
out
}
fn clear_values(&mut self) {
Expand Down
1 change: 1 addition & 0 deletions rust/otap-dataflow/crates/telemetry/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -49,3 +49,4 @@ tracing-subscriber = { workspace = true, features = ["env-filter","registry", "s
[dev-dependencies]
tower = { workspace = true }
tempfile = { workspace = true }
opentelemetry_sdk = { workspace = true, features = ["testing"] }
26 changes: 13 additions & 13 deletions rust/otap-dataflow/crates/telemetry/src/collector.rs
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,7 @@ mod tests {
MetricsDescriptor, MetricsField, Temporality,
};
use crate::metrics::MetricSetHandler;
use crate::metrics::MetricValue;
use crate::metrics::SnapshotValue;
use crate::registry::MetricSetKey;
use std::collections::HashMap;
use std::fmt::Debug;
Expand All @@ -102,13 +102,13 @@ mod tests {

#[derive(Debug)]
struct MockMetricSet {
values: Vec<MetricValue>,
values: Vec<SnapshotValue>,
}

impl MockMetricSet {
fn new() -> Self {
Self {
values: vec![MetricValue::U64(0), MetricValue::U64(0)],
values: vec![SnapshotValue::from(0u64), SnapshotValue::from(0u64)],
}
}
}
Expand Down Expand Up @@ -154,11 +154,11 @@ mod tests {
fn descriptor(&self) -> &'static MetricsDescriptor {
&MOCK_METRICS_DESCRIPTOR
}
fn snapshot_values(&self) -> Vec<MetricValue> {
fn snapshot_values(&self) -> Vec<SnapshotValue> {
self.values.clone()
}
fn clear_values(&mut self) {
self.values.iter_mut().for_each(MetricValue::reset);
self.values.iter_mut().for_each(SnapshotValue::reset);
}
fn needs_flush(&self) -> bool {
self.values.iter().any(|&v| !v.is_zero())
Expand Down Expand Up @@ -207,7 +207,7 @@ mod tests {
}
}

fn create_test_snapshot(key: MetricSetKey, values: Vec<MetricValue>) -> MetricSetSnapshot {
fn create_test_snapshot(key: MetricSetKey, values: Vec<SnapshotValue>) -> MetricSetSnapshot {
MetricSetSnapshot {
key,
metrics: values,
Expand Down Expand Up @@ -249,14 +249,14 @@ mod tests {
reporter
.report_snapshot(create_test_snapshot(
key,
vec![MetricValue::U64(10), MetricValue::U64(20)],
vec![SnapshotValue::from(10u64), SnapshotValue::from(20u64)],
))
.await
.unwrap();
reporter
.report_snapshot(create_test_snapshot(
key,
vec![MetricValue::U64(5), MetricValue::U64(15)],
vec![SnapshotValue::from(5u64), SnapshotValue::from(15u64)],
))
.await
.unwrap();
Expand All @@ -275,9 +275,9 @@ mod tests {
assert_eq!(collected.len(), 2);
// Order follows descriptor order
assert_eq!(collected[0].0, "counter1");
assert_eq!(collected[0].1, MetricValue::U64(15));
assert_eq!(collected[0].1, SnapshotValue::from(15u64));
assert_eq!(collected[1].0, "counter2");
assert_eq!(collected[1].1, MetricValue::U64(35));
assert_eq!(collected[1].1, SnapshotValue::from(35u64));

// Close the channel and ensure loop ends returning None
drop(reporter);
Expand All @@ -298,7 +298,7 @@ mod tests {
reporter
.report_snapshot(create_test_snapshot(
key,
vec![MetricValue::U64(7), MetricValue::U64(0)],
vec![SnapshotValue::from(7u64), SnapshotValue::from(0u64)],
))
.await
.unwrap();
Expand All @@ -314,8 +314,8 @@ mod tests {
assert_eq!(
first,
vec![
("counter1", MetricValue::U64(7)),
("counter2", MetricValue::U64(0))
("counter1", SnapshotValue::from(7u64)),
("counter2", SnapshotValue::from(0u64))
]
);

Expand Down
5 changes: 5 additions & 0 deletions rust/otap-dataflow/crates/telemetry/src/descriptor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,11 @@ pub enum Instrument {
Gauge,
/// Distribution of recorded values, used for latencies or request sizes
Histogram,
/// Pre-aggregated min/max/sum/count summary.
///
/// Internally tracked as an `Mmsc` instrument; the dispatcher exports the
/// aggregated snapshot as a synthetic OTel histogram without bucket counts.
Mmsc,
}

/// Aggregation temporality for sum-like instruments.
Expand Down
Loading
Loading