| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604 |
- #include <regex.h>
- #include "common.h"
- #include "cJSON.h"
- #include "mqtt.h"
- #include "thread.h"
- #include "mongoose.h"
- #include "elog.h"
- #include "lock.h"
- #include "cfg.h"
- #if 0
- #define LOGD log_d
- #define LOGE log_e
- #define LOGW log_w
- #else
- #define LOGD printf
- #define LOGE printf
- #define LOGW printf
- #endif
- #ifdef USE_MQTT
- char *topic_sub[MQTT_SUB_MAX]={
- "/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[MQTT_PUB_MAX]={
- "/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;
- typedef struct {
- uint8_t qos;
- uint8_t ver;
- uint8_t clean;
- uint8_t retain;
- }mqtt_para_t;
- typedef struct _mqtt_conn_t{
- mg_conn_t *c;
- mg_opts_t opts;
- mqtt_para_t para;
- mqtt_info_t *info;
- bool isover;
- struct _mqtt_conn_t *next;
- }mqtt_conn_t;
- typedef struct {
- mgr_t mgr;
- mqtt_conn_t *conn;
- int inited;
- int prod_id;
- GlobalDeviceManager *dm;
- }mqtt_handle_t;
- static mqtt_handle_t mqHandle={0};
- static int my_conn(mqtt_handle_t *h, mqtt_info_t *info);
- static int my_recv(char *topic, char *data);
- static int info_cmp(mqtt_info_t *a, mqtt_info_t *b)
- {
- return memcmp(a, b, sizeof(mqtt_info_t)-sizeof(mqtt_info_t*));
- }
- 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.user.url)) {
- struct mg_tls_opts opts = {.ca = mg_unpacked("/certs/ca.pem"),
- .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));
- my_recv(mm->topic.buf, mm->data.buf);
- }
- if (ev == MG_EV_ERROR || ev == MG_EV_CLOSE) {
- MG_INFO(("got event %d, stopping...", ev));
- *(bool *) c->fn_data = true; // Signal that we're done
- }
- }
- static void* mqtt_thread(void *arg)
- {
- int r;
- thread_handle_t *th=(thread_handle_t*)arg;
- mqtt_handle_t *h=(mqtt_handle_t*)th->arg;
-
- while(th->quit==0) {
- if(h->inited) {
-
- my_conn_check(h);
- lock_s_hold(LOCK_ID_MQTT);
- mg_mgr_poll(&h->mgr, 300);
- lock_s_release(LOCK_ID_MQTT);
- }
- }
- pthread_exit(NULL);
- }
- static int my_sub(mg_conn_t *c, int prod_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], prod_id);
- opts.topic = mg_str(temp);
- mg_mqtt_sub(c, &opts);
- }
-
- return 0;
- }
- static int find_conn(mqtt_handle_t *h, mqtt_info_t *info)
- {
- //
- return 0;
- }
- static int my_conn(mqtt_handle_t *h, mqtt_info_t *info)
- {
- int r=-1;
- mqtt_conn_t *conn,*c;
-
- h->opts.clean = true,
- h->opts.qos = h->para.conn.qos,
- h->opts.topic = mg_str(""),
- h->opts.version = h->para.conn.ver,
- h->opts.message = mg_str("bye");
- conn = h->conn;
- while(info) {
- if(!conn==NULL) {
- conn = calloc(1, sizeof(mqtt_conn_t));
- if(!conn) {
- LOGE("___my_conn calloc failed\n");
- break;
- }
- }
- conn->c = mg_mqtt_connect(&h->mgr, conn->para.user.url, &conn->opts, mqtt_fn, &conn->isover);
- if(conn>c) {
- my_sub(conn>c, h->prod_id, 1);
- }
- info = info->next;
- if(!info) {
- r = 0; break;
- }
- }
- return r;
- }
- static int my_disconn(mqtt_handle_t *h, mqtt_conn_t *c)
- {
- mqtt_conn_t *c1,*c2;
-
- c1 = c2 = h->conn;
- while(c1) {
- if(c1==c) {
- mg_mqtt_disconnect(c1, &c1->opts);
- break;
- }
- c1 = c1->next;
- }
- h->conn = NULL;
- }
- static int my_disconn_all(mqtt_handle_t *h)
- {
- mqtt_conn_t *c,*conn=h->conn;
- while(conn) {
- c = conn;
- mg_mqtt_disconnect(c, &c->opts);
- conn = conn->next;
- free(c);
- }
- h->conn = NULL;
- }
- static int my_conn_check(mqtt_handle_t *h)
- {
- int r=-1;
- mqtt_conn_t *conn=h->conn;
- while(conn) {
- if(conn->c && conn->isover) {
- conn->c = mg_mqtt_connect(&h->mgr, conn->para.user.url, &conn->opts, mqtt_fn, &conn->isover);
- if(conn>c) {
- my_sub(conn>c, h->prod_id, 1); r = 0;
- }
- }
- conn = conn->next;
- }
- return r;
- }
- #endif
- //////////////////////////////////////////////////////////////////////////
- int mqtt_init(void)
- {
- int r=-1;
- #ifdef USE_MQTT
- mqtt_handle_t *h=&mqHandle;
-
- memset(h, 0, sizeof(mqtt_handle_t));
- mg_mgr_init(&h->mgr);
-
- h->dm = &__globalDeviceManage;
- h->prod_id = h->dm->_globalDevInfo.product_id;
- //h->para.user = ;
- h->para.conn.ver = 4;
- h->para.conn.qos = 1;
-
- thread_start(THREAD_ID_MQTT, mqtt_thread, h, 4*MB, 0);
- h->inited = 1;
- r = 0;
- #endif
- return r;
- }
- int mqtt_deinit(void)
- {
- int r=-1;
- #ifdef USE_MQTT
- mqtt_handle_t *h=&mqHandle;
-
- thread_stop(THREAD_ID_MQTT);
- mg_mgr_free(&h->mgr);
- h->inited = 0;
- r = 0;
- #endif
- return r;
- }
- int mqtt_conn(void)
- {
- int r=-1;
- #ifdef USE_MQTT
- mqtt_handle_t *h=&mqHandle;
-
- lock_s_hold(LOCK_ID_MQTT);
- if(!h->inited) {
- goto quit;
- }
-
- if(h->conn) {
- mg_mqtt_disconnect(h->conn, NULL);
- }
- my_conn(h);
-
- quit:
- lock_s_release(LOCK_ID_MQTT);
- r = h->conn?0:-1;
- #endif
- return r;
- }
- int mqtt_disconn(void)
- {
- int r=-1;
- #ifdef USE_MQTT
- mqtt_handle_t *h=&mqHandle;
- lock_s_hold(LOCK_ID_MQTT);
- if(!h->inited) {
- r = -1;
- goto quit;
- }
- mg_mqtt_disconnect(h->conn, NULL);
- r = 0;
- quit:
- lock_s_release(LOCK_ID_MQTT);
- #endif
- return r;
- }
- int mqtt_sub(char *topic, int qos)
- {
- int r=-1;
- #ifdef USE_MQTT
- mg_opts_t opts;
- mqtt_handle_t *h=&mqHandle;
- lock_s_hold(LOCK_ID_MQTT);
- if(!h->inited || !h->conn) {
- r = -1;
- goto quit;
- }
-
- memset(&opts, 0, sizeof(opts));
- opts.topic = mg_str(topic);
- opts.qos = qos;
- mg_mqtt_sub(h->conn, &opts);
- r = 0;
- quit:
- lock_s_release(LOCK_ID_MQTT);
- #endif
- return r;
- }
- int mqtt_pub(char *topic, int qos, char *data, int dlen)
- {
- int r=-1;
- #ifdef USE_MQTT
- mg_opts_t opts;
- mqtt_handle_t *h=&mqHandle;
- if(!h->inited) {
- return -1;
- }
- memset(&opts, 0, sizeof(opts));
- opts.topic = mg_str(topic);
- opts.qos = qos;
- opts.message.buf = data;
- opts.message.len = dlen;
- opts.retain = false;
- while(1) {
- lock_s_hold(LOCK_ID_MQTT);
- mg_mqtt_pub(h->conn, &opts);
- lock_s_release(LOCK_ID_MQTT);
- }
- #endif
- return r;
- }
- int mqtt_send(int type, int id, void *data)
- {
- int r=-1;
- #ifdef USE_MQTT
- char topic[512];
- char content[4096];
- mqtt_handle_t *h=&mqHandle;
- if(type<0 || type>=MQTT_PUB_MAX || !data) {
- return -1;
- }
- if(type==MQTT_PUB_STAT_POWER_CHN || type==MQTT_PUB_STAT_SENSOR) {
- snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id, id);
- }
- else {
- snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id);
- }
- switch(type) {
- case MQTT_PUB_INFO_DEVICE:
- {
- //
- }
- break;
- case MQTT_PUB_STAT_NETWORK:
- {
- }
- break;
- case MQTT_PUB_STAT_POWER_ALL:
- {
-
- }
- break;
- case MQTT_PUB_STAT_POWER_CHN:
- {
-
- }
- break;
- case MQTT_PUB_STAT_SENSOR:
- {
-
- }
- break;
- case MQTT_PUB_ALARM_NETWORK:
- {
-
- }
- break;
- case MQTT_PUB_ALARM_POWER:
- {
-
- }
- break;
- case MQTT_PUB_ALARM_SENSOR:
- {
-
- }
- break;
- }
- //r = mqtt_pub(topic, );
- #endif
- return r;
- }
- //////////////////////////////////////////////////////////
- static int get_cmd(char *json)
- {
- int cmd=-1;
- cJSON* cjson=cJSON_Parse(json);
- if(cjson) {
- cJSON* order=cJSON_GetObjectItem(cjson,"order");
- if(order && order->valuestring) {
- cmd = atoi(order->valuestring);
- }
- cJSON_Delete(cjson);
- }
- return cmd;
- }
- enum {
- SERV_NTP=0,
- SERV_SMTP,
- SERV_MESG,
- SERV_MQTT,
- SERV_CLOUD,
- SERV_TELNET,
- SERV_SNMP_V1,
- SERV_SNMP_V2C,
- SERV_SNMP_V3,
- SERV_SNMP_TRAP,
-
- SERV_MAX
- };
- const char *serv_str[SERV_MAX]={
- "ntp",
- "smtp",
- "message",
- "mqtt",
- "cloud",
- "telnet",
- "snmp_v1",
- "snmp_v2c",
- "snmp_v3",
- "snmp_trap",
- };
- static uint32_t get_serv(char *json)
- {
- int i;
- cJSON* tmp=NULL;
- uint32_t flag=0;
- cJSON* cjson=cJSON_Parse(json);
- if(cjson) {
- for(i=0; i<SERV_MAX; i++) {
- tmp = cJSON_GetObjectItem(cjson, serv_str[i]);
- if(tmp && tmp->valuestring) {
- flag |= (atoi(tmp->valuestring)<<i);
- }
- }
- cJSON_Delete(cjson);
- }
- return flag;
- }
- static int my_recv(char *topic, char *data)
- {
- int i,r,cmd;
- #ifdef USE_MQTT
- char temp[2000];
- mqtt_handle_t *h=&mqHandle;
- for(i=0; i<MQTT_SUB_MAX; i++) {
- snprintf(temp, sizeof(temp), topic_sub[i], h->prod_id);
- if(strstr(topic, temp)) {
- switch(i) {
- case MQTT_SUB_CMD_POWER_ALL:
- {
- LOGD("___ MQTT_SUB_CMD_POWER_ALL\n");
- cmd = get_cmd(data);
- }
- break;
- case MQTT_SUB_CMD_POWER_CHN:
- {
- int ch=atoi(topic+strlen(temp));
- cmd = get_cmd(data);
- LOGD("___ MQTT_SUB_CMD_POWER_CHN, %d\n", ch);
- }
- break;
- case MQTT_SUB_CMD_SENSOR:
- {
- int id=atoi(topic+strlen(temp));
- cmd = get_cmd(data);
- LOGD("___ MQTT_SUB_CMD_POWER_CHN, %d\n", id);
- }
- break;
- case MQTT_SUB_CMD_RESTART:
- {
- LOGD("___ MQTT_SUB_CMD_RESTART\n");
- cmd = get_cmd(data);
-
- }
- break;
- case MQTT_SUB_CMD_RESET:
- {
- LOGD("___ MQTT_SUB_CMD_RESET\n");
- cmd = get_cmd(data);
- }
- break;
- case MQTT_SUB_CMD_SERVICE:
- {
- LOGD("___ MQTT_SUB_CMD_SERVICE\n");
- uint32_t flag=get_serv(data);
- }
- break;
- }
- }
- }
- #endif
- return 0;
- }
|