#include #include #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 #include #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; ic) { 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; iconn[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; iconn[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; ich[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; ivaluestring) { flag |= (atoi(tmp->valuestring)<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; }