소스 검색

mqtt后台添加群发功能

guohui 2 년 전
부모
커밋
7c3e132627
7개의 변경된 파일과 206개의 추가작업 그리고 225개의 파일을 삭제
  1. 10 5
      pro/config/cfg.h
  2. 2 2
      pro/src/common.h
  3. 162 215
      pro/src/mqtt.c
  4. 0 3
      pro/src/mqtt.h
  5. 6 0
      pro/src/smtp.c
  6. 24 0
      pro/src/sys.c
  7. 2 0
      pro/src/sys.h

+ 10 - 5
pro/config/cfg.h

@@ -19,12 +19,15 @@
 #define PHASE_GROUP_CNT                 8
 #define CHANNEL_DELAY_SEC               8
 
+#define USE_SMTP
+
 
 /////////////customer demannd//////////////////////
 //#define CUST_GOM037     //xian
 //#define CUST_USASL
 #define CUST_ANDERSON
 //#define CUST_HYPERTEC
+//#define CUST_ITK
 
 
 #ifdef CUST_GOM037      //xian
@@ -45,9 +48,13 @@
 #endif
 
 
+#ifdef CUST_ITK
+    #define SMTP_SEND_BUILTIN
+#endif
+
+
 #ifdef CUST_ANDERSON
-    #define USE_SMTP
-    //#define USE_MQTT
+    #define USE_MQTT
     #define USE_NETSWITCH
 
     #define SW_STATIC_PORT_START
@@ -60,9 +67,7 @@
         #define SW_PDU_IP               "192.168.234.132"
     #endif
 
-    #ifdef USE_SMTP
-        #define SMTP_SEND_BUILTIN
-    #endif
+    #define SMTP_SEND_BUILTIN
 
     #ifdef USE_MQTT
         #define MQTT_SERVER_BUILTIN

+ 2 - 2
pro/src/common.h

@@ -645,8 +645,8 @@ typedef struct {
     char                password[32];
 }mqtt_server_t;
 typedef struct {
-    #define MQTT_SERV_MAX   10
-    mqtt_server_t       serv[MQTT_SERV_MAX];
+    #define MQTT_SER_MAX   10
+    mqtt_server_t       ser[MQTT_SER_MAX];
 }mqtt_info_t;
 
 typedef struct 

+ 162 - 215
pro/src/mqtt.c

@@ -41,54 +41,171 @@ char *topic_pub[MQTT_PUB_MAX]={
     "/PDU/%d/status/service",
 };
 
-
 typedef struct mg_mgr mgr_t;
 typedef struct mg_mqtt_opts mg_opts_t;
 typedef struct mg_connection mg_conn_t;
-typedef struct {
-    uint8_t     qos;
-    uint8_t     ver;
-    uint8_t     clean;
-    uint8_t     retain;
-}mqtt_para_t;
+
 typedef struct _mqtt_conn_t{
-    mg_conn_t     *c;
-    mg_opts_t     opts;
-    mqtt_para_t   para;
-    mqtt_info_t   *info;
-    bool          isover;
-    struct _mqtt_conn_t *next;
+    mg_conn_t       *c;
+    mqtt_server_t   ser;
 }mqtt_conn_t;
 typedef struct {
-    mgr_t         mgr;
-    mqtt_conn_t   *conn;
     int           inited;
     int           prod_id;
     GlobalDeviceManager *dm;
+
+    mgr_t         mgr;
+    mg_opts_t     opts;
+    mqtt_conn_t   conn[MQTT_SER_MAX];
 }mqtt_handle_t;
 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 int ser_cmp(mqtt_server_t *a, mqtt_server_t *b)
+{
+    if(strcmp(a->server, b->server) || 
+       strcmp(a->ip, b->ip) || 
+       strcmp(a->port, b->port) || 
+       strcmp(a->user, b->user) ||
+       strcmp(a->password, b->password) ||
+       (a->mode!=b->mode)) {
+        return 1;
+    }
 
-static int my_conn(mqtt_handle_t *h, mqtt_info_t *info);
-static int my_recv(char *topic, char *data);
+    return 0;
+}
+static int get_id(mqtt_handle_t *h, mqtt_server_t *ser)
+{
+    int i,r=-1;
+    mqtt_conn_t *conn=h->conn;
 
-static int info_cmp(mqtt_info_t *a, mqtt_info_t *b)
+    for(i=0; i<MQTT_SER_MAX; i++) {
+        if(conn[i].c && ser_cmp(ser, &conn[i].ser)==0) {
+            return i;
+        }
+    }
+    return -1;
+}
+static mqtt_conn_t* get_conn(mqtt_handle_t *h, mg_conn_t *c)
 {
-    return memcmp(a, b, sizeof(mqtt_info_t)-sizeof(mqtt_info_t*));
+    int i;
+    for(i=0; i<MQTT_SER_MAX; i++) {
+        if(h->conn[i].c==c) {
+            return &h->conn[i];
+        }
+    }
+    return NULL;
 }
