vector/sinks/splunk_hec/metrics/
sink.rs1use std::{fmt, sync::Arc};
2
3use serde::Serialize;
4use vector_lib::event::{Metric, MetricValue};
5use vrl::path::OwnedValuePath;
6
7use super::request_builder::HecMetricsRequestBuilder;
8use crate::{
9 internal_events::{SplunkInvalidMetricReceivedError, TemplateRenderingError},
10 sinks::{
11 prelude::*,
12 splunk_hec::common::request::HecRequest,
13 util::{encode_namespace, processed_event::ProcessedEvent},
14 },
15};
16
17pub struct HecMetricsSink<S> {
18 pub service: S,
19 pub batch_settings: BatcherSettings,
20 pub request_builder: HecMetricsRequestBuilder,
21 pub sourcetype: Option<ConfinedTemplate>,
22 pub source: Option<ConfinedTemplate>,
23 pub index: Option<ConfinedTemplate>,
24 pub host_key: Option<OwnedValuePath>,
25 pub default_namespace: Option<String>,
26}
27
28impl<S> HecMetricsSink<S>
29where
30 S: Service<HecRequest> + Send + 'static,
31 S::Future: Send + 'static,
32 S::Response: DriverResponse + Send + 'static,
33 S::Error: fmt::Debug + Into<crate::Error> + Send,
34{
35 async fn run_inner(self: Box<Self>, input: BoxStream<'_, Event>) -> Result<(), ()> {
36 let sourcetype = self.sourcetype.as_ref();
37 let source = self.source.as_ref();
38 let index = self.index.as_ref();
39 let host_key = self.host_key.as_ref();
40 let default_namespace = self.default_namespace.as_deref();
41 let batch_settings = self.batch_settings;
42
43 input
44 .map(|event| (event.size_of(), event.into_metric()))
45 .filter_map(move |(event_byte_size, metric)| {
46 future::ready(process_metric(
47 metric,
48 event_byte_size,
49 sourcetype,
50 source,
51 index,
52 host_key,
53 default_namespace,
54 ))
55 })
56 .batched_partitioned(EventPartitioner, batch_settings.timeout, |_| {
57 batch_settings.as_byte_size_config()
58 })
59 .request_builder(
60 default_request_builder_concurrency_limit(),
61 self.request_builder,
62 )
63 .filter_map(|request| async move {
64 match request {
65 Err(e) => {
66 error!("Failed to build HEC Metrics request: {:?}.", e);
67 None
68 }
69 Ok(req) => Some(req),
70 }
71 })
72 .into_driver(self.service)
73 .run()
74 .await
75 }
76}
77
78#[async_trait]
79impl<S> StreamSink<Event> for HecMetricsSink<S>
80where
81 S: Service<HecRequest> + Send + 'static,
82 S::Future: Send + 'static,
83 S::Response: DriverResponse + Send + 'static,
84 S::Error: fmt::Debug + Into<crate::Error> + Send,
85{
86 async fn run(self: Box<Self>, input: BoxStream<'_, Event>) -> Result<(), ()> {
87 self.run_inner(input).await
88 }
89}
90
91#[derive(Default)]
92struct EventPartitioner;
93
94impl Partitioner for EventPartitioner {
95 type Item = HecProcessedEvent;
96 type Key = Option<Arc<str>>;
97
98 fn partition(&self, item: &Self::Item) -> Self::Key {
99 item.event.metadata().splunk_hec_token()
100 }
101}
102
103#[derive(Serialize)]
104pub struct HecMetricsProcessedEventMetadata {
105 pub event_byte_size: usize,
106 pub sourcetype: Option<String>,
107 pub source: Option<String>,
108 pub index: Option<String>,
109 pub host: Option<String>,
110 pub metric_name: String,
111 pub metric_value: f64,
112}
113
114impl ByteSizeOf for HecMetricsProcessedEventMetadata {
115 fn allocated_bytes(&self) -> usize {
116 self.sourcetype.allocated_bytes()
117 + self.source.allocated_bytes()
118 + self.index.allocated_bytes()
119 + self.host.allocated_bytes()
120 + self.metric_name.allocated_bytes()
121 }
122}
123
124impl HecMetricsProcessedEventMetadata {
125 fn extract_metric_name(metric: &Metric, default_namespace: Option<&str>) -> String {
126 encode_namespace(metric.namespace().or(default_namespace), '.', metric.name())
127 }
128
129 fn extract_metric_value(metric: &Metric) -> Option<f64> {
130 match *metric.value() {
131 MetricValue::Counter { value } => Some(value),
132 MetricValue::Gauge { value } => Some(value),
133 _ => {
134 emit!(SplunkInvalidMetricReceivedError {
135 value: metric.value(),
136 kind: &metric.kind(),
137 error: "Metric kind not supported.".into(),
138 });
139 None
140 }
141 }
142 }
143}
144
145pub type HecProcessedEvent = ProcessedEvent<Metric, HecMetricsProcessedEventMetadata>;
146
147pub fn process_metric(
148 metric: Metric,
149 event_byte_size: usize,
150 sourcetype: Option<&ConfinedTemplate>,
151 source: Option<&ConfinedTemplate>,
152 index: Option<&ConfinedTemplate>,
153 host_key: Option<&OwnedValuePath>,
154 default_namespace: Option<&str>,
155) -> Option<HecProcessedEvent> {
156 let metric_name =
157 HecMetricsProcessedEventMetadata::extract_metric_name(&metric, default_namespace);
158 let metric_value = HecMetricsProcessedEventMetadata::extract_metric_value(&metric)?;
159
160 let render = |tpl: Option<&ConfinedTemplate>, field: &'static str| -> Option<Option<String>> {
164 let Some(t) = tpl else { return Some(None) };
165 match t.render_string(&metric) {
166 Ok(s) => Some(Some(s)),
167 Err(error @ crate::template::TemplateRenderingError::Confined { .. }) => {
168 emit!(TemplateRenderingError {
169 error,
170 field: Some(field),
171 drop_event: true,
172 });
173 None
174 }
175 Err(error) => {
176 emit!(TemplateRenderingError {
177 error,
178 field: Some(field),
179 drop_event: false,
180 });
181 Some(None)
182 }
183 }
184 };
185
186 let sourcetype = render(sourcetype, "sourcetype")?;
187 let source = render(source, "source")?;
188 let index = render(index, "index")?;
189 let host = host_key.and_then(|key| metric.tag_value(key.to_string().as_str()));
190
191 let metadata = HecMetricsProcessedEventMetadata {
192 event_byte_size,
193 sourcetype,
194 source,
195 index,
196 host,
197 metric_name,
198 metric_value,
199 };
200
201 Some(HecProcessedEvent {
202 event: metric,
203 metadata,
204 })
205}
206
207impl EventCount for HecProcessedEvent {
208 fn event_count(&self) -> usize {
209 1
211 }
212}