opendal_layer_otelmetrics/
lib.rs1#![doc = include_str!("../README.md")]
19#![cfg_attr(docsrs, feature(doc_cfg))]
20#![cfg_attr(docsrs, doc(auto_cfg))]
21#![deny(missing_docs)]
22use opendal_core::OperationContext;
23use opendal_core::raw::*;
24use opendal_layer_observe_metrics_common as observe;
25use opentelemetry::KeyValue;
26use opentelemetry::metrics::Counter;
27use opentelemetry::metrics::Histogram;
28use opentelemetry::metrics::Meter;
29use opentelemetry::metrics::UpDownCounter;
30
31#[derive(Clone, Debug)]
53pub struct OtelMetricsLayer {
54 interceptor: OtelMetricsInterceptor,
55}
56
57impl OtelMetricsLayer {
58 pub fn builder() -> OtelMetricsLayerBuilder {
76 OtelMetricsLayerBuilder::default()
77 }
78}
79
80pub struct OtelMetricsLayerBuilder {
82 bytes_boundaries: Vec<f64>,
83 bytes_rate_boundaries: Vec<f64>,
84 entries_boundaries: Vec<f64>,
85 entries_rate_boundaries: Vec<f64>,
86 duration_seconds_boundaries: Vec<f64>,
87 ttfb_boundaries: Vec<f64>,
88}
89
90impl Default for OtelMetricsLayerBuilder {
91 fn default() -> Self {
92 Self {
93 bytes_boundaries: observe::DEFAULT_BYTES_BUCKETS.to_vec(),
94 bytes_rate_boundaries: observe::DEFAULT_BYTES_RATE_BUCKETS.to_vec(),
95 entries_boundaries: observe::DEFAULT_ENTRIES_BUCKETS.to_vec(),
96 entries_rate_boundaries: observe::DEFAULT_ENTRIES_RATE_BUCKETS.to_vec(),
97 duration_seconds_boundaries: observe::DEFAULT_DURATION_SECONDS_BUCKETS.to_vec(),
98 ttfb_boundaries: observe::DEFAULT_TTFB_BUCKETS.to_vec(),
99 }
100 }
101}
102
103impl OtelMetricsLayerBuilder {
104 pub fn bytes_boundaries(mut self, boundaries: Vec<f64>) -> Self {
106 if !boundaries.is_empty() {
107 self.bytes_boundaries = boundaries;
108 }
109 self
110 }
111
112 pub fn bytes_rate_boundaries(mut self, boundaries: Vec<f64>) -> Self {
114 if !boundaries.is_empty() {
115 self.bytes_rate_boundaries = boundaries;
116 }
117 self
118 }
119
120 pub fn entries_boundaries(mut self, boundaries: Vec<f64>) -> Self {
122 if !boundaries.is_empty() {
123 self.entries_boundaries = boundaries;
124 }
125 self
126 }
127
128 pub fn entries_rate_boundaries(mut self, boundaries: Vec<f64>) -> Self {
130 if !boundaries.is_empty() {
131 self.entries_rate_boundaries = boundaries;
132 }
133 self
134 }
135
136 pub fn duration_seconds_boundaries(mut self, boundaries: Vec<f64>) -> Self {
138 if !boundaries.is_empty() {
139 self.duration_seconds_boundaries = boundaries;
140 }
141 self
142 }
143
144 pub fn ttfb_boundaries(mut self, boundaries: Vec<f64>) -> Self {
146 if !boundaries.is_empty() {
147 self.ttfb_boundaries = boundaries;
148 }
149 self
150 }
151
152 pub fn register(self, meter: &Meter) -> OtelMetricsLayer {
170 let operation_bytes = {
171 let metric = observe::MetricValue::OperationBytes(0);
172 register_u64_histogram_meter(
173 meter,
174 "opendal.operation.bytes",
175 metric,
176 self.bytes_boundaries.clone(),
177 )
178 };
179 let operation_bytes_rate = {
180 let metric = observe::MetricValue::OperationBytesRate(0.0);
181 register_f64_histogram_meter(
182 meter,
183 "opendal.operation.bytes_rate",
184 metric,
185 self.bytes_rate_boundaries.clone(),
186 )
187 };
188 let operation_entries = {
189 let metric = observe::MetricValue::OperationEntries(0);
190 register_u64_histogram_meter(
191 meter,
192 "opendal.operation.entries",
193 metric,
194 self.entries_boundaries.clone(),
195 )
196 };
197 let operation_entries_rate = {
198 let metric = observe::MetricValue::OperationEntriesRate(0.0);
199 register_f64_histogram_meter(
200 meter,
201 "opendal.operation.entries_rate",
202 metric,
203 self.entries_rate_boundaries.clone(),
204 )
205 };
206 let operation_duration_seconds = {
207 let metric = observe::MetricValue::OperationDurationSeconds(Duration::default());
208 register_f64_histogram_meter(
209 meter,
210 "opendal.operation.duration",
211 metric,
212 self.duration_seconds_boundaries.clone(),
213 )
214 };
215 let operation_errors_total = {
216 let metric = observe::MetricValue::OperationErrorsTotal;
217 meter
218 .u64_counter("opendal.operation.errors")
219 .with_description(metric.help())
220 .build()
221 };
222 let operation_executing = {
223 let metric = observe::MetricValue::OperationExecuting(0);
224 meter
225 .i64_up_down_counter("opendal.operation.executing")
226 .with_description(metric.help())
227 .build()
228 };
229 let operation_ttfb_seconds = {
230 let metric = observe::MetricValue::OperationTtfbSeconds(Duration::default());
231 register_f64_histogram_meter(
232 meter,
233 "opendal.operation.ttfb",
234 metric,
235 self.ttfb_boundaries.clone(),
236 )
237 };
238
239 let http_executing = {
240 let metric = observe::MetricValue::HttpExecuting(0);
241 meter
242 .i64_up_down_counter("opendal.http.executing")
243 .with_description(metric.help())
244 .build()
245 };
246 let http_request_bytes = {
247 let metric = observe::MetricValue::HttpRequestBytes(0);
248 register_u64_histogram_meter(
249 meter,
250 "opendal.http.request.bytes",
251 metric,
252 self.bytes_boundaries.clone(),
253 )
254 };
255 let http_request_bytes_rate = {
256 let metric = observe::MetricValue::HttpRequestBytesRate(0.0);
257 register_f64_histogram_meter(
258 meter,
259 "opendal.http.request.bytes_rate",
260 metric,
261 self.bytes_rate_boundaries.clone(),
262 )
263 };
264 let http_request_duration_seconds = {
265 let metric = observe::MetricValue::HttpRequestDurationSeconds(Duration::default());
266 register_f64_histogram_meter(
267 meter,
268 "opendal.http.request.duration",
269 metric,
270 self.duration_seconds_boundaries.clone(),
271 )
272 };
273 let http_response_bytes = {
274 let metric = observe::MetricValue::HttpResponseBytes(0);
275 register_u64_histogram_meter(
276 meter,
277 "opendal.http.response.bytes",
278 metric,
279 self.bytes_boundaries.clone(),
280 )
281 };
282 let http_response_bytes_rate = {
283 let metric = observe::MetricValue::HttpResponseBytesRate(0.0);
284 register_f64_histogram_meter(
285 meter,
286 "opendal.http.response.bytes_rate",
287 metric,
288 self.bytes_rate_boundaries.clone(),
289 )
290 };
291 let http_response_duration_seconds = {
292 let metric = observe::MetricValue::HttpResponseDurationSeconds(Duration::default());
293 register_f64_histogram_meter(
294 meter,
295 "opendal.http.response.duration",
296 metric,
297 self.duration_seconds_boundaries.clone(),
298 )
299 };
300 let http_connection_errors_total = {
301 let metric = observe::MetricValue::HttpConnectionErrorsTotal;
302 meter
303 .u64_counter("opendal.http.connection_errors")
304 .with_description(metric.help())
305 .build()
306 };
307 let http_status_errors_total = {
308 let metric = observe::MetricValue::HttpStatusErrorsTotal;
309 meter
310 .u64_counter("opendal.http.status_errors")
311 .with_description(metric.help())
312 .build()
313 };
314
315 OtelMetricsLayer {
316 interceptor: OtelMetricsInterceptor {
317 operation_bytes,
318 operation_bytes_rate,
319 operation_entries,
320 operation_entries_rate,
321 operation_duration_seconds,
322 operation_errors_total,
323 operation_executing,
324 operation_ttfb_seconds,
325
326 http_executing,
327 http_request_bytes,
328 http_request_bytes_rate,
329 http_request_duration_seconds,
330 http_response_bytes,
331 http_response_bytes_rate,
332 http_response_duration_seconds,
333 http_connection_errors_total,
334 http_status_errors_total,
335 },
336 }
337 }
338}
339
340impl Layer for OtelMetricsLayer {
341 fn apply_service(&self, inner: Servicer) -> Servicer {
343 observe::MetricsLayer::new(self.interceptor.clone()).apply_service(inner)
344 }
345
346 fn apply_context(&self, srv: Servicer, inner: OperationContext) -> OperationContext {
347 observe::MetricsLayer::new(self.interceptor.clone()).apply_context(srv, inner)
348 }
349}
350
351#[doc(hidden)]
352#[derive(Clone, Debug)]
353pub struct OtelMetricsInterceptor {
354 operation_bytes: Histogram<u64>,
355 operation_bytes_rate: Histogram<f64>,
356 operation_entries: Histogram<u64>,
357 operation_entries_rate: Histogram<f64>,
358 operation_duration_seconds: Histogram<f64>,
359 operation_errors_total: Counter<u64>,
360 operation_executing: UpDownCounter<i64>,
361 operation_ttfb_seconds: Histogram<f64>,
362
363 http_executing: UpDownCounter<i64>,
364 http_request_bytes: Histogram<u64>,
365 http_request_bytes_rate: Histogram<f64>,
366 http_request_duration_seconds: Histogram<f64>,
367 http_response_bytes: Histogram<u64>,
368 http_response_bytes_rate: Histogram<f64>,
369 http_response_duration_seconds: Histogram<f64>,
370 http_connection_errors_total: Counter<u64>,
371 http_status_errors_total: Counter<u64>,
372}
373
374impl observe::MetricsIntercept for OtelMetricsInterceptor {
375 fn observe(&self, labels: observe::MetricLabels, value: observe::MetricValue) {
376 let attributes = self.create_attributes(labels);
377
378 match value {
379 observe::MetricValue::OperationBytes(v) => self.operation_bytes.record(v, &attributes),
380 observe::MetricValue::OperationBytesRate(v) => {
381 self.operation_bytes_rate.record(v, &attributes)
382 }
383 observe::MetricValue::OperationEntries(v) => {
384 self.operation_entries.record(v, &attributes)
385 }
386 observe::MetricValue::OperationEntriesRate(v) => {
387 self.operation_entries_rate.record(v, &attributes)
388 }
389 observe::MetricValue::OperationDurationSeconds(v) => self
390 .operation_duration_seconds
391 .record(v.as_secs_f64(), &attributes),
392 observe::MetricValue::OperationErrorsTotal => {
393 self.operation_errors_total.add(1, &attributes)
394 }
395 observe::MetricValue::OperationExecuting(v) => {
396 self.operation_executing.add(v as i64, &attributes)
397 }
398 observe::MetricValue::OperationTtfbSeconds(v) => self
399 .operation_ttfb_seconds
400 .record(v.as_secs_f64(), &attributes),
401
402 observe::MetricValue::HttpExecuting(v) => {
403 self.http_executing.add(v as i64, &attributes)
404 }
405 observe::MetricValue::HttpRequestBytes(v) => {
406 self.http_request_bytes.record(v, &attributes)
407 }
408 observe::MetricValue::HttpRequestBytesRate(v) => {
409 self.http_request_bytes_rate.record(v, &attributes)
410 }
411 observe::MetricValue::HttpRequestDurationSeconds(v) => self
412 .http_request_duration_seconds
413 .record(v.as_secs_f64(), &attributes),
414 observe::MetricValue::HttpResponseBytes(v) => {
415 self.http_response_bytes.record(v, &attributes)
416 }
417 observe::MetricValue::HttpResponseBytesRate(v) => {
418 self.http_response_bytes_rate.record(v, &attributes)
419 }
420 observe::MetricValue::HttpResponseDurationSeconds(v) => self
421 .http_response_duration_seconds
422 .record(v.as_secs_f64(), &attributes),
423 observe::MetricValue::HttpConnectionErrorsTotal => {
424 self.http_connection_errors_total.add(1, &attributes)
425 }
426 observe::MetricValue::HttpStatusErrorsTotal => {
427 self.http_status_errors_total.add(1, &attributes)
428 }
429 _ => {}
430 }
431 }
432}
433
434impl OtelMetricsInterceptor {
435 fn create_attributes(&self, attrs: observe::MetricLabels) -> Vec<KeyValue> {
436 let mut attributes = Vec::with_capacity(6);
437
438 attributes.extend([
439 KeyValue::new(observe::LABEL_SCHEME, attrs.scheme),
440 KeyValue::new(observe::LABEL_NAMESPACE, attrs.namespace),
441 KeyValue::new(observe::LABEL_ROOT, attrs.root),
442 KeyValue::new(observe::LABEL_OPERATION, attrs.operation),
443 ]);
444
445 if let Some(error) = attrs.error {
446 attributes.push(KeyValue::new(observe::LABEL_ERROR, error.into_static()));
447 }
448
449 if let Some(status_code) = attrs.status_code {
450 attributes.push(KeyValue::new(
451 observe::LABEL_STATUS_CODE,
452 status_code.as_u16() as i64,
453 ));
454 }
455
456 if let Some(service_operation) = attrs.service_operation {
457 attributes.push(KeyValue::new(
458 observe::LABEL_SERVICE_OPERATION,
459 service_operation,
460 ));
461 }
462
463 attributes
464 }
465}
466
467fn register_u64_histogram_meter(
468 meter: &Meter,
469 name: &'static str,
470 metric: observe::MetricValue,
471 boundaries: Vec<f64>,
472) -> Histogram<u64> {
473 let (_name, unit) = metric.name_with_unit();
474 let description = metric.help();
475
476 let builder = meter
477 .u64_histogram(name)
478 .with_description(description)
479 .with_boundaries(boundaries);
480
481 if let Some(unit) = unit {
482 builder.with_unit(unit).build()
483 } else {
484 builder.build()
485 }
486}
487
488fn register_f64_histogram_meter(
489 meter: &Meter,
490 name: &'static str,
491 metric: observe::MetricValue,
492 boundaries: Vec<f64>,
493) -> Histogram<f64> {
494 let (_name, unit) = metric.name_with_unit();
495 let description = metric.help();
496
497 let builder = meter
498 .f64_histogram(name)
499 .with_description(description)
500 .with_boundaries(boundaries);
501
502 if let Some(unit) = unit {
503 builder.with_unit(unit).build()
504 } else {
505 builder.build()
506 }
507}