|
|
@@ -13,8 +13,33 @@
|
|
|
#define LOGW printf
|
|
|
#endif
|
|
|
|
|
|
+
|
|
|
#define MQTT_SERVER_URL "XXXXXXXX"
|
|
|
|
|
|
+
|
|
|
+char *topic_sub[]={
|
|
|
+ "/pdu/%d/control/power/all_channel",
|
|
|
+ "/pdu/%d/control/power/%d",
|
|
|
+ "/pdu/%d/control/sensor/%d",
|
|
|
+ "/pdu/%d/control/device/restart",
|
|
|
+ "/pdu/%d/control/device/reset",
|
|
|
+ "/PDU/%d/control/service",
|
|
|
+};
|
|
|
+char *topic_pub[]={
|
|
|
+ "/pdu/%d/info/device",
|
|
|
+ "/PDU/%d/status/network",
|
|
|
+ "/pdu/%d/status/power/all_channel",
|
|
|
+ "/pdu/%d/status/power/%d",
|
|
|
+ "/pdu/%d/status/sensor/%d",
|
|
|
+ "/pdu/%d/alarm/network",
|
|
|
+ "/pdu/%d/alarm/power",
|
|
|
+ "/pdu/%d/alarm/sensor",
|
|
|
+ "/PDU/%d/status/service",
|
|
|
+};
|
|
|
+
|
|
|
+
|
|
|
+
|
|
|
+
|
|
|
typedef struct mg_mgr mgr_t;
|
|
|
typedef struct mg_mqtt_opts mg_opts_t;
|
|
|
typedef struct mg_connection mg_conn_t;
|
|
|
@@ -22,7 +47,7 @@ typedef struct {
|
|
|
mgr_t mgr;
|
|
|
mg_opts_t opts;
|
|
|
mg_conn_t *conn;
|
|
|
- conn_para_t para;
|
|
|
+ mqtt_para_t para;
|
|
|
|
|
|
pthread_mutex_t mutex;
|
|
|
int inited;
|
|
|
@@ -31,30 +56,38 @@ typedef struct {
|
|
|
}mqtt_handle_t;
|
|
|
static mqtt_handle_t mqHandle;
|
|
|
|
|
|
-static void mqtt_fn(struct mg_connection *c, int ev, void *ev_data)
|
|
|
+static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
|
|
|
{
|
|
|
mqtt_handle_t *h=&mqHandle;
|
|
|
|
|
|
if (ev == MG_EV_OPEN) {
|
|
|
// c->is_hexdumping = 1;
|
|
|
} else if (ev == MG_EV_CONNECT) {
|
|
|
- if (mg_url_is_ssl(h->para.url)) {
|
|
|
+ if (mg_url_is_ssl(h->para.user.url)) {
|
|
|
struct mg_tls_opts opts = {.ca = mg_unpacked("/certs/ca.pem"),
|
|
|
- .name = mg_url_host(h->para.url)};
|
|
|
+ .name = mg_url_host(h->para.user.url)};
|
|
|
mg_tls_init(c, &opts);
|
|
|
}
|
|
|
} else if (ev == MG_EV_ERROR) {
|
|
|
// On error, log error message
|
|
|
MG_ERROR(("%p %s", c->fd, (char *) ev_data));
|
|
|
} else if (ev == MG_EV_MQTT_OPEN) {
|
|
|
-
|
|
|
+ mg_opts_t opts={
|
|
|
+ .client_id = mg_str(h->para.user.url),
|
|
|
+ .user = mg_str(h->para.user.name),
|
|
|
+ .pass = mg_str(h->para.user.pass),
|
|
|
+ };
|
|
|
+ size_t len=c->send.len;
|
|
|
+
|
|
|
+ mg_mqtt_login(c, &opts);
|
|
|
+ mg_ws_wrap(c, c->send.len - len, WEBSOCKET_OP_BINARY);
|
|
|
} else if (ev == MG_EV_MQTT_MSG) {
|
|
|
// When we receive MQTT message, print it
|
|
|
struct mg_mqtt_message *mm = (struct mg_mqtt_message *) ev_data;
|
|
|
- MG_INFO(("Received on %.*s : %.*s", (int) mm->topic.len, mm->topic.buf,
|
|
|
- (int) mm->data.len, mm->data.buf));
|
|
|
- } else if (ev == MG_EV_POLL && c->data[0] == 'X') {
|
|
|
+ //MG_INFO(("Received on %.*s : %.*s", (int) mm->topic.len, mm->topic.buf, (int) mm->data.len, mm->data.buf));
|
|
|
|
|
|
+
|
|
|
+
|
|
|
}
|
|
|
|
|
|
if (ev == MG_EV_ERROR || ev == MG_EV_CLOSE) {
|
|
|
@@ -65,20 +98,22 @@ static void mqtt_fn(struct mg_connection *c, int ev, void *ev_data)
|
|
|
static void* mqtt_thread(void *arg)
|
|
|
{
|
|
|
int r;
|
|
|
- thread_handle_t *h=(thread_handle_t*)arg;
|
|
|
- mqtt_handle_t *mh=(mqtt_handle_t*)h->arg;
|
|
|
+ thread_handle_t *th=(thread_handle_t*)arg;
|
|
|
+ mqtt_handle_t *h=(mqtt_handle_t*)th->arg;
|
|
|
|
|
|
- while(h->quit==0) {
|
|
|
- pthread_mutex_lock(&mh->mutex);
|
|
|
- if (mh->isover) {
|
|
|
- mh->conn = mg_mqtt_connect(&mh->mgr, mh->para.url, &mh->opts, mqtt_fn, &mh->isover);
|
|
|
+ while(th->quit==0) {
|
|
|
+ if(h->inited) {
|
|
|
+ pthread_mutex_lock(&h->mutex);
|
|
|
+ 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);
|
|
|
}
|
|
|
- else {
|
|
|
- mg_mgr_poll(&mh->mgr, 300);
|
|
|
- }
|
|
|
- pthread_mutex_unlock(&mh->mutex);
|
|
|
|
|
|
- if(mh->isover) usleep(1000);
|
|
|
+ if(h->isover) usleep(1000);
|
|
|
}
|
|
|
pthread_exit(NULL);
|
|
|
}
|
|
|
@@ -91,9 +126,10 @@ int mqtt_init(void)
|
|
|
memset(h, 9, sizeof(mqtt_handle_t));
|
|
|
mg_mgr_init(&h->mgr);
|
|
|
|
|
|
- h->para.ver = 4;
|
|
|
- h->para.qos = 1;
|
|
|
- strcpy(h->para.url, MQTT_SERVER_URL);
|
|
|
+ //h->para.user = ;
|
|
|
+
|
|
|
+ h->para.conn.ver = 4;
|
|
|
+ h->para.conn.qos = 1;
|
|
|
|
|
|
r = pthread_mutex_init(&h->mutex, NULL);
|
|
|
if(r) {
|
|
|
@@ -120,11 +156,17 @@ int mqtt_deinit(void)
|
|
|
return 0;
|
|
|
}
|
|
|
|
|
|
-int mqtt_set(conn_para_t *para)
|
|
|
+
|
|
|
+int mqtt_set(mqtt_para_t *para)
|
|
|
{
|
|
|
mqtt_handle_t *h=&mqHandle;
|
|
|
|
|
|
- if(!para || para->ver<3 || para->ver>5) {
|
|
|
+ if(!para) {
|
|
|
+ return -1;
|
|
|
+ }
|
|
|
+
|
|
|
+ //para check
|
|
|
+ if(para->conn.ver<3 || para->conn.ver>5) {
|
|
|
return -1;
|
|
|
}
|
|
|
h->para = *para;
|
|
|
@@ -144,9 +186,9 @@ int mqtt_conn(void)
|
|
|
|
|
|
memset(&h->opts, 0, sizeof(h->opts));
|
|
|
h->opts.clean = true,
|
|
|
- h->opts.qos = h->para.qos,
|
|
|
+ h->opts.qos = h->para.conn.qos,
|
|
|
h->opts.topic = mg_str(""),
|
|
|
- h->opts.version = h->para.ver,
|
|
|
+ h->opts.version = h->para.conn.ver,
|
|
|
h->opts.message = mg_str("bye");
|
|
|
|
|
|
if(h->conn) {
|
|
|
@@ -154,7 +196,7 @@ int mqtt_conn(void)
|
|
|
}
|
|
|
|
|
|
h->isover = false;
|
|
|
- h->conn = mg_mqtt_connect(&h->mgr, h->para.url, &h->opts, mqtt_fn, &h->isover);
|
|
|
+ h->conn = mg_mqtt_connect(&h->mgr, h->para.user.url, &h->opts, mqtt_fn, &h->isover);
|
|
|
|
|
|
quit:
|
|
|
pthread_mutex_unlock(&h->mutex);
|
|
|
@@ -221,7 +263,7 @@ int mqtt_pub(char *topic, int qos, char *data, int dlen)
|
|
|
opts.message.buf = data;
|
|
|
opts.message.len = dlen;
|
|
|
opts.retain = false;
|
|
|
- mg_mqtt_pub(h->conn, &opts);
|
|
|
+ r = mg_mqtt_pub(h->conn, &opts);
|
|
|
|
|
|
quit:
|
|
|
pthread_mutex_unlock(&h->mutex);
|