Sfoglia il codice sorgente

1、增加MQTT Broker ping连接状态检测,完善MQTT重连机制。2、完善多网卡切换功能。

liyuezong 1 anno fa
parent
commit
22d39c6530
3 ha cambiato i file con 85 aggiunte e 17 eliminazioni
  1. 36 2
      pro/src/Cellular.c
  2. 0 2
      pro/src/log_json.cpp
  3. 49 13
      pro/src/mqtt.c

+ 36 - 2
pro/src/Cellular.c

@@ -58,6 +58,8 @@
 
 
 #define UNLINK_MAX                    3
 #define UNLINK_MAX                    3
 
 
+#define MAX_CONNECTIONS 6  // EC200N最多支持6个连接
+
 int UnLinkCount=0;
 int UnLinkCount=0;
 
 
 enum {
 enum {
@@ -1017,21 +1019,51 @@ static int cell_sock_send(sock_conn_t *conn, void *data, int len)
 
 
     return r;
     return r;
 }
 }
-static int cell_sock_stat(sock_conn_t *conn)
+static int cell_sock_disconn_all(cell_handle_t *h)
 {
 {
     int r;
     int r;
     char buf[100];
     char buf[100];
-    cell_handle_t *h=(cell_handle_t*)conn->h;
 
 
     if(h->stat==AT_STAT_UNINIT) {
     if(h->stat==AT_STAT_UNINIT) {
         return -1;
         return -1;
     }
     }
 
 
+    int count = 0;
+    uint8_t conn_ids[MAX_CONNECTIONS];
+    memset(conn_ids, 0, sizeof(conn_ids));
     r = send_cmd(h, AT_CMD_QISTATE, NULL, buf, sizeof(buf));
     r = send_cmd(h, AT_CMD_QISTATE, NULL, buf, sizeof(buf));
     if(r!=AT_OK) {
     if(r!=AT_OK) {
         return -1;
         return -1;
     }
     }
 
 
+    char *ptr = buf;
+    while ((ptr = strstr(ptr, "+QISTATE: ")) != NULL && count < MAX_CONNECTIONS) {
+        int conn_id;
+        char conn_type[100];
+        char conn_IP[100];
+        if (sscanf(ptr, "+QISTATE: %d,%s,%s,", &conn_id,conn_type,conn_IP) == 1) {
+            conn_ids[count++] = (uint8_t)conn_id;
+            log_i("4G Cellular Found active connection: ID=%d, Type=%s, IP=%s", conn_id,conn_type,conn_IP);
+            ptr += 10; // 跳过已解析的部分
+        }
+        else break;
+    }
+
+    char closebuf[200];
+    for (int i = 0; i < count; i++)
+    {
+        char ids[10] = { 0 };
+        sprintf(ids, "%d", conn_ids[i]);
+        sprintf(closebuf, "=%s,%s", ids, CONNECT_TIMEOUT);
+        r = send_cmd(h, AT_CMD_QICLOSE, closebuf, NULL, 0);
+        if (r != AT_OK)
+        {
+            log_e("4G Cellular active connection: ID=%d close Failed!!!!", conn_ids[i]);
+        }
+        else
+            log_i("4G Cellular active connection: ID=%d close Successful!!!!", conn_ids[i]);
+    }
+
     return 0;
     return 0;
 }
 }
 static int cell_sock_connected(cell_handle_t *h)
 static int cell_sock_connected(cell_handle_t *h)
@@ -1571,6 +1603,7 @@ static void* cellular_thread(void *arg)
         {
         {
             if( h->fd>=0) 
             if( h->fd>=0) 
             {
             {
+                //cell_sock_disconn_all(h);
                 cell_deinit(h);
                 cell_deinit(h);
                 port_deinit(h);
                 port_deinit(h);
             }
             }
@@ -1579,6 +1612,7 @@ static void* cellular_thread(void *arg)
         UpdateCellularManager(&h->all);
         UpdateCellularManager(&h->all);
         sleep(1);
         sleep(1);
     }
     }
+    //cell_sock_disconn_all(h);
     reset_var(h);
     reset_var(h);
     UpdateCellularManager(&h->all);
     UpdateCellularManager(&h->all);
     cell_deinit(h);
     cell_deinit(h);

+ 0 - 2
pro/src/log_json.cpp

@@ -1,2 +0,0 @@
-#include <iostream>
-

+ 49 - 13
pro/src/mqtt.c

@@ -32,8 +32,8 @@
     #define LOGW            printf
     #define LOGW            printf
 #endif
 #endif
 
 
-
-#define CONN_PERIOD         5000    //ms
+#define KEEPLIVE         30    //s
+#define CONN_PERIOD         6000    //ms
 #define SEND_PERIOD         1000    //ms
 #define SEND_PERIOD         1000    //ms
 #define SEND_PERIOD2        5000    //ms
 #define SEND_PERIOD2        5000    //ms
 
 
@@ -59,6 +59,7 @@ typedef struct {
     int            quit;
     int            quit;
 
 
     void           *h;
     void           *h;
+   time_t last_recv_time;  
 }mqtt_conn_t;
 }mqtt_conn_t;
 typedef struct {
 typedef struct {
     int           inited;
     int           inited;
@@ -156,7 +157,12 @@ static int get_account(mqtt_account_t *host, user_account_t *user)
 }
 }
 
 
 
 
-
+static int  get_status(mqtt_conn_t *conn)
+{
+    mqtt_handle_t *h=(mqtt_handle_t*)conn->h;
+    int id=conn->ser.id;
+    return h->dm->mqttInfo.ser[id].status;
+}
 static void set_status(mqtt_conn_t *conn, int flag)
 static void set_status(mqtt_conn_t *conn, int flag)
 {
 {
     mqtt_handle_t *h=(mqtt_handle_t*)conn->h;
     mqtt_handle_t *h=(mqtt_handle_t*)conn->h;
@@ -230,7 +236,7 @@ static int pub_one(mg_conn_t *c, char *topic, char *data)
 
 
 static int my_pub(mqtt_conn_t *conn, char *topic, char *data)
 static int my_pub(mqtt_conn_t *conn, char *topic, char *data)
 {
 {
-    if(conn->c) {
+    if(conn->c && get_status(conn)) {
         pub_one(conn->c, topic, data);
         pub_one(conn->c, topic, data);
     }
     }
     return 0;
     return 0;
@@ -387,12 +393,15 @@ static int conn_one(mqtt_conn_t *conn)
         conn->c = mg_mqtt_connect(&conn->mgr, url, &opts, mqtt_fn, conn);
         conn->c = mg_mqtt_connect(&conn->mgr, url, &opts, mqtt_fn, conn);
         if(conn->c) {
         if(conn->c) {
            //mg_mqtt_bind(conn->c);
            //mg_mqtt_bind(conn->c);
-
-            sprintf(conn->ser.cid, "smartPDU_%s\n", h->prod_id);
-            my_sub(conn->c, h->prod_id);
+           log_i("Trying connect MQTT Broker server ID: %d : %s...",conn->ser.id, conn->ser.server);
+            //my_sub(conn->c, h->prod_id);
             r = 0;
             r = 0;
         }
         }
     }
     }
+    else 
+    {
+       if(get_status(conn)) mg_mqtt_ping(conn->c);
+    }
 
 
     return r;
     return r;
 }
 }
@@ -506,26 +515,49 @@ static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
 
 
         case MG_EV_MQTT_OPEN:
         case MG_EV_MQTT_OPEN:
         {
         {
+            conn->c = c;
             send_once(conn);
             send_once(conn);
             set_status(conn, 1);
             set_status(conn, 1);
+            log_i("MQTT Broker server ID: %d : %s connected successful...",conn->ser.id, conn->ser.server);
+            conn->last_recv_time = time(NULL);  // 初始化接收时间
         }
         }
         break;
         break;
         
         
         case MG_EV_POLL:
         case MG_EV_POLL:
         {
         {
             my_send(conn);
             my_send(conn);
+
+            // 检查连接状态,如果超过KEEPLIVE时间没有收到PINGRESP响应则断开连接.
+            //当前配置为每6s发送一次心跳包,总共5次30s超时未收到心跳包则断开连接。
+            time_t now = time(NULL);
+            if (now - conn->last_recv_time > (time_t)(KEEPLIVE) && get_status(conn)) {
+                log_e("MQTT Broker server ID: %d : %s  PINGRESP timeout! Disconnecting...",conn->ser.id, conn->ser.server);
+                conn->c = NULL;
+                set_status(conn, 0);
+            }
         }
         }
         break;
         break;
 
 
         case MG_EV_MQTT_MSG:
         case MG_EV_MQTT_MSG:
         {
         {
-            struct mg_mqtt_message *mm=(struct mg_mqtt_message*)ev_data;
-            if(mm && mm->topic.buf && mm->data.buf) {
-                my_recv(mm->topic.buf, mm->data.buf);
+            struct mg_mqtt_message *msg=(struct mg_mqtt_message*)ev_data;
+            if(msg && msg->topic.buf && msg->data.buf) {
+                my_recv(msg->topic.buf, msg->data.buf);
             }
             }
         }
         }
         break;
         break;
     
     
+         case MG_EV_MQTT_CMD: 
+         {
+             struct mg_mqtt_message *msg = (struct mg_mqtt_message *)ev_data;
+             if (msg && msg->cmd == MQTT_CMD_PINGRESP)
+             {
+                 conn->last_recv_time = time(NULL); // 更新接收时间
+                 log_i("Received PINGRESP from MQTT Broker server %d : %s !", conn->ser.id, conn->ser.server);
+             }
+        }
+        break;
+
         case MG_EV_ERROR:
         case MG_EV_ERROR:
         {
         {
             //MG_ERROR(("___MG_EV_ERROR, %p %s", c->fd, (char *) ev_data));
             //MG_ERROR(("___MG_EV_ERROR, %p %s", c->fd, (char *) ev_data));
@@ -535,8 +567,12 @@ static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
 
 
         case MG_EV_CLOSE:
         case MG_EV_CLOSE:
         {
         {
-            conn->c = NULL;
-            set_status(conn, 0);
+            if (conn->c && conn->c == c)
+            {
+                conn->c = NULL;
+                set_status(conn, 0);
+                log_e("MQTT Broker server ID: %d : %s connected failed...", conn->ser.id, conn->ser.server);
+            }
         }
         }
         break;
         break;
     }
     }
@@ -640,7 +676,7 @@ static int send_stat(mqtt_conn_t *conn, int type)
     char buf[20],date[40],time[40];
     char buf[20],date[40],time[40];
     mqtt_handle_t *h=(mqtt_handle_t*)conn->h;
     mqtt_handle_t *h=(mqtt_handle_t*)conn->h;
 
 
-    if((!conn->c) || (type<0) || (type>=MQTT_PUB_MAX)) {
+    if((!conn->c)|| !get_status(conn) || (type<0) || (type>=MQTT_PUB_MAX)) {
         return -1;
         return -1;
     }
     }