瀏覽代碼

简化、统一线程接口
解决mqtt部分bug

guohui 2 年之前
父節點
當前提交
721a38d1ae
共有 12 個文件被更改,包括 331 次插入 和 224 次删除
  1. 9 10
      pro/src/app.c
  2. 2 2
      pro/src/appweb_handle.c
  3. 5 5
      pro/src/cascade.c
  4. 17 4
      pro/src/language_hashmap.c
  5. 222 176
      pro/src/mqtt.c
  6. 6 4
      pro/src/netswitch.c
  7. 7 3
      pro/src/smtp.c
  8. 2 2
      pro/src/sqlite_handle.c
  9. 2 2
      pro/src/sys.c
  10. 34 11
      pro/src/thread.c
  11. 23 3
      pro/src/thread.h
  12. 2 2
      pro/src/websocket_handle.c

+ 9 - 10
pro/src/app.c

@@ -69,8 +69,6 @@ static void db_ow_test(sqlite3 *db, int *index, int *index3)
 }
 
 
-
-
 void* power_thread(void* arg)
 {
     int ret = 0 ;
@@ -1729,6 +1727,8 @@ int _global_device_manage_init(GlobalDeviceManager* _globalDeviceManager)
     }
     log_i("database %s open ok.\n", path);
 
+
+#if 1
     if(_globalDeviceManager->_globalDevInfo.product.pwr_type == SmartPDU_AC ||
         _globalDeviceManager->_globalDevInfo.product.pwr_type == SmartPDU_DC)       // single AC
     {
@@ -1922,7 +1922,6 @@ int _global_device_manage_init(GlobalDeviceManager* _globalDeviceManager)
     log_d("GCPDU init begin!");
     ResetChmData(0);
 
-#if 1
     ///todo:此信息需要和数据库更新以后在关联
     ret = dev_get_sensor_manage_info(_globalDeviceManager->db,
                                     _globalDeviceManager->_globalDevInfo.product.id,
@@ -1946,12 +1945,12 @@ int _global_device_manage_init(GlobalDeviceManager* _globalDeviceManager)
         log_d("get breaker info ok.");
 
     _globalDeviceManager->dev_samp_flag = true;
-    thread_start(THREAD_ID_POWER, power_thread, _globalDeviceManager, 20*MB, 0);
-    thread_start(THREAD_ID_SENSOR, sensor_thread, _globalDeviceManager, 10*MB, 0);
-    thread_start(THREAD_ID_SHMW, shm_thread, _globalDeviceManager, 10*MB, 0);
+    thread_start(THREAD_ID_POWER, _globalDeviceManager);
+    thread_start(THREAD_ID_SENSOR, _globalDeviceManager);
+    thread_start(THREAD_ID_SHMW, _globalDeviceManager);
 
-    thread_start(THREAD_ID_BREAKER_SCANNER,breaker_scanner_thread,_globalDeviceManager,1*MB,0);
-    thread_start(THREAD_ID_BREAKER, breaker_thread, _globalDeviceManager, 5*MB, 0);
+    thread_start(THREAD_ID_BREAKER_SCANNER,_globalDeviceManager);
+    thread_start(THREAD_ID_BREAKER, _globalDeviceManager);
 
     dev_Alarm_Run_message(_globalDeviceManager,language_alarm_Init_Success[0],"Server Thread");
 
@@ -1983,11 +1982,11 @@ int _global_device_manage_init(GlobalDeviceManager* _globalDeviceManager)
         log_d("get ntp info from database succeeded");
     }
 
-    thread_start(THREAD_ID_NTP, ntp_thread, NULL, 10 * MB, 0);
+    thread_start(THREAD_ID_NTP, NULL);
     dev_Alarm_Run_message(_globalDeviceManager, language_alarm_Init_Success[0], "NTP");
 
 #ifndef USE_NETSWITCH
-    thread_start(THREAD_ID_SNMP, snmp_thread, _globalDeviceManager, 10 * MB, 0);
+    thread_start(THREAD_ID_SNMP, _globalDeviceManager);
 #endif
 
     dev_Alarm_Run_message(_globalDeviceManager, language_alarm_Init_Success[0], "SNMP");

+ 2 - 2
pro/src/appweb_handle.c

@@ -205,7 +205,7 @@ static void http_mg_cb(struct mg_connection *c, int ev, void *ev_data)
         }
     }
 }
-static void *web_thread(void *arg)
+void *web_thread(void *arg)
 {
     struct mg_mgr mgr;
     struct mg_connection *c;
@@ -5237,7 +5237,7 @@ static int webSetCallback(void)
 
 int appweb_init(void)
 {
-    thread_start(THREAD_ID_WEB, web_thread, &__globalDeviceManage, 16*MB, 0);
+    thread_start(THREAD_ID_WEB, &__globalDeviceManage);
     
     return 0;
 }

+ 5 - 5
pro/src/cascade.c

@@ -1544,7 +1544,7 @@ static int slave_receive(cascade_handle_t *cas)
     return mb_receive(cas);
 }
 
