Skip to main content

opendal_layer_otelmetrics/
lib.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18#![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/// `OtelMetricsLayer` records OpenDAL operation and HTTP fetch metrics with
32/// [opentelemetry::metrics].
33///
34/// This layer records operation metrics from OpenDAL API calls and HTTP metrics
35/// from requests made through OpenDAL's HTTP fetcher.
36///
37/// # Examples
38///
39/// ```no_run
40/// # use opendal_core::services;
41/// # use opendal_core::Operator;
42/// # use opendal_core::Result;
43/// # use opendal_layer_otelmetrics::OtelMetricsLayer;
44/// #
45/// # fn main() -> Result<()> {
46/// let meter = opentelemetry::global::meter("opendal");
47/// let _ = Operator::new(services::Memory::default())?
48///     .layer(OtelMetricsLayer::builder().register(&meter));
49/// # Ok(())
50/// # }
51/// ```
52#[derive(Clone, Debug)]
53pub struct OtelMetricsLayer {
54    interceptor: OtelMetricsInterceptor,
55}
56
57impl OtelMetricsLayer {
58    /// Create a [`OtelMetricsLayerBuilder`] to set the configuration of metrics.
59    ///
60    /// # Examples
61    ///
62    /// ```no_run
63    /// # use opendal_core::services;
64    /// # use opendal_core::Operator;
65    /// # use opendal_core::Result;
66    /// # use opendal_layer_otelmetrics::OtelMetricsLayer;
67    /// #
68    /// # fn main() -> Result<()> {
69    /// let meter = opentelemetry::global::meter("opendal");
70    /// let op = Operator::new(services::Memory::default())?
71    ///     .layer(OtelMetricsLayer::builder().register(&meter));
72    /// # Ok(())
73    /// # }
74    /// ```
75    pub fn builder() -> OtelMetricsLayerBuilder {
76        OtelMetricsLayerBuilder::default()
77    }
78}
79
80/// [`OtelMetricsLayerBuilder`] is a config builder to build a [`OtelMetricsLayer`].
81pub 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    /// Set boundaries for bytes histograms, including operation bytes and HTTP body sizes.
105    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    /// Set boundaries for bytes rate histograms, including operation and HTTP body rates.
113    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    /// Set boundaries for entries related histogram like `operation_entries`.
121    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    /// Set boundaries for entries rate related histogram like `operation_entries_rate`.
129    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    /// Set boundaries for duration histograms, including operation and HTTP request/response durations.
137    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    /// Set boundaries for ttfb related histogram like `operation_ttfb_seconds`.
145    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    /// Register the metrics and return a [`OtelMetricsLayer`].
153    ///
154    /// # Examples
155    ///
156    /// ```no_run
157    /// # use opendal_core::services;
158    /// # use opendal_core::Operator;
159    /// # use opendal_core::Result;
160    /// # use opendal_layer_otelmetrics::OtelMetricsLayer;
161    /// #
162    /// # fn main() -> Result<()> {
163    /// let meter = opentelemetry::global::meter("opendal");
164    /// let op = Operator::new(services::Memory::default())?
165    ///     .layer(OtelMetricsLayer::builder().register(&meter));
166    /// # Ok(())
167    /// # }
168    /// ```
169    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    // Operation and HTTP metrics share the same registered OpenTelemetry instruments.
342    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}