ygndotgg commented on code in PR #4303:
URL: https://github.com/apache/iggy/pull/4303#discussion_r4158425362


##########
core/connectors/sources/mqtt_source/src/driver.rs:
##########
@@ -0,0 +1,1157 @@
+// 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.
+
+use super::{MqttProtocol, MqttSourceConfig, Qos, qos_for_subscription};
+use iggy_common::{HeaderKey, HeaderValue};
+use rumqttc::tokio_rustls::rustls::{self, ClientConfig, RootCertStore};
+use rumqttc::v5::{
+    AsyncClient as Mqtt5Client, Event as Mqtt5Event, EventLoop as 
Mqtt5EventLoop,
+    Incoming as Mqtt5Incoming, MqttOptions as Mqtt5Options,
+    mqttbytes::{QoS as Mqtt5Qos, v5::Publish as Mqtt5Publish},
+};
+use rumqttc::{
+    AsyncClient as Mqtt311Client, Event as Mqtt311Event, EventLoop as 
Mqtt311EventLoop,
+    Incoming as Mqtt311Incoming, MqttOptions as Mqtt311Options, 
TlsConfiguration, Transport,
+    mqttbytes::{QoS as Mqtt311Qos, v4::Publish as Mqtt311Publish},
+};
+use rustls_native_certs::load_native_certs;
+use rustls_pemfile::{certs, private_key};
+use secrecy::ExposeSecret;
+use std::{
+    collections::{BTreeMap, VecDeque},
+    io::{BufReader, Cursor},
+    sync::Arc,
+    time::Duration,
+};
+use tokio::time::timeout;
+use url::Url;
+
+const ACK_RETRY_ATTEMPTS: usize = 5;
+const ACK_RETRY_DELAY: Duration = Duration::from_millis(10);
+
+// Metadata is kept separate from the payload so protocol details can be
+// preserved in Iggy headers without changing the application bytes.
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub(crate) struct MqttMessageMetadata {
+    pub(crate) qos: Qos,
+    pub(crate) packet_id: Option<u16>,
+    pub(crate) dup: bool,
+    pub(crate) retain: bool,
+}
+
+/// MQTT payload plus normalized topic metadata ready for Iggy headers.
+#[derive(Debug)]
+pub(crate) struct MqttMessage {
+    pub(crate) payload: Vec<u8>,
+    pub(crate) topic: String,
+    pub(crate) headers: BTreeMap<HeaderKey, HeaderValue>,
+    pub(crate) metadata: MqttMessageMetadata,
+}
+
+/// A normalized MQTT message and its optional deferred acknowledgement token.
+#[derive(Debug)]
+pub(crate) struct ReceivedMessage {
+    pub(crate) message: MqttMessage,
+    pub(crate) ack_token: Option<AckToken>,
+}
+
+/// The original publish packet needed by rumqttc to acknowledge QoS 1 or QoS 
2.
+#[derive(Debug)]
+pub(crate) struct AckToken(AckTokenKind);
+
+#[derive(Debug)]
+enum AckTokenKind {
+    Mqtt311(Mqtt311Publish),
+    Mqtt5(Mqtt5Publish),
+}
+
+enum MqttConnection {
+    Mqtt311 {
+        client: Box<Mqtt311Client>,
+        event_loop: Box<Mqtt311EventLoop>,
+    },
+    Mqtt5 {
+        client: Box<Mqtt5Client>,
+        event_loop: Box<Mqtt5EventLoop>,
+    },
+}
+
+/// Protocol-specific MQTT client and event loop behind one common driver API.
+pub(crate) struct MqttDriver {
+    // Messages received while flushing acknowledgements are retained here 
rather
+    // than dropped. The source batch size bounds this queue.
+    connection: MqttConnection,
+    buffered_messages: VecDeque<ReceivedMessage>,
+}
+
+impl MqttDriver {
+    pub(crate) async fn connect(
+        id: u32,
+        config: &MqttSourceConfig,
+        qos: Qos,
+        keep_alive: Duration,
+        poll_timeout: Duration,
+        request_capacity: usize,
+    ) -> Result<Self, iggy_connector_sdk::Error> {
+        install_rustls_provider();
+        let broker_url = broker_url_with_client_id(config, id)?;
+        match config.protocol {
+            MqttProtocol::Mqtt311 => {
+                let mut options = 
Mqtt311Options::parse_url(&broker_url).map_err(|error| {
+                    
iggy_connector_sdk::Error::InvalidConfigValue(format!("broker_url: {error}"))
+                })?;
+                options
+                    .set_client_id(client_id(config, id))
+                    .set_clean_session(config.clean_start)
+                    .set_keep_alive(keep_alive)
+                    .set_request_channel_capacity(request_capacity)
+                    // Manual acknowledgements let Iggy persistence happen 
before
+                    // PUBACK/PUBREC is sent to the broker.
+                    .set_manual_acks(true);
+                set_mqtt311_credentials(&mut options, config);
+                if let Some(transport) = tls_transport(config)? {
+                    options.set_transport(transport);
+                }
+
+                let (client, mut event_loop) = Mqtt311Client::new(options, 
request_capacity);
+                for topic in &config.subscriptions {
+                    // Each filter may select its own QoS override; the global
+                    // value is used only when no exact override is configured.
+                    let subscription_qos = qos_for_subscription(config, topic, 
qos)?;
+                    client
+                        .subscribe(topic, subscription_qos.into())
+                        .await
+                        .map_err(|error| {
+                            
iggy_connector_sdk::Error::Connection(error.to_string())
+                        })?;
+                }
+                poll_mqtt311(&mut event_loop, poll_timeout).await?;
+
+                Ok(Self {
+                    connection: MqttConnection::Mqtt311 {
+                        client: Box::new(client),
+                        event_loop: Box::new(event_loop),
+                    },
+                    buffered_messages: VecDeque::new(),
+                })
+            }
+            MqttProtocol::Mqtt5 => {
+                let mut options = 
Mqtt5Options::parse_url(&broker_url).map_err(|error| {
+                    
iggy_connector_sdk::Error::InvalidConfigValue(format!("broker_url: {error}"))
+                })?;
+                options
+                    .set_client_id(client_id(config, id))
+                    .set_clean_start(config.clean_start)
+                    .set_keep_alive(keep_alive)
+                    .set_request_channel_capacity(request_capacity)
+                    // MQTT 5 uses the same deferred-acknowledgement lifecycle;
+                    // rumqttc maps the token to PUBREC/PUBREL/PUBCOMP 
internally.
+                    .set_manual_acks(true)
+                    
.set_session_expiry_interval(config.session_expiry_interval);
+                set_mqtt5_credentials(&mut options, config);
+                if let Some(transport) = tls_transport(config)? {
+                    options.set_transport(transport);
+                }
+
+                let (client, mut event_loop) = Mqtt5Client::new(options, 
request_capacity);
+                for topic in &config.subscriptions {
+                    // Keep subscription QoS resolution identical across 
protocol
+                    // versions so the configuration has one predictable 
meaning.
+                    let subscription_qos = qos_for_subscription(config, topic, 
qos)?;
+                    client
+                        .subscribe(topic, subscription_qos.into())
+                        .await
+                        .map_err(|error| {
+                            
iggy_connector_sdk::Error::Connection(error.to_string())
+                        })?;
+                }
+                poll_mqtt5(&mut event_loop, poll_timeout).await?;
+
+                Ok(Self {
+                    connection: MqttConnection::Mqtt5 {
+                        client: Box::new(client),
+                        event_loop: Box::new(event_loop),
+                    },
+                    buffered_messages: VecDeque::new(),
+                })
+            }
+        }
+    }
+
+    pub(crate) async fn next_message(
+        &mut self,
+        poll_timeout: Duration,
+    ) -> Result<Option<ReceivedMessage>, iggy_connector_sdk::Error> {
+        // Consume buffered messages first. They arrived while another 
message's
+        // acknowledgement was being flushed and are already valid source data.
+        if let Some(message) = self.buffered_messages.pop_front() {
+            return Ok(Some(message));
+        }
+
+        // Bound each driver poll so the source can flush partial batches and
+        // respond to runtime shutdown instead of waiting indefinitely.
+        let deadline = tokio::time::Instant::now() + poll_timeout;
+        loop {
+            let remaining = 
deadline.saturating_duration_since(tokio::time::Instant::now());
+            if remaining.is_zero() {
+                return Ok(None);
+            }
+            let event = match &mut self.connection {
+                MqttConnection::Mqtt311 { event_loop, .. } => {
+                    match timeout(remaining, event_loop.poll()).await {
+                        Ok(Ok(event)) => match event {
+                            
Mqtt311Event::Incoming(Mqtt311Incoming::Publish(publish)) => {
+                                Some(normalize_mqtt311(publish)?)
+                            }
+                            // Connection and subscription events advance the
+                            // event loop but are not source messages.
+                            Mqtt311Event::Outgoing(_) | 
Mqtt311Event::Incoming(_) => None,
+                        },
+                        Ok(Err(error)) => {
+                            return 
Err(iggy_connector_sdk::Error::Connection(error.to_string()));
+                        }
+                        Err(_) => return Ok(None),
+                    }
+                }
+                MqttConnection::Mqtt5 { event_loop, .. } => {
+                    match timeout(remaining, event_loop.poll()).await {
+                        Ok(Ok(event)) => match event {
+                            
Mqtt5Event::Incoming(Mqtt5Incoming::Publish(publish)) => {
+                                Some(normalize_mqtt5(publish)?)
+                            }
+                            Mqtt5Event::Outgoing(_) | Mqtt5Event::Incoming(_) 
=> None,
+                        },
+                        Ok(Err(error)) => {
+                            return 
Err(iggy_connector_sdk::Error::Connection(error.to_string()));
+                        }
+                        Err(_) => return Ok(None),
+                    }
+                }
+            };
+            if event.is_some() {
+                return Ok(event);
+            }
+        }
+    }
+
+    pub(crate) async fn acknowledge_batch(
+        &mut self,
+        ack_tokens: &mut Vec<AckToken>,
+        poll_timeout: Duration,
+        max_buffered_messages: usize,
+    ) -> Result<(), iggy_connector_sdk::Error> {
+        // Process tokens in order. On failure, remove only tokens already
+        // acknowledged so the remaining suffix can be retried.
+        let mut acknowledged = 0;
+        let retry_delay = ACK_RETRY_DELAY.min(poll_timeout);
+        while acknowledged < ack_tokens.len() {
+            let mut last_error = None;
+            let mut acknowledged_token = false;
+            for attempt in 0..=ACK_RETRY_ATTEMPTS {
+                match self.try_acknowledge(&ack_tokens[acknowledged]) {
+                    Ok(()) => {
+                        acknowledged_token = true;
+                        break;
+                    }
+                    Err(error) => {
+                        last_error = Some(error);
+                        if attempt == ACK_RETRY_ATTEMPTS {
+                            break;
+                        }
+                        // rumqttc may need event-loop progress before try_ack 
can
+                        // enqueue the acknowledgement.
+                        if let Err(error) = self
+                            .poll_for_ack_progress(poll_timeout, 
max_buffered_messages)
+                            .await
+                        {
+                            retain_unacknowledged_tokens(ack_tokens, 
acknowledged);
+                            return Err(error);
+                        }
+                        if !retry_delay.is_zero() {
+                            tokio::time::sleep(retry_delay).await;
+                        }
+                    }
+                }
+            }
+
+            if !acknowledged_token {
+                retain_unacknowledged_tokens(ack_tokens, acknowledged);
+                return 
Err(last_error.unwrap_or(iggy_connector_sdk::Error::InvalidState));
+            }
+            acknowledged += 1;
+        }
+        ack_tokens.clear();
+        Ok(())
+    }
+
+    fn try_acknowledge(&self, ack_token: &AckToken) -> Result<(), 
iggy_connector_sdk::Error> {
+        // A token must be acknowledged by the same protocol client that 
created
+        // it. Mixing MQTT 3.1.1 and MQTT 5 tokens is an invalid internal 
state.
+        match (&self.connection, &ack_token.0) {
+            (MqttConnection::Mqtt311 { client, .. }, 
AckTokenKind::Mqtt311(publish)) => client
+                .try_ack(publish)
+                .map_err(|error| 
iggy_connector_sdk::Error::Connection(error.to_string())),
+            (MqttConnection::Mqtt5 { client, .. }, 
AckTokenKind::Mqtt5(publish)) => client
+                .try_ack(publish)
+                .map_err(|error| 
iggy_connector_sdk::Error::Connection(error.to_string())),
+            _ => Err(iggy_connector_sdk::Error::InvalidState),
+        }
+    }
+
+    async fn poll_for_ack_progress(
+        &mut self,
+        poll_timeout: Duration,
+        max_buffered_messages: usize,
+    ) -> Result<(), iggy_connector_sdk::Error> {
+        // A publish can arrive while the event loop is being driven for an 
ACK.
+        // Preserve it for the next source poll, subject to the batch bound.
+        let received = match &mut self.connection {
+            MqttConnection::Mqtt311 { event_loop, .. } => {
+                let event = timeout(poll_timeout, event_loop.poll())
+                    .await
+                    .map_err(|_| {
+                        iggy_connector_sdk::Error::Connection(
+                            "timed out while flushing MQTT 
acknowledgements".to_string(),
+                        )
+                    })?
+                    .map_err(|error| 
iggy_connector_sdk::Error::Connection(error.to_string()))?;
+                match event {
+                    Mqtt311Event::Incoming(Mqtt311Incoming::Publish(publish)) 
=> {
+                        Some(normalize_mqtt311(publish)?)
+                    }
+                    Mqtt311Event::Outgoing(_) | Mqtt311Event::Incoming(_) => 
None,
+                }
+            }
+            MqttConnection::Mqtt5 { event_loop, .. } => {
+                let event = timeout(poll_timeout, event_loop.poll())
+                    .await
+                    .map_err(|_| {
+                        iggy_connector_sdk::Error::Connection(
+                            "timed out while flushing MQTT 
acknowledgements".to_string(),
+                        )
+                    })?
+                    .map_err(|error| 
iggy_connector_sdk::Error::Connection(error.to_string()))?;
+                match event {
+                    Mqtt5Event::Incoming(Mqtt5Incoming::Publish(publish)) => {
+                        Some(normalize_mqtt5(publish)?)
+                    }
+                    Mqtt5Event::Outgoing(_) | Mqtt5Event::Incoming(_) => None,
+                }
+            }
+        };
+
+        if let Some(received) = received {
+            if self.buffered_messages.len() >= max_buffered_messages {
+                return Err(iggy_connector_sdk::Error::Connection(
+                    "MQTT message buffer reached batch_size while flushing 
acknowledgements"
+                        .to_string(),
+                ));
+            }
+            self.buffered_messages.push_back(received);
+        }
+        Ok(())
+    }
+}
+
+fn retain_unacknowledged_tokens(ack_tokens: &mut Vec<AckToken>, acknowledged: 
usize) {
+    // The prefix has already been accepted by the broker. Keep the suffix so a
+    // retry does not repeat successful acknowledgements or discard failures.
+    ack_tokens.drain(..acknowledged);
+}
+
+fn install_rustls_provider() {
+    let _ = rustls::crypto::ring::default_provider().install_default();
+}
+
+fn tls_transport(
+    config: &MqttSourceConfig,
+) -> Result<Option<Transport>, iggy_connector_sdk::Error> {
+    // No TLS table means plaintext MQTT. When present, build a complete rustls
+    // transport from either the supplied CA or the system trust store.
+    let Some(tls) = &config.tls else {
+        return Ok(None);
+    };
+
+    let ca = tls
+        .ca_file
+        .as_deref()
+        .map(|path| read_certificate_file(path, "tls.ca_file"))
+        .transpose()?;
+    let client_auth = match (&tls.client_cert_file, &tls.client_key_file) {
+        (Some(cert_path), Some(key_path)) => Some((
+            read_certificate_file(cert_path, "tls.client_cert_file")?,
+            read_private_key_file(key_path, "tls.client_key_file")?,
+        )),
+        (None, None) => None,
+        _ => unreachable!("TLS client certificate and key validation occurs 
before connect"),
+    };
+
+    let tls_config = TlsConfiguration::Rustls(Arc::new(build_tls_config(ca, 
client_auth)?));
+
+    Ok(Some(Transport::tls_with_config(tls_config)))
+}
+
+fn read_file(path: &str, field: &str) -> Result<Vec<u8>, 
iggy_connector_sdk::Error> {
+    std::fs::read(path).map_err(|error| {
+        iggy_connector_sdk::Error::InvalidConfigValue(format!("{field} could 
not be read: {error}"))
+    })
+}
+
+fn read_certificate_file(path: &str, field: &str) -> Result<Vec<u8>, 
iggy_connector_sdk::Error> {
+    // Validate PEM contents during initialization so bad certificates produce 
an
+    // actionable configuration error before the broker connection starts.
+    let contents = read_file(path, field)?;
+    let mut reader = BufReader::new(Cursor::new(&contents));
+    let certificates = certs(&mut reader)
+        .collect::<Result<Vec<_>, _>>()
+        .map_err(|error| {
+            iggy_connector_sdk::Error::InvalidConfigValue(format!("{field} is 
invalid: {error}"))
+        })?;
+    if certificates.is_empty() {
+        return Err(iggy_connector_sdk::Error::InvalidConfigValue(format!(
+            "{field} does not contain a PEM certificate"
+        )));
+    }
+    Ok(contents)
+}
+
+fn read_private_key_file(path: &str, field: &str) -> Result<Vec<u8>, 
iggy_connector_sdk::Error> {
+    let contents = read_file(path, field)?;
+    let mut reader = BufReader::new(Cursor::new(&contents));
+    if private_key(&mut reader)
+        .map_err(|error| {
+            iggy_connector_sdk::Error::InvalidConfigValue(format!("{field} is 
invalid: {error}"))
+        })?
+        .is_none()
+    {
+        return Err(iggy_connector_sdk::Error::InvalidConfigValue(format!(
+            "{field} does not contain a PEM private key"
+        )));
+    }
+    Ok(contents)
+}
+
+fn system_root_cert_store() -> Result<RootCertStore, 
iggy_connector_sdk::Error> {
+    // This path is used only when TLS is enabled without a custom CA bundle.
+    let mut root_cert_store = RootCertStore::empty();
+    let native_certificates = load_native_certs();
+    for certificate in native_certificates.certs {
+        root_cert_store.add(certificate).map_err(|error| {
+            iggy_connector_sdk::Error::InvalidConfigValue(format!(
+                "system TLS certificate is invalid: {error}"
+            ))
+        })?;
+    }
+    if root_cert_store.is_empty() {
+        return Err(iggy_connector_sdk::Error::InvalidConfigValue(
+            "no system TLS root certificates were found".to_string(),
+        ));
+    }
+    Ok(root_cert_store)
+}
+
+fn build_tls_config(
+    ca: Option<Vec<u8>>,
+    client_auth: Option<(Vec<u8>, Vec<u8>)>,
+) -> Result<ClientConfig, iggy_connector_sdk::Error> {
+    // Client authentication is optional, but certificate and key parsing stays
+    // here so mismatches fail before rumqttc starts its reconnect loop.
+    let root_cert_store = match ca {
+        Some(ca) => {
+            let mut reader = BufReader::new(Cursor::new(ca));
+            let certificates =
+                certs(&mut reader)
+                    .collect::<Result<Vec<_>, _>>()
+                    .map_err(|error| {
+                        iggy_connector_sdk::Error::InvalidConfigValue(format!(
+                            "tls.ca_file is invalid: {error}"
+                        ))
+                    })?;
+            let mut root_cert_store = RootCertStore::empty();
+            for certificate in certificates {
+                root_cert_store.add(certificate).map_err(|error| {
+                    iggy_connector_sdk::Error::InvalidConfigValue(format!(
+                        "tls.ca_file contains an invalid certificate: {error}"
+                    ))
+                })?;
+            }
+            if root_cert_store.is_empty() {
+                return Err(iggy_connector_sdk::Error::InvalidConfigValue(
+                    "tls.ca_file does not contain a usable 
certificate".to_string(),
+                ));
+            }
+            root_cert_store
+        }
+        None => system_root_cert_store()?,
+    };
+    let client_config = 
ClientConfig::builder().with_root_certificates(root_cert_store);
+
+    let Some((client_certificate, client_key)) = client_auth else {
+        return Ok(client_config.with_no_client_auth());
+    };
+    let mut reader = BufReader::new(Cursor::new(client_certificate));
+    let certificates = certs(&mut reader)
+        .collect::<Result<Vec<_>, _>>()
+        .map_err(|error| {
+            iggy_connector_sdk::Error::InvalidConfigValue(format!(
+                "tls.client_cert_file is invalid: {error}"
+            ))
+        })?;
+    let mut key_reader = BufReader::new(Cursor::new(client_key));
+    let key = private_key(&mut key_reader)
+        .map_err(|error| {
+            iggy_connector_sdk::Error::InvalidConfigValue(format!(
+                "tls.client_key_file is invalid: {error}"
+            ))
+        })?
+        .ok_or_else(|| {
+            iggy_connector_sdk::Error::InvalidConfigValue(
+                "tls.client_key_file does not contain a PEM private 
key".to_string(),
+            )
+        })?;
+
+    client_config
+        .with_client_auth_cert(certificates, key)
+        .map_err(|error| {
+            iggy_connector_sdk::Error::InvalidConfigValue(format!(
+                "TLS client certificate and key do not match: {error}"
+            ))
+        })
+}
+
+fn broker_url_with_client_id(
+    config: &MqttSourceConfig,
+    id: u32,
+) -> Result<String, iggy_connector_sdk::Error> {
+    // rumqttc reads the client ID from the URL query when parsing options. 
Add a
+    // generated value only when the operator did not configure one explicitly.
+    let mut url = Url::parse(&config.broker_url).map_err(|error| {
+        iggy_connector_sdk::Error::InvalidConfigValue(format!("broker_url: 
{error}"))
+    })?;
+    let existing_query = url
+        .query_pairs()
+        .filter(|(key, _)| key != "client_id")
+        .map(|(key, value)| (key.into_owned(), value.into_owned()))
+        .collect::<Vec<_>>();
+    let mut query = url::form_urlencoded::Serializer::new(String::new());
+    for (key, value) in existing_query {
+        query.append_pair(&key, &value);
+    }
+    query.append_pair("client_id", &client_id(config, id));
+    url.set_query(Some(&query.finish()));
+    Ok(url.into())
+}
+
+fn client_id(config: &MqttSourceConfig, id: u32) -> String {
+    config
+        .client_id
+        .clone()
+        .unwrap_or_else(|| format!("iggy-mqtt-source-{id}"))
+}
+
+fn set_mqtt311_credentials(options: &mut Mqtt311Options, config: 
&MqttSourceConfig) {
+    if let (Some(username), Some(password)) = (&config.username, 
&config.password) {
+        options.set_credentials(username.clone(), 
password.expose_secret().to_string());
+    }
+}
+
+fn set_mqtt5_credentials(options: &mut Mqtt5Options, config: 
&MqttSourceConfig) {
+    if let (Some(username), Some(password)) = (&config.username, 
&config.password) {
+        options.set_credentials(username.clone(), 
password.expose_secret().to_string());
+    }
+}
+
+async fn poll_mqtt311(
+    event_loop: &mut Mqtt311EventLoop,
+    poll_timeout: Duration,
+) -> Result<(), iggy_connector_sdk::Error> {
+    timeout(poll_timeout, event_loop.poll())
+        .await
+        .map_err(|_| {
+            iggy_connector_sdk::Error::Connection(
+                "timed out connecting to MQTT 3.1.1 broker".to_string(),
+            )
+        })?
+        .map_err(|error| 
iggy_connector_sdk::Error::Connection(error.to_string()))?;
+    Ok(())
+}
+
+async fn poll_mqtt5(
+    event_loop: &mut Mqtt5EventLoop,
+    poll_timeout: Duration,
+) -> Result<(), iggy_connector_sdk::Error> {
+    timeout(poll_timeout, event_loop.poll())
+        .await
+        .map_err(|_| {
+            iggy_connector_sdk::Error::Connection(
+                "timed out connecting to MQTT 5 broker".to_string(),
+            )
+        })?
+        .map_err(|error| 
iggy_connector_sdk::Error::Connection(error.to_string()))?;
+    Ok(())
+}
+
+fn normalize_mqtt311(
+    publish: Mqtt311Publish,
+) -> Result<ReceivedMessage, iggy_connector_sdk::Error> {
+    // Normalize both MQTT protocol versions into one source representation. 
Only
+    // the acknowledgement token type remains protocol-specific.
+    let qos = publish.qos.into();
+    let metadata = MqttMessageMetadata {
+        qos,
+        packet_id: (qos != Qos::Zero).then_some(publish.pkid),
+        dup: publish.dup,
+        retain: publish.retain,
+    };
+    let headers = metadata_headers("mqtt311", &publish.topic, metadata)?;
+    let message = MqttMessage {
+        payload: publish.payload.to_vec(),
+        topic: publish.topic.clone(),
+        headers,
+        metadata,
+    };
+    let ack_token = (metadata.qos != 
Qos::Zero).then_some(AckToken(AckTokenKind::Mqtt311(publish)));
+    Ok(ReceivedMessage { message, ack_token })
+}
+
+pub(crate) fn normalize_mqtt5(
+    publish: Mqtt5Publish,
+) -> Result<ReceivedMessage, iggy_connector_sdk::Error> {
+    // MQTT 5 adds publish properties; common metadata headers are built first
+    // and the optional properties are appended afterward.
+    let qos = publish.qos.into();
+    let metadata = MqttMessageMetadata {
+        qos,
+        packet_id: (qos != Qos::Zero).then_some(publish.pkid),
+        dup: publish.dup,
+        retain: publish.retain,
+    };
+    let topic = topic_string(&publish.topic)?;
+    let mut headers = metadata_headers("mqtt5", &topic, metadata)?;
+    if let Some(properties) = publish.properties.as_ref() {
+        insert_mqtt5_properties(&mut headers, properties)?;
+    }
+    let message = MqttMessage {
+        payload: publish.payload.to_vec(),
+        topic,
+        headers,
+        metadata,
+    };
+    let ack_token = (metadata.qos != 
Qos::Zero).then_some(AckToken(AckTokenKind::Mqtt5(publish)));
+    Ok(ReceivedMessage { message, ack_token })
+}
+
+fn metadata_headers(
+    protocol: &str,
+    topic: &str,
+    metadata: MqttMessageMetadata,
+) -> Result<BTreeMap<HeaderKey, HeaderValue>, iggy_connector_sdk::Error> {
+    // Stable lowercase header names let consumers inspect MQTT metadata 
without
+    // decoding the original MQTT packet again.
+    let mut headers = BTreeMap::new();
+    insert_string_header(&mut headers, "mqtt.protocol", protocol)?;
+    insert_string_header(&mut headers, "mqtt.topic", topic)?;
+    headers.insert(header_key("mqtt.qos")?, qos_value(metadata.qos).into());
+    headers.insert(header_key("mqtt.dup")?, metadata.dup.into());
+    headers.insert(header_key("mqtt.retain")?, metadata.retain.into());
+    if let Some(packet_id) = metadata.packet_id {
+        headers.insert(header_key("mqtt.packet_id")?, packet_id.into());
+    }
+    Ok(headers)
+}
+
+fn insert_mqtt5_properties(
+    headers: &mut BTreeMap<HeaderKey, HeaderValue>,
+    properties: &rumqttc::v5::mqttbytes::v5::PublishProperties,
+) -> Result<(), iggy_connector_sdk::Error> {
+    // MQTT 5 properties are copied into headers because the payload remains 
the
+    // original application bytes. Repeated user properties use distinct keys.
+    if let Some(value) = properties.payload_format_indicator {
+        headers.insert(header_key("mqtt.payload_format_indicator")?, 
value.into());
+    }
+    if let Some(value) = properties.message_expiry_interval {
+        headers.insert(header_key("mqtt.message_expiry_interval")?, 
value.into());
+    }
+    if let Some(value) = properties.topic_alias {
+        headers.insert(header_key("mqtt.topic_alias")?, value.into());
+    }
+    if let Some(value) = properties.response_topic.as_deref() {
+        insert_string_header(headers, "mqtt.response_topic", value)?;
+    }
+    if let Some(value) = properties.correlation_data.as_ref() {
+        insert_raw_header(headers, "mqtt.correlation_data", value)?;
+    }
+    for (index, (key, value)) in properties.user_properties.iter().enumerate() 
{
+        insert_string_header(headers, 
&format!("mqtt.user_property.{index}.key"), key)?;
+        insert_string_header(headers, 
&format!("mqtt.user_property.{index}.value"), value)?;
+    }
+    for (index, value) in 
properties.subscription_identifiers.iter().enumerate() {
+        headers.insert(
+            header_key(&format!("mqtt.subscription_identifier.{index}"))?,
+            (*value as u64).into(),
+        );
+    }
+    if let Some(value) = properties.content_type.as_deref() {
+        insert_string_header(headers, "mqtt.content_type", value)?;
+    }
+    Ok(())
+}
+
+fn insert_string_header(

Review Comment:
   Fixed by ommiting empty or oversized iggy header values, log each omission, 
preserve the mqtt payload, and stored only common publish properties and when 
include_metadata = true, store the full mqtt 5 topic and publish properties in 
the JSON envelope. to preserve properties > 255 bytes



-- 
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]

Reply via email to