1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
//! Opt-in compute-duration timing for processors.
//!
//! Processors that perform meaningful synchronous compute can add a
//! [`ComputeDuration`] field and call [`ComputeDuration::timed`] to
//! measure the wall-clock duration of that work. Timing is gated on
//! the `NODE_LOCAL_DURATION` interest at the detailed metric level or through
//! the node's `policies.telemetry.duration` opt-in.
//!
//! Duration is grouped by outcome so operators can distinguish compute time
//! from error-path time without defining separate instruments.
//!
//! This API complements, but is not required for, the engine's automatic
//! per-message flow_metric. [`ComputeDuration::timed`] provides the
//! outcome split for `processor.compute.duration`,
//! while the engine's `Instant`-marker timing on the EffectHandler captures
//! total wall-clock compute between sends for flow_metrics without processor
//! cooperation.
//!
//! The closure-based API structurally prevents timing from spanning
//! `.await` points.
use std::cell::RefCell;
use crate::Interests;
use otel_arrow_dfe_telemetry::common_attributes::{Outcome, OutcomeAttributes};
use otel_arrow_dfe_telemetry::instrument::{HistogramNormal, Timer};
use otel_arrow_dfe_telemetry::metrics::MeasurementMetricSet;
use otel_arrow_dfe_telemetry::reporter::MetricsReporter;
use otel_arrow_dfe_telemetry_macros::metric_set;
use crate::context::PipelineContext;
/// Metric set containing processor compute duration grouped by outcome.
#[metric_set(
name = "processor.compute",
measurement_attributes = OutcomeAttributes
)]
#[derive(Debug, Default, Clone)]
pub struct ComputeDurationMetrics {
/// Wall-clock duration of processor-local synchronous compute.
#[metric(unit = "s")]
pub duration: HistogramNormal,
}
/// Wrapper providing interests-gated duration recording and reporting.
pub struct ComputeDuration {
metrics: MeasurementMetricSet<ComputeDurationMetrics>,
/// Accumulator for successful durations.
/// Uses `RefCell` for interior mutability so `timed` can take `&self`,
/// allowing callers to hold shared borrows of sibling fields in the
/// closure.
acc_success: RefCell<HistogramNormal>,
/// Accumulator for failed durations.
acc_failed: RefCell<HistogramNormal>,
}
impl ComputeDuration {
/// Register a new compute-duration metric set on the given pipeline.
#[must_use]
pub fn new(pipeline_ctx: &PipelineContext) -> Self {
Self {
metrics: ComputeDurationMetrics::register(pipeline_ctx),
acc_success: RefCell::new(HistogramNormal::default()),
acc_failed: RefCell::new(HistogramNormal::default()),
}
}
/// Time a synchronous, fallible closure for the process-duration outcome
/// split if interests includes `NODE_LOCAL_DURATION`, otherwise just call
/// `f` directly.
///
/// The elapsed time is recorded into the `success` or `failed`
/// accumulator based on the closure's `Result` outcome. This feeds only
/// the `processor.compute.duration` metric; flow_metric
/// participation is handled separately by the engine's `Instant`-marker
/// timing on the EffectHandler.
///
/// The closure-based API structurally prevents the timer from
/// being held across `.await` -- the closure is `FnOnce`, not
/// async, so the compiler enforces that only synchronous work is
/// measured.
#[inline]
pub fn timed<T, E>(
&self,
interests: Interests,
f: impl FnOnce() -> Result<T, E>,
) -> Result<T, E> {
if interests.contains(Interests::NODE_LOCAL_DURATION) {
let timer = Timer::start();
let result = f();
let elapsed_seconds = timer.elapsed_nanos() / 1e9;
let acc = if result.is_ok() {
&self.acc_success
} else {
&self.acc_failed
};
acc.borrow_mut().record(elapsed_seconds);
result
} else {
f()
}
}
/// Report accumulated duration metrics to the collector.
///
/// Drains both accumulators into the metric set, then reports
/// and resets as usual.
pub fn report(&mut self, reporter: &mut MetricsReporter) {
let success = self.acc_success.replace(HistogramNormal::default());
self.metrics
.with(OutcomeAttributes {
outcome: Outcome::Success,
})
.duration
.merge(success);
let failed = self.acc_failed.replace(HistogramNormal::default());
self.metrics
.with(OutcomeAttributes {
outcome: Outcome::Failure,
})
.duration
.merge(failed);
let _ = reporter.report_measurement(&mut self.metrics);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::testing::test_pipeline_ctx;
use otel_arrow_dfe_telemetry::reporter::MetricsReporter;
/// Scenario: timed processor compute succeeds twice and fails once.
/// Guarantees: duration observations are accumulated in the corresponding outcome buckets.
#[test]
fn timed_splits_by_outcome() {
let (ctx, _) = test_pipeline_ctx();
let cd = ComputeDuration::new(&ctx);
let active = Interests::NODE_LOCAL_DURATION;
// Two Ok results and one Err.
let _ = cd.timed(active, || Ok::<_, &str>(std::hint::black_box(42)));
let _ = cd.timed(active, || Ok::<_, &str>(std::hint::black_box(43)));
let _ = cd.timed(active, || Err::<i32, _>("fail"));
let (success_count, _, success_min, _) = cd.acc_success.borrow().get().summary();
assert_eq!(success_count, 2);
assert!(success_min >= 0.0);
let (failed_count, _, failed_min, _) = cd.acc_failed.borrow().get().summary();
assert_eq!(failed_count, 1);
assert!(failed_min >= 0.0);
}
/// Scenario: processor compute timing runs without the node-duration interest.
/// Guarantees: the closure executes without recording any duration observations.
#[test]
fn timed_noop_when_disabled() {
let (ctx, _) = test_pipeline_ctx();
let cd = ComputeDuration::new(&ctx);
let _ = cd.timed(Interests::empty(), || Ok::<_, &str>(1));
let _ = cd.timed(Interests::empty(), || Err::<i32, _>("fail"));
assert_eq!(cd.acc_success.borrow().get().count(), 0);
assert_eq!(cd.acc_failed.borrow().get().count(), 0);
}
/// Scenario: successful and failed processor compute observations are reported.
/// Guarantees: one seconds-based processor.compute.duration histogram is emitted with bounded outcome dimensions.
#[test]
fn report_emits_expected_metric_names() {
let (ctx, _) = test_pipeline_ctx();
let mut cd = ComputeDuration::new(&ctx);
let active = Interests::NODE_LOCAL_DURATION;
let _ = cd.timed(active, || Ok::<_, &str>(1));
let _ = cd.timed(active, || Err::<i32, _>("fail"));
let (rx, mut reporter) = MetricsReporter::create_new_and_receiver(4);
cd.report(&mut reporter);
let mut snapshots: Vec<_> = rx.try_iter().collect();
snapshots.sort_by_key(|snapshot| {
snapshot
.measurement_attribute_value("outcome")
.unwrap_or_default()
});
assert_eq!(snapshots.len(), 2);
assert_eq!(
snapshots
.iter()
.map(|snapshot| snapshot.measurement_attribute_value("outcome"))
.collect::<Vec<_>>(),
vec![Some("failure"), Some("success")]
);
for snapshot in snapshots {
assert_eq!(snapshot.descriptor().name, "processor.compute");
assert_eq!(snapshot.descriptor().metrics.len(), 1);
assert_eq!(snapshot.descriptor().metrics[0].name, "duration");
assert_eq!(snapshot.descriptor().metrics[0].unit, "s");
}
}
}