| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249 |
- #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 "cfg.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 5000 //ms
- 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
- };
- #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;
- uint32_t prod_id;
-
- 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/%d/control/power/all_channel",
- "/pdu/%d/control/power/+",
- "/pdu/%d/control/sensor/+",
- "/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",
- };
- 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, int 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 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 = 60,
- .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) {
- sprintf(conn->ser.cid, "smartPDU_%d\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);
- send_stat(conn, MQTT_PUB_STAT_SENSOR);
- send_stat(conn, MQTT_PUB_STAT_SERVICE);
- 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);
- }
- 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);
- }
- 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;
- h->prod_id = h->dm->_globalDevInfo.product.id;
- //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) {
- cJSON_AddStringToObject(root,"id", get_int_str(h->prod_id));
- cJSON_AddStringToObject(root,"name", "Gowone smartPDU");
- cJSON_AddStringToObject(root,"type", "smartPDU AC");
- cJSON_AddStringToObject(root,"firmware", VERSION);
- cJSON_AddStringToObject(root,"channel_num", "8");
- cJSON_AddStringToObject(root,"sensor_num", "10");
- 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;
- }
- for(i=0; i<cnt; i++) {
- _OverAllPwrAckInfo *info=&pall->ch[i];
- cJSON* root=cJSON_CreateObject();
- if(root) {
- sprintf(buf, "L%d", i+1);
- cJSON_AddStringToObject(root,"phase", buf);
- cJSON_AddStringToObject(root,"voltage", get_float_str(info->voltage));
- cJSON_AddStringToObject(root,"current", get_float_str(info->current));
- cJSON_AddStringToObject(root,"power", get_float_str(info->power/1000));
- cJSON_AddStringToObject(root,"consumption", get_float_str(info->consumption));
- p1 = info->power;
- p3 = (info->factor?(p1/info->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));
- 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) {
- sprintf(buf, "CH%d", tmp->product_ch_id);
- cJSON_AddStringToObject(root,"name", buf);
- 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,"power_freq", get_float_str(tmp->_PowerInfo.freq));
- cJSON_AddStringToObject(root,"consumption", get_float_str(tmp->_PowerInfo.consumption));
- cJSON_AddStringToObject(root,"power_factor", get_float_str(tmp->_PowerInfo.factor));
-
- p1 = tmp->_PowerInfo.power;
- p3 = p1 / tmp->_PowerInfo.factor;
- 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));
- 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", "temp sensor");
- 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));
-
- 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) {
- cJSON_AddStringToObject(root,"num", get_int_str(para->nalarm)); //alarm number, how to get?
- cJSON_AddStringToObject(root,"context", content);
- 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;
- }
|