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}