|
@@ -1,231 +0,0 @@
|
|
|
-#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;
|
|
|
|
|
-}
|
|
|
|
|
-
|
|
|
|
|
-
|
|
|