-static void* cascade_thread(void *arg)
+void* cascade_thread(void *arg)
 {
     int r;
     thread_handle_t *h=(thread_handle_t*)arg;
@@ -1568,7 +1568,7 @@ static void* cascade_thread(void *arg)
     
     pthread_exit(NULL);
 }
-static void* scan_thread(void *arg)
+void* cascade_scan_thread(void *arg)
 {
     int r;
     thread_handle_t *h=(thread_handle_t*)arg;
@@ -1738,8 +1738,8 @@ int cascade_init(void)
     set_modbus(cas, get_mb());
     slave_add(cas, 0);
 
-    thread_start(THREAD_ID_CASCADE, cascade_thread, cas, 4*MB, 0);
-    thread_start(THREAD_ID_SCAN, scan_thread, cas, 4*MB, 0);
+    thread_start(THREAD_ID_CASCADE, cas);
+    thread_start(THREAD_ID_CASCADE_SCAN, cas);
         
     return 0;
 }
@@ -1750,7 +1750,7 @@ int cascade_deinit(void)
     cascade_handle_t *cas=&casHandle;
 
     thread_stop(THREAD_ID_CASCADE);
-    thread_stop(THREAD_ID_SCAN);
+    thread_stop(THREAD_ID_CASCADE_SCAN);
 
     pthread_mutex_destroy(&cas->mutex);
     pthread_mutex_destroy(&cas->lock);

+ 17 - 4
pro/src/language_hashmap.c

@@ -1,7 +1,20 @@
-#include "language_hashmap.h"
-#include "common.h"
 #include <stdlib.h>
 #include <string.h>
