Ver código fonte

1、更新mongoose到最新版
2、增加全局静态锁,简化调用
3、添加mqtt,已打通,还需完善

guohui 2 anos atrás
pai
commit
ba9953719e

BIN
cust_data/db/ac_i3_o3_russian.db


+ 9 - 2
pro/config/cfg.h

@@ -20,7 +20,7 @@
 #define CHANNEL_DELAY_SEC               8
 
 #define USE_SMTP
-
+#define SMTP_SEND_BUILTIN
 
 /////////////customer demannd//////////////////////
 //#define CUST_GOM037     //xian
@@ -28,6 +28,13 @@
 #define CUST_ANDERSON
 //#define CUST_HYPERTEC
 //#define CUST_ITK
+//#define CUST_TEST
+
+
+#ifdef CUST_TEST
+    #undef MAX_CHN_COUNT
+    #define MAX_CHN_COUNT               24
+#endif
 
 
 #ifdef CUST_GOM037      //xian
@@ -54,7 +61,7 @@
 
 
 #ifdef CUST_ANDERSON
-    #define USE_MQTT
+    //#define USE_MQTT
     #define USE_NETSWITCH
 
     #define SW_STATIC_PORT_START

Diferenças do arquivo suprimidas por serem muito extensas
+ 1804 - 457
pro/mongoose/mongoose.c


+ 142 - 12
pro/mongoose/mongoose.h

