guohui лет назад: 2
Родитель
Сommit
1edab10def
6 измененных файлов с 4334 добавлено и 3 удалено
  1. 5 3
      pro/CMakeLists.txt
  2. 231 0
      pro/src/mqtt.c
  3. 25 0
      pro/src/mqtt.h
  4. 4044 0
      pro/src/smtp.c
  5. 28 0
      pro/src/smtp.h
  6. 1 0
      pro/src/thread.h

+ 5 - 3
pro/CMakeLists.txt

@@ -52,7 +52,7 @@ set(PDU_SRC_LIST
 	${PROJECT_SOURCE_DIR}/easylogger/src/elog_buf.c	
 	${PROJECT_SOURCE_DIR}/easylogger/src/elog_utils.c
 
-    ${PROJECT_SOURCE_DIR}/ini/dictionary.c
+        ${PROJECT_SOURCE_DIR}/ini/dictionary.c
 	${PROJECT_SOURCE_DIR}/ini/iniparser.c
 
 	${PROJECT_SOURCE_DIR}/mongoose/mongoose.c
@@ -77,9 +77,11 @@ set(PDU_SRC_LIST
 	${PROJECT_SOURCE_DIR}/src/xlist.c
 	${PROJECT_SOURCE_DIR}/src/thread.c
 	${PROJECT_SOURCE_DIR}/src/cascade.c
-    ${PROJECT_SOURCE_DIR}/src/cascade_slave.c
+        ${PROJECT_SOURCE_DIR}/src/cascade_slave.c
 	${PROJECT_SOURCE_DIR}/src/upg.c
 	${PROJECT_SOURCE_DIR}/src/file.c
+	${PROJECT_SOURCE_DIR}/src/smtp.c
+	${PROJECT_SOURCE_DIR}/src/mqtt.c
 	${PROJECT_SOURCE_DIR}/src/netswitch.c
 	${PROJECT_SOURCE_DIR}/src/main.c
 
@@ -100,7 +102,7 @@ set(CMAKE_CXX_FLAGS "${CMAKE_CXX_FLAGS} ${NO_WARNNING_FLAGS} -Werror")
 
 
 add_executable(smartPDU ${PDU_SRC_LIST})
-target_link_libraries(smartPDU -lpcre -lsqlite3 -lmodbus -lc -ldl -lpthread -lnetsnmp -lnetsnmpagent -lrt)
+target_link_libraries(smartPDU -lpcre -lsqlite3 -lssl -lcrypto -lmodbus -lc -ldl -lpthread -lnetsnmp -lnetsnmpagent -lrt)
 
 
 set(SRC_DIR ${PROJECT_BINARY_DIR})

+ 231 - 0
pro/src/mqtt.c

@@ -0,0 +1,231 @@
+#include "mqtt.h"
+#include "thread.h"
+#include "mongoose.h"
+#include "elog.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
+
+#define MQTT_SERVER_URL     "XXXXXXXX"
+
+typedef struct mg_mgr mgr_t;
+typedef struct mg_mqtt_opts mg_opts_t;
+typedef struct mg_connection mg_conn_t;
+typedef struct {
+    mgr_t         mgr;
+    mg_opts_t     opts;
+    mg_conn_t     *conn;
+    conn_para_t   para;
+    
+    pthread_mutex_t mutex;
+    int           inited;
+    bool          isover;
+    
+}mqtt_handle_t;
+static mqtt_handle_t mqHandle;
+
+static void mqtt_fn(struct mg_connection *c, int ev, void *ev_data)
+{
+    mqtt_handle_t *h=&mqHandle;
+    
+    if (ev == MG_EV_OPEN) {
+        // c->is_hexdumping = 1;
+    } else if (ev == MG_EV_CONNECT) {
+        if (mg_url_is_ssl(h->para.url)) {
+            struct mg_tls_opts opts = {.ca = mg_unpacked("/certs/ca.pem"),
+                                     .name = mg_url_host(h->para.url)};
+            mg_tls_init(c, &opts);
+        }
+    } else if (ev == MG_EV_ERROR) {
+        // On error, log error message
+        MG_ERROR(("%p %s", c->fd, (char *) ev_data));
+    } else if (ev == MG_EV_MQTT_OPEN) {
+        
+    } else if (ev == MG_EV_MQTT_MSG) {
+        // When we receive MQTT message, print it
+        struct mg_mqtt_message *mm = (struct mg_mqtt_message *) ev_data;
+        MG_INFO(("Received on %.*s : %.*s", (int) mm->topic.len, mm->topic.buf,
+             (int) mm->data.len, mm->data.buf));
+    } else if (ev == MG_EV_POLL && c->data[0] == 'X') {
+        
+    }
+
+    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
+    }
+}
+static void* mqtt_thread(void *arg)
+{
+    int r;
+    thread_handle_t *h=(thread_handle_t*)arg;
+    mqtt_handle_t   *mh=(mqtt_handle_t*)h->arg;
+    
+    while(h->quit==0) {
+        pthread_mutex_lock(&mh->mutex);
+        if (mh->isover) {
+            mh->conn = mg_mqtt_connect(&mh->mgr, mh->para.url, &mh->opts, mqtt_fn, &mh->isover);
+        }
+        else {
+            mg_mgr_poll(&mh->mgr, 300);
+        }
+        pthread_mutex_unlock(&mh->mutex);
+        
+        if(mh->isover) usleep(1000);
+    }
+    pthread_exit(NULL);
+}
+//////////////////////////////////////////////////////////////////////////
+int mqtt_init(void)
+{
+    int r;
+    mqtt_handle_t *h=&mqHandle;
+    
+    memset(h, 9, sizeof(mqtt_handle_t));
+    mg_mgr_init(&h->mgr);
+    
+    h->para.ver = 4;
+    h->para.qos = 1;
+    strcpy(h->para.url, MQTT_SERVER_URL);
+    
+    r = pthread_mutex_init(&h->mutex, NULL);
+    if(r) {
+        LOGE("___ mqtt mutex init failed\n");
+        return -1;
+    }
+    
+    thread_start(THREAD_ID_MQTT, mqtt_thread, h, 4*MB, 0);
+    h->inited = 1;
+    
+    return 0;
+}
+
+
+int mqtt_deinit(void)
+{
+    mqtt_handle_t *h=&mqHandle;
+    
+    thread_stop(THREAD_ID_MQTT);
+    mg_mgr_free(&h->mgr);
+    pthread_mutex_destroy(&h->mutex);
+    h->inited = 0;
+    
+    return 0;
+}
+
+int mqtt_set(conn_para_t *para)
+{
+    mqtt_handle_t *h=&mqHandle;
+    
+    if(!para || para->ver<3 || para->ver>5) {
+        return -1;
+    }
+    h->para = *para;
+    
+    return 0;
+}
+
+
+int mqtt_conn(void)
+{
+    mqtt_handle_t *h=&mqHandle;
+    
+    pthread_mutex_lock(&h->mutex);
+    if(!h->inited) {
+        goto quit;
+    }
+    
+    memset(&h->opts, 0, sizeof(h->opts));
+    h->opts.clean = true,
+    h->opts.qos = h->para.qos,
+    h->opts.topic = mg_str(""),
+    h->opts.version = h->para.ver,
+    h->opts.message = mg_str("bye");
+    
+    if(h->conn) {
+        mg_mqtt_disconnect(h->conn, NULL);
+    }
+    
+    h->isover = false;
+    h->conn = mg_mqtt_connect(&h->mgr, h->para.url, &h->opts, mqtt_fn, &h->isover);
+    
+quit:
+    pthread_mutex_unlock(&h->mutex);
+    return h->conn?0:-1;
+}
+
+
+int mqtt_disconn(void)
+{
+    int r=0;
+    mqtt_handle_t *h=&mqHandle;
+
+    pthread_mutex_lock(&h->mutex);
+    if(!h->inited) {
+        r = -1;
+        goto quit;
+    }
+    mg_mqtt_disconnect(h->conn, NULL);
+    
+quit:
+    pthread_mutex_unlock(&h->mutex);
+    return r;
+}
+
+
+int mqtt_sub(char *topic, int qos)
+{
+    int r=0;
+    mg_opts_t opts;
+    mqtt_handle_t *h=&mqHandle;
+
+    pthread_mutex_lock(&h->mutex);
+    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);
+    
+quit:
+    pthread_mutex_unlock(&h->mutex);
+    return r;
+}
+
+
+int mqtt_pub(char *topic, int qos, char *data, int dlen)
+{
+    int r=0;
+    mg_opts_t opts;
+    mqtt_handle_t *h=&mqHandle;
+
+    pthread_mutex_lock(&h->mutex);
+    if(!h->inited || !h->conn) {
+        r = -1;
+        goto quit;
+    }
+    
+    memset(&opts, 0, sizeof(opts));
+    opts.topic = mg_str(topic);
+    opts.qos = qos;
+    opts.message.buf = data;
+    opts.message.len = dlen;
+    opts.retain = false;
+    mg_mqtt_pub(h->conn, &opts);
+    
+quit:
+    pthread_mutex_unlock(&h->mutex);
+    return r;
+}
+
+