+#include "elog.h"
+#include "common.h"
+#include "language_hashmap.h"
+
+
+#if 1
+    #define LOGD            log_d
+    #define LOGE            log_e
+    #define LOGW            log_w
+#else
+    #define LOGD            printf
+    #define LOGE            printf
+    #define LOGW            printf
+#endif
+
 
 unsigned int _hash(HashTable *table,const char *key)
 {  
@@ -31,7 +44,7 @@ int insert_Newhash(HashTable *table, const char *key, const char value[][64],int
     unsigned int index = _hash(table,key);
     
     if(table->ptable[index]) {
-        printf("hash value is repeat, index: %d, ___%s___, ___%s___\n", index, table->ptable[index]->key, key);
+        LOGW("hash value is repeat, index: %d, ___%s___, ___%s___\n", index, table->ptable[index]->key, key);
         return HASH_REPEAT;
     }
 
@@ -55,7 +68,7 @@ int insert_Newhash(HashTable *table, const char *key, const char value[][64],int
     entry->next = table->ptable[index];  
     table->ptable[index] = entry;
     
-    //printf("___key: %s, index: %d\n", key, index);
+    //LOGD("___key: %s, index: %d\n", key, index);
     return HASH_OK;
 }  
   

+ 222 - 176
pro/src/mqtt.c

@@ -14,7 +14,7 @@
 #include "switch_ctrl.h"
 #include "sqlite_handle.h"
 
-#if 0
+#if 1
     #define LOGD            log_d
     #define LOGE            log_e
     #define LOGW            log_w
@@ -29,8 +29,55 @@
 #define SEND_PERIOD         5000    //ms
 
 
-#ifdef USE_MQTT
+enum {
+    SERV_NTP=0,
+    SERV_SMTP,
+    SERV_MESG,
+    SERV_MQTT,
+    SERV_CLOUD,
+    SERV_TELNET,
+    SERV_SNMP_V1,
+    SERV_SNMP_V2C,
+    SERV_SNMP_V3,
+    SERV_SNMP_TRAP,
+    
+    SERV_MAX
+};
+#define BGET(flag,mask)  ((flag)&(1<<(mask)))
+#define BSET(flag,mask)  ((flag)|=(1<<(mask)))
 
+
+typedef struct mg_mgr mgr_t;
+typedef struct mg_timer mg_timer_t;
+typedef struct mg_mqtt_opts mg_opts_t;
+typedef struct mg_connection mg_conn_t;
+
+typedef struct {
+    int  id;
+    char topic[256];
+    char content[2048];
+}mqtt_pkt_t;
+typedef struct {
+    mgr_t          mgr;
+    mg_conn_t      *c;
+    mqtt_server_t  ser;
+    pthread_t      tid;
+    int            quit;
+
+    void           *h;
+}mqtt_conn_t;
+typedef struct {
+    int           inited;
+    uint32_t      prod_id;
+    
+    mqtt_conn_t   conn[MQTT_SER_MAX];
+    handle_t      list;
+    mg_timer_t    *timer;
+
+    GlobalDeviceManager *dm;
+}mqtt_handle_t;
+
+#ifdef USE_MQTT
 char *topic_sub[MQTT_SUB_MAX]={
     "/pdu/%d/control/power/all_channel",
     "/pdu/%d/control/power/+",
@@ -50,60 +97,25 @@ char *topic_pub[MQTT_PUB_MAX]={
     "/pdu/%d/alarm/sensor",
     "/pdu/%d/status/service",
 };
-
-typedef struct mg_mgr mgr_t;
-typedef struct mg_timer mg_timer_t;
-typedef struct mg_mqtt_opts mg_opts_t;
-typedef struct mg_connection mg_conn_t;
-
-typedef struct {
-    int  id;
-    char topic[256];
-    char content[2048];
-}mqtt_pkt_t;
-typedef struct _mqtt_conn_t{
-    mg_conn_t       *c;
-    mqtt_server_t   ser;
-    //mg_opts_t       opts;
-    uint8_t         flag[MQTT_PUB_MAX];
-}mqtt_conn_t;
-typedef struct {
-    int           inited;
-    int           prod_id;
-    GlobalDeviceManager *dm;
-
-    mgr_t         mgr;
-    mg_opts_t     opts;
-    mqtt_conn_t   conn[MQTT_SER_MAX];
-
-    handle_t      list;
-    mg_timer_t    *timer;
-}mqtt_handle_t;
+const char *serv_str[SERV_MAX]={
+    "ntp",
+    "smtp",
+    "message",
+    "mqtt",
+    "cloud",
+    "telnet",
+    "snmp_v1",
+    "snmp_v2c",
+    "snmp_v3",
+    "snmp_trap",
+};
 static mqtt_handle_t mqHandle={0};
 static int my_recv(char *topic, char *data);
 static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data);
-static void timer_start(mqtt_handle_t *h);
-static int send_stat(mqtt_handle_t *h, int type);
-static uint32_t get_serv(char *json);
-static void send_once(mqtt_handle_t *h);
-
-enum {
-    SERV_NTP=0,
-    SERV_SMTP,
-    SERV_MESG,
-    SERV_MQTT,
-    SERV_CLOUD,
-    SERV_TELNET,
-    SERV_SNMP_V1,
-    SERV_SNMP_V2C,
-    SERV_SNMP_V3,
-    SERV_SNMP_TRAP,
-    
-    SERV_MAX
-};
-#define BGET(flag,mask)  ((flag)&(1<<(mask)))
-#define BSET(flag,mask)  ((flag)|=(1<<(mask)))
-
+static int send_stat(mqtt_conn_t *conn, int type);
+static int get_serv(char *json);
+static void* client_thread(void *arg);
+/////////////////////////////////////////////////////////
 
 static char *get_int_str(int n)
 {
@@ -186,32 +198,14 @@ static int is_period(int id)
     }
     return 1;
 }
-static int my_pub(mqtt_handle_t *h, int id, char *topic, char *data)
+static int my_pub(mqtt_conn_t *conn, char *topic, char *data)
 {
-    int i;
-    for(i=0; i<MQTT_SER_MAX; i++) {
-        if(h->conn[i].c) {
-            if(is_period(id) || !h->conn[i].flag[id]) {
-                pub_one(h->conn[i].c, topic, data);
-                h->conn[i].flag[id] = 1;
-            }
-        }
+    if(conn->c) {
+        pub_one(conn->c, topic, data);
     }
     return 0;
 }
-static int my_disconn(mqtt_handle_t *h)
-{
-    int i,r=-1;
 
-    for(i=0; i<MQTT_SER_MAX; i++) {
-        if(h->conn[i].c) {
-            mg_mqtt_disconnect(h->conn[i].c, &h->opts);
-            h->conn[i].c = NULL;
-        }
-    }
-
-    return 0;
-}
 static GlobalPowerManger* get_power(mqtt_handle_t *h, int ch)
 {
     GlobalPowerManger *tmp=NULL;
@@ -226,12 +220,14 @@ static GlobalPowerManger* get_power(mqtt_handle_t *h, int ch)
 }
 static int set_ch(mqtt_handle_t *h, int ch, int flag)
 {
-    int r,sendAddr=0;
+    int r;
     GlobalPowerManger *tmp=NULL;
+
     list_for_each_entry(tmp, &h->dm->_globalPowerManger.list, list)
     {
-        if(tmp->product_saddr==sendAddr) continue;
-        sendAddr = tmp->product_saddr;
+        if(tmp==NULL) {
+            break;
+        }
 
         if(tmp->product_ch_id==ch) {
             r = g_switch_set_all_chn_ctrl(&h->dm->_globalRelaySampManger, tmp, tmp->product_saddr, tmp->product_ch_addr, flag, false);
@@ -245,6 +241,7 @@ static int set_ch(mqtt_handle_t *h, int ch, int flag)
             g_switch_set_all_ctrl(&h->dm->_globalRelaySampManger, tmp->product_ch_type, tmp->product_saddr, flag);
         }
     }
+
     return 0;
 }
 static GlobalSensorManger* get_sensor(mqtt_handle_t *h, int id)
@@ -290,7 +287,66 @@ static int get_url(mqtt_server_t *ser, char *url)
     sprintf(url, "%s://%s:%d", head, ser->server, port);
     return 0;
 }
-static int my_conn(mqtt_handle_t *h)
+
+static int start_one(mqtt_conn_t *conn)
+{
+    int r;
+    if(conn->tid==0) {
+        conn->quit = 0;
+        r = pthread_create(&conn->tid, NULL, client_thread, conn);
+    }
+    return r;
+}
+static int stop_one(mqtt_conn_t *conn)
+{
+    if(conn->tid) {
+        conn->quit = 1;
+        pthread_join(conn->tid, NULL);
+        conn->tid = 0;
+    }
+    return 0;
+}
+static int stop_all(mqtt_handle_t *h)
+{
+    int i;
+    for(i=0; i<MQTT_SER_MAX; i++) {
+        stop_one(&h->conn[i]);
+    }
+    return 0;
+}
+
+static int conn_one(mqtt_conn_t *conn)
+{
+    int r=-1;
+    mqtt_handle_t *h=&mqHandle;
+
+    if(!conn->c) {
+        char url[1024];
+        mg_opts_t opts={
+            .clean = true,
+            .qos = 1,
+            .version = 4,
+            .keepalive = 60,
+            .topic = mg_str("hello"),
+            .message = mg_str("bye"),
+            .client_id = mg_str(conn->ser.cid),
+            .user = mg_str(conn->ser.user),
+            .pass = mg_str(conn->ser.password),
+        };
+
+        get_url(&conn->ser, url);
+        conn->c = mg_mqtt_connect(&conn->mgr, url, &opts, mqtt_fn, conn);
+        if(conn->c) {
+            sprintf(conn->ser.cid, "smartPDU_%d\n", h->prod_id);
+            my_sub(conn->c, h->prod_id);
+            r = 0;
+        }
+    }
+
+    return r;
+}
+
+static int my_check(mqtt_handle_t *h)
 {
     int i,j,r=-1;
     char url[1024];
@@ -298,24 +354,16 @@ static int my_conn(mqtt_handle_t *h)
 
     for(i=0; i<MQTT_SER_MAX; i++) {
         r = get_url(&ser[i], url);
-        if((r==0) && (!h->conn[i].c)) {
-            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),
-            };
-
-            h->conn[i].c = mg_mqtt_connect(&h->mgr, url, &opts, mqtt_fn, &h->conn[i]);
-            if(h->conn[i].c) {
+        if((r==0)) {
+            if(h->conn[i].tid==0) {
                 h->conn[i].ser = ser[i];
-                sprintf(h->conn[i].ser.cid, "smartPDU_%d\n", h->prod_id);
-                my_sub(h->conn[i].c, h->prod_id);
+                h->conn[i].h = h;
+                start_one(&h->conn[i]);
+            }
+        }
+        else {
+            if(h->conn[i].tid) {
+                stop_one(&h->conn[i]);
             }
         }
     }
@@ -323,53 +371,53 @@ static int my_conn(mqtt_handle_t *h)
     return 0;
 }
 
-static int my_send(mqtt_handle_t *h)
+static int my_send(mqtt_conn_t *conn)
 {
     int r=-1;
+    mqtt_handle_t *h=(mqtt_handle_t*)conn->h;
     list_node_t *ln=NULL;
 
     r = xlist_get_node(h->list, &ln, 0);
     if(r==0) {
         mqtt_pkt_t *pkt=(mqtt_pkt_t*)ln->data.buf;
-        my_pub(h, pkt->id, pkt->topic, pkt->content);
+        my_pub(conn, pkt->topic, pkt->content);
     }
 
     return r;
 }
-static void send_once(mqtt_handle_t *h)
+static int send_once(mqtt_conn_t *conn)
 {
-    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);
+    send_stat(conn, MQTT_PUB_INFO_DEVICE);
+    send_stat(conn, MQTT_PUB_STAT_NETWORK);
+    send_stat(conn, MQTT_PUB_STAT_SENSOR);
+    send_stat(conn, MQTT_PUB_STAT_SERVICE);
+    return 0;
 }