@@ -186,7 +186,7 @@ extern "C" {
 #define calloc(a, b) mg_calloc(a, b)
 #define free(a) vPortFree(a)
 #define malloc(a) pvPortMalloc(a)
-#define strdup(s) mg_mprintf("%s", s)
+#define strdup(s) ((char *) mg_strdup(mg_str(s)).buf)
 
 // Re-route calloc/free to the FreeRTOS's functions, don't use stdlib
 static inline void *mg_calloc(size_t cnt, size_t size) {
@@ -288,7 +288,7 @@ extern uint32_t rt_time_get(void);
 #include "cmsis_os2.h"  // keep this include
 #endif
 
-#define strdup(s) mg_mprintf("%s", s)
+#define strdup(s) ((char *) mg_strdup(mg_str(s)).buf)
 
 #if defined(__ARMCC_VERSION)
 #define mode_t size_t
@@ -549,7 +549,6 @@ int sscanf(const char *, const char *, ...);
 
 #include <FreeRTOS_IP.h>
 #include <FreeRTOS_Sockets.h>
-#include <FreeRTOS_errno_TCP.h>  // contents to be moved and file removed, some day
 
 #define MG_SOCKET_TYPE Socket_t
 #define MG_INVALID_SOCKET FREERTOS_INVALID_SOCKET
@@ -581,6 +580,7 @@ int sscanf(const char *, const char *, ...);
 
 #define sockaddr_in freertos_sockaddr
 #define sockaddr freertos_sockaddr
+#define sin_addr sin_address.ulIP_IPv4
 #define accept(a, b, c) FreeRTOS_accept((a), (b), (c))
 #define connect(a, b, c) FreeRTOS_connect((a), (b), (c))
 #define bind(a, b, c) FreeRTOS_bind((a), (b), (c))
@@ -707,7 +707,7 @@ struct timeval {
 #endif
 
 #ifndef MG_ENABLE_IPV6
-#define MG_ENABLE_IPV6 1
+#define MG_ENABLE_IPV6 0
 #endif
 
 #ifndef MG_IPV6_V6ONLY
@@ -861,6 +861,7 @@ struct mg_str mg_str_n(const char *s, size_t n);
 int mg_casecmp(const char *s1, const char *s2);
 int mg_strcmp(const struct mg_str str1, const struct mg_str str2);
 int mg_strcasecmp(const struct mg_str str1, const struct mg_str str2);
+struct mg_str mg_strdup(const struct mg_str s);
 bool mg_match(struct mg_str str, struct mg_str pattern, struct mg_str *caps);
 bool mg_span(struct mg_str s, struct mg_str *a, struct mg_str *b, char delim);
 
@@ -1051,7 +1052,6 @@ uint16_t mg_ntohs(uint16_t net);
 uint32_t mg_ntohl(uint32_t net);
 uint32_t mg_crc32(uint32_t crc, const char *buf, size_t len);
 uint64_t mg_millis(void);  // Return milliseconds since boot
-uint64_t mg_now(void);     // Return milliseconds since Epoch
 bool mg_path_is_sane(const struct mg_str path);
 
 #define mg_htons(x) mg_ntohs(x)
@@ -1935,6 +1935,116 @@ typedef uint64_t mg_uecc_word_t;
 
 #endif /* _UECC_TYPES_H_ */
 // End of uecc BSD-2
+// portable8439 v1.0.1
+// Source: https://github.com/DavyLandman/portable8439
+// Licensed under CC0-1.0
+// Contains poly1305-donna e6ad6e091d30d7f4ec2d4f978be1fcfcbce72781 (Public
+// Domain)
+
+
+
+
+#ifndef __PORTABLE_8439_H
+#define __PORTABLE_8439_H
+#if defined(__cplusplus)
+extern "C" {
+#endif
+
+// provide your own decl specificier like -DPORTABLE_8439_DECL=ICACHE_RAM_ATTR
+#ifndef PORTABLE_8439_DECL
+#define PORTABLE_8439_DECL
+#endif
+
+/*
+ This library implements RFC 8439 a.k.a. ChaCha20-Poly1305 AEAD
+
+ You can use this library to avoid attackers mutating or reusing your
+ encrypted messages. This does assume you never reuse a nonce+key pair and,
+ if possible, carefully pick your associated data.
+*/
+
+// Make sure we are either nested in C++ or running in a C99+ compiler
+#if !defined(__cplusplus) && !defined(_MSC_VER) && \
+    (!defined(__STDC_VERSION__) || __STDC_VERSION__ < 199901L)
+#error "C99 or newer required"
+#endif
+
+// #if CHAR_BIT > 8
+// #    error "Systems without native octals not suppoted"
+// #endif
+
+#if defined(_MSC_VER) || defined(__cplusplus)
+// add restrict support is possible
+#if (defined(_MSC_VER) && _MSC_VER >= 1900) || defined(__clang__) || \
+    defined(__GNUC__)
+#define restrict __restrict
+#else
+#define restrict
+#endif
+#endif
+
+#define RFC_8439_TAG_SIZE (16)
+#define RFC_8439_KEY_SIZE (32)
+#define RFC_8439_NONCE_SIZE (12)
+
+/*
+    Encrypt/Seal plain text bytes into a cipher text that can only be
+    decrypted by knowing the key, nonce and associated data.
+
+    input:
+        - key: RFC_8439_KEY_SIZE bytes that all parties have agreed
+            upon beforehand
+        - nonce: RFC_8439_NONCE_SIZE bytes that should never be repeated
+            for the same key. A counter or a pseudo-random value are fine.
+        - ad: associated data to include with calculating the tag of the
+            cipher text. Can be null for empty.
+        - plain_text: data to be encrypted, pointer + size should not overlap
+            with cipher_text pointer
+
+    output:
+        - cipher_text: encrypted plain_text with a tag appended. Make sure to
+            allocate at least plain_text_size + RFC_8439_TAG_SIZE
+
+    returns:
+        - size of bytes written to cipher_text, can be -1 if overlapping
+            pointers are passed for plain_text and cipher_text
+*/
+PORTABLE_8439_DECL size_t mg_chacha20_poly1305_encrypt(
+    uint8_t *restrict cipher_text, const uint8_t key[RFC_8439_KEY_SIZE],
+    const uint8_t nonce[RFC_8439_NONCE_SIZE], const uint8_t *restrict ad,
+    size_t ad_size, const uint8_t *restrict plain_text, size_t plain_text_size);
+
+/*
+    Decrypt/unseal cipher text given the right key, nonce, and additional data.
+
+    input:
+        - key: RFC_8439_KEY_SIZE bytes that all parties have agreed
+            upon beforehand
+        - nonce: RFC_8439_NONCE_SIZE bytes that should never be repeated for
+            the same key. A counter or a pseudo-random value are fine.
+        - ad: associated data to include with calculating the tag of the
+            cipher text. Can be null for empty.
+        - cipher_text: encrypted message.
+
+    output:
+        - plain_text: data to be encrypted, pointer + size should not overlap
+            with cipher_text pointer, leave at least enough room for
+            cipher_text_size - RFC_8439_TAG_SIZE
+
+    returns:
+        - size of bytes written to plain_text, -1 signals either:
+            - incorrect key/nonce/ad
+            - corrupted cipher_text
+            - overlapping pointers are passed for plain_text and cipher_text
+*/
+PORTABLE_8439_DECL size_t mg_chacha20_poly1305_decrypt(
+    uint8_t *restrict plain_text, const uint8_t key[RFC_8439_KEY_SIZE],
+    const uint8_t nonce[RFC_8439_NONCE_SIZE],
+    const uint8_t *restrict cipher_text, size_t cipher_text_size);
+#if defined(__cplusplus)
+}
+#endif
+#endif
 
 
 struct mg_connection;
@@ -2155,7 +2265,7 @@ void mg_http_serve_ssi(struct mg_connection *c, const char *root,
 #define MG_TLS_CUSTOM 4   // Custom implementation
 
 #ifndef MG_TLS
-#define MG_TLS MG_TLS_NONE
+#define MG_TLS MG_TLS_BUILTIN
 #endif
 
 
@@ -2265,6 +2375,8 @@ struct mg_connection *mg_sntp_connect(struct mg_mgr *mgr, const char *url,
 void mg_sntp_request(struct mg_connection *c);
 int64_t mg_sntp_parse(const unsigned char *buf, size_t len);
 
+uint64_t mg_now(void);     // Return milliseconds since Epoch
+
 
 
 
@@ -2571,6 +2683,16 @@ struct mg_tcpip_driver {
   bool (*up)(struct mg_tcpip_if *);                           // Up/down status
 };
 
+typedef void (*mg_tcpip_event_handler_t)(struct mg_tcpip_if *ifp, int ev,
+                                         void *ev_data);
+
+enum {
+  MG_TCPIP_EV_ST_CHG,     // state change             uint8_t * (&ifp->state)
+  MG_TCPIP_EV_DHCP_DNS,   // DHCP DNS assignment      uint32_t *ipaddr
+  MG_TCPIP_EV_DHCP_SNTP,  // DHCP SNTP assignment     uint32_t *ipaddr
+  MG_TCPIP_EV_USER        // Starting ID for user events
+};
+
 // Network interface
 struct mg_tcpip_if {
   uint8_t mac[6];                  // MAC address. Must be set to a valid MAC
@@ -2579,10 +2701,13 @@ struct mg_tcpip_if {
   bool enable_dhcp_client;         // Enable DCHP client
   bool enable_dhcp_server;         // Enable DCHP server
   bool enable_get_gateway;         // DCHP server sets client as gateway
+  bool enable_req_dns;             // DCHP client requests DNS server
+  bool enable_req_sntp;            // DCHP client requests SNTP server
   bool enable_crc32_check;         // Do a CRC check on RX frames and strip it
   bool enable_mac_check;           // Do a MAC check on RX frames
   struct mg_tcpip_driver *driver;  // Low level driver
   void *driver_data;               // Driver-specific data
+  mg_tcpip_event_handler_t fn;     // User-specified event handler function
   struct mg_mgr *mgr;              // Mongoose event manager
   struct mg_queue recv_queue;      // Receive queue
   uint16_t mtu;                    // Interface MTU
@@ -2852,8 +2977,14 @@ struct mg_tcpip_driver_stm32f_data {
 #endif
 
 
-#if MG_ENABLE_TCPIP && defined(MG_ENABLE_DRIVER_STM32H) && \
-    MG_ENABLE_DRIVER_STM32H
+#if MG_ENABLE_TCPIP
+#if !defined(MG_ENABLE_DRIVER_STM32H)
+#define MG_ENABLE_DRIVER_STM32H 0
+#endif
+#if !defined(MG_ENABLE_DRIVER_MCXN)
+#define MG_ENABLE_DRIVER_MCXN 0
+#endif
+#if MG_ENABLE_DRIVER_STM32H || MG_ENABLE_DRIVER_MCXN
 
 struct mg_tcpip_driver_stm32h_data {
   // MDC clock divider. MDC clock is derived from HCLK, must not exceed 2.5MHz
@@ -2866,7 +2997,8 @@ struct mg_tcpip_driver_stm32h_data {
   //    35-60 MHz     HCLK/26        3
   //    150-250 MHz   HCLK/102       4  <-- value for max speed HSI
   //    250-300 MHz   HCLK/124       5  <-- value for Nucleo-H* on CSI
-  //    110, 111 Reserved
+  //    300-500 MHz   HCLK/204       6
+  //    500-800 MHz   HCLK/324       7
   int mdc_cr;  // Valid values: -1, 0, 1, 2, 3, 4, 5
 
   uint8_t phy_addr;  // PHY address
@@ -2898,6 +3030,7 @@ struct mg_tcpip_driver_stm32h_data {
   } while (0)
 
 #endif
+#endif
 
 
 #if MG_ENABLE_TCPIP && defined(MG_ENABLE_DRIVER_TM4C) && MG_ENABLE_DRIVER_TM4C
@@ -2924,9 +3057,6 @@ struct mg_tcpip_driver_tm4c_data {
 
 #if MG_ENABLE_TCPIP && defined(MG_ENABLE_DRIVER_W5500) && MG_ENABLE_DRIVER_W5500
 
-#undef MG_ENABLE_TCPIP_DRIVER_INIT
-#define MG_ENABLE_TCPIP_DRIVER_INIT 0
-
 #endif
 
 

+ 15 - 11
pro/src/app.c

@@ -18,6 +18,7 @@
 #include "cfg.h"
 #include "lock.h"
 #include "smtp.h"
+#include "mqtt.h"
 #include "netswitch.h"
 #include "breaker_detection.h"
 #include <sys/mman.h>
@@ -165,9 +166,8 @@ void* power_thread(void* arg)
         int _chn = 0;
         memset(Vavg,0.0,sizeof(Vavg));
         memset(VavgTimes,0,sizeof(VavgTimes));
-        //数据刷新
-        pthread_mutex_lock(&__globalDeviceManage._power_update);
-                //清空三项总数据
+        
+        lock_s_hold(LOCK_ID_POWER_UPDATE);
         if (__globalDeviceManage._globalDevInfo.product_pwr_type == SmartPDU_Tree_AC_One_B)
         {
             list_for_each_entry(pTreeAPowerInfo, &_globalDeviceManager->_globalPowerManger.list_Tree_AC, list_Tree_AC)
@@ -970,7 +970,7 @@ void* power_thread(void* arg)
                 log_d("ch%d insert data error.\n", _globalTotalPowerInfo.product_ch_id);
             }
         }
-        pthread_mutex_unlock(&__globalDeviceManage._power_update);
+        lock_s_release(LOCK_ID_POWER_UPDATE);
   
         //判断阈值状态
 
@@ -1200,12 +1200,11 @@ void* breaker_thread(void *arg)
     int r;
     GlobalDeviceManager *_globalDeviceManager = (GlobalDeviceManager *)&__globalDeviceManage;
     GlobalBreakerManager *temp = NULL;
-    pthread_mutex_init(&_globalDeviceManager->_breaker_mutex,NULL);
     uint16_t value = 0;
     breaker_scanner_t* breaker_scann = get_breaker_scanner(); 
     while (_globalDeviceManager->dev_samp_flag)
     {
-        pthread_mutex_lock(&__globalDeviceManage._breaker_mutex);
+        lock_s_hold(LOCK_ID_BREAKER);
         list_for_each_entry(temp, &_globalDeviceManager->g_new_global_breaker.list,list)
         {
             // if(temp->breaker_gather_addr != 0)
@@ -1234,7 +1233,7 @@ void* breaker_thread(void *arg)
                 }
             }
         }
-        pthread_mutex_unlock(&__globalDeviceManage._breaker_mutex);
+        lock_s_release(LOCK_ID_BREAKER);
         sleep(1);
     }
     pthread_exit(NULL);
@@ -1296,7 +1295,8 @@ void* sensor_thread(void* arg)
         {
             dev_change_last_Index(_globalDeviceManager->db, &nSaveIndex, nDeleteNumber, strTable);
         }
-        pthread_mutex_lock(&__globalDeviceManage._sensor_update);
+        
+        lock_s_hold(LOCK_ID_SENSOR);
         list_for_each_entry(_globalSensorMangerTemp, &_globalDeviceManager->_globalSensorManger.list, list)
         {
             switch (_globalSensorMangerTemp->sensor_type)
@@ -1441,7 +1441,7 @@ void* sensor_thread(void* arg)
                 sensor_alarm(_globalSensorMangerTemp, SENSOR_VAL2_LOWER);
             }
         }
-        pthread_mutex_unlock(&__globalDeviceManage._sensor_update);
+        lock_s_release(LOCK_ID_SENSOR);
         sleep(1);
     }
 #endif
@@ -1868,7 +1868,7 @@ int _global_device_manage_init(GlobalDeviceManager* _globalDeviceManager)
     log_d("GCPDU init begin!");
     ResetChmData(0);
 
-#if 1
+
     //初始化传感器列表
     INIT_LIST_HEAD(&_globalDeviceManager->_globalSensorManger.list);
     INIT_LIST_HEAD(&_globalDeviceManager->g_new_global_breaker.list);
@@ -1895,6 +1895,8 @@ int _global_device_manage_init(GlobalDeviceManager* _globalDeviceManager)
     }
     else
         log_d("get breaker info ok.");
+
+#if 1
     //创建modbus采集线程 继电器控制线程
     #if(OS_HANDWARE==1)
     _globalDeviceManager->dev_samp_flag = true;
@@ -1949,7 +1951,6 @@ int _global_device_manage_init(GlobalDeviceManager* _globalDeviceManager)
 #endif
 
     dev_Alarm_Run_message(_globalDeviceManager, language_alarm_Init_Success[0], "SNMP");
- #endif
 
     // 初始化web
     websocket_init();
@@ -1958,6 +1959,9 @@ int _global_device_manage_init(GlobalDeviceManager* _globalDeviceManager)
     cascade_init();
     netswitch_init();
     smtp_init();
+#endif
+
+    mqtt_init();
 
     //通知smartUPG,app已经运行起来
     upg_init();

+ 90 - 35
pro/src/appweb_handle.c

@@ -9,6 +9,7 @@
 #include "cfg.h"
 #include "cascade.h"
 #include "upg.h"
+#include "lock.h"
 #include "file.h"
 #include "mongoose.h"
 #include "language_common.h"
@@ -1184,9 +1185,9 @@ static void channelStatusMonitoring(void* conn)
         else {
             if(product_chn_id==0)
             {
-                pthread_mutex_lock(&__globalDeviceManage._power_update); 
+                lock_s_hold(LOCK_ID_POWER_UPDATE);
                 dev_search_latest_power_All_info(&__globalDeviceManage,&_over_chn_pwr_back_info, cur_dev_addr);
-                pthread_mutex_unlock(&__globalDeviceManage._power_update);
+                lock_s_release(LOCK_ID_POWER_UPDATE);
             }
             else //指定通道
             {
@@ -1844,8 +1845,8 @@ static void sensorStatusMonitoring(void* conn)
         log_d("start listQuery handle.");
         
         #if(1)
-        //链表初始化
-        pthread_mutex_lock(&__globalDeviceManage._sensor_update);
+        lock_s_hold(LOCK_ID_SENSOR);
+
         INIT_LIST_HEAD(&_over_all_sensor_back_info.list);
         //获取传感器当前的最新值
         GlobalSensorManger* _globalSensorManagerTemp = NULL;
@@ -1918,7 +1919,8 @@ static void sensorStatusMonitoring(void* conn)
             //入队
             list_add_tail(&_overSensorBackInfoTemp->list,&_over_all_sensor_back_info.list);
         }    
-        pthread_mutex_unlock(&__globalDeviceManage._sensor_update);
+        lock_s_release(LOCK_ID_SENSOR);
+
         //参数转换json
         ack_str = over_all_sensor_info_to_json(&_over_all_sensor_back_info,0);
 
@@ -1945,7 +1947,8 @@ static void sensorStatusMonitoring(void* conn)
         else
         {
             // 维护列表也需要删除
-            pthread_mutex_lock(&__globalDeviceManage._sensor_update);
+            lock_s_hold(LOCK_ID_SENSOR);
+
             GlobalSensorManger *node;
             GlobalSensorManger *next;
             list_for_each_entry_safe(node, next, &__globalDeviceManage._globalSensorManger.list, list)
@@ -1957,7 +1960,7 @@ static void sensorStatusMonitoring(void* conn)
                     break;
                 }
             }
-            pthread_mutex_unlock(&__globalDeviceManage._sensor_update);
+            lock_s_release(LOCK_ID_SENSOR);
         }
 
 
@@ -2195,9 +2198,9 @@ static void sensorStatusMonitoring(void* conn)
             log_e("dev insert error.");
         }
         //全局接口体中也需要添加相关参数
-        pthread_mutex_lock(&__globalDeviceManage._sensor_update);
+        lock_s_hold(LOCK_ID_SENSOR);
         list_add_tail(&_globalSensorMangerTemp->list,&__globalDeviceManage._globalSensorManger.list);
-        pthread_mutex_unlock(&__globalDeviceManage._sensor_update);
+        lock_s_release(LOCK_ID_SENSOR);
 
         char temp[100];
        // sprintf(temp,"$新增$|$传感器$|%d",_over_all_sensor_request.sensorId);
@@ -3367,8 +3370,48 @@ static void serviceManagement(void *conn)
                 dev_smtp_update_recv(gdm->db, &recv[i]);
             }
         }
-    }   
-    //查询snmp信息
+    }
+    else if(strcmp(_serviceManageRequestInfo.type,"listMqttQuery")==0)
+    {
+        ack_str = service_mqtt_ack_to_json(&gdm->mqttInfo);
+    }
+    else if(strcmp(_serviceManageRequestInfo.type,"mqttSave")==0)
+    {
+        int save_flag=0;
+        _MQTT_ServerRequestInfo *req=&_serviceManageRequestInfo._mqtt_info;
+        mqtt_server_t *ser=&gdm->mqttInfo.ser[req->idx];
+
+        if(req->idx |= ser->id) { ser->id = req->idx; 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->user) { strcpy(ser->user, req->user); save_flag=1; }
+        if(req->password) { strcpy(ser->password, req->password); save_flag=1; }
+        if(req->cid) { strcpy(ser->cid, req->cid); save_flag=1; }
+
+        if(save_flag) {
+            dev_mqtt_update(gdm->db, ser);
+        }
+    }
+    else if(strcmp(_serviceManageRequestInfo.type,"mqttOpen")==0 ||
+            strcmp(_serviceManageRequestInfo.type,"mqttClose")==0)
+    {
+        int save_flag=0;
+        _MQTT_ServerRequestInfo *req=&_serviceManageRequestInfo._mqtt_info;
+        mqtt_server_t *ser=&gdm->mqttInfo.ser[req->idx];
+        if(req->mode!=ser->mode) {
+            ser->mode = req->mode;
+            dev_mqtt_update(gdm->db, ser);
+        }
+    }
+    else if(strcmp(_serviceManageRequestInfo.type,"mqttZtQuery")==0)
+    {
+        //
+    }
+    else if(strcmp(_serviceManageRequestInfo.type,"mqttZtSave")==0)
+    {
+        //
+    }
     else if(strcmp(_serviceManageRequestInfo.type,"listSnmpQuery")==0)
     {
         // 查询SNMP信息
@@ -3575,16 +3618,16 @@ static void serviceManagement(void *conn)
         }
         */
         // 返回定时信息
