|
|
@@ -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;
|
|
|
+}
|
|
|
+
|
|
|
+
|