-static void send_period(mqtt_handle_t *h)
+static void send_period(mqtt_conn_t *conn)
 {
-    send_stat(h, MQTT_PUB_STAT_POWER_ALL);
-    send_stat(h, MQTT_PUB_STAT_POWER_CHN);
+    send_stat(conn, MQTT_PUB_STAT_POWER_ALL);
+    send_stat(conn, MQTT_PUB_STAT_POWER_CHN);
 }
 
 static void timer_conn_fn(void *arg)
 {
-    mqtt_handle_t *h=(mqtt_handle_t*)arg;
-    my_conn(h);
+    mqtt_conn_t *conn=(mqtt_conn_t*)arg;
+    conn_one(conn);
 }
 static void timer_period_fn(void *arg)
 {
-    mqtt_handle_t *h=(mqtt_handle_t*)arg;
-    send_period(h);
+    mqtt_conn_t *conn=(mqtt_conn_t*)arg;
+    send_period(conn);
 }
-
-static void timer_start(mqtt_handle_t *h)
+static void timer_add(mqtt_conn_t *conn)
 {
-    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);
+    mg_timer_add(&conn->mgr, CONN_PERIOD, MG_TIMER_REPEAT | MG_TIMER_RUN_NOW, timer_conn_fn, conn);
+    mg_timer_add(&conn->mgr, SEND_PERIOD, MG_TIMER_REPEAT | MG_TIMER_RUN_NOW, timer_period_fn, conn);
 }
 
 static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
 {
-    mqtt_handle_t *h=&mqHandle;
-    mqtt_conn_t *mc=((mqtt_conn_t*)(c->fn_data));
+    mqtt_conn_t *conn=((mqtt_conn_t*)(c->fn_data));
     
     switch(ev) {
         case MG_EV_OPEN:
@@ -381,10 +429,11 @@ static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
         case MG_EV_CONNECT:
         {
             char url[1024];
-            get_url(&mc->ser, url);
+
+            get_url(&conn->ser, url);
             if (mg_url_is_ssl(url)) {
-                struct mg_tls_opts opts = {.ca = mg_str(mc->ser.cert),
-                                           .name = mg_url_host(mc->ser.server)};
+                struct mg_tls_opts opts = {.ca = mg_str(conn->ser.cert),
+                                           .name = mg_url_host(conn->ser.server)};
                 mg_tls_init(c, &opts);
             }
         }
@@ -392,13 +441,13 @@ static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
 
         case  MG_EV_MQTT_OPEN:
         {
-            send_once(h);
+            send_once(conn);
         }
         break;
         
         case MG_EV_POLL:
         {
-            my_send(h);
+            my_send(conn);
         }
         break;
 
@@ -413,32 +462,51 @@ static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
     
         case MG_EV_ERROR:
         {
-            MG_ERROR(("___MG_EV_ERROR, %p %s", c->fd, (char *) ev_data));
-            memset(mc->flag, 0, sizeof(mc->flag));
+            //MG_ERROR(("___MG_EV_ERROR, %p %s", c->fd, (char *) ev_data));
+            //memset(mc->flag, 0, sizeof(mc->flag));
         }
         break;
 
         case MG_EV_CLOSE:
         {
-            LOGD("____ MG_EV_CLOSE\n");
-            mc->c = NULL;
-            memset(mc->flag, 0, sizeof(mc->flag));
+            conn->c = NULL;
         }
         break;
     }
 }