-        pthread_mutex_lock(&__globalDeviceManage._power_ds);             
+        lock_s_hold(LOCK_ID_POWER_TIMER);
         ack_str = service_ds_all_ack_to_json(_serviceManageRequestInfo._ds_info.channelDate,
                                         _serviceManageRequestInfo._ds_info.productChName,
                                         _serviceManageRequestInfo._ds_info.productChId,
                                         &__globalDeviceManage._globalPowerManger);
-        pthread_mutex_unlock(&__globalDeviceManage._power_ds);
+        lock_s_release(LOCK_ID_POWER_TIMER);
     }
     else if (strcmp(_serviceManageRequestInfo.type, "dsSave") == 0)
     {
-        pthread_mutex_lock(&__globalDeviceManage._power_ds);         
+        lock_s_hold(LOCK_ID_POWER_TIMER);
         int chn_id = -1;
         for (int i = 0; i < 64; i++)
         {
@@ -3640,7 +3683,7 @@ static void serviceManagement(void *conn)
                 }
             }
         }
-        pthread_mutex_unlock(&__globalDeviceManage._power_ds);
+        lock_s_release(LOCK_ID_POWER_TIMER);
 
         char temp[100];
         //sprintf(temp, "$设置$|定时$");
@@ -3652,13 +3695,15 @@ static void serviceManagement(void *conn)
     //批量删除
     else if (strcmp(_serviceManageRequestInfo.type, "deleteDs")==0)
     {
-        pthread_mutex_lock(&__globalDeviceManage._power_ds);     
+        
         int nChID = 0;
         char pcDate[16];
         char pcTime[16];
         GlobalPowerManger *_globalPowerMangerTemp = NULL;
         _PowerDSManage_t* _powerDsManageTemp=NULL;;
         char *pCharData;
+
+        lock_s_hold(LOCK_ID_POWER_TIMER);
         for (int i = 0; i < 64; i++)
         {
             pCharData = _serviceManageRequestInfo._ds_info.idx[i];
@@ -3684,7 +3729,7 @@ static void serviceManagement(void *conn)
                 }                
             }
         }
-        pthread_mutex_unlock(&__globalDeviceManage._power_ds);
+        lock_s_release(LOCK_ID_POWER_TIMER);
 
         char temp[100];
        // sprintf(temp, "$删除$|定时$");
@@ -4429,7 +4474,7 @@ static void equipmentMintenance(void *conn)
             ack_str=ctrl_sts_ack_to_json();
             // GlobalBreakerManager *temp = NULL;
             // GlobalDeviceManager *dm = &__globalDeviceManage2;
-            // pthread_mutex_lock(&dm->_breaker_mutex);
+            // lock_s_hold(LOCK_ID_BREAKER);
             // list_for_each_entry(temp, &dm->g_new_global_breaker.list, list)
             // {
             //     char buff[120];
@@ -4458,7 +4503,7 @@ static void equipmentMintenance(void *conn)
             //     info->status,info->cjfs,info->comkh,info->cjdz,info->jdh);
             //     list_add_tail(&info->list,&_over_breaker_switch_info.list);
             // }
-            // pthread_mutex_unlock(&dm->_breaker_mutex);
+            // lock_s_release(LOCK_ID_BREAKER);
             // ack_str = over_all_breaker_switch_info_to_json(&_over_breaker_switch_info,0);
             // _OverBreakerSwitchInfo* node;
             // _OverBreakerSwitchInfo* next;
