Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions plugins/mqtt/azure_iot_plugin.c
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
1 change: 1 addition & 0 deletions plugins/mqtt/mqtt_config.c
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
62 changes: 19 additions & 43 deletions plugins/mqtt/mqtt_plugin_intf.c
Original file line number Diff line number Diff line change
Expand Up @@ -387,66 +387,39 @@ 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) {
plog_error(plugin, "neu_mqtt_config_parse fail");
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);
}
} else if (neu_mqtt_client_check_version_change(plugin->client,
config.version)) {
// plugin stopped and version changed
if (plugin->client != NULL) {
neu_mqtt_client_close(plugin->client);
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) {
rv = NEU_ERR_MQTT_INIT_FAILURE;
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) {
Expand Down Expand Up @@ -487,7 +460,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) {
Expand Down