-static void* mqtt_thread(void *arg)
+static void* client_thread(void *arg)
+{
+    int r;
+    mqtt_conn_t *conn=(mqtt_conn_t*)arg;
+    
+    mg_mgr_init(&conn->mgr);
+    conn->mgr.dns4.url = "udp://114.114.114.114:53";
+    //conn->mgr.dns6.url = "udp://114.114.114.114:53";
+
+    timer_add(conn);
+    while(conn->quit==0) {
+        mg_mgr_poll(&conn->mgr, 1000);
+    }
+    mg_mqtt_disconnect(conn->c, NULL);
+    mg_mgr_free(&conn->mgr);
+    conn->c = NULL;
+
+    pthread_exit(NULL);
+}
+
+
+
+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;
+    mqtt_handle_t *h=(mqtt_handle_t*)th->arg;
     
-    timer_start(h);
     while(th->quit==0) {
-        mg_mgr_poll(&h->mgr, 2000);
+        my_check(h);
+        sleep(1);
     }
-    my_disconn(h);
+    stop_all(h);
 
     pthread_exit(NULL);
 }
@@ -454,7 +522,6 @@ int mqtt_init(void)
     mqtt_handle_t *h=&mqHandle;
     
     memset(h, 0, sizeof(mqtt_handle_t));
-    mg_mgr_init(&h->mgr);
 
     //mg_log_set(MG_LL_DEBUG);
     
@@ -462,9 +529,6 @@ int mqtt_init(void)
     h->prod_id = h->dm->_globalDevInfo.product.id;
     //sys_get_chip_id(&h->prod_id);
 
-    h->mgr.dns4.url = "udp://114.114.114.114:53";
-    //h->mgr.dns6.url = "udp://114.114.114.114:53";
-
     list_cfg_t lc;
     lc.mode = LIST_FULL_FIFO;
     lc.max = 100;
@@ -472,8 +536,7 @@ int mqtt_init(void)
     h->list = xlist_init(&lc);
 
     dev_mqtt_init(h->dm->db, &h->dm->mqttInfo);
-
-    thread_start(THREAD_ID_MQTT, mqtt_thread, h, 4*MB, 0);
+    thread_start(THREAD_ID_MQTT, h);
     h->inited = 1;
     r = 0;
 #endif
@@ -490,7 +553,6 @@ int mqtt_deinit(void)
     mqtt_handle_t *h=&mqHandle;
     
     thread_stop(THREAD_ID_MQTT);
-    mg_mgr_free(&h->mgr);
     xlist_free(h->list);
     h->inited = 0;
     r = 0;
@@ -499,20 +561,23 @@ int mqtt_deinit(void)
     return r;
 }
 