@@ -4474,7 +4519,7 @@ static void equipmentMintenance(void *conn)
         }else   //current 
         {
             GlobalBreakerManager *temp = NULL;
-            pthread_mutex_lock(&__globalDeviceManage._breaker_mutex);
+            lock_s_hold(LOCK_ID_BREAKER);
             list_for_each_entry(temp, &__globalDeviceManage.g_new_global_breaker.list, list)
             {
                 char buff[20];
@@ -4502,7 +4547,8 @@ static void equipmentMintenance(void *conn)
                 // info->status,info->cjfs,info->comkh,info->cjdz,info->jdh);
                 list_add_tail(&info->list,&_over_breaker_switch_info.list);
             }
-            pthread_mutex_unlock(&__globalDeviceManage._breaker_mutex);
+            lock_s_release(LOCK_ID_BREAKER);
+
             printf("_glos_asdsd  addr=0x%x\n",(uint32_t)(&__globalDeviceManage));
             ack_str = over_all_breaker_switch_info_to_json(&_over_breaker_switch_info,0);
 
@@ -4539,12 +4585,13 @@ static void equipmentMintenance(void *conn)
             //bianli shujujiegou
             GlobalBreakerManager *temp = NULL;
             int i = 1;
-             pthread_mutex_lock(&__globalDeviceManage._breaker_mutex);
+
+            lock_s_hold(LOCK_ID_BREAKER);
             list_for_each_entry(temp, &__globalDeviceManage.g_new_global_breaker.list, list)
             {
                 id[temp->breaker_id] = 1; 
             }
-            pthread_mutex_unlock(&__globalDeviceManage._breaker_mutex);
+            lock_s_release(LOCK_ID_BREAKER);
             for(;i < MAX_BREAKER_ID;i++)
             {
                 if(id[i] == 0)
@@ -4569,7 +4616,8 @@ static void equipmentMintenance(void *conn)
                 {
                     char jdh_flag = 0;
                     char chn[MAX_BREAKER_CHN] = {0};
-                    pthread_mutex_lock(&__globalDeviceManage._breaker_mutex);
+
+                    lock_s_hold(LOCK_ID_BREAKER);
                     list_for_each_entry(temp, &__globalDeviceManage.g_new_global_breaker.list, list)
                     {
                         if(temp->breaker_gather_addr == breakerManager->breaker_gather_addr)
@@ -4577,7 +4625,8 @@ static void equipmentMintenance(void *conn)
                             chn[temp->breaker_chn-1] = 1;
                         }
                     }
-                    pthread_mutex_unlock(&__globalDeviceManage._breaker_mutex);
+                    lock_s_release(LOCK_ID_BREAKER);
+
                     if(chn[jdh-1] != 1)
                     {
                         breakerManager->breaker_chn = jdh;
@@ -4594,9 +4643,11 @@ static void equipmentMintenance(void *conn)
                         
                     }else{
                         dev_insert_breaker_genera_manage(__globalDeviceManage.db,breakerManager);
-                        pthread_mutex_lock(&__globalDeviceManage._breaker_mutex);
+
+                        lock_s_hold(LOCK_ID_BREAKER);
                         list_add_tail(&breakerManager->list,&__globalDeviceManage.g_new_global_breaker.list);
-                        pthread_mutex_unlock(&__globalDeviceManage._breaker_mutex);
+                        lock_s_release(LOCK_ID_BREAKER);
+
                         ack_str=ctrl_sts_ack_to_json();
                     }
                 }else {
@@ -4629,7 +4680,8 @@ static void equipmentMintenance(void *conn)
         {
             GlobalBreakerManager *temp = NULL;
             int switch_id = _Over_equipment_mintenance.switchId;
-            pthread_mutex_lock(&__globalDeviceManage._breaker_mutex);
+
+            lock_s_hold(LOCK_ID_BREAKER);
             list_for_each_entry(temp, &__globalDeviceManage.g_new_global_breaker.list, list)
             {
                 if(temp->breaker_id == switch_id)
@@ -4641,7 +4693,8 @@ static void equipmentMintenance(void *conn)
                     break;
                 }
             }
-            pthread_mutex_unlock(&__globalDeviceManage._breaker_mutex);
+            lock_s_release(LOCK_ID_BREAKER);
+
             ack_str=ctrl_sts_ack_to_json();
         }
     }else if(strcmp(_Over_equipment_mintenance.type,"switchDelete") == 0)
@@ -4657,7 +4710,8 @@ static void equipmentMintenance(void *conn)
             cascade_request(&cmd);
             GlobalBreakerManager* node = NULL;
             GlobalBreakerManager* next = NULL;
-            pthread_mutex_lock(&__globalDeviceManage._breaker_mutex);
+
+            lock_s_hold(LOCK_ID_BREAKER);
             list_for_each_entry_safe(node,next,&__globalDeviceManage2.g_new_global_breaker.list, list)
             {
                 if(node->breaker_id == _Over_equipment_mintenance.switchId)
@@ -4669,7 +4723,7 @@ static void equipmentMintenance(void *conn)
                     break;
                 }
             }
-            pthread_mutex_unlock(&__globalDeviceManage._breaker_mutex);
+            lock_s_release(LOCK_ID_BREAKER);
 
             ack_str=ctrl_sts_ack_to_json();
         }else
@@ -4678,7 +4732,8 @@ static void equipmentMintenance(void *conn)
             //GlobalBreakerManager *temp = NULL;
             GlobalBreakerManager* node;
             GlobalBreakerManager* next;
-            pthread_mutex_lock(&__globalDeviceManage._breaker_mutex);
+
+            lock_s_hold(LOCK_ID_BREAKER);
             list_for_each_entry_safe(node,next,&__globalDeviceManage.g_new_global_breaker.list, list)
             {
                 if(node->breaker_id == switch_id)
@@ -4690,7 +4745,7 @@ static void equipmentMintenance(void *conn)
                     break;
                 }
             }
-            pthread_mutex_unlock(&__globalDeviceManage._breaker_mutex);
+            lock_s_release(LOCK_ID_BREAKER);
             ack_str=ctrl_sts_ack_to_json();
         }
 
@@ -4717,9 +4772,9 @@ static void equipmentMintenance(void *conn)
         http_reply(conn, 200, ack_str);
 
         //__globalDeviceManage._ReloadChm_flag=true;
-        pthread_mutex_lock(&__globalDeviceManage._power_update);
+        lock_s_hold(LOCK_ID_POWER_UPDATE);
         ResetChmData(1);
-        pthread_mutex_unlock(&__globalDeviceManage._power_update);
+        lock_s_release(LOCK_ID_POWER_UPDATE);
         //__globalDeviceManage._ReloadChm_flag=false;
     }
     else if (strcmp(_Over_equipment_mintenance.type,"cleanUp") == 0)

+ 7 - 6
pro/src/breaker_detection.c

@@ -8,6 +8,7 @@
 #include "common.h"
 #include "thread.h"
 #include "string.h"
+#include "lock.h"
 
 static breaker_scanner_t breaker_scanner;
 
@@ -202,12 +203,12 @@ void* breaker_scanner_thread(void* argv)
             {
 
                     breaker_485_get_all_chn_status((void*)&gm->_globalRelaySampManger,BREADER_485_ADDR+i,status);
-                    pthread_mutex_lock(&gm->_breaker_mutex);
+                    lock_s_hold(LOCK_ID_BREAKER);
                     memcpy(breaker->board[i].ch_status,status,sizeof(status));
-                    pthread_mutex_unlock(&gm->_breaker_mutex);
+                    lock_s_release(LOCK_ID_BREAKER);
                 // ret = breaker_485_get_chn_status((void*)&gm->_globalRelaySampManger,
                 //     BREADER_485_ADDR+i,BREADER_485_ADDR_CHN_REG,&value);
-                // //pthread_mutex_lock(&gm->_breaker_mutex);
+                // //lock_s_hold(LOCK_ID_BREAKER);
                 // if(ret == 0)
                 // {
                 //     breaker_scanner.breaker_flag[i] = 1;  // board is online
@@ -219,14 +220,14 @@ void* breaker_scanner_thread(void* argv)
                 // }
                 
                 
-                // // pthread_mutex_unlock(&gm->_breaker_mutex);
+                // // lock_s_release(LOCK_ID_BREAKER);
                 memset(status,0,sizeof(status));
             }else
             {
                 memset(status,0,sizeof(status));
-                pthread_mutex_lock(&gm->_breaker_mutex);
+                lock_s_hold(LOCK_ID_BREAKER);
                 memcpy(breaker->board[i].ch_status,status,sizeof(status));
-                pthread_mutex_unlock(&gm->_breaker_mutex);
+                lock_s_release(LOCK_ID_BREAKER);
             }
         }
         sleep(1);

+ 19 - 12
pro/src/cascade.c

@@ -2,6 +2,7 @@
 #include "common.h"
 #include "elog.h"
 #include "cfg.h"
+#include "lock.h"
 #include "thread.h"
 #include "switch_ctrl.h"
 #include "modbus_handle.h"
@@ -211,7 +212,7 @@ static int slave_breaker_info(cascade_handle_t *cas)
     breaker_info_t *breaker=&cas->sBreaker;
     int cnt = 0;
     //list_for_each_entry(temp, &__globalDeviceManage.g_new_global_breaker.list, list)
-    pthread_mutex_lock(&dm->_breaker_mutex);
+    lock_s_hold(LOCK_ID_BREAKER);
     if(!list_empty(&dm->g_new_global_breaker.list)) {
         printf("isnot empty!!!!!\n");
         list_for_each_entry(temp,&dm->g_new_global_breaker.list, list)
@@ -231,7 +232,7 @@ static int slave_breaker_info(cascade_handle_t *cas)
         
     }
     breaker->cnt = cnt;
-    pthread_mutex_unlock(&dm->_breaker_mutex);
+    lock_s_release(LOCK_ID_BREAKER);
     
     LOGD("___ slave_breaker_info, breaker_all: saddr=0x%x %d\n",(uint32_t)dm,breaker->cnt);
 
@@ -272,7 +273,8 @@ static int slave_breaker_update(cascade_breaker_update_t *data)
         GlobalDeviceManager *dm=get_dm();
         GlobalBreakerManager *temp= NULL;
         GlobalBreakerManager *pos= NULL;
-        pthread_mutex_lock(&dm->_breaker_mutex);
+
+        lock_s_hold(LOCK_ID_BREAKER);
         if(!list_empty(&dm->g_new_global_breaker.list)) {
             
             list_for_each_entry_safe(temp,pos,&dm->g_new_global_breaker.list,list)
@@ -288,7 +290,7 @@ static int slave_breaker_update(cascade_breaker_update_t *data)
                 }
             }
         }
-        pthread_mutex_unlock(&dm->_breaker_mutex);
+        lock_s_release(LOCK_ID_BREAKER);
         
     }
     return ret;
@@ -302,7 +304,8 @@ static int slave_breaker_delete(cascade_braeker_delete_t *data)
         GlobalDeviceManager *dm=get_dm();
         GlobalBreakerManager *temp= NULL;
         GlobalBreakerManager *pos= NULL;
-        pthread_mutex_lock(&dm->_breaker_mutex);
+
+        lock_s_hold(LOCK_ID_BREAKER);
         if(!list_empty(&dm->g_new_global_breaker.list)) {
             
             list_for_each_entry_safe(temp,pos,&dm->g_new_global_breaker.list,list)
@@ -328,7 +331,7 @@ static int slave_breaker_delete(cascade_braeker_delete_t *data)
             }
             
         }
-        pthread_mutex_unlock(&dm->_breaker_mutex);
+        lock_s_release(LOCK_ID_BREAKER);
     }
     return ret;
 }
@@ -349,12 +352,14 @@ static int slave_breaker_add(cascade_breaker_add_t *data)
             //bianli shujujiegou
     GlobalBreakerManager *temp = NULL;
     int i = 1;