+static int sub_one(mg_conn_t *c, char *topic)
+{
+    int r=-1;
+    mg_opts_t opts={0};
+    
+    opts.topic = mg_str(topic);
+    opts.qos = 1;
+    mg_mqtt_sub(c, &opts);
 
+    return 0;
+}
+static int my_sub(mg_conn_t *c, int prod_id)
+{
+    int i;
+    char topic[1024];
+    
+    for(i=0; topic_sub[i]; i++) {
+        snprintf(topic, sizeof(topic), topic_sub[i], prod_id);
+        sub_one(c, topic);
+    }
+    
+    return 0;
+}
+static int pub_one(mg_conn_t *c, char *topic, char *data)
+{
+    int r=-1;
+    mg_opts_t opts={0};
+    
+    opts.topic = mg_str(topic);
+    opts.message = mg_str(data);
+    opts.qos = 1;
+    opts.retain = false;
+    mg_mqtt_pub(c, &opts);
+
+    return 0;
+}
+static int my_pub(mqtt_handle_t *h, char *topic, char *data)
+{
+    int i;
+    for(i=0; i<MQTT_SER_MAX; i++) {
+        if(h->conn[i].c) {
+            pub_one(h->conn[i].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 int my_conn(mqtt_handle_t *h, mqtt_info_t *info)
+{
+    int i,j,r=-1;
+    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);
+                }
+            }
+        }
+    }
+
+    return 0;
+}
+static int my_check(mqtt_handle_t *h)
+{
+    int i,r=-1;
+    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]);
+        }
+    }
+
+    return 0;
+}
+static void my_poll(mqtt_handle_t *h, int ms)
+{
+    my_check(h);
+    mg_mgr_poll(&h->mgr, ms);
+}
 
 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(h->para.user.url)) {
+        if (mg_url_is_ssl(mc->ser.server)) {
             struct mg_tls_opts opts = {.ca = mg_unpacked("/certs/ca.pem"),
-                                       .name = mg_url_host(h->para.user.url)};
+                                       .name = mg_url_host(mc->ser.server)};
             mg_tls_init(c, &opts);
         }
     } else if (ev == MG_EV_ERROR) {
@@ -96,9 +213,9 @@ static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
         MG_ERROR(("%p %s", c->fd, (char *) ev_data));
     } else if (ev == MG_EV_MQTT_OPEN) {
         mg_opts_t opts={
-            .client_id = mg_str(h->para.user.url),
-            .user = mg_str(h->para.user.name),
-            .pass = mg_str(h->para.user.pass),
+            .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;
 
@@ -113,7 +230,7 @@ static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
 
     if (ev == MG_EV_ERROR || ev == MG_EV_CLOSE) {
         MG_INFO(("got event %d, stopping...", ev));
-        *(bool *) c->fn_data = true;  // Signal that we're done
+        ((mqtt_conn_t *)(c->fn_data))->c = NULL;  // Signal that we're done
     }
 }
 static void* mqtt_thread(void *arg)
@@ -124,122 +241,13 @@ static void* mqtt_thread(void *arg)
     
     while(th->quit==0) {
         if(h->inited) {
-            
-            my_conn_check(h);
-
             lock_s_hold(LOCK_ID_MQTT);
-            mg_mgr_poll(&h->mgr, 300);
+            my_poll(h, 300);
             lock_s_release(LOCK_ID_MQTT);
         }
     }
     pthread_exit(NULL);
 }
-
-static int my_sub(mg_conn_t *c, int prod_id, int qos)
-{
-    int i;
-    mg_opts_t opts;
-    char temp[512];
-    
-    memset(&opts, 0, sizeof(opts));
-    opts.qos = qos;
-    for(i=0; topic_sub[i]; i++) {
-        sprintf(temp, topic_sub[i], prod_id);
-        opts.topic = mg_str(temp);
-        mg_mqtt_sub(c, &opts);
-    }
-    
-    return 0;
-}
-
-static int find_conn(mqtt_handle_t *h, mqtt_info_t *info)
-{
-    //
-    return 0;
-}
-static int my_conn(mqtt_handle_t *h, mqtt_info_t *info)
-{
-    int r=-1;
-    mqtt_conn_t *conn,*c;
-    
-    h->opts.clean = true,
-    h->opts.qos = h->para.conn.qos,
-    h->opts.topic = mg_str(""),
-    h->opts.version = h->para.conn.ver,
-    h->opts.message = mg_str("bye");
-
-    conn = h->conn;
-    while(info) {
-        if(!conn==NULL) {
-            conn = calloc(1, sizeof(mqtt_conn_t));
-            if(!conn) {
-                LOGE("___my_conn calloc failed\n");
-                break;
-            }
-        }
-
-        conn->c = mg_mqtt_connect(&h->mgr, conn->para.user.url, &conn->opts, mqtt_fn, &conn->isover);
-        if(conn>c) {
-            my_sub(conn>c, h->prod_id, 1);
-        }
-
-        info = info->next;
-        if(!info) {
-            r = 0; break;
-        }
-    }
-
-    return r;
-}
-static int my_disconn(mqtt_handle_t *h, mqtt_conn_t *c)
-{
-    mqtt_conn_t *c1,*c2;
-    
-    c1 = c2 = h->conn;
-    while(c1) {
-        if(c1==c) {
-            mg_mqtt_disconnect(c1, &c1->opts);
-            break;
-        }
-        c1 = c1->next;
-    }
-    h->conn = NULL;
-}
-
-static int my_disconn_all(mqtt_handle_t *h)
-{
-    mqtt_conn_t *c,*conn=h->conn;
-    while(conn) {
-        c = conn;
-        mg_mqtt_disconnect(c, &c->opts);
-        conn = conn->next;
-        free(c);
-    }
-    h->conn = NULL;
-}
-
-
-
-static int my_conn_check(mqtt_handle_t *h)
-{
-    int r=-1;
-    mqtt_conn_t *conn=h->conn;
-
-    while(conn) {
-        if(conn->c && conn->isover) {
-            conn->c = mg_mqtt_connect(&h->mgr, conn->para.user.url, &conn->opts, mqtt_fn, &conn->isover);
-            if(conn>c) {
-                my_sub(conn>c, h->prod_id, 1); r = 0;
-            }
-        }
-        conn = conn->next;
-    }
-    return r;
-}
-
-
-
-
 #endif
 
 
@@ -256,11 +264,12 @@ int mqtt_init(void)
     
     h->dm = &__globalDeviceManage;
     h->prod_id = h->dm->_globalDevInfo.product_id;
-    //h->para.user = ;
 
-    h->para.conn.ver = 4;
-    h->para.conn.qos = 1;
-    
+    h->opts.clean = true,
+    h->opts.qos = 1,
+    h->opts.version = 4,
+    h->opts.topic = mg_str(""),
+    h->opts.message = mg_str("bye");
 
     thread_start(THREAD_ID_MQTT, mqtt_thread, h, 4*MB, 0);
     h->inited = 1;
@@ -290,25 +299,16 @@ int mqtt_deinit(void)
 
 int mqtt_conn(void)
 {
-    int r=-1;
+    int r=0;
 
 #ifdef USE_MQTT
     mqtt_handle_t *h=&mqHandle;
     
     lock_s_hold(LOCK_ID_MQTT);
-    if(!h->inited) {
-        goto quit;
+    if(h->inited) {
+        r = my_conn(h, &h->dm->mqttInfo);
     }
-    
-    if(h->conn) {
-        mg_mqtt_disconnect(h->conn, NULL);
-    }
-    my_conn(h);
-    
-quit:
     lock_s_release(LOCK_ID_MQTT);
-
-    r = h->conn?0:-1;
 #endif
 
     return r;
@@ -323,40 +323,9 @@ int mqtt_disconn(void)
     mqtt_handle_t *h=&mqHandle;
 
     lock_s_hold(LOCK_ID_MQTT);
-    if(!h->inited) {
-        r = -1;
-        goto quit;
+    if(h->inited) {
+        r = my_disconn(h);
     }
-    mg_mqtt_disconnect(h->conn, NULL);
-    r = 0;
-quit:
-    lock_s_release(LOCK_ID_MQTT);
-#endif
-
-    return r;
-}
-
-
-int mqtt_sub(char *topic, int qos)
-{
-    int r=-1;
-
-#ifdef USE_MQTT
-    mg_opts_t opts;
-    mqtt_handle_t *h=&mqHandle;
-
-    lock_s_hold(LOCK_ID_MQTT);
-    if(!h->inited || !h->conn) {
-        r = -1;
-        goto quit;
-    }
-    
-    memset(&opts, 0, sizeof(opts));
-    opts.topic = mg_str(topic);
-    opts.qos = qos;
-    mg_mqtt_sub(h->conn, &opts);
-    r = 0;
-quit:
     lock_s_release(LOCK_ID_MQTT);
 #endif
 
@@ -364,45 +333,18 @@ quit:
 }
 
 
-int mqtt_pub(char *topic, int qos, char *data, int dlen)
-{
-    int r=-1;
-
-#ifdef USE_MQTT
-    mg_opts_t opts;
-    mqtt_handle_t *h=&mqHandle;
-
-    if(!h->inited) {
-        return -1;
-    }
-
-    memset(&opts, 0, sizeof(opts));
-    opts.topic = mg_str(topic);
-    opts.qos = qos;
-    opts.message.buf = data;
-    opts.message.len = dlen;
-    opts.retain = false;
-    while(1) {
-        lock_s_hold(LOCK_ID_MQTT);
-        mg_mqtt_pub(h->conn, &opts);
-        lock_s_release(LOCK_ID_MQTT);
-    }
-#endif
-
-    return r;
-}
-
-
 int mqtt_send(int type, int id, void *data)
 {
     int r=-1;
+    int datalen=0;
+
 
 #ifdef USE_MQTT
     char topic[512];
     char content[4096];
     mqtt_handle_t *h=&mqHandle;
 
-    if(type<0 || type>=MQTT_PUB_MAX || !data) {
+    if(!h->inited || type<0 || type>=MQTT_PUB_MAX || !data) {
         return -1;
     }
 
@@ -461,9 +403,14 @@ int mqtt_send(int type, int id, void *data)
             
         }
         break;
+
+        default:
+        return -1;
     }
 
-    //r = mqtt_pub(topic, );
+    lock_s_hold(LOCK_ID_MQTT);
+    r = my_pub(h, topic, content);
+    lock_s_release(LOCK_ID_MQTT);
 #endif
 
     return r;

+ 0 - 3
pro/src/mqtt.h

@@ -33,9 +33,6 @@ int mqtt_deinit(void);
 int mqtt_conn(void);
 int mqtt_disconn(void);
 
-int mqtt_sub(char *topic, int qos);
-int mqtt_pub(char *topic, int qos, char *data, int dlen);
-
 int mqtt_send(int type, int id, void *data);
 
 #endif

+ 6 - 0
pro/src/smtp.c

@@ -3979,6 +3979,12 @@ smtp_attachment_clear_all(struct smtp *const smtp){
 #include "lock.h"
 #include "sqlite_handle.h"
 
+
+//https://bgithub.xyz/embeddedmz/mailclient-cpp/tree/master/MAIL
+
+
+
+
 #if 1
     #define LOGD            log_d
     #define LOGE            log_e

+ 24 - 0
pro/src/sys.c

@@ -1014,6 +1014,30 @@ quit:
     pthread_exit(NULL);
 }
 
+int sys_is_running(char *prog)
+{
+    FILE *fp;
+    int pid=0;
+    char buf[20]={0};
+    char command[200];
+
+    snprintf(command, sizeof(command), "ps -ef | grep -v grep | grep -w -c %s", prog);
+    fp = popen(command, "r");
+    if (fp == NULL) {
+        LOGE("execute %s failed: %s", command, strerror(errno));
+        return 0;
+    }
+
+    if (fgets(buf, sizeof(buf), fp)) {
+        pid = atoi(buf);
+    }
+    pclose(fp);
+    
+    return (pid)?1:0;
+}
+
+
+
 
 int sys_init(void)
 {

+ 2 - 0
pro/src/sys.h

@@ -20,6 +20,8 @@ enum {
 int sys_init(void);
 int sys_reboot(void);
 int sys_set_factory(void);
+
+int sys_is_running(char *prog);
 int sys_get_path(char *path, char *name);
 int sys_get_net(NetworkInfo_t *para, int ipVersion);
 int sys_set_net(NetworkInfo_t *para, int ipVersion);