github-actions[bot] commented on code in PR #4303:
URL: https://github.com/apache/iggy/pull/4303#discussion_r4158587801


##########
core/server/config.toml:
##########
@@ -67,7 +67,7 @@ max_request_size = "2 MB"
 # a warning will be logged at startup but the server will continue to run.
 # `true` enables the embedded Web UI (requires server built with 'iggy-web' 
feature).
 # `false` disables the embedded Web UI (default).
-web_ui = false
+web_ui = true

Review Comment:
   warning: `web_ui` flips to `true` in the shipped server configuration, which 
is unrelated to the MQTT source and contradicts the `false` default on line 69. 
Revert it to `false`, or move the change to its own pull request.



##########
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:
   warning: `session_expiry_interval` defaults to `None`, and MQTT 5 reads an 
absent value as zero. A disconnect then discards the session and its 
unacknowledged QoS 1/2 messages, although `clean_start = false`. Default the 
field to a non-zero value, because `README.md:403` promises redelivery after a 
restart.



##########
core/connectors/sources/mqtt_source/src/driver.rs:
##########
@@ -0,0 +1,1226 @@
+// 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 base64::Engine;
+use iggy_common::{HeaderKey, HeaderValue};
+use rumqttc::tokio_rustls::rustls::{
+    self, ClientConfig, RootCertStore,
+    pki_types::{CertificateDer, PrivateKeyDer},
+};
+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 serde::Serialize;
+use std::{
+    collections::{BTreeMap, VecDeque},
+    io::{BufReader, Cursor},
+    sync::{Arc, LazyLock},
+    time::Duration,
+};
+use tokio::time::timeout;
+use tracing::warn;
+use url::Url;
+
+const ACK_RETRY_DELAY: Duration = Duration::from_millis(10);
+const MAX_HEADER_VALUE_LENGTH: usize = 255;
+const MQTT_PROTOCOL_HEADER: &str = "mqtt.protocol";
+const MQTT_TOPIC_HEADER: &str = "mqtt.topic";
+
+static MQTT_PROTOCOL_HEADER_KEY: LazyLock<HeaderKey> = LazyLock::new(|| {
+    HeaderKey::try_from(MQTT_PROTOCOL_HEADER).expect("MQTT protocol header key 
is valid")
+});
+static MQTT_TOPIC_HEADER_KEY: LazyLock<HeaderKey> = LazyLock::new(|| {
+    HeaderKey::try_from(MQTT_TOPIC_HEADER).expect("MQTT topic header key is 
valid")
+});
+static MQTT_QOS_HEADER_KEY: LazyLock<HeaderKey> =
+    LazyLock::new(|| HeaderKey::try_from("mqtt.qos").expect("MQTT QoS header 
key is valid"));
+static MQTT_DUP_HEADER_KEY: LazyLock<HeaderKey> =
+    LazyLock::new(|| HeaderKey::try_from("mqtt.dup").expect("MQTT duplicate 
header key is valid"));
+static MQTT_RETAIN_HEADER_KEY: LazyLock<HeaderKey> =
+    LazyLock::new(|| HeaderKey::try_from("mqtt.retain").expect("MQTT retain 
header key is valid"));
+#[cfg(test)]
+static MQTT_PACKET_ID_HEADER_KEY: LazyLock<HeaderKey> = LazyLock::new(|| {
+    HeaderKey::try_from("mqtt.packet_id").expect("MQTT packet ID header key is 
valid")
+});
+#[cfg(test)]
+static MQTT_RESPONSE_TOPIC_HEADER_KEY: LazyLock<HeaderKey> = LazyLock::new(|| {
+    HeaderKey::try_from("mqtt.response_topic").expect("MQTT response topic 
header key is valid")
+});
+
+// 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,
+    pub(crate) mqtt5_properties: Option<Mqtt5EnvelopeProperties>,
+}
+
+#[derive(Debug, Clone, Serialize)]
+pub(crate) struct Mqtt5EnvelopeProperties {
+    pub(crate) payload_format_indicator: Option<u8>,
+    pub(crate) message_expiry_interval: Option<u32>,
+    pub(crate) topic_alias: Option<u16>,
+    pub(crate) response_topic: Option<String>,
+    pub(crate) correlation_data_base64: Option<String>,
+    pub(crate) user_properties: Vec<Mqtt5UserProperty>,
+    pub(crate) subscription_identifiers: Vec<usize>,
+    pub(crate) content_type: Option<String>,
+}
+
+#[derive(Debug, Clone, Serialize)]
+pub(crate) struct Mqtt5UserProperty {
+    pub(crate) key: String,
+    pub(crate) value: String,
+}
+
+/// 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())
+                        })?;
+                }
+                let buffered_messages =
+                    poll_mqtt311(&mut event_loop, config.subscriptions.len(), 
poll_timeout).await?;
+
+                Ok(Self {
+                    connection: MqttConnection::Mqtt311 {
+                        client: Box::new(client),
+                        event_loop: Box::new(event_loop),
+                    },
+                    buffered_messages,
+                })
+            }
+            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())
+                        })?;
+                }
+                let buffered_messages =
+                    poll_mqtt5(&mut event_loop, config.subscriptions.len(), 
poll_timeout).await?;
+
+                Ok(Self {
+                    connection: MqttConnection::Mqtt5 {
+                        client: Box::new(client),
+                        event_loop: Box::new(event_loop),
+                    },
+                    buffered_messages,
+                })
+            }
+        }
+    }
+
+    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,
+        max_retries: u32,
+    ) -> 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..=max_retries {
+                match self.try_acknowledge(&ack_tokens[acknowledged]) {
+                    Ok(()) => {
+                        acknowledged_token = true;
+                        break;
+                    }
+                    Err(error) => {
+                        last_error = Some(error);
+                        if attempt == max_retries {
+                            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 {

Review Comment:
   warning: `poll_for_ack_progress` drops the publish it just read when the 
buffer is full, so a QoS 0 payload is lost and the rest of the flush is 
abandoned. Push the message into `buffered_messages` before the error return.



##########
core/connectors/sources/mqtt_source/Cargo.toml:
##########
@@ -0,0 +1,56 @@
+# 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.
+
+[package]
+name = "iggy_connector_mqtt_source"
+version = "0.5.0-edge.4"

Review Comment:
   nit: The crate version is `0.5.0-edge.4` while every other connector crate 
uses `0.5.0`. Set `version = "0.5.0"` so the value reported by 
`iggy_source_version` matches its peers.



##########
core/connectors/sources/mqtt_source/README.md:
##########
@@ -0,0 +1,425 @@
+# Apache Iggy MQTT Source Connector (`iggy_connector_mqtt_source`)
+
+The **MQTT Source Connector** is a dynamically loaded shared library plugin 
(`.so` / `.dylib` / `.dll`) for Apache Iggy. Built using the 
`iggy_connector_sdk::source_connector!` C-FFI ABI and the 
[`rumqttc`](https://crates.io/crates/rumqttc) asynchronous client, it ingests 
telemetry and event streams from external MQTT brokers (such as **EMQX** or 
Mosquitto) directly into persistent **Apache Iggy** streams, topics, and 
partitions.
+
+---
+
+## Architecture Overview
+
+```text
+  [ IoT Edge Devices ]
+           │
+           │ MQTT 3.1.1 / MQTT 5 (Sensor Telemetry & Events)
+           ▼
+┌──────────────────────────────────────────────────────────────────────────┐
+│                           EMQX BROKER                                    │
+│  • Client Authentication, Authorizations, & Session Queues              │
+│  • Topic Subscriptions, Wildcards (+, #), & QoS Management              │
+└────────────────────────────────────┬─────────────────────────────────────┘
+                                     │
+                                     │ MQTT / TLS Connection
+                                     ▼
+┌──────────────────────────────────────────────────────────────────────────┐
+│                   IGGY CONNECTORS RUNTIME PROCESS                        │
+│                                                                          │
+│  ┌────────────────────────────────────────────────────────────────────┐  │
+│  │  MQTT Source Plugin (libiggy_connector_mqtt_source.so)             │  │
+│  │  • Embedded `rumqttc` Client (`AsyncClient` + `EventLoop`)         │  │
+│  │  • Normalizes MQTT Packets ──► Iggy Payload & Metadata Headers     │  │
+│  │  • Holds QoS 1/2 PUBACK until Iggy Confirms Persistence            │  │
+│  └─────────────────────────────────┬──────────────────────────────────┘  │
+│                                    │                                     │
+│                                    │ C-ABI FFI Boundary                  │
+│                                    │ (iggy_source_open, poll, ack, etc)  │
+│                                    ▼                                     │
+│  ┌────────────────────────────────────────────────────────────────────┐  │
+│  │  Iggy Connector Runtime Core                                       │  │
+│  │  • Schema Decoders / Encoders (`Schema::Raw`)                      │  │
+│  │  • Optional Field Transformations (Add/Delete/Rename)              │  │
+│  │  • State Checkpoint Storage (File / HTTP endpoint)                 │  │
+│  └─────────────────────────────────┬──────────────────────────────────┘  │
+└────────────────────────────────────┼─────────────────────────────────────┘
+                                     │
+                                     │ Native Binary Client Protocol
+                                     ▼ (TCP / QUIC)
+┌──────────────────────────────────────────────────────────────────────────┐
+│                            IGGY SERVER                                   │
+│  • Thread-Per-Core Shared-Nothing Storage Engine (`io_uring`)            │
+│  • Appends Messages to Stream ──► Topic ──► Partition Disk Logs          │
+└──────────────────────────────────────────────────────────────────────────┘
+```
+
+---
+
+## Features
+
+* **Multi-Protocol Support**: Supports **MQTT 3.1.1** (`mqtt311`) and **MQTT 
5** (`mqtt5`) protocols.
+* **Quality of Service (QoS)**: Ingests QoS 0, QoS 1, and QoS 2 messages.
+* **Deferred Acknowledgment (Manual ACKs)**: Configures `rumqttc` with manual 
acknowledgments enabled. For QoS 1/2 messages, broker acknowledgments (`PUBACK` 
/ `PUBREC`) are held pending and only sent after `iggy-server` confirms 
persistence.
+* **Per-Subscription QoS Overrides**: Configure a global fallback QoS 
alongside per-subscription overrides for specific topic filters.
+* **TLS & Mutual TLS (mTLS)**: Supports secure broker URLs (`mqtts://` / 
`ssl://`) with system trust roots or custom CA bundles (`ca_file`), client 
certificates (`client_cert_file`), and client keys (`client_key_file`) via 
Rustls.
+* **Source-Side Micro-Batching**: Micro-batches incoming messages up to 
`batch_size` or until `batch_timeout` fires, minimizing FFI serialization and 
network overhead.
+* **Stable MQTT Headers**: Preserves protocol, topic, QoS, retain, and dup 
metadata as Iggy message headers. MQTT 5 extended properties are ignored by 
default and are available through the optional metadata envelope.
+* **MQTT 5 Metadata Envelope**: Optionally stores MQTT 5 properties and the 
exact application payload in a versioned JSON envelope.
+* **Route Isolation**: Run multiple connector instances in parallel to route 
distinct MQTT topic filters into separate Iggy streams and topics.
+* **Process Isolation**: Operates inside the `iggy-connectors` process memory 
space via C-ABI FFI, keeping the core `iggy-server` decoupled and untouched.
+
+---
+
+## Configuration Architecture
+
+Configuration operates in **two distinct layers**:
+
+1. **Runtime Layer**: Configured in the connector TOML file (e.g., 
`mqtt_source.toml`). Defines the connector identity, shared library path, and 
destination Iggy stream/topic mapping.
+2. **Plugin Layer (`[plugin_config]`)**: Defines MQTT-specific connection 
parameters, credentials, subscriptions, QoS levels, and TLS configurations. The 
runtime passes this table to the plugin serialized as JSON across the FFI 
boundary.
+
+---
+
+## Configuration Options (`[plugin_config]`)
+
+| Field | Type | Required | Default | Description |
+| :--- | :--- | :---: | :--- | :--- |
+| `broker_url` | String | **Yes** | — | Broker URL scheme (`mqtt://`, 
`mqtts://`, or `ssl://`) and host/port. |
+| `subscriptions` | Array[String] | **Yes** | — | Non-empty list of unique 
MQTT topic filters (wildcards `+` and `#` supported). |
+| `protocol` | String | No | `"mqtt5"` | MQTT protocol version: `"mqtt5"` or 
`"mqtt311"`. |
+| `include_metadata` | Boolean | No | `false` | Enables the MQTT 5 metadata 
envelope. Requires `protocol = "mqtt5"` and a stream schema of `"json"`. |
+| `qos` | Integer | No | `1` | Default global subscription QoS (`0`, `1`, or 
`2`). |
+| `subscription_qos` | Table | No | `{}` | Map of topic-filter strings to 
explicit QoS overrides (`0`, `1`, or `2`). |
+| `client_id` | String | No | `iggy-mqtt-source-{id}` | MQTT client 
identifier. |
+| `username` | String | No | None | Broker authentication username (must be 
supplied together with `password`). |
+| `password` | String | No | None | Broker authentication password (wrapped as 
a secret, never logged). |
+| `clean_start` | Boolean | No | `false` | MQTT clean start flag (or clean 
session for MQTT 3.1.1). |
+| `session_expiry_interval` | Integer | No | None | MQTT 5 session expiry 
interval in seconds. |
+| `keep_alive` | Duration | No | `"30s"` | Ping interval string (minimum `1s` 
for MQTT 3.1.1, minimum `5s` for MQTT 5). |
+| `poll_timeout` | Duration | No | `"1s"` | Maximum duration spent waiting on 
network events per driver poll tick. |
+| `request_capacity` | Integer | No | `32` | Bounded capacity for `rumqttc` 
internal request channel (must be `> 0`). |
+| `batch_size` | Integer | No | `100` | Maximum messages accumulated into a 
single source batch. |
+| `batch_timeout` | Duration | No | `"10ms"` | Maximum wait duration after the 
first message arrives before flushing a batch. |
+| `max_retries` | Integer | No | `5` | Number of retries after the initial 
MQTT acknowledgement attempt. After retries are exhausted, the source abandons 
the broker acknowledgement and continues; the broker may redeliver the message 
later, producing duplicates. |
+| `verbose_logging` | Boolean | No | `false` | Enables additional debug 
logging inside the plugin driver. |
+| `tls.ca_file` | String | No | None | File path to custom CA root 
certificates in PEM format. |
+| `tls.client_cert_file` | String | No | None | File path to mTLS client 
certificate in PEM format (must pair with `client_key_file`). |
+| `tls.client_key_file` | String | No | None | File path to mTLS client 
private key in PEM format (must pair with `client_cert_file`). |
+
+---
+
+## Configuration Examples
+
+### 1. Basic Production Telemetry (QoS 1 with Credentials)
+
+```toml
+type = "source"
+key = "mqtt_telemetry"
+enabled = true
+version = 1
+name = "MQTT Telemetry Source"
+path = "target/release/libiggy_connector_mqtt_source"
+plugin_config_format = "json"
+
+[[streams]]
+stream = "iot"
+topic = "telemetry"
+schema = "raw"
+batch_length = 100
+linger_time = "5ms"
+
+[plugin_config]
+broker_url = "mqtt://127.0.0.1:1883"
+subscriptions = ["devices/+/telemetry"]
+protocol = "mqtt5"
+qos = 1
+client_id = "iggy-telemetry-source"
+username = "iggy_app"
+password = "local-secret-password"
+clean_start = false
+session_expiry_interval = 3600
+keep_alive = "30s"
+batch_size = 100
+batch_timeout = "10ms"
+max_retries = 5
+```
+
+### 2. MQTT 5 Metadata Envelope
+
+Use envelope mode when MQTT 5 publish properties must be preserved. This mode
+has three required rules:
+
+1. `protocol` must be `"mqtt5"`.
+2. `include_metadata` must be `true`.
+3. Every destination stream used by this connector must use `schema = "json"`.
+
+MQTT 3.1.1 cannot use envelope mode because it has no MQTT 5 extended
+PUBLISH properties. The connector rejects `include_metadata = true` with
+`protocol = "mqtt311"` during initialization.
+
+```toml
+type = "source"
+key = "mqtt5_metadata"
+enabled = true
+version = 1
+name = "MQTT 5 source with metadata"
+path = "target/release/libiggy_connector_mqtt_source"
+plugin_config_format = "json"
+
+[[streams]]
+stream = "iot"
+topic = "telemetry"
+schema = "json"
+batch_length = 100
+linger_time = "5ms"
+
+[plugin_config]
+broker_url = "mqtt://127.0.0.1:1883"
+subscriptions = ["devices/+/telemetry"]
+protocol = "mqtt5"
+include_metadata = true
+qos = 1
+client_id = "iggy-mqtt-source-metadata"
+clean_start = false
+session_expiry_interval = 3600
+keep_alive = "30s"
+batch_size = 100
+batch_timeout = "10ms"
+```
+
+The message payload becomes a JSON envelope containing `version`, `topic`,
+`properties`, and `payload_base64`. The original MQTT payload is recovered by
+Base64-decoding `payload_base64`. The stable routing headers remain present:
+`mqtt.protocol`, `mqtt.topic`, `mqtt.qos`, `mqtt.retain`, and `mqtt.dup`.
+
+### 3. Per-Subscription QoS Overrides
+
+```toml
+[plugin_config]
+broker_url = "mqtt://127.0.0.1:1883"
+subscriptions = [
+  "devices/+/telemetry",
+  "devices/+/alerts",
+  "devices/+/diagnostics"
+]
+qos = 1 # Fallback for devices/+/alerts
+
+[plugin_config.subscription_qos]
+"devices/+/telemetry" = 2   # High-importance telemetry via QoS 2
+"devices/+/diagnostics" = 0 # Disposable diagnostic metrics via QoS 0
+```
+
+### 4. Secure TLS & Mutual TLS (mTLS) Setup
+
+```toml
+[plugin_config]
+broker_url = "mqtts://emqx.example.com:8883"
+subscriptions = ["factory/+/metrics"]
+protocol = "mqtt5"
+qos = 1
+
+[plugin_config.tls]
+ca_file = "/etc/iggy/certs/ca.pem"
+client_cert_file = "/etc/iggy/certs/client-cert.pem"
+client_key_file = "/etc/iggy/certs/client-key.pem"
+```
+
+### 5. Multi-Instance Route Isolation
+
+To route distinct MQTT topic filters to different Iggy streams or topics, run
+one source connector instance per route. The runtime loads every connector TOML
+from its configured `config_dir`. For the standard example layout, create these
+files locally:
+
+`core/connectors/runtime/example_config/connectors/mqtt_site_a.toml`:
+
+```toml
+type = "source"
+key = "mqtt_site_a"
+enabled = true
+version = 1
+name = "MQTT source - site A"
+path = "target/release/libiggy_connector_mqtt_source"
+plugin_config_format = "json"
+
+[[streams]]
+stream = "site_a"
+topic = "telemetry"
+schema = "raw"
+batch_length = 100
+linger_time = "5ms"
+
+[plugin_config]
+broker_url = "mqtt://127.0.0.1:1883"
+subscriptions = ["devices/site-a/#"]
+protocol = "mqtt5"
+qos = 1
+client_id = "iggy-mqtt-source-site-a"
+clean_start = false
+session_expiry_interval = 3600
+keep_alive = "30s"
+poll_timeout = "1s"
+request_capacity = 32
+batch_size = 100
+batch_timeout = "10ms"
+verbose_logging = false
+```
+
+`core/connectors/runtime/example_config/connectors/mqtt_site_b.toml`:
+
+```toml
+type = "source"
+key = "mqtt_site_b"
+enabled = true
+version = 1
+name = "MQTT source - site B"
+path = "target/release/libiggy_connector_mqtt_source"
+plugin_config_format = "json"
+
+[[streams]]
+stream = "site_b"
+topic = "telemetry"
+schema = "raw"
+batch_length = 100
+linger_time = "5ms"
+
+[plugin_config]
+broker_url = "mqtt://127.0.0.1:1883"
+subscriptions = ["devices/site-b/#"]
+protocol = "mqtt5"
+qos = 1
+client_id = "iggy-mqtt-source-site-b"
+clean_start = false
+session_expiry_interval = 3600
+keep_alive = "30s"
+poll_timeout = "1s"
+request_capacity = 32
+batch_size = 100
+batch_timeout = "10ms"
+verbose_logging = false
+```
+
+Use unique connector keys and MQTT client IDs for every route. Keep topic
+filters non-overlapping unless duplicate delivery to multiple Iggy destinations
+is intentional. Build the plugin first, copy these files into the active
+connector `config_dir`, and restart `iggy-connectors`.
+
+---
+
+## Environment Variable Overrides
+
+Any property inside `[plugin_config]` can be overridden at runtime using 
environment variables without modifying configuration files:
+
+```bash
+export 
IGGY_CONNECTORS_SOURCE_MQTT_PLUGIN_CONFIG_BROKER_URL="mqtt://127.0.0.1:1883"
+export IGGY_CONNECTORS_SOURCE_MQTT_PLUGIN_CONFIG_USERNAME="emqx-user"
+export IGGY_CONNECTORS_SOURCE_MQTT_PLUGIN_CONFIG_PASSWORD="secret-password"
+export IGGY_CONNECTORS_SOURCE_MQTT_PLUGIN_CONFIG_QOS=1
+export 
IGGY_CONNECTORS_SOURCE_MQTT_PLUGIN_CONFIG_SUBSCRIPTIONS='["devices/+/telemetry"]'
+```
+
+---
+
+## Message Header Reference
+
+The connector attaches MQTT metadata directly to `ProducedMessage.headers`:
+
+The connector always writes the stable MQTT headers (`mqtt.protocol`,
+`mqtt.topic`, `mqtt.qos`, `mqtt.dup`, and `mqtt.retain`). MQTT 5 extended
+properties are not written as Iggy headers. They are ignored when
+`include_metadata = false` and stored in the metadata envelope when
+`include_metadata = true` with MQTT 5. MQTT 3.1.1 has no extended PUBLISH
+properties and remains in raw-payload mode.
+
+When metadata envelope mode is enabled, the stream schema must be `"json"`.
+The envelope contains the complete MQTT topic, all supported MQTT 5
+properties, and the original payload as Base64. The topic is intentionally
+available both in the envelope and as `mqtt.topic` for consumers that route by
+headers. Envelope mode is not available for MQTT 3.1.1.
+
+Iggy header values are limited to 255 bytes and cannot be empty. When a stable
+header value is empty or exceeds 255 bytes, the connector omits only that
+header, logs the omission with its name, reason, and byte length, and continues
+persisting the MQTT payload and other valid headers. This also applies to
+`mqtt.topic`; it is omitted when it cannot be represented by an Iggy header.
+
+| Header Key | Type | Description |
+| :--- | :--- | :--- |
+| `mqtt.protocol` | String | `"mqtt311"` or `"mqtt5"`. |
+| `mqtt.topic` | String | Topic filter on which the broker delivered the 
message. |
+| `mqtt.qos` | Integer | Delivered MQTT QoS (`0`, `1`, or `2`). |
+| `mqtt.dup` | Boolean | Duplicate delivery flag from broker. |
+| `mqtt.retain` | Boolean | Retained message flag. |
+
+---
+
+## Redelivery & Acknowledgment Mechanism
+
+The connector normally provides **At-Least-Once Delivery** for QoS 1 and QoS 2
+messages using manual protocol acknowledgments. `max_retries` bounds the
+acknowledgement attempts so a transient broker failure does not stop the 
source.
+After `max_retries` failures, the MQTT source abandons the broker
+acknowledgement and continues. The broker may redeliver the message later,
+producing duplicates.
+
+```text
+  MQTT Broker             MQTT Source Plugin            Connector Runtime      
     Iggy Server
+       │                         │                              │              
          │
+       │ 1. MQTT PUBLISH (QoS 1) │                              │              
          │
+       ├────────────────────────►│                              │              
          │
+       │                         │ [Message held pending]       │              
          │
+       │                         │ [NO PUBACK SENT YET]         │              
          │
+       │                         │                              │              
          │
+       │                         │ 2. Source::poll() batch      │              
          │
+       │                         ├─────────────────────────────►│              
          │
+       │                         │                              │ 3. Send over 
TCP/QUIC  │
+       │                         │                              
├───────────────────────►│
+       │                         │                              │              
          │
+       │                         │                              │ 4. Persists 
to log     │
+       │                         │                              
│◄───────────────────────┤
+       │                         │                              │    Returns 
Iggy ACK    │
+       │                         │                              │              
          │
+       │                         │ 5. SourceBatchResult::Ack    │              
          │
+       │                         │◄─────────────────────────────┤              
          │
+       │                         │                              │              
          │
+       │ 6. MQTT PUBACK          │                              │              
          │
+       │◄────────────────────────┤                              │              
          │
+       │                         │ [Pending message cleared]    │              
          │
+```
+
+### Transaction Steps
+
+1. **Inbound Staging**: When `rumqttc` receives a QoS 1 or QoS 2 `PUBLISH` 
packet, the driver stages the message into an internal pending buffer along 
with its deferred `AckToken`. **No `PUBACK` or `PUBREC` is sent to the broker 
yet**.
+2. **Polling & FFI Handoff**: The runtime calls `Source::poll()`. The plugin 
packages staged messages into a batch, records `AckTokens` inside 
`pending_batch`, serializes candidate `ConnectorState`, and returns 
`ProducedMessages` with `Schema::Raw`.
+3. **Iggy Persistence & State Save**: The runtime sends the batch to 
`iggy-server` over TCP/QUIC and saves candidate state checkpoints to disk or 
HTTP storage.
+4. **Ack Callback**: Upon successful Iggy write, the runtime calls 
`Source::on_batch_result(Ack)`.
+5. **Broker Acknowledgment**: The plugin retrieves in-flight `AckTokens` and 
executes `client.try_ack(&publish)`. `rumqttc` transmits the `PUBACK` (for QoS 
1) or initiates `PUBREC` (for QoS 2) to EMQX. The plugin retries failed 
acknowledgements according to `max_retries`.
+6. **Bounded failure**: If all acknowledgement retries fail, the plugin 
commits the already-persisted candidate state, discards the remaining 
acknowledgement tokens, and continues polling. The broker may redeliver the 
abandoned message, so downstream consumers must tolerate duplicates.
+7. **Nack Rollback**: If Iggy delivery fails or times out, the runtime returns 
`SourceBatchResult::Nack`. The plugin drops the candidate state **without 
acknowledging the MQTT tokens**. The broker's QoS 1/2 retry loop will 
subsequently redeliver the unacknowledged messages.
+
+---
+
+## Failure Modes & Observed Redelivery Behaviors
+
+| Failure Scenario | Current Observation | Explanation & Root Cause |
+| :--- | :--- | :--- |
+| **Broker terminates / restarts** | Message may disappear | EMQX may lose its 
in-flight MQTT session state upon restart. Even with persistent storage 
enabled, a standard Docker data volume does not guarantee the preservation of 
every unacknowledged packet across container recreations. |
+| **Iggy terminates / restarts only** | No immediate redelivery | The MQTT 
source plugin and runtime remain connected to the broker. The plugin holds 
`PUBACK`/`PUBREC` pending while waiting for Iggy to recover. Because the active 
TCP/MQTT session never drops, the broker does not trigger an immediate 
redelivery while the session timer remains active. |
+| **Iggy and Connector restart** | Message redelivered | Restarting the 
connector drops the underlying TCP socket and initiates a new MQTT connection. 
Upon reconnecting (with `clean_start = false`), EMQX detects the unacknowledged 
QoS 1/2 message in the session queue and redelivers it. |
+| **Connector terminates / restarts** | Message redelivered | EMQX remains 
online and retains the unacknowledged packet in its active MQTT session. Once 
the connector restarts and re-establishes its session, EMQX immediately 
redelivers the unacknowledged packet. |
+
+---
+
+## Critical End-User & Operational Notes
+
+1. **At-Least-Once Delivery & Duplicate Tolerance**:
+   * The persisted `ConnectorState` records an `acknowledged_messages: u64` 
count for runtime tracking—**it is not a durable MQTT broker cursor**.

Review Comment:
   warning: The README states that the connector persists `ConnectorState` with 
an `acknowledged_messages` count, but `poll()` returns `state: None` on every 
path and no such counter exists. Replace those claims with the actual recovery 
path, where the broker redelivers unacknowledged QoS 1/2 messages. Also at 
lines 388 and 392.



##########
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>,
+    /// MQTT keep-alive interval.
+    pub keep_alive: Option<String>,
+    /// Maximum time to wait for a message during batch collection.
+    pub poll_timeout: Option<String>,
+    /// Capacity of rumqttc's request channel.
+    pub request_capacity: Option<usize>,
+    /// Maximum number of messages in one Iggy source batch.
+    pub batch_size: Option<usize>,
+    /// Maximum time to wait after the first message before flushing a batch.
+    pub batch_timeout: Option<String>,
+    /// Number of retries for a broker acknowledgement after the initial 
attempt.
+    /// After retries are exhausted, the source abandons the broker 
acknowledgement
+    /// and continues. The broker may redeliver the message later, producing 
duplicates.
+    #[serde(default)]
+    pub max_retries: Option<u32>,
+    /// Enables per-message MQTT logging for troubleshooting.
+    pub verbose_logging: Option<bool>,
+}
+
+#[derive(Debug, Deserialize)]
+pub struct MqttTlsConfig {
+    /// Optional custom CA bundle. System roots are used when this is absent.
+    pub ca_file: Option<String>,
+    /// Client certificate for mutual TLS, paired with `client_key_file`.
+    pub client_cert_file: Option<String>,
+    /// Client private key for mutual TLS, paired with `client_cert_file`.
+    pub client_key_file: Option<String>,
+}
+
+pub struct MqttSource {
+    id: u32,
+    config: MqttSourceConfig,
+    poll_timeout: Duration,
+    batch_size: usize,
+    batch_timeout: Duration,
+    // The driver is moved out briefly during event-loop I/O to avoid holding a
+    // mutex guard across await points.
+    driver: Mutex<Option<MqttDriver>>,
+    // At most one batch is staged because the SDK reports one batch result at 
a
+    // time. It holds the MQTT acknowledgement tokens for that batch.
+    pending_batch: Mutex<Option<PendingBatch>>,
+}
+
+impl fmt::Debug for MqttSource {
+    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+        formatter
+            .debug_struct("MqttSource")
+            .field("id", &self.id)
+            .field("config", &self.config)
+            .field("poll_timeout", &self.poll_timeout)
+            .finish_non_exhaustive()
+    }
+}
+
+#[derive(Debug)]
+struct PendingBatch {
+    replay_messages: Vec<ReplayMessage>,
+    schema: Schema,
+    // Tokens are kept separately because the runtime acknowledges Iggy before
+    // the driver acknowledges MQTT.
+    ack_tokens: Vec<AckToken>,
+    // The SDK permits one batch in flight. A NACK leaves this false so poll()
+    // can replay the batch before reading newer MQTT messages.
+    in_flight: bool,
+}
+
+#[derive(Debug, Clone)]
+struct ReplayMessage {
+    headers: Option<BTreeMap<HeaderKey, HeaderValue>>,
+    payload: Vec<u8>,
+}
+
+impl PendingBatch {
+    fn replay(&self) -> ProducedMessages {
+        ProducedMessages {
+            schema: self.schema,
+            messages: self
+                .replay_messages
+                .iter()
+                .cloned()
+                .map(|message| ProducedMessage {
+                    id: None,
+                    checksum: None,
+                    timestamp: None,
+                    origin_timestamp: None,
+                    headers: message.headers,
+                    payload: message.payload,
+                })
+                .collect(),
+            state: None,
+        }
+    }
+}
+
+impl MqttSource {
+    pub fn new(id: u32, config: MqttSourceConfig, _state: 
Option<ConnectorState>) -> Self {
+        Self {
+            id,
+            config,
+            poll_timeout: Duration::from_secs(1),
+            batch_size: DEFAULT_BATCH_SIZE,
+            batch_timeout: Duration::from_millis(10),
+            driver: Mutex::new(None),
+            pending_batch: Mutex::new(None),
+        }
+    }
+
+    fn validate_config(&self) -> Result<(Qos, Duration, Duration, usize, 
usize, Duration), Error> {
+        // Validate static values before opening the network connection so an
+        // operator sees configuration errors during initialization.
+        if self.config.broker_url.trim().is_empty() {
+            return Err(Error::InvalidConfigValue(
+                "broker_url must not be empty".to_string(),
+            ));
+        }
+        if self.config.include_metadata && self.config.protocol == 
MqttProtocol::Mqtt311 {
+            return Err(Error::InvalidConfigValue(
+                "include_metadata is only supported with protocol = 
\"mqtt5\"".to_string(),
+            ));
+        }
+        if self.config.subscriptions.is_empty()
+            || self
+                .config
+                .subscriptions
+                .iter()
+                .any(|topic| topic.trim().is_empty())
+        {
+            return Err(Error::InvalidConfigValue(
+                "subscriptions must contain at least one non-empty 
topic".to_string(),
+            ));
+        }
+        let mut unique_subscriptions = 
HashSet::with_capacity(self.config.subscriptions.len());
+        for subscription in &self.config.subscriptions {
+            if !unique_subscriptions.insert(subscription) {
+                return Err(Error::InvalidConfigValue(format!(
+                    "subscriptions contains duplicate topic filter: 
{subscription}"
+                )));
+            }
+        }
+        // An override is keyed by the exact configured filter. An unknown key
+        // would otherwise be accepted but never used.
+        for subscription in self.config.subscription_qos.keys() {
+            if !self
+                .config
+                .subscriptions
+                .iter()
+                .any(|topic| topic == subscription)
+            {
+                return Err(Error::InvalidConfigValue(format!(
+                    "subscription_qos contains unknown topic filter: 
{subscription}"
+                )));
+            }
+        }
+        if self.config.username.is_some() != self.config.password.is_some() {
+            return Err(Error::InvalidConfigValue(
+                "username and password must be configured 
together".to_string(),
+            ));
+        }
+        // Certificate and hostname checks happen before rumqttc is 
constructed,
+        // so malformed TLS configuration cannot become a reconnect loop.
+        validate_tls_config(&self.config)?;
+
+        let qos = Qos::try_from(self.config.qos)?;
+        for subscription in &self.config.subscriptions {
+            qos_for_subscription(&self.config, subscription, qos)?;
+        }
+        let keep_alive = parse_duration(
+            self.config.keep_alive.as_deref(),
+            DEFAULT_KEEP_ALIVE,
+            "keep_alive",
+        )?;
+        let poll_timeout = parse_duration(
+            self.config.poll_timeout.as_deref(),
+            DEFAULT_POLL_TIMEOUT,
+            "poll_timeout",
+        )?;
+        if poll_timeout.is_zero() {
+            return Err(Error::InvalidConfigValue(
+                "poll_timeout must be greater than zero".to_string(),
+            ));
+        }
+        let request_capacity = self
+            .config
+            .request_capacity
+            .unwrap_or(DEFAULT_REQUEST_CAPACITY);
+        if request_capacity == 0 {
+            return Err(Error::InvalidConfigValue(
+                "request_capacity must be greater than zero".to_string(),
+            ));
+        }
+        if request_capacity < self.config.subscriptions.len() {
+            return Err(Error::InvalidConfigValue(format!(
+                "request_capacity must be at least the number of subscriptions 
({})",
+                self.config.subscriptions.len()
+            )));
+        }
+
+        // These bounds protect both the plugin-owned batch and the MQTT 
request
+        // channel from configurations that would otherwise never make 
progress.
+        let batch_size = self.config.batch_size.unwrap_or(DEFAULT_BATCH_SIZE);
+        if batch_size == 0 {
+            return Err(Error::InvalidConfigValue(
+                "batch_size must be greater than zero".to_string(),
+            ));
+        }
+        let batch_timeout = parse_duration(
+            self.config.batch_timeout.as_deref(),
+            DEFAULT_BATCH_TIMEOUT,
+            "batch_timeout",
+        )?;
+
+        let minimum_keep_alive = match self.config.protocol {
+            MqttProtocol::Mqtt311 => Duration::from_secs(1),
+            MqttProtocol::Mqtt5 => Duration::from_secs(5),
+        };
+        if !keep_alive.is_zero() && keep_alive < minimum_keep_alive {
+            return Err(Error::InvalidConfigValue(format!(
+                "keep_alive must be at least {:?} for {:?}",
+                minimum_keep_alive, self.config.protocol
+            )));
+        }
+
+        Ok((
+            qos,
+            keep_alive,
+            poll_timeout,
+            request_capacity,
+            batch_size,
+            batch_timeout,
+        ))
+    }
+
+    fn max_retries(&self) -> u32 {
+        self.config.max_retries.unwrap_or(DEFAULT_MAX_RETRIES)
+    }
+
+    async fn collect_batch(
+        &self,
+        driver: &mut MqttDriver,
+    ) -> Result<Vec<driver::ReceivedMessage>, Error> {
+        // Start the batch timeout only after the first message arrives. A 
quiet
+        // connector therefore remains cheap while partial batches still flush.
+        let Some(first) = driver.next_message(self.poll_timeout).await? else {
+            return Ok(Vec::new());
+        };
+
+        let mut messages = Vec::with_capacity(self.batch_size);

Review Comment:
   nit: `collect_batch` reserves `batch_size` slots on every batch, although a 
partial batch usually holds one message. Start from `Vec::new()`, as 
`http_source` does for the same buffer.



##########
core/connectors/sources/mqtt_source/src/driver.rs:
##########
@@ -0,0 +1,1226 @@
+// 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 base64::Engine;
+use iggy_common::{HeaderKey, HeaderValue};
+use rumqttc::tokio_rustls::rustls::{
+    self, ClientConfig, RootCertStore,
+    pki_types::{CertificateDer, PrivateKeyDer},
+};
+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 serde::Serialize;
+use std::{
+    collections::{BTreeMap, VecDeque},
+    io::{BufReader, Cursor},
+    sync::{Arc, LazyLock},
+    time::Duration,
+};
+use tokio::time::timeout;
+use tracing::warn;
+use url::Url;
+
+const ACK_RETRY_DELAY: Duration = Duration::from_millis(10);
+const MAX_HEADER_VALUE_LENGTH: usize = 255;
+const MQTT_PROTOCOL_HEADER: &str = "mqtt.protocol";
+const MQTT_TOPIC_HEADER: &str = "mqtt.topic";
+
+static MQTT_PROTOCOL_HEADER_KEY: LazyLock<HeaderKey> = LazyLock::new(|| {
+    HeaderKey::try_from(MQTT_PROTOCOL_HEADER).expect("MQTT protocol header key 
is valid")
+});
+static MQTT_TOPIC_HEADER_KEY: LazyLock<HeaderKey> = LazyLock::new(|| {
+    HeaderKey::try_from(MQTT_TOPIC_HEADER).expect("MQTT topic header key is 
valid")
+});
+static MQTT_QOS_HEADER_KEY: LazyLock<HeaderKey> =
+    LazyLock::new(|| HeaderKey::try_from("mqtt.qos").expect("MQTT QoS header 
key is valid"));
+static MQTT_DUP_HEADER_KEY: LazyLock<HeaderKey> =
+    LazyLock::new(|| HeaderKey::try_from("mqtt.dup").expect("MQTT duplicate 
header key is valid"));
+static MQTT_RETAIN_HEADER_KEY: LazyLock<HeaderKey> =
+    LazyLock::new(|| HeaderKey::try_from("mqtt.retain").expect("MQTT retain 
header key is valid"));
+#[cfg(test)]
+static MQTT_PACKET_ID_HEADER_KEY: LazyLock<HeaderKey> = LazyLock::new(|| {
+    HeaderKey::try_from("mqtt.packet_id").expect("MQTT packet ID header key is 
valid")
+});
+#[cfg(test)]
+static MQTT_RESPONSE_TOPIC_HEADER_KEY: LazyLock<HeaderKey> = LazyLock::new(|| {
+    HeaderKey::try_from("mqtt.response_topic").expect("MQTT response topic 
header key is valid")
+});
+
+// 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,
+    pub(crate) mqtt5_properties: Option<Mqtt5EnvelopeProperties>,
+}
+
+#[derive(Debug, Clone, Serialize)]
+pub(crate) struct Mqtt5EnvelopeProperties {
+    pub(crate) payload_format_indicator: Option<u8>,
+    pub(crate) message_expiry_interval: Option<u32>,
+    pub(crate) topic_alias: Option<u16>,
+    pub(crate) response_topic: Option<String>,
+    pub(crate) correlation_data_base64: Option<String>,
+    pub(crate) user_properties: Vec<Mqtt5UserProperty>,
+    pub(crate) subscription_identifiers: Vec<usize>,
+    pub(crate) content_type: Option<String>,
+}
+
+#[derive(Debug, Clone, Serialize)]
+pub(crate) struct Mqtt5UserProperty {
+    pub(crate) key: String,
+    pub(crate) value: String,
+}
+
+/// 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())
+                        })?;
+                }
+                let buffered_messages =
+                    poll_mqtt311(&mut event_loop, config.subscriptions.len(), 
poll_timeout).await?;
+
+                Ok(Self {
+                    connection: MqttConnection::Mqtt311 {
+                        client: Box::new(client),
+                        event_loop: Box::new(event_loop),
+                    },
+                    buffered_messages,
+                })
+            }
+            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())
+                        })?;
+                }
+                let buffered_messages =
+                    poll_mqtt5(&mut event_loop, config.subscriptions.len(), 
poll_timeout).await?;
+
+                Ok(Self {
+                    connection: MqttConnection::Mqtt5 {
+                        client: Box::new(client),
+                        event_loop: Box::new(event_loop),
+                    },
+                    buffered_messages,
+                })
+            }
+        }
+    }
+
+    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,
+        max_retries: u32,
+    ) -> 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..=max_retries {
+                match self.try_acknowledge(&ack_tokens[acknowledged]) {
+                    Ok(()) => {
+                        acknowledged_token = true;
+                        break;
+                    }
+                    Err(error) => {
+                        last_error = Some(error);
+                        if attempt == max_retries {
+                            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) {

Review Comment:
   simplification: `retain_unacknowledged_tokens` drains a prefix that the 
caller discards when `on_batch_result` drops the batch on error, so the 
retained suffix never retries. Delete the helper and its unit test, or retry 
the suffix from `on_batch_result`.



##########
core/connectors/sources/mqtt_source/README.md:
##########
@@ -0,0 +1,425 @@
+# Apache Iggy MQTT Source Connector (`iggy_connector_mqtt_source`)
+
+The **MQTT Source Connector** is a dynamically loaded shared library plugin 
(`.so` / `.dylib` / `.dll`) for Apache Iggy. Built using the 
`iggy_connector_sdk::source_connector!` C-FFI ABI and the 
[`rumqttc`](https://crates.io/crates/rumqttc) asynchronous client, it ingests 
telemetry and event streams from external MQTT brokers (such as **EMQX** or 
Mosquitto) directly into persistent **Apache Iggy** streams, topics, and 
partitions.
+
+---
+
+## Architecture Overview
+
+```text
+  [ IoT Edge Devices ]
+           │
+           │ MQTT 3.1.1 / MQTT 5 (Sensor Telemetry & Events)
+           ▼
+┌──────────────────────────────────────────────────────────────────────────┐
+│                           EMQX BROKER                                    │
+│  • Client Authentication, Authorizations, & Session Queues              │
+│  • Topic Subscriptions, Wildcards (+, #), & QoS Management              │
+└────────────────────────────────────┬─────────────────────────────────────┘
+                                     │
+                                     │ MQTT / TLS Connection
+                                     ▼
+┌──────────────────────────────────────────────────────────────────────────┐
+│                   IGGY CONNECTORS RUNTIME PROCESS                        │
+│                                                                          │
+│  ┌────────────────────────────────────────────────────────────────────┐  │
+│  │  MQTT Source Plugin (libiggy_connector_mqtt_source.so)             │  │
+│  │  • Embedded `rumqttc` Client (`AsyncClient` + `EventLoop`)         │  │
+│  │  • Normalizes MQTT Packets ──► Iggy Payload & Metadata Headers     │  │
+│  │  • Holds QoS 1/2 PUBACK until Iggy Confirms Persistence            │  │
+│  └─────────────────────────────────┬──────────────────────────────────┘  │
+│                                    │                                     │
+│                                    │ C-ABI FFI Boundary                  │
+│                                    │ (iggy_source_open, poll, ack, etc)  │
+│                                    ▼                                     │
+│  ┌────────────────────────────────────────────────────────────────────┐  │
+│  │  Iggy Connector Runtime Core                                       │  │
+│  │  • Schema Decoders / Encoders (`Schema::Raw`)                      │  │
+│  │  • Optional Field Transformations (Add/Delete/Rename)              │  │
+│  │  • State Checkpoint Storage (File / HTTP endpoint)                 │  │
+│  └─────────────────────────────────┬──────────────────────────────────┘  │
+└────────────────────────────────────┼─────────────────────────────────────┘
+                                     │
+                                     │ Native Binary Client Protocol
+                                     ▼ (TCP / QUIC)
+┌──────────────────────────────────────────────────────────────────────────┐
+│                            IGGY SERVER                                   │
+│  • Thread-Per-Core Shared-Nothing Storage Engine (`io_uring`)            │
+│  • Appends Messages to Stream ──► Topic ──► Partition Disk Logs          │
+└──────────────────────────────────────────────────────────────────────────┘
+```
+
+---
+
+## Features
+
+* **Multi-Protocol Support**: Supports **MQTT 3.1.1** (`mqtt311`) and **MQTT 
5** (`mqtt5`) protocols.
+* **Quality of Service (QoS)**: Ingests QoS 0, QoS 1, and QoS 2 messages.
+* **Deferred Acknowledgment (Manual ACKs)**: Configures `rumqttc` with manual 
acknowledgments enabled. For QoS 1/2 messages, broker acknowledgments (`PUBACK` 
/ `PUBREC`) are held pending and only sent after `iggy-server` confirms 
persistence.
+* **Per-Subscription QoS Overrides**: Configure a global fallback QoS 
alongside per-subscription overrides for specific topic filters.
+* **TLS & Mutual TLS (mTLS)**: Supports secure broker URLs (`mqtts://` / 
`ssl://`) with system trust roots or custom CA bundles (`ca_file`), client 
certificates (`client_cert_file`), and client keys (`client_key_file`) via 
Rustls.
+* **Source-Side Micro-Batching**: Micro-batches incoming messages up to 
`batch_size` or until `batch_timeout` fires, minimizing FFI serialization and 
network overhead.
+* **Stable MQTT Headers**: Preserves protocol, topic, QoS, retain, and dup 
metadata as Iggy message headers. MQTT 5 extended properties are ignored by 
default and are available through the optional metadata envelope.
+* **MQTT 5 Metadata Envelope**: Optionally stores MQTT 5 properties and the 
exact application payload in a versioned JSON envelope.
+* **Route Isolation**: Run multiple connector instances in parallel to route 
distinct MQTT topic filters into separate Iggy streams and topics.
+* **Process Isolation**: Operates inside the `iggy-connectors` process memory 
space via C-ABI FFI, keeping the core `iggy-server` decoupled and untouched.
+
+---
+
+## Configuration Architecture
+
+Configuration operates in **two distinct layers**:
+
+1. **Runtime Layer**: Configured in the connector TOML file (e.g., 
`mqtt_source.toml`). Defines the connector identity, shared library path, and 
destination Iggy stream/topic mapping.
+2. **Plugin Layer (`[plugin_config]`)**: Defines MQTT-specific connection 
parameters, credentials, subscriptions, QoS levels, and TLS configurations. The 
runtime passes this table to the plugin serialized as JSON across the FFI 
boundary.
+
+---
+
+## Configuration Options (`[plugin_config]`)
+
+| Field | Type | Required | Default | Description |
+| :--- | :--- | :---: | :--- | :--- |
+| `broker_url` | String | **Yes** | — | Broker URL scheme (`mqtt://`, 
`mqtts://`, or `ssl://`) and host/port. |
+| `subscriptions` | Array[String] | **Yes** | — | Non-empty list of unique 
MQTT topic filters (wildcards `+` and `#` supported). |
+| `protocol` | String | No | `"mqtt5"` | MQTT protocol version: `"mqtt5"` or 
`"mqtt311"`. |
+| `include_metadata` | Boolean | No | `false` | Enables the MQTT 5 metadata 
envelope. Requires `protocol = "mqtt5"` and a stream schema of `"json"`. |
+| `qos` | Integer | No | `1` | Default global subscription QoS (`0`, `1`, or 
`2`). |
+| `subscription_qos` | Table | No | `{}` | Map of topic-filter strings to 
explicit QoS overrides (`0`, `1`, or `2`). |
+| `client_id` | String | No | `iggy-mqtt-source-{id}` | MQTT client 
identifier. |
+| `username` | String | No | None | Broker authentication username (must be 
supplied together with `password`). |
+| `password` | String | No | None | Broker authentication password (wrapped as 
a secret, never logged). |
+| `clean_start` | Boolean | No | `false` | MQTT clean start flag (or clean 
session for MQTT 3.1.1). |
+| `session_expiry_interval` | Integer | No | None | MQTT 5 session expiry 
interval in seconds. |
+| `keep_alive` | Duration | No | `"30s"` | Ping interval string (minimum `1s` 
for MQTT 3.1.1, minimum `5s` for MQTT 5). |
+| `poll_timeout` | Duration | No | `"1s"` | Maximum duration spent waiting on 
network events per driver poll tick. |
+| `request_capacity` | Integer | No | `32` | Bounded capacity for `rumqttc` 
internal request channel (must be `> 0`). |
+| `batch_size` | Integer | No | `100` | Maximum messages accumulated into a 
single source batch. |
+| `batch_timeout` | Duration | No | `"10ms"` | Maximum wait duration after the 
first message arrives before flushing a batch. |
+| `max_retries` | Integer | No | `5` | Number of retries after the initial 
MQTT acknowledgement attempt. After retries are exhausted, the source abandons 
the broker acknowledgement and continues; the broker may redeliver the message 
later, producing duplicates. |
+| `verbose_logging` | Boolean | No | `false` | Enables additional debug 
logging inside the plugin driver. |
+| `tls.ca_file` | String | No | None | File path to custom CA root 
certificates in PEM format. |
+| `tls.client_cert_file` | String | No | None | File path to mTLS client 
certificate in PEM format (must pair with `client_key_file`). |
+| `tls.client_key_file` | String | No | None | File path to mTLS client 
private key in PEM format (must pair with `client_cert_file`). |
+
+---
+
+## Configuration Examples
+
+### 1. Basic Production Telemetry (QoS 1 with Credentials)
+
+```toml
+type = "source"
+key = "mqtt_telemetry"
+enabled = true
+version = 1
+name = "MQTT Telemetry Source"
+path = "target/release/libiggy_connector_mqtt_source"
+plugin_config_format = "json"
+
+[[streams]]
+stream = "iot"
+topic = "telemetry"
+schema = "raw"
+batch_length = 100
+linger_time = "5ms"
+
+[plugin_config]
+broker_url = "mqtt://127.0.0.1:1883"
+subscriptions = ["devices/+/telemetry"]
+protocol = "mqtt5"
+qos = 1
+client_id = "iggy-telemetry-source"
+username = "iggy_app"
+password = "local-secret-password"
+clean_start = false
+session_expiry_interval = 3600
+keep_alive = "30s"
+batch_size = 100
+batch_timeout = "10ms"
+max_retries = 5
+```
+
+### 2. MQTT 5 Metadata Envelope
+
+Use envelope mode when MQTT 5 publish properties must be preserved. This mode
+has three required rules:
+
+1. `protocol` must be `"mqtt5"`.
+2. `include_metadata` must be `true`.
+3. Every destination stream used by this connector must use `schema = "json"`.
+
+MQTT 3.1.1 cannot use envelope mode because it has no MQTT 5 extended
+PUBLISH properties. The connector rejects `include_metadata = true` with
+`protocol = "mqtt311"` during initialization.
+
+```toml
+type = "source"
+key = "mqtt5_metadata"
+enabled = true
+version = 1
+name = "MQTT 5 source with metadata"
+path = "target/release/libiggy_connector_mqtt_source"
+plugin_config_format = "json"
+
+[[streams]]
+stream = "iot"
+topic = "telemetry"
+schema = "json"
+batch_length = 100
+linger_time = "5ms"
+
+[plugin_config]
+broker_url = "mqtt://127.0.0.1:1883"
+subscriptions = ["devices/+/telemetry"]
+protocol = "mqtt5"
+include_metadata = true
+qos = 1
+client_id = "iggy-mqtt-source-metadata"
+clean_start = false
+session_expiry_interval = 3600
+keep_alive = "30s"
+batch_size = 100
+batch_timeout = "10ms"
+```
+
+The message payload becomes a JSON envelope containing `version`, `topic`,
+`properties`, and `payload_base64`. The original MQTT payload is recovered by
+Base64-decoding `payload_base64`. The stable routing headers remain present:
+`mqtt.protocol`, `mqtt.topic`, `mqtt.qos`, `mqtt.retain`, and `mqtt.dup`.
+
+### 3. Per-Subscription QoS Overrides
+
+```toml
+[plugin_config]
+broker_url = "mqtt://127.0.0.1:1883"
+subscriptions = [
+  "devices/+/telemetry",
+  "devices/+/alerts",
+  "devices/+/diagnostics"
+]
+qos = 1 # Fallback for devices/+/alerts
+
+[plugin_config.subscription_qos]
+"devices/+/telemetry" = 2   # High-importance telemetry via QoS 2
+"devices/+/diagnostics" = 0 # Disposable diagnostic metrics via QoS 0
+```
+
+### 4. Secure TLS & Mutual TLS (mTLS) Setup
+
+```toml
+[plugin_config]
+broker_url = "mqtts://emqx.example.com:8883"
+subscriptions = ["factory/+/metrics"]
+protocol = "mqtt5"
+qos = 1
+
+[plugin_config.tls]
+ca_file = "/etc/iggy/certs/ca.pem"
+client_cert_file = "/etc/iggy/certs/client-cert.pem"
+client_key_file = "/etc/iggy/certs/client-key.pem"
+```
+
+### 5. Multi-Instance Route Isolation
+
+To route distinct MQTT topic filters to different Iggy streams or topics, run
+one source connector instance per route. The runtime loads every connector TOML
+from its configured `config_dir`. For the standard example layout, create these
+files locally:
+
+`core/connectors/runtime/example_config/connectors/mqtt_site_a.toml`:
+
+```toml
+type = "source"
+key = "mqtt_site_a"
+enabled = true
+version = 1
+name = "MQTT source - site A"
+path = "target/release/libiggy_connector_mqtt_source"
+plugin_config_format = "json"
+
+[[streams]]
+stream = "site_a"
+topic = "telemetry"
+schema = "raw"
+batch_length = 100
+linger_time = "5ms"
+
+[plugin_config]
+broker_url = "mqtt://127.0.0.1:1883"
+subscriptions = ["devices/site-a/#"]
+protocol = "mqtt5"
+qos = 1
+client_id = "iggy-mqtt-source-site-a"
+clean_start = false
+session_expiry_interval = 3600
+keep_alive = "30s"
+poll_timeout = "1s"
+request_capacity = 32
+batch_size = 100
+batch_timeout = "10ms"
+verbose_logging = false
+```
+
+`core/connectors/runtime/example_config/connectors/mqtt_site_b.toml`:
+
+```toml
+type = "source"
+key = "mqtt_site_b"
+enabled = true
+version = 1
+name = "MQTT source - site B"
+path = "target/release/libiggy_connector_mqtt_source"
+plugin_config_format = "json"
+
+[[streams]]
+stream = "site_b"
+topic = "telemetry"
+schema = "raw"
+batch_length = 100
+linger_time = "5ms"
+
+[plugin_config]
+broker_url = "mqtt://127.0.0.1:1883"
+subscriptions = ["devices/site-b/#"]
+protocol = "mqtt5"
+qos = 1
+client_id = "iggy-mqtt-source-site-b"
+clean_start = false
+session_expiry_interval = 3600
+keep_alive = "30s"
+poll_timeout = "1s"
+request_capacity = 32
+batch_size = 100
+batch_timeout = "10ms"
+verbose_logging = false
+```
+
+Use unique connector keys and MQTT client IDs for every route. Keep topic
+filters non-overlapping unless duplicate delivery to multiple Iggy destinations
+is intentional. Build the plugin first, copy these files into the active
+connector `config_dir`, and restart `iggy-connectors`.
+
+---
+
+## Environment Variable Overrides
+
+Any property inside `[plugin_config]` can be overridden at runtime using 
environment variables without modifying configuration files:

Review Comment:
   nit: The README promises an environment override for any `[plugin_config]​` 
property, but the runtime builds one flat key from the suffix. 
`[plugin_config.tls]​` and `subscription_qos` overrides are silently ignored, 
so document the flat-only limit as the `http_source` README does.



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