opentelemetry_sdk/metrics/
periodic_reader.rs1use 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
24pub const OTEL_METRIC_EXPORT_INTERVAL: &str = "OTEL_METRIC_EXPORT_INTERVAL";
27
28pub const OTEL_METRIC_EXPORT_INTERVAL_DEFAULT: Duration = Duration::from_secs(60);
30
31#[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 pub fn with_interval(mut self, interval: Duration) -> Self {
59 if !interval.is_zero() {
60 self.interval = interval;
61 }
62 self
63 }
64
65 pub fn build(self) -> PeriodicReader<E> {
67 PeriodicReader::new(self.exporter, self.interval)
68 }
69}
70
71pub 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 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 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 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 interval_start = Instant::now();
243 remaining_interval = Duration::ZERO;
244 }
245 }
246 Ok(Message::Shutdown(response_sender)) => {
247 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 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 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 otel_debug!(
325 name: "PeriodReaderThreadExiting",
326 reason = "MessageSenderDisconnected"
327 );
328 break;
329 }
330 }
331 }
332 otel_debug!(
333 name: "PeriodReaderThreadStopped"
334 );
335 });
336
337 #[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 futures_executor::block_on(self.exporter.export(rm))
423 }
424
425 fn force_flush(&self) -> OTelSdkResult {
426 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 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 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 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 fn shutdown_with_timeout(&self, _timeout: Duration) -> OTelSdkResult {
506 self.inner.shutdown()
507 }
508
509 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 #[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 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 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 std::thread::sleep(interval * 5 * 20);
644
645 assert!(i.load(Ordering::Relaxed) >= 5);
647 }
648
649 #[test]
650 fn shutdown_repeat() {
651 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 let result = meter_provider.shutdown();
661 assert!(result.is_err());
662 assert!(matches!(result, Err(OTelSdkError::AlreadyShutdown)));
663
664 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 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 let result = meter_provider.force_flush();
685 assert!(result.is_err());
686 }
687
688 #[test]
689 fn flush_repeat() {
690 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 let result = meter_provider.force_flush();
700 assert!(result.is_ok());
701 }
702
703 #[test]
704 fn periodic_reader_without_pipeline() {
705 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 let result = reader.collect(rm);
715 assert!(result.is_err());
716
717 let result = reader.force_flush();
719 assert!(result.is_err());
720
721 let meter_provider = SdkMeterProvider::builder()
724 .with_reader(reader.clone())
725 .build();
726
727 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 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 std::thread::sleep(Duration::from_millis(500));
764
765 assert!(exporter.get_count() >= 2);
767 }
768
769 #[test]
770 fn shutdown_passed_to_exporter() {
771 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 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 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 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 trigger(meter_provider);
865
866 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 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 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::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] 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 let value = rt.block_on(some_tokio_async_function());
977 observer.observe(value, &[]);
978 })
979 .build();
980 } 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 let value = rt.block_on(some_tokio_async_function());
989 observer.observe(value, &[]);
990 })
991 .build();
992 };
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}