-#ifdef USE_MQTT
-static int send_stat(mqtt_handle_t *h, int type)
+
+static int send_stat(mqtt_conn_t *conn, int type)
 {
     int i,r=-1;
+
+#ifdef USE_MQTT
     float p1,p2,p3;
     char topic[256];
     char* content=NULL;
     char buf[20],date[40],time[40];
+    mqtt_handle_t *h=(mqtt_handle_t*)conn->h;
 
-    if(!h->inited || type<0 || type>=MQTT_PUB_MAX) {
+    if((!conn->c) || (type<0) || (type>=MQTT_PUB_MAX)) {
         return -1;
     }
 
-    if(type!=MQTT_PUB_STAT_POWER_CHN && type!=MQTT_PUB_STAT_SENSOR) {
+    if((type!=MQTT_PUB_STAT_POWER_CHN) && (type!=MQTT_PUB_STAT_SENSOR)) {
         snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id);
     }
 
@@ -533,7 +598,7 @@ static int send_stat(mqtt_handle_t *h, int type)
                 cJSON_AddStringToObject(root,"time", time);
 
                 content = cJSON_Print(root);    
-                my_pub(h, type, topic, content);
+                my_pub(conn, topic, content);
 
                 cJSON_free(content);
                 cJSON_Delete(root);
@@ -592,7 +657,7 @@ static int send_stat(mqtt_handle_t *h, int type)
                 cJSON_AddStringToObject(root,"time", time);
 
                 content = cJSON_Print(root);    
-                my_pub(h, type, topic, content);
+                my_pub(conn, topic, content);
 
                 cJSON_free(content);
                 cJSON_Delete(root);
@@ -633,7 +698,7 @@ static int send_stat(mqtt_handle_t *h, int type)
                     cJSON_AddStringToObject(root,"time", time);
 
                     content = cJSON_Print(root);    
-                    my_pub(h, type, topic, content);
+                    my_pub(conn, topic, content);
                     cJSON_free(content);
 
                     cJSON_Delete(root);
@@ -680,7 +745,7 @@ static int send_stat(mqtt_handle_t *h, int type)
 
                     snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id, tmp->product_ch_id);
                     content = cJSON_Print(root);    
-                    my_pub(h, type, topic, content);
+                    my_pub(conn, topic, content);
 
                     cJSON_free(content);
                     cJSON_Delete(root);
@@ -718,7 +783,7 @@ static int send_stat(mqtt_handle_t *h, int type)
 
                     snprintf(topic, sizeof(topic), topic_pub[type], h->prod_id, tmp->sensor_id);
                     content = cJSON_Print(root);    
-                    my_pub(h, type, topic, content);
+                    my_pub(conn, topic, content);
 
                     cJSON_free(content);
                     cJSON_Delete(root);
@@ -733,7 +798,7 @@ static int send_stat(mqtt_handle_t *h, int type)
             cJSON* root=cJSON_CreateObject();
             if(root) {
     
-                uint32_t flag = get_serv(NULL);
+                int flag = get_serv(NULL);
                 cJSON_AddStringToObject(root,"telnet", BGET(flag, SERV_TELNET)?"1":"0");
                 cJSON_AddStringToObject(root,"smtp", BGET(flag, SERV_SMTP)?"1":"0");
                 cJSON_AddStringToObject(root,"snmp_v1", BGET(flag, SERV_SNMP_V1)?"1":"0");
@@ -750,7 +815,7 @@ static int send_stat(mqtt_handle_t *h, int type)
                 cJSON_AddStringToObject(root,"time", time);
 
                 content = cJSON_Print(root);
-                my_pub(h, type, topic, content);
+                my_pub(conn, topic, content);
 
                 cJSON_free(content);
                 cJSON_Delete(root);
@@ -761,11 +826,10 @@ static int send_stat(mqtt_handle_t *h, int type)
         default:
         return -1;
     }
+#endif
 
     return r;
 }
-#endif
-
 
 static int post_alarm(int type, int subtype, char *content, alarm_para_t *para)
 {
@@ -838,6 +902,7 @@ int mqtt_post_alarm(int type, int subtype, char *content, alarm_para_t *para)
 }
 
 
+
 #ifdef USE_MQTT
 //////////////////////////////////////////////////////////
 static int get_cmd(char *json)
@@ -856,24 +921,12 @@ static int get_cmd(char *json)
 
     return cmd;
 }
-const char *serv_str[SERV_MAX]={
-    "ntp",
-    "smtp",
-    "message",
-    "mqtt",
-    "cloud",
-    "telnet",
-    "snmp_v1",
-    "snmp_v2c",
-    "snmp_v3",
-    "snmp_trap",
-};
 