+ 25 - 0
pro/src/mqtt.h

@@ -0,0 +1,25 @@
+#ifndef __MQTT_Hx__
+#define __MQTT_Hx__
+
+#include <stdint.h>
+
+
+typedef struct {
+    uint8_t qos;
+    uint8_t ver;
+    uint8_t clean;
+    char    url[1024];
+}conn_para_t;
+
+
+int mqtt_init(void);
+int mqtt_deinit(void);
+
+int mqtt_set(conn_para_t *para);
+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);
+
+#endif

Разница между файлами не показана из-за своего большого размера
+ 4044 - 0
pro/src/smtp.c


+ 28 - 0
pro/src/smtp.h

@@ -0,0 +1,28 @@
+#ifndef __SMTP_H__
+#define __SMTP_H__
+
+#include <stdint.h>
+
+typedef struct {
+    char server[128];
+    char port[10];
+    
+    char user[128];
+    char passwd[128];
+    
+    char name_fr[128];
+    char mail_fr[128];
+    
+    char name_to[128];
+    char mail_to[128];
+    
+    char subject[128];
+    int  bodylen;
+    char *body;
+}smtp_data_t;
+
+
+int smtp_send(smtp_data_t *sd);
+
+#endif /* SMTP_H */
+

+ 1 - 0
pro/src/thread.h

@@ -27,6 +27,7 @@ enum {
     THREAD_ID_BREAKER,
     THREAD_ID_SWITCH,
     THREAD_ID_BREAKER_SCANNER,
+    THREAD_ID_MQTT,
     
     THREAD_ID_MAX=30
 };