ygndotgg commented on code in PR #4303: URL: https://github.com/apache/iggy/pull/4303#discussion_r4162920926
########## core/connectors/sources/mqtt_source/src/lib.rs: ########## @@ -0,0 +1,1177 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +mod driver; + +use async_trait::async_trait; +use base64::Engine; +use humantime::Duration as HumanDuration; +use iggy_common::{HeaderKey, HeaderValue}; +use iggy_connector_sdk::{ + ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source, + source::SourceBatchResult, source_connector, +}; +use secrecy::SecretString; +use serde::{Deserialize, Serialize}; +use std::{ + collections::{BTreeMap, HashSet}, + fmt, + str::FromStr, + time::Duration, +}; +use tokio::{sync::Mutex, time::sleep}; +use tracing::{debug, error, info, warn}; +use url::Url; + +use driver::{AckToken, MqttDriver}; +use driver::{Mqtt5EnvelopeProperties, MqttMessage}; + +source_connector!(MqttSource); + +const CONNECTOR_NAME: &str = "MQTT source"; +const DEFAULT_KEEP_ALIVE: &str = "30s"; +const DEFAULT_POLL_TIMEOUT: &str = "1s"; +const DEFAULT_REQUEST_CAPACITY: usize = 32; +const DEFAULT_BATCH_SIZE: usize = 100; +const DEFAULT_BATCH_TIMEOUT: &str = "10ms"; +const DEFAULT_MAX_RETRIES: u32 = 5; + +/// MQTT wire protocol selected for the broker connection. +#[derive(Debug, Clone, Copy, Default, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "lowercase")] +pub enum MqttProtocol { + #[serde(rename = "mqtt311")] + Mqtt311, + #[serde(rename = "mqtt5")] + #[default] + Mqtt5, +} + +/// Subscription delivery level requested from the broker. +/// +/// QoS 0 has no acknowledgement token. QoS 1 and QoS 2 produce a token that +/// remains pending until the corresponding Iggy batch is acknowledged. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum Qos { + Zero, + One, + Two, +} + +impl TryFrom<u8> for Qos { + type Error = Error; + + fn try_from(value: u8) -> Result<Self, Self::Error> { + match value { + 0 => Ok(Self::Zero), + 1 => Ok(Self::One), + 2 => Ok(Self::Two), + value => Err(Error::InvalidConfigValue(format!( + "qos {value} is unsupported; expected 0, 1, or 2" + ))), + } + } +} + +#[derive(Debug, Deserialize)] +pub struct MqttSourceConfig { + /// MQTT broker URL, using `mqtt://`, `mqtts://`, or `ssl://`. + pub broker_url: String, + /// Topic filters subscribed to by this connector instance. + pub subscriptions: Vec<String>, + /// Exact topic-filter QoS overrides applied on top of `qos`. + #[serde(default)] + pub subscription_qos: BTreeMap<String, u8>, + /// MQTT protocol version used for the connection. + #[serde(default)] + pub protocol: MqttProtocol, + /// Preserves MQTT 5 extended properties in a JSON payload envelope. + #[serde(default)] + pub include_metadata: bool, + /// Default subscription QoS when no per-filter override exists. + #[serde(default = "default_qos")] + pub qos: u8, + /// Optional CA and client-authentication material for TLS connections. + #[serde(default)] + pub tls: Option<MqttTlsConfig>, + /// Explicit broker client ID. A connector-specific ID is generated when absent. + pub client_id: Option<String>, + /// Broker username, which must be paired with `password`. + pub username: Option<String>, + /// Broker password, kept secret in memory and never included in debug output. + pub password: Option<SecretString>, + /// Whether the broker should discard the previous session on connect. + #[serde(default)] + pub clean_start: bool, + /// MQTT 5 session expiry interval in seconds. + pub session_expiry_interval: Option<u32>, Review Comment: Defaulted session_expiry_interval to 3600 seconds so MQTT 5 preserves unacknowledged QoS messages across restarts when clean_start = false. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
