Procházet zdrojové kódy

1、继续完善mqtt,准备测试

guohui před 2 roky
rodič
revize
3722576c94

+ 1 - 1
pro/src/appweb_handle.c

@@ -3530,8 +3530,8 @@ static void serviceManagement(void *conn)
         mqtt_server_t *ser=&gdm->mqttInfo.ser[req->idx];
         mqtt_server_t *ser=&gdm->mqttInfo.ser[req->idx];
 
 
         if(req->idx |= ser->id) { ser->id = req->idx; save_flag=1; }
         if(req->idx |= ser->id) { ser->id = req->idx; save_flag=1; }
+        if(req->name) { strcpy(ser->name, req->name); save_flag=1; }
         if(req->server) { strcpy(ser->server, req->server); save_flag=1; }
         if(req->server) { strcpy(ser->server, req->server); save_flag=1; }
-        if(req->ip) { strcpy(ser->ip, req->ip); save_flag=1; }
         if(req->port) { strcpy(ser->port, req->port); save_flag=1; }
         if(req->port) { strcpy(ser->port, req->port); save_flag=1; }
         if(req->user) { strcpy(ser->user, req->user); save_flag=1; }
         if(req->user) { strcpy(ser->user, req->user); save_flag=1; }
         if(req->password) { strcpy(ser->password, req->password); save_flag=1; }
         if(req->password) { strcpy(ser->password, req->password); save_flag=1; }

+ 3 - 3
pro/src/common.h

