vector/sinks/mqtt/
config.rs1use std::time::Duration;
2
3use rand::Rng;
4use rumqttc::{MqttOptions, QoS, TlsConfiguration, Transport};
5use snafu::ResultExt;
6use vector_lib::codecs::JsonSerializerConfig;
7
8use crate::{
9 codecs::EncodingConfig,
10 common::mqtt::{
11 ConfigurationError, ConfigurationSnafu, MqttCommonConfig, MqttConnector, MqttError,
12 TlsSnafu,
13 },
14 config::{AcknowledgementsConfig, Input, SinkConfig, SinkContext},
15 sinks::{Healthcheck, VectorSink, mqtt::sink::MqttSink, prelude::*},
16 template::{ConfinementConfig, Template},
17 tls::MaybeTlsSettings,
18};
19
20#[configurable_component(sink("mqtt"))]
22#[derive(Clone, Debug)]
23pub struct MqttSinkConfig {
24 #[serde(flatten)]
25 pub common: MqttCommonConfig,
26
27 #[serde(default = "default_clean_session")]
29 pub clean_session: bool,
30
31 pub topic: Template,
33
34 #[serde(default = "default_retain")]
36 pub retain: bool,
37
38 #[configurable(derived)]
39 pub encoding: EncodingConfig,
40
41 #[configurable(derived)]
42 #[serde(
43 default,
44 deserialize_with = "crate::serde::bool_or_struct",
45 skip_serializing_if = "crate::serde::is_default"
46 )]
47 pub acknowledgements: AcknowledgementsConfig,
48
49 #[configurable(derived)]
50 #[serde(default = "default_qos")]
51 pub quality_of_service: MqttQoS,
52
53 #[configurable(derived)]
54 #[serde(flatten)]
55 pub confinement: ConfinementConfig,
56}
57
58#[configurable_component]
60#[derive(Clone, Copy, Debug, Default)]
61#[serde(rename_all = "lowercase")]
62#[allow(clippy::enum_variant_names)]
63pub enum MqttQoS {
64 #[default]
66 AtLeastOnce,
67
68 AtMostOnce,
70
71 ExactlyOnce,
73}
74
75impl From<MqttQoS> for QoS {
76 fn from(value: MqttQoS) -> Self {
77 match value {
78 MqttQoS::AtLeastOnce => QoS::AtLeastOnce,
79 MqttQoS::AtMostOnce => QoS::AtMostOnce,
80 MqttQoS::ExactlyOnce => QoS::ExactlyOnce,
81 }
82 }
83}
84
85const fn default_clean_session() -> bool {
86 false
87}
88
89const fn default_qos() -> MqttQoS {
90 MqttQoS::AtLeastOnce
91}
92
93const fn default_retain() -> bool {
94 false
95}
96
97impl Default for MqttSinkConfig {
98 fn default() -> Self {
99 Self {
100 common: MqttCommonConfig::default(),
101 clean_session: default_clean_session(),
102
103 topic: Template::try_from("vector").expect("Cannot parse as a template"),
104 retain: default_retain(),
105 encoding: JsonSerializerConfig::default().into(),
106 acknowledgements: AcknowledgementsConfig::default(),
107 quality_of_service: MqttQoS::default(),
108 confinement: ConfinementConfig::default(),
109 }
110 }
111}
112
113impl_generate_config_from_default!(MqttSinkConfig);
114
115#[async_trait::async_trait]
116#[typetag::serde(name = "mqtt")]
117impl SinkConfig for MqttSinkConfig {
118 async fn build(&self, _cx: SinkContext) -> crate::Result<(VectorSink, Healthcheck)> {
119 let mut config = self.clone();
120 config.topic = config
121 .topic
122 .confine(&self.confinement, Self::NAME, "topic")?;
123 let connector = config.build_connector()?;
124 let sink = MqttSink::new(&config, connector.clone())?;
125
126 self.confinement.set_confinement_gauge("sink", Self::NAME);
127 Ok((
128 VectorSink::from_event_streamsink(sink),
129 Box::pin(async move { connector.healthcheck().await }),
130 ))
131 }
132
133 fn input(&self) -> Input {
134 Input::log()
135 }
136
137 fn acknowledgements(&self) -> &AcknowledgementsConfig {
138 &self.acknowledgements
139 }
140}
141
142impl MqttSinkConfig {
143 fn build_connector(&self) -> Result<MqttConnector, MqttError> {
144 let client_id = self.common.client_id.clone().unwrap_or_else(|| {
145 let hash = rand::rng()
146 .sample_iter(&rand_distr::Alphanumeric)
147 .take(6)
148 .map(char::from)
149 .collect::<String>();
150 format!("vectorSink{hash}")
151 });
152
153 if client_id.is_empty() {
154 return Err(ConfigurationError::EmptyClientId).context(ConfigurationSnafu);
155 }
156 let tls =
157 MaybeTlsSettings::from_config(self.common.tls.as_ref(), false).context(TlsSnafu)?;
158 let mut options = MqttOptions::new(&client_id, &self.common.host, self.common.port);
159 options.set_keep_alive(Duration::from_secs(self.common.keep_alive.into()));
160 options.set_max_packet_size(self.common.max_packet_size, self.common.max_packet_size);
161 options.set_clean_session(self.clean_session);
162 match (&self.common.user, &self.common.password) {
163 (Some(user), Some(password)) => {
164 options.set_credentials(user, password);
165 }
166 (None, None) => {}
167 _ => {
168 return Err(MqttError::Configuration {
169 source: ConfigurationError::InvalidCredentials,
170 });
171 }
172 }
173 if let Some(tls) = tls.tls() {
174 let ca = tls.authorities_pem().flatten().collect();
175 let client_auth = tls.identity_pem();
176 let alpn = Some(vec!["mqtt".into()]);
177 options.set_transport(Transport::Tls(TlsConfiguration::Simple {
178 ca,
179 client_auth,
180 alpn,
181 }));
182 }
183 Ok(MqttConnector::new(options))
184 }
185}
186
187#[cfg(test)]
188mod test {
189 use super::*;
190 use crate::template::{ConfinementConfig, Template};
191
192 #[test]
193 fn generate_config() {
194 crate::test_util::test_generate_config::<MqttSinkConfig>();
195 }
196
197 #[test]
198 fn confinement_rejects_unconfined_topic() {
199 let template = Template::try_from("{{ topic }}").unwrap();
200 let config = ConfinementConfig::default();
201 let result = template.confine(&config, "mqtt", "topic");
202 assert!(result.is_err());
203 }
204
205 #[test]
206 fn confinement_opt_out_allows_unconfined_topic() {
207 let template = Template::try_from("{{ topic }}").unwrap();
208 let config = ConfinementConfig {
209 dangerously_allow_unconfined_template_resolution: true,
210 };
211 let result = template.confine(&config, "mqtt", "topic");
212 assert!(result.is_ok());
213 }
214
215 #[test]
216 fn confinement_allows_prefixed_topic() {
217 let template = Template::try_from("events-{{ env }}").unwrap();
218 let config = ConfinementConfig::default();
219 let result = template.confine(&config, "mqtt", "topic");
220 assert!(result.is_ok());
221 }
222}