-    pthread_mutex_lock(&dm->_breaker_mutex);
+
+    lock_s_hold(LOCK_ID_BREAKER);
     list_for_each_entry(temp, &dm->g_new_global_breaker.list, list)
     {           
         id[temp->breaker_id] = 1; 
     }
-    pthread_mutex_unlock(&dm->_breaker_mutex);
+    lock_s_release(LOCK_ID_BREAKER);
+
     for(;i < MAX_BREAKER_ID;i++)
     {
         if(id[i] == 0)
@@ -387,7 +392,8 @@ static int slave_breaker_add(cascade_breaker_add_t *data)
         char jdh_flag = 0;
 
         char chn[MAX_BREAKER_CHN] = {0};
-        pthread_mutex_lock(&dm->_breaker_mutex);
+
+        lock_s_hold(LOCK_ID_BREAKER);
         list_for_each_entry(temp, &dm->g_new_global_breaker.list, list)
         {
             if(temp->breaker_gather_addr == breakerManager->breaker_gather_addr)
@@ -395,7 +401,7 @@ static int slave_breaker_add(cascade_breaker_add_t *data)
                 chn[temp->breaker_chn-1] = 1;
             }
         }
-        pthread_mutex_unlock(&dm->_breaker_mutex);
+        lock_s_release(LOCK_ID_BREAKER);
 
         if(chn[jdh-1] != 1)
         {
@@ -410,9 +416,10 @@ static int slave_breaker_add(cascade_breaker_add_t *data)
             goto quit2;
         }else{
             dev_insert_breaker_genera_manage(dm->db,breakerManager);
-            pthread_mutex_lock(&dm->_breaker_mutex);
+
+            lock_s_hold(LOCK_ID_BREAKER);
             list_add_tail(&breakerManager->list,&dm->g_new_global_breaker.list);
-            pthread_mutex_unlock(&dm->_breaker_mutex);
+            lock_s_release(LOCK_ID_BREAKER);
             r = 0;
         }
 

+ 5 - 5
pro/src/dflt.c

@@ -55,11 +55,11 @@ smtp_send_t DFLT_SMTP_SEND={
 mqtt_server_t DFLT_MQTT={
     .id = 0,
     .mode = 1,
-    .server = "smtp.163.com",
-    .ip = "rcp064867@163.com",
+    .server = "192.168.234.10",
+    .ip = "192.168.234.10",
     .port = "25",
-    .cid = ",",
-    .user = "UODRIFWDFBTJTVLW",
-    .password = "PLAIN",
+    .cid = "mqtlx_928507fe",
+    .user = "root",
+    .password = "root",
 };
 

+ 77 - 38
pro/src/json_handle.c

@@ -7,6 +7,7 @@
 #include "list.h"
 #include "file.h"
 #include "sys.h"
+#include "lock.h"
 #include "cascade.h"
 #include "netswitch.h"
 
@@ -2786,7 +2787,7 @@ char *breaker_status_ack_to_json(GlobalDeviceManager *gm,int slave)
     _OverBreakerSwitchInfo _over_breaker_switch_info;
     INIT_LIST_HEAD(&_over_breaker_switch_info.list);
     
-    pthread_mutex_lock(&gm->_breaker_mutex);
+    lock_s_hold(LOCK_ID_BREAKER);
     list_for_each_entry(temp, &_globalbreakerManger->list, list)
     {
         char buff[120] = {0};
@@ -2817,7 +2818,8 @@ char *breaker_status_ack_to_json(GlobalDeviceManager *gm,int slave)
                 //printf("push %s %d %d %s %s %s %s\n",info->productChids,info->switch_id,info->status,info->cjfs,info->comkh,info->cjdz,info->jdh);
                 list_add_tail(&info->list,&_over_breaker_switch_info.list);
     }
-    pthread_mutex_unlock(&gm->_breaker_mutex);
+    lock_s_release(LOCK_ID_BREAKER);
+
     strRet = over_all_breaker_switch_info_to_json(&_over_breaker_switch_info,1);
                 //释放资源
     _OverBreakerSwitchInfo* node;
@@ -4350,29 +4352,71 @@ cJSON* json_to_service_manage(const char* str,_ServiceManageRequestInfo* _servic
             }
         }
     }
-    else if(strcmp(_serviceManageRequestInfo->type,"listMqttQuery")==0)
-    {
-        //
-    }
     else if(strcmp(_serviceManageRequestInfo->type,"mqttSave")==0)
     {
-        //
+        cJSON *tmp;
+        _MQTT_ServerRequestInfo *info=&_serviceManageRequestInfo->_mqtt_info;
+
+        tmp = cJSON_GetObjectItem(cjson, "idx");
+        if(tmp) {
+            info->idx = atoi(tmp->valuestring);
+        }
+
+        tmp = cJSON_GetObjectItem(cjson, "serverName");
+        if(tmp) {
+            info->server = tmp->valuestring;
+        }
+
+        tmp = cJSON_GetObjectItem(cjson, "ip");
+        if(tmp) {
+            info->ip = tmp->valuestring;
+        }
+
+        tmp = cJSON_GetObjectItem(cjson, "port");
+        if(tmp) {
+            info->port = tmp->valuestring;
+        }
+
+        tmp = cJSON_GetObjectItem(cjson, "kfdId");
+        if(tmp) {
+            info->cid = tmp->valuestring;
+        }
+
+        tmp = cJSON_GetObjectItem(cjson, "yhm");
+        if(tmp) {
+            info->user = tmp->valuestring;
+        }
+
+        tmp = cJSON_GetObjectItem(cjson, "mm");
+        if(tmp) {
+            info->password = tmp->valuestring;
+        }
     }
     else if(strcmp(_serviceManageRequestInfo->type,"mqttOpen")==0)
     {
-        //
+        _MQTT_ServerRequestInfo *info=&_serviceManageRequestInfo->_mqtt_info;
+        cJSON *tmp=cJSON_GetObjectItem(cjson, "idx");
+        if(tmp) {
+            info->idx = atoi(tmp->valuestring);
+        }
+        info->mode = 1;
     }
     else if(strcmp(_serviceManageRequestInfo->type,"mqttClose")==0)
     {
-        //
-    }
-    else if(strcmp(_serviceManageRequestInfo->type,"mqttZtQuery")==0)
-    {
-        //
+        _MQTT_ServerRequestInfo *info=&_serviceManageRequestInfo->_mqtt_info;
+        cJSON *tmp=cJSON_GetObjectItem(cjson, "idx");
+        if(tmp) {
+            info->idx = atoi(tmp->valuestring);
+        }
+        info->mode = 0;
     }
     else if(strcmp(_serviceManageRequestInfo->type,"mqttZtSave")==0)
     {
-        //
+        _MQTT_ServerRequestInfo *info=&_serviceManageRequestInfo->_mqtt_info;
+        cJSON *tmp=cJSON_GetObjectItem(cjson, "mode");
+        if(tmp) {
+            //info->mode = atoi(tmp->valuestring);
+        }
     }
     //SNMP上传参数
     else if(strcmp(_serviceManageRequestInfo->type,"SnmpSave")==0)
@@ -4912,34 +4956,30 @@ char* service_mqtt_ack_to_json(mqtt_info_t* info)
 
     cJSON* data_filed = cJSON_CreateObject();
 
-#if 0
-    cJSON* array=NULL;
-    for(int i=0; i<MQTT_SERV_MAX; i++) {
-
-            cJSON* tmp = cJSON_CreateObject();
+#if 1
+    cJSON *tmp,*array=cJSON_CreateArray();
+    if(array) {
+        
+        for(int i=0; i<MQTT_SER_MAX; i++) {
+            mqtt_server_t *ser=&info->ser[i];
+            cJSON *tmp=cJSON_CreateObject();
             if(tmp) {
-                cJSON_AddStringToObject(data_filed,"id", ser->id);
+                sprintf(buf, "%d", ser->id);
+                cJSON_AddStringToObject(tmp,"idx", buf);
 
-                sprintf(buf, "%d", info->mode);
-                cJSON_AddStringToObject(data_filed,"mode", buf);
-                cJSON_AddStringToObject(data_filed,"serverName", ser->server);
-                cJSON_AddStringToObject(data_filed,"ip", ser->ip);
-                cJSON_AddStringToObject(data_filed,"port", ser->port);
-                cJSON_AddStringToObject(data_filed,"zh", ser->user);
-                cJSON_AddStringToObject(data_filed,"mm", ser->password);
-                cJSON_AddStringToObject(data_filed,"rzfs", ser->cid);
+                sprintf(buf, "%d", ser->mode);
+                cJSON_AddStringToObject(tmp,"mode", buf);
 
+                cJSON_AddStringToObject(tmp,"serverName", ser->server);
+                cJSON_AddStringToObject(tmp,"ip", ser->ip);
+                cJSON_AddStringToObject(tmp,"port", ser->port);
+                cJSON_AddStringToObject(tmp,"yhm", ser->user);
+                cJSON_AddStringToObject(tmp,"mm", ser->password);
+                cJSON_AddStringToObject(tmp,"kfdId", ser->cid);
 
-                if(array==NULL) {
-                    array = cJSON_CreateArray();
-                }
-
-                if(array) {
-                    cJSON_AddItemToArray(array, tmp);
-                }
+                cJSON_AddItemToArray(array, tmp);
             }
-    }
-    if(array) {
+        }
         cJSON_AddItemToObject(data_filed, "data", array);
     }
 #endif
@@ -4971,7 +5011,6 @@ char* service_group_ack_to_json(GroupInfo* _globalgroupInfo)
     cJSON_AddNumberToObject(root, "code", 200); 
     //创建数据
     cJSON* data_root_array = cJSON_CreateArray();   
-
     
     //循环添加数据
     list_for_each_entry(_GroupInfoTemp, &_globalgroupInfo->list, list)

+ 13 - 0
pro/src/json_handle.h

@@ -353,6 +353,18 @@ typedef struct
     char* recvAccount[SMTP_RECV_MAX];
 }_SMTP_ServerRequestInfo;
 