@@ -638,9 +638,9 @@ typedef struct
 typedef struct {
 typedef struct {
     int                 id;    		                    //id
     int                 id;    		                    //id
     int                 mode;    				        //0: open   1: close
     int                 mode;    				        //0: open   1: close
-    char                proto[10];                      //mqtt, mqtts, ws, wss
-    char                server[512];                    //server
-    char                ip[48]; 		                //server ip
+    char                proto[16];                      //mqtt, mqtts, ws, wss
+    char                name[32];                       //server name
+    char                server[512]; 		            //server ip
     char                port[10]; 			            //server port
     char                port[10]; 			            //server port
     char                cid[32];    		            //client ID
     char                cid[32];    		            //client ID
     char                user[32];    			        //认证方式  
     char                user[32];    			        //认证方式  

+ 2 - 2
pro/src/dflt.c

@@ -59,8 +59,8 @@ mqtt_info_t DFLT_MQTT={
         .mode = 1,
         .mode = 1,
         .cid = "",
         .cid = "",
         .proto = "mqtt",
         .proto = "mqtt",
+        .name  = "server1",
         .server = "192.168.1.12",
         .server = "192.168.1.12",
-        .ip = "",
         .port = "1883",
         .port = "1883",
         .cert = "",
         .cert = "",
         .user = "gowone100",
         .user = "gowone100",
@@ -72,8 +72,8 @@ mqtt_info_t DFLT_MQTT={
         .mode = 1,
         .mode = 1,
         .cid = "",
         .cid = "",
         .proto = "mqtts",
         .proto = "mqtts",
+        .name  = "server2",
         .server = "gaae4d7b.ala.cn-hangzhou.emqxsl.cn",
         .server = "gaae4d7b.ala.cn-hangzhou.emqxsl.cn",
-        .ip = "",
         .port = "8883",
         .port = "8883",
         .cert = "-----BEGIN CERTIFICATE-----\r\n"
         .cert = "-----BEGIN CERTIFICATE-----\r\n"
                 "MIIDrzCCApegAwIBAgIQCDvgVpBCRrGhdWrJWZHHSjANBgkqhkiG9w0BAQUFADBh\r\n"
                 "MIIDrzCCApegAwIBAgIQCDvgVpBCRrGhdWrJWZHHSjANBgkqhkiG9w0BAQUFADBh\r\n"

+ 4 - 4
pro/src/json_handle.c

@@ -4440,12 +4440,12 @@ cJSON* json_to_service_manage(const char* str,_ServiceManageRequestInfo* _servic
 
 
         tmp = cJSON_GetObjectItem(cjson, "serverName");
         tmp = cJSON_GetObjectItem(cjson, "serverName");
         if(tmp) {
         if(tmp) {
-            info->server = tmp->valuestring;
+            info->name = tmp->valuestring;
         }
         }
 
 
         tmp = cJSON_GetObjectItem(cjson, "ip");
         tmp = cJSON_GetObjectItem(cjson, "ip");
         if(tmp) {
         if(tmp) {
-            info->ip = tmp->valuestring;
+            info->server = tmp->valuestring;
         }
         }
 
 
         tmp = cJSON_GetObjectItem(cjson, "port");
         tmp = cJSON_GetObjectItem(cjson, "port");
@@ -5044,8 +5044,8 @@ char* service_mqtt_ack_to_json(mqtt_info_t* info)
                 sprintf(buf, "%d", ser->mode);
                 sprintf(buf, "%d", ser->mode);
                 cJSON_AddStringToObject(tmp,"mode", buf);
                 cJSON_AddStringToObject(tmp,"mode", buf);
 
 
-                cJSON_AddStringToObject(tmp,"serverName", ser->server);
-                cJSON_AddStringToObject(tmp,"ip", ser->ip);
+                cJSON_AddStringToObject(tmp,"serverName", ser->name);
+                cJSON_AddStringToObject(tmp,"ip", ser->server);
                 cJSON_AddStringToObject(tmp,"port", ser->port);
                 cJSON_AddStringToObject(tmp,"port", ser->port);
                 cJSON_AddStringToObject(tmp,"yhm", ser->user);
                 cJSON_AddStringToObject(tmp,"yhm", ser->user);
                 cJSON_AddStringToObject(tmp,"mm", ser->password);
                 cJSON_AddStringToObject(tmp,"mm", ser->password);

+ 1 - 1
pro/src/json_handle.h

@@ -357,8 +357,8 @@ typedef struct
 {
 {
     int   idx;
     int   idx;
     int   mode;
     int   mode;
+    char* name;
     char* server;
     char* server;
-    char* ip;
     char* port;
     char* port;
     char* cid;
     char* cid;
     char* user;
     char* user;

+ 146 - 104
pro/src/mqtt.c

@@ -10,6 +10,7 @@
 #include "sys.h"
 #include "sys.h"
 #include "cfg.h"
 #include "cfg.h"
 #include "xlist.h"
 #include "xlist.h"
+#include "websocket_handle.h"
 #include "switch_ctrl.h"
 #include "switch_ctrl.h"
 #include "sqlite_handle.h"
 #include "sqlite_handle.h"
 
 
@@ -24,8 +25,10 @@
 #endif
 #endif
 
 
 
 
+#define CONN_PERIOD         5000    //ms
 #define SEND_PERIOD         5000    //ms
 #define SEND_PERIOD         5000    //ms
 
 
+
 #ifdef USE_MQTT
 #ifdef USE_MQTT
 
 
 char *topic_sub[MQTT_SUB_MAX]={
 char *topic_sub[MQTT_SUB_MAX]={
@@ -80,9 +83,9 @@ static mqtt_handle_t mqHandle={0};
 static int my_recv(char *topic, char *data);
 static int my_recv(char *topic, char *data);
 static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data);
 static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data);
 static void timer_start(mqtt_handle_t *h);
 static void timer_start(mqtt_handle_t *h);
-static int post_stat(int type, int id);
+static int send_stat(mqtt_handle_t *h, int type);
 static uint32_t get_serv(char *json);
 static uint32_t get_serv(char *json);
-static void post_once(void);
+static void send_once(mqtt_handle_t *h);
 
 
 enum {
 enum {
     SERV_NTP=0,
     SERV_NTP=0,
@@ -101,31 +104,7 @@ enum {
 #define BGET(flag,mask)  ((flag)&(1<<(mask)))
 #define BGET(flag,mask)  ((flag)&(1<<(mask)))
 #define BSET(flag,mask)  ((flag)|=(1<<(mask)))
 #define BSET(flag,mask)  ((flag)|=(1<<(mask)))
 
 
-static int ser_cmp(mqtt_server_t *a, mqtt_server_t *b)
-{
-    if(strcmp(a->server, b->server) || 
-       strcmp(a->ip, b->ip) || 
-       strcmp(a->port, b->port) || 
-       strcmp(a->user, b->user) ||
-       strcmp(a->password, b->password) ||
-       (a->mode!=b->mode)) {
-        return 1;
-    }
 
 
-    return 0;
-}
-static int get_id(mqtt_handle_t *h, mqtt_server_t *ser)
-{
-    int i,r=-1;
-    mqtt_conn_t *conn=h->conn;
-
-    for(i=0; i<MQTT_SER_MAX; i++) {
-        if(conn[i].c && ser_cmp(ser, &conn[i].ser)==0) {
-            return i;
-        }
-    }
-    return -1;
-}
 static char *get_int_str(int n)
 static char *get_int_str(int n)
 {
 {
     static char tmp[32];
     static char tmp[32];
@@ -135,7 +114,7 @@ static char *get_int_str(int n)
 static char *get_float_str(float n)
 static char *get_float_str(float n)
 {
 {
     static char tmp[32];
     static char tmp[32];
-    snprintf(tmp, sizeof(tmp), "%f", n);
+    snprintf(tmp, sizeof(tmp), "%0.3f", n);
     return tmp;
     return tmp;
 }
 }
 static int get_date_time(char *d, char *t)
 static int get_date_time(char *d, char *t)
@@ -287,7 +266,7 @@ static int get_url(mqtt_server_t *ser, char *url)
     int port=1883;
     int port=1883;
     char *head="mqtts";
     char *head="mqtts";
 
 
-    if((!ser->server[0] && !ser->ip[0]) || !ser->port[0]) {
+    if(!ser->server[0] || !ser->port[0]) {
         return -1;
         return -1;
     }
     }
 
 
@@ -305,8 +284,7 @@ static int get_url(mqtt_server_t *ser, char *url)
         head = "wss";
         head = "wss";
     }
     }
 
 
-    char *server=(ser->server[0]?ser->server:ser->ip);
-    sprintf(url, "%s://%s:%d", head, server, port);
+    sprintf(url, "%s://%s:%d", head, ser->server, port);
     return 0;
     return 0;
 }
 }
 static int my_conn(mqtt_handle_t *h)
 static int my_conn(mqtt_handle_t *h)
@@ -323,6 +301,8 @@ static int my_conn(mqtt_handle_t *h)
                 .qos = 1,
                 .qos = 1,
                 .version = 4,
                 .version = 4,
                 .keepalive = 60,
                 .keepalive = 60,
+                .topic = mg_str("hello"),
+                .message = mg_str("bye"),
                 .client_id = mg_str(ser[i].cid),
                 .client_id = mg_str(ser[i].cid),
                 .user = mg_str(ser[i].user),
                 .user = mg_str(ser[i].user),
                 .pass = mg_str(ser[i].password),
                 .pass = mg_str(ser[i].password),
@@ -353,17 +333,17 @@ static int my_send(mqtt_handle_t *h)
 
 
     return r;
     return r;
 }
 }
-static void post_once(void)
+static void send_once(mqtt_handle_t *h)
 {
 {
-    post_stat(MQTT_PUB_INFO_DEVICE, 0);
-    post_stat(MQTT_PUB_STAT_NETWORK, 0);
-    post_stat(MQTT_PUB_STAT_SENSOR, 0);
-    post_stat(MQTT_PUB_STAT_SERVICE, 0);
+    send_stat(h, MQTT_PUB_INFO_DEVICE);
+    send_stat(h, MQTT_PUB_STAT_NETWORK);
+    send_stat(h, MQTT_PUB_STAT_SENSOR);
+    send_stat(h, MQTT_PUB_STAT_SERVICE);
 }
 }
-static void post_period(void)
+static void send_period(mqtt_handle_t *h)
 {
 {
-    post_stat(MQTT_PUB_STAT_POWER_ALL, 0);
-    post_stat(MQTT_PUB_STAT_POWER_CHN, 0);
+    send_stat(h, MQTT_PUB_STAT_POWER_ALL);
+    send_stat(h, MQTT_PUB_STAT_POWER_CHN);
 }
 }
 
 
 static void timer_conn_fn(void *arg)
 static void timer_conn_fn(void *arg)
@@ -372,19 +352,15 @@ static void timer_conn_fn(void *arg)
     my_conn(h);
     my_conn(h);
 }
 }
 static void timer_period_fn(void *arg)
 static void timer_period_fn(void *arg)
-{
-    post_period();
-}
-static void timer_send_fn(void *arg)
 {
 {
     mqtt_handle_t *h=(mqtt_handle_t*)arg;
     mqtt_handle_t *h=(mqtt_handle_t*)arg;
-    my_send(h);
+    send_period(h);
 }
 }
+
 static void timer_start(mqtt_handle_t *h)
 static void timer_start(mqtt_handle_t *h)
 {
 {
-    mg_timer_add(&h->mgr, 3000, MG_TIMER_REPEAT | MG_TIMER_RUN_NOW, timer_conn_fn, h);
-    mg_timer_add(&h->mgr, 5000, MG_TIMER_REPEAT | MG_TIMER_RUN_NOW, timer_period_fn, h);
-    mg_timer_add(&h->mgr, 2000, MG_TIMER_REPEAT | MG_TIMER_RUN_NOW, timer_send_fn, h);
+    mg_timer_add(&h->mgr, CONN_PERIOD, MG_TIMER_REPEAT | MG_TIMER_RUN_NOW, timer_conn_fn, h);
+    mg_timer_add(&h->mgr, SEND_PERIOD, MG_TIMER_REPEAT | MG_TIMER_RUN_NOW, timer_period_fn, h);
 }
 }
 
 
 static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
 static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
@@ -401,10 +377,9 @@ static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
 
 
         case MG_EV_CONNECT:
         case MG_EV_CONNECT:
         {
         {
-            //char path[200];
-            //sys_get_path(path, "emqxsl-ca.crt");
-            if (mg_url_is_ssl(mc->ser.server)) {
-
+            char url[1024];
+            get_url(&mc->ser, url);
+            if (mg_url_is_ssl(url)) {
                 struct mg_tls_opts opts = {.ca = mg_str(mc->ser.cert),
                 struct mg_tls_opts opts = {.ca = mg_str(mc->ser.cert),
                                            .name = mg_url_host(mc->ser.server)};
                                            .name = mg_url_host(mc->ser.server)};
                 mg_tls_init(c, &opts);
                 mg_tls_init(c, &opts);
@@ -414,10 +389,16 @@ static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
 
 
         case  MG_EV_MQTT_OPEN:
         case  MG_EV_MQTT_OPEN:
         {
         {
-            post_once();
+            send_once(h);
         }
         }
         break;
         break;
         
         
+        case MG_EV_POLL:
+        {
+            my_send(h);
+        }
+        break;
+
         case MG_EV_MQTT_MSG:
         case MG_EV_MQTT_MSG:
         {
         {
             struct mg_mqtt_message *mm=(struct mg_mqtt_message*)ev_data;
             struct mg_mqtt_message *mm=(struct mg_mqtt_message*)ev_data;
@@ -516,43 +497,44 @@ int mqtt_deinit(void)
 }
 }
 
 
 
 
-static int post_stat(int type, int id)
+static int send_stat(mqtt_handle_t *h, int type)
 {
 {
-    int r=-1;
+    int i,r=-1;
 
 
 #ifdef USE_MQTT
 #ifdef USE_MQTT
+    float p1,p2,p3;
     char topic[256];
     char topic[256];
     char* content=NULL;
     char* content=NULL;
-    mqtt_handle_t *h=&mqHandle;
-    char date[40],time[40];
+    char buf[20],date[40],time[40];
 
 
     if(!h->inited || type<0 || type>=MQTT_PUB_MAX) {
     if(!h->inited || type<0 || type>=MQTT_PUB_MAX) {
         return -1;
         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 {
+    if(type!=MQTT_PUB_STAT_POWER_CHN && type!=MQTT_PUB_STAT_SENSOR) {
         snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id);
         snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id);
     }
     }
-    get_date_time(date, time);
 
 
     switch(type) {
     switch(type) {
         case MQTT_PUB_INFO_DEVICE:
         case MQTT_PUB_INFO_DEVICE:
         {
         {
             cJSON* root=cJSON_CreateObject();
             cJSON* root=cJSON_CreateObject();
             if(root) {
             if(root) {
-                cJSON_AddStringToObject(root,"id", get_int_str(id));
+                cJSON_AddStringToObject(root,"id", get_int_str(h->prod_id));
                 cJSON_AddStringToObject(root,"name", "Gowone smartPDU");
                 cJSON_AddStringToObject(root,"name", "Gowone smartPDU");
                 cJSON_AddStringToObject(root,"type", "smartPDU AC");
                 cJSON_AddStringToObject(root,"type", "smartPDU AC");
                 cJSON_AddStringToObject(root,"firmware", VERSION);
                 cJSON_AddStringToObject(root,"firmware", VERSION);
                 cJSON_AddStringToObject(root,"channel_num", "8");
                 cJSON_AddStringToObject(root,"channel_num", "8");
                 cJSON_AddStringToObject(root,"sensor_num", "10");
                 cJSON_AddStringToObject(root,"sensor_num", "10");
+
+                get_date_time(date, time);
                 cJSON_AddStringToObject(root,"date", date);
                 cJSON_AddStringToObject(root,"date", date);
                 cJSON_AddStringToObject(root,"time", time);
                 cJSON_AddStringToObject(root,"time", time);
 
 
                 content = cJSON_Print(root);    
                 content = cJSON_Print(root);    
+                my_pub(h, type, topic, content);
+
+                cJSON_free(content);
                 cJSON_Delete(root);
                 cJSON_Delete(root);
             }
             }
         }
         }
@@ -603,10 +585,15 @@ static int post_stat(int type, int id)
                 cJSON_AddStringToObject(root,"modbus_baud", get_int_str(h->dm->_globalDevInfo._gmodbus_info.product_modbus_baud));
                 cJSON_AddStringToObject(root,"modbus_baud", get_int_str(h->dm->_globalDevInfo._gmodbus_info.product_modbus_baud));
                 cJSON_AddStringToObject(root,"modbus_mode", get_int_str(h->dm->_globalDevInfo._gmodbus_info.product_modbus_type));
                 cJSON_AddStringToObject(root,"modbus_mode", get_int_str(h->dm->_globalDevInfo._gmodbus_info.product_modbus_type));
                 cJSON_AddStringToObject(root,"vpn_enable", "0");
                 cJSON_AddStringToObject(root,"vpn_enable", "0");
+
+                get_date_time(date, time);
                 cJSON_AddStringToObject(root,"date", date);
                 cJSON_AddStringToObject(root,"date", date);
                 cJSON_AddStringToObject(root,"time", time);
                 cJSON_AddStringToObject(root,"time", time);
 
 
                 content = cJSON_Print(root);    
                 content = cJSON_Print(root);    
+                my_pub(h, type, topic, content);
+
+                cJSON_free(content);
                 cJSON_Delete(root);
                 cJSON_Delete(root);
             }
             }
         }
         }
@@ -616,18 +603,38 @@ static int post_stat(int type, int id)
         {
         {
             cJSON* root=cJSON_CreateObject();
             cJSON* root=cJSON_CreateObject();
             if(root) {
             if(root) {
-                //h->dm->_globalPowerManger
-                cJSON_AddStringToObject(root,"voltage", "");
-                cJSON_AddStringToObject(root,"current", "");
-                cJSON_AddStringToObject(root,"power", "");
-                cJSON_AddStringToObject(root,"consumption", "");
-                cJSON_AddStringToObject(root,"pactive_power", "");
-                cJSON_AddStringToObject(root,"reactive_power", "");
-                cJSON_AddStringToObject(root,"apparent_power", "");
-                cJSON_AddStringToObject(root,"date", date);
-                cJSON_AddStringToObject(root,"time", time);
+                int cnt=1;
+                pwrall_info_t *pall=websocket_get_pwrall();
 
 
-                content = cJSON_Print(root);    
+                if(pall->ph3) {
+                    cnt = 3;
+                }
+
+                for(i=0; i<cnt; i++) {
+                    _OverAllPwrAckInfo *info=&pall->ch[i];
+
+                    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(h, type, topic, content);
+                    cJSON_free(content);
+                }
                 cJSON_Delete(root);
                 cJSON_Delete(root);
             }
             }
         }
         }
@@ -635,11 +642,20 @@ static int post_stat(int type, int id)
 
 
         case MQTT_PUB_STAT_POWER_CHN:
         case MQTT_PUB_STAT_POWER_CHN:
         {
         {
-            cJSON* root=cJSON_CreateObject();
-            if(root) {
-                GlobalPowerManger *tmp=get_power(h, id);
-                if(tmp) {
-                    cJSON_AddStringToObject(root,"name", "app server channel");
+            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,"status", get_int_str(tmp->_PowerInfo.status));
                     cJSON_AddStringToObject(root,"voltage", get_float_str(tmp->_PowerInfo.voltage));
                     cJSON_AddStringToObject(root,"voltage", get_float_str(tmp->_PowerInfo.voltage));
                     cJSON_AddStringToObject(root,"current", get_float_str(tmp->_PowerInfo.current));
                     cJSON_AddStringToObject(root,"current", get_float_str(tmp->_PowerInfo.current));
@@ -647,25 +663,44 @@ static int post_stat(int type, int id)
                     cJSON_AddStringToObject(root,"power_freq", get_float_str(tmp->_PowerInfo.freq));
                     cJSON_AddStringToObject(root,"power_freq", get_float_str(tmp->_PowerInfo.freq));
                     cJSON_AddStringToObject(root,"consumption", get_float_str(tmp->_PowerInfo.consumption));
                     cJSON_AddStringToObject(root,"consumption", get_float_str(tmp->_PowerInfo.consumption));
                     cJSON_AddStringToObject(root,"power_factor", get_float_str(tmp->_PowerInfo.factor));
                     cJSON_AddStringToObject(root,"power_factor", get_float_str(tmp->_PowerInfo.factor));
-                    cJSON_AddStringToObject(root,"pactive_power", "");
-                    cJSON_AddStringToObject(root,"reactive_power", "");
-                    cJSON_AddStringToObject(root,"apparent_power", "");
-                }
-                cJSON_AddStringToObject(root,"date", date);
-                cJSON_AddStringToObject(root,"time", time);
 
 
-                content = cJSON_Print(root);    
-                cJSON_Delete(root);
+                    
+                    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(h, type, topic, content);
+
+                    cJSON_free(content);
+                    cJSON_Delete(root);
+                }
             }
             }
+            lock_s_release(LOCK_ID_POWER_UPDATE);
         }
         }
         break;
         break;
 
 
         case MQTT_PUB_STAT_SENSOR:
         case MQTT_PUB_STAT_SENSOR:
         {
         {
-            cJSON* root=cJSON_CreateObject();
-            if(root) {
-                GlobalSensorManger *tmp=get_sensor(h, id);
-                if(tmp) {
+            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,"name", "temp sensor");
                     cJSON_AddStringToObject(root,"type", get_int_str(tmp->sensor_type));
                     cJSON_AddStringToObject(root,"type", get_int_str(tmp->sensor_type));
                     cJSON_AddStringToObject(root,"modbus_address", get_int_str(tmp->sensor_addr));
                     cJSON_AddStringToObject(root,"modbus_address", get_int_str(tmp->sensor_addr));
@@ -674,13 +709,20 @@ static int post_stat(int type, int id)
                     cJSON_AddStringToObject(root,"value_num", get_int_str(tmp->sensor_val_count));
                     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,"value1", get_float_str(tmp->Cur_sensor_info.val1));
                     cJSON_AddStringToObject(root,"value2", get_float_str(tmp->Cur_sensor_info.val2));
                     cJSON_AddStringToObject(root,"value2", get_float_str(tmp->Cur_sensor_info.val2));
-                }
-                cJSON_AddStringToObject(root,"date", date);
-                cJSON_AddStringToObject(root,"time", time);
+                
+                    get_date_time(date, time);
+                    cJSON_AddStringToObject(root,"date", date);
+                    cJSON_AddStringToObject(root,"time", time);
 
 
-                content = cJSON_Print(root);    
-                cJSON_Delete(root);
+                    snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id, tmp->sensor_id);
+                    content = cJSON_Print(root);    
+                    my_pub(h, type, topic, content);
+
+                    cJSON_free(content);
+                    cJSON_Delete(root);
+                }
             }
             }
+            lock_s_release(LOCK_ID_SENSOR);
         }
         }
         break;
         break;
 
 
