Перейти к содержанию
Материалы раздела

uMQTT

Простой MQTT клиент на базе протокола MQTT 3.1.1.

Особенности

  • Минималистичная реализация MQTT 3.1.1
  • Поддержка QoS 0 и QoS 1
  • SSL/TLS для защищённых соединений
  • Last Will and Testament (завещание)
  • Аутентификация по логину/паролю
  • Блокирующий и неблокирующий режимы приёма сообщений
  • Низкое потребление памяти

Быстрый старт

-- Создание клиента
local client = umqtt.new("esp32_device", "broker.hivemq.com", 1883)

-- Установка callback для входящих сообщений
client:set_callback(function(topic, msg)
    print("Получено:", topic, msg)
end)

-- Подключение
client:connect(true, 5000)

-- Подписка
client:subscribe("test/topic", 0)

-- Публикация
client:publish("test/status", "online", false, 0)

-- Основной цикл
while client:connected() do
    client:check_msg()  -- неблокирующая проверка
    tmr.delayms(100)
end

-- Отключение
client:disconnect()

client = umqtt.new(client_id, server, port, [user], [password], [ssl], [keepalive], [ca_file_or_pem])

Создаёт новый экземпляр MQTT клиента.

Аргументы:

  • client_id (строка): уникальный идентификатор клиента.
  • server (строка): адрес MQTT брокера (IP или доменное имя).
  • port (целое число, опционально): порт брокера. По умолчанию 1883 (8883 для SSL).
  • user (строка, опционально): имя пользователя для аутентификации.
  • password (строка, опционально): пароль для аутентификации.
  • ssl (булево, опционально): использовать SSL/TLS соединение. По умолчанию false.
  • keepalive (целое число, опционально): интервал keepalive в секундах. По умолчанию 0 (отключён).
  • ca_file_or_pem (строка, опционально): путь к PEM-файлу доверенного CA или сам PEM-сертификат, не более 32 КиБ. При включённой проверке TLS этот аргумент обязателен для SSL-клиента.

Возвращает: экземпляр клиента.

-- Простое подключение
local client = umqtt.new("device_001", "192.168.1.100", 1883)

-- С аутентификацией
local client = umqtt.new("device_001", "broker.example.com", 1883, "user", "password")

-- С SSL/TLS и проверкой сертификата
local client = umqtt.new("device_001", "broker.example.com", 8883,
    "user", "password", true, 60, "/certs/broker-root-ca.pem")

-- SSL без аутентификации
local client = umqtt.new("device_001", "broker.example.com", 8883,
    nil, nil, true, 0, "/certs/broker-root-ca.pem")

-- С keepalive 60 секунд
local client = umqtt.new("device_001", "broker.example.com", 1883, "user", "pass", false, 60)

client:set_callback(callback)

Устанавливает функцию обратного вызова для обработки входящих сообщений.

Важно: callback должен быть установлен до вызова subscribe().

Аргументы:

  • callback (функция или nil): функция с двумя аргументами либо nil для удаления ранее установленного callback:
    • topic (строка): название топика
    • message (строка): содержимое сообщения

Возвращает: ничего.

client:set_callback(function(topic, msg)
    print("Topic: " .. topic)
    print("Message: " .. msg)

    -- Обработка JSON
    local ok, data = pcall(cjson.decode, msg)
    if ok then
        print("Temperature: " .. (data.temp or "N/A"))
    end
end)

-- Удалить callback
client:set_callback(nil)

Callback может замыкать сам client: такая ссылка участвует в обычной сборке циклов Lua и не удерживает клиент навсегда. Если callback завершился Lua-ошибкой, ошибка передаётся из check_msg(), wait_msg() либо операции, ожидающей ACK. Для входящего QoS 1 сообщения PUBACK в этом случае не отправляется, поэтому брокер может повторить доставку. Если ошибка возникла во время другой операции, ожидающей ACK, клиент закрывает соединение, чтобы запоздавший ACK не был принят следующей операцией.

client:set_last_will(topic, msg, [retain], [qos])

Устанавливает Last Will and Testament (LWT) — сообщение, которое брокер опубликует при неожиданном отключении клиента.

Важно: должен быть вызван до connect().

