Bladeren bron

1、修改了MQTT协议topic定义,将所有通道、传感器合并为一个包,降低网络延迟;2、调整电源信息、传感器、网络状态、设备信息发送频率,不同状态类型发送频率不一致

liyuezong 1 jaar geleden
bovenliggende
commit
d59114053f
1 gewijzigde bestanden met toevoegingen van 144 en 118 verwijderingen
  1. 144 118
      pro/src/mqtt.c

+ 144 - 118
pro/src/mqtt.c

@@ -34,8 +34,8 @@
 
 
 #define CONN_PERIOD         5000    //ms
-#define SEND_PERIOD         3000    //ms
-
+#define SEND_PERIOD         1000    //ms
+#define SEND_PERIOD2        5000    //ms
 
 #define BGET(flag,mask)  ((flag)&(1<<(mask)))
 #define BSET(flag,mask)  ((flag)|=(1<<(mask)))
@@ -83,9 +83,9 @@ char *topic_sub[MQTT_SUB_MAX]={
 char *topic_pub[MQTT_PUB_MAX] = {
     "/pdu/%s/info/device",
     "/pdu/%s/status/network",
-    "/pdu/%s/status/power/all_channel",
-    "/pdu/%s/status/power/%d",
-    "/pdu/%s/status/sensor/%d",
+    "/pdu/%s/status/power/all",
+    "/pdu/%s/status/power/channel",
+    "/pdu/%s/status/sensor",
     "/pdu/%s/status/service",
     "/pdu/%s/alarm/network",
     "/pdu/%s/alarm/power",
@@ -221,7 +221,7 @@ static int pub_one(mg_conn_t *c, char *topic, char *data)
     
     opts.topic = mg_str(topic);
     opts.message = mg_str(data);
-    opts.qos = 1;
+    opts.qos = 0;
     opts.retain = false;
     mg_mqtt_pub(c, &opts);
 
@@ -372,7 +372,7 @@ static int conn_one(mqtt_conn_t *conn)
 
         mg_opts_t opts={
             .clean = true,
-            .qos = 1,
+            .qos = 0,
             .version = 4,
             .keepalive = 30,
             .topic = mg_str("hello"),
@@ -448,11 +448,16 @@ static void send_period(mqtt_conn_t *conn)
     send_stat(conn, MQTT_PUB_STAT_POWER_ALL);
     send_stat(conn, MQTT_PUB_STAT_POWER_CHN);
     send_stat(conn, MQTT_PUB_STAT_SENSOR);
+}
+
+static void send_period2(mqtt_conn_t *conn)
+{
     send_stat(conn, MQTT_PUB_STAT_SERVICE);
     send_stat(conn, MQTT_PUB_INFO_DEVICE);
     send_stat(conn, MQTT_PUB_STAT_NETWORK);
 }
 
+
 static void timer_conn_fn(void *arg)
 {
     mqtt_conn_t *conn=(mqtt_conn_t*)arg;
@@ -463,10 +468,16 @@ static void timer_period_fn(void *arg)
     mqtt_conn_t *conn=(mqtt_conn_t*)arg;
     send_period(conn);
 }
+static void timer_period2_fn(void *arg)
+{
+    mqtt_conn_t *conn=(mqtt_conn_t*)arg;
+    send_period2(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);
+    mg_timer_add(&conn->mgr, SEND_PERIOD2, MG_TIMER_REPEAT | MG_TIMER_RUN_NOW, timer_period2_fn, conn);
 }
 
 static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
@@ -633,9 +644,7 @@ static int send_stat(mqtt_conn_t *conn, int type)
         return -1;
     }
 
-    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);
 
     switch(type) {
         case MQTT_PUB_INFO_DEVICE:
@@ -803,126 +812,143 @@ static int send_stat(mqtt_conn_t *conn, int type)
             int quit=0;
             GlobalPowerManger *tmp=NULL;
 
-            lock_s_hold(LOCK_ID_POWER_UPDATE);
-            list_for_each_entry(tmp, &h->dm->_globalPowerManger.list, list)
-            {
-                if(tmp==NULL) {
-                    break;
-                }
-
-                cJSON* root=cJSON_CreateObject();
-                if(root) {
-                     cJSON_AddStringToObject(root,"id", get_int_str(tmp->product_ch_id));
-                    cJSON_AddStringToObject(root,"name", tmp->product_ch_name);
-                    cJSON_AddStringToObject(root,"status", get_int_str(tmp->_PowerInfo.status));
-                    cJSON_AddStringToObject(root,"voltage", get_float_str(tmp->_PowerInfo.voltage));
-                    cJSON_AddStringToObject(root, "current", get_float_str(tmp->_PowerInfo.current));
-                    cJSON_AddStringToObject(root, "power", get_float_str(tmp->_PowerInfo.power));
-                    cJSON_AddStringToObject(root, "consumption", get_float_str(tmp->_PowerInfo.consumption));
-
-                    cJSON_AddStringToObject(root, "voltage_over", get_int_str(tmp->_PowerWarninginfo.w_voltage_up));
-                    cJSON_AddStringToObject(root, "voltage_low", get_int_str(tmp->_PowerWarninginfo.w_voltage_down));
-                    cJSON_AddStringToObject(root, "current_over", get_int_str(tmp->_PowerWarninginfo.w_current));
-                    cJSON_AddStringToObject(root, "power_over", get_int_str(tmp->_PowerWarninginfo.w_power));
-                    cJSON_AddStringToObject(root, "consumption_over", get_int_str(tmp->_PowerWarninginfo.w_consumption));
-
-                    if (h->dm->_globalDevInfo.product.pwr_type == SmartPDU_Tree_AC_Tree)
-                    {
-                        cJSON_AddStringToObject(root, "phase_loss_L1", get_int_str(tmp->_PowerWarninginfo.w_phase_lossA));
-                        cJSON_AddStringToObject(root, "phase_loss_L2", get_int_str(tmp->_PowerWarninginfo.w_phase_lossB));
-                        cJSON_AddStringToObject(root, "phase_loss_L3", get_int_str(tmp->_PowerWarninginfo.w_phase_lossC));
-                    }
-
-                    if (h->dm->_globalDevInfo.product.pwr_type != SmartPDU_DC)
-                    {
-                        cJSON_AddStringToObject(root, "power_freq", get_float_str(tmp->_PowerInfo.freq));
-                        cJSON_AddStringToObject(root, "power_factor", get_float_str(tmp->_PowerInfo.factor));
-                        p1 = tmp->_PowerInfo.power;
-                        p3 = tmp->_PowerInfo.factor>0 ? (p1 / tmp->_PowerInfo.factor) : 0.0f;
-                        p2 = p3 - p1;
-                        cJSON_AddStringToObject(root, "pactive_power", get_float_str(p1));
-                        cJSON_AddStringToObject(root, "reactive_power", get_float_str(p2));
-                        cJSON_AddStringToObject(root, "apparent_power", get_float_str(p3));
-                    }
-
-                     cJSON* data_root_array = cJSON_CreateArray();
-                    if (h->dm->_globalDevInfo.product.pwr_type == TREE_AC_TYPE|| h->dm->_globalDevInfo.product.pwr_type==DOUBLE_AC_TYPE|| h->dm->_globalDevInfo.product.pwr_type==SmartPDU_Tree_AC_One_B)
-                    {
-                        GlobalTreeACManager *_TreeACTemp = NULL;
-                        // 判断是否三相单输出情况下
-                        int phnum = 0;
-                        list_for_each_entry(_TreeACTemp, &tmp->list_Tree_AC, list_Tree_AC)
-                        {
-                            if(_TreeACTemp->product_ph_outputStatus == 1 &&_TreeACTemp->product_ph_type >= 0 && _TreeACTemp->product_ph_type < 3)
-                            {
-                                cJSON *data_filed = cJSON_CreateObject();
-                                cJSON_AddStringToObject(data_filed, "phase_type", g_product_t_ac_type_str[_TreeACTemp->product_ph_type]);
-                                cJSON_AddStringToObject(data_filed, "phase_voltage", get_float_str(_TreeACTemp->_PowerInfo.voltage));
-                                cJSON_AddStringToObject(data_filed, "phase_current", get_float_str(_TreeACTemp->_PowerInfo.current));
-                                cJSON_AddStringToObject(data_filed, "phase_power", get_float_str(_TreeACTemp->_PowerInfo.power));
-                                cJSON_AddStringToObject(data_filed, "phase_consumption", get_float_str(_TreeACTemp->_PowerInfo.consumption));
-                                cJSON_AddItemToObject(data_root_array, "phase_data", data_filed);
-                            }
-                        }
-                    }
-                     cJSON_AddItemToObject(root, "phase_data", data_root_array);
-
-                    get_date_time(date, time);
-                    cJSON_AddStringToObject(root,"date", date);
-                    cJSON_AddStringToObject(root,"time", time);
-
-                    snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id, tmp->product_ch_id);
-                    content = cJSON_Print(root);    
-                    my_pub(conn, topic, content);
-
-                    cJSON_free(content);
-                    cJSON_Delete(root);
-                }
-            }
+             lock_s_hold(LOCK_ID_POWER_UPDATE);       
+             cJSON *root = cJSON_CreateObject();
+             cJSON* sub_data_root_array = cJSON_CreateArray();
+             if(root && sub_data_root_array) 
+             {
+                 list_for_each_entry(tmp, &h->dm->_globalPowerManger.list, list)
+                 {
+                     if (tmp == NULL)
+                     {
+                         break;
+                     }
+                    cJSON* sub_data_filed = cJSON_CreateObject();
+
+                    if(sub_data_filed)
+                     {
+                         cJSON_AddStringToObject(sub_data_filed, "id", get_int_str(tmp->product_ch_id));
+                         cJSON_AddStringToObject(sub_data_filed, "name", tmp->product_ch_name);
+                         cJSON_AddStringToObject(sub_data_filed, "status", get_int_str(tmp->_PowerInfo.status));
+                         cJSON_AddStringToObject(sub_data_filed, "voltage", get_float_str(tmp->_PowerInfo.voltage));
+                         cJSON_AddStringToObject(sub_data_filed, "current", get_float_str(tmp->_PowerInfo.current));
+                         cJSON_AddStringToObject(sub_data_filed, "power", get_float_str(tmp->_PowerInfo.power));
+                         cJSON_AddStringToObject(sub_data_filed, "consumption", get_float_str(tmp->_PowerInfo.consumption));
+
+                         cJSON_AddStringToObject(sub_data_filed, "voltage_over", get_int_str(tmp->_PowerWarninginfo.w_voltage_up));
+                         cJSON_AddStringToObject(sub_data_filed, "voltage_low", get_int_str(tmp->_PowerWarninginfo.w_voltage_down));
+                         cJSON_AddStringToObject(sub_data_filed, "current_over", get_int_str(tmp->_PowerWarninginfo.w_current));
+                         cJSON_AddStringToObject(sub_data_filed, "power_over", get_int_str(tmp->_PowerWarninginfo.w_power));
+                         cJSON_AddStringToObject(sub_data_filed, "consumption_over", get_int_str(tmp->_PowerWarninginfo.w_consumption));
+
+                         if (h->dm->_globalDevInfo.product.pwr_type == SmartPDU_Tree_AC_Tree)
+                         {
+                             cJSON_AddStringToObject(sub_data_filed, "phase_loss_L1", get_int_str(tmp->_PowerWarninginfo.w_phase_lossA));
+                             cJSON_AddStringToObject(sub_data_filed, "phase_loss_L2", get_int_str(tmp->_PowerWarninginfo.w_phase_lossB));
+                             cJSON_AddStringToObject(sub_data_filed, "phase_loss_L3", get_int_str(tmp->_PowerWarninginfo.w_phase_lossC));
+                         }
+
+                         if (h->dm->_globalDevInfo.product.pwr_type != SmartPDU_DC)
+                         {
+                             cJSON_AddStringToObject(sub_data_filed, "power_freq", get_float_str(tmp->_PowerInfo.freq));
+                             cJSON_AddStringToObject(sub_data_filed, "power_factor", get_float_str(tmp->_PowerInfo.factor));
+                             p1 = tmp->_PowerInfo.power;
+                             p3 = tmp->_PowerInfo.factor > 0 ? (p1 / tmp->_PowerInfo.factor) : 0.0f;
+                             p2 = p3 - p1;
+                             cJSON_AddStringToObject(sub_data_filed, "pactive_power", get_float_str(p1));
+                             cJSON_AddStringToObject(sub_data_filed, "reactive_power", get_float_str(p2));
+                             cJSON_AddStringToObject(sub_data_filed, "apparent_power", get_float_str(p3));
+                         }
+
+                         if (h->dm->_globalDevInfo.product.pwr_type == TREE_AC_TYPE || h->dm->_globalDevInfo.product.pwr_type == DOUBLE_AC_TYPE || h->dm->_globalDevInfo.product.pwr_type == SmartPDU_Tree_AC_One_B)
+                         {
+                             cJSON *data_root_array = cJSON_CreateArray();
+                             GlobalTreeACManager *_TreeACTemp = NULL;
+                             // 判断是否三相单输出情况下
+                             int phnum = 0;
+                             list_for_each_entry(_TreeACTemp, &tmp->list_Tree_AC, list_Tree_AC)
+                             {
+                                 if (_TreeACTemp->product_ph_outputStatus == 1 && _TreeACTemp->product_ph_type >= 0 && _TreeACTemp->product_ph_type < 3)
+                                 {
+                                     cJSON *data_filed = cJSON_CreateObject();
+                                     cJSON_AddStringToObject(data_filed, "phase_type", g_product_t_ac_type_str[_TreeACTemp->product_ph_type]);
+                                     cJSON_AddStringToObject(data_filed, "phase_voltage", get_float_str(_TreeACTemp->_PowerInfo.voltage));
+                                     cJSON_AddStringToObject(data_filed, "phase_current", get_float_str(_TreeACTemp->_PowerInfo.current));
+                                     cJSON_AddStringToObject(data_filed, "phase_power", get_float_str(_TreeACTemp->_PowerInfo.power));
+                                     cJSON_AddStringToObject(data_filed, "phase_consumption", get_float_str(_TreeACTemp->_PowerInfo.consumption));
+                                     cJSON_AddItemToObject(data_root_array, "phase_data", data_filed);
+                                 }
+                             }
+                             cJSON_AddItemToObject(sub_data_filed, "phase_data", data_root_array);
+                         }
+                     }
+                     cJSON_AddItemToObject(sub_data_root_array, "data", sub_data_filed);
+                 }
+                 cJSON_AddItemToObject(root, "data", sub_data_root_array);
+
+                 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);
+             }
             lock_s_release(LOCK_ID_POWER_UPDATE);
         }
         break;
 
         case MQTT_PUB_STAT_SENSOR:
         {
-            GlobalSensorManger *tmp=NULL;
+            GlobalSensorManger *tmp = NULL;
 
             lock_s_hold(LOCK_ID_SENSOR);
-            list_for_each_entry(tmp, &h->dm->_globalSensorManger.list, list)
+            cJSON *root = cJSON_CreateObject();
+            cJSON *sub_data_root_array = cJSON_CreateArray();
+            if (root && sub_data_root_array)
             {
-                if(tmp==NULL) {
-                    break;
-                }
+                list_for_each_entry(tmp, &h->dm->_globalSensorManger.list, list)
+                {
+                    if (tmp == NULL)
+                    {
+                        break;
+                    }
 
-                cJSON* root=cJSON_CreateObject();
-                if(root) {
-                    cJSON_AddStringToObject(root,"name", tmp->sensor_name);
-                    cJSON_AddStringToObject(root,"type", get_int_str(tmp->sensor_type));
-                    cJSON_AddStringToObject(root,"modbus_address", get_int_str(tmp->sensor_addr));
-                    cJSON_AddStringToObject(root,"node", get_int_str(tmp->sensor_node_number));
-                    cJSON_AddStringToObject(root,"status", get_int_str(tmp->sensor_status));
-                    cJSON_AddStringToObject(root,"value_num", get_int_str(tmp->sensor_val_count));
-                    cJSON_AddStringToObject(root,"value1", get_float_str(tmp->Cur_sensor_info.val1));
-                    cJSON_AddStringToObject(root,"value2", get_float_str(tmp->Cur_sensor_info.val2));
-
-                    cJSON_AddStringToObject(root, "value1_over", get_int_str(tmp->warning_status.sensor_val1_upper));
-                    cJSON_AddStringToObject(root, "value1_low", get_int_str(tmp->warning_status.sensor_val1_lower));
-                    cJSON_AddStringToObject(root, "value2_over", get_int_str(tmp->warning_status.sensor_val2_upper));
-                    cJSON_AddStringToObject(root, "value2_low", get_int_str(tmp->warning_status.sensor_val2_lower));
-
-                    get_date_time(date, time);
-                    cJSON_AddStringToObject(root,"date", date);
-                    cJSON_AddStringToObject(root,"time", time);
-
-                    snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id, tmp->sensor_id);
-                    content = cJSON_Print(root);    
-                    my_pub(conn, topic, content);
-
-                    cJSON_free(content);
-                    cJSON_Delete(root);
+                    cJSON *sub_data_filed = cJSON_CreateObject();
+                    if (sub_data_filed)
+                    {
+                        cJSON_AddStringToObject(sub_data_filed, "name", tmp->sensor_name);
+                        cJSON_AddStringToObject(sub_data_filed, "type", get_int_str(tmp->sensor_type));
+                        cJSON_AddStringToObject(sub_data_filed, "modbus_address", get_int_str(tmp->sensor_addr));
+                        cJSON_AddStringToObject(sub_data_filed, "node", get_int_str(tmp->sensor_node_number));
+                        cJSON_AddStringToObject(sub_data_filed, "status", get_int_str(tmp->sensor_status));
+                        cJSON_AddStringToObject(sub_data_filed, "value_num", get_int_str(tmp->sensor_val_count));
+                        cJSON_AddStringToObject(sub_data_filed, "value1", get_float_str(tmp->Cur_sensor_info.val1));
+                        cJSON_AddStringToObject(sub_data_filed, "value2", get_float_str(tmp->Cur_sensor_info.val2));
+
+                        cJSON_AddStringToObject(sub_data_filed, "value1_over", get_int_str(tmp->warning_status.sensor_val1_upper));
+                        cJSON_AddStringToObject(sub_data_filed, "value1_low", get_int_str(tmp->warning_status.sensor_val1_lower));
+                        cJSON_AddStringToObject(sub_data_filed, "value2_over", get_int_str(tmp->warning_status.sensor_val2_upper));
+                        cJSON_AddStringToObject(sub_data_filed, "value2_low", get_int_str(tmp->warning_status.sensor_val2_lower));
+
+                         cJSON_AddItemToObject(sub_data_root_array, "data", sub_data_filed);
+                    }
                 }
+                cJSON_AddItemToObject(root, "data", sub_data_root_array);
+
+                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);
+                lock_s_release(LOCK_ID_SENSOR);
             }
-            lock_s_release(LOCK_ID_SENSOR);
         }
         break;