-static uint32_t get_serv(char *json)
+
+static int get_serv(char *json)
 {
-    int i;
+    int i,flag=0;
     cJSON* tmp=NULL;
-    uint32_t flag=0;
 
     if(json) {
         cJSON* cjson=cJSON_Parse(json);
@@ -927,7 +980,7 @@ static uint32_t get_serv(char *json)
 
     return flag;
 }
-static int set_serv(uint32_t flag)
+static int set_serv(int flag)
 {
     if(BGET(flag,SERV_NTP)) {
         if(!thread_is_running(THREAD_ID_NTP)) {
@@ -997,7 +1050,7 @@ static int my_recv(char *topic, char *data)
 
 #ifdef USE_MQTT
     char *p,temp[512];
-    uint32_t flag=0;
+    int flag=0;
     mqtt_handle_t *h=&mqHandle;
 
     for(i=0; i<MQTT_SUB_MAX; i++) {
@@ -1078,9 +1131,7 @@ static int my_recv(char *topic, char *data)
                     cmd = get_cmd(data);
                     LOGD("___ MQTT_SUB_CMD_RESET, %d\n", cmd);
                     if(cmd==1) {
-                        //sys_set_factory();
-                        char temp[200];
-                        sys_get_path(temp, "");
+                        sys_set_factory();
                     }
                 }
                 break;
@@ -1127,11 +1178,6 @@ int mqtt_test(void)
     mInfo.ser[0] = ser;
     h->dm->mqttInfo = mInfo;
 
-    my_conn(h);
-    while(1) {
-        send_period(h);
-        sleep(3);
-    }
 
 #endif
     return 0;

+ 6 - 4
pro/src/netswitch.c

@@ -857,10 +857,11 @@ static int sw_alarm(netswitch_evts_t *evts, netswitch_info_t *old, netswitch_inf
 
     return r;
 }
+#endif
 
-
-static void* sw_thread(void *arg)
+void* netswitch_thread(void *arg)
 {
+#ifdef USE_NETSWITCH
     int r;
     netswitch_evts_t   evts;
     thread_handle_t    *h=(thread_handle_t*)arg;
@@ -883,10 +884,11 @@ static void* sw_thread(void *arg)
         sleep(1);
     }
     SOCK_CLEANUP;
+#endif
 
     pthread_exit(NULL);
 }
-#endif
+
 
 int netswitch_init(void)
 {
@@ -898,7 +900,7 @@ int netswitch_init(void)
 
     dev_switch_init(get_db(), &h->sw);
     pthread_mutex_init(&h->mutex, NULL);
-    thread_start(THREAD_ID_SWITCH, sw_thread, h, 10*MB, 0);
+    thread_start(THREAD_ID_SWITCH, h);
 #endif
 
     return 0;

+ 7 - 3
pro/src/smtp.c

@@ -23,6 +23,7 @@
  */
 
 #include "cfg.h"
+#include <pthread.h>
 
 #ifdef USE_SMTP
 
@@ -4106,8 +4107,11 @@ static void sleep_ms(int ms)
 {
     usleep(ms*1000);
 }
-static void* smtp_thread(void *arg)
+#endif
+
+void* smtp_thread(void *arg)
 {
+#ifdef USE_SMTP
     int r;
     thread_handle_t *th=(thread_handle_t*)arg;
     smtp_handle_t   *h=(smtp_handle_t*)th->arg;
@@ -4125,10 +4129,10 @@ static void* smtp_thread(void *arg)
         }
         sleep_ms(100);
     }
+#endif
 
     pthread_exit(NULL);
 }
-#endif
 
 
 int smtp_init(void)
@@ -4147,7 +4151,7 @@ int smtp_init(void)
     h->info = &h->dm->smtpInfo;
 
     dev_smtp_init(dm->db, h->info);
-    thread_start(THREAD_ID_SMTP, smtp_thread, h, 5*MB, 0);
+    thread_start(THREAD_ID_SMTP, h);
     h->inited = 1;
 
     //smtp_test();

+ 2 - 2
pro/src/sqlite_handle.c

@@ -4479,11 +4479,11 @@ int dev_mqtt_update_ser(sqlite3 *db, mqtt_server_t *ser)
 
     cnt = get_count(db, table);
     if(cnt>0) {
-        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;", 
+        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);
         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\n\n", errmsg);
             sqlite3_free(errmsg);
         }
         else {

+ 2 - 2
pro/src/sys.c

@@ -1189,7 +1189,7 @@ static void mem_check(void)
     }
 }
 
-static void* polling_thread(void *arg)
+void* polling_thread(void *arg)
 {
     thread_handle_t *h=(thread_handle_t*)arg;
 
@@ -1208,7 +1208,7 @@ int sys_init(void)
     lock_s_init();
     paras_init();
     
-    thread_start(THREAD_ID_POLLING, polling_thread, NULL, 4*MB, 0);
+    thread_start(THREAD_ID_POLLING, NULL);
 
     return 0;
 }

+ 34 - 11
pro/src/thread.c

@@ -1,32 +1,55 @@
 #include "thread.h"
 
-static thread_attr_t thdAttrs[THREAD_ID_MAX]={0};
+static thread_attr_t thdAttrs[THREAD_ID_MAX]={
+    //name          fn                          stacksize       prio
+    {"io",          NULL,                       8*MB,           1},
+    {"ws",          ws_thread,                  8*MB,           1},
+    {"web",         web_thread,                 8*MB,           1},
+    {"ntp",         ntp_thread,                 8*MB,           1},
+    {"snmp",        snmp_thread,                8*MB,           1},
+    {"power",       power_thread,               8*MB,           1},
+    {"sensor",      sensor_thread,              8*MB,           1},
+    {"polling",     polling_thread,             8*MB,           1},
+    {"cascade",     cascade_thread,             8*MB,           1},
+    {"casscan",     cascade_scan_thread,        8*MB,           1},
+    {"shmw",        shm_thread,                 8*MB,           1},
+    {"breaker",     breaker_thread,             8*MB,           1},
+    {"switch",      netswitch_thread,           8*MB,           1},
+    {"brkscan",     breaker_scanner_thread,     8*MB,           1},
+    {"mqtt",        mqtt_thread,                8*MB,           1},
+    {"smtp",        smtp_thread,                8*MB,           1},
+};
 static thread_handle_t thdHandles[THREAD_ID_MAX]={0};
 
 
-int thread_reg(int id, thread_attr_t *attr)
+int thread_start(int id, void *arg)
 {
-    if(id<0 || id>=THREAD_ID_MAX || !attr) {
+    if(id<0 || id>=THREAD_ID_MAX) {
         return -1;
     }
-    thdAttrs[id] = *attr;
-    return 0;
-}
+    thread_attr_t *attr=&thdAttrs[id];
 
+    return thread_startEx(id, attr->name, attr->fn, arg, attr->stksz, attr->prio);
+}
 
-int thread_start(int id, thread_fn fn, void *arg, int stksz, int prio)
+int thread_startEx(int id, const char *name, thread_fn fn, void *arg, int stksz, int prio)
 {
     int r;
     pthread_attr_t attr;
+    thread_handle_t *h=NULL;
 
-    if(id<0 || id>=THREAD_ID_MAX) {
+    if(id<0 || id>=THREAD_ID_MAX || !fn) {
         return -1;
     }
-    thread_handle_t *h=&thdHandles[id];
+    h = &thdHandles[id];
 
+    if(h->running) {
+        return -1;
+    }
+
+    h->name = name;
     h->arg  = arg;
     h->quit = 0;
-
     r = pthread_attr_init(&attr);
     if(r) {
         printf("____ thread %d attr init failed\n", id);
@@ -139,7 +162,7 @@ int thread_print_stacksize(void)
         printf("___ print stacksize, attr init failed\n");
     }
     else {
-        printf("_____ stacksize: %dMB\n", size/MB);
+        printf("_____ stacksize: %dMB\n", size/(1024*1024));
     }
     pthread_attr_destroy(&attr);
     return r;

+ 23 - 3
pro/src/thread.h

@@ -22,7 +22,7 @@ enum {
     THREAD_ID_SENSOR,
     THREAD_ID_POLLING,
     THREAD_ID_CASCADE,
-    THREAD_ID_SCAN,
+    THREAD_ID_CASCADE_SCAN,
     THREAD_ID_SHMW,
     THREAD_ID_BREAKER,
     THREAD_ID_SWITCH,
@@ -37,13 +37,14 @@ typedef void *(*thread_fn)(void *arg);
 
 
 typedef struct {
+    const char *name;
     thread_fn  fn;
-    void       *arg;
     int        stksz;
     int        prio;
 }thread_attr_t;
 
 typedef struct {
+    const char *name;
     pthread_t  tid;
     int        quit;
     int        running;
@@ -52,7 +53,26 @@ typedef struct {
     void       *stk;
 }thread_handle_t;
 
-int thread_start(int id, thread_fn fn, void *arg, int stksz, int prio);
+
+void* ws_thread(void *arg);
+void* web_thread(void *arg);
+void* ntp_thread(void *arg);
+void* snmp_thread(void *arg);
+void* power_thread(void *arg);
+void* sensor_thread(void *arg);
+void* polling_thread(void *arg);
+void* cascade_thread(void *arg);
+void* cascade_scan_thread(void *arg);
+void* shm_thread(void *arg);
+void* breaker_thread(void *arg);
+void* netswitch_thread(void *arg);
+void* breaker_scanner_thread(void *arg);
+void* mqtt_thread(void *arg);
+void* smtp_thread(void *arg);
+
+
+int thread_start(int id, void *arg);
+int thread_startEx(int id, const char *name, thread_fn fn, void *arg, int stksz, int prio);
 int thread_stop(int id);
 int thread_stop_all(void);
 

+ 2 - 2
pro/src/websocket_handle.c

@@ -271,7 +271,7 @@ static void fn(mg_conn_t *c, int ev, void *ev_data)
   }
 
 }
-void* websocket_thread(void* arg)
+void* ws_thread(void* arg)
 {
     int r;
     thread_handle_t *h=(thread_handle_t*)arg;
@@ -305,7 +305,7 @@ int websocket_init(void)
     memset(wh, 0, sizeof(ws_handle_t));
     //read_ipaddr(wh);
 
-    thread_start(THREAD_ID_WS, websocket_thread, wh, 30*MB, 0);
+    thread_start(THREAD_ID_WS, wh);
     return 0;
 }