Аргументы:

  • topic (строка): топик для LWT сообщения.
  • msg (строка): содержимое LWT сообщения.
  • retain (булево, опционально): флаг сохранения. По умолчанию false.
  • qos (целое число, опционально): качество обслуживания (0 или 1). По умолчанию 0.

Возвращает: ничего.

-- Установка LWT перед подключением
client:set_last_will("devices/esp32/status", "offline", true, 0)
client:connect()

-- Теперь при неожиданном отключении брокер опубликует "offline"

client:connect([clean_session], [timeout_ms])

Подключается к MQTT брокеру.

Аргументы:

  • clean_session (булево, опционально): очистить сессию при подключении. По умолчанию true.
  • timeout_ms (целое число, опционально): таймаут подключения в миллисекундах. DNS и неблокирующий TCP connect ограничены общим монотонным сроком; оставшееся время передаётся TLS и MQTT I/O. По умолчанию 10000.

Возвращает session_present (булево) — флаг наличия сохранённой сессии на брокере. Сетевые/TLS ошибки вызывают Lua-исключение; отказ CONNACK возвращает nil, code.

-- Простое подключение
local session = client:connect()

-- С сохранением сессии
local session = client:connect(false)
if session then
    print("Восстановлена существующая сессия")
end

-- С увеличенным таймаутом
local session = client:connect(true, 30000)

client:reconnect()

Переподключается к брокеру с теми же параметрами, что использовались при последнем успешном connect().

Полезно для восстановления соединения после разрыва.

Аргументы: нет.

Возвращает:

  • true — при успешном переподключении
  • false, error_code — при ошибке (отрицательный код = ошибка сети, положительный = код CONNACK)

Коды ошибок сети:

  • -1: не удалось разрешить имя хоста
  • -2: не удалось создать сокет
  • -3: не удалось подключиться
  • -4: неверный CONNACK
  • -5: ошибка отправки
  • -6: ошибка создания SSL контекста
  • -7: ошибка создания SSL сессии
  • -8: ошибка SSL рукопожатия
  • -9: CA-сертификат отсутствует или не прошёл разбор
  • -10: истёк таймаут подключения

Коды CONNACK (положительные):

  • 1: неприемлемая версия протокола
  • 2: идентификатор клиента отклонён
  • 3: сервер недоступен
  • 4: неверный логин/пароль
  • 5: не авторизован
-- Простой reconnect
local ok, err = client:reconnect()
if ok then
    print("Переподключен!")
    -- Нужно заново подписаться на топики
    client:subscribe("my/topic", 0)
else
    print("Ошибка переподключения:", err)
end

-- Reconnect с повторными попытками
local function reconnect_loop(client, max_attempts, delay_ms)
    for i = 1, max_attempts do
        local ok, err = client:reconnect()
        if ok then
            return true
        end
        print("Попытка", i, "неудачна:", err)
        tmr.delayms(delay_ms)
    end
    return false
end

client:disconnect()

Корректно отключается от брокера, отправляя пакет DISCONNECT.

Аргументы: нет.

Возвращает: ничего.

client:disconnect()

client:connected()

Проверяет состояние подключения.

Аргументы: нет.

Возвращает: булево значение.

if client:connected() then
    print("Подключен к брокеру")
else
    print("Отключен")
end

client:publish(topic, msg, [retain], [qos])

Публикует сообщение в указанный топик.

Аргументы:

  • topic (строка): название топика.
  • msg (строка): содержимое сообщения.
  • retain (булево, опционально): флаг сохранения сообщения на брокере. По умолчанию false.
  • qos (целое число, опционально): качество обслуживания (0 или 1). По умолчанию 0.

Возвращает:

  • Для QoS 0: true при успешной отправке
  • Для QoS 1: message_id (целое число) — идентификатор сообщения, подтверждённый брокером (PUBACK получен)
  • При ошибке: nil, error_message
-- Простая публикация (QoS 0)
local ok = client:publish("sensors/temp", "25.5")
if ok then print("Отправлено") end

-- С retain флагом (последнее сообщение сохраняется на брокере)
client:publish("devices/status", "online", true)

-- С QoS 1 (гарантированная доставка с подтверждением)
local msg_id, err = client:publish("commands/important", "reboot", false, 1)
if msg_id then
    print("Сообщение доставлено, ID:", msg_id)
else
    print("Ошибка:", err)
