Skip to main content

vector/sinks/mqtt/
config.rs

1use 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/// Configuration for the `mqtt` sink
21#[configurable_component(sink("mqtt"))]
22#[derive(Clone, Debug)]
23pub struct MqttSinkConfig {
24    #[serde(flatten)]
25    pub common: MqttCommonConfig,
26
27    /// If set to true, the MQTT session is cleaned on login.
28    #[serde(default = "default_clean_session")]
29    pub clean_session: bool,
30
31    /// MQTT publish topic (templates allowed)
32    pub topic: Template,
33
34    /// Whether the messages should be retained by the server
35    #[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/// Supported Quality of Service types for MQTT.
59#[configurable_component]
60#[derive(Clone, Copy, Debug, Default)]
61#[serde(rename_all = "lowercase")]
62#[allow(clippy::enum_variant_names)]
63pub enum MqttQoS {
64    /// AtLeastOnce.
65    #[default]
66    AtLeastOnce,
67
68    /// AtMostOnce.
69    AtMostOnce,
70
71    /// ExactlyOnce.
72    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}