Skip to main content

relay_statsd/
lib.rs

1//! A high-level StatsD metric client built on cadence.
2//!
3//! ## Defining Metrics
4//!
5//! In order to use metrics, one needs to first define one of the metric traits on a custom enum.
6//! The following types of metrics are available: `counter`, `timer`, `gauge`, `distribution`, and
7//! `set`. For explanations on what that means see [Metric Types].
8//!
9//! The metric traits serve only to provide a type safe metric name. All metric types have exactly
10//! the same form, they are different only to ensure that a metric can only be used for the type for
11//! which it was defined, (e.g. a counter metric cannot be used as a timer metric). See the traits
12//! for more detailed examples.
13//!
14//! ## Initializing the Client
15//!
16//! Metrics can be used without initializing a statsd client. In that case, invoking `with_client`
17//! or the [`metric!`] macro will become a noop. Only when configured, metrics will actually be
18//! collected.
19//!
20//! To initialize the client, use [`init`] to create a default client with known arguments:
21//!
22//! ```no_run
23//! # use std::collections::BTreeMap;
24//! # use relay_statsd::MetricsConfig;
25//!
26//! relay_statsd::init(MetricsConfig {
27//!     prefix: "myprefix".to_owned(),
28//!     host: "localhost:8125".to_owned(),
29//!     buffer_size: None,
30//!     default_tags: BTreeMap::new(),
31//! });
32//! ```
33//!
34//! ## Macro Usage
35//!
36//! The recommended way to record metrics is by using the [`metric!`] macro. See the trait docs
37//! for more information on how to record each type of metric.
38//!
39//! ```
40//! use relay_statsd::{metric, CounterMetric};
41//!
42//! struct MyCounter;
43//!
44//! impl CounterMetric for MyCounter {
45//!     fn name(&self) -> &'static str {
46//!         "counter"
47//!     }
48//! }
49//!
50//! metric!(counter(MyCounter) += 1);
51//! ```
52//! [Metric Types]: https://github.com/statsd/statsd/blob/master/docs/metric_types.md
53use metrics_exporter_dogstatsd::{AggregationMode, BuildError, DogStatsDBuilder};
54use metrics_util::MetricKindMask;
55use metrics_util::layers::{Layer, PrefixLayer, RouterBuilder};
56
57use std::{collections::BTreeMap, fmt, sync::Arc};
58
59use crate::mock::MockRecorder;
60
61mod mock;
62
63/// Custom namespaces which are never prefixed.
64const CUSTOM_NAMESPACES: &[&str] = &["arroyo.", "datadog.dogstatsd.client."];
65
66#[doc(hidden)]
67pub mod _metrics {
68    pub use ::metrics::*;
69}
70
71/// Client configuration used for initialization of the metrics sub-system.
72#[derive(Debug)]
73pub struct MetricsConfig {
74    /// Prefix which is appended to all metric names, except for DogStatsD and Arroyo metrics.
75    pub prefix: String,
76    /// Host of the metrics upstream.
77    pub host: String,
78    /// The buffer size to use for the socket.
79    pub buffer_size: Option<usize>,
80    /// Tags that are added to all metrics.
81    pub default_tags: BTreeMap<String, String>,
82}
83
84/// Error returned from [`init`].
85#[derive(Debug)]
86pub struct Error(BuildError);
87
88impl fmt::Display for Error {
89    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
90        self.0.fmt(f)
91    }
92}
93
94impl std::error::Error for Error {}
95
96impl From<BuildError> for Error {
97    fn from(value: BuildError) -> Self {
98        Self(value)
99    }
100}
101
102/// Set a test client for the period of the called function (only affects the current thread).
103pub fn with_capturing_test_client(f: impl FnOnce()) -> Vec<String> {
104    let recorder = MockRecorder::default();
105    metrics::with_local_recorder(&recorder, f);
106    recorder.consume()
107}
108
109/// Tell the metrics system to report to statsd.
110pub fn init(config: MetricsConfig) -> Result<(), Error> {
111    relay_log::info!("reporting metrics to statsd at {}", config.host);
112
113    let default_labels = config
114        .default_tags
115        .into_iter()
116        .map(|(key, value)| metrics::Label::new(key, value))
117        .collect();
118
119    let mut statsd = DogStatsDBuilder::default()
120        .with_remote_address(&config.host)?
121        .with_telemetry(true)
122        .with_aggregation_mode(AggregationMode::Aggressive)
123        .send_histograms_as_distributions(true)
124        .with_histogram_sampling(true)
125        .with_global_labels(default_labels);
126
127    if let Some(buffer_size) = config.buffer_size {
128        statsd = statsd.with_maximum_payload_length(buffer_size)?;
129    };
130
131    let recorder = Arc::new(statsd.build()?);
132
133    // Metrics routed through this layer have the configured prefix prepended.
134    let prefix_layer = PrefixLayer::new(config.prefix).layer(Arc::clone(&recorder));
135
136    // Router decides if metrics go through the prefix layer or directly to the DogStatsD recorder.
137    let mut router = RouterBuilder::from_recorder(prefix_layer);
138    for namespace in CUSTOM_NAMESPACES {
139        router.add_route(MetricKindMask::ALL, namespace, Arc::clone(&recorder));
140    }
141
142    metrics::set_global_recorder(router.build()).map_err(|_| BuildError::FailedToInstall)?;
143
144    Ok(())
145}
146
147/// A metric for capturing timings.
148///
149/// Timings are a positive number of milliseconds between a start and end time. Examples include
150/// time taken to render a web page or time taken for a database call to return.
151///
152/// ## Example
153///
154/// ```
155/// use relay_statsd::{metric, TimerMetric};
156///
157/// enum MyTimer {
158///     ProcessA,
159///     ProcessB,
160/// }
161///
162/// impl TimerMetric for MyTimer {
163///     fn name(&self) -> &'static str {
164///         match self {
165///             Self::ProcessA => "process_a",
166///             Self::ProcessB => "process_b",
167///         }
168///     }
169/// }
170///
171/// # fn process_a() {}
172///
173/// // measure time by explicitly setting a std::timer::Duration
174/// # use std::time::Instant;
175/// let start_time = Instant::now();
176/// process_a();
177/// metric!(timer(MyTimer::ProcessA) = start_time.elapsed());
178///
179/// // provide tags to a timer
180/// metric!(
181///     timer(MyTimer::ProcessA) = start_time.elapsed(),
182///     server = "server1",
183///     host = "host1",
184/// );
185///
186/// // measure time implicitly by enclosing a code block in a metric
187/// metric!(timer(MyTimer::ProcessA), {
188///     process_a();
189/// });
190///
191/// // measure block and also provide tags
192/// metric!(
193///     timer(MyTimer::ProcessB),
194///     server = "server1",
195///     host = "host1",
196///     {
197///         process_a();
198///     }
199/// );
200/// ```
201pub trait TimerMetric {
202    /// Returns the timer metric name that will be sent to statsd.
203    fn name(&self) -> &'static str;
204}
205
206/// A metric for capturing counters.
207///
208/// Counters are simple values incremented or decremented by a client. The rates at which these
209/// events occur or average values will be determined by the server receiving them. Examples of
210/// counter uses include number of logins to a system or requests received.
211///
212/// ## Example
213///
214/// ```
215/// use relay_statsd::{metric, CounterMetric};
216///
217/// enum MyCounter {
218///     TotalRequests,
219///     TotalBytes,
220/// }
221///
222/// impl CounterMetric for MyCounter {
223///     fn name(&self) -> &'static str {
224///         match self {
225///             Self::TotalRequests => "total_requests",
226///             Self::TotalBytes => "total_bytes",
227///         }
228///     }
229/// }
230///
231/// # let buffer = &[(), ()];
232///
233/// // add to the counter
234/// metric!(counter(MyCounter::TotalRequests) += 1);
235/// metric!(counter(MyCounter::TotalBytes) += buffer.len() as u64);
236///
237/// // add to the counter and provide tags
238/// metric!(
239///     counter(MyCounter::TotalRequests) += 1,
240///     server = "s1",
241///     host = "h1"
242/// );
243/// ```
244pub trait CounterMetric {
245    /// Returns the counter metric name that will be sent to statsd.
246    fn name(&self) -> &'static str;
247}
248
249/// A metric for capturing distributions.
250///
251/// A distribution is often similar to timers. Distributions can be thought of as a
252/// more general (not limited to timing things) form of timers.
253///
254/// ## Example
255///
256/// ```
257/// use relay_statsd::{metric, DistributionMetric};
258///
259/// struct QueueSize;
260///
261/// impl DistributionMetric for QueueSize {
262///     fn name(&self) -> &'static str {
263///         "queue_size"
264///     }
265/// }
266///
267/// # use std::collections::VecDeque;
268/// let queue = VecDeque::new();
269/// # let _hint: &VecDeque<()> = &queue;
270///
271/// // record a distribution value (uses global sample rate)
272/// metric!(distribution(QueueSize) = queue.len() as u64);
273///
274/// // record with tags
275/// metric!(
276///     distribution(QueueSize) = queue.len() as u64,
277///     server = "server1",
278///     host = "host1",
279/// );
280/// ```
281pub trait DistributionMetric {
282    /// Returns the distribution metric name that will be sent to statsd.
283    fn name(&self) -> &'static str;
284}
285
286/// A metric for capturing gauges.
287///
288/// Gauge values are an instantaneous measurement of a value determined by the client. They do not
289/// change unless changed by the client. Examples include things like load average or how many
290/// connections are active.
291///
292/// ## Example
293///
294/// ```
295/// use relay_statsd::{metric, GaugeMetric};
296///
297/// struct QueueSize;
298///
299/// impl GaugeMetric for QueueSize {
300///     fn name(&self) -> &'static str {
301///         "queue_size"
302///     }
303/// }
304///
305/// # use std::collections::VecDeque;
306/// let queue = VecDeque::new();
307/// # let _hint: &VecDeque<()> = &queue;
308///
309/// // a simple gauge value
310/// metric!(gauge(QueueSize) = queue.len() as u64);
311///
312/// // a gauge with tags
313/// metric!(
314///     gauge(QueueSize) = queue.len() as u64,
315///     server = "server1",
316///     host = "host1"
317/// );
318///
319/// // subtract from the gauge
320/// metric!(gauge(QueueSize) -= 1);
321///
322/// // subtract from the gauge and provide tags
323/// metric!(
324///     gauge(QueueSize) -= 1,
325///     server = "s1",
326///     host = "h1"
327/// );
328/// ```
329pub trait GaugeMetric {
330    /// Returns the gauge metric name that will be sent to statsd.
331    fn name(&self) -> &'static str;
332}
333
334#[doc(hidden)]
335#[macro_export]
336macro_rules! key_var {
337    ($id:expr $(,)*) => {{
338        let name = $crate::_metrics::KeyName::from_const_str($id);
339        $crate::_metrics::Key::from_static_labels(name, &[])
340    }};
341    ($id:expr $(, $k:expr => $v:expr)* $(,)?) => {{
342        let name = $crate::_metrics::KeyName::from_const_str($id);
343        let labels = ::std::vec![
344            $($crate::_metrics::Label::new(
345                $crate::_metrics::SharedString::const_str($k),
346                $crate::_metrics::SharedString::from_owned($v.into())
347            )),*
348        ];
349
350        $crate::_metrics::Key::from_parts(name, labels)
351    }};
352}
353
354/// Emits a metric.
355///
356/// See [crate-level documentation](self) for examples.
357#[macro_export]
358macro_rules! metric {
359    // counter increment
360    (counter($id:expr) += $value:expr $(, $($k:ident).* = $v:expr)* $(,)?) => {{
361        match $value {
362            value if value != 0 => {
363                let key = $crate::key_var!($crate::CounterMetric::name(&$id) $(, stringify!($($k).*) => $v)*);
364                let metadata = $crate::_metrics::metadata_var!(::std::module_path!(), $crate::_metrics::Level::INFO);
365                $crate::_metrics::with_recorder(|recorder| recorder.register_counter(&key, metadata))
366                    .increment(value);
367            }
368            _ => {}
369        }
370    }};
371
372    // gauge set
373    (gauge($id:expr) = $value:expr $(, $($k:ident).* = $v:expr)* $(,)?) => {{
374        let key = $crate::key_var!($crate::GaugeMetric::name(&$id) $(, stringify!($($k).*) => $v)*);
375        let metadata = $crate::_metrics::metadata_var!(::std::module_path!(), $crate::_metrics::Level::INFO);
376        $crate::_metrics::with_recorder(|recorder| recorder.register_gauge(&key, metadata))
377            .set($value as f64);
378    }};
379    // gauge increment
380    (gauge($id:expr) += $value:expr $(, $($k:ident).* = $v:expr)* $(,)?) => {{
381        let key = $crate::key_var!($crate::GaugeMetric::name(&$id) $(, stringify!($($k).*) => $v)*);
382        let metadata = $crate::_metrics::metadata_var!(::std::module_path!(), $crate::_metrics::Level::INFO);
383        $crate::_metrics::with_recorder(|recorder| recorder.register_gauge(&key, metadata))
384            .increment($value as f64);
385    }};
386    // gauge decrement
387    (gauge($id:expr) -= $value:expr $(, $($k:ident).* = $v:expr)* $(,)?) => {{
388        let key = $crate::key_var!($crate::GaugeMetric::name(&$id) $(, stringify!($($k).*) => $v)*);
389        let metadata = $crate::_metrics::metadata_var!(::std::module_path!(), $crate::_metrics::Level::INFO);
390        $crate::_metrics::with_recorder(|recorder| recorder.register_gauge(&key, metadata))
391            .decrement($value as f64);
392    }};
393
394    // distribution
395    (distribution($id:expr) = $value:expr $(, $($k:ident).* = $v:expr)* $(,)?) => {{
396        let key = $crate::key_var!($crate::DistributionMetric::name(&$id) $(, stringify!($($k).*) => $v)*);
397        let metadata = $crate::_metrics::metadata_var!(::std::module_path!(), $crate::_metrics::Level::INFO);
398        $crate::_metrics::with_recorder(|recorder| recorder.register_histogram(&key, metadata))
399            .record($value as f64);
400    }};
401
402    // timer value
403    (timer($id:expr) = $value:expr $(, $($k:ident).* = $v:expr)* $(,)?) => {{
404        let key = $crate::key_var!($crate::TimerMetric::name(&$id) $(, stringify!($($k).*) => $v)*);
405        let metadata = $crate::_metrics::metadata_var!(::std::module_path!(), $crate::_metrics::Level::INFO);
406        $crate::_metrics::with_recorder(|recorder| recorder.register_histogram(&key, metadata))
407            .record($value.as_nanos() as f64 / 1e6);
408    }};
409
410    // timed block
411    (timer($id:expr), $($($k:ident).* = $v:expr,)* $block:block) => {{
412        struct TimerOnDrop {
413            start: std::time::Instant,
414            key: $crate::_metrics::Key,
415        }
416        impl Drop for TimerOnDrop {
417            fn drop(&mut self) {
418                let elapsed = self.start.elapsed();
419                let metadata = $crate::_metrics::metadata_var!(::std::module_path!(), $crate::_metrics::Level::INFO);
420                 $crate::_metrics::with_recorder(|recorder| recorder.register_histogram(&self.key, metadata))
421                    .record(elapsed.as_nanos() as f64 / 1e6);
422            }
423        }
424
425        let _timer = TimerOnDrop {
426            start: std::time::Instant::now(),
427            key: $crate::key_var!($crate::TimerMetric::name(&$id) $(, stringify!($($k).*) => $v)*),
428        };
429        {$block}
430    }};
431}
432
433#[cfg(test)]
434mod tests {
435    use std::time::Duration;
436
437    use super::*;
438
439    enum TestGauges {
440        Foo,
441        Bar,
442    }
443
444    impl GaugeMetric for TestGauges {
445        fn name(&self) -> &'static str {
446            match self {
447                Self::Foo => "foo",
448                Self::Bar => "bar",
449            }
450        }
451    }
452
453    struct TestCounter;
454
455    impl CounterMetric for TestCounter {
456        fn name(&self) -> &'static str {
457            "counter"
458        }
459    }
460
461    struct TestDistribution;
462
463    impl DistributionMetric for TestDistribution {
464        fn name(&self) -> &'static str {
465            "distribution"
466        }
467    }
468
469    struct TestTimer;
470
471    impl TimerMetric for TestTimer {
472        fn name(&self) -> &'static str {
473            "timer"
474        }
475    }
476
477    #[test]
478    fn test_capturing_client() {
479        let captures = with_capturing_test_client(|| {
480            metric!(
481                gauge(TestGauges::Foo) = 123,
482                server = "server1",
483                host = "host1"
484            );
485            metric!(
486                gauge(TestGauges::Bar) = 456,
487                server = "server2",
488                host = "host2"
489            );
490        });
491
492        assert_eq!(
493            captures,
494            [
495                "foo:123|g|#server:server1,host:host1",
496                "bar:456|g|#server:server2,host:host2"
497            ]
498        )
499    }
500
501    #[test]
502    fn test_counter_tags_with_dots() {
503        let captures = with_capturing_test_client(|| {
504            metric!(
505                counter(TestCounter) += 10,
506                hc.project_id = "567",
507                server = "server1",
508            );
509            metric!(
510                counter(TestCounter) += 5,
511                hc.project_id = "567",
512                server = "server1",
513            );
514        });
515        assert_eq!(
516            captures,
517            [
518                "counter:10|c|#hc.project_id:567,server:server1",
519                "counter:5|c|#hc.project_id:567,server:server1"
520            ]
521        );
522    }
523
524    #[test]
525    fn test_gauge_tags_with_dots() {
526        let captures = with_capturing_test_client(|| {
527            metric!(
528                gauge(TestGauges::Foo) = 123,
529                hc.project_id = "567",
530                server = "server1",
531            );
532        });
533        assert_eq!(captures, ["foo:123|g|#hc.project_id:567,server:server1"]);
534    }
535
536    #[test]
537    fn test_distribution_tags_with_dots() {
538        let captures = with_capturing_test_client(|| {
539            metric!(
540                distribution(TestDistribution) = 123,
541                hc.project_id = "567",
542                server = "server1",
543            );
544        });
545        assert_eq!(
546            captures,
547            ["distribution:123|d|#hc.project_id:567,server:server1"]
548        );
549    }
550
551    #[test]
552    fn test_timer_tags_with_dots() {
553        let captures = with_capturing_test_client(|| {
554            let duration = Duration::from_secs(100);
555            metric!(
556                timer(TestTimer) = duration,
557                hc.project_id = "567",
558                server = "server1",
559            );
560        });
561        assert_eq!(
562            captures,
563            ["timer:100000|d|#hc.project_id:567,server:server1"]
564        );
565    }
566
567    #[test]
568    fn test_timed_block_tags_with_dots() {
569        let captures = with_capturing_test_client(|| {
570            metric!(
571                timer(TestTimer),
572                hc.project_id = "567",
573                server = "server1",
574                {
575                    // your code could be here
576                }
577            )
578        });
579        // just check the tags to not make this flaky
580        assert!(captures[0].ends_with("|d|#hc.project_id:567,server:server1"));
581    }
582
583    #[test]
584    fn test_timer_nanos_rounding_error() {
585        let one_day = Duration::from_secs(60 * 60 * 24);
586        let captures = with_capturing_test_client(|| {
587            metric!(timer(TestTimer) = one_day + Duration::from_nanos(1),);
588        });
589
590        // for "short" durations, precision is preserved:
591        assert_eq!(captures, ["timer:86400000.000001|d|#"]);
592
593        let one_year = Duration::from_secs(60 * 60 * 24 * 365);
594        let captures = with_capturing_test_client(|| {
595            metric!(timer(TestTimer) = one_year + Duration::from_nanos(1),);
596        });
597
598        // for very long durations, precision is lost:
599        assert_eq!(captures, ["timer:31536000000|d|#"]);
600    }
601
602    #[test]
603    fn test_timer_early_return() {
604        #[expect(unreachable_code)]
605        let captures = with_capturing_test_client(|| {
606            metric!(timer(TestTimer), {
607                // An early return still needs to capture the timer.
608                return;
609            });
610            // This shouldn't be reached due to the `return`.
611            unreachable!()
612        });
613        assert_eq!(captures.len(), 1);
614    }
615}