mirror of
https://github.com/cesanta/mongoose.git
synced 2026-10-01 05:52:32 +07:00
Add MQTT module to STM32 pack
This commit is contained in:
@@ -1,11 +1,11 @@
|
||||
PROG ?= example # Program we are building
|
||||
DELETE = rm -rf # Command to remove files
|
||||
OUT ?= -o $(PROG) # Compiler argument for output file
|
||||
SOURCES = main.c mongoose.c packed_fs.c # Source code files, packed_fs.c contains ca.pem, which contains CA certs for TLS
|
||||
SOURCES = main.c mongoose.c mongoose_mqtt.c packed_fs.c # Source code files, packed_fs.c contains ca.pem, which contains CA certs for TLS
|
||||
CFLAGS = -W -Wall -Wextra -g -I. # Build options
|
||||
|
||||
# Mongoose build options. See https://mongoose.ws/documentation/#build-options
|
||||
CFLAGS_MONGOOSE += -DMG_ENABLE_LINES=1
|
||||
# Mongoose build options. See https://mongoose.ws/docs/getting-started/build-options/
|
||||
CFLAGS_MONGOOSE += -DMG_ENABLE_LINES=1 -DMG_TLS=MG_TLS_BUILTIN
|
||||
|
||||
ifeq ($(OS),Windows_NT) # Windows settings. Assume MinGW compiler. To use VC: make CC=cl CFLAGS=/MD OUT=/Feprog.exe
|
||||
PROG ?= example.exe # Use .exe suffix for the binary
|
||||
@@ -24,9 +24,3 @@ $(PROG): $(SOURCES) # Build program from sources
|
||||
|
||||
clean: # Cleanup. Delete built program and all build artifacts
|
||||
$(DELETE) $(PROG) *.o *.obj *.exe *.dSYM mbedtls
|
||||
|
||||
# see https://mongoose.ws/tutorials/tls/#how-to-build for TLS build options
|
||||
|
||||
mbedtls: # Pull and build mbedTLS library
|
||||
git clone --depth 1 -b v2.28.2 https://github.com/mbed-tls/mbedtls $@
|
||||
$(MAKE) -C mbedtls/library
|
||||
|
||||
@@ -1,122 +1,19 @@
|
||||
// Copyright (c) 2023-2025 Cesanta Software Limited
|
||||
// All rights reserved
|
||||
//
|
||||
// Example MQTT client. It performs the following steps:
|
||||
// 1. Connects to the MQTT server specified by `s_url` variable
|
||||
// 2. When connected, subscribes to the topic `s_sub_topic`
|
||||
// 3. When it receives a message, echoes it back to `s_pub_topic`
|
||||
// 4. Timer-based reconnection logic revives the connection when it is down
|
||||
// 5. Ping server periodically. When disconnected, a last will is published
|
||||
//
|
||||
// To enable SSL/TLS, see https://mongoose.ws/tutorials/tls/#how-to-build
|
||||
|
||||
#include "mongoose.h"
|
||||
|
||||
static const char *s_url = "mqtt://broker.hivemq.com:1883";
|
||||
static const char *s_sub_topic = "mg/123/rx"; // Subscribe topic
|
||||
static const char *s_pub_topic = "mg/123/tx"; // Publish topic
|
||||
static uint8_t s_qos = 1; // MQTT QoS
|
||||
static struct mg_connection *s_mqtt_conn; // Client connection
|
||||
|
||||
static void subscribe(struct mg_connection *c, struct mg_str topic) {
|
||||
struct mg_mqtt_opts opts = {};
|
||||
memset(&opts, 0, sizeof(opts));
|
||||
opts.topic = topic;
|
||||
opts.qos = s_qos;
|
||||
mg_mqtt_sub(c, &opts);
|
||||
MG_INFO(("%lu SUBSCRIBED to %.*s", c->id, topic.len, topic.buf));
|
||||
}
|
||||
|
||||
static void publish(struct mg_connection *c, struct mg_str topic,
|
||||
struct mg_str message) {
|
||||
struct mg_mqtt_opts opts = {};
|
||||
memset(&opts, 0, sizeof(opts));
|
||||
opts.topic = topic;
|
||||
opts.message = message;
|
||||
opts.qos = s_qos;
|
||||
mg_mqtt_pub(c, &opts);
|
||||
MG_INFO(("%lu PUBLISHED %.*s -> %.*s", c->id, topic.len, topic.buf,
|
||||
message.len, message.buf));
|
||||
}
|
||||
|
||||
static void mqtt_ev_handler(struct mg_connection *c, int ev, void *ev_data) {
|
||||
if (ev == MG_EV_OPEN) {
|
||||
MG_INFO(("%lu CREATED", c->id));
|
||||
// c->is_hexdumping = 1;
|
||||
} else if (ev == MG_EV_CONNECT) {
|
||||
if (c->is_tls) {
|
||||
struct mg_tls_opts opts = {.ca = mg_unpacked("/certs/ca.pem"),
|
||||
.name = mg_url_host(s_url)};
|
||||
mg_tls_init(c, &opts);
|
||||
}
|
||||
} else if (ev == MG_EV_ERROR) {
|
||||
// On error, log error message
|
||||
MG_ERROR(("%lu ERROR %s", c->id, (char *) ev_data));
|
||||
} else if (ev == MG_EV_MQTT_OPEN) {
|
||||
// MQTT connect is successful
|
||||
MG_INFO(("%lu CONNECTED to %s", c->id, s_url));
|
||||
subscribe(c, mg_str(s_sub_topic));
|
||||
} else if (ev == MG_EV_MQTT_MSG) {
|
||||
// When we get echo response, print it
|
||||
char response[100];
|
||||
struct mg_mqtt_message *mm = (struct mg_mqtt_message *) ev_data;
|
||||
mg_snprintf(response, sizeof(response), "Got %.*s -> %.*s", mm->topic.len,
|
||||
mm->topic.buf, mm->data.len, mm->data.buf);
|
||||
publish(c, mg_str(s_pub_topic), mg_str(response));
|
||||
} else if (ev == MG_EV_MQTT_CMD) {
|
||||
struct mg_mqtt_message *mm = (struct mg_mqtt_message *) ev_data;
|
||||
if (mm->cmd == MQTT_CMD_PINGREQ) mg_mqtt_pong(c);
|
||||
} else if (ev == MG_EV_CLOSE) {
|
||||
MG_INFO(("%lu CLOSED", c->id));
|
||||
s_mqtt_conn = NULL; // Mark that we're closed
|
||||
}
|
||||
}
|
||||
|
||||
// Timer function - recreate client connection if it is closed
|
||||
static void timer_fn(void *arg) {
|
||||
if (s_mqtt_conn == NULL) {
|
||||
struct mg_mgr *mgr = (struct mg_mgr *) arg;
|
||||
struct mg_mqtt_opts opts = {.clean = true,
|
||||
.qos = s_qos,
|
||||
.topic = mg_str(s_pub_topic),
|
||||
.keepalive = 5,
|
||||
.version = 4,
|
||||
.message = mg_str("bye")};
|
||||
s_mqtt_conn = mg_mqtt_connect(mgr, s_url, &opts, mqtt_ev_handler, NULL);
|
||||
} else {
|
||||
mg_mqtt_ping(s_mqtt_conn);
|
||||
}
|
||||
}
|
||||
|
||||
int main(int argc, char *argv[]) {
|
||||
int main(void) {
|
||||
struct mg_mgr mgr;
|
||||
int i;
|
||||
|
||||
// Parse command-line flags
|
||||
for (i = 1; i < argc; i++) {
|
||||
if (strcmp(argv[i], "-u") == 0 && argv[i + 1] != NULL) {
|
||||
s_url = argv[++i];
|
||||
} else if (strcmp(argv[i], "-p") == 0 && argv[i + 1] != NULL) {
|
||||
s_pub_topic = argv[++i];
|
||||
} else if (strcmp(argv[i], "-s") == 0 && argv[i + 1] != NULL) {
|
||||
s_sub_topic = argv[++i];
|
||||
} else if (strcmp(argv[i], "-v") == 0 && argv[i + 1] != NULL) {
|
||||
mg_log_set(atoi(argv[++i]));
|
||||
} else {
|
||||
MG_ERROR(("Unknown option: %s. Usage:", argv[i]));
|
||||
MG_ERROR(
|
||||
("%s [-u mqtts://SERVER:PORT] [-p PUB_TOPIC] [-s SUB_TOPIC] "
|
||||
"[-v DEBUG_LEVEL]",
|
||||
argv[0], argv[i]));
|
||||
return 1;
|
||||
}
|
||||
}
|
||||
|
||||
mg_mgr_init(&mgr);
|
||||
mg_timer_add(&mgr, 3000, MG_TIMER_REPEAT | MG_TIMER_RUN_NOW, timer_fn, &mgr);
|
||||
mg_mqtt_init(&mgr);
|
||||
|
||||
for (;;) {
|
||||
mg_mgr_poll(&mgr, 1000);
|
||||
mg_mgr_poll(&mgr, 100);
|
||||
mg_mqtt_poll(&mgr);
|
||||
}
|
||||
|
||||
mg_mgr_free(&mgr);
|
||||
|
||||
return 0;
|
||||
|
||||
@@ -0,0 +1,123 @@
|
||||
// Copyright (c) 2026 Cesanta Software Limited
|
||||
// All rights reserved
|
||||
//
|
||||
// Example MQTT client. It performs the following steps:
|
||||
// 1. Connects to the MQTT server specified by MQTT_SERVER_URL
|
||||
// 2. When connected, subscribes to the topic MQTT_SUBSCRIBE_TOPIC
|
||||
// 3. When it receives a message, echoes it back to MQTT_PUBLISH_TOPIC
|
||||
// 4. Timer-based reconnection logic revives the connection when it is down
|
||||
// 5. Ping server periodically. When disconnected, a last will is published
|
||||
// 6. Implements "ota.update" for OTA updates, see https://mongoose.ws/mqtt/
|
||||
|
||||
#include "mongoose.h"
|
||||
|
||||
#define MQTT_SERVER_URL "mqtt://broker.hivemq.com:1883"
|
||||
#define MQTT_PUBLISH_TOPIC "mg/123/tx"
|
||||
#define MQTT_SUBSCRIBE_TOPIC "mg/123/rx"
|
||||
#define MQTT_QOS 1
|
||||
#define RECONNECT_PERIOD_MS 3000
|
||||
|
||||
static struct mg_connection *s_mqtt_conn; // Client connection
|
||||
static struct mg_rpc *s_rpc = NULL; // List of registered RPC methods
|
||||
|
||||
static void subscribe(struct mg_connection *c, struct mg_str topic) {
|
||||
struct mg_mqtt_opts opts = {};
|
||||
memset(&opts, 0, sizeof(opts));
|
||||
opts.topic = topic;
|
||||
opts.qos = MQTT_QOS;
|
||||
mg_mqtt_sub(c, &opts);
|
||||
MG_DEBUG(("%lu SUBSCRIBED to %.*s", c->id, topic.len, topic.buf));
|
||||
}
|
||||
|
||||
static void publish(struct mg_connection *c, struct mg_str topic,
|
||||
struct mg_str message) {
|
||||
struct mg_mqtt_opts opts = {};
|
||||
memset(&opts, 0, sizeof(opts));
|
||||
opts.topic = topic;
|
||||
opts.message = message;
|
||||
opts.qos = MQTT_QOS;
|
||||
mg_mqtt_pub(c, &opts);
|
||||
MG_DEBUG(("%lu PUBLISHED %.*s -> %.*s", c->id, topic.len, topic.buf,
|
||||
message.len, message.buf));
|
||||
}
|
||||
|
||||
static void rpc_ota_update(struct mg_rpc_req *r) {
|
||||
long ofs = mg_json_get_long(r->frame, "$.params.offset", -1);
|
||||
long tot = mg_json_get_long(r->frame, "$.params.total", -1);
|
||||
int len = 0;
|
||||
char *buf = mg_json_get_b64(r->frame, "$.params.chunk", &len);
|
||||
if (buf == NULL) {
|
||||
mg_rpc_err(r, 1, "%m", MG_ESC("Chunk decoding error"));
|
||||
} else if (ofs < 0 || tot < 0) {
|
||||
mg_rpc_err(r, 1, "%m", MG_ESC("offset and total not set"));
|
||||
} else if (ofs == 0 && mg_ota_begin((size_t) tot) == false) {
|
||||
mg_rpc_err(r, 1, "\"mg_ota_begin(%ld) failed\"", tot);
|
||||
} else if (len > 0 && mg_ota_write(buf, len) == false) {
|
||||
mg_rpc_err(r, 1, "\"mg_ota_write(%lu) @%ld failed\"", len, ofs);
|
||||
mg_ota_end();
|
||||
} else if (len == 0 && mg_ota_end() == false) {
|
||||
mg_rpc_err(r, 1, "\"mg_ota_end() failed\"", tot);
|
||||
} else {
|
||||
mg_rpc_ok(r, "%m", MG_ESC("ok"));
|
||||
}
|
||||
mg_free(buf);
|
||||
}
|
||||
|
||||
static void mqtt_ev_handler(struct mg_connection *c, int ev, void *ev_data) {
|
||||
if (ev == MG_EV_OPEN) {
|
||||
MG_DEBUG(("%lu CREATED, %s", c->id, ev_data));
|
||||
// c->is_hexdumping = 1;
|
||||
} else if (ev == MG_EV_CONNECT) {
|
||||
if (c->is_tls) {
|
||||
struct mg_tls_opts opts = {.ca = mg_unpacked("/certs/ca.pem"),
|
||||
.name = mg_url_host(MQTT_SERVER_URL)};
|
||||
mg_tls_init(c, &opts);
|
||||
}
|
||||
} else if (ev == MG_EV_ERROR) {
|
||||
// On error, log error message
|
||||
MG_ERROR(("%lu ERROR %s", c->id, (char *) ev_data));
|
||||
} else if (ev == MG_EV_MQTT_OPEN) {
|
||||
// MQTT connect is successful
|
||||
MG_DEBUG(("%lu CONNECTED to %s", c->id, MQTT_SERVER_URL));
|
||||
subscribe(c, mg_str(MQTT_SUBSCRIBE_TOPIC));
|
||||
} else if (ev == MG_EV_MQTT_MSG) {
|
||||
// When we get echo response, print it
|
||||
char response[100];
|
||||
struct mg_mqtt_message *mm = (struct mg_mqtt_message *) ev_data;
|
||||
mg_snprintf(response, sizeof(response), "Got %.*s -> %.*s", mm->topic.len,
|
||||
mm->topic.buf, mm->data.len, mm->data.buf);
|
||||
publish(c, mg_str(MQTT_PUBLISH_TOPIC), mg_str(response));
|
||||
} else if (ev == MG_EV_MQTT_CMD) {
|
||||
struct mg_mqtt_message *mm = (struct mg_mqtt_message *) ev_data;
|
||||
if (mm->cmd == MQTT_CMD_PINGREQ) mg_mqtt_pong(c);
|
||||
} else if (ev == MG_EV_CLOSE) {
|
||||
MG_ERROR(("%lu CLOSED", c->id));
|
||||
s_mqtt_conn = NULL; // Mark that we're closed
|
||||
}
|
||||
}
|
||||
|
||||
void mg_mqtt_init(struct mg_mgr *mgr) {
|
||||
(void) mgr;
|
||||
if (!s_rpc) mg_rpc_add(&s_rpc, mg_str("ota.update"), rpc_ota_update, NULL);
|
||||
}
|
||||
|
||||
void mg_mqtt_poll(struct mg_mgr *mgr) {
|
||||
static uint64_t timer = 1; // 1 triggers expiration on first poll
|
||||
|
||||
// Reconnect if connection is closed, and send MQTT PINGs to keep
|
||||
// the connection alive or to detect connection loss
|
||||
if (mg_timer_expired(&timer, RECONNECT_PERIOD_MS, mg_now())) {
|
||||
if (s_mqtt_conn == NULL) {
|
||||
struct mg_mqtt_opts opts = {.clean = true,
|
||||
.qos = MQTT_QOS,
|
||||
.topic = mg_str(MQTT_PUBLISH_TOPIC),
|
||||
.keepalive = 5,
|
||||
.version = 4,
|
||||
.message = mg_str("bye")};
|
||||
s_mqtt_conn =
|
||||
mg_mqtt_connect(mgr, MQTT_SERVER_URL, &opts, mqtt_ev_handler, NULL);
|
||||
} else {
|
||||
mg_mqtt_ping(s_mqtt_conn);
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user