end

-- Проверка доставки критичных сообщений
local function reliable_publish(client, topic, msg)
    local result, err = client:publish(topic, msg, false, 1)
    if not result then
        -- Попытка переподключения и повторной отправки
        if client:reconnect() then
            result, err = client:publish(topic, msg, false, 1)
        end
    end
    return result, err
end

client:subscribe(topic, [qos])

Подписывается на топик.

Важно: перед подпиской должен быть установлен callback через set_callback().

Аргументы:

  • topic (строка): название топика. Поддерживаются wildcard символы:
    • + — один уровень (например, sensors/+/temp)
    • # — все подуровни (например, sensors/#)
  • qos (целое число, опционально): качество обслуживания (0 или 1). По умолчанию 0.

Возвращает: ничего или исключение.

-- Подписка на конкретный топик
client:subscribe("devices/esp32/commands", 0)

-- Подписка с wildcard
client:subscribe("sensors/+/temperature", 0)
client:subscribe("home/#", 1)

client:unsubscribe(topic)

Отписывается от топика.

Аргументы:

  • topic (строка): название топика.

Возвращает: ничего или исключение.

client:unsubscribe("sensors/temp")

client:check_msg()

Проверяет наличие входящих сообщений в неблокирующем режиме.

Если сообщение доступно, вызывается callback. Если сообщений нет — функция немедленно возвращает управление.

Аргументы: нет.

Возвращает: тип пакета (целое число) или nil если сообщений нет.

Если callback вызвал ошибку, check_msg() повторно вызывает это Lua-исключение.

-- Типичное использование в основном цикле
while true do
    if client:connected() then
        client:check_msg()
        -- Другая логика
    end
    tmr.delayms(100)
end

client:wait_msg()

Ожидает входящее сообщение в блокирующем режиме.

Функция блокирует выполнение до получения сообщения или ошибки.

Аргументы: нет.

Возвращает: тип пакета (целое число) или nil и сообщение об ошибке.

Если callback вызвал ошибку, wait_msg() повторно вызывает это Lua-исключение.

-- Блокирующее ожидание
while true do
    local result, err = client:wait_msg()
    if not result then
        print("Ошибка: " .. (err or "unknown"))
        break
    end
end

client:ping()

Отправляет PING запрос брокеру для поддержания соединения.

Используется при ручном управлении keepalive.

Аргументы: нет.

Возвращает: ничего или исключение.

-- Ручной keepalive каждые 30 секунд
local last_ping = os.time()

while client:connected() do
    client:check_msg()

    if os.time() - last_ping > 30 then
        client:ping()
        last_ping = os.time()
    end

    tmr.delayms(100)
end

Пример: датчик температуры

-- Конфигурация
local BROKER = "192.168.1.100"
local PORT = 1883
local CLIENT_ID = "temp_sensor_" .. os.flashEUI()
local TOPIC_STATUS = "devices/" .. CLIENT_ID .. "/status"
local TOPIC_TEMP = "devices/" .. CLIENT_ID .. "/temperature"
local TOPIC_CMD = "devices/" .. CLIENT_ID .. "/command"

-- Создание клиента
local client = umqtt.new(CLIENT_ID, BROKER, PORT)

-- Callback для команд
client:set_callback(function(topic, msg)
    print("Command received: " .. msg)
    if msg == "get_temp" then
        -- Здесь читаем реальную температуру
        local temp = 25.5
        client:publish(TOPIC_TEMP, tostring(temp), false, 0)
    elseif msg == "reboot" then
        client:publish(TOPIC_STATUS, "rebooting", false, 0)
        client:disconnect()
        os.reboot()
    end
end)

-- Last Will - брокер опубликует при отключении
client:set_last_will(TOPIC_STATUS, "offline", true, 0)

-- Подключение
local ok, err = pcall(client.connect, client, true, 5000)
if not ok then
    print("Connection failed: " .. tostring(err))
    return
end

print("Connected to broker")

-- Публикуем статус online
client:publish(TOPIC_STATUS, "online", true, 0)

-- Подписка на команды
client:subscribe(TOPIC_CMD, 0)

-- Основной цикл
local last_temp_time = 0
while client:connected() do
    -- Проверка входящих сообщений
    client:check_msg()

    -- Публикация температуры каждые 60 секунд
    local now = os.time()
    if now - last_temp_time >= 60 then
        local temp = 25.5 + math.random() * 2  -- Симуляция
        client:publish(TOPIC_TEMP, string.format("%.1f", temp), false, 0)
        last_temp_time = now
    end

    tmr.delayms(100)
end

print("Disconnected from broker")

Полный пример: умная розетка

local client = umqtt.new("smart_plug_1", "mqtt.local", 1883, "user", "pass")

local relay_state = false
local TOPIC_STATE = "home/plug1/state"
local TOPIC_SET = "home/plug1/set"

client:set_callback(function(topic, msg)
    if msg == "ON" or msg == "1" then
        relay_state = true
        pio.pin.sethigh(pio.GPIO26)
    elseif msg == "OFF" or msg == "0" then
        relay_state = false
        pio.pin.setlow(pio.GPIO26)
    end

    -- Публикуем актуальное состояние с retain
    client:publish(TOPIC_STATE, relay_state and "ON" or "OFF", true, 0)
end)

client:set_last_will(TOPIC_STATE, "OFFLINE", true, 0)
client:connect(true)

-- Публикуем начальное состояние
client:publish(TOPIC_STATE, relay_state and "ON" or "OFF", true, 0)

client:subscribe(TOPIC_SET, 1)

while client:connected() do
    client:check_msg()
    tmr.delayms(50)
end

Пример: SSL/TLS подключение

-- Подключение к брокеру с SSL/TLS
local CA = "/certs/broker-root-ca.pem"
local client = umqtt.new(
    "secure_device",
    "broker.example.com",
    8883,
    nil,      -- user (опционально)
    nil,      -- password (опционально)
    true,     -- SSL включён
    60,       -- keepalive
    CA        -- доверенный CA
)

-- Callback для сообщений
client:set_callback(function(topic, msg)
    print("Received:", topic, msg)
end)

-- Подключение
local ok, err = pcall(function()
    client:connect(true, 10000)
end)

if not ok then
    print("SSL connection failed:", err)
    return
end

print("Connected via SSL/TLS!")

-- Подписка и публикация работают как обычно
client:subscribe("test/secure", 0)
client:publish("test/secure", "Hello over SSL!", false, 0)

-- Основной цикл
while client:connected() do
    client:check_msg()
    tmr.delayms(100)
end

client:disconnect()

Пример: SSL с аутентификацией

-- Подключение к защищённому брокеру с логином и паролем
local client = umqtt.new(
    "iot_device_001",
    "mqtt.example.com",
    8883,
    "username",
    "password",
    true,         -- SSL
    60,           -- keepalive 60 секунд
    "/certs/broker-root-ca.pem"
)

client:set_last_will("devices/iot_device_001/status", "offline", true, 0)

client:set_callback(function(topic, msg)
    print("Command:", msg)
end)

local ok = pcall(function()
    client:connect(true, 15000)
end)

if ok then
    client:publish("devices/iot_device_001/status", "online", true, 0)
    client:subscribe("devices/iot_device_001/commands", 1)
end

Обработка ошибок

local client = umqtt.new("test", "broker.local", 1883)

-- Подключение с обработкой ошибок
local ok, err = pcall(function()
    client:connect(true, 5000)
end)

if not ok then
    print("Ошибка подключения: " .. tostring(err))
    return
end

-- Публикация с обработкой ошибок
local function safe_publish(topic, msg)
    if not client:connected() then
        return false, "not connected"
    end

    local ok, err = pcall(client.publish, client, topic, msg, false, 0)
    return ok, err
end

local ok, err = safe_publish("test", "hello")
if not ok then
    print("Publish failed: " .. tostring(err))
end

Ограничения

  • Нет QoS 2: поддерживаются только QoS 0 и QoS 1
  • Ручное переподключение: используйте client:reconnect() — автоматического переподключения нет
  • Блокирующие операции: connect(), subscribe(), publish() с QoS 1 блокируют выполнение
  • Подписки не сохраняются: после reconnect() нужно заново подписаться на топики
  • SSL/TLS: проверка сертификата включена по умолчанию; передайте точный доверенный CA через ca_file_or_pem. Глобальное отключение проверки через net.skip_ssl_verify(true) допустимо только для диагностики и не защищает от подмены сервера.