vector/sinks/splunk_hec/metrics/
config.rs1use 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#[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 #[configurable(metadata(docs::examples = "service"))]
43 pub default_namespace: Option<String>,
44
45 #[serde(alias = "token")]
49 #[configurable(metadata(
50 docs::examples = "${SPLUNK_HEC_TOKEN}",
51 docs::examples = "A94A8FE5CCB19BA61C4C08"
52 ))]
53 pub default_token: SensitiveString,
54
55 #[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 #[configurable(metadata(docs::advanced))]
75 #[serde(default = "config_host_key")]
76 pub host_key: OptionalValuePath,
77
78 #[configurable(metadata(docs::examples = "{{ host }}", docs::examples = "custom_index"))]
82 pub index: Option<Template>,
83
84 #[configurable(metadata(docs::advanced))]
88 #[configurable(metadata(docs::examples = "{{ sourcetype }}", docs::examples = "_json",))]
89 pub sourcetype: Option<Template>,
90
91 #[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}