|
|
@@ -2,8 +2,9 @@
|
|
|
#include "thread.h"
|
|
|
#include "mongoose.h"
|
|
|
#include "elog.h"
|
|
|
+#include "lock.h"
|
|
|
|
|
|
-#if 1
|
|
|
+#if 0
|
|
|
#define LOGD log_d
|
|
|
#define LOGE log_e
|
|
|
#define LOGW log_w
|
|
|
@@ -24,6 +25,7 @@ char *topic_sub[]={
|
|
|
"/pdu/%d/control/device/restart",
|
|
|
"/pdu/%d/control/device/reset",
|
|
|
"/PDU/%d/control/service",
|
|
|
+ NULL,
|
|
|
};
|
|
|
char *topic_pub[]={
|
|
|
"/pdu/%d/info/device",
|
|
|
@@ -35,6 +37,7 @@ char *topic_pub[]={
|
|
|
"/pdu/%d/alarm/power",
|
|
|
"/pdu/%d/alarm/sensor",
|
|
|
"/PDU/%d/status/service",
|
|
|
+ NULL,
|
|
|
};
|
|
|
|
|
|
|
|
|
@@ -49,7 +52,7 @@ typedef struct {
|
|
|
mg_conn_t *conn;
|
|
|
mqtt_para_t para;
|
|
|
|
|
|
- pthread_mutex_t mutex;
|
|
|
+ lock_t lck;
|
|
|
int inited;
|
|
|
bool isover;
|
|
|
|
|
|
@@ -91,7 +94,7 @@ static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
|
|
|
}
|
|
|
|
|
|
if (ev == MG_EV_ERROR || ev == MG_EV_CLOSE) {
|
|
|
- MG_INFO(("Got event %d, stopping...", ev));
|
|
|
+ MG_INFO(("got event %d, stopping...", ev));
|
|
|
*(bool *) c->fn_data = true; // Signal that we're done
|
|
|
}
|
|
|
}
|
|
|
@@ -103,20 +106,39 @@ static void* mqtt_thread(void *arg)
|
|
|
|
|
|
while(th->quit==0) {
|
|
|
if(h->inited) {
|
|
|
- pthread_mutex_lock(&h->mutex);
|
|
|
+ lock_s_hold(LOCK_ID_MQTT);
|
|
|
if (h->isover) {
|
|
|
h->conn = mg_mqtt_connect(&h->mgr, h->para.user.url, &h->opts, mqtt_fn, &h->isover);
|
|
|
}
|
|
|
else {
|
|
|
mg_mgr_poll(&h->mgr, 300);
|
|
|
}
|
|
|
- pthread_mutex_unlock(&h->mutex);
|
|
|
+ lock_s_release(LOCK_ID_MQTT);
|
|
|
+ lock_s_release(LOCK_ID_MQTT);
|
|
|
}
|
|
|
|
|
|
if(h->isover) usleep(1000);
|
|
|
}
|
|
|
pthread_exit(NULL);
|
|
|
}
|
|
|
+
|
|
|
+static int sub_all(mqtt_handle_t *h, int id, int qos)
|
|
|
+{
|
|
|
+ int i;
|
|
|
+ mg_opts_t opts;
|
|
|
+ char temp[512];
|
|
|
+
|
|
|
+ memset(&opts, 0, sizeof(opts));
|
|
|
+
|
|
|
+ opts.qos = qos;
|
|
|
+ for(i=0; topic_sub[i]; i++) {
|
|
|
+ sprintf(temp, topic_sub[i], id);
|
|
|
+ opts.topic = mg_str(temp);
|
|
|
+ mg_mqtt_sub(h->conn, &opts);
|
|
|
+ }
|
|
|
+
|
|
|
+ return 0;
|
|
|
+}
|
|
|
//////////////////////////////////////////////////////////////////////////
|
|
|
int mqtt_init(void)
|
|
|
{
|
|
|
@@ -131,12 +153,6 @@ int mqtt_init(void)
|
|
|
h->para.conn.ver = 4;
|
|
|
h->para.conn.qos = 1;
|
|
|
|
|
|
- r = pthread_mutex_init(&h->mutex, NULL);
|
|
|
- if(r) {
|
|
|
- LOGE("___ mqtt mutex init failed\n");
|
|
|
- return -1;
|
|
|
- }
|
|
|
-
|
|
|
thread_start(THREAD_ID_MQTT, mqtt_thread, h, 4*MB, 0);
|
|
|
h->inited = 1;
|
|
|
|
|
|
@@ -150,7 +166,6 @@ int mqtt_deinit(void)
|
|
|
|
|
|
thread_stop(THREAD_ID_MQTT);
|
|
|
mg_mgr_free(&h->mgr);
|
|
|
- pthread_mutex_destroy(&h->mutex);
|
|
|
h->inited = 0;
|
|
|
|
|
|
return 0;
|
|
|
@@ -159,17 +174,21 @@ int mqtt_deinit(void)
|
|
|
|
|
|
int mqtt_set(mqtt_para_t *para)
|
|
|
{
|
|
|
+ int r=0;
|
|
|
mqtt_handle_t *h=&mqHandle;
|
|
|
|
|
|
if(!para) {
|
|
|
return -1;
|
|
|
}
|
|
|
|
|
|
- //para check
|
|
|
+ lock_s_hold(LOCK_ID_MQTT);
|
|
|
if(para->conn.ver<3 || para->conn.ver>5) {
|
|
|
- return -1;
|
|
|
+ r = -1;
|
|
|
+ }
|
|
|
+ else {
|
|
|
+ h->para = *para;
|
|
|
}
|
|
|
- h->para = *para;
|
|
|
+ lock_s_release(LOCK_ID_MQTT);
|
|
|
|
|
|
return 0;
|
|
|
}
|
|
|
@@ -179,7 +198,7 @@ int mqtt_conn(void)
|
|
|
{
|
|
|
mqtt_handle_t *h=&mqHandle;
|
|
|
|
|
|
- pthread_mutex_lock(&h->mutex);
|
|
|
+ lock_s_hold(LOCK_ID_MQTT);
|
|
|
if(!h->inited) {
|
|
|
goto quit;
|
|
|
}
|
|
|
@@ -199,7 +218,7 @@ int mqtt_conn(void)
|
|
|
h->conn = mg_mqtt_connect(&h->mgr, h->para.user.url, &h->opts, mqtt_fn, &h->isover);
|
|
|
|
|
|
quit:
|
|
|
- pthread_mutex_unlock(&h->mutex);
|
|
|
+ lock_s_release(LOCK_ID_MQTT);
|
|
|
return h->conn?0:-1;
|
|
|
}
|
|
|
|
|
|
@@ -209,7 +228,7 @@ int mqtt_disconn(void)
|
|
|
int r=0;
|
|
|
mqtt_handle_t *h=&mqHandle;
|
|
|
|
|
|
- pthread_mutex_lock(&h->mutex);
|
|
|
+ lock_s_hold(LOCK_ID_MQTT);
|
|
|
if(!h->inited) {
|
|
|
r = -1;
|
|
|
goto quit;
|
|
|
@@ -217,7 +236,7 @@ int mqtt_disconn(void)
|
|
|
mg_mqtt_disconnect(h->conn, NULL);
|
|
|
|
|
|
quit:
|
|
|
- pthread_mutex_unlock(&h->mutex);
|
|
|
+ lock_s_release(LOCK_ID_MQTT);
|
|
|
return r;
|
|
|
}
|
|
|
|
|
|
@@ -228,7 +247,7 @@ int mqtt_sub(char *topic, int qos)
|
|
|
mg_opts_t opts;
|
|
|
mqtt_handle_t *h=&mqHandle;
|
|
|
|
|
|
- pthread_mutex_lock(&h->mutex);
|
|
|
+ lock_s_hold(LOCK_ID_MQTT);
|
|
|
if(!h->inited || !h->conn) {
|
|
|
r = -1;
|
|
|
goto quit;
|
|
|
@@ -240,7 +259,7 @@ int mqtt_sub(char *topic, int qos)
|
|
|
mg_mqtt_sub(h->conn, &opts);
|
|
|
|
|
|
quit:
|
|
|
- pthread_mutex_unlock(&h->mutex);
|
|
|
+ lock_s_release(LOCK_ID_MQTT);
|
|
|
return r;
|
|
|
}
|
|
|
|
|
|
@@ -251,7 +270,7 @@ int mqtt_pub(char *topic, int qos, char *data, int dlen)
|
|
|
mg_opts_t opts;
|
|
|
mqtt_handle_t *h=&mqHandle;
|
|
|
|
|
|
- pthread_mutex_lock(&h->mutex);
|
|
|
+ lock_s_hold(LOCK_ID_MQTT);
|
|
|
if(!h->inited || !h->conn) {
|
|
|
r = -1;
|
|
|
goto quit;
|
|
|
@@ -266,7 +285,7 @@ int mqtt_pub(char *topic, int qos, char *data, int dlen)
|
|
|
r = mg_mqtt_pub(h->conn, &opts);
|
|
|
|
|
|
quit:
|
|
|
- pthread_mutex_unlock(&h->mutex);
|
|
|
+ lock_s_release(LOCK_ID_MQTT);
|
|
|
return r;
|
|
|
}
|
|
|
|