Un cliente MQTT es una aplicación que publica mensajes y se suscribe a temas a través de un broker. Zephyr incluye una librería en C (<zephyr/net/mqtt.h>) en la que configuramos de forma explícita los buffers y el transporte.
Configuración (prj.conf)
Añadimos a lo anterior:
CONFIG_MQTT_LIB=yAnatomía de un cliente MQTT en Zephyr
Necesitamos definir tres cosas:
Buffers: Donde se guardan los datos antes de enviarse (TX) y al recibirse (RX).
Broker: La dirección IP y puerto del servidor.
Cliente: La estructura que guarda el estado de la conexión.
El código
El ejemplo usa la dirección 192.0.2.10, perteneciente al rango reservado para documentación. Sustitúyela por la IP de tu broker dentro de la red de pruebas.
#include <zephyr/kernel.h>
#include <zephyr/net/socket.h>
#include <zephyr/net/mqtt.h>
#include <errno.h>
#include <stdbool.h>
#include <string.h>
/* Buffers para MQTT */
static uint8_t rx_buffer[128];
static uint8_t tx_buffer[128];
static struct mqtt_client client_ctx;
static struct sockaddr_in broker;
static bool connected;
static void mqtt_evt_handler(struct mqtt_client *const client,
const struct mqtt_evt *evt)
{
switch (evt->type) {
case MQTT_EVT_CONNACK:
if (evt->result == 0) {
connected = true;
printk("Sesión MQTT establecida\n");
} else {
printk("El broker rechazó la conexión: %d\n", evt->result);
}
break;
case MQTT_EVT_PUBLISH: {
size_t remaining = evt->param.publish.message.payload.len;
uint8_t payload[64];
printk("Mensaje recibido: ");
while (remaining > 0) {
int len = mqtt_read_publish_payload(client, payload,
MIN(sizeof(payload), remaining));
if (len <= 0) {
break;
}
printk("%.*s", len, payload);
remaining -= len;
}
printk("\n");
break;
}
case MQTT_EVT_DISCONNECT:
connected = false;
printk("Sesión MQTT cerrada: %d\n", evt->result);
break;
default:
break;
}
}
/* Función auxiliar para preparar la estructura del cliente */
static int prepare_mqtt_client(void)
{
mqtt_client_init(&client_ctx);
/* Dirección del broker de laboratorio */
broker.sin_family = AF_INET;
broker.sin_port = htons(1883);
if (zsock_inet_pton(AF_INET, "192.0.2.10", &broker.sin_addr) != 1) {
return -EINVAL;
}
client_ctx.broker = (struct sockaddr *)&broker;
client_ctx.evt_cb = mqtt_evt_handler;
client_ctx.client_id.utf8 = (uint8_t *)"zephyr_device_001";
client_ctx.client_id.size = strlen("zephyr_device_001");
client_ctx.protocol_version = MQTT_VERSION_3_1_1;
client_ctx.transport.type = MQTT_TRANSPORT_NON_SECURE;
/* Buffers manuales (Obligatorio en Zephyr) */
client_ctx.rx_buf = rx_buffer;
client_ctx.rx_buf_size = sizeof(rx_buffer);
client_ctx.tx_buf = tx_buffer;
client_ctx.tx_buf_size = sizeof(tx_buffer);
return 0;
}
static int process_mqtt(int timeout_ms)
{
struct zsock_pollfd fds = {
.fd = client_ctx.transport.tcp.sock,
.events = ZSOCK_POLLIN,
};
int rc = zsock_poll(&fds, 1, timeout_ms);
if (rc < 0) {
return -errno;
}
if (rc > 0 && (fds.revents & (ZSOCK_POLLERR | ZSOCK_POLLHUP))) {
return -ECONNRESET;
}
if (rc > 0 && (fds.revents & ZSOCK_POLLIN)) {
rc = mqtt_input(&client_ctx);
if (rc != 0) {
return rc;
}
}
return mqtt_live(&client_ctx);
}
int main(void)
{
/* ... (Asumimos que el Wi-Fi ya conectó y tenemos IP) ... */
int rc = prepare_mqtt_client();
if (rc != 0) {
return rc;
}
/* 1. Conexión */
rc = mqtt_connect(&client_ctx);
if (rc != 0) {
printk("Error conectando MQTT: %d\n", rc);
return rc;
}
/* mqtt_connect envía CONNECT; procesamos la entrada hasta recibir CONNACK. */
while (!connected) {
rc = process_mqtt(1000);
if (rc != 0) {
return rc;
}
}
/* 2. Suscribirse a comandos */
struct mqtt_topic topics[] = {{
.topic = {
.utf8 = (uint8_t *)"zephyr/curso/comandos",
.size = strlen("zephyr/curso/comandos"),
},
.qos = MQTT_QOS_0_AT_MOST_ONCE,
}};
const struct mqtt_subscription_list subscriptions = {
.list = topics,
.list_count = ARRAY_SIZE(topics),
.message_id = 1,
};
rc = mqtt_subscribe(&client_ctx, &subscriptions);
if (rc != 0) {
return rc;
}
/* 3. Publicar telemetría */
char payload[] = "Hola desde Zephyr OS";
struct mqtt_publish_param param = {
.message = {
.topic = {
.topic = {
.utf8 = (uint8_t *)"zephyr/curso/telemetria",
.size = strlen("zephyr/curso/telemetria"),
},
.qos = MQTT_QOS_0_AT_MOST_ONCE,
},
.payload = {
.data = payload,
.len = strlen(payload),
},
},
.message_id = 2,
};
rc = mqtt_publish(&client_ctx, ¶m);
if (rc == 0) {
printk("Mensaje publicado correctamente.\n");
}
/* 4. Procesar entrada y mantener viva la conexión */
while (connected) {
rc = process_mqtt(1000);
if (rc != 0) {
printk("Error en la sesión MQTT: %d\n", rc);
break;
}
}
mqtt_disconnect(&client_ctx);
return rc;
}El puerto 1883 del ejemplo no cifra ni autentica el servidor. Úsalo solo en una red de laboratorio. Para producción, configura MQTT_TRANSPORT_SECURE, registra las credenciales TLS y valida el nombre del broker.