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)допустимо только для диагностики и не защищает от подмены сервера.