| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332 |
- #include <regex.h>
- #include <time.h>
- #include "common.h"
- #include "cJSON.h"
- #include "mqtt.h"
- #include "thread.h"
- #include "mongoose.h"
- #include "elog.h"
- #include "lock.h"
- #include "sys.h"
- #include "Cellular.h"
- #include "xlist.h"
- #include "websocket_handle.h"
- #include "switch_ctrl.h"
- #include "sqlite_handle.h"
- #if(CHIP_TYPE == CHIP_T113s)
- #include <openssl/evp.h>
- #include <openssl/hmac.h>
- #endif
- #if 1
- #define LOGD log_d
- #define LOGE log_e
- #define LOGW log_w
- #else
- #define LOGD printf
- #define LOGE printf
- #define LOGW printf
- #endif
- #define CONN_PERIOD 5000 //ms
- #define SEND_PERIOD 3000 //ms
- #define BGET(flag,mask) ((flag)&(1<<(mask)))
- #define BSET(flag,mask) ((flag)|=(1<<(mask)))
- typedef struct mg_mgr mgr_t;
- typedef struct mg_timer mg_timer_t;
- typedef struct mg_mqtt_opts mg_opts_t;
- typedef struct mg_connection mg_conn_t;
- typedef struct {
- int id;
- char topic[256];
- char content[2048];
- }mqtt_pkt_t;
- typedef struct {
- mgr_t mgr;
- mg_conn_t *c;
- mqtt_server_t ser;
- pthread_t tid;
- int quit;
- void *h;
- }mqtt_conn_t;
- typedef struct {
- int inited;
- char prod_id[32];
-
- mqtt_conn_t conn[MQTT_SER_MAX];
- handle_t list;
- mg_timer_t *timer;
- GlobalDeviceManager *dm;
- }mqtt_handle_t;
- #ifdef USE_MQTT
- char *topic_sub[MQTT_SUB_MAX]={
- "/pdu/%s/control/power/all_channel",
- "/pdu/%s/control/power/+",
- "/pdu/%s/control/sensor/+",
- "/pdu/%s/control/device/restart",
- "/pdu/%s/control/device/reset",
- "/pdu/%s/control/service",
- };
- char *topic_pub[MQTT_PUB_MAX] = {
- "/pdu/%s/info/device",
- "/pdu/%s/status/network",
- "/pdu/%s/status/power/all_channel",
- "/pdu/%s/status/power/%d",
- "/pdu/%s/status/sensor/%d",
- "/pdu/%s/status/service",
- "/pdu/%s/alarm/network",
- "/pdu/%s/alarm/power",
- "/pdu/%s/alarm/sensor",
- };
- const char *serv_str[SERV_MAX]={
- "ntp",
- "smtp",
- "message",
- "mqtt",
- "cloud",
- "telnet",
- "snmp_v1",
- "snmp_v2c",
- "snmp_v3",
- "snmp_trap",
- };
- static mqtt_handle_t mqHandle={0};
- static int my_recv(char *topic, char *data);
- static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data);
- static int send_stat(mqtt_conn_t *conn, int type);
- static int get_serv(char *json);
- static void* client_thread(void *arg);
- /////////////////////////////////////////////////////////
- typedef struct {
- char user[64];
- char passwd[64];
- }user_account_t;
- typedef struct {
- char *instanceId;
- char *host;
- char *topic;
- char *groupId;
- char *clientId;
- char *accessKey;
- char *secretKey;
- uint16_t port;
- }mqtt_account_t;
- static mqtt_account_t mqtt_aliyun={
- "post-cn-0w73xnchz01",
- "post-cn-0w73xnchz01.mqtt.aliyuncs.com",
- "SmartPDU",
- "GID_PDU",
- "smartPDU_23542352",
- "LTAI5tAfuRj1JkB8ZjDS2p3i",
- "p4FrWlrFim9zM9A5mYxjK99EqKrSHh",
- 8883,
- };
- static int get_account(mqtt_account_t *host, user_account_t *user)
- {
- unsigned int len=0;
- char tempData[100];
- char clientIdUrl[64];
- //username和 Password 签名模式下的设置方法,参考文档 https://help.aliyun.com/document_detail/48271.html?spm=a2c4g.11186623.6.553.217831c3BSFry7
- sprintf(clientIdUrl, "%s@@@%s", host->groupId, host->clientId);
- HMAC(EVP_sha1(), host->secretKey, strlen(host->secretKey), (const unsigned char*)clientIdUrl, strlen(clientIdUrl), (unsigned char*)tempData, &len);
- int passwdLen = EVP_EncodeBlock((unsigned char *) user->passwd, (const unsigned char*)tempData, len);
- user->passwd[passwdLen] = '\0';
-
- sprintf(user->user,"Signature|%s|%s", host->accessKey, host->instanceId);
- return 0;
- }
- static void set_status(mqtt_conn_t *conn, int flag)
- {
- mqtt_handle_t *h=(mqtt_handle_t*)conn->h;
- int id=conn->ser.id;
- h->dm->mqttInfo.ser[id].status = flag;
- }
- static char *get_int_str(int n)
- {
- static char tmp[32];
- snprintf(tmp, sizeof(tmp), "%d", n);
- return tmp;
- }
- static char *get_float_str(float n)
- {
- static char tmp[32];
- snprintf(tmp, sizeof(tmp), "%0.3f", n);
- return tmp;
- }
- static int get_date_time(char *d, char *t)
- {
- time_t now;
- struct tm* tm;
- time(&now);
- tm = localtime(&now);
- if(!tm) {
- return -1;
- }
- sprintf(d, "%04d-%02d-%02d", tm->tm_year+1900, tm->tm_mon+1, tm->tm_mday);
- sprintf(t, "%02d:%02d:%02d", tm->tm_hour, tm->tm_min, tm->tm_sec);
- return 0;
- }
- static int sub_one(mg_conn_t *c, char *topic)
- {
- int r=-1;
- mg_opts_t opts={0};
-
- opts.topic = mg_str(topic);
- opts.qos = 1;
- mg_mqtt_sub(c, &opts);
- return 0;
- }
- static int my_sub(mg_conn_t *c, char * prod_id)
- {
- int i;
- char topic[1024];
-
- for(i=0; i<MQTT_SUB_MAX; i++) {
- snprintf(topic, sizeof(topic), topic_sub[i], prod_id);
- sub_one(c, topic);
- }
-
- return 0;
- }
- static int pub_one(mg_conn_t *c, char *topic, char *data)
- {
- int r=-1;
- mg_opts_t opts={0};
-
- opts.topic = mg_str(topic);
- opts.message = mg_str(data);
- opts.qos = 1;
- opts.retain = false;
- mg_mqtt_pub(c, &opts);
- return 0;
- }
- static int my_pub(mqtt_conn_t *conn, char *topic, char *data)
- {
- if(conn->c) {
- pub_one(conn->c, topic, data);
- }
- return 0;
- }
- static GlobalPowerManger* get_power(mqtt_handle_t *h, int ch)
- {
- GlobalPowerManger *tmp=NULL;
- list_for_each_entry(tmp, &h->dm->_globalPowerManger.list, list)
- {
- if(tmp->product_ch_id==ch) {
- return tmp;
- }
- }
- return NULL;
- }
- static int set_ch(mqtt_handle_t *h, int ch, int flag)
- {
- int r;
- GlobalPowerManger *tmp=NULL;
- list_for_each_entry(tmp, &h->dm->_globalPowerManger.list, list)
- {
- if(tmp==NULL) {
- break;
- }
- if(tmp->product_ch_id==ch) {
- r = g_switch_set_all_chn_ctrl(&h->dm->_globalRelaySampManger, tmp, tmp->product_saddr, tmp->product_ch_addr, flag, false);
- if(r<0) {
- LOGE("___ g_switch_set_all_chn_ctrl failed, saddr: %d, ch_addr: %d\n", tmp->product_saddr, tmp->product_ch_addr);
- }
- break;
- }
- if(ch==0xff) {
- g_switch_set_all_ctrl(&h->dm->_globalRelaySampManger, tmp->product_ch_type, tmp->product_saddr, flag);
- }
- }
- return 0;
- }
- static GlobalSensorManger* get_sensor(mqtt_handle_t *h, int id)
- {
- GlobalSensorManger *tmp=NULL;
- lock_s_hold(LOCK_ID_SENSOR);
- list_for_each_entry(tmp, &h->dm->_globalSensorManger.list, list)
- {
- if(tmp->sensor_id==id) {
- return tmp;
- }
- }
- lock_s_release(LOCK_ID_SENSOR);
- return NULL;
- }
- //////////////////////////////////////////////////////////////
- static int get_url(mqtt_server_t *ser, char *url)
- {
- int port=1883;
- char *head="mqtts";
- if(!ser->server[0] || !ser->port[0] || !ser->mode) {
- return -1;
- }
- port = atoi(ser->port);
- if(port==1883) {
- head = "mqtt";
- }
- else if(port==8883) {
- head = "mqtts";
- }
- else if(port==8083) {
- head = "ws";
- }
- else if(port==8884) {
- head = "wss";
- }
- sprintf(url, "%s://%s:%d", head, ser->server, port);
- return 0;
- }
- static int start_one(mqtt_conn_t *conn)
- {
- int r;
- if(conn->tid==0) {
- conn->quit = 0;
- r = pthread_create(&conn->tid, NULL, client_thread, conn);
- }
- return r;
- }
- static int stop_one(mqtt_conn_t *conn)
- {
- if(conn->tid) {
- conn->quit = 1;
- pthread_join(conn->tid, NULL);
- conn->tid = 0;
- }
- return 0;
- }
- static int stop_all(mqtt_handle_t *h)
- {
- int i;
- for(i=0; i<MQTT_SER_MAX; i++) {
- stop_one(&h->conn[i]);
- }
- return 0;
- }
- static int mg_mqtt_bind(mg_conn_t *c)
- {
- int r=0;
- if(cellular_is_connected()) {
- r = sys_bind_nic(NIC_USB, (int)c->fd);
- }
- return r;
- }
- static int conn_one(mqtt_conn_t *conn)
- {
- int r=-1;
- mqtt_handle_t *h=&mqHandle;
- if(!conn->c) {
- char url[1024];
- if(conn->ser.plat==MQTT_PLAT_ALIYUN) {
- user_account_t user;
- mqtt_aliyun.clientId = conn->ser.cid;
- get_account(&mqtt_aliyun, &user);
- strcpy(conn->ser.user, user.user);
- strcpy(conn->ser.password, user.passwd);
- }
- mg_opts_t opts={
- .clean = true,
- .qos = 1,
- .version = 4,
- .keepalive = 30,
- .topic = mg_str("hello"),
- .message = mg_str("bye"),
- .client_id = mg_str(conn->ser.cid),
- .user = mg_str(conn->ser.user),
- .pass = mg_str(conn->ser.password),
- };
- get_url(&conn->ser, url);
- conn->c = mg_mqtt_connect(&conn->mgr, url, &opts, mqtt_fn, conn);
- if(conn->c) {
- //mg_mqtt_bind(conn->c);
- sprintf(conn->ser.cid, "smartPDU_%s\n", h->prod_id);
- my_sub(conn->c, h->prod_id);
- r = 0;
- }
- }
- return r;
- }
- static int my_check(mqtt_handle_t *h)
- {
- int i,j,r=-1;
- char url[1024];
- mqtt_server_t *ser=h->dm->mqttInfo.ser;
- for(i=0; i<MQTT_SER_MAX; i++) {
- r = get_url(&ser[i], url);
- if((r==0)) {
- if(h->conn[i].tid==0) {
- h->conn[i].ser = ser[i];
- h->conn[i].h = h;
- start_one(&h->conn[i]);
- }
- }
- else {
- if(h->conn[i].tid) {
- stop_one(&h->conn[i]);
- }
- }
- }
- return 0;
- }
- static int my_send(mqtt_conn_t *conn)
- {
- int r=-1;
- mqtt_handle_t *h=(mqtt_handle_t*)conn->h;
- list_node_t *ln=NULL;
- r = xlist_take_node(h->list, &ln, 0);
- if(r==0) {
- mqtt_pkt_t *pkt=(mqtt_pkt_t*)ln->data.buf;
- my_pub(conn, pkt->topic, pkt->content);
- xlist_back_node(h->list, ln);
- }
- return r;
- }
- static int send_once(mqtt_conn_t *conn)
- {
- send_stat(conn, MQTT_PUB_INFO_DEVICE);
- send_stat(conn, MQTT_PUB_STAT_NETWORK);
- return 0;
- }
- static void send_period(mqtt_conn_t *conn)
- {
- send_stat(conn, MQTT_PUB_STAT_POWER_ALL);
- send_stat(conn, MQTT_PUB_STAT_POWER_CHN);
- send_stat(conn, MQTT_PUB_STAT_SENSOR);
- send_stat(conn, MQTT_PUB_STAT_SERVICE);
- send_stat(conn, MQTT_PUB_INFO_DEVICE);
- send_stat(conn, MQTT_PUB_STAT_NETWORK);
- }
- static void timer_conn_fn(void *arg)
- {
- mqtt_conn_t *conn=(mqtt_conn_t*)arg;
- conn_one(conn);
- }
- static void timer_period_fn(void *arg)
- {
- mqtt_conn_t *conn=(mqtt_conn_t*)arg;
- send_period(conn);
- }
- static void timer_add(mqtt_conn_t *conn)
- {
- mg_timer_add(&conn->mgr, CONN_PERIOD, MG_TIMER_REPEAT | MG_TIMER_RUN_NOW, timer_conn_fn, conn);
- mg_timer_add(&conn->mgr, SEND_PERIOD, MG_TIMER_REPEAT | MG_TIMER_RUN_NOW, timer_period_fn, conn);
- }
- static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
- {
- mqtt_conn_t *conn=((mqtt_conn_t*)(c->fn_data));
-
- switch(ev) {
- case MG_EV_OPEN:
- {
- // c->is_hexdumping = 1;
- }
- break;
- case MG_EV_CONNECT:
- {
- char url[1024];
- get_url(&conn->ser, url);
- if (mg_url_is_ssl(url)) {
- struct mg_tls_opts opts = {.ca = mg_str(conn->ser.cert),
- .name = mg_url_host(conn->ser.server)};
- mg_tls_init(c, &opts);
- }
- }
- break;
- case MG_EV_MQTT_OPEN:
- {
- send_once(conn);
- set_status(conn, 1);
- }
- break;
-
- case MG_EV_POLL:
- {
- my_send(conn);
- }
- break;
- case MG_EV_MQTT_MSG:
- {
- struct mg_mqtt_message *mm=(struct mg_mqtt_message*)ev_data;
- if(mm && mm->topic.buf && mm->data.buf) {
- my_recv(mm->topic.buf, mm->data.buf);
- }
- }
- break;
-
- case MG_EV_ERROR:
- {
- //MG_ERROR(("___MG_EV_ERROR, %p %s", c->fd, (char *) ev_data));
- //memset(mc->flag, 0, sizeof(mc->flag));
- }
- break;
- case MG_EV_CLOSE:
- {
- conn->c = NULL;
- set_status(conn, 0);
- }
- break;
- }
- }
- static void* client_thread(void *arg)
- {
- int r;
- mqtt_conn_t *conn=(mqtt_conn_t*)arg;
-
- mg_mgr_init(&conn->mgr);
- conn->mgr.dns4.url = "udp://114.114.114.114:53";
- //conn->mgr.dns6.url = "udp://114.114.114.114:53";
- timer_add(conn);
- while(conn->quit==0) {
- mg_mgr_poll(&conn->mgr, 1000);
- }
- if(conn->c)mg_mqtt_disconnect(conn->c, NULL);
- mg_mgr_free(&conn->mgr);
- conn->c = NULL;
- pthread_exit(NULL);
- }
- #endif
- static void* mqtt_thread(void *arg)
- {
- #ifdef USE_MQTT
- int r;
- thread_handle_t *th=(thread_handle_t*)arg;
- mqtt_handle_t *h=(mqtt_handle_t*)th->arg;
-
- while(th->quit==0) {
- my_check(h);
- sleep(1);
- }
- stop_all(h);
- #endif
- pthread_exit(NULL);
- }
- //////////////////////////////////////////////////////////////////////////
- int mqtt_init(void)
- {
- int r=-1;
- #ifdef USE_MQTT
- mqtt_handle_t *h=&mqHandle;
-
- memset(h, 0, sizeof(mqtt_handle_t));
- //mg_log_set(MG_LL_DEBUG);
-
- h->dm = &__globalDeviceManage;
- strcpy( h->prod_id,h->dm->_globalDevInfo.product.number);
- //sys_get_chip_id(&h->prod_id);
- list_cfg_t lc;
- lc.mode = LIST_FULL_FIFO;
- lc.max = 100;
- lc.log = 0;
- h->list = xlist_init(&lc);
- dev_mqtt_init(h->dm->db, &h->dm->mqttInfo);
- thread_start(THREAD_ID_MQTT, mqtt_thread, h);
- 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);
- xlist_free(h->list);
- h->inited = 0; r = 0;
- #endif
- return r;
- }
- static int send_stat(mqtt_conn_t *conn, int type)
- {
- int i,r=-1;
- #ifdef USE_MQTT
- float p1,p2,p3;
- char topic[256];
- char* content=NULL;
- char buf[20],date[40],time[40];
- mqtt_handle_t *h=(mqtt_handle_t*)conn->h;
- if((!conn->c) || (type<0) || (type>=MQTT_PUB_MAX)) {
- return -1;
- }
- if((type!=MQTT_PUB_STAT_POWER_CHN) && (type!=MQTT_PUB_STAT_SENSOR)) {
- snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id);
- }
- switch(type) {
- case MQTT_PUB_INFO_DEVICE:
- {
- cJSON* root=cJSON_CreateObject();
- if(root) {
- int nch = 0;
- GlobalPowerManger *chiterator;
- list_for_each_entry(chiterator, &h->dm->_globalPowerManger.list, list)
- {
- nch++;
- }
- int nsen = 0;
- GlobalSensorManger *sensoriterator;
- list_for_each_entry(sensoriterator, &h->dm->_globalSensorManger.list, list)
- {
- nsen++;
- }
- cJSON_AddStringToObject(root,"id", h->prod_id);
- cJSON_AddStringToObject(root,"name",h->dm->_globalDevInfo.product.name);
- cJSON_AddStringToObject(root,"type",g_product_type_str[h->dm->_globalDevInfo.product.pwr_type]);
- cJSON_AddStringToObject(root,"firmware", VERSION);
- cJSON_AddStringToObject(root,"channel_num", get_int_str(nch));
- cJSON_AddStringToObject(root,"sensor_num", get_int_str(nsen));
- get_date_time(date, time);
- cJSON_AddStringToObject(root,"date", date);
- cJSON_AddStringToObject(root,"time", time);
- content = cJSON_Print(root);
- my_pub(conn, topic, content);
- cJSON_free(content);
- cJSON_Delete(root);
- }
- }
- break;
- case MQTT_PUB_STAT_NETWORK:
- {
- NetworkInfo_t nw4,nw6;
- int r4 = sys_get_net(&nw4, IP_V4);
- int r6 = sys_get_net(&nw6, IP_V6);
- cJSON* root=cJSON_CreateObject();
- if(root) {
-
- if(r4==0) {
- cJSON_AddStringToObject(root,"lan_ipv4_link", nw4.mode?"1":"0");
- cJSON_AddStringToObject(root,"lan_ipv4_address", nw4.ip_address);
- cJSON_AddStringToObject(root,"lan_ipv4_mask", nw4.mask);
- }
- else {
- cJSON_AddStringToObject(root,"lan_ipv4_link", "");
- cJSON_AddStringToObject(root,"lan_ipv4_address", "");
- cJSON_AddStringToObject(root,"lan_ipv4_mask", "");
- }
- if(r4==0) {
- cJSON_AddStringToObject(root,"lan_ipv6_link", nw6.mode?"1":"0");
- cJSON_AddStringToObject(root,"lan_ipv6_address", nw6.ip_address);
- cJSON_AddStringToObject(root,"lan_ipv6_subnet_length", nw6.mask);
- }
- else {
- cJSON_AddStringToObject(root,"lan_ipv6_link", "");
- cJSON_AddStringToObject(root,"lan_ipv6_address", "");
- cJSON_AddStringToObject(root,"lan_ipv6_subnet_length", "");
- }
- if(0) {
- cJSON_AddStringToObject(root,"wifi_ipv4_link", "");
- cJSON_AddStringToObject(root,"wifi_ipv4_address", "");
- cJSON_AddStringToObject(root,"wifi_ipv4_mask", "");
- cJSON_AddStringToObject(root,"wifi_ipv6_link", "");
- cJSON_AddStringToObject(root,"wifi_ipv6_address", "");
- cJSON_AddStringToObject(root,"wifi_ipv6_mask", "");
- }
- cJSON_AddStringToObject(root,"modbus_address", get_int_str(h->dm->_globalDevInfo.cascade.addr));
- cJSON_AddStringToObject(root,"modbus_baud", get_int_str(h->dm->_globalDevInfo.cascade.baudrate));
- cJSON_AddStringToObject(root,"modbus_mode", get_int_str(h->dm->_globalDevInfo.cascade.mode));
- cJSON_AddStringToObject(root,"vpn_enable", "0");
- get_date_time(date, time);
- cJSON_AddStringToObject(root,"date", date);
- cJSON_AddStringToObject(root,"time", time);
- content = cJSON_Print(root);
- my_pub(conn, topic, content);
- cJSON_free(content);
- cJSON_Delete(root);
- }
- }
- break;
- case MQTT_PUB_STAT_POWER_ALL:
- {
- int cnt=1;
- pwrall_info_t *pall=websocket_get_pwrall();
- if (pall->ph3)
- {
- cnt = 3;
- }
- cJSON *root = cJSON_CreateObject();
- if (root)
- {
- cJSON *data_root_array = cJSON_CreateArray();
- for (i = 0; i < cnt; i++)
- {
- _OverAllPwrAckInfo *info = &pall->ch[i];
- cJSON *data_filed = cJSON_CreateObject();
- if (data_filed)
- {
- sprintf(buf, "L%d", i + 1);
- if (cnt == 3)
- cJSON_AddStringToObject(data_filed, "phase", buf);
- cJSON_AddStringToObject(data_filed, "voltage", get_float_str(info->voltage));
- cJSON_AddStringToObject(data_filed, "current", get_float_str(info->current));
- cJSON_AddStringToObject(data_filed, "power", get_float_str(info->power / 1000));
- cJSON_AddStringToObject(data_filed, "consumption", get_float_str(info->consumption));
- if (h->dm->_globalDevInfo.product.pwr_type != SmartPDU_DC)
- {
- cJSON_AddStringToObject(data_filed, "power_factor", get_float_str(info->factor));
- cJSON_AddStringToObject(data_filed, "power_freq", get_float_str(info->freq));
- p1 = info->power;
- p3 = info->factor>0 ? (p1 / info->factor) : 0.0f;
- p2 = p3 - p1;
- cJSON_AddStringToObject(data_filed, "pactive_power", get_float_str(p1));
- cJSON_AddStringToObject(data_filed, "reactive_power", get_float_str(p2));
- cJSON_AddStringToObject(data_filed, "apparent_power", get_float_str(p3));
- }
- cJSON_AddItemToObject(data_root_array, "power_status", data_filed);
- }
- }
- cJSON_AddItemToObject(root, "power_status", data_root_array);
- cJSON_AddStringToObject(root, "voltage_over", get_int_str(h->dm->_globalPowerManger._PowerWarninginfo.w_voltage_up));
- cJSON_AddStringToObject(root, "voltage_low", get_int_str(h->dm->_globalPowerManger._PowerWarninginfo.w_voltage_down));
- cJSON_AddStringToObject(root, "current_over", get_int_str(h->dm->_globalPowerManger._PowerWarninginfo.w_current));
- cJSON_AddStringToObject(root, "power_over", get_int_str(h->dm->_globalPowerManger._PowerWarninginfo.w_power));
- cJSON_AddStringToObject(root, "consumption_over", get_int_str(h->dm->_globalPowerManger._PowerWarninginfo.w_consumption));
- get_date_time(date, time);
- cJSON_AddStringToObject(root, "date", date);
- cJSON_AddStringToObject(root, "time", time);
- content = cJSON_Print(root);
- my_pub(conn, topic, content);
- cJSON_free(content);
- cJSON_Delete(root);
- }
- }
- break;
- case MQTT_PUB_STAT_POWER_CHN:
- {
- int quit=0;
- GlobalPowerManger *tmp=NULL;
- lock_s_hold(LOCK_ID_POWER_UPDATE);
- list_for_each_entry(tmp, &h->dm->_globalPowerManger.list, list)
- {
- if(tmp==NULL) {
- break;
- }
- cJSON* root=cJSON_CreateObject();
- if(root) {
- cJSON_AddStringToObject(root,"id", get_int_str(tmp->product_ch_id));
- cJSON_AddStringToObject(root,"name", tmp->product_ch_name);
- cJSON_AddStringToObject(root,"status", get_int_str(tmp->_PowerInfo.status));
- cJSON_AddStringToObject(root,"voltage", get_float_str(tmp->_PowerInfo.voltage));
- cJSON_AddStringToObject(root, "current", get_float_str(tmp->_PowerInfo.current));
- cJSON_AddStringToObject(root, "power", get_float_str(tmp->_PowerInfo.power));
- cJSON_AddStringToObject(root, "consumption", get_float_str(tmp->_PowerInfo.consumption));
- cJSON_AddStringToObject(root, "voltage_over", get_int_str(tmp->_PowerWarninginfo.w_voltage_up));
- cJSON_AddStringToObject(root, "voltage_low", get_int_str(tmp->_PowerWarninginfo.w_voltage_down));
- cJSON_AddStringToObject(root, "current_over", get_int_str(tmp->_PowerWarninginfo.w_current));
- cJSON_AddStringToObject(root, "power_over", get_int_str(tmp->_PowerWarninginfo.w_power));
- cJSON_AddStringToObject(root, "consumption_over", get_int_str(tmp->_PowerWarninginfo.w_consumption));
- if (h->dm->_globalDevInfo.product.pwr_type == SmartPDU_Tree_AC_Tree)
- {
- cJSON_AddStringToObject(root, "phase_loss_L1", get_int_str(tmp->_PowerWarninginfo.w_phase_lossA));
- cJSON_AddStringToObject(root, "phase_loss_L2", get_int_str(tmp->_PowerWarninginfo.w_phase_lossB));
- cJSON_AddStringToObject(root, "phase_loss_L3", get_int_str(tmp->_PowerWarninginfo.w_phase_lossC));
- }
- if (h->dm->_globalDevInfo.product.pwr_type != SmartPDU_DC)
- {
- cJSON_AddStringToObject(root, "power_freq", get_float_str(tmp->_PowerInfo.freq));
- cJSON_AddStringToObject(root, "power_factor", get_float_str(tmp->_PowerInfo.factor));
- p1 = tmp->_PowerInfo.power;
- p3 = tmp->_PowerInfo.factor>0 ? (p1 / tmp->_PowerInfo.factor) : 0.0f;
- p2 = p3 - p1;
- cJSON_AddStringToObject(root, "pactive_power", get_float_str(p1));
- cJSON_AddStringToObject(root, "reactive_power", get_float_str(p2));
- cJSON_AddStringToObject(root, "apparent_power", get_float_str(p3));
- }
- cJSON* data_root_array = cJSON_CreateArray();
- if (h->dm->_globalDevInfo.product.pwr_type == TREE_AC_TYPE|| h->dm->_globalDevInfo.product.pwr_type==DOUBLE_AC_TYPE|| h->dm->_globalDevInfo.product.pwr_type==SmartPDU_Tree_AC_One_B)
- {
- GlobalTreeACManager *_TreeACTemp = NULL;
- // 判断是否三相单输出情况下
- int phnum = 0;
- list_for_each_entry(_TreeACTemp, &tmp->list_Tree_AC, list_Tree_AC)
- {
- if(_TreeACTemp->product_ph_outputStatus == 1 &&_TreeACTemp->product_ph_type >= 0 && _TreeACTemp->product_ph_type < 3)
- {
- cJSON *data_filed = cJSON_CreateObject();
- cJSON_AddStringToObject(data_filed, "phase_type", g_product_t_ac_type_str[_TreeACTemp->product_ph_type]);
- cJSON_AddStringToObject(data_filed, "phase_voltage", get_float_str(_TreeACTemp->_PowerInfo.voltage));
- cJSON_AddStringToObject(data_filed, "phase_current", get_float_str(_TreeACTemp->_PowerInfo.current));
- cJSON_AddStringToObject(data_filed, "phase_power", get_float_str(_TreeACTemp->_PowerInfo.power));
- cJSON_AddStringToObject(data_filed, "phase_consumption", get_float_str(_TreeACTemp->_PowerInfo.consumption));
- cJSON_AddItemToObject(data_root_array, "phase_data", data_filed);
- }
- }
- }
- cJSON_AddItemToObject(root, "phase_data", data_root_array);
- get_date_time(date, time);
- cJSON_AddStringToObject(root,"date", date);
- cJSON_AddStringToObject(root,"time", time);
- snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id, tmp->product_ch_id);
- content = cJSON_Print(root);
- my_pub(conn, topic, content);
- cJSON_free(content);
- cJSON_Delete(root);
- }
- }
- lock_s_release(LOCK_ID_POWER_UPDATE);
- }
- break;
- case MQTT_PUB_STAT_SENSOR:
- {
- GlobalSensorManger *tmp=NULL;
- lock_s_hold(LOCK_ID_SENSOR);
- list_for_each_entry(tmp, &h->dm->_globalSensorManger.list, list)
- {
- if(tmp==NULL) {
- break;
- }
- cJSON* root=cJSON_CreateObject();
- if(root) {
- cJSON_AddStringToObject(root,"name", tmp->sensor_name);
- cJSON_AddStringToObject(root,"type", get_int_str(tmp->sensor_type));
- cJSON_AddStringToObject(root,"modbus_address", get_int_str(tmp->sensor_addr));
- cJSON_AddStringToObject(root,"node", get_int_str(tmp->sensor_node_number));
- cJSON_AddStringToObject(root,"status", get_int_str(tmp->sensor_status));
- cJSON_AddStringToObject(root,"value_num", get_int_str(tmp->sensor_val_count));
- cJSON_AddStringToObject(root,"value1", get_float_str(tmp->Cur_sensor_info.val1));
- cJSON_AddStringToObject(root,"value2", get_float_str(tmp->Cur_sensor_info.val2));
- cJSON_AddStringToObject(root, "value1_over", get_int_str(tmp->warning_status.sensor_val1_upper));
- cJSON_AddStringToObject(root, "value1_low", get_int_str(tmp->warning_status.sensor_val1_lower));
- cJSON_AddStringToObject(root, "value2_over", get_int_str(tmp->warning_status.sensor_val2_upper));
- cJSON_AddStringToObject(root, "value2_low", get_int_str(tmp->warning_status.sensor_val2_lower));
- get_date_time(date, time);
- cJSON_AddStringToObject(root,"date", date);
- cJSON_AddStringToObject(root,"time", time);
- snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id, tmp->sensor_id);
- content = cJSON_Print(root);
- my_pub(conn, topic, content);
- cJSON_free(content);
- cJSON_Delete(root);
- }
- }
- lock_s_release(LOCK_ID_SENSOR);
- }
- break;
- case MQTT_PUB_STAT_SERVICE:
- {
- cJSON* root=cJSON_CreateObject();
- if(root) {
-
- int flag = get_serv(NULL);
- cJSON_AddStringToObject(root,"telnet", BGET(flag, SERV_TELNET)?"1":"0");
- cJSON_AddStringToObject(root,"smtp", BGET(flag, SERV_SMTP)?"1":"0");
- cJSON_AddStringToObject(root,"snmp_v1", BGET(flag, SERV_SNMP_V1)?"1":"0");
- cJSON_AddStringToObject(root,"snmp_v2c", BGET(flag, SERV_SNMP_V2C)?"1":"0");
- cJSON_AddStringToObject(root,"snmp_v3", BGET(flag, SERV_SNMP_V3)?"1":"0");
- cJSON_AddStringToObject(root,"snmp_trap", BGET(flag, SERV_SNMP_TRAP)?"1":"0");
- cJSON_AddStringToObject(root,"message", BGET(flag, SERV_MESG)?"1":"0");
- cJSON_AddStringToObject(root,"mqtt", BGET(flag, SERV_MQTT)?"1":"0");
- cJSON_AddStringToObject(root,"cloud", BGET(flag, SERV_CLOUD)?"1":"0");
- cJSON_AddStringToObject(root,"ntp", BGET(flag, SERV_NTP)?"1":"0");
- get_date_time(date, time);
- cJSON_AddStringToObject(root,"date", date);
- cJSON_AddStringToObject(root,"time", time);
- content = cJSON_Print(root);
- my_pub(conn, topic, content);
- cJSON_free(content);
- cJSON_Delete(root);
- }
- }
- break;
- default:
- return -1;
- }
- #endif
- return r;
- }
- static int post_alarm(int type, int subtype, char *content, alarm_para_t *para)
- {
- int r=-1;
- #ifdef USE_MQTT
- char topic[256];
- mqtt_handle_t *h=&mqHandle;
- char date[40],time[40];
- int pub_type;
- char *json_str=NULL;
- if(!h->inited) {
- return -1;
- }
-
- if(type==ALARM_TYPE_POWER) {
- pub_type = MQTT_PUB_ALARM_POWER;
- }
- else if(type==ALARM_TYPE_SENSOR) {
- pub_type = MQTT_PUB_ALARM_SENSOR;
- }
- else if(type==ALARM_TYPE_NETWORK) {
- pub_type = MQTT_PUB_ALARM_NETWORK;
- }
- else {
- return -1;
- }
- snprintf(topic, sizeof(topic), topic_pub[pub_type], h->prod_id);
- get_date_time(date, time);
- switch(pub_type) {
- case MQTT_PUB_ALARM_NETWORK:
- case MQTT_PUB_ALARM_SENSOR:
- case MQTT_PUB_ALARM_POWER:
- {
- cJSON* root=cJSON_CreateObject();
- if(root) {
-
- if(MQTT_PUB_ALARM_SENSOR==pub_type)cJSON_AddStringToObject(root,"sensor_id", get_int_str(para->id));
- else if(MQTT_PUB_ALARM_POWER==pub_type)cJSON_AddStringToObject(root,"channel_id", get_int_str(para->id));
- else cJSON_AddStringToObject(root,"id", get_int_str(para->id));
- cJSON_AddStringToObject(root, "context", content);
- if (MQTT_PUB_ALARM_POWER == pub_type)
- {
- cJSON_AddStringToObject(root, "action", get_int_str(para->action));
- cJSON_AddStringToObject(root, "action_para", get_int_str(para->actionId));
- }
- cJSON_AddStringToObject(root, "date", date);
- cJSON_AddStringToObject(root, "time", time);
- json_str = cJSON_Print(root);
- cJSON_Delete(root);
- }
- }
- break;
- }
- mqtt_pkt_t pkt;
- pkt.id = pub_type;
- snprintf(pkt.topic, sizeof(pkt.topic), "%s", topic);
- snprintf(pkt.content, sizeof(pkt.content), "%s", json_str);
- r = xlist_append(h->list, 0, &pkt, sizeof(pkt));
- cJSON_free(json_str);
- #endif
- return r;
- }
- int mqtt_post_alarm(int type, int subtype, char *content, alarm_para_t *para)
- {
- return post_alarm(type, subtype, content, para);
- }
- #ifdef USE_MQTT
- //////////////////////////////////////////////////////////
- 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;
- }
- static int get_serv(char *json)
- {
- int i,flag=0;
- cJSON* tmp=NULL;
- if(json) {
- 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);
- }
- }
- else {
- if(thread_is_running(THREAD_ID_NTP)) {
- BSET(flag, SERV_NTP);
- }
- if(thread_is_running(THREAD_ID_MAIL)) {
- BSET(flag, SERV_SMTP);
- }
- if(thread_is_running(THREAD_ID_MQTT)) {
- BSET(flag, SERV_MQTT);
- }
- #if 0
- if(thread_is_running(SERV_MESG)) {
- BSET(flag, SERV_MESG);
- }
- if(thread_is_running(THREAD_ID_CLOUD)) {
- BSET(flag, SERV_CLOUD);
- }
- if(thread_is_running(THREAD_ID_TELNET)) {
- BSET(flag, SERV_TELNET);
- }
- #endif
- if(sys_is_running("snmpd")) {
- BSET(flag, SERV_NTP);
- }
- if(thread_is_running(THREAD_ID_NTP)) {
- BSET(flag, SERV_NTP);
- }
- }
- return flag;
- }
- static int set_serv(int flag)
- {
- if(BGET(flag,SERV_NTP)) {
- if(!thread_is_running(THREAD_ID_NTP)) {
- //thread_stop(THREAD_ID_NTP);
- }
- }
- else {
- if(thread_is_running(THREAD_ID_NTP)) {
- thread_stop(THREAD_ID_NTP);
- }
- }
- if(BGET(flag,SERV_SMTP)) {
- if(!thread_is_running(THREAD_ID_MAIL)) {
- //thread_stop(THREAD_ID_MAIL);
- }
- }
- else {
- if(thread_is_running(THREAD_ID_MAIL)) {
- thread_stop(THREAD_ID_MAIL);
- }
- }
- if(BGET(flag,SERV_MQTT)) {
- }
- else {
-
- }
- if(BGET(flag,SERV_MESG)) {
- }
- else {
-
- }
- if(BGET(flag,SERV_CLOUD)) {
- }
- else {
-
- }
- if(BGET(flag,SERV_TELNET)) {
- }
- else {
-
- }
- if(BGET(flag,SERV_SNMP_V1) || BGET(flag,SERV_SNMP_V2C) || BGET(flag,SERV_SNMP_V3) || BGET(flag,SERV_SNMP_TRAP)) {
- }
- else {
-
- }
- return 0;
- }
- static int my_recv(char *topic, char *data)
- {
- int i,r,cmd;
- #ifdef USE_MQTT
- char *p,temp[512];
- int flag=0;
- mqtt_handle_t *h=&mqHandle;
- for(i=0; i<MQTT_SUB_MAX; i++) {
- snprintf(temp, sizeof(temp), topic_sub[i], h->prod_id);
- p = strrchr(temp, '+');
- if(p) p[0] = 0;
-
- if(strstr(topic, temp)) {
- switch(i) {
- case MQTT_SUB_CMD_POWER_ALL:
- {
- cmd = get_cmd(data);
- LOGD("___ MQTT_SUB_CMD_POWER_ALL, %d\n", cmd);
- if(cmd==0 || cmd==1) {
- set_ch(h, 0xff, cmd);
- }
- }
- break;
- case MQTT_SUB_CMD_POWER_CHN:
- {
- int ch=atoi(topic+strlen(temp));
- cmd = get_cmd(data);
- LOGD("___ MQTT_SUB_CMD_POWER_CHN, %d, %d\n", ch, cmd);
- if(cmd==0 || cmd==1) {
- set_ch(h, ch, cmd);
- }
- }
- break;
- case MQTT_SUB_CMD_SENSOR:
- {
- int id=atoi(topic+strlen(temp));
- cmd = get_cmd(data);
- switch(id) {
- case SENSOR_TYPE_TEMP_HUMI:
- break;
- case SENSOR_TYPE_SMOKE:
- break;
- case SENSOR_TYPE_WATER:
- break;
- case SENSOR_TYPE_ACCESS:
- break;
- case SENSOR_TYPE_GAS:
- break;
- case SENSOR_TYPE_AIR_PRESS:
- break;
- case SENSOR_TYPE_TEMP_DOUBLE:
- break;
- case SENSOR_TYPE_TEMP:
- break;
- case SENSOR_TYPE_LEAK:
- break;
- }
- LOGD("___ MQTT_SUB_CMD_POWER_CHN, %d\n", id);
- }
- break;
- case MQTT_SUB_CMD_RESTART:
- {
- cmd = get_cmd(data);
- LOGD("___ MQTT_SUB_CMD_RESTART, %d\n", cmd);
- if(cmd==0) {
- system("shutdown");
- }
- else if(cmd==1) {
- system("reboot");
- }
- }
- break;
- case MQTT_SUB_CMD_RESET:
- {
- cmd = get_cmd(data);
- LOGD("___ MQTT_SUB_CMD_RESET, %d\n", cmd);
- if(cmd==1) {
- sys_set_factory();
- }
- }
- break;
- case MQTT_SUB_CMD_SERVICE:
- {
- flag = get_serv(data);
- LOGD("___ MQTT_SUB_CMD_SERVICE, 0x%08x\n", flag);
- set_serv(flag);
- }
- break;
- }
- }
- }
- #endif
- return 0;
- }
- #endif
- int mqtt_test(void)
- {
- #ifdef USE_MQTT
- int cnt=0;
- mqtt_handle_t *h=&mqHandle;
- mqtt_info_t mInfo={0};
- mqtt_server_t ser={
- .id = 0,
- .mode = 1,
- .cid = "mqtlx_e92334545tet3465",
- .name = "serverX",
- .server = "192.168.1.12",
- .port = "1883",
- //root
- //root0219107X
- .user = "gowone100",
- .password = "gowone100",
- //.user = "gowone101",
- //.password = "gowone101",
- };
- mInfo.ser[0] = ser;
- h->dm->mqttInfo = mInfo;
- #endif
- return 0;
- }
|