From e677d6a9545ace654deba7006bd418ffdca343e7 Mon Sep 17 00:00:00 2001 From: Jonathan Naylor Date: Fri, 7 Jul 2023 17:57:36 +0100 Subject: [PATCH] Convert the gateway to log via MQTT. --- Conf.cpp | 62 +++++++------ Conf.h | 97 +++++++++++---------- DAPNETGateway.cpp | 26 +++--- DAPNETGateway.h | 1 - DAPNETGateway.ini | 11 ++- Log.cpp | 135 ++++++----------------------- Log.h | 8 +- MQTTConnection.cpp | 211 +++++++++++++++++++++++++++++++++++++++++++++ MQTTConnection.h | 63 ++++++++++++++ Makefile | 7 +- Utils.cpp | 76 +--------------- Utils.h | 11 +-- Version.h | 4 +- 13 files changed, 423 insertions(+), 289 deletions(-) create mode 100644 MQTTConnection.cpp create mode 100644 MQTTConnection.h diff --git a/Conf.cpp b/Conf.cpp index 77b7e5b..7c51223 100644 --- a/Conf.cpp +++ b/Conf.cpp @@ -1,5 +1,5 @@ /* - * Copyright (C) 2018,2020 by Jonathan Naylor G4KLX + * Copyright (C) 2018,2020,2023 by Jonathan Naylor G4KLX * * This program is free software; you can redistribute it and/or modify * it under the terms of the GNU General Public License as published by @@ -27,10 +27,11 @@ const int BUFFER_SIZE = 500; enum SECTION { - SECTION_NONE, - SECTION_GENERAL, - SECTION_LOG, - SECTION_DAPNET + SECTION_NONE, + SECTION_GENERAL, + SECTION_LOG, + SECTION_MQTT, + SECTION_DAPNET }; CConf::CConf(const std::string& file) : @@ -45,10 +46,11 @@ m_myAddress(), m_myPort(0U), m_daemon(false), m_logDisplayLevel(0U), -m_logFileLevel(0U), -m_logFilePath(), -m_logFileRoot(), -m_logFileRotate(true), +m_logMQTTLevel(0U), +m_mqttAddress("127.0.0.1"), +m_mqttPort(1883U), +m_mqttKeepalive(60U), +m_mqttName("dapnet-gateway"), m_dapnetAddress(), m_dapnetPort(0U), m_dapnetAuthKey(), @@ -80,6 +82,8 @@ bool CConf::read() section = SECTION_GENERAL; else if (::strncmp(buffer, "[Log]", 5U) == 0) section = SECTION_LOG; + else if (::strncmp(buffer, "[MQTT]", 6U) == 0) + section = SECTION_MQTT; else if (::strncmp(buffer, "[DAPNET]", 8U) == 0) section = SECTION_DAPNET; else @@ -153,16 +157,19 @@ bool CConf::read() else if (::strcmp(key, "Daemon") == 0) m_daemon = ::atoi(value) == 1; } else if (section == SECTION_LOG) { - if (::strcmp(key, "FilePath") == 0) - m_logFilePath = value; - else if (::strcmp(key, "FileRoot") == 0) - m_logFileRoot = value; - else if (::strcmp(key, "FileLevel") == 0) - m_logFileLevel = (unsigned int)::atoi(value); + if (::strcmp(key, "MQTTLevel") == 0) + m_logMQTTLevel = (unsigned int)::atoi(value); else if (::strcmp(key, "DisplayLevel") == 0) m_logDisplayLevel = (unsigned int)::atoi(value); - else if (::strcmp(key, "FileRotate") == 0) - m_logFileRotate = ::atoi(value) == 1; + } else if (section == SECTION_MQTT) { + if (::strcmp(key, "Address") == 0) + m_mqttAddress = value; + else if (::strcmp(key, "Port") == 0) + m_mqttPort = (unsigned short)::atoi(value); + else if (::strcmp(key, "Keepalive") == 0) + m_mqttKeepalive = (unsigned int)::atoi(value); + else if (::strcmp(key, "Name") == 0) + m_mqttName = value; } else if (section == SECTION_DAPNET) { if (::strcmp(key, "Address") == 0) m_dapnetAddress = value; @@ -239,24 +246,29 @@ unsigned int CConf::getLogDisplayLevel() const return m_logDisplayLevel; } -unsigned int CConf::getLogFileLevel() const +unsigned int CConf::getLogMQTTLevel() const { - return m_logFileLevel; + return m_logMQTTLevel; } -std::string CConf::getLogFilePath() const +std::string CConf::getMQTTAddress() const { - return m_logFilePath; + return m_mqttAddress; } -std::string CConf::getLogFileRoot() const +unsigned short CConf::getMQTTPort() const { - return m_logFileRoot; + return m_mqttPort; } -bool CConf::getLogFileRotate() const +unsigned int CConf::getMQTTKeepalive() const { - return m_logFileRotate; + return m_mqttKeepalive; +} + +std::string CConf::getMQTTName() const +{ + return m_mqttName; } std::string CConf::getDAPNETAddress() const diff --git a/Conf.h b/Conf.h index 4faf12d..c2cc26c 100644 --- a/Conf.h +++ b/Conf.h @@ -1,5 +1,5 @@ /* - * Copyright (C) 2018 by Jonathan Naylor G4KLX + * Copyright (C) 2018,2023 by Jonathan Naylor G4KLX * * This program is free software; you can redistribute it and/or modify * it under the terms of the GNU General Public License as published by @@ -25,61 +25,66 @@ class CConf { public: - CConf(const std::string& file); - ~CConf(); + CConf(const std::string& file); + ~CConf(); - bool read(); + bool read(); - // The General section - std::string getCallsign() const; - std::vector getWhiteList() const; - std::vector getBlackList() const; - std::string getblacklistRegexfile() const; - std::string getwhitelistRegexfile() const; - std::string getRptAddress() const; - unsigned short getRptPort() const; - std::string getMyAddress() const; - unsigned short getMyPort() const; - bool getDaemon() const; + // The General section + std::string getCallsign() const; + std::vector getWhiteList() const; + std::vector getBlackList() const; + std::string getblacklistRegexfile() const; + std::string getwhitelistRegexfile() const; + std::string getRptAddress() const; + unsigned short getRptPort() const; + std::string getMyAddress() const; + unsigned short getMyPort() const; + bool getDaemon() const; - // The Log section - unsigned int getLogDisplayLevel() const; - unsigned int getLogFileLevel() const; - std::string getLogFilePath() const; - std::string getLogFileRoot() const; - bool getLogFileRotate() const; + // The Log section + unsigned int getLogDisplayLevel() const; + unsigned int getLogMQTTLevel() const; - // The DAPNET section - std::string getDAPNETAddress() const; - unsigned short getDAPNETPort() const; - std::string getDAPNETAuthKey() const; - bool getDAPNETDebug() const; + // The MQTT section + std::string getMQTTAddress() const; + unsigned short getMQTTPort() const; + unsigned int getMQTTKeepalive() const; + std::string getMQTTName() const; + + // The DAPNET section + std::string getDAPNETAddress() const; + unsigned short getDAPNETPort() const; + std::string getDAPNETAuthKey() const; + bool getDAPNETDebug() const; private: - std::string m_file; + std::string m_file; - std::string m_callsign; - std::vector m_whiteList; - std::vector m_blackList; + std::string m_callsign; + std::vector m_whiteList; + std::vector m_blackList; - std::string m_blacklistRegexfile; - std::string m_whitelistRegexfile; - std::string m_rptAddress; - unsigned short m_rptPort; - std::string m_myAddress; - unsigned short m_myPort; - bool m_daemon; + std::string m_blacklistRegexfile; + std::string m_whitelistRegexfile; + std::string m_rptAddress; + unsigned short m_rptPort; + std::string m_myAddress; + unsigned short m_myPort; + bool m_daemon; - unsigned int m_logDisplayLevel; - unsigned int m_logFileLevel; - std::string m_logFilePath; - std::string m_logFileRoot; - bool m_logFileRotate; + unsigned int m_logDisplayLevel; + unsigned int m_logMQTTLevel; - std::string m_dapnetAddress; - unsigned short m_dapnetPort; - std::string m_dapnetAuthKey; - bool m_dapnetDebug; + std::string m_mqttAddress; + unsigned short m_mqttPort; + unsigned int m_mqttKeepalive; + std::string m_mqttName; + + std::string m_dapnetAddress; + unsigned short m_dapnetPort; + std::string m_dapnetAuthKey; + bool m_dapnetDebug; }; #endif diff --git a/DAPNETGateway.cpp b/DAPNETGateway.cpp index 834c5ba..196b0f1 100644 --- a/DAPNETGateway.cpp +++ b/DAPNETGateway.cpp @@ -1,5 +1,5 @@ /* -* Copyright (C) 2018,2020 by Jonathan Naylor G4KLX +* Copyright (C) 2018,2020,2023 by Jonathan Naylor G4KLX * * This program is free software; you can redistribute it and/or modify * it under the terms of the GNU General Public License as published by @@ -16,6 +16,7 @@ * Foundation, Inc., 675 Mass Ave, Cambridge, MA 02139, USA. */ +#include "MQTTConnection.h" #include "DAPNETGateway.h" #include "StopWatch.h" #include "Version.h" @@ -41,6 +42,9 @@ const char* DEFAULT_INI_FILE = "DAPNETGateway.ini"; const char* DEFAULT_INI_FILE = "/etc/DAPNETGateway.ini"; #endif +// In Log.cpp +extern CMQTTConnection* m_mqtt; + #include #include @@ -192,16 +196,6 @@ int CDAPNETGateway::run() } #endif -#if !defined(_WIN32) && !defined(_WIN64) - ret = ::LogInitialise(m_daemon, m_conf.getLogFilePath(), m_conf.getLogFileRoot(), m_conf.getLogFileLevel(), m_conf.getLogDisplayLevel(), m_conf.getLogFileRotate()); -#else - ret = ::LogInitialise(false, m_conf.getLogFilePath(), m_conf.getLogFileRoot(), m_conf.getLogFileLevel(), m_conf.getLogDisplayLevel(), m_conf.getLogFileRotate()); -#endif - if (!ret) { - ::fprintf(stderr, "DAPNETGateway: unable to open the log file\n"); - return 1; - } - #if !defined(_WIN32) && !defined(_WIN64) if (m_daemon) { ::close(STDIN_FILENO); @@ -210,6 +204,16 @@ int CDAPNETGateway::run() } #endif + ::LogInitialise(m_conf.getLogDisplayLevel(), m_conf.getLogMQTTLevel()); + + std::vector> subscriptions; + m_mqtt = new CMQTTConnection(m_conf.getMQTTAddress(), m_conf.getMQTTPort(), m_conf.getMQTTName(), subscriptions, m_conf.getMQTTKeepalive()); + ret = m_mqtt->open(); + if (!ret) { + delete m_mqtt; + return -1; + } + bool debug = m_conf.getDAPNETDebug(); std::string rptAddress = m_conf.getRptAddress(); diff --git a/DAPNETGateway.h b/DAPNETGateway.h index d3dbb17..28ed57a 100644 --- a/DAPNETGateway.h +++ b/DAPNETGateway.h @@ -60,7 +60,6 @@ private: unsigned int calculateCodewords(const CPOCSAGMessage* message) const; void loadSchedule(); bool sendMessage(CPOCSAGMessage* message) const; - }; #endif diff --git a/DAPNETGateway.ini b/DAPNETGateway.ini index 1610a0d..cd3e282 100644 --- a/DAPNETGateway.ini +++ b/DAPNETGateway.ini @@ -13,10 +13,13 @@ Daemon=0 [Log] # Logging levels, 0=No logging DisplayLevel=1 -FileLevel=1 -FilePath=. -FileRoot=DAPNETGateway -FileRotate=1 +MQTTLevel=1 + +[MQTT] +Address=127.0.0.1 +Port=1883 +Keepalive=60 +Name=dapnet-gateway [DAPNET] Address=dapnet.afu.rwth-aachen.de diff --git a/Log.cpp b/Log.cpp index 752601e..4537544 100644 --- a/Log.cpp +++ b/Log.cpp @@ -1,5 +1,5 @@ /* - * Copyright (C) 2015,2016,2020 by Jonathan Naylor G4KLX + * Copyright (C) 2015,2016,2020,2022,2023 by Jonathan Naylor G4KLX * * This program is free software; you can redistribute it and/or modify * it under the terms of the GNU General Public License as published by @@ -17,6 +17,7 @@ */ #include "Log.h" +#include "MQTTConnection.h" #if defined(_WIN32) || defined(_WIN64) #include @@ -32,117 +33,27 @@ #include #include -static unsigned int m_fileLevel = 2U; -static std::string m_filePath; -static std::string m_fileRoot; -static bool m_fileRotate = true; +CMQTTConnection* m_mqtt = NULL; -static FILE* m_fpLog = NULL; -static bool m_daemon = false; +static unsigned int m_mqttLevel = 2U; static unsigned int m_displayLevel = 2U; -static struct tm m_tm; - static char LEVELS[] = " DMIWEF"; -static bool logOpenRotate() +void LogInitialise(unsigned int displayLevel, unsigned int mqttLevel) { - bool status = false; - - if (m_fileLevel == 0U) - return true; - - time_t now; - ::time(&now); - - struct tm* tm = ::gmtime(&now); - - if (tm->tm_mday == m_tm.tm_mday && tm->tm_mon == m_tm.tm_mon && tm->tm_year == m_tm.tm_year) { - if (m_fpLog != NULL) - return true; - } else { - if (m_fpLog != NULL) - ::fclose(m_fpLog); - } - - char filename[200U]; -#if defined(_WIN32) || defined(_WIN64) - ::sprintf(filename, "%s\\%s-%04d-%02d-%02d.log", m_filePath.c_str(), m_fileRoot.c_str(), tm->tm_year + 1900, tm->tm_mon + 1, tm->tm_mday); -#else - ::sprintf(filename, "%s/%s-%04d-%02d-%02d.log", m_filePath.c_str(), m_fileRoot.c_str(), tm->tm_year + 1900, tm->tm_mon + 1, tm->tm_mday); -#endif - - if ((m_fpLog = ::fopen(filename, "a+t")) != NULL) { - status = true; - -#if !defined(_WIN32) && !defined(_WIN64) - if (m_daemon) - dup2(fileno(m_fpLog), fileno(stderr)); -#endif - } - - m_tm = *tm; - - return status; -} - -static bool logOpenNoRotate() -{ - bool status = false; - - if (m_fileLevel == 0U) - return true; - - if (m_fpLog != NULL) - return true; - - char filename[200U]; -#if defined(_WIN32) || defined(_WIN64) - ::sprintf(filename, "%s\\%s.log", m_filePath.c_str(), m_fileRoot.c_str()); -#else - ::sprintf(filename, "%s/%s.log", m_filePath.c_str(), m_fileRoot.c_str()); -#endif - - if ((m_fpLog = ::fopen(filename, "a+t")) != NULL) { - status = true; - -#if !defined(_WIN32) && !defined(_WIN64) - if (m_daemon) - dup2(fileno(m_fpLog), fileno(stderr)); -#endif - } - - return status; -} - -bool LogOpen() -{ - if (m_fileRotate) - return logOpenRotate(); - else - return logOpenNoRotate(); -} - -bool LogInitialise(bool daemon, const std::string& filePath, const std::string& fileRoot, unsigned int fileLevel, unsigned int displayLevel, bool rotate) -{ - m_filePath = filePath; - m_fileRoot = fileRoot; - m_fileLevel = fileLevel; + m_mqttLevel = mqttLevel; m_displayLevel = displayLevel; - m_daemon = daemon; - m_fileRotate = rotate; - - if (m_daemon) - m_displayLevel = 0U; - - return ::LogOpen(); } void LogFinalise() { - if (m_fpLog != NULL) - ::fclose(m_fpLog); + if (m_mqtt != NULL) { + m_mqtt->close(); + delete m_mqtt; + m_mqtt = NULL; + } } void Log(unsigned int level, const char* fmt, ...) @@ -171,22 +82,26 @@ void Log(unsigned int level, const char* fmt, ...) va_end(vl); - if (level >= m_fileLevel && m_fileLevel != 0U) { - bool ret = ::LogOpen(); - if (!ret) - return; - - ::fprintf(m_fpLog, "%s\n", buffer); - ::fflush(m_fpLog); - } + if (m_mqtt != NULL && level >= m_mqttLevel && m_mqttLevel != 0U) + m_mqtt->publish("log", buffer); if (level >= m_displayLevel && m_displayLevel != 0U) { ::fprintf(stdout, "%s\n", buffer); ::fflush(stdout); } - if (level == 6U) { // Fatal - ::fclose(m_fpLog); + if (level == 6U) // Fatal exit(1); +} + +void WriteJSON(const std::string& topLevel, nlohmann::json& json) +{ + if (m_mqtt != NULL) { + nlohmann::json top; + + top[topLevel] = json; + + m_mqtt->publish("json", top.dump()); } } + diff --git a/Log.h b/Log.h index ae95b60..7392164 100644 --- a/Log.h +++ b/Log.h @@ -1,5 +1,5 @@ /* - * Copyright (C) 2015,2016,2020 by Jonathan Naylor G4KLX + * Copyright (C) 2015,2016,2020,2022,2023 by Jonathan Naylor G4KLX * * This program is free software; you can redistribute it and/or modify * it under the terms of the GNU General Public License as published by @@ -21,6 +21,8 @@ #include +#include + #define LogDebug(fmt, ...) Log(1U, fmt, ##__VA_ARGS__) #define LogMessage(fmt, ...) Log(2U, fmt, ##__VA_ARGS__) #define LogInfo(fmt, ...) Log(3U, fmt, ##__VA_ARGS__) @@ -30,7 +32,9 @@ extern void Log(unsigned int level, const char* fmt, ...); -extern bool LogInitialise(bool daemon, const std::string& filePath, const std::string& fileRoot, unsigned int fileLevel, unsigned int displayLevel, bool rotate); +extern void LogInitialise(unsigned int displayLevel, unsigned int mqttLevel); extern void LogFinalise(); +extern void WriteJSON(const std::string& topLevel, nlohmann::json& json); + #endif diff --git a/MQTTConnection.cpp b/MQTTConnection.cpp new file mode 100644 index 0000000..fa951e5 --- /dev/null +++ b/MQTTConnection.cpp @@ -0,0 +1,211 @@ +/* + * Copyright (C) 2022,2023 by Jonathan Naylor G4KLX + * + * This program is free software; you can redistribute it and/or modify + * it under the terms of the GNU General Public License as published by + * the Free Software Foundation; either version 2 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU General Public License for more details. + * + * You should have received a copy of the GNU General Public License + * along with this program; if not, write to the Free Software + * Foundation, Inc., 675 Mass Ave, Cambridge, MA 02139, USA. + */ + +#include "MQTTConnection.h" + +#include +#include +#include + + +CMQTTConnection::CMQTTConnection(const std::string& host, unsigned short port, const std::string& name, const std::vector>& subs, unsigned int keepalive, MQTT_QOS qos) : +m_host(host), +m_port(port), +m_name(name), +m_subs(subs), +m_keepalive(keepalive), +m_qos(qos), +m_mosq(NULL), +m_connected(false) +{ + assert(!host.empty()); + assert(port > 0U); + assert(!name.empty()); + assert(keepalive >= 5U); + + ::mosquitto_lib_init(); +} + +CMQTTConnection::~CMQTTConnection() +{ + ::mosquitto_lib_cleanup(); +} + +bool CMQTTConnection::open() +{ + m_mosq = ::mosquitto_new(m_name.c_str(), true, this); + if (m_mosq == NULL){ + ::fprintf(stderr, "MQTT Error newing: Out of memory.\n"); + return false; + } + + ::mosquitto_connect_callback_set(m_mosq, onConnect); + ::mosquitto_subscribe_callback_set(m_mosq, onSubscribe); + ::mosquitto_message_callback_set(m_mosq, onMessage); + ::mosquitto_disconnect_callback_set(m_mosq, onDisconnect); + + int rc = ::mosquitto_connect(m_mosq, m_host.c_str(), m_port, m_keepalive); + if (rc != MOSQ_ERR_SUCCESS) { + ::mosquitto_destroy(m_mosq); + m_mosq = NULL; + ::fprintf(stderr, "MQTT Error connecting: %s\n", ::mosquitto_strerror(rc)); + return false; + } + + rc = ::mosquitto_loop_start(m_mosq); + if (rc != MOSQ_ERR_SUCCESS) { + ::mosquitto_disconnect(m_mosq); + ::mosquitto_destroy(m_mosq); + m_mosq = NULL; + ::fprintf(stderr, "MQTT Error loop starting: %s\n", ::mosquitto_strerror(rc)); + return false; + } + + return true; +} + +bool CMQTTConnection::publish(const char* topic, const char* text) +{ + assert(topic != NULL); + assert(text != NULL); + + return publish(topic, (unsigned char*)text, ::strlen(text)); +} + +bool CMQTTConnection::publish(const char* topic, const std::string& text) +{ + assert(topic != NULL); + + return publish(topic, (unsigned char*)text.c_str(), text.size()); +} + +bool CMQTTConnection::publish(const char* topic, const unsigned char* data, unsigned int len) +{ + assert(topic != NULL); + assert(data != NULL); + + if (!m_connected) + return false; + + if (::strchr(topic, '/') == NULL) { + char topicEx[100U]; + ::sprintf(topicEx, "%s/%s", m_name.c_str(), topic); + + int rc = ::mosquitto_publish(m_mosq, NULL, topicEx, len, data, static_cast(m_qos), false); + if (rc != MOSQ_ERR_SUCCESS) { + ::fprintf(stderr, "MQTT Error publishing: %s\n", ::mosquitto_strerror(rc)); + return false; + } + } else { + int rc = ::mosquitto_publish(m_mosq, NULL, topic, len, data, static_cast(m_qos), false); + if (rc != MOSQ_ERR_SUCCESS) { + ::fprintf(stderr, "MQTT Error publishing: %s\n", ::mosquitto_strerror(rc)); + return false; + } + } + + return true; +} + +void CMQTTConnection::close() +{ + if (m_mosq != NULL) { + ::mosquitto_disconnect(m_mosq); + ::mosquitto_destroy(m_mosq); + m_mosq = NULL; + } +} + +void CMQTTConnection::onConnect(mosquitto* mosq, void* obj, int rc) +{ + assert(mosq != NULL); + assert(obj != NULL); + + ::fprintf(stdout, "MQTT: on_connect: %s\n", ::mosquitto_connack_string(rc)); + if (rc != 0) { + ::mosquitto_disconnect(mosq); + return; + } + + CMQTTConnection* p = static_cast(obj); + p->m_connected = true; + + for (std::vector>::const_iterator it = p->m_subs.cbegin(); it != p->m_subs.cend(); ++it) { + std::string topic = (*it).first; + + if (topic.find_first_of('/') == std::string::npos) { + char topicEx[100U]; + ::sprintf(topicEx, "%s/%s", p->m_name.c_str(), topic.c_str()); + + rc = ::mosquitto_subscribe(mosq, NULL, topicEx, static_cast(p->m_qos)); + if (rc != MOSQ_ERR_SUCCESS) { + ::fprintf(stderr, "MQTT: error subscribing to %s - %s\n", topicEx, ::mosquitto_strerror(rc)); + ::mosquitto_disconnect(mosq); + } + } else { + rc = ::mosquitto_subscribe(mosq, NULL, topic.c_str(), static_cast(p->m_qos)); + if (rc != MOSQ_ERR_SUCCESS) { + ::fprintf(stderr, "MQTT: error subscribing to %s - %s\n", topic.c_str(), ::mosquitto_strerror(rc)); + ::mosquitto_disconnect(mosq); + } + } + } +} + +void CMQTTConnection::onSubscribe(mosquitto* mosq, void* obj, int mid, int qosCount, const int* grantedQOS) +{ + assert(mosq != NULL); + assert(obj != NULL); + assert(grantedQOS != NULL); + + for (int i = 0; i < qosCount; i++) + ::fprintf(stdout, "MQTT: on_subscribe: %d:%d\n", i, grantedQOS[i]); +} + +void CMQTTConnection::onMessage(mosquitto* mosq, void* obj, const mosquitto_message* message) +{ + assert(mosq != NULL); + assert(obj != NULL); + assert(message != NULL); + + CMQTTConnection* p = static_cast(obj); + + for (std::vector>::const_iterator it = p->m_subs.cbegin(); it != p->m_subs.cend(); ++it) { + std::string topic = (*it).first; + + char topicEx[100U]; + ::sprintf(topicEx, "%s/%s", p->m_name.c_str(), topic.c_str()); + + if (::strcmp(topicEx, message->topic) == 0) { + (*it).second((unsigned char*)message->payload, message->payloadlen); + break; + } + } +} + +void CMQTTConnection::onDisconnect(mosquitto* mosq, void* obj, int rc) +{ + assert(mosq != NULL); + assert(obj != NULL); + + ::fprintf(stdout, "MQTT: on_disconnect: %s\n", ::mosquitto_reason_string(rc)); + + CMQTTConnection* p = static_cast(obj); + p->m_connected = false; +} + diff --git a/MQTTConnection.h b/MQTTConnection.h new file mode 100644 index 0000000..8fc98ce --- /dev/null +++ b/MQTTConnection.h @@ -0,0 +1,63 @@ +/* + * Copyright (C) 2022,2023 by Jonathan Naylor G4KLX + * + * This program is free software; you can redistribute it and/or modify + * it under the terms of the GNU General Public License as published by + * the Free Software Foundation; either version 2 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU General Public License for more details. + * + * You should have received a copy of the GNU General Public License + * along with this program; if not, write to the Free Software + * Foundation, Inc., 675 Mass Ave, Cambridge, MA 02139, USA. + */ + +#if !defined(MQTTPUBLISHER_H) +#define MQTTPUBLISHER_H + +#include + +#include +#include + +enum MQTT_QOS { + MQTT_QOS_AT_MODE_ONCE = 0U, + MQTT_QOS_AT_LEAST_ONCE = 1U, + MQTT_QOS_EXACTLY_ONCE = 2U +}; + +class CMQTTConnection { +public: + CMQTTConnection(const std::string& host, unsigned short port, const std::string& name, const std::vector>& subs, unsigned int keepalive, MQTT_QOS qos = MQTT_QOS_EXACTLY_ONCE); + ~CMQTTConnection(); + + bool open(); + + bool publish(const char* topic, const char* text); + bool publish(const char* topic, const std::string& text); + bool publish(const char* topic, const unsigned char* data, unsigned int len); + + void close(); + +private: + std::string m_host; + unsigned short m_port; + std::string m_name; + std::vector> m_subs; + unsigned int m_keepalive; + MQTT_QOS m_qos; + mosquitto* m_mosq; + bool m_connected; + + static void onConnect(mosquitto* mosq, void* obj, int rc); + static void onSubscribe(mosquitto* mosq, void* obj, int mid, int qosCount, const int* grantedQOS); + static void onMessage(mosquitto* mosq, void* obj, const mosquitto_message* message); + static void onDisconnect(mosquitto* mosq, void* obj, int rc); +}; + +#endif + diff --git a/Makefile b/Makefile index a789596..e311d90 100644 --- a/Makefile +++ b/Makefile @@ -1,10 +1,11 @@ CC = cc CXX = c++ -CFLAGS = -g -O3 -Wall -DHAVE_LOG_H -std=c++0x -pthread -LIBS = -lm -lpthread +CFLAGS = -g -O3 -Wall -std=c++0x -pthread +LIBS = -lm -lpthread -lmosquitto LDFLAGS = -g -OBJECTS = Conf.o DAPNETGateway.o DAPNETNetwork.o Log.o POCSAGMessage.o POCSAGNetwork.o StopWatch.o TCPSocket.o Thread.o Timer.o UDPSocket.o Utils.o REGEX.o +OBJECTS = Conf.o DAPNETGateway.o DAPNETNetwork.o Log.o MQTTConnection.o POCSAGMessage.o POCSAGNetwork.o REGEX.o StopWatch.o \ + TCPSocket.o Thread.o Timer.o UDPSocket.o Utils.o all: DAPNETGateway diff --git a/Utils.cpp b/Utils.cpp index 49ded13..7276d60 100644 --- a/Utils.cpp +++ b/Utils.cpp @@ -1,5 +1,5 @@ /* - * Copyright (C) 2009,2014,2015,2016 Jonathan Naylor, G4KLX + * Copyright (C) 2009,2014,2015,2016,2023 Jonathan Naylor, G4KLX * * This program is free software; you can redistribute it and/or modify * it under the terms of the GNU General Public License as published by @@ -70,77 +70,3 @@ void CUtils::dump(int level, const std::string& title, const unsigned char* data } } -void CUtils::dump(const std::string& title, const bool* bits, unsigned int length) -{ - assert(bits != NULL); - - dump(2U, title, bits, length); -} - -void CUtils::dump(int level, const std::string& title, const bool* bits, unsigned int length) -{ - assert(bits != NULL); - - unsigned char bytes[100U]; - unsigned int nBytes = 0U; - for (unsigned int n = 0U; n < length; n += 8U, nBytes++) - bitsToByteBE(bits + n, bytes[nBytes]); - - dump(level, title, bytes, nBytes); -} - -void CUtils::byteToBitsBE(unsigned char byte, bool* bits) -{ - assert(bits != NULL); - - bits[0U] = (byte & 0x80U) == 0x80U; - bits[1U] = (byte & 0x40U) == 0x40U; - bits[2U] = (byte & 0x20U) == 0x20U; - bits[3U] = (byte & 0x10U) == 0x10U; - bits[4U] = (byte & 0x08U) == 0x08U; - bits[5U] = (byte & 0x04U) == 0x04U; - bits[6U] = (byte & 0x02U) == 0x02U; - bits[7U] = (byte & 0x01U) == 0x01U; -} - -void CUtils::byteToBitsLE(unsigned char byte, bool* bits) -{ - assert(bits != NULL); - - bits[0U] = (byte & 0x01U) == 0x01U; - bits[1U] = (byte & 0x02U) == 0x02U; - bits[2U] = (byte & 0x04U) == 0x04U; - bits[3U] = (byte & 0x08U) == 0x08U; - bits[4U] = (byte & 0x10U) == 0x10U; - bits[5U] = (byte & 0x20U) == 0x20U; - bits[6U] = (byte & 0x40U) == 0x40U; - bits[7U] = (byte & 0x80U) == 0x80U; -} - -void CUtils::bitsToByteBE(const bool* bits, unsigned char& byte) -{ - assert(bits != NULL); - - byte = bits[0U] ? 0x80U : 0x00U; - byte |= bits[1U] ? 0x40U : 0x00U; - byte |= bits[2U] ? 0x20U : 0x00U; - byte |= bits[3U] ? 0x10U : 0x00U; - byte |= bits[4U] ? 0x08U : 0x00U; - byte |= bits[5U] ? 0x04U : 0x00U; - byte |= bits[6U] ? 0x02U : 0x00U; - byte |= bits[7U] ? 0x01U : 0x00U; -} - -void CUtils::bitsToByteLE(const bool* bits, unsigned char& byte) -{ - assert(bits != NULL); - - byte = bits[0U] ? 0x01U : 0x00U; - byte |= bits[1U] ? 0x02U : 0x00U; - byte |= bits[2U] ? 0x04U : 0x00U; - byte |= bits[3U] ? 0x08U : 0x00U; - byte |= bits[4U] ? 0x10U : 0x00U; - byte |= bits[5U] ? 0x20U : 0x00U; - byte |= bits[6U] ? 0x40U : 0x00U; - byte |= bits[7U] ? 0x80U : 0x00U; -} diff --git a/Utils.h b/Utils.h index ade28c0..752d4dc 100644 --- a/Utils.h +++ b/Utils.h @@ -1,5 +1,5 @@ /* - * Copyright (C) 2009,2014,2015 by Jonathan Naylor, G4KLX + * Copyright (C) 2009,2014,2015,2023 by Jonathan Naylor, G4KLX * * This program is free software; you can redistribute it and/or modify * it under the terms of the GNU General Public License as published by @@ -21,15 +21,6 @@ public: static void dump(const std::string& title, const unsigned char* data, unsigned int length); static void dump(int level, const std::string& title, const unsigned char* data, unsigned int length); - static void dump(const std::string& title, const bool* bits, unsigned int length); - static void dump(int level, const std::string& title, const bool* bits, unsigned int length); - - static void byteToBitsBE(unsigned char byte, bool* bits); - static void byteToBitsLE(unsigned char byte, bool* bits); - - static void bitsToByteBE(const bool* bits, unsigned char& byte); - static void bitsToByteLE(const bool* bits, unsigned char& byte); - private: }; diff --git a/Version.h b/Version.h index e41719c..acbd087 100644 --- a/Version.h +++ b/Version.h @@ -1,5 +1,5 @@ /* - * Copyright (C) 2018 by Jonathan Naylor G4KLX + * Copyright (C) 2018,2020,2023 by Jonathan Naylor G4KLX * * This program is free software; you can redistribute it and/or modify * it under the terms of the GNU General Public License as published by @@ -19,6 +19,6 @@ #if !defined(VERSION_H) #define VERSION_H -const char* VERSION = "20201031"; +const char* VERSION = "20230707"; #endif