@@ -700,10 +742,15 @@ static int post_stat(int type, int id)
                 cJSON_AddStringToObject(root,"mqtt", BGET(flag, SERV_MQTT)?"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,"cloud", BGET(flag, SERV_CLOUD)?"1":"0");
                 cJSON_AddStringToObject(root,"ntp", BGET(flag, SERV_NTP)?"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,"date", date);
                 cJSON_AddStringToObject(root,"time", time);
                 cJSON_AddStringToObject(root,"time", time);
 
 
                 content = cJSON_Print(root);
                 content = cJSON_Print(root);
+                my_pub(h, type, topic, content);
+
+                cJSON_free(content);
                 cJSON_Delete(root);
                 cJSON_Delete(root);
             }
             }
         }
         }
@@ -713,12 +760,8 @@ static int post_stat(int type, int id)
         return -1;
         return -1;
     }
     }
 
 
-    mqtt_pkt_t pkt;
-    pkt.id = type;
-    snprintf(pkt.topic, sizeof(pkt.topic), "%s", topic);
-    snprintf(pkt.content, sizeof(pkt.content), "%s", content);
-    r = xlist_append(h->list, 0, &pkt, sizeof(pkt));
-    cJSON_free(content);
+    
+    
 #endif
 #endif
 
 
     return r;
     return r;
@@ -888,7 +931,7 @@ static int set_serv(uint32_t flag)
 {
 {
     if(BGET(flag,SERV_NTP)) {
     if(BGET(flag,SERV_NTP)) {
         if(!thread_is_running(THREAD_ID_NTP)) {
         if(!thread_is_running(THREAD_ID_NTP)) {
-            
+
         }
         }
     }
     }
     else {
     else {
@@ -1057,9 +1100,8 @@ int mqtt_test(void)
             .id = 0,
             .id = 0,
             .mode = 1,
             .mode = 1,
             .cid = "mqtlx_e92334545tet3465",
             .cid = "mqtlx_e92334545tet3465",
-
+            .name = "serverX",
             .server = "192.168.1.12",
             .server = "192.168.1.12",
-            .ip = "",
             .port = "1883",
             .port = "1883",
 
 
             //root
             //root
@@ -1076,7 +1118,7 @@ int mqtt_test(void)
 
 
     my_conn(h);
     my_conn(h);
     while(1) {
     while(1) {
-        post_period();
+        send_period(h);
         sleep(3);
         sleep(3);
     }
     }
 
 

+ 8 - 8
pro/src/sqlite_handle.c

@@ -4365,8 +4365,8 @@ int dev_mqtt_init(sqlite3 *db, mqtt_info_t *info)
         sprintf(temp, "CREATE TABLE IF NOT EXISTS %s("
         sprintf(temp, "CREATE TABLE IF NOT EXISTS %s("
                             "id INTEGER PRIMARY KEY,"
                             "id INTEGER PRIMARY KEY,"
                             "mode INTEGER,"
                             "mode INTEGER,"
+                            "name TEXT,"
                             "server TEXT,"
                             "server TEXT,"
-                            "ip TEXT,"
                             "port TEXT,"
                             "port TEXT,"
                             "cid TEXT,"
                             "cid TEXT,"
                             "user TEXT,"
                             "user TEXT,"
@@ -4399,10 +4399,10 @@ int dev_mqtt_init(sqlite3 *db, mqtt_info_t *info)
         ser[idx].mode = sqlite3_column_int(stmt, 1);
         ser[idx].mode = sqlite3_column_int(stmt, 1);
 
 
         p = (char*)sqlite3_column_text(stmt, 2);
         p = (char*)sqlite3_column_text(stmt, 2);
-        if(p) strcpy(ser[idx].server, p);
+        if(p) strcpy(ser[idx].name, p);
 
 
         p = (char*)sqlite3_column_text(stmt, 3);
         p = (char*)sqlite3_column_text(stmt, 3);
-        if(p) strcpy(ser[idx].ip, p);
+        if(p) strcpy(ser[idx].server, p);
 
 
         p = (char*)sqlite3_column_text(stmt, 4);
         p = (char*)sqlite3_column_text(stmt, 4);
         if(p) strcpy(ser[idx].port, p);
         if(p) strcpy(ser[idx].port, p);
@@ -4452,8 +4452,8 @@ int dev_mqtt_update_ser(sqlite3 *db, mqtt_server_t *ser)
 
 
     cnt = get_count(db, table);
     cnt = get_count(db, table);
     if(cnt>0) {
     if(cnt>0) {
-        sprintf(pbuf, "UPDATE %s SET id=%d, mode=%d, server='%s', ip='%s', port='%s', cid='%s', user='%s', password='%s, cert='%s', WHERE id=%d;", 
-                table, ser->id, ser->mode, ser->server, ser->ip, ser->port, ser->cid, ser->user, ser->password, ser->cert, ser->id);
+        sprintf(pbuf, "UPDATE %s SET id=%d, mode=%d, name='%s', server='%s', port='%s', cid='%s', user='%s', password='%s, cert='%s', WHERE id=%d;", 
+                table, ser->id, ser->mode, ser->name, ser->server, ser->port, ser->cid, ser->user, ser->password, ser->cert?ser->cert:"", ser->id);
         r = sqlite3_exec(db, pbuf, NULL, 0, &errmsg);
         r = sqlite3_exec(db, pbuf, NULL, 0, &errmsg);
         if(r != SQLITE_OK) {
         if(r != SQLITE_OK) {
             log_e("___dev_mqtt_update_ser, sqlite3_exec failed, %s, %s\n\n", pbuf, errmsg);
             log_e("___dev_mqtt_update_ser, sqlite3_exec failed, %s, %s\n\n", pbuf, errmsg);
@@ -4465,7 +4465,7 @@ int dev_mqtt_update_ser(sqlite3 *db, mqtt_server_t *ser)
     }
     }
 
 
     if(exist==0) {
     if(exist==0) {
-        sprintf(pbuf, "INSERT INTO %s (id,mode,server,ip,port,cid,user,password,cert) VALUES (?,?,?,?,?,?,?,?,?);", table);
+        sprintf(pbuf, "INSERT INTO %s (id,mode,name,server,port,cid,user,password,cert) VALUES (?,?,?,?,?,?,?,?,?);", table);
         r = sqlite3_prepare_v2(db, pbuf, -1, &stmt, 0);
         r = sqlite3_prepare_v2(db, pbuf, -1, &stmt, 0);
         if (r != SQLITE_OK) {
         if (r != SQLITE_OK) {
             log_e("___dev_mqtt_update_ser, sqlite3_prepare_v2 failed, %s\n", sqlite3_errmsg(db));
             log_e("___dev_mqtt_update_ser, sqlite3_prepare_v2 failed, %s\n", sqlite3_errmsg(db));
@@ -4474,8 +4474,8 @@ int dev_mqtt_update_ser(sqlite3 *db, mqtt_server_t *ser)
 
 
         sqlite3_bind_int(stmt,  1, ser->id);
         sqlite3_bind_int(stmt,  1, ser->id);
         sqlite3_bind_int(stmt,  2, ser->mode);
         sqlite3_bind_int(stmt,  2, ser->mode);
-        sqlite3_bind_text(stmt, 3, ser->server,   -1, NULL);
-        sqlite3_bind_text(stmt, 4, ser->ip,       -1, NULL);
+        sqlite3_bind_text(stmt, 3, ser->name,     -1, NULL);
+        sqlite3_bind_text(stmt, 4, ser->server,   -1, NULL);
         sqlite3_bind_text(stmt, 5, ser->port,     -1, NULL);
         sqlite3_bind_text(stmt, 5, ser->port,     -1, NULL);
         sqlite3_bind_text(stmt, 6, ser->cid,      -1, NULL);
         sqlite3_bind_text(stmt, 6, ser->cid,      -1, NULL);
         sqlite3_bind_text(stmt, 7, ser->user,     -1, NULL);
         sqlite3_bind_text(stmt, 7, ser->user,     -1, NULL);

+ 30 - 25
pro/src/websocket_handle.c

@@ -31,7 +31,9 @@ typedef struct {
     mg_conn_t *c;
     mg_conn_t *c;
     int       inited;
     int       inited;
     int       ws_cnt;
     int       ws_cnt;
+    pwrall_info_t pwrall;
 }ws_handle_t;
 }ws_handle_t;
+
 static ws_handle_t wsHandle;
 static ws_handle_t wsHandle;
 
 
 static int ws_is_online(ws_handle_t *wh) 
 static int ws_is_online(ws_handle_t *wh) 
@@ -86,6 +88,7 @@ static void read_ipaddr(ws_handle_t *h)
   sprintf(h->wpath,"ws://%s:6785", getLocalIpAddress("eth0"));
   sprintf(h->wpath,"ws://%s:6785", getLocalIpAddress("eth0"));
   printf("____ ws path: %s\n", h->wpath);
   printf("____ ws path: %s\n", h->wpath);
 }
 }
+
 static void web_update(ws_handle_t *wh)
 static void web_update(ws_handle_t *wh)
 {
 {
     GlobalSensorManger* _globalSensorMangerTemp;
     GlobalSensorManger* _globalSensorMangerTemp;
@@ -97,63 +100,62 @@ static void web_update(ws_handle_t *wh)
     char* json_str = NULL ;
     char* json_str = NULL ;
      GlobalPowerInfo _powerInfo;
      GlobalPowerInfo _powerInfo;
     GlobalPowerManger* _globalPowerMangerTemp = NULL ;
     GlobalPowerManger* _globalPowerMangerTemp = NULL ;
-    _OverAllPwrAckInfo allInfo; //设备电源信息
     _OverChnPwrAckInfo chInfo;
     _OverChnPwrAckInfo chInfo;
     GlobalDeviceManager* gdm= &__globalDeviceManage;
     GlobalDeviceManager* gdm= &__globalDeviceManage;
     GlobalDeviceManager* gdm2=&__globalDeviceManage2;
     GlobalDeviceManager* gdm2=&__globalDeviceManage2;
     GlobalDeviceManager* pdm=NULL;
     GlobalDeviceManager* pdm=NULL;
+    int r,r1,r2,r3;
 
 
     if(ws_is_online(wh))
     if(ws_is_online(wh))
     {
     {
         //overall info
         //overall info
+        pwrall_info_t *pall=&wh->pwrall;
         if(gdm->_globalDevInfo.product_pwr_type==SmartPDU_Tree_AC_Tree ||
         if(gdm->_globalDevInfo.product_pwr_type==SmartPDU_Tree_AC_Tree ||
             gdm->_globalDevInfo.product_pwr_type==SmartPDU_Tree_AC_One ||
             gdm->_globalDevInfo.product_pwr_type==SmartPDU_Tree_AC_One ||
 			gdm->_globalDevInfo.product_pwr_type==SmartPDU_Tree_AC_One_B)
 			gdm->_globalDevInfo.product_pwr_type==SmartPDU_Tree_AC_One_B)
 		{
 		{
-            int nRet_L1 = -1;
-            int nRet_L2 = -1;
-            int nRet_L3 = -1;
-            _OverAllPwrAckInfo l1, l2, l3;
+            pall->ph3 = 1;
+
             if (cur_dev_addr == 0)
             if (cur_dev_addr == 0)
             {
             {
                 lock_s_hold(LOCK_ID_POWER_UPDATE);
                 lock_s_hold(LOCK_ID_POWER_UPDATE);
-                nRet_L1 = dev_search_latest_t_ac_power_statistic_info(0, 0, &l1, gdm);
-                nRet_L2 = dev_search_latest_t_ac_power_statistic_info(0, 1, &l2, gdm);
-                nRet_L3 = dev_search_latest_t_ac_power_statistic_info(0, 2, &l3, gdm);
+                r1 = dev_search_latest_t_ac_power_statistic_info(0, 0, &pall->ch[0], gdm);
+                r2 = dev_search_latest_t_ac_power_statistic_info(0, 1, &pall->ch[1], gdm);
+                r3 = dev_search_latest_t_ac_power_statistic_info(0, 2, &pall->ch[2], gdm);
                 lock_s_release(LOCK_ID_POWER_UPDATE);
                 lock_s_release(LOCK_ID_POWER_UPDATE);
             }
             }
             else
             else
             {
             {
                 cascade_lock();
                 cascade_lock();
-                nRet_L1 = dev_search_latest_t_ac_power_statistic_info(0, 0, &l1, gdm2);
-                nRet_L2 = dev_search_latest_t_ac_power_statistic_info(0, 1, &l2, gdm2);
-                nRet_L3 = dev_search_latest_t_ac_power_statistic_info(0, 2, &l3, gdm2);
+                r1 = dev_search_latest_t_ac_power_statistic_info(0, 0, &pall->ch[0], gdm2);
+                r2 = dev_search_latest_t_ac_power_statistic_info(0, 1, &pall->ch[1], gdm2);
+                r3 = dev_search_latest_t_ac_power_statistic_info(0, 2, &pall->ch[2], gdm2);
                 cascade_unlock();
                 cascade_unlock();
             }
             }
-            if (0 == nRet_L1 && 0 == nRet_L2 && 0 == nRet_L3)
+            if (0 == r1 && 0 == r2 && 0 == r3)
             {
             {
-                json_str = over_all_pwr_monitor_Tree_AC_to_json(&l1, &l2, &l3);
-
+                json_str = over_all_pwr_monitor_Tree_AC_to_json(&pall->ch[0], &pall->ch[1], &pall->ch[2]);
                 ws_send_data(wh, json_str, strlen(json_str));
                 ws_send_data(wh, json_str, strlen(json_str));
                 cJSON_free((void *)json_str);
                 cJSON_free((void *)json_str);
             }
             }
         }
         }
         else
         else
         {
         {
-            int nRetTotal = -1;
+            pall->ph3 = 0;
+
             if(cur_dev_addr==0) {
             if(cur_dev_addr==0) {
                 lock_s_hold(LOCK_ID_POWER_UPDATE);
                 lock_s_hold(LOCK_ID_POWER_UPDATE);
-                nRetTotal=dev_search_latest_power_statistic_info(gdm, &allInfo);
+                r = dev_search_latest_power_statistic_info(gdm, &pall->ch[0]);
                 lock_s_release(LOCK_ID_POWER_UPDATE);
                 lock_s_release(LOCK_ID_POWER_UPDATE);
             }
             }
             else {
             else {
                 cascade_lock();
                 cascade_lock();
-                nRetTotal=dev_search_latest_power_statistic_info(gdm2, &allInfo);
+                r = dev_search_latest_power_statistic_info(gdm2, &pall->ch[0]);
                 cascade_unlock();
                 cascade_unlock();
             }
             }
-            if (0==nRetTotal)
+            if (0==r)
             {
             {
-                json_str = ws_over_status_ack_to_json(0, &allInfo);
+                json_str = ws_over_status_ack_to_json(0, &pall->ch[0]);
                 ws_send_data(wh, json_str, strlen(json_str));
                 ws_send_data(wh, json_str, strlen(json_str));
                 cJSON_free(json_str);
                 cJSON_free(json_str);
             }
             }
@@ -162,21 +164,20 @@ static void web_update(ws_handle_t *wh)
         //channel info
         //channel info
         {
         {
             INIT_LIST_HEAD(&chInfo.list);
             INIT_LIST_HEAD(&chInfo.list);
-             int nRetCh = -1;
 
 
             if(cur_dev_addr==0) {
             if(cur_dev_addr==0) {
                 lock_s_hold(LOCK_ID_POWER_UPDATE);
                 lock_s_hold(LOCK_ID_POWER_UPDATE);
-                nRetCh=dev_search_latest_power_All_info(gdm,&chInfo, cur_dev_addr);
+                r = dev_search_latest_power_All_info(gdm,&chInfo, cur_dev_addr);
                 lock_s_release(LOCK_ID_POWER_UPDATE);
                 lock_s_release(LOCK_ID_POWER_UPDATE);
                 pdm = gdm;
                 pdm = gdm;
             }
             }
             else {
             else {
                 cascade_lock();
                 cascade_lock();
-                nRetCh=dev_search_latest_power_All_info(gdm2,&chInfo, cur_dev_addr);
+                r = dev_search_latest_power_All_info(gdm2,&chInfo, cur_dev_addr);
                 cascade_unlock();
                 cascade_unlock();
                 pdm = gdm2;
                 pdm = gdm2;
             }
             }
-            if (0 == nRetCh)
+            if (0 == r)
             {
             {
                 json_str = ws_chn_status_ack_to_json(1, &chInfo, pdm);
                 json_str = ws_chn_status_ack_to_json(1, &chInfo, pdm);
                 ws_send_data(wh, json_str, strlen(json_str));
                 ws_send_data(wh, json_str, strlen(json_str));
@@ -199,7 +200,8 @@ static void web_update(ws_handle_t *wh)
             ws_send_data(wh, json_str, strlen(json_str));
             ws_send_data(wh, json_str, strlen(json_str));
             cJSON_free(json_str);
             cJSON_free(json_str);
         }
         }
-		        // // breaker info
+		
+        // breaker info
         {
         {
             if(cur_dev_addr==0) { 
             if(cur_dev_addr==0) { 
                 json_str = breaker_status_ack_to_json(gdm,0);
                 json_str = breaker_status_ack_to_json(gdm,0);
@@ -308,4 +310,7 @@ int websocket_init(void)
 
 
 
 
 
 
-  
+pwrall_info_t *websocket_get_pwrall(void)
+{
+    return &wsHandle.pwrall;
+}

+ 8 - 0
pro/src/websocket_handle.h

@@ -1,6 +1,8 @@
 #ifndef __WEBSOCKET_HANDLE_H
 #ifndef __WEBSOCKET_HANDLE_H
 #define __WEBSOCKET_HANDLE_H
 #define __WEBSOCKET_HANDLE_H
 
 
+#include "json_handle.h"
+
 #define WS_EVT_DATA   0x11
 #define WS_EVT_DATA   0x11
 
 
 typedef struct {
 typedef struct {
@@ -9,7 +11,13 @@ typedef struct {
     int     len;
     int     len;
 }ws_data_t;
 }ws_data_t;
 
 
+typedef struct {
+    int ph3;
+    _OverAllPwrAckInfo ch[3];
+}pwrall_info_t;
+
 int websocket_init(void);
 int websocket_init(void);
+pwrall_info_t *websocket_get_pwrall(void);
 
 
 #endif
 #endif