guohui лет назад: 2
Родитель
Сommit
795cfd3e8e
7 измененных файлов с 157 добавлено и 90 удалено
  1. 5 2
      pro/src/appweb_handle.c
  2. 2 2
      pro/src/common.h
  3. 48 40
      pro/src/dflt.c
  4. 1 0
      pro/src/json_handle.h
  5. 17 12
      pro/src/mqtt.c
  6. 81 32
      pro/src/sqlite_handle.c
  7. 3 2
      pro/src/sqlite_handle.h

+ 5 - 2
pro/src/appweb_handle.c

@@ -3463,9 +3463,12 @@ static void serviceManagement(void *conn)
         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(req->cert) { 
+            strcpy(ser->cert, req->cert); save_flag=1;
+        }
 
         if(save_flag) {
-            dev_mqtt_update(gdm->db, ser);
+            dev_mqtt_update_ser(gdm->db, ser);
         }
     }
     else if(strcmp(_serviceManageRequestInfo.type,"mqttOpen")==0 ||
@@ -3476,7 +3479,7 @@ static void serviceManagement(void *conn)
         mqtt_server_t *ser=&gdm->mqttInfo.ser[req->idx];
         if(req->mode!=ser->mode) {
             ser->mode = req->mode;
-            dev_mqtt_update(gdm->db, ser);
+            dev_mqtt_update_ser(gdm->db, ser);
         }
     }
     else if(strcmp(_serviceManageRequestInfo.type,"mqttZtQuery")==0)

+ 2 - 2
pro/src/common.h

@@ -645,10 +645,10 @@ typedef struct {
     char                cid[32];    		            //client ID
     char                user[32];    			        //认证方式  
     char                password[32];
-    char                *cert;
+    char                cert[8192];
 }mqtt_server_t;
 typedef struct {
-    #define MQTT_SER_MAX   10
+    #define MQTT_SER_MAX   5
     mqtt_server_t       ser[MQTT_SER_MAX];
 }mqtt_info_t;
 

+ 48 - 40
pro/src/dflt.c

@@ -52,46 +52,54 @@ smtp_send_t DFLT_SMTP_SEND={
     .auth = "PLAIN",
 };
 
-mqtt_server_t DFLT_MQTT={
+mqtt_info_t DFLT_MQTT={
+
+    .ser[0] = {
+        .id = 0,
+        .mode = 1,
+        .cid = "",
+        .proto = "mqtt",
+        .server = "192.168.1.12",
+        .ip = "",
+        .port = "1883",
+        .cert = "",
+        .user = "gowone100",
+        .password = "gowone100",
+    },
+
+    .ser[1] = {
+        .id = 1,
+        .mode = 1,
+        .cid = "",
+        .proto = "mqtts",
+        .server = "gaae4d7b.ala.cn-hangzhou.emqxsl.cn",
+        .ip = "",
+        .port = "8883",
+        .cert = "-----BEGIN CERTIFICATE-----\r\n"
+                "MIIDrzCCApegAwIBAgIQCDvgVpBCRrGhdWrJWZHHSjANBgkqhkiG9w0BAQUFADBh\r\n"
+                "MQswCQYDVQQGEwJVUzEVMBMGA1UEChMMRGlnaUNlcnQgSW5jMRkwFwYDVQQLExB3\r\n"
+                "d3cuZGlnaWNlcnQuY29tMSAwHgYDVQQDExdEaWdpQ2VydCBHbG9iYWwgUm9vdCBD\r\n"
+                "QTAeFw0wNjExMTAwMDAwMDBaFw0zMTExMTAwMDAwMDBaMGExCzAJBgNVBAYTAlVT\r\n"
+                "MRUwEwYDVQQKEwxEaWdpQ2VydCBJbmMxGTAXBgNVBAsTEHd3dy5kaWdpY2VydC5j\r\n"
+                "b20xIDAeBgNVBAMTF0RpZ2lDZXJ0IEdsb2JhbCBSb290IENBMIIBIjANBgkqhkiG\r\n"
+                "9w0BAQEFAAOCAQ8AMIIBCgKCAQEA4jvhEXLeqKTTo1eqUKKPC3eQyaKl7hLOllsB\r\n"
+                "CSDMAZOnTjC3U/dDxGkAV53ijSLdhwZAAIEJzs4bg7/fzTtxRuLWZscFs3YnFo97\r\n"
+                "nh6Vfe63SKMI2tavegw5BmV/Sl0fvBf4q77uKNd0f3p4mVmFaG5cIzJLv07A6Fpt\r\n"
+                "43C/dxC//AH2hdmoRBBYMql1GNXRor5H4idq9Joz+EkIYIvUX7Q6hL+hqkpMfT7P\r\n"
+                "T19sdl6gSzeRntwi5m3OFBqOasv+zbMUZBfHWymeMr/y7vrTC0LUq7dBMtoM1O/4\r\n"
+                "gdW7jVg/tRvoSSiicNoxBN33shbyTApOB6jtSj1etX+jkMOvJwIDAQABo2MwYTAO\r\n"
+                "BgNVHQ8BAf8EBAMCAYYwDwYDVR0TAQH/BAUwAwEB/zAdBgNVHQ4EFgQUA95QNVbR\r\n"
+                "TLtm8KPiGxvDl7I90VUwHwYDVR0jBBgwFoAUA95QNVbRTLtm8KPiGxvDl7I90VUw\r\n"
+                "DQYJKoZIhvcNAQEFBQADggEBAMucN6pIExIK+t1EnE9SsPTfrgT1eXkIoyQY/Esr\r\n"
+                "hMAtudXH/vTBH1jLuG2cenTnmCmrEbXjcKChzUyImZOMkXDiqw8cvpOp/2PV5Adg\r\n"
+                "06O/nVsJ8dWO41P0jmP6P6fbtGbfYmbW0W5BjfIttep3Sp+dWOIrWcBAI+0tKIJF\r\n"
+                "PnlUkiaY4IBIqDfv8NZ5YBberOgOzW6sRBc4L0na4UU+Krk2U886UAb3LujEV0ls\r\n"
+                "YSEY1QSteDwsOoBrp+uvFRTp2InBuThs4pFsiv9kuXclVzDAGySj4dzp30d8tbQk\r\n"
+                "CAUw7C29C79Fv1C5qfPrmAESrciIxpg0X40KPMbp1ZWVbd4=\r\n"
+                "-----END CERTIFICATE-----\r\n",
+        .user = "gowone100",
+        .password = "gowone100",
+    },
 
-    .id = 0,
-    .mode = 1,
-    .cid = "",
-#if 0
-    .proto = "mqtt",
-    .server = "192.168.1.12",
-    .ip = "",
-    .port = "1883",
-    .cert = "",
-#else
-    .proto = "mqtts",
-    .server = "gaae4d7b.ala.cn-hangzhou.emqxsl.cn",
-    .ip = "",
-    .port = "8883",
-    .cert = "-----BEGIN CERTIFICATE-----\r\n"
-            "MIIDrzCCApegAwIBAgIQCDvgVpBCRrGhdWrJWZHHSjANBgkqhkiG9w0BAQUFADBh\r\n"
-            "MQswCQYDVQQGEwJVUzEVMBMGA1UEChMMRGlnaUNlcnQgSW5jMRkwFwYDVQQLExB3\r\n"
-            "d3cuZGlnaWNlcnQuY29tMSAwHgYDVQQDExdEaWdpQ2VydCBHbG9iYWwgUm9vdCBD\r\n"
-            "QTAeFw0wNjExMTAwMDAwMDBaFw0zMTExMTAwMDAwMDBaMGExCzAJBgNVBAYTAlVT\r\n"
-            "MRUwEwYDVQQKEwxEaWdpQ2VydCBJbmMxGTAXBgNVBAsTEHd3dy5kaWdpY2VydC5j\r\n"
-            "b20xIDAeBgNVBAMTF0RpZ2lDZXJ0IEdsb2JhbCBSb290IENBMIIBIjANBgkqhkiG\r\n"
-            "9w0BAQEFAAOCAQ8AMIIBCgKCAQEA4jvhEXLeqKTTo1eqUKKPC3eQyaKl7hLOllsB\r\n"
-            "CSDMAZOnTjC3U/dDxGkAV53ijSLdhwZAAIEJzs4bg7/fzTtxRuLWZscFs3YnFo97\r\n"
-            "nh6Vfe63SKMI2tavegw5BmV/Sl0fvBf4q77uKNd0f3p4mVmFaG5cIzJLv07A6Fpt\r\n"
-            "43C/dxC//AH2hdmoRBBYMql1GNXRor5H4idq9Joz+EkIYIvUX7Q6hL+hqkpMfT7P\r\n"
-            "T19sdl6gSzeRntwi5m3OFBqOasv+zbMUZBfHWymeMr/y7vrTC0LUq7dBMtoM1O/4\r\n"
-            "gdW7jVg/tRvoSSiicNoxBN33shbyTApOB6jtSj1etX+jkMOvJwIDAQABo2MwYTAO\r\n"
-            "BgNVHQ8BAf8EBAMCAYYwDwYDVR0TAQH/BAUwAwEB/zAdBgNVHQ4EFgQUA95QNVbR\r\n"
-            "TLtm8KPiGxvDl7I90VUwHwYDVR0jBBgwFoAUA95QNVbRTLtm8KPiGxvDl7I90VUw\r\n"
-            "DQYJKoZIhvcNAQEFBQADggEBAMucN6pIExIK+t1EnE9SsPTfrgT1eXkIoyQY/Esr\r\n"
-            "hMAtudXH/vTBH1jLuG2cenTnmCmrEbXjcKChzUyImZOMkXDiqw8cvpOp/2PV5Adg\r\n"
-            "06O/nVsJ8dWO41P0jmP6P6fbtGbfYmbW0W5BjfIttep3Sp+dWOIrWcBAI+0tKIJF\r\n"
-            "PnlUkiaY4IBIqDfv8NZ5YBberOgOzW6sRBc4L0na4UU+Krk2U886UAb3LujEV0ls\r\n"
-            "YSEY1QSteDwsOoBrp+uvFRTp2InBuThs4pFsiv9kuXclVzDAGySj4dzp30d8tbQk\r\n"
-            "CAUw7C29C79Fv1C5qfPrmAESrciIxpg0X40KPMbp1ZWVbd4=\r\n"
-            "-----END CERTIFICATE-----\r\n",
-#endif
-    .user = "gowone100",
-    .password = "gowone100",
 };
 

+ 1 - 0
pro/src/json_handle.h

@@ -363,6 +363,7 @@ typedef struct
     char* cid;
     char* user;
     char* password;
+    char* cert;
 }_MQTT_ServerRequestInfo;
 
 /// @brief SNMP信息

+ 17 - 12
pro/src/mqtt.c

@@ -79,7 +79,7 @@ typedef struct {
 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, int ms);
+static void timer_start(mqtt_handle_t *h);
 static int post_stat(int type, int id);
 static uint32_t get_serv(char *json);
 static void post_once(void);
@@ -366,17 +366,25 @@ static void post_period(void)
     post_stat(MQTT_PUB_STAT_POWER_CHN, 0);
 }
 
-static void timer_fn(void *arg)
+static void timer_conn_fn(void *arg)
 {
     mqtt_handle_t *h=(mqtt_handle_t*)arg;
-
     my_conn(h);
+}
+static void timer_period_fn(void *arg)
+{
     post_period();
 }
-
-static void timer_start(mqtt_handle_t *h, int ms)
+static void timer_send_fn(void *arg)
 {
-    mg_timer_add(&h->mgr, ms, MG_TIMER_REPEAT | MG_TIMER_RUN_NOW, timer_fn, h);
+    mqtt_handle_t *h=(mqtt_handle_t*)arg;
+    my_send(h);
+}
+static void timer_start(mqtt_handle_t *h)
+{
+    mg_timer_add(&h->mgr, 3000, MG_TIMER_REPEAT | MG_TIMER_RUN_NOW, timer_conn_fn, h);
+    mg_timer_add(&h->mgr, 5000, MG_TIMER_REPEAT | MG_TIMER_RUN_NOW, timer_period_fn, h);
+    mg_timer_add(&h->mgr, 2000, MG_TIMER_REPEAT | MG_TIMER_RUN_NOW, timer_send_fn, h);
 }
 
 static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
@@ -442,9 +450,8 @@ static void* mqtt_thread(void *arg)
     mqtt_handle_t   *h=(mqtt_handle_t*)th->arg;
     mqtt_info_t     *info=&h->dm->mqttInfo;
     
-    //timer_start(h, 10000);
+    timer_start(h);
     while(th->quit==0) {
-        //my_send(h);
         mg_mgr_poll(&h->mgr, 2000);
     }
     my_disconn(h);
@@ -480,15 +487,13 @@ int mqtt_init(void)
     lc.log = 0;
     h->list = xlist_init(&lc);
 
-    dev_mqtt_init(h->dm->db, h->dm->mqttInfo.ser);
+    dev_mqtt_init(h->dm->db, &h->dm->mqttInfo);
 
     thread_start(THREAD_ID_MQTT, mqtt_thread, h, 4*MB, 0);
     h->inited = 1;
     r = 0;
 #endif
 
-    mqtt_test();
-
     return r;
 }
 
@@ -1066,7 +1071,7 @@ int mqtt_test(void)
 
     mInfo.ser[0] = ser;
     h->dm->mqttInfo = mInfo;
-return 0;
+
     my_conn(h);
     while(1) {
         post_period();

+ 81 - 32
pro/src/sqlite_handle.c

@@ -4349,16 +4349,19 @@ quit:
 
 ////////////////mqtt/////////////////////////////////////////////
 #define MQTT_MANAGER_TABLE   "Table_MQTTManage"
-int dev_mqtt_init(sqlite3 *db, mqtt_server_t *ser)
+int dev_mqtt_init(sqlite3 *db, mqtt_info_t *info)
 {
     int i,r;
 
 #ifdef USE_MQTT
+    int idx=0;
     char temp[1024];
     char *errmsg=NULL;
     char *table=NULL;
+    char** pResult=NULL;
     int nrow,ncol;
-    char** pResult = NULL;
+    sqlite3_stmt *stmt=NULL;
+    mqtt_server_t *ser=info->ser;
     
     table = MQTT_MANAGER_TABLE;
     if(!tab_exists(db, table)) {
@@ -4370,7 +4373,8 @@ int dev_mqtt_init(sqlite3 *db, mqtt_server_t *ser)
                             "port TEXT,"
                             "cid TEXT,"
                             "user TEXT,"
-                            "password TEXT);", table);
+                            "password TEXT,"
+                            "cert TEXT);", table);
         r = sqlite3_exec(db, temp, 0, 0, &errmsg);
         if(r!=SQLITE_OK) {
             log_e("___ %s create failed, %s\n%s\n", table, temp, errmsg); sqlite3_free(errmsg);
@@ -4378,47 +4382,72 @@ int dev_mqtt_init(sqlite3 *db, mqtt_server_t *ser)
         }
 
 #ifdef MQTT_SERVER_BUILTIN
-        extern mqtt_server_t DFLT_MQTT;
-        dev_mqtt_update(db, &DFLT_MQTT);
+        extern mqtt_info_t DFLT_MQTT;
+        dev_mqtt_update_info(db, &DFLT_MQTT);
 #endif
     }
 
-    nrow = ncol = 0;
+    
     sprintf(temp, "SELECT * FROM %s;", table);
-    r = sqlite3_get_table(db, temp, &pResult, &nrow, &ncol, &errmsg);
-    if(r != SQLITE_OK) {
-        log_e("sqlite handle error: %s \n%s\n", temp, errmsg); sqlite3_free(errmsg);
-        return -1 ;
+    r = sqlite3_prepare_v2(db, temp, -1, &stmt, 0);
+    if (r != SQLITE_OK) {
+        log_e("___mqtt_load, sqlite3_prepare_v2 failed\n");
+        return -1;
     }
+    
+    memset(info, 0, sizeof(mqtt_info_t));
+    char *p=NULL;
+    while (sqlite3_step(stmt) == SQLITE_ROW) {
+        ser[idx].id = sqlite3_column_int(stmt, 0);
+        ser[idx].mode = sqlite3_column_int(stmt, 1);
 
-    if(ser && ncol>0) {
-        memset(ser, 0, sizeof(mqtt_server_t));
-        
-        ser->id = atoi(pResult[ncol+0]);
-        ser->mode = atoi(pResult[ncol+1]);
-        strcpy(ser->server, pResult[ncol+2]);
-        strcpy(ser->ip, pResult[ncol+3]);
-        strcpy(ser->port, pResult[ncol+4]);
-        strcpy(ser->cid, pResult[ncol+5]);
-        strcpy(ser->user, pResult[ncol+6]);
-        strcpy(ser->password, pResult[ncol+7]);
+        p = (char*)sqlite3_column_text(stmt, 2);
+        if(p) strcpy(ser[idx].server, p);
+
+        p = (char*)sqlite3_column_text(stmt, 3);
+        if(p) strcpy(ser[idx].ip, p);
+
+        p = (char*)sqlite3_column_text(stmt, 4);
+        if(p) strcpy(ser[idx].port, p);
+
+        p = (char*)sqlite3_column_text(stmt, 5);
+        if(p) strcpy(ser[idx].cid, p);
+
+        p = (char*)sqlite3_column_text(stmt, 6);
+        if(p) strcpy(ser[idx].user, p);
+
+        p = (char*)sqlite3_column_text(stmt, 7);
+        if(p) strcpy(ser[idx].password, p);
+
+        p = (char*)sqlite3_column_text(stmt, 8);
+        if(p) strcpy(ser[idx].cert, p);
+
+        idx++;
+        if(idx>=MQTT_SER_MAX) {
+            break;
+        }
     }
+    sqlite3_finalize(stmt);
 #endif
 
     return 0;
 }
 
 
-int dev_mqtt_update(sqlite3 *db, mqtt_server_t *ser)
+int dev_mqtt_update_ser(sqlite3 *db, mqtt_server_t *ser)
 {
     int i,r=0,exist=0,cnt=0;
     char mac[10];
 
 #ifdef USE_MQTT
-    char temp[1024];
     sqlite3_stmt *stmt;
     char *errmsg=NULL;
     char *table=MQTT_MANAGER_TABLE;
+    char *pbuf=malloc(sizeof(mqtt_server_t)+1024);
+
+    if(!pbuf) {
+        return -1;
+    }
 
     if(sys_get_mac(mac, sizeof(mac))==0) {
         sprintf(ser->cid, "smartPDU_%s", mac);
@@ -4426,11 +4455,11 @@ int dev_mqtt_update(sqlite3 *db, mqtt_server_t *ser)
 
     cnt = get_count(db, table);
     if(cnt>0) {
-        sprintf(temp, "UPDATE %s SET id=%d, mode=%d, server='%s', ip='%s', port='%s', cid='%s', user='%s', password='%s', WHERE id=%d;", 
-                table, ser->id, ser->mode, ser->server, ser->ip, ser->port, ser->cid, ser->user, ser->password, ser->id);
-        r = sqlite3_exec(db, temp, NULL, 0, &errmsg);
+        sprintf(pbuf, "UPDATE %s SET id=%d, mode=%d, server='%s', ip='%s', port='%s', cid='%s', user='%s', password='%s, cert='%s', WHERE id=%d;", 
+                table, ser->id, ser->mode, ser->server, ser->ip, ser->port, ser->cid, ser->user, ser->password, ser->cert, ser->id);
+        r = sqlite3_exec(db, pbuf, NULL, 0, &errmsg);
         if(r != SQLITE_OK) {
-            log_e("___dev_mqtt_update, sqlite3_exec failed, %s, %s\n\n", temp, errmsg);
+            log_e("___dev_mqtt_update_ser, sqlite3_exec failed, %s, %s\n\n", pbuf, errmsg);
             sqlite3_free(errmsg);
         }
         else {
@@ -4439,10 +4468,10 @@ int dev_mqtt_update(sqlite3 *db, mqtt_server_t *ser)
     }
 
     if(exist==0) {
-        sprintf(temp, "INSERT INTO %s (id,mode,server,ip,port,cid,user,password) VALUES (?,?,?,?,?,?,?,?);", table);
-        r = sqlite3_prepare_v2(db, temp, -1, &stmt, 0);
+        sprintf(pbuf, "INSERT INTO %s (id,mode,server,ip,port,cid,user,password,cert) VALUES (?,?,?,?,?,?,?,?,?);", table);
+        r = sqlite3_prepare_v2(db, pbuf, -1, &stmt, 0);
         if (r != SQLITE_OK) {
-            log_e("___dev_mqtt_update, sqlite3_prepare_v2 failed, %s\n", sqlite3_errmsg(db));
+            log_e("___dev_mqtt_update_ser, sqlite3_prepare_v2 failed, %s\n", sqlite3_errmsg(db));
             r = -1; goto quit;
         }
 
@@ -4454,16 +4483,36 @@ int dev_mqtt_update(sqlite3 *db, mqtt_server_t *ser)
         sqlite3_bind_text(stmt, 6, ser->cid,      -1, NULL);
         sqlite3_bind_text(stmt, 7, ser->user,     -1, NULL);
         sqlite3_bind_text(stmt, 8, ser->password, -1, NULL);
+        sqlite3_bind_text(stmt, 9, ser->cert,     -1, NULL);
 
         r = sqlite3_step(stmt);
         sqlite3_finalize(stmt);
-
         if (r != SQLITE_DONE) {
-            log_e("___dev_mqtt_update, sqlite3_step failed\n");
+            log_e("___dev_mqtt_update_ser, sqlite3_step failed\n");
             r = -1;
         }
+        r = 0;
     }
 quit:
+    free(pbuf);
+#endif
+
+    return r;
+}
+
+
+int dev_mqtt_update_info(sqlite3 *db, mqtt_info_t *info)
+{
+    int i,r=-1;
+
+#ifdef USE_MQTT
+    for(i=0; i<MQTT_SER_MAX; i++) {
+        mqtt_server_t *ser=&info->ser[i];
+        r = dev_mqtt_update_ser(db, ser);
+        if(r) {
+            break;
+        }
+    }
 #endif
 
     return r;

+ 3 - 2
pro/src/sqlite_handle.h

@@ -149,8 +149,9 @@ int dev_smtp_init(sqlite3 *db, smtp_info_t *info);
 int dev_smtp_update_send(sqlite3 *db, smtp_send_t *send);
 int dev_smtp_update_recv(sqlite3 *db, smtp_recv_t *recv);
 
-int dev_mqtt_init(sqlite3 *db, mqtt_server_t *ser);
-int dev_mqtt_update(sqlite3 *db, mqtt_server_t *ser);
+int dev_mqtt_init(sqlite3 *db, mqtt_info_t *info);
+int dev_mqtt_update_ser(sqlite3 *db, mqtt_server_t *ser);
+int dev_mqtt_update_info(sqlite3 *db, mqtt_info_t *info);
 
 int dev_sql_VACCUM(sqlite3 *db);
 int dev_sql_set_overwrite(int flag);