From 43aa8288ca884fdf2aa5e92e30f385fd8ac2b51c Mon Sep 17 00:00:00 2001 From: fengzero Date: Mon, 16 Jun 2025 02:29:50 +0000 Subject: [PATCH 1/2] fix: mqtt disable subscribe topic --- plugins/mqtt/mqtt_config.c | 1 + plugins/mqtt/mqtt_plugin_intf.c | 62 +++++++++++---------------------- 2 files changed, 21 insertions(+), 42 deletions(-) diff --git a/plugins/mqtt/mqtt_config.c b/plugins/mqtt/mqtt_config.c index fa52a04f4..422cfcde6 100644 --- a/plugins/mqtt/mqtt_config.c +++ b/plugins/mqtt/mqtt_config.c @@ -495,6 +495,7 @@ int mqtt_config_parse(neu_plugin_t *plugin, const char *setting, plog_notice(plugin, "config qos : %d", config->qos); plog_notice(plugin, "config format : %s", mqtt_upload_format_str(config->format)); + plog_notice(plugin, "config enable-topic : %d", config->enable_topic); plog_notice(plugin, "config write-req-topic : %s", config->write_req_topic); plog_notice(plugin, "config write-resp-topic: %s", config->write_resp_topic); diff --git a/plugins/mqtt/mqtt_plugin_intf.c b/plugins/mqtt/mqtt_plugin_intf.c index 970c8c2d4..6ed6a920d 100644 --- a/plugins/mqtt/mqtt_plugin_intf.c +++ b/plugins/mqtt/mqtt_plugin_intf.c @@ -387,7 +387,6 @@ int mqtt_plugin_config(neu_plugin_t *plugin, const char *setting) int rv = 0; const char * plugin_name = neu_plugin_module.module_name; mqtt_config_t config = { 0 }; - bool started = false; rv = plugin->parse_config(plugin, setting, &config); if (0 != rv) { @@ -395,35 +394,13 @@ int mqtt_plugin_config(neu_plugin_t *plugin, const char *setting) return NEU_ERR_NODE_SETTING_INVALID; } - if (NULL == plugin->client) { - plugin->client = neu_mqtt_client_new(config.version); - if (NULL == plugin->client) { - plog_error(plugin, "neu_mqtt_client_new fail"); - rv = NEU_ERR_EINTERNAL; - goto error; - } - } else if (neu_mqtt_client_is_open(plugin->client)) { - started = true; - if (plugin->config.enable_topic) { - plugin->unsubscribe(plugin, &plugin->config); - } - rv = neu_mqtt_client_close(plugin->client); - if (0 != rv) { - plog_error(plugin, "neu_mqtt_client_close fail"); - rv = NEU_ERR_EINTERNAL; - goto error; - } - if (neu_mqtt_client_check_version_change(plugin->client, - config.version)) { - neu_mqtt_client_free(plugin->client); - plugin->client = neu_mqtt_client_new(config.version); + if (plugin->client != NULL) { + if (neu_mqtt_client_is_open(plugin->client)) { + rv = neu_mqtt_client_close(plugin->client); } - } else if (neu_mqtt_client_check_version_change(plugin->client, - config.version)) { - // plugin stopped and version changed neu_mqtt_client_free(plugin->client); - plugin->client = neu_mqtt_client_new(config.version); } + plugin->client = neu_mqtt_client_new(config.version); rv = config_mqtt_client(plugin, plugin->client, &config); if (0 != rv) { @@ -431,22 +408,20 @@ int mqtt_plugin_config(neu_plugin_t *plugin, const char *setting) goto error; } - if (started) { - if (0 != neu_mqtt_client_open(plugin->client)) { - plog_error(plugin, "neu_mqtt_client_open fail"); - rv = NEU_ERR_MQTT_CONNECT_FAILURE; - goto error; - } - if (0 != start_hearbeat_timer(plugin, config.heartbeat_interval)) { - plog_error(plugin, "start hearbeat_timer failed"); - rv = NEU_ERR_EINTERNAL; + if (0 != neu_mqtt_client_open(plugin->client)) { + plog_error(plugin, "neu_mqtt_client_open fail"); + rv = NEU_ERR_MQTT_CONNECT_FAILURE; + goto error; + } + if (0 != start_hearbeat_timer(plugin, config.heartbeat_interval)) { + plog_error(plugin, "start hearbeat_timer failed"); + rv = NEU_ERR_EINTERNAL; + goto error; + } + if (config.enable_topic) { + if (0 != (rv = plugin->subscribe(plugin, &config))) { goto error; } - if (plugin->config.enable_topic) { - if (0 != (rv = plugin->subscribe(plugin, &config))) { - goto error; - } - } } if (plugin->config.host) { @@ -487,7 +462,10 @@ int mqtt_plugin_start(neu_plugin_t *plugin) goto end; } - rv = plugin->subscribe(plugin, &plugin->config); + rv = 0; + if (plugin->config.enable_topic) { + rv = plugin->subscribe(plugin, &plugin->config); + } end: if (0 == rv) { From fba83e6519b4bb8938e591d47e766871970aaa20 Mon Sep 17 00:00:00 2001 From: fengzero Date: Mon, 16 Jun 2025 02:39:03 +0000 Subject: [PATCH 2/2] fix cppcheck --- plugins/mqtt/azure_iot_plugin.c | 1 + plugins/mqtt/mqtt_plugin_intf.c | 4 +--- 2 files changed, 2 insertions(+), 3 deletions(-) diff --git a/plugins/mqtt/azure_iot_plugin.c b/plugins/mqtt/azure_iot_plugin.c index c5caa55ab..0c7c5a43f 100644 --- a/plugins/mqtt/azure_iot_plugin.c +++ b/plugins/mqtt/azure_iot_plugin.c @@ -184,6 +184,7 @@ static int azure_parse_config(neu_plugin_t *plugin, const char *setting, config->cert = cert.v.val_str; config->key = key.v.val_str; config->upload_drv_state = false; + config->enable_topic = true; plog_notice(plugin, "config MQTT version : %d", config->version); plog_notice(plugin, "config client-id : %s", config->client_id); diff --git a/plugins/mqtt/mqtt_plugin_intf.c b/plugins/mqtt/mqtt_plugin_intf.c index 6ed6a920d..00e36e491 100644 --- a/plugins/mqtt/mqtt_plugin_intf.c +++ b/plugins/mqtt/mqtt_plugin_intf.c @@ -395,9 +395,7 @@ int mqtt_plugin_config(neu_plugin_t *plugin, const char *setting) } if (plugin->client != NULL) { - if (neu_mqtt_client_is_open(plugin->client)) { - rv = neu_mqtt_client_close(plugin->client); - } + neu_mqtt_client_close(plugin->client); neu_mqtt_client_free(plugin->client); } plugin->client = neu_mqtt_client_new(config.version);