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]