+typedef struct 
+{
+    int   idx;
+    int   mode;
+    char* server;
+    char* ip;
+    char* port;
+    char* cid;
+    char* user;
+    char* password;
+}_MQTT_ServerRequestInfo;
+
 /// @brief SNMP信息
 typedef struct
 {
@@ -425,6 +437,7 @@ typedef struct
 {
     char* user_token;
     char* type;
+    _MQTT_ServerRequestInfo _mqtt_info;
     _SMTP_ServerRequestInfo _smtp_info ;
     _SNMP_ServerRequestInfo _snmp_info ;
     _SNMPTRAP_ServerRequestInfo _snmp_trap_info;

+ 5 - 1
pro/src/lock.h

@@ -7,7 +7,11 @@ enum {
     LOCK_ID_MQTT,
     LOCK_ID_SMTP,
     LOCK_ID_PARAS,
-    
+    LOCK_ID_SENSOR,
+    LOCK_ID_BREAKER,
+    LOCK_ID_NETSWITCH,
+    LOCK_ID_POWER_TIMER,
+    LOCK_ID_POWER_UPDATE,
     
     LOCK_ID_MAX
 };

+ 374 - 111
pro/src/mqtt.c

@@ -6,7 +6,9 @@
 #include "mongoose.h"
 #include "elog.h"
 #include "lock.h"
+#include "sys.h"
 #include "cfg.h"
+#include "xlist.h"
 
 #if 0
     #define LOGD            log_d
@@ -23,8 +25,8 @@
 
 char *topic_sub[MQTT_SUB_MAX]={
     "/pdu/%d/control/power/all_channel",
-    "/pdu/%d/control/power/%d",
-    "/pdu/%d/control/sensor/%d",
+    "/pdu/%d/control/power/+",
+    "/pdu/%d/control/sensor/+",
     "/pdu/%d/control/device/restart",
     "/pdu/%d/control/device/reset",
     "/PDU/%d/control/service",
@@ -45,9 +47,14 @@ typedef struct mg_mgr mgr_t;
 typedef struct mg_mqtt_opts mg_opts_t;
 typedef struct mg_connection mg_conn_t;
 
+typedef struct {
+    char topic[256];
+    char content[2048];
+}mqtt_pkt_t;
 typedef struct _mqtt_conn_t{
     mg_conn_t       *c;
     mqtt_server_t   ser;
+    mg_opts_t       opts;
 }mqtt_conn_t;
 typedef struct {
     int           inited;
@@ -57,6 +64,8 @@ typedef struct {
     mgr_t         mgr;
     mg_opts_t     opts;
     mqtt_conn_t   conn[MQTT_SER_MAX];
+
+    handle_t      list;
 }mqtt_handle_t;
 static mqtt_handle_t mqHandle={0};
 static int my_recv(char *topic, char *data);
@@ -88,6 +97,12 @@ static int get_id(mqtt_handle_t *h, mqtt_server_t *ser)
     }
     return -1;
 }
+static char *get_str(int n)
+{
+    static char tmp[32];
+    snprintf(tmp, sizeof(tmp), "%d", n);
+    return tmp;
+}
 static mqtt_conn_t* get_conn(mqtt_handle_t *h, mg_conn_t *c)
 {
     int i;
@@ -114,7 +129,7 @@ static int my_sub(mg_conn_t *c, int prod_id)
     int i;
     char topic[1024];
     
-    for(i=0; topic_sub[i]; i++) {
+    for(i=0; i<MQTT_SUB_MAX; i++) {
         snprintf(topic, sizeof(topic), topic_sub[i], prod_id);
         sub_one(c, topic);
     }
@@ -157,18 +172,39 @@ static int my_disconn(mqtt_handle_t *h)
 
     return 0;
 }
+static char *get_url(mqtt_server_t *ser)
+{
+    char *url=ser->server[0]?ser->server:(ser->ip[0]?ser->ip:NULL);
+    return url;
+}
 static int my_conn(mqtt_handle_t *h, mqtt_info_t *info)
 {
     int i,j,r=-1;
+    char *url=NULL;
     mqtt_server_t *ser=info->ser;
 
     for(i=0; i<MQTT_SER_MAX; i++) {
-        for(j=0; j<MQTT_SER_MAX; j++) {
-            if(h->conn[j].c==NULL) {
-                h->conn[j].c = mg_mqtt_connect(&h->mgr, ser[i].server, &h->opts, mqtt_fn, &h->conn[j]);
-                if(h->conn[j].c) {
-                    h->conn[j].ser = ser[i];
-                    my_sub(h->conn[j].c, h->prod_id);
+        url = get_url(&ser[i]);
+        if(url) {
+            if(h->conn[i].c==NULL) {
+                mg_opts_t opts={
+                    .clean = true,
+                    .qos = 1,
+                    .version = 4,
+                    .keepalive = 60,
+                    //.topic = mg_str("hello"),
+                    //.message = mg_str("bye"),
+                    .client_id = mg_str(ser[i].cid),
+                    .user = mg_str(ser[i].user),
+                    .pass = mg_str(ser[i].password),
+                };
+                char *url=get_url(&ser[i]);
+
+                h->conn[i].c = mg_mqtt_connect(&h->mgr, url, &opts, mqtt_fn, &h->conn[i]);
+                if(h->conn[i].c) {
+                    h->conn[i].ser = ser[i];
+                    h->conn[i].opts = opts;
+                    my_sub(h->conn[i].c, h->prod_id);
                 }
             }
         }
@@ -179,20 +215,33 @@ static int my_conn(mqtt_handle_t *h, mqtt_info_t *info)
 static int my_check(mqtt_handle_t *h)
 {
     int i,r=-1;
+    char *url=NULL;
     mqtt_server_t *ser=h->dm->mqttInfo.ser;
 
     for(i=0; i<MQTT_SER_MAX; i++) {
-        if(h->conn[i].ser.server[0] && h->conn[i].c==NULL) {
-            h->conn[i].c = mg_mqtt_connect(&h->mgr, h->conn[i].ser.server, &h->opts, mqtt_fn, &h->conn[i]);
+        url = get_url(&h->conn[i].ser);
+        if(h->conn[i].c==NULL && url) {
+            h->conn[i].c = mg_mqtt_connect(&h->mgr, url, &h->conn[i].opts, mqtt_fn, &h->conn[i]);
+            if(h->conn[i].c) {
+                my_sub(h->conn[i].c, h->prod_id);
+            }
         }
     }
 
     return 0;
 }
-static void my_poll(mqtt_handle_t *h, int ms)
+static int my_send(mqtt_handle_t *h)
 {
-    my_check(h);
-    mg_mgr_poll(&h->mgr, ms);
+    int r=-1;
+    list_node_t *ln=NULL;
+
+    r = xlist_get_node(h->list, &ln, 0);
+    if(r==0) {
+        mqtt_pkt_t *pkt=(mqtt_pkt_t*)ln->data.buf;
+        my_pub(h, pkt->topic, pkt->content);
+    }
+
+    return r;
 }
 
 static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
@@ -200,37 +249,64 @@ static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
     mqtt_handle_t *h=&mqHandle;
     mqtt_conn_t *mc=get_conn(h,c);
     
-    if (ev == MG_EV_OPEN) {
-        // c->is_hexdumping = 1;
-    } else if (ev == MG_EV_CONNECT) {
-        if (mg_url_is_ssl(mc->ser.server)) {
-            struct mg_tls_opts opts = {.ca = mg_unpacked("/certs/ca.pem"),
-                                       .name = mg_url_host(mc->ser.server)};
-            mg_tls_init(c, &opts);
+    switch(ev) {
+        case MG_EV_OPEN:
+        {
+            // c->is_hexdumping = 1;
         }
-    } else if (ev == MG_EV_ERROR) {
-        // On error, log error message
-        MG_ERROR(("%p %s", c->fd, (char *) ev_data));
-    } else if (ev == MG_EV_MQTT_OPEN) {
-        mg_opts_t opts={
-            .client_id = mg_str(mc->ser.cid),
-            .user = mg_str(mc->ser.user),
-            .pass = mg_str(mc->ser.password),
-        };
-        size_t len=c->send.len;
-
-        mg_mqtt_login(c, &opts);
-        mg_ws_wrap(c, c->send.len - len, WEBSOCKET_OP_BINARY);
-    } else if (ev == MG_EV_MQTT_MSG) {
-        // When we receive MQTT message, print it
-        struct mg_mqtt_message *mm = (struct mg_mqtt_message *) ev_data;
-        //MG_INFO(("Received on %.*s : %.*s", (int) mm->topic.len, mm->topic.buf, (int) mm->data.len, mm->data.buf));
-        my_recv(mm->topic.buf, mm->data.buf);
-    }
+        break;
 
-    if (ev == MG_EV_ERROR || ev == MG_EV_CLOSE) {
-        MG_INFO(("got event %d, stopping...", ev));
-        ((mqtt_conn_t *)(c->fn_data))->c = NULL;  // Signal that we're done
+        case MG_EV_CONNECT:
+        {
+            char path[200];
+            sys_get_path(path, "emqxsl-ca.crt");
+            if (mg_url_is_ssl(mc->ser.server)) {
+                struct mg_tls_opts opts = {.ca = mg_str(path),
+                                           //.cert = mg_str(path),
+                                           //.key = mg_str(path),
+                                           .name = mg_url_host(mc->ser.server)};
+                mg_tls_init(c, &opts);
+            }
+        }
+        break;
+
+        case  MG_EV_MQTT_OPEN:
+        {
+            
+        }
+        break;
+
+        case  MG_EV_POLL:
+        {
+            my_send(h);
+        }
+        break;
+        
+        case MG_EV_MQTT_MSG:
+        {
+            if(ev_data) {
+                // When we receive MQTT message, print it
+                struct mg_mqtt_message *mm = (struct mg_mqtt_message *) ev_data;
+                if(mm->topic.buf && mm->data.buf) {
+                    MG_INFO(("Received on %.*s : %.*s", (int) mm->topic.len, mm->topic.buf, (int) mm->data.len, 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));
+        }
+        break;
+
+        case MG_EV_CLOSE:
+        {
+            LOGD("____ MG_EV_CLOSE\n");
+            ((mqtt_conn_t *)(c->fn_data))->c = NULL;
+        }
+        break;
     }
 }
 static void* mqtt_thread(void *arg)
@@ -238,14 +314,15 @@ static void* mqtt_thread(void *arg)
     int r;
     thread_handle_t *th=(thread_handle_t*)arg;
     mqtt_handle_t   *h=(mqtt_handle_t*)th->arg;
-    
+    mqtt_info_t     *info=&h->dm->mqttInfo;
+
+    my_conn(h, info);
     while(th->quit==0) {
-        if(h->inited) {
-            lock_s_hold(LOCK_ID_MQTT);
-            my_poll(h, 300);
-            lock_s_release(LOCK_ID_MQTT);
-        }
+        my_check(h);
+        mg_mgr_poll(&h->mgr, 500);
     }
+    my_disconn(h);
+
     pthread_exit(NULL);
 }
 #endif
@@ -264,18 +341,22 @@ int mqtt_init(void)
     
     h->dm = &__globalDeviceManage;
     h->prod_id = h->dm->_globalDevInfo.product_id;
+    h->mgr.dns4.url = "udp://114.114.114.114:53";
+    //h->mgr.dns6.url = "udp://114.114.114.114:53";
 
-    h->opts.clean = true,
-    h->opts.qos = 1,
-    h->opts.version = 4,
-    h->opts.topic = mg_str(""),
-    h->opts.message = mg_str("bye");
+    list_cfg_t lc;
+    lc.mode = LIST_FULL_FIFO;
+    lc.max = 100;
+    lc.log = 0;
+    h->list = xlist_init(&lc);
 
     thread_start(THREAD_ID_MQTT, mqtt_thread, h, 4*MB, 0);
     h->inited = 1;
     r = 0;
 #endif
 
+    mqtt_test();
+
     return r;
 }
 
@@ -289,6 +370,7 @@ int mqtt_deinit(void)
     
     thread_stop(THREAD_ID_MQTT);
     mg_mgr_free(&h->mgr);
+    xlist_free(h->list);
     h->inited = 0;
     r = 0;
 #endif
@@ -297,54 +379,16 @@ int mqtt_deinit(void)
 }
 
 
-int mqtt_conn(void)
-{
-    int r=0;
-
-#ifdef USE_MQTT
-    mqtt_handle_t *h=&mqHandle;
-    
-    lock_s_hold(LOCK_ID_MQTT);
-    if(h->inited) {
-        r = my_conn(h, &h->dm->mqttInfo);
-    }
-    lock_s_release(LOCK_ID_MQTT);
-#endif
-
-    return r;
-}
-
-
-int mqtt_disconn(void)
-{
-    int r=-1;
-
-#ifdef USE_MQTT
-    mqtt_handle_t *h=&mqHandle;
-
-    lock_s_hold(LOCK_ID_MQTT);
-    if(h->inited) {
-        r = my_disconn(h);
-    }
-    lock_s_release(LOCK_ID_MQTT);
-#endif
-
-    return r;
-}
-
-
-int mqtt_send(int type, int id, void *data)
+int mqtt_post(int type, int id, void *data)
 {
     int r=-1;
-    int datalen=0;
-
 
 #ifdef USE_MQTT
     char topic[512];
-    char content[4096];
+    char* content=NULL;
     mqtt_handle_t *h=&mqHandle;
 
-    if(!h->inited || type<0 || type>=MQTT_PUB_MAX || !data) {
+    if(!h->inited || type<0 || type>=MQTT_PUB_MAX) {
         return -1;
     }
 
@@ -358,49 +402,216 @@ int mqtt_send(int type, int id, void *data)
     switch(type) {
         case MQTT_PUB_INFO_DEVICE:
         {
-            //
+            cJSON* root=cJSON_CreateObject();
+            if(root) {
+                cJSON_AddStringToObject(root,"Topic", topic);
+                cJSON* data=cJSON_CreateObject();
+                if(data) {
+                    cJSON_AddStringToObject(data,"id", get_str(id));
+                    cJSON_AddStringToObject(data,"name", "Gowone smartPDU");
+                    cJSON_AddStringToObject(data,"type", "smartPDU AC");
+                    cJSON_AddStringToObject(data,"firmware", VERSION);
+                    cJSON_AddStringToObject(data,"channel_num", "8");
+                    cJSON_AddStringToObject(data,"sensor_num", "10");
+
+                    cJSON_AddItemToObject(root, "Data", data);
+                }
+
+                content = cJSON_Print(root);    
+                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) {
+                cJSON_AddStringToObject(root,"Topic", topic);
+                cJSON* data=cJSON_CreateObject();
+                if(data) {
+                    if(r4==0) {
+                        cJSON_AddStringToObject(data,"lan_ipv4_link", nw4.mode?"1":"0");
+                        cJSON_AddStringToObject(data,"lan_ipv4_address", nw4.ip_address);
+                        cJSON_AddStringToObject(data,"lan_ipv4_mask", nw4.mask);
+                    }
+                    else {
+                        cJSON_AddStringToObject(data,"lan_ipv4_link", "");
+                        cJSON_AddStringToObject(data,"lan_ipv4_address", "");
+                        cJSON_AddStringToObject(data,"lan_ipv4_mask", "");
+                    }
+
+                    if(r4==0) {
+                        cJSON_AddStringToObject(data,"lan_ipv6_link", nw6.mode?"1":"0");
+                        cJSON_AddStringToObject(data,"lan_ipv6_address", nw6.ip_address);
+                        cJSON_AddStringToObject(data,"lan_ipv6_subnet_length", nw6.mask);
+                    }
+                    else {
+                        cJSON_AddStringToObject(data,"lan_ipv6_link", "");
+                        cJSON_AddStringToObject(data,"lan_ipv6_address", "");
+                        cJSON_AddStringToObject(data,"lan_ipv6_subnet_length", "");
+                    }
+
+                    if(0) {
+                        cJSON_AddStringToObject(data,"wifi_ipv4_link", "");
+                        cJSON_AddStringToObject(data,"wifi_ipv4_address", "");
+                        cJSON_AddStringToObject(data,"wifi_ipv4_mask", "");
+
+                        cJSON_AddStringToObject(data,"wifi_ipv6_link", "");
+                        cJSON_AddStringToObject(data,"wifi_ipv6_address", "");
+                        cJSON_AddStringToObject(data,"wifi_ipv6_mask", "");
+                    }
+
+                    cJSON_AddStringToObject(data,"modbus_address", get_str(h->dm->_globalDevInfo._gmodbus_info.product_modbus_addr));
+                    cJSON_AddStringToObject(data,"modbus_baud", get_str(h->dm->_globalDevInfo._gmodbus_info.product_modbus_baud));
+                    cJSON_AddStringToObject(data,"modbus_mode", get_str(h->dm->_globalDevInfo._gmodbus_info.product_modbus_type));
+                    cJSON_AddStringToObject(data,"vpn_enable", "0");
+
+                    cJSON_AddItemToObject(root, "Data", data);
+                }
 
+                content = cJSON_Print(root);    
+                cJSON_Delete(root);
+            }
         }
         break;
 
         case MQTT_PUB_STAT_POWER_ALL:
         {
-            
+            cJSON* root=cJSON_CreateObject();
+            if(root) {
+                cJSON_AddStringToObject(root,"Topic", topic);
+                cJSON* data=cJSON_CreateObject();
+                if(data) {
+                    cJSON_AddStringToObject(data,"voltage", "");
+                    cJSON_AddStringToObject(data,"current", "");
+                    cJSON_AddStringToObject(data,"power", "");
+                    cJSON_AddStringToObject(data,"consumption", "");
+                    cJSON_AddStringToObject(data,"pactive_power", "");
+                    cJSON_AddStringToObject(data,"reactive_power", "");
+                    cJSON_AddStringToObject(data,"apparent_power", "");
+                    cJSON_AddStringToObject(data,"date", "");
+                    cJSON_AddStringToObject(data,"time", "");
+
+                    cJSON_AddItemToObject(root, "Data", data);
+                }
+
+                content = cJSON_Print(root);    
+                cJSON_Delete(root);
+            }
         }
         break;
 
         case MQTT_PUB_STAT_POWER_CHN:
         {
-            
+            cJSON* root=cJSON_CreateObject();
+            if(root) {
+                cJSON_AddStringToObject(root,"Topic", topic);
+                cJSON* data=cJSON_CreateObject();
+                if(data) {
+                    cJSON_AddStringToObject(data,"name", "app server channel");
+                    cJSON_AddStringToObject(data,"status", "1");
+                    cJSON_AddStringToObject(data,"voltage", "");
+                    cJSON_AddStringToObject(data,"current", "");
+                    cJSON_AddStringToObject(data,"power", "");
+                    cJSON_AddStringToObject(data,"power_freq", "");
+                    cJSON_AddStringToObject(data,"consumption", "");
+                    cJSON_AddStringToObject(data,"power_factor", "");
+                    cJSON_AddStringToObject(data,"pactive_power", "");
+                    cJSON_AddStringToObject(data,"reactive_power", "");
+                    cJSON_AddStringToObject(data,"apparent_power", "");
+                    cJSON_AddStringToObject(data,"date", "");
+                    cJSON_AddStringToObject(data,"time", "");
+
+                    cJSON_AddItemToObject(root, "Data", data);
+                }
+
+                content = cJSON_Print(root);    
+                cJSON_Delete(root);
+            }
         }
         break;
 
         case MQTT_PUB_STAT_SENSOR:
         {
-            
-        }
-        break;
+            cJSON* root=cJSON_CreateObject();
+            if(root) {
+                cJSON_AddStringToObject(root,"Topic", topic);
+                cJSON* data=cJSON_CreateObject();
+                if(data) {
+                    cJSON_AddStringToObject(data,"name", "temp sensor");
+                    cJSON_AddStringToObject(data,"type", "0");
+                    cJSON_AddStringToObject(data,"modbus_address", "");
+                    cJSON_AddStringToObject(data,"node", "-1");
+                    cJSON_AddStringToObject(data,"status", "1");
+                    cJSON_AddStringToObject(data,"value_num", "");
+                    cJSON_AddStringToObject(data,"value1", "");
+                    cJSON_AddStringToObject(data,"value2", "");
+                    cJSON_AddStringToObject(data,"date", "");
+                    cJSON_AddStringToObject(data,"time", "");
+
+                    cJSON_AddItemToObject(root, "Data", data);
+                }
 
-        case MQTT_PUB_ALARM_NETWORK:
-        {
-            
+                content = cJSON_Print(root);    
+                cJSON_Delete(root);
+            }
         }
         break;
 
-        case MQTT_PUB_ALARM_POWER:
+        case MQTT_PUB_STAT_SERVICE:
         {
-            
+            cJSON* root=cJSON_CreateObject();
+            if(root) {
+                cJSON_AddStringToObject(root,"Topic", topic);
+                cJSON* data=cJSON_CreateObject();
+                if(data) {
+                    cJSON_AddStringToObject(data,"telnet", "0");
+                    cJSON_AddStringToObject(data,"smtp", "1");
+                    cJSON_AddStringToObject(data,"snmp_v1", "1");
+                    cJSON_AddStringToObject(data,"snmp_v2c", "1");
+                    cJSON_AddStringToObject(data,"snmp_v3", "1");
+                    cJSON_AddStringToObject(data,"snmp_trap", "1");
+                    cJSON_AddStringToObject(data,"message", "");
+                    cJSON_AddStringToObject(data,"mqtt", "1");
+                    cJSON_AddStringToObject(data,"cloud", "1");
+                    cJSON_AddStringToObject(data,"ntp", "1");
+
+                    cJSON_AddItemToObject(root, "Data", data);
+                }
+
+                content = cJSON_Print(root);
+                cJSON_Delete(root);
+            }
         }
         break;
 
+        case MQTT_PUB_ALARM_NETWORK:
         case MQTT_PUB_ALARM_SENSOR:
+        case MQTT_PUB_ALARM_POWER:
         {
-            
+            cJSON* root=cJSON_CreateObject();
+            if(root) {
+                cJSON_AddStringToObject(root,"Topic", topic);
+                cJSON* data=cJSON_CreateObject();
+                if(data) {
+                    cJSON_AddStringToObject(data,"num", "temp sensor");
+                    cJSON_AddStringToObject(data,"context", "");
+                    cJSON_AddStringToObject(data,"action", "1");
+                    cJSON_AddStringToObject(data,"action_para", "7");
+                    cJSON_AddStringToObject(data,"date", "");
+                    cJSON_AddStringToObject(data,"time", "");
+
+                    cJSON_AddItemToObject(root, "Data", data);
+                }
+
+                content = cJSON_Print(root);    
+                cJSON_Delete(root);
+            }
         }
         break;
 
@@ -408,9 +619,11 @@ int mqtt_send(int type, int id, void *data)
         return -1;
     }
 
-    lock_s_hold(LOCK_ID_MQTT);
-    r = my_pub(h, topic, content);
-    lock_s_release(LOCK_ID_MQTT);
+    mqtt_pkt_t pkt;
+    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
 
     return r;
@@ -487,7 +700,7 @@ static int my_recv(char *topic, char *data)
     int i,r,cmd;
 
 #ifdef USE_MQTT
-    char temp[2000];
+    char temp[512];
     mqtt_handle_t *h=&mqHandle;
 
     for(i=0; i<MQTT_SUB_MAX; i++) {
@@ -548,3 +761,53 @@ static int my_recv(char *topic, char *data)
     return 0;
 }
 
+
+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,
+#if 0
+            .server = "mqtts://gaae4d7b.ala.cn-hangzhou.emqxsl.cn:8883",
+            .ip = "",
+            .port = "8883",
+            .cid = "mqtlx_e92334545",
+            .user = "AABB123",
+            .password = "asdfasd3225",
+#else
+            .server = "192.168.1.12:1883",
+            .ip = "",
+            .port = "1883",
+            .cid = "mqtlx_asdjflsd",
+
+            .user = "gowone100",
+            .password = "gowone100",
+
+            //.user = "gowone101",
+            //.password = "gowone101",
+#endif
+
+            //.user = "adsg33434",
+            //.password = "aafweertzv",
+        }; 
+
+    mInfo.ser[0] = ser;
+
+    h->dm->mqttInfo = mInfo;
+    //mg_log_set(MG_LL_VERBOSE);
+
+    while(1) {
+        mqtt_post(MQTT_PUB_INFO_DEVICE, 0, 0);
+        sleep(2);
+    }
+
+
+
+#endif
+    return 0;
+}
+

+ 2 - 4
pro/src/mqtt.h

@@ -30,9 +30,7 @@ enum {
 int mqtt_init(void);
 int mqtt_deinit(void);
 
-int mqtt_conn(void);
-int mqtt_disconn(void);
-
-int mqtt_send(int type, int id, void *data);
+int mqtt_post(int type, int id, void *data);
+int mqtt_test(void);
 
 #endif

+ 9 - 0
pro/src/thread.c

@@ -39,6 +39,15 @@ int thread_start(int id, thread_fn fn, void *arg, int stksize, int prio)
 #endif
     
     r = pthread_create(&h->tid, &attr, fn, h);
+#if 0
+    if(r==0) {
+        r = pthread_setschedprio(h->tid, prio);
+        if(r) {
+            printf("____ set thread prio %d failed\n", prio);
+        }
+    }
+#endif
+
     pthread_attr_destroy(&attr);
 
     return r;

+ 12 - 10
pro/src/websocket_handle.c

@@ -4,13 +4,14 @@
 #include "elog.h"
 #include "common.h"
 #include "sys.h"
+#include "lock.h"
 #include "thread.h"
 #include "sqlite_handle.h"
 #include "json_handle.h"
 #include "cascade.h"
 #include "cfg.h"
 
-//#define USE_WS2
+#define USE_WS2
 #define WS_SEND_TLEN_MAX  (1024*1024*10)
 
 #define WS_MAX    10
@@ -115,11 +116,11 @@ static void web_update(ws_handle_t *wh)
             _OverAllPwrAckInfo l1, l2, l3;
             if (cur_dev_addr == 0)
             {
-                pthread_mutex_lock(&gdm->_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);
-                pthread_mutex_unlock(&gdm->_power_update);
+                lock_s_release(LOCK_ID_POWER_UPDATE);
             }
             else
             {
@@ -141,9 +142,9 @@ static void web_update(ws_handle_t *wh)
         {
             int nRetTotal = -1;
             if(cur_dev_addr==0) {
-                pthread_mutex_lock(&gdm->_power_update); 
+                lock_s_hold(LOCK_ID_POWER_UPDATE);
                 nRetTotal=dev_search_latest_power_statistic_info(gdm, &allInfo);
-                pthread_mutex_unlock(&gdm->_power_update);
+                lock_s_release(LOCK_ID_POWER_UPDATE);
             }
             else {
                 cascade_lock();
@@ -164,9 +165,9 @@ static void web_update(ws_handle_t *wh)
              int nRetCh = -1;
 
             if(cur_dev_addr==0) {
-                pthread_mutex_lock(&gdm->_power_update); 
+                lock_s_hold(LOCK_ID_POWER_UPDATE);
                 nRetCh=dev_search_latest_power_All_info(gdm,&chInfo, cur_dev_addr);
-                pthread_mutex_unlock(&gdm->_power_update);
+                lock_s_release(LOCK_ID_POWER_UPDATE);
                 pdm = gdm;
             }
             else {
@@ -191,9 +192,10 @@ static void web_update(ws_handle_t *wh)
 
         //sensor info
         {
-            //pthread_mutex_lock(&__globalDeviceManage._sensor_update);
+            //lock_s_hold(LOCK_ID_SENSOR);
             json_str = ws_sensor_status_ack_to_json(&gdm->_globalSensorManger);
-           // pthread_mutex_unlock(&__globalDeviceManage._sensor_update);
+            //lock_s_release(LOCK_ID_SENSOR);
+
             ws_send_data(wh, json_str, strlen(json_str));
             cJSON_free(json_str);
         }
@@ -273,7 +275,7 @@ void* websocket_thread(void* arg)
     ws_handle_t *wh=(ws_handle_t*)h->arg;
     const char *listen_addr="ws://[::]:6785";
 
-    mg_log_set(MG_LL_DEBUG);
+    //mg_log_set(MG_LL_DEBUG);
     mg_mgr_init(&wh->mgr);  // Initialise event manager
     mg_http_listen(&wh->mgr, listen_addr, fn, NULL);  // Create HTTP listener
     wh->inited = 1;

+ 1 - 1
pro/src/xlist.c

@@ -36,7 +36,7 @@ static void* lock_init(void)
 {
     int r;
     
-    pthread_mutex_t *m=MALLOC(sizeof(pthread_mutex_t));
+    pthread_mutex_t *m=(pthread_mutex_t*)MALLOC(sizeof(pthread_mutex_t));
     if(!m) {
         return NULL;
     }