1use std::collections::BTreeMap;
29use std::hash::{Hash, Hasher};
30use std::io;
31use std::io::IsTerminal;
32use std::str::FromStr;
33use std::sync::LazyLock;
34use std::sync::atomic::{AtomicU64, Ordering};
35use std::sync::{Arc, Mutex, OnceLock};
36use std::time::Duration;
37
38#[cfg(feature = "tokio-console")]
39use console_subscriber::ConsoleLayer;
40use derivative::Derivative;
41use http::HeaderMap;
42use hyper_tls::HttpsConnector;
43use hyper_util::client::legacy::connect::HttpConnector;
44use opentelemetry::propagation::{Extractor, Injector};
45use opentelemetry::trace::TracerProvider;
46use opentelemetry::{KeyValue, global};
47use opentelemetry_otlp::WithTonicConfig;
48use opentelemetry_sdk::propagation::TraceContextPropagator;
49use opentelemetry_sdk::runtime::Tokio;
50use opentelemetry_sdk::trace::span_processor_with_async_runtime::BatchSpanProcessor;
51use opentelemetry_sdk::{Resource, trace};
52use prometheus::IntCounter;
53use tonic::metadata::MetadataMap;
54use tonic::transport::Endpoint;
55use tracing::{Event, Level, Span, Subscriber, warn};
56#[cfg(feature = "capture")]
57use tracing_capture::{CaptureLayer, SharedStorage};
58use tracing_opentelemetry::OpenTelemetrySpanExt;
59use tracing_subscriber::filter::Directive;
60use tracing_subscriber::fmt::format::{Writer, format};
61use tracing_subscriber::fmt::{self, FmtContext, FormatEvent, FormatFields};
62use tracing_subscriber::layer::{Layer, SubscriberExt};
63use tracing_subscriber::registry::LookupSpan;
64use tracing_subscriber::util::SubscriberInitExt;
65use tracing_subscriber::{EnvFilter, Registry, reload};
66
67use crate::metric;
68use crate::metrics::MetricsRegistry;
69#[cfg(feature = "tokio-console")]
70use crate::netio::SocketAddr;
71use crate::now::{EpochMillis, NowFn, SYSTEM_TIME};
72
73#[derive(Derivative)]
77#[derivative(Debug)]
78pub struct TracingConfig<F> {
79 pub service_name: &'static str,
81 pub stderr_log: StderrLogConfig,
83 pub opentelemetry: Option<OpenTelemetryConfig>,
85 #[cfg_attr(nightly_doc_features, doc(cfg(feature = "tokio-console")))]
89 #[cfg(feature = "tokio-console")]
90 pub tokio_console: Option<TokioConsoleConfig>,
91 #[cfg(feature = "capture")]
93 #[derivative(Debug = "ignore")]
94 pub capture: Option<SharedStorage>,
95 pub sentry: Option<SentryConfig<F>>,
97 pub build_version: &'static str,
99 pub build_sha: &'static str,
101 pub registry: MetricsRegistry,
103}
104
105#[derive(Debug, Clone)]
107pub struct SentryConfig<F> {
108 pub dsn: String,
110 pub environment: Option<String>,
115 pub tags: BTreeMap<String, String>,
117 pub event_filter: F,
119}
120
121#[derive(Debug)]
123pub struct StderrLogConfig {
124 pub format: StderrLogFormat,
126 pub filter: EnvFilter,
128}
129
130#[derive(Debug, Clone)]
132pub enum StderrLogFormat {
133 Text {
137 prefix: Option<String>,
139 },
140 Json,
144}
145
146#[derive(Debug)]
148pub struct OpenTelemetryConfig {
149 pub endpoint: String,
153 pub headers: HeaderMap,
155 pub filter: EnvFilter,
157 pub max_batch_queue_size: usize,
159 pub max_export_batch_size: usize,
161 pub max_concurrent_exports: usize,
164 pub batch_scheduled_delay: Duration,
166 pub max_export_timeout: Duration,
168 pub resource: Resource,
171}
172
173#[cfg_attr(nightly_doc_features, doc(cfg(feature = "tokio-console")))]
177#[cfg(feature = "tokio-console")]
178#[derive(Debug, Clone)]
179pub struct TokioConsoleConfig {
180 pub listen_addr: SocketAddr,
184 pub publish_interval: Duration,
188 pub retention: Duration,
192}
193
194type Reloader = Arc<dyn Fn(EnvFilter, Vec<Directive>) -> Result<(), anyhow::Error> + Send + Sync>;
195type DirectiveReloader = Arc<dyn Fn(Vec<Directive>) -> Result<(), anyhow::Error> + Send + Sync>;
196
197#[derive(Clone)]
199pub struct TracingHandle {
200 stderr_log: Reloader,
201 opentelemetry: Reloader,
202 sentry: DirectiveReloader,
203}
204
205impl TracingHandle {
206 pub fn disabled() -> TracingHandle {
210 TracingHandle {
211 stderr_log: Arc::new(|_, _| Ok(())),
212 opentelemetry: Arc::new(|_, _| Ok(())),
213 sentry: Arc::new(|_| Ok(())),
214 }
215 }
216
217 pub fn reload_stderr_log_filter(
219 &self,
220 filter: EnvFilter,
221 defaults: Vec<Directive>,
222 ) -> Result<(), anyhow::Error> {
223 (self.stderr_log)(filter, defaults)
224 }
225
226 pub fn reload_opentelemetry_filter(
228 &self,
229 filter: EnvFilter,
230 defaults: Vec<Directive>,
231 ) -> Result<(), anyhow::Error> {
232 (self.opentelemetry)(filter, defaults)
233 }
234
235 pub fn reload_sentry_directives(&self, defaults: Vec<Directive>) -> Result<(), anyhow::Error> {
237 (self.sentry)(defaults)
238 }
239}
240
241impl std::fmt::Debug for TracingHandle {
242 fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
243 f.debug_struct("TracingHandle").finish_non_exhaustive()
244 }
245}
246
247pub const LOGGING_DEFAULTS_STR: [&str; 2] = [
257 "kube_client::client::builder=off",
258 "aws_config::profile::credentials=off",
261];
262pub static LOGGING_DEFAULTS: LazyLock<Vec<Directive>> = LazyLock::new(|| {
264 LOGGING_DEFAULTS_STR
265 .into_iter()
266 .map(|directive| Directive::from_str(directive).expect("valid directive"))
267 .collect()
268});
269pub const OPENTELEMETRY_DEFAULTS_STR: [&str; 2] = ["h2=off", "hyper=off"];
275pub static OPENTELEMETRY_DEFAULTS: LazyLock<Vec<Directive>> = LazyLock::new(|| {
277 OPENTELEMETRY_DEFAULTS_STR
278 .into_iter()
279 .map(|directive| Directive::from_str(directive).expect("valid directive"))
280 .collect()
281});
282
283pub const SENTRY_DEFAULTS_STR: [&str; 2] =
286 ["kube_client::client::builder=off", "mysql_async::conn=off"];
287pub static SENTRY_DEFAULTS: LazyLock<Vec<Directive>> = LazyLock::new(|| {
289 SENTRY_DEFAULTS_STR
290 .into_iter()
291 .map(|directive| Directive::from_str(directive).expect("valid directive"))
292 .collect()
293});
294
295type GlobalSubscriber = Arc<dyn Subscriber + Send + Sync + 'static>;
297
298pub static GLOBAL_SUBSCRIBER: OnceLock<GlobalSubscriber> = OnceLock::new();
301
302#[allow(clippy::unused_async)]
323pub async fn configure<F>(config: TracingConfig<F>) -> Result<TracingHandle, anyhow::Error>
324where
325 F: Fn(&tracing::Metadata<'_>) -> sentry_tracing::EventFilter + Send + Sync + 'static,
326{
327 let stderr_log_layer: Box<dyn Layer<Registry> + Send + Sync> = match config.stderr_log.format {
328 StderrLogFormat::Text { prefix } => {
329 let no_color = std::env::var_os("NO_COLOR").unwrap_or_else(|| "".into()) != "";
331 Box::new(
332 fmt::layer()
333 .with_writer(io::stderr)
334 .event_format(PrefixFormat {
335 inner: format(),
336 prefix,
337 })
338 .with_ansi(!no_color && io::stderr().is_terminal()),
339 )
340 }
341 StderrLogFormat::Json => Box::new(
342 fmt::layer()
343 .with_writer(io::stderr)
344 .json()
345 .with_current_span(true),
346 ),
347 };
348 let (stderr_log_filter, stderr_log_filter_reloader) = reload::Layer::new({
349 let mut filter = config.stderr_log.filter;
350 for directive in LOGGING_DEFAULTS.iter() {
351 filter = filter.add_directive(directive.clone());
352 }
353 filter
354 });
355 let otel_rate_limit_filter = OpenTelemetryRateLimitingFilter::new(Duration::from_secs(
359 OPENTELEMETRY_RATE_LIMIT_BACKOFF_SECS,
360 ));
361 let stderr_log_layer = stderr_log_layer
366 .with_filter(otel_rate_limit_filter)
367 .with_filter(stderr_log_filter);
368 let stderr_log_reloader = Arc::new(move |mut filter: EnvFilter, defaults: Vec<Directive>| {
369 for directive in &defaults {
370 filter = filter.add_directive(directive.clone());
371 }
372 Ok(stderr_log_filter_reloader.reload(filter)?)
373 });
374
375 let (otel_layer, otel_reloader): (_, Reloader) = if let Some(otel_config) = config.opentelemetry
376 {
377 opentelemetry::global::set_text_map_propagator(TraceContextPropagator::new());
378
379 let channel = Endpoint::from_shared(otel_config.endpoint)?
383 .timeout(opentelemetry_otlp::OTEL_EXPORTER_OTLP_TIMEOUT_DEFAULT)
384 .connect_with_connector_lazy({
386 let mut http = HttpConnector::new();
387 http.enforce_http(false);
388 HttpsConnector::from((
389 http,
390 tokio_native_tls::TlsConnector::from(
392 native_tls::TlsConnector::builder()
393 .request_alpns(&["h2"])
394 .build()
395 .unwrap(),
396 ),
397 ))
398 });
399 let exporter = opentelemetry_otlp::SpanExporter::builder()
400 .with_tonic()
401 .with_channel(channel)
402 .with_metadata(MetadataMap::from_headers(otel_config.headers))
403 .build()?;
404 let batch_config = opentelemetry_sdk::trace::BatchConfigBuilder::default()
405 .with_max_queue_size(otel_config.max_batch_queue_size)
406 .with_max_export_batch_size(otel_config.max_export_batch_size)
407 .with_max_concurrent_exports(otel_config.max_concurrent_exports)
408 .with_scheduled_delay(otel_config.batch_scheduled_delay)
409 .with_max_export_timeout(otel_config.max_export_timeout)
410 .build();
411 let batch_span_processor = BatchSpanProcessor::builder(exporter, Tokio)
412 .with_batch_config(batch_config)
413 .build();
414 let tracer = trace::SdkTracerProvider::builder()
415 .with_resource(
416 Resource::builder()
417 .with_service_name(config.service_name.to_string())
418 .with_attributes(
419 otel_config
420 .resource
421 .iter()
422 .map(|(k, v)| KeyValue::new(k.clone(), v.clone())),
423 )
425 .build(),
426 )
427 .with_span_processor(batch_span_processor)
428 .with_max_events_per_span(2048)
429 .build()
430 .tracer(config.service_name);
431
432 let (filter, filter_handle) = reload::Layer::new({
433 let mut filter = otel_config.filter;
434 for directive in OPENTELEMETRY_DEFAULTS.iter() {
435 filter = filter.add_directive(directive.clone());
436 }
437 filter
438 });
439 let metrics_layer = MetricsLayer::new(&config.registry);
440 let layer = tracing_opentelemetry::layer()
441 .with_tracer(tracer)
447 .and_then(metrics_layer)
448 .with_filter(filter);
452 let reloader = Arc::new(move |mut filter: EnvFilter, defaults: Vec<Directive>| {
453 for directive in &defaults {
455 filter = filter.add_directive(directive.clone());
456 }
457 Ok(filter_handle.reload(filter)?)
458 });
459 (Some(layer), reloader)
460 } else {
461 let reloader = Arc::new(|_, _| Ok(()));
462 (None, reloader)
463 };
464
465 #[cfg(feature = "tokio-console")]
466 let tokio_console_layer = if let Some(console_config) = config.tokio_console.clone() {
467 let builder = ConsoleLayer::builder()
468 .publish_interval(console_config.publish_interval)
469 .retention(console_config.retention);
470 let builder = match console_config.listen_addr {
471 SocketAddr::Inet(addr) => builder.server_addr(addr),
472 SocketAddr::Unix(addr) => {
473 let path = addr.as_pathname().unwrap().as_ref();
474 builder.server_addr(path)
475 }
476 SocketAddr::Turmoil(_) => unimplemented!(),
477 };
478 Some(builder.spawn())
479 } else {
480 None
481 };
482
483 let (sentry_layer, sentry_reloader): (_, DirectiveReloader) =
484 if let Some(sentry_config) = config.sentry {
485 let mut options = sentry::ClientOptions::default();
486 options.attach_stacktrace = true;
487 options.release = Some(format!("materialize@{0}", config.build_version).into());
488 options.environment = sentry_config.environment.map(Into::into);
489 let guard = sentry::init((sentry_config.dsn, options));
490
491 std::mem::forget(guard);
494
495 sentry::configure_scope(|scope| {
496 scope.set_tag("service_name", config.service_name);
497 scope.set_tag("build_sha", config.build_sha.to_string());
498 for (k, v) in sentry_config.tags {
499 scope.set_tag(&k, v);
500 }
501 });
502
503 let (filter, filter_handle) = reload::Layer::new({
504 let mut filter = EnvFilter::new("info");
506 for directive in SENTRY_DEFAULTS.iter() {
507 filter = filter.add_directive(directive.clone());
508 }
509 filter
510 });
511 let layer = sentry_tracing::layer()
512 .event_filter(sentry_config.event_filter)
513 .with_filter(filter);
538 let reloader = Arc::new(move |defaults: Vec<Directive>| {
539 let mut filter = EnvFilter::new("info");
541 for directive in &defaults {
543 filter = filter.add_directive(directive.clone());
544 }
545 Ok(filter_handle.reload(filter)?)
546 });
547 (Some(layer), reloader)
548 } else {
549 let reloader = Arc::new(|_| Ok(()));
550 (None, reloader)
551 };
552
553 #[cfg(feature = "capture")]
554 let capture = config.capture.map(|storage| CaptureLayer::new(&storage));
555
556 let stack = tracing_subscriber::registry();
557 let stack = stack.with(stderr_log_layer);
558 #[cfg(feature = "capture")]
559 let stack = stack.with(capture);
560 let stack = stack.with(otel_layer);
561 #[cfg(feature = "tokio-console")]
562 let stack = stack.with(tokio_console_layer);
563 let stack = stack.with(sentry_layer);
564
565 assert!(GLOBAL_SUBSCRIBER.set(Arc::new(stack)).is_ok());
567 Arc::clone(GLOBAL_SUBSCRIBER.get().unwrap()).init();
569
570 #[cfg(feature = "tokio-console")]
571 if let Some(console_config) = config.tokio_console {
572 let endpoint = match console_config.listen_addr {
573 SocketAddr::Inet(addr) => format!("http://{addr}"),
574 SocketAddr::Unix(addr) => format!("file://localhost{addr}"),
575 SocketAddr::Turmoil(_) => unimplemented!(),
576 };
577 tracing::info!("starting tokio console on {endpoint}");
578 }
579
580 let handle = TracingHandle {
581 stderr_log: stderr_log_reloader,
582 opentelemetry: otel_reloader,
583 sentry: sentry_reloader,
584 };
585
586 Ok(handle)
587}
588
589pub fn crate_level(filter: &EnvFilter, crate_name: &'static str) -> Level {
592 let mut default_level = Level::ERROR;
598 for directive in format!("{}", filter).split(',') {
601 match directive.split('=').collect::<Vec<_>>().as_slice() {
602 [target, level] => {
603 if *target == crate_name {
604 match Level::from_str(*level) {
605 Ok(level) => return level,
606 Err(err) => warn!("invalid level for {}: {}", target, err),
607 }
608 }
609 }
610 [token] => match Level::from_str(*token) {
611 Ok(level) => default_level = default_level.max(level),
612 Err(_) => {
613 if *token == crate_name {
615 default_level = default_level.max(Level::TRACE);
616 }
617 }
618 },
619 _ => {}
620 }
621 }
622
623 default_level
624}
625
626#[derive(Debug)]
629pub struct PrefixFormat<F> {
630 inner: F,
631 prefix: Option<String>,
632}
633
634impl<F, C, N> FormatEvent<C, N> for PrefixFormat<F>
635where
636 C: Subscriber + for<'a> LookupSpan<'a>,
637 N: for<'a> FormatFields<'a> + 'static,
638 F: FormatEvent<C, N>,
639{
640 fn format_event(
641 &self,
642 ctx: &FmtContext<'_, C, N>,
643 mut writer: Writer<'_>,
644 event: &Event<'_>,
645 ) -> std::fmt::Result {
646 match &self.prefix {
647 None => self.inner.format_event(ctx, writer, event)?,
648 Some(prefix) => {
649 let mut prefix = yansi::Paint::new(prefix);
650 if writer.has_ansi_escapes() {
651 prefix = prefix.bold();
652 }
653 write!(writer, "{}: ", prefix)?;
654 self.inner.format_event(ctx, writer, event)?;
655 }
656 }
657 Ok(())
658 }
659}
660
661#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
665pub struct OpenTelemetryContext {
666 inner: BTreeMap<String, String>,
667}
668
669impl OpenTelemetryContext {
670 pub fn attach_as_parent(&self) {
677 self.attach_as_parent_to(&tracing::Span::current())
678 }
679
680 pub fn attach_as_parent_to(&self, span: &Span) {
686 let parent_cx = global::get_text_map_propagator(|prop| prop.extract(self));
687 let _ = span.set_parent(parent_cx);
688 }
689
690 pub fn obtain() -> Self {
692 let mut context = Self::empty();
693 global::get_text_map_propagator(|propagator| {
694 propagator.inject_context(&tracing::Span::current().context(), &mut context)
695 });
696
697 context
698 }
699
700 pub fn empty() -> Self {
702 Self {
703 inner: BTreeMap::new(),
704 }
705 }
706}
707
708impl Extractor for OpenTelemetryContext {
709 fn get(&self, key: &str) -> Option<&str> {
710 self.inner.get(&key.to_lowercase()).map(|v| v.as_str())
711 }
712
713 fn keys(&self) -> Vec<&str> {
714 self.inner.keys().map(|k| k.as_str()).collect::<Vec<_>>()
715 }
716}
717
718impl Injector for OpenTelemetryContext {
719 fn set(&mut self, key: &str, value: String) {
720 self.inner.insert(key.to_lowercase(), value);
721 }
722}
723
724impl From<OpenTelemetryContext> for BTreeMap<String, String> {
725 fn from(ctx: OpenTelemetryContext) -> Self {
726 ctx.inner
727 }
728}
729
730impl From<BTreeMap<String, String>> for OpenTelemetryContext {
731 fn from(map: BTreeMap<String, String>) -> Self {
732 Self { inner: map }
733 }
734}
735
736struct MetricsLayer {
737 on_close: IntCounter,
738}
739
740impl MetricsLayer {
741 fn new(registry: &MetricsRegistry) -> Self {
742 MetricsLayer {
743 on_close: registry.register(metric!(
744 name: "mz_otel_on_close",
745 help: "count of on_close events sent to otel",
746 )),
747 }
748 }
749}
750
751impl<S: tracing::Subscriber> Layer<S> for MetricsLayer {
752 fn on_close(&self, _id: tracing::span::Id, _ctx: tracing_subscriber::layer::Context<'_, S>) {
753 self.on_close.inc()
754 }
755}
756
757const OPENTELEMETRY_TARGET_PREFIX: &str = "opentelemetry";
761
762const OPENTELEMETRY_RATE_LIMIT_BACKOFF_SECS: u64 = 30;
764
765#[derive(Debug)]
775pub struct OpenTelemetryRateLimitingFilter {
776 backoff_duration: Duration,
778 last_logged: Mutex<BTreeMap<u64, EpochMillis>>,
781 suppressed_count: AtomicU64,
783 now_fn: NowFn,
785}
786
787impl OpenTelemetryRateLimitingFilter {
788 pub fn new(backoff_duration: Duration) -> Self {
793 Self {
794 backoff_duration,
795 last_logged: Mutex::new(BTreeMap::new()),
796 suppressed_count: AtomicU64::new(0),
797 now_fn: SYSTEM_TIME.clone(),
798 }
799 }
800
801 #[cfg(test)]
803 fn with_now_fn(mut self, now_fn: NowFn) -> Self {
804 self.now_fn = now_fn;
805 self
806 }
807
808 fn compute_key(metadata: &tracing::Metadata<'_>) -> u64 {
810 let mut hasher = std::hash::DefaultHasher::new();
811 metadata.target().hash(&mut hasher);
812 metadata.name().hash(&mut hasher);
813 if let Some(file) = metadata.file() {
815 file.hash(&mut hasher);
816 }
817 if let Some(line) = metadata.line() {
818 line.hash(&mut hasher);
819 }
820 hasher.finish()
821 }
822
823 fn should_log(&self, metadata: &tracing::Metadata<'_>) -> bool {
825 if !metadata.target().starts_with(OPENTELEMETRY_TARGET_PREFIX) {
827 return true;
828 }
829
830 let key = Self::compute_key(metadata);
831 let now = (self.now_fn)();
832
833 let mut last_logged = self.last_logged.lock().unwrap();
834
835 if let Some(last_time) = last_logged.get(&key) {
836 if Duration::from_millis(now - last_time) < self.backoff_duration {
837 self.suppressed_count.fetch_add(1, Ordering::Relaxed);
838 return false;
839 }
840 }
841
842 last_logged.insert(key, now);
843
844 if last_logged.len() > 1000 {
847 last_logged
848 .retain(|_, time| Duration::from_millis(now - *time) < self.backoff_duration);
849 }
850
851 true
852 }
853
854 pub fn suppressed_count(&self) -> u64 {
856 self.suppressed_count.load(Ordering::Relaxed)
857 }
858}
859
860impl<S> tracing_subscriber::layer::Filter<S> for OpenTelemetryRateLimitingFilter
861where
862 S: tracing::Subscriber + for<'lookup> LookupSpan<'lookup>,
863{
864 fn enabled(
870 &self,
871 _metadata: &tracing::Metadata<'_>,
872 _ctx: &tracing_subscriber::layer::Context<'_, S>,
873 ) -> bool {
874 true
875 }
876
877 fn event_enabled(
878 &self,
879 event: &Event<'_>,
880 _ctx: &tracing_subscriber::layer::Context<'_, S>,
881 ) -> bool {
882 self.should_log(event.metadata())
883 }
884}
885
886#[cfg(test)]
887mod tests {
888 use std::str::FromStr;
889 use tracing::Level;
890 use tracing_subscriber::filter::{EnvFilter, LevelFilter, Targets};
891
892 #[crate::test]
893 fn overriding_targets() {
894 let user_defined = Targets::new().with_target("my_crate", Level::INFO);
895
896 let default = Targets::new().with_target("my_crate", LevelFilter::OFF);
897 assert!(!default.would_enable("my_crate", &Level::INFO));
898
899 let filters = Targets::new()
901 .with_targets(default)
902 .with_targets(user_defined);
903 assert!(filters.would_enable("my_crate", &Level::INFO));
904 }
905
906 #[crate::test]
907 fn crate_level() {
908 let filter = EnvFilter::from_str("abc=trace,def=debug").expect("valid");
910 assert_eq!(super::crate_level(&filter, "abc"), Level::TRACE);
911 assert_eq!(super::crate_level(&filter, "def"), Level::DEBUG);
912 assert_eq!(super::crate_level(&filter, "def"), Level::DEBUG);
913 assert_eq!(
914 super::crate_level(&filter, "abc::doesnt::exist"),
915 Level::ERROR
916 );
917 assert_eq!(super::crate_level(&filter, "doesnt::exist"), Level::ERROR);
918
919 let filter = EnvFilter::from_str("abc=trace,def=debug,info").expect("valid");
921 assert_eq!(super::crate_level(&filter, "abc"), Level::TRACE);
922 assert_eq!(
923 super::crate_level(&filter, "abc::doesnt:exist"),
924 Level::INFO
925 );
926 assert_eq!(super::crate_level(&filter, "def"), Level::DEBUG);
927 assert_eq!(super::crate_level(&filter, "nan"), Level::INFO);
928
929 let filter = EnvFilter::from_str("abc::def::ghi=trace,debug").expect("valid");
931 assert_eq!(super::crate_level(&filter, "abc"), Level::DEBUG);
932 assert_eq!(super::crate_level(&filter, "def"), Level::DEBUG);
933 assert_eq!(
934 super::crate_level(&filter, "gets_the_default"),
935 Level::DEBUG
936 );
937
938 let filter =
940 EnvFilter::from_str("abc[s]=trace,def[s{g=h}]=debug,[{s2}]=debug,info").expect("valid");
941 assert_eq!(super::crate_level(&filter, "abc"), Level::INFO);
942 assert_eq!(super::crate_level(&filter, "def"), Level::INFO);
943 assert_eq!(super::crate_level(&filter, "gets_the_default"), Level::INFO);
944
945 let filter = EnvFilter::from_str("abc,info").expect("valid");
947 assert_eq!(super::crate_level(&filter, "abc"), Level::TRACE);
948 assert_eq!(super::crate_level(&filter, "gets_the_default"), Level::INFO);
949 assert_eq!(super::crate_level(&filter, "abc::def"), Level::INFO);
953 }
954
955 #[crate::test]
956 fn otel_rate_limiting_filter_backoff() {
957 use std::sync::Arc;
958 use std::sync::atomic::{AtomicU64, Ordering};
959 use std::time::Duration;
960 use tracing::Callsite;
961
962 use crate::now::NowFn;
963
964 let current_time = Arc::new(AtomicU64::new(0));
966 let time_for_closure = Arc::clone(¤t_time);
967 let now_fn: NowFn = NowFn::from(move || time_for_closure.load(Ordering::SeqCst));
968
969 let filter = super::OpenTelemetryRateLimitingFilter::new(Duration::from_millis(100))
970 .with_now_fn(now_fn);
971
972 static OTEL_CALLSITE: tracing::callsite::DefaultCallsite =
974 tracing::callsite::DefaultCallsite::new(&tracing::Metadata::new(
975 "test_event",
976 "opentelemetry_sdk::trace",
977 Level::WARN,
978 Some(file!()),
979 Some(line!()),
980 Some(module_path!()),
981 tracing::field::FieldSet::new(&[], tracing::callsite::Identifier(&OTEL_CALLSITE)),
982 tracing::metadata::Kind::EVENT,
983 ));
984 let otel_meta = OTEL_CALLSITE.metadata();
985
986 static OTHER_CALLSITE: tracing::callsite::DefaultCallsite =
987 tracing::callsite::DefaultCallsite::new(&tracing::Metadata::new(
988 "test_event",
989 "my_app::module",
990 Level::WARN,
991 Some(file!()),
992 Some(line!()),
993 Some(module_path!()),
994 tracing::field::FieldSet::new(&[], tracing::callsite::Identifier(&OTHER_CALLSITE)),
995 tracing::metadata::Kind::EVENT,
996 ));
997 let other_meta = OTHER_CALLSITE.metadata();
998
999 assert!(filter.should_log(other_meta));
1001 assert!(filter.should_log(other_meta));
1002 assert!(filter.should_log(other_meta));
1003 assert_eq!(filter.suppressed_count(), 0);
1004
1005 assert!(filter.should_log(otel_meta));
1007 assert_eq!(filter.suppressed_count(), 0);
1008
1009 current_time.store(50, Ordering::SeqCst); assert!(!filter.should_log(otel_meta));
1012 assert!(!filter.should_log(otel_meta));
1013 assert_eq!(filter.suppressed_count(), 2);
1014
1015 current_time.store(150, Ordering::SeqCst); assert!(filter.should_log(otel_meta));
1018 assert_eq!(filter.suppressed_count(), 2);
1019 }
1020}