Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
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
509 changes: 394 additions & 115 deletions rust/otap-dataflow/crates/admin/src/telemetry.rs

Large diffs are not rendered by default.

148 changes: 82 additions & 66 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,7 +262,7 @@ 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() {
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"] }
16 changes: 8 additions & 8 deletions rust/otap-dataflow/crates/telemetry/src/collector.rs
Original file line number Diff line number Diff line change
Expand Up @@ -108,7 +108,7 @@ mod tests {
impl MockMetricSet {
fn new() -> Self {
Self {
values: vec![MetricValue::U64(0), MetricValue::U64(0)],
values: vec![MetricValue::from(0u64), MetricValue::from(0u64)],
}
}
}
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![MetricValue::from(10u64), MetricValue::from(20u64)],
))
.await
.unwrap();
reporter
.report_snapshot(create_test_snapshot(
key,
vec![MetricValue::U64(5), MetricValue::U64(15)],
vec![MetricValue::from(5u64), MetricValue::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, MetricValue::from(15u64));
assert_eq!(collected[1].0, "counter2");
assert_eq!(collected[1].1, MetricValue::U64(35));
assert_eq!(collected[1].1, MetricValue::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![MetricValue::from(7u64), MetricValue::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", MetricValue::from(7u64)),
("counter2", MetricValue::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