Skip to main content

opentelemetry_sdk/metrics/
periodic_reader.rs

1use std::{
2    env, fmt,
3    sync::{
4        mpsc::{self, Receiver, Sender},
5        Arc, Mutex, Weak,
6    },
7    thread,
8    time::{Duration, Instant},
9};
10
11use opentelemetry::{otel_debug, otel_error, otel_info, otel_warn, Context};
12
13use crate::{
14    error::{OTelSdkError, OTelSdkResult},
15    metrics::{exporter::PushMetricExporter, reader::SdkProducer},
16    Resource,
17};
18
19use super::{
20    data::ResourceMetrics, instrument::InstrumentKind, pipeline::Pipeline, reader::MetricReader,
21    Temporality,
22};
23
24/// Environment variable for configuring the delay interval (in milliseconds)
25/// between two consecutive exports for [PeriodicReader].
26pub const OTEL_METRIC_EXPORT_INTERVAL: &str = "OTEL_METRIC_EXPORT_INTERVAL";
27
28/// Default delay interval between two consecutive exports for [PeriodicReader].
29pub const OTEL_METRIC_EXPORT_INTERVAL_DEFAULT: Duration = Duration::from_secs(60);
30
31/// Configuration options for [PeriodicReader].
32#[derive(Debug)]
33pub struct PeriodicReaderBuilder<E> {
34    interval: Duration,
35    exporter: E,
36}
37
38impl<E> PeriodicReaderBuilder<E>
39where
40    E: PushMetricExporter,
41{
42    fn new(exporter: E) -> Self {
43        let interval = env::var(OTEL_METRIC_EXPORT_INTERVAL)
44            .ok()
45            .and_then(|v| v.parse().map(Duration::from_millis).ok())
46            .unwrap_or(OTEL_METRIC_EXPORT_INTERVAL_DEFAULT);
47
48        PeriodicReaderBuilder { interval, exporter }
49    }
50
51    /// Configures the intervening time between exports for a [PeriodicReader].
52    ///
53    /// This option overrides any value set for the `OTEL_METRIC_EXPORT_INTERVAL`
54    /// environment variable.
55    ///
56    /// If this option is not used or `interval` is equal to zero, 60 seconds is
57    /// used as the default.
58    pub fn with_interval(mut self, interval: Duration) -> Self {
59        if !interval.is_zero() {
60            self.interval = interval;
61        }
62        self
63    }
64
65    /// Create a [PeriodicReader] with the given config.
66    pub fn build(self) -> PeriodicReader<E> {
67        PeriodicReader::new(self.exporter, self.interval)
68    }
69}
70
71/// A `MetricReader` that periodically collects and exports metrics at a configurable interval.
72///
73/// By default, [`PeriodicReader`] collects and exports metrics every **60 seconds**.
74/// The time taken for export is **not** included in the interval. Use [`PeriodicReaderBuilder`]
75/// to customize the interval.
76///
77/// [`PeriodicReader`] spawns a background thread to handle metric collection and export.
78/// This thread remains active until [`shutdown()`] is called.
79///
80/// ## Collection Process
81/// "Collection" refers to gathering aggregated metrics from the SDK's internal storage.
82/// During this phase, callbacks from observable instruments are also triggered.
83///
84/// [`PeriodicReader`] does **not** enforce a timeout for collection. If an
85/// observable callback takes too long, it may delay the next collection cycle.
86/// If a callback never returns, it **will stall** all metric collection (and exports)
87/// indefinitely.
88///
89/// ## Exporter Compatibility
90/// When used with the [`OTLP Exporter`](https://docs.rs/opentelemetry-otlp), the following
91/// transport options are supported:
92///
93/// - **`grpc-tonic`**: Requires [`MeterProvider`] to be initialized within a `tokio` runtime.
94/// - **`reqwest-blocking-client`**: Works with both a standard (`main`) function and `tokio::main`.
95///
96/// Async HTTP clients such as `reqwest-client` and `hyper-client` are not
97/// supported by this default reader. The OTLP HTTP exporter chooses its default
98/// HTTP client from enabled crate features and cannot tell which reader will
99/// drive it. If your dependency graph enables async HTTP client features, either
100/// pass an explicit blocking client for this reader or use the experimental
101/// async-runtime periodic reader.
102///
103/// [`PeriodicReader`] does **not** enforce a timeout for exports either. Instead,
104/// the configured exporter is responsible for enforcing timeouts. If an export operation
105/// never returns, [`PeriodicReader`] will **stop exporting new metrics**, stalling
106/// metric collection.
107///
108/// ## Manual Export & Shutdown
109/// Users can manually trigger an export via [`force_flush()`]. Calling [`shutdown()`]
110/// exports any remaining metrics and should be done before application exit to ensure
111/// all data is sent.
112///
113/// **Warning**: If using **tokio’s current-thread runtime**, calling [`shutdown()`]
114/// from the main thread may cause a deadlock. To prevent this, call [`shutdown()`]
115/// from a separate thread or use tokio's `spawn_blocking`.
116///
117/// [`PeriodicReader`]: crate::metrics::PeriodicReader
118/// [`PeriodicReaderBuilder`]: crate::metrics::PeriodicReaderBuilder
119/// [`MeterProvider`]: crate::metrics::SdkMeterProvider
120/// [`shutdown()`]: crate::metrics::SdkMeterProvider::shutdown
121/// [`force_flush()`]: crate::metrics::SdkMeterProvider::force_flush
122///
123/// # Example
124///
125/// ```no_run
126/// use opentelemetry_sdk::metrics::PeriodicReader;
127/// # fn example<E>(get_exporter: impl Fn() -> E)
128/// # where
129/// #     E: opentelemetry_sdk::metrics::exporter::PushMetricExporter,
130/// # {
131///
132/// let exporter = get_exporter(); // set up a push exporter
133///
134/// let reader = PeriodicReader::builder(exporter).build();
135/// # drop(reader);
136/// # }
137/// ```
138pub struct PeriodicReader<E: PushMetricExporter> {
139    inner: Arc<PeriodicReaderInner<E>>,
140}
141
142impl<E: PushMetricExporter> Clone for PeriodicReader<E> {
143    fn clone(&self) -> Self {
144        Self {
145            inner: Arc::clone(&self.inner),
146        }
147    }
148}
149
150impl<E: PushMetricExporter> PeriodicReader<E> {
151    /// Configuration options for a periodic reader with own thread
152    pub fn builder(exporter: E) -> PeriodicReaderBuilder<E> {
153        PeriodicReaderBuilder::new(exporter)
154    }
155
156    fn new(exporter: E, interval: Duration) -> Self {
157        let (message_sender, message_receiver): (Sender<Message>, Receiver<Message>) =
158            mpsc::channel();
159        let exporter_arc = Arc::new(exporter);
160        let reader = PeriodicReader {
161            inner: Arc::new(PeriodicReaderInner {
162                message_sender,
163                producer: Mutex::new(None),
164                exporter: exporter_arc.clone(),
165            }),
166        };
167        let cloned_reader = reader.clone();
168
169        let mut rm = ResourceMetrics {
170            resource: Resource::empty(),
171            scope_metrics: Vec::new(),
172        };
173
174        let result_thread_creation = thread::Builder::new()
175            .name("OpenTelemetry.Metrics.PeriodicReader".to_string())
176            .spawn(move || {
177                let _suppress_guard = Context::enter_telemetry_suppressed_scope();
178                let mut interval_start = Instant::now();
179                let mut remaining_interval = interval;
180                otel_debug!(
181                    name: "PeriodReaderThreadStarted",
182                    interval_in_millisecs = interval.as_millis(),
183                );
184                loop {
185                    otel_debug!(
186                        name: "PeriodReaderThreadLoopAlive", message = "Next export will happen after interval, unless flush or shutdown is triggered.", interval_in_millisecs = remaining_interval.as_millis()
187                    );
188                    match message_receiver.recv_timeout(remaining_interval) {
189                        Ok(Message::Flush(response_sender)) => {
190                            otel_debug!(
191                                name: "PeriodReaderThreadExportingDueToFlush"
192                            );
193                            let export_result = cloned_reader.collect_and_export(&mut rm);
194                            otel_debug!(
195                                name: "PeriodReaderInvokedExport",
196                                export_result = format!("{:?}", export_result)
197                            );
198
199                            // If response_sender is disconnected, we can't send
200                            // the result back. This occurs when the thread that
201                            // initiated flush gave up due to timeout.
202                            // Gracefully handle that with internal logs. The
203                            // internal errors are of Info level, as this is
204                            // useful for user to know whether the flush was
205                            // successful or not, when flush() itself merely
206                            // tells that it timed out.
207
208                            if export_result.is_err() {
209                                if response_sender.send(false).is_err() {
210                                    otel_debug!(
211                                        name: "PeriodReader.Flush.ResponseSendError",
212                                        message = "PeriodicReader's flush has failed, but unable to send this info back to caller.
213                                        This occurs when the caller has timed out waiting for the response. If you see this occuring frequently, consider increasing the flush timeout."
214                                    );
215                                }
216                            } else if response_sender.send(true).is_err() {
217                                otel_debug!(
218                                    name: "PeriodReader.Flush.ResponseSendError",
219                                    message = "PeriodicReader's flush has completed successfully, but unable to send this info back to caller.
220                                    This occurs when the caller has timed out waiting for the response. If you see this occuring frequently, consider increasing the flush timeout."
221                                );
222                            }
223
224                            // Adjust the remaining interval after the flush
225                            let elapsed = interval_start.elapsed();
226                            if elapsed < interval {
227                                remaining_interval = interval - elapsed;
228                                otel_debug!(
229                                    name: "PeriodReaderThreadAdjustingRemainingIntervalAfterFlush",
230                                    remaining_interval = remaining_interval.as_secs()
231                                );
232                            } else {
233                                otel_debug!(
234                                    name: "PeriodReaderThreadAdjustingExportAfterFlush",
235                                );
236                                // Reset the interval if the flush finishes after the expected export time
237                                // effectively missing the normal export.
238                                // Should we attempt to do the missed export immediately?
239                                // Or do the next export at the next interval?
240                                // Currently this attempts the next export immediately.
241                                // i.e calling Flush can affect the regularity.
242                                interval_start = Instant::now();
243                                remaining_interval = Duration::ZERO;
244                            }
245                        }
246                        Ok(Message::Shutdown(response_sender)) => {
247                            // Perform final export and break out of loop and exit the thread
248                            otel_debug!(name: "PeriodReaderThreadExportingDueToShutdown");
249                            let export_result = cloned_reader.collect_and_export(&mut rm);
250                            otel_debug!(
251                                name: "PeriodReaderInvokedExport",
252                                export_result = format!("{:?}", export_result)
253                            );
254                            let shutdown_result = exporter_arc.shutdown();
255                            otel_debug!(
256                                name: "PeriodReaderInvokedExporterShutdown",
257                                shutdown_result = format!("{:?}", shutdown_result)
258                            );
259
260                            // If response_sender is disconnected, we can't send
261                            // the result back. This occurs when the thread that
262                            // initiated shutdown gave up due to timeout.
263                            // Gracefully handle that with internal logs and
264                            // continue with shutdown (i.e exit thread) The
265                            // internal errors are of Info level, as this is
266                            // useful for user to know whether the shutdown was
267                            // successful or not, when shutdown() itself merely
268                            // tells that it timed out.
269                            if export_result.is_err() || shutdown_result.is_err() {
270                                if response_sender.send(false).is_err() {
271                                    otel_info!(
272                                        name: "PeriodReaderThreadShutdown.ResponseSendError",
273                                        message = "PeriodicReader's shutdown has failed, but unable to send this info back to caller.
274                                        This occurs when the caller has timed out waiting for the response. If you see this occuring frequently, consider increasing the shutdown timeout."
275                                    );
276                                }
277                            } else if response_sender.send(true).is_err() {
278                                otel_debug!(
279                                    name: "PeriodReaderThreadShutdown.ResponseSendError",
280                                    message = "PeriodicReader completed its shutdown, but unable to send this info back to caller.
281                                    This occurs when the caller has timed out waiting for the response. If you see this occuring frequently, consider increasing the shutdown timeout."
282                                );
283                            }
284
285                            otel_debug!(
286                                name: "PeriodReaderThreadExiting",
287                                reason = "ShutdownRequested"
288                            );
289                            break;
290                        }
291                        Err(mpsc::RecvTimeoutError::Timeout) => {
292                            let export_start = Instant::now();
293                            otel_debug!(
294                                name: "PeriodReaderThreadExportingDueToTimer"
295                            );
296
297                            let export_result = cloned_reader.collect_and_export(&mut rm);
298                            otel_debug!(
299                                name: "PeriodReaderInvokedExport",
300                                export_result = format!("{:?}", export_result)
301                            );
302
303                            let time_taken_for_export = export_start.elapsed();
304                            if time_taken_for_export > interval {
305                                otel_debug!(
306                                    name: "PeriodReaderThreadExportTookLongerThanInterval"
307                                );
308                                // if export took longer than interval, do the
309                                // next export immediately.
310                                // Alternatively, we could skip the next export
311                                // and wait for the next interval.
312                                // Or enforce that export timeout is less than interval.
313                                // What is the desired behavior?
314                                interval_start = Instant::now();
315                                remaining_interval = Duration::ZERO;
316                            } else {
317                                remaining_interval = interval - time_taken_for_export;
318                                interval_start = Instant::now();
319                            }
320                        }
321                        Err(mpsc::RecvTimeoutError::Disconnected) => {
322                            // Channel disconnected, only thing to do is break
323                            // out (i.e exit the thread)
324                            otel_debug!(
325                                name: "PeriodReaderThreadExiting",
326                                reason = "MessageSenderDisconnected"
327                            );
328                            break;
329                        }
330                    }
331                }
332                otel_debug!(
333                    name: "PeriodReaderThreadStopped"
334                );
335            });
336
337        // TODO: Should we fail-fast here and bubble up the error to user?
338        #[allow(unused_variables)]
339        if let Err(e) = result_thread_creation {
340            otel_error!(
341                name: "PeriodReaderThreadStartError",
342                message = "Failed to start PeriodicReader thread. Metrics will not be exported.",
343                error = format!("{:?}", e)
344            );
345        }
346        reader
347    }
348
349    fn collect_and_export(&self, rm: &mut ResourceMetrics) -> OTelSdkResult {
350        self.inner.collect_and_export(rm)
351    }
352}
353
354impl<E: PushMetricExporter> fmt::Debug for PeriodicReader<E> {
355    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
356        f.debug_struct("PeriodicReader").finish()
357    }
358}
359
360struct PeriodicReaderInner<E: PushMetricExporter> {
361    exporter: Arc<E>,
362    message_sender: mpsc::Sender<Message>,
363    producer: Mutex<Option<Weak<dyn SdkProducer>>>,
364}
365
366impl<E: PushMetricExporter> PeriodicReaderInner<E> {
367    fn register_pipeline(&self, producer: Weak<dyn SdkProducer>) {
368        let mut inner = self.producer.lock().expect("lock poisoned");
369        *inner = Some(producer);
370    }
371
372    fn temporality(&self, _kind: InstrumentKind) -> Temporality {
373        self.exporter.temporality()
374    }
375
376    fn collect(&self, rm: &mut ResourceMetrics) -> OTelSdkResult {
377        let producer = self.producer.lock().expect("lock poisoned");
378        if let Some(p) = producer.as_ref() {
379            p.upgrade()
380                .ok_or(OTelSdkError::AlreadyShutdown)?
381                .produce(rm)?;
382            Ok(())
383        } else {
384            otel_warn!(
385            name: "PeriodReader.MeterProviderNotRegistered",
386            message = "PeriodicReader is not registered with MeterProvider. Metrics will not be collected. \
387                   This occurs when a periodic reader is created but not associated with a MeterProvider \
388                   by calling `.with_reader(reader)` on MeterProviderBuilder."
389            );
390            Err(OTelSdkError::InternalFailure(
391                "MeterProvider is not registered".into(),
392            ))
393        }
394    }
395
396    fn collect_and_export(&self, rm: &mut ResourceMetrics) -> OTelSdkResult {
397        let current_time = Instant::now();
398        let collect_result = self.collect(rm);
399        let time_taken_for_collect = current_time.elapsed();
400
401        #[allow(clippy::question_mark)]
402        if let Err(e) = collect_result {
403            otel_warn!(
404                name: "PeriodReaderCollectError",
405                error = format!("{:?}", e)
406            );
407            return Err(OTelSdkError::InternalFailure(e.to_string()));
408        }
409
410        if rm.scope_metrics.is_empty() {
411            otel_debug!(name: "NoMetricsCollected");
412            return Ok(());
413        }
414
415        let metrics_count = rm.scope_metrics.iter().fold(0, |count, scope_metrics| {
416            count + scope_metrics.metrics.len()
417        });
418        otel_debug!(name: "PeriodicReaderMetricsCollected", count = metrics_count, time_taken_in_millis = time_taken_for_collect.as_millis());
419
420        // Relying on futures executor to execute async call.
421        // TODO: Pass timeout to exporter
422        futures_executor::block_on(self.exporter.export(rm))
423    }
424
425    fn force_flush(&self) -> OTelSdkResult {
426        // TODO: Better message for this scenario.
427        // Flush and Shutdown called from 2 threads Flush check shutdown
428        // flag before shutdown thread sets it. Both threads attempt to send
429        // message to the same channel. Case1: Flush thread sends message first,
430        // shutdown thread sends message next. Flush would succeed, as
431        // background thread won't process shutdown message until flush
432        // triggered export is done. Case2: Shutdown thread sends message first,
433        // flush thread sends message next. Shutdown would succeed, as
434        // background thread would process shutdown message first. The
435        // background exits so it won't receive the flush message. ForceFlush
436        // returns Failure, but we could indicate specifically that shutdown has
437        // completed. TODO is to see if this message can be improved.
438
439        let (response_tx, response_rx) = mpsc::channel();
440        self.message_sender
441            .send(Message::Flush(response_tx))
442            .map_err(|e| OTelSdkError::InternalFailure(e.to_string()))?;
443
444        if let Ok(response) = response_rx.recv() {
445            // TODO: call exporter's force_flush method.
446            if response {
447                Ok(())
448            } else {
449                Err(OTelSdkError::InternalFailure("Failed to flush".into()))
450            }
451        } else {
452            Err(OTelSdkError::InternalFailure("Failed to flush".into()))
453        }
454    }
455
456    fn shutdown(&self) -> OTelSdkResult {
457        // TODO: See if this is better to be created upfront.
458        let (response_tx, response_rx) = mpsc::channel();
459        self.message_sender
460            .send(Message::Shutdown(response_tx))
461            .map_err(|e| OTelSdkError::InternalFailure(e.to_string()))?;
462
463        // TODO: Make this timeout configurable.
464        match response_rx.recv_timeout(Duration::from_secs(5)) {
465            Ok(response) => {
466                if response {
467                    Ok(())
468                } else {
469                    Err(OTelSdkError::InternalFailure("Failed to shutdown".into()))
470                }
471            }
472            Err(mpsc::RecvTimeoutError::Timeout) => {
473                Err(OTelSdkError::Timeout(Duration::from_secs(5)))
474            }
475            Err(mpsc::RecvTimeoutError::Disconnected) => {
476                Err(OTelSdkError::InternalFailure("Failed to shutdown".into()))
477            }
478        }
479    }
480}
481
482#[derive(Debug)]
483enum Message {
484    Flush(Sender<bool>),
485    Shutdown(Sender<bool>),
486}
487
488impl<E: PushMetricExporter> MetricReader for PeriodicReader<E> {
489    fn register_pipeline(&self, pipeline: Weak<Pipeline>) {
490        self.inner.register_pipeline(pipeline);
491    }
492
493    fn collect(&self, rm: &mut ResourceMetrics) -> OTelSdkResult {
494        self.inner.collect(rm)
495    }
496
497    fn force_flush(&self) -> OTelSdkResult {
498        self.inner.force_flush()
499    }
500
501    // TODO: Offer an async version of shutdown so users can await the shutdown
502    // completion, and avoid blocking the thread. The default shutdown on drop
503    // can still use blocking call. If user already explicitly called shutdown,
504    // drop won't call shutdown again.
505    fn shutdown_with_timeout(&self, _timeout: Duration) -> OTelSdkResult {
506        self.inner.shutdown()
507    }
508
509    /// To construct a [MetricReader][metric-reader] when setting up an SDK,
510    /// The output temporality (optional), a function of instrument kind.
511    /// This function SHOULD be obtained from the exporter.
512    ///
513    /// If not configured, the Cumulative temporality SHOULD be used.
514    ///
515    /// [metric-reader]: https://github.com/open-telemetry/opentelemetry-specification/blob/0a78571045ca1dca48621c9648ec3c832c3c541c/specification/metrics/sdk.md#metricreader
516    fn temporality(&self, kind: InstrumentKind) -> Temporality {
517        kind.temporality_preference(self.inner.temporality(kind))
518    }
519}
520
521#[cfg(all(test, feature = "testing"))]
522mod tests {
523    use super::PeriodicReader;
524    use crate::{
525        error::{OTelSdkError, OTelSdkResult},
526        metrics::{
527            data::ResourceMetrics, exporter::PushMetricExporter, reader::MetricReader,
528            InMemoryMetricExporter, SdkMeterProvider, Temporality,
529        },
530        Resource,
531    };
532    use opentelemetry::metrics::MeterProvider;
533    use std::{
534        sync::{
535            atomic::{AtomicBool, AtomicUsize, Ordering},
536            mpsc, Arc,
537        },
538        time::Duration,
539    };
540
541    // use below command to run all tests
542    // cargo test metrics::periodic_reader::tests --features=testing,spec_unstable_metrics_views -- --nocapture
543
544    #[derive(Debug, Clone)]
545    struct MetricExporterThatFailsOnlyOnFirst {
546        count: Arc<AtomicUsize>,
547    }
548
549    impl Default for MetricExporterThatFailsOnlyOnFirst {
550        fn default() -> Self {
551            MetricExporterThatFailsOnlyOnFirst {
552                count: Arc::new(AtomicUsize::new(0)),
553            }
554        }
555    }
556
557    impl MetricExporterThatFailsOnlyOnFirst {
558        fn get_count(&self) -> usize {
559            self.count.load(Ordering::Relaxed)
560        }
561    }
562
563    impl PushMetricExporter for MetricExporterThatFailsOnlyOnFirst {
564        async fn export(&self, _metrics: &ResourceMetrics) -> OTelSdkResult {
565            if self.count.fetch_add(1, Ordering::Relaxed) == 0 {
566                Err(OTelSdkError::InternalFailure("export failed".into()))
567            } else {
568                Ok(())
569            }
570        }
571
572        fn force_flush(&self) -> OTelSdkResult {
573            Ok(())
574        }
575
576        fn shutdown(&self) -> OTelSdkResult {
577            Ok(())
578        }
579
580        fn shutdown_with_timeout(&self, _timeout: Duration) -> OTelSdkResult {
581            Ok(())
582        }
583
584        fn temporality(&self) -> Temporality {
585            Temporality::Cumulative
586        }
587    }
588
589    #[derive(Debug, Clone, Default)]
590    struct MockMetricExporter {
591        is_shutdown: Arc<AtomicBool>,
592    }
593
594    impl PushMetricExporter for MockMetricExporter {
595        async fn export(&self, _metrics: &ResourceMetrics) -> OTelSdkResult {
596            Ok(())
597        }
598
599        fn force_flush(&self) -> OTelSdkResult {
600            Ok(())
601        }
602
603        fn shutdown(&self) -> OTelSdkResult {
604            self.shutdown_with_timeout(Duration::from_secs(5))
605        }
606
607        fn shutdown_with_timeout(&self, _timeout: Duration) -> OTelSdkResult {
608            self.is_shutdown.store(true, Ordering::Relaxed);
609            Ok(())
610        }
611
612        fn temporality(&self) -> Temporality {
613            Temporality::Cumulative
614        }
615    }
616
617    #[test]
618    fn collection_triggered_by_interval_multiple() {
619        // Arrange
620        let interval = std::time::Duration::from_millis(1);
621        let exporter = InMemoryMetricExporter::default();
622        let reader = PeriodicReader::builder(exporter.clone())
623            .with_interval(interval)
624            .build();
625        let i = Arc::new(AtomicUsize::new(0));
626        let i_clone = i.clone();
627
628        // Act
629        let meter_provider = SdkMeterProvider::builder().with_reader(reader).build();
630        let meter = meter_provider.meter("test");
631        let _counter = meter
632            .u64_observable_counter("testcounter")
633            .with_callback(move |_| {
634                i_clone.fetch_add(1, Ordering::Relaxed);
635            })
636            .build();
637
638        // Sleep for a duration 5X (plus liberal buffer to account for potential
639        // CI slowness) the interval to ensure multiple collection.
640        // Not a fan of such tests, but this seems to be the only way to test
641        // if periodic reader is doing its job.
642        // TODO: Decide if this should be ignored in CI
643        std::thread::sleep(interval * 5 * 20);
644
645        // Assert
646        assert!(i.load(Ordering::Relaxed) >= 5);
647    }
648
649    #[test]
650    fn shutdown_repeat() {
651        // Arrange
652        let exporter = InMemoryMetricExporter::default();
653        let reader = PeriodicReader::builder(exporter.clone()).build();
654
655        let meter_provider = SdkMeterProvider::builder().with_reader(reader).build();
656        let result = meter_provider.shutdown();
657        assert!(result.is_ok());
658
659        // calling shutdown again should return Err
660        let result = meter_provider.shutdown();
661        assert!(result.is_err());
662        assert!(matches!(result, Err(OTelSdkError::AlreadyShutdown)));
663
664        // calling shutdown again should return Err
665        let result = meter_provider.shutdown();
666        assert!(result.is_err());
667        assert!(matches!(result, Err(OTelSdkError::AlreadyShutdown)));
668    }
669
670    #[test]
671    fn flush_after_shutdown() {
672        // Arrange
673        let exporter = InMemoryMetricExporter::default();
674        let reader = PeriodicReader::builder(exporter.clone()).build();
675
676        let meter_provider = SdkMeterProvider::builder().with_reader(reader).build();
677        let result = meter_provider.force_flush();
678        assert!(result.is_ok());
679
680        let result = meter_provider.shutdown();
681        assert!(result.is_ok());
682
683        // calling force_flush after shutdown should return Err
684        let result = meter_provider.force_flush();
685        assert!(result.is_err());
686    }
687
688    #[test]
689    fn flush_repeat() {
690        // Arrange
691        let exporter = InMemoryMetricExporter::default();
692        let reader = PeriodicReader::builder(exporter.clone()).build();
693
694        let meter_provider = SdkMeterProvider::builder().with_reader(reader).build();
695        let result = meter_provider.force_flush();
696        assert!(result.is_ok());
697
698        // calling force_flush again should return Ok
699        let result = meter_provider.force_flush();
700        assert!(result.is_ok());
701    }
702
703    #[test]
704    fn periodic_reader_without_pipeline() {
705        // Arrange
706        let exporter = InMemoryMetricExporter::default();
707        let reader = PeriodicReader::builder(exporter.clone()).build();
708
709        let rm = &mut ResourceMetrics {
710            resource: Resource::empty(),
711            scope_metrics: Vec::new(),
712        };
713        // Pipeline is not registered, so collect should return an error
714        let result = reader.collect(rm);
715        assert!(result.is_err());
716
717        // Pipeline is not registered, so flush should return an error
718        let result = reader.force_flush();
719        assert!(result.is_err());
720
721        // Adding reader to meter provider should register the pipeline
722        // TODO: This part might benefit from a different design.
723        let meter_provider = SdkMeterProvider::builder()
724            .with_reader(reader.clone())
725            .build();
726
727        // Now collect and flush should succeed
728        let result = reader.collect(rm);
729        assert!(result.is_ok());
730
731        let result = meter_provider.force_flush();
732        assert!(result.is_ok());
733    }
734
735    #[test]
736    fn exporter_failures_are_handled() {
737        // create a mock exporter that fails 1st time and succeeds 2nd time
738        // Validate using this exporter that periodic reader can handle exporter failure
739        // and continue to export metrics.
740        // Arrange
741        let interval = std::time::Duration::from_millis(10);
742        let exporter = MetricExporterThatFailsOnlyOnFirst::default();
743        let reader = PeriodicReader::builder(exporter.clone())
744            .with_interval(interval)
745            .build();
746
747        let meter_provider = SdkMeterProvider::builder().with_reader(reader).build();
748        let meter = meter_provider.meter("test");
749        let counter = meter.u64_counter("sync_counter").build();
750        counter.add(1, &[]);
751        let _obs_counter = meter
752            .u64_observable_counter("testcounter")
753            .with_callback(move |observer| {
754                observer.observe(1, &[]);
755            })
756            .build();
757
758        // Sleep for a duration much longer than the interval to trigger
759        // multiple exports, including failures.
760        // Not a fan of such tests, but this seems to be the
761        // only way to test if periodic reader is doing its job. TODO: Decide if
762        // this should be ignored in CI
763        std::thread::sleep(Duration::from_millis(500));
764
765        // Assert that atleast 2 exports are attempted given the 1st one fails.
766        assert!(exporter.get_count() >= 2);
767    }
768
769    #[test]
770    fn shutdown_passed_to_exporter() {
771        // Arrange
772        let exporter = MockMetricExporter::default();
773        let reader = PeriodicReader::builder(exporter.clone()).build();
774
775        let meter_provider = SdkMeterProvider::builder().with_reader(reader).build();
776        let meter = meter_provider.meter("test");
777        let counter = meter.u64_counter("sync_counter").build();
778        counter.add(1, &[]);
779
780        // shutdown the provider, which should call shutdown on periodic reader
781        // which in turn should call shutdown on exporter.
782        let result = meter_provider.shutdown();
783        assert!(result.is_ok());
784        assert!(exporter.is_shutdown.load(Ordering::Relaxed));
785    }
786
787    #[test]
788    fn collection() {
789        collection_triggered_by_interval_helper();
790        collection_triggered_by_flush_helper();
791        collection_triggered_by_shutdown_helper();
792        collection_triggered_by_drop_helper();
793    }
794
795    #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
796    async fn collection_from_tokio_multi_with_one_worker() {
797        collection_triggered_by_interval_helper();
798        collection_triggered_by_flush_helper();
799        collection_triggered_by_shutdown_helper();
800        collection_triggered_by_drop_helper();
801    }
802
803    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
804    async fn collection_from_tokio_with_two_worker() {
805        collection_triggered_by_interval_helper();
806        collection_triggered_by_flush_helper();
807        collection_triggered_by_shutdown_helper();
808        collection_triggered_by_drop_helper();
809    }
810
811    #[tokio::test(flavor = "current_thread")]
812    async fn collection_from_tokio_current() {
813        collection_triggered_by_interval_helper();
814        collection_triggered_by_flush_helper();
815        collection_triggered_by_shutdown_helper();
816        collection_triggered_by_drop_helper();
817    }
818
819    fn collection_triggered_by_interval_helper() {
820        collection_helper(|_| {
821            // Sleep for a duration longer than the interval to ensure at least one collection
822            // Not a fan of such tests, but this seems to be the only way to test
823            // if periodic reader is doing its job.
824            // TODO: Decide if this should be ignored in CI
825            std::thread::sleep(Duration::from_millis(500));
826        });
827    }
828
829    fn collection_triggered_by_flush_helper() {
830        collection_helper(|meter_provider| {
831            meter_provider.force_flush().expect("flush should succeed");
832        });
833    }
834
835    fn collection_triggered_by_shutdown_helper() {
836        collection_helper(|meter_provider| {
837            meter_provider.shutdown().expect("shutdown should succeed");
838        });
839    }
840
841    fn collection_triggered_by_drop_helper() {
842        collection_helper(|meter_provider| {
843            drop(meter_provider);
844        });
845    }
846
847    fn collection_helper(trigger: fn(SdkMeterProvider)) {
848        // Arrange
849        let exporter = InMemoryMetricExporter::default();
850        let reader = PeriodicReader::builder(exporter.clone()).build();
851        let (sender, receiver) = mpsc::channel();
852
853        let meter_provider = SdkMeterProvider::builder().with_reader(reader).build();
854        let meter = meter_provider.meter("test");
855        let _counter = meter
856            .u64_observable_counter("testcounter")
857            .with_callback(move |observer| {
858                observer.observe(1, &[]);
859                sender.send(()).expect("channel should still be open");
860            })
861            .build();
862
863        // Act
864        trigger(meter_provider);
865
866        // Assert
867        receiver
868            .recv_timeout(Duration::ZERO)
869            .expect("message should be available in channel, indicating a collection occurred, which should trigger observable callback");
870
871        let exported_metrics = exporter
872            .get_finished_metrics()
873            .expect("this should not fail");
874        assert!(
875            !exported_metrics.is_empty(),
876            "Metrics should be available in exporter."
877        );
878    }
879
880    async fn some_async_function() -> u64 {
881        // No dependency on any particular async runtime.
882        std::thread::sleep(std::time::Duration::from_millis(1));
883        1
884    }
885
886    #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
887    async fn async_inside_observable_callback_from_tokio_multi_with_one_worker() {
888        async_inside_observable_callback_helper();
889    }
890
891    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
892    async fn async_inside_observable_callback_from_tokio_multi_with_two_worker() {
893        async_inside_observable_callback_helper();
894    }
895
896    #[tokio::test(flavor = "current_thread")]
897    async fn async_inside_observable_callback_from_tokio_current_thread() {
898        async_inside_observable_callback_helper();
899    }
900
901    #[test]
902    fn async_inside_observable_callback_from_regular_main() {
903        async_inside_observable_callback_helper();
904    }
905
906    fn async_inside_observable_callback_helper() {
907        let interval = std::time::Duration::from_millis(10);
908        let exporter = InMemoryMetricExporter::default();
909        let reader = PeriodicReader::builder(exporter.clone())
910            .with_interval(interval)
911            .build();
912
913        let meter_provider = SdkMeterProvider::builder().with_reader(reader).build();
914        let meter = meter_provider.meter("test");
915        let _gauge = meter
916            .u64_observable_gauge("my_observable_gauge")
917            .with_callback(|observer| {
918                // using futures_executor::block_on intentionally and avoiding
919                // any particular async runtime.
920                let value = futures_executor::block_on(some_async_function());
921                observer.observe(value, &[]);
922            })
923            .build();
924
925        meter_provider.force_flush().expect("flush should succeed");
926        let exported_metrics = exporter
927            .get_finished_metrics()
928            .expect("this should not fail");
929        assert!(
930            !exported_metrics.is_empty(),
931            "Metrics should be available in exporter."
932        );
933    }
934
935    async fn some_tokio_async_function() -> u64 {
936        // Tokio specific async function
937        tokio::time::sleep(Duration::from_millis(1)).await;
938        1
939    }
940
941    #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
942
943    async fn tokio_async_inside_observable_callback_from_tokio_multi_with_one_worker() {
944        tokio_async_inside_observable_callback_helper(true);
945    }
946
947    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
948    async fn tokio_async_inside_observable_callback_from_tokio_multi_with_two_worker() {
949        tokio_async_inside_observable_callback_helper(true);
950    }
951
952    #[tokio::test(flavor = "current_thread")]
953    #[ignore] //TODO: Investigate if this can be fixed.
954    async fn tokio_async_inside_observable_callback_from_tokio_current_thread() {
955        tokio_async_inside_observable_callback_helper(true);
956    }
957
958    #[test]
959    fn tokio_async_inside_observable_callback_from_regular_main() {
960        tokio_async_inside_observable_callback_helper(false);
961    }
962
963    fn tokio_async_inside_observable_callback_helper(use_current_tokio_runtime: bool) {
964        let exporter = InMemoryMetricExporter::default();
965        let reader = PeriodicReader::builder(exporter.clone()).build();
966
967        let meter_provider = SdkMeterProvider::builder().with_reader(reader).build();
968        let meter = meter_provider.meter("test");
969
970        if use_current_tokio_runtime {
971            let rt = tokio::runtime::Handle::current().clone();
972            let _gauge = meter
973                .u64_observable_gauge("my_observable_gauge")
974                .with_callback(move |observer| {
975                    // call tokio specific async function from here
976                    let value = rt.block_on(some_tokio_async_function());
977                    observer.observe(value, &[]);
978                })
979                .build();
980            // rt here is a reference to the current tokio runtime.
981            // Dropping it occurs when the tokio::main itself ends.
982        } else {
983            let rt = tokio::runtime::Runtime::new().unwrap();
984            let _gauge = meter
985                .u64_observable_gauge("my_observable_gauge")
986                .with_callback(move |observer| {
987                    // call tokio specific async function from here
988                    let value = rt.block_on(some_tokio_async_function());
989                    observer.observe(value, &[]);
990                })
991                .build();
992            // rt is not dropped here as it is moved to the closure,
993            // and is dropped only when MeterProvider itself is dropped.
994            // This works when called from normal main.
995        };
996
997        meter_provider.force_flush().expect("flush should succeed");
998        let exported_metrics = exporter
999            .get_finished_metrics()
1000            .expect("this should not fail");
1001        assert!(
1002            !exported_metrics.is_empty(),
1003            "Metrics should be available in exporter."
1004        );
1005    }
1006}