Skip to main content

vector/sinks/splunk_hec/metrics/
config.rs

1use std::sync::Arc;
2
3use futures_util::FutureExt;
4use tower::ServiceBuilder;
5use vector_lib::{
6    configurable::configurable_component, lookup::lookup_v2::OptionalValuePath,
7    sensitive_string::SensitiveString, sink::VectorSink,
8};
9
10use super::{request_builder::HecMetricsRequestBuilder, sink::HecMetricsSink};
11use crate::{
12    config::{AcknowledgementsConfig, GenerateConfig, Input, SinkConfig, SinkContext},
13    http::HttpClient,
14    sinks::{
15        Healthcheck,
16        splunk_hec::common::{
17            EndpointTarget, SplunkHecDefaultBatchSettings,
18            acknowledgements::HecClientAcknowledgementsConfig,
19            build_healthcheck, build_http_batch_service, config_host_key, create_client,
20            service::{HecService, HttpRequestBuilder},
21        },
22        util::{
23            BatchConfig, Compression, ServiceBuilderExt, TowerRequestConfig, http::HttpRetryLogic,
24        },
25    },
26    template::Template,
27    tls::TlsConfig,
28};
29
30/// Configuration of the `splunk_hec_metrics` sink.
31#[configurable_component(sink(
32    "splunk_hec_metrics",
33    "Deliver metric data to Splunk's HTTP Event Collector."
34))]
35#[derive(Clone, Debug)]
36#[serde(deny_unknown_fields)]
37pub struct HecMetricsSinkConfig {
38    /// Sets the default namespace for any metrics sent.
39    ///
40    /// This namespace is only used if a metric has no existing namespace. When a namespace is
41    /// present, it is used as a prefix to the metric name, and separated with a period (`.`).
42    #[configurable(metadata(docs::examples = "service"))]
43    pub default_namespace: Option<String>,
44
45    /// Default Splunk HEC token.
46    ///
47    /// If an event has a token set in its metadata, it prevails over the one set here.
48    #[serde(alias = "token")]
49    #[configurable(metadata(
50        docs::examples = "${SPLUNK_HEC_TOKEN}",
51        docs::examples = "A94A8FE5CCB19BA61C4C08"
52    ))]
53    pub default_token: SensitiveString,
54
55    /// The base URL of the Splunk instance.
56    ///
57    /// The scheme (`http` or `https`) must be specified. No path should be included since the paths defined
58    /// by the [`Splunk`][splunk] API are used.
59    ///
60    /// [splunk]: https://docs.splunk.com/Documentation/Splunk/8.0.0/Data/HECRESTendpoints
61    #[configurable(metadata(
62        docs::examples = "https://http-inputs-hec.splunkcloud.com",
63        docs::examples = "https://hec.splunk.com:8088",
64        docs::examples = "http://example.com"
65    ))]
66    #[configurable(validation(format = "uri"))]
67    pub endpoint: String,
68
69    /// Overrides the name of the log field used to retrieve the hostname to send to Splunk HEC.
70    ///
71    /// By default, the [global `log_schema.host_key` option][global_host_key] is used.
72    ///
73    /// [global_host_key]: https://vector.dev/docs/reference/configuration/global-options/#log_schema.host_key
74    #[configurable(metadata(docs::advanced))]
75    #[serde(default = "config_host_key")]
76    pub host_key: OptionalValuePath,
77
78    /// The name of the index where to send the events to.
79    ///
80    /// If not specified, the default index defined within Splunk is used.
81    #[configurable(metadata(docs::examples = "{{ host }}", docs::examples = "custom_index"))]
82    pub index: Option<Template>,
83
84    /// The sourcetype of events sent to this sink.
85    ///
86    /// If unset, Splunk defaults to `httpevent`.
87    #[configurable(metadata(docs::advanced))]
88    #[configurable(metadata(docs::examples = "{{ sourcetype }}", docs::examples = "_json",))]
89    pub sourcetype: Option<Template>,
90
91    /// The source of events sent to this sink.
92    ///
93    /// This is typically the filename the logs originated from.
94    ///
95    /// If unset, the Splunk collector sets it.
96    #[configurable(metadata(docs::advanced))]
97    #[configurable(metadata(
98        docs::examples = "{{ file }}",
99        docs::examples = "/var/log/syslog",
100        docs::examples = "UDP:514"
101    ))]
102    pub source: Option<Template>,
103
104    #[configurable(derived)]
105    #[serde(default)]
106    pub compression: Compression,
107
108    #[configurable(derived)]
109    #[serde(default)]
110    pub batch: BatchConfig<SplunkHecDefaultBatchSettings>,
111
112    #[configurable(derived)]
113    #[serde(default)]
114    pub request: TowerRequestConfig,
115
116    #[configurable(derived)]
117    pub tls: Option<TlsConfig>,
118
119    #[configurable(derived)]
120    #[serde(default)]
121    pub acknowledgements: HecClientAcknowledgementsConfig,
122
123    #[configurable(derived)]
124    #[serde(flatten)]
125    pub confinement: crate::template::ConfinementConfig,
126}
127
128impl GenerateConfig for HecMetricsSinkConfig {
129    fn generate_config() -> toml::Value {
130        toml::Value::try_from(Self {
131            default_namespace: None,
132            default_token: "${VECTOR_SPLUNK_HEC_TOKEN}".to_owned().into(),
133            endpoint: "http://localhost:8088".to_owned(),
134            host_key: config_host_key(),
135            index: None,
136            sourcetype: None,
137            source: None,
138            compression: Compression::default(),
139            batch: BatchConfig::default(),
140            request: TowerRequestConfig::default(),
141            tls: None,
142            acknowledgements: Default::default(),
143            confinement: Default::default(),
144        })
145        .unwrap()
146    }
147}
148
149#[async_trait::async_trait]
150#[typetag::serde(name = "splunk_hec_metrics")]
151impl SinkConfig for HecMetricsSinkConfig {
152    async fn build(&self, cx: SinkContext) -> crate::Result<(VectorSink, Healthcheck)> {
153        let mut config = self.clone();
154        config.sourcetype = config
155            .sourcetype
156            .map(|t| t.confine(&self.confinement, Self::NAME, "sourcetype"))
157            .transpose()?;
158        config.source = config
159            .source
160            .map(|t| t.confine(&self.confinement, Self::NAME, "source"))
161            .transpose()?;
162        config.index = config
163            .index
164            .map(|t| t.confine(&self.confinement, Self::NAME, "index"))
165            .transpose()?;
166
167        let client = create_client(self.tls.as_ref(), cx.proxy())?;
168        let healthcheck = build_healthcheck(
169            self.endpoint.clone(),
170            self.default_token.inner().to_owned(),
171            client.clone(),
172        )
173        .boxed();
174        let sink = config.build_processor(client, cx)?;
175        self.confinement.set_confinement_gauge("sink", Self::NAME);
176        Ok((sink, healthcheck))
177    }
178
179    fn input(&self) -> Input {
180        Input::metric()
181    }
182
183    fn acknowledgements(&self) -> &AcknowledgementsConfig {
184        &self.acknowledgements.inner
185    }
186}
187
188impl HecMetricsSinkConfig {
189    pub fn build_processor(&self, client: HttpClient, _: SinkContext) -> crate::Result<VectorSink> {
190        let ack_client = if self.acknowledgements.indexer_acknowledgements_enabled {
191            Some(client.clone())
192        } else {
193            None
194        };
195
196        let request_builder = HecMetricsRequestBuilder {
197            compression: self.compression,
198        };
199
200        let request_settings = self.request.into_settings();
201        let http_request_builder = Arc::new(HttpRequestBuilder::new(
202            self.endpoint.clone(),
203            EndpointTarget::default(),
204            self.default_token.inner().to_owned(),
205            self.compression,
206        ));
207        let http_service = ServiceBuilder::new()
208            .settings(request_settings, HttpRetryLogic::default())
209            .service(build_http_batch_service(
210                client,
211                Arc::clone(&http_request_builder),
212                EndpointTarget::Event,
213                false,
214            ));
215
216        let service = HecService::new(
217            http_service,
218            ack_client,
219            http_request_builder,
220            self.acknowledgements.clone(),
221        );
222
223        let batch_settings = self.batch.into_batcher_settings()?;
224
225        let sink = HecMetricsSink {
226            service,
227            batch_settings,
228            request_builder,
229            sourcetype: self.sourcetype.clone(),
230            source: self.source.clone(),
231            index: self.index.clone(),
232            host_key: self.host_key.path.clone(),
233            default_namespace: self.default_namespace.clone(),
234        };
235
236        Ok(VectorSink::from_event_streamsink(sink))
237    }
238}