Parcourir la source

1、修复MQTT报警多个服务器只推送一个的问题。2、调整订阅指令发送方式,改为连接成功后订阅。

liyuezong il y a 1 an
Parent
commit
ae1c438956
1 fichiers modifiés avec 53 ajouts et 27 suppressions
  1. 53 27
      pro/src/mqtt.c

+ 53 - 27
pro/src/mqtt.c

@@ -58,15 +58,15 @@ typedef struct {
     pthread_t      tid;
     int            quit;
 
-    void           *h;
-   time_t last_recv_time;  
+    void *h;
+    time_t last_recv_time;
 }mqtt_conn_t;
 typedef struct {
     int           inited;
     char      prod_id[32];
     
     mqtt_conn_t   conn[MQTT_SER_MAX];
-    handle_t      list;
+    handle_t list[MQTT_SER_MAX];
     mg_timer_t    *timer;
 
     GlobalDeviceManager *dm;
@@ -161,7 +161,7 @@ static int  get_status(mqtt_conn_t *conn)
 {
     mqtt_handle_t *h=(mqtt_handle_t*)conn->h;
     int id=conn->ser.id;
-    if (id >= 0 && id < MQTT_SER_MAX)
+    if (id >= 0 && id < MQTT_SER_MAX && h)
     {
         return h->dm->mqttInfo.ser[id].status;
     }
@@ -171,7 +171,7 @@ static void set_status(mqtt_conn_t *conn, int flag)
 {
     mqtt_handle_t *h=(mqtt_handle_t*)conn->h;
     int id=conn->ser.id;
-    if (id >= 0 && id < MQTT_SER_MAX)
+    if (id >= 0 && id < MQTT_SER_MAX && h)
     {
         h->dm->mqttInfo.ser[id].status = flag;
     }
@@ -398,10 +398,11 @@ static int conn_one(mqtt_conn_t *conn)
         get_url(&conn->ser, url);
 
         conn->c = mg_mqtt_connect(&conn->mgr, url, &opts, mqtt_fn, conn);
-        if(conn->c) {
-           //mg_mqtt_bind(conn->c);
-           log_i("Trying connect MQTT Broker server ID: %d : %s...",conn->ser.id, conn->ser.server);
-            //my_sub(conn->c, h->prod_id);
+        if (conn->c)
+        {
+            // mg_mqtt_bind(conn->c);
+            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;
         }
     }
@@ -444,11 +445,16 @@ static int my_send(mqtt_conn_t *conn)
     mqtt_handle_t *h=(mqtt_handle_t*)conn->h;
     list_node_t *ln=NULL;
 
-    r = xlist_take_node(h->list, &ln, 0);
-    if(r==0) {
-        mqtt_pkt_t *pkt=(mqtt_pkt_t*)ln->data.buf;
-        my_pub(conn, pkt->topic, pkt->content);
-        xlist_back_node(h->list, ln);
+    int id=conn->ser.id;
+    if (id >= 0 && id < MQTT_SER_MAX)
+    {
+        r = xlist_take_node(h->list[id], &ln, 0);
+        if (r == 0)
+        {
+            mqtt_pkt_t *pkt = (mqtt_pkt_t *)ln->data.buf;
+            my_pub(conn, pkt->topic, pkt->content);
+            xlist_back_node(h->list[id], ln);
+        }
     }
 
     return r;
@@ -529,6 +535,8 @@ static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
             conn->c = c;
             send_once(conn);
             set_status(conn, 1);
+            mqtt_handle_t *h=(mqtt_handle_t*)conn->h;
+            if(h)my_sub(conn->c, h->prod_id);
             log_i("MQTT Broker server ID: %d : %s connected successful...",conn->ser.id, conn->ser.server);
             conn->last_recv_time = time(NULL);  // 初始化接收时间
         }
@@ -554,6 +562,7 @@ static void mqtt_fn(mg_conn_t *c, int ev, void *ev_data)
             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);
+                log_i("Received Topic from MQTT Broker server %d : %s !", conn->ser.id, conn->ser.server);
             }
         }
         break;
@@ -645,11 +654,14 @@ int mqtt_init(void)
     strcpy( h->prod_id,h->dm->_globalDevInfo.product.number);
     //sys_get_chip_id(&h->prod_id);
 
-    list_cfg_t lc;
-    lc.mode = LIST_FULL_FIFO;
-    lc.max = 100;
-    lc.log = 0;
-    h->list = xlist_init(&lc);
+    for (int i = 0; i < MQTT_SER_MAX; i++)
+    {
+        list_cfg_t lc;
+        lc.mode = LIST_FULL_FIFO;
+        lc.max = 100;
+        lc.log = 0;
+        h->list[i] = xlist_init(&lc);
+    }
 
     dev_mqtt_init(h->dm->db, &h->dm->mqttInfo);
     thread_start(THREAD_ID_MQTT, mqtt_thread, h);
@@ -666,10 +678,14 @@ int mqtt_deinit(void)
 
 #ifdef USE_MQTT
     mqtt_handle_t *h=&mqHandle;
-    
+
     thread_stop(THREAD_ID_MQTT);
-    xlist_free(h->list);
-    h->inited = 0; r = 0;
+    for (int i = 0; i < MQTT_SER_MAX; i++)
+    {
+        xlist_free(h->list[i]);
+    }
+    h->inited = 0;
+    r = 0;
 #endif
 
     return r;
@@ -1096,12 +1112,22 @@ static int post_alarm(int type, int subtype, char *content, alarm_para_t *para)
         break;
     }
 
-    mqtt_pkt_t pkt;
+    if (h->inited)
+    {
+        for (int i = 0; i < MQTT_SER_MAX; i++)
+        {
+            if (get_status(&(h->conn[i])))
+            {
+                mqtt_pkt_t pkt;
+
+                pkt.id = pub_type;
+                snprintf(pkt.topic, sizeof(pkt.topic), "%s", topic);
+                snprintf(pkt.content, sizeof(pkt.content), "%s", json_str);
+                xlist_append(h->list[i], 0, &pkt, sizeof(pkt));
+            }
+        }
+    }
 
-    pkt.id = pub_type;
-    snprintf(pkt.topic, sizeof(pkt.topic), "%s", topic);
-    snprintf(pkt.content, sizeof(pkt.content), "%s", json_str);
-    r = xlist_append(h->list, 0, &pkt, sizeof(pkt));
     cJSON_free(json_str);
 #endif