extract request related parts
This commit is contained in:
@@ -3,6 +3,7 @@ add_definitions(-Wconversion -Wno-unused-parameter)
|
|||||||
set(ebusd_SOURCES
|
set(ebusd_SOURCES
|
||||||
bushandler.h bushandler.cpp
|
bushandler.h bushandler.cpp
|
||||||
datahandler.h datahandler.cpp
|
datahandler.h datahandler.cpp
|
||||||
|
request.h request.cpp
|
||||||
network.h network.cpp
|
network.h network.cpp
|
||||||
mainloop.h mainloop.cpp
|
mainloop.h mainloop.cpp
|
||||||
scan.h scan.cpp
|
scan.h scan.cpp
|
||||||
|
|||||||
+28
-4
@@ -31,6 +31,7 @@
|
|||||||
#include <map>
|
#include <map>
|
||||||
#include <vector>
|
#include <vector>
|
||||||
#include "ebusd/mainloop.h"
|
#include "ebusd/mainloop.h"
|
||||||
|
#include "ebusd/network.h"
|
||||||
#include "lib/utils/log.h"
|
#include "lib/utils/log.h"
|
||||||
#include "lib/utils/httpclient.h"
|
#include "lib/utils/httpclient.h"
|
||||||
#include "ebusd/scan.h"
|
#include "ebusd/scan.h"
|
||||||
@@ -143,9 +144,15 @@ static MessageMap* s_messageMap = nullptr;
|
|||||||
/** the @a ScanHelper instance, or nullptr. */
|
/** the @a ScanHelper instance, or nullptr. */
|
||||||
static ScanHelper* s_scanHelper = nullptr;
|
static ScanHelper* s_scanHelper = nullptr;
|
||||||
|
|
||||||
|
/** the @a Request @a Queue instance, or nullptr. */
|
||||||
|
static Queue<Request*>* s_requestQueue = nullptr;
|
||||||
|
|
||||||
/** the @a MainLoop instance, or nullptr. */
|
/** the @a MainLoop instance, or nullptr. */
|
||||||
static MainLoop* s_mainLoop = nullptr;
|
static MainLoop* s_mainLoop = nullptr;
|
||||||
|
|
||||||
|
/** the @a Network instance, or nullptr. */
|
||||||
|
static Network* s_network = nullptr;
|
||||||
|
|
||||||
/** the (optionally corrected) config path for retrieving configuration files from. */
|
/** the (optionally corrected) config path for retrieving configuration files from. */
|
||||||
static string s_configPath = CONFIG_PATH;
|
static string s_configPath = CONFIG_PATH;
|
||||||
|
|
||||||
@@ -735,15 +742,27 @@ void daemonize() {
|
|||||||
* Clean up all dynamically allocated and stop main loop and all dependent components.
|
* Clean up all dynamically allocated and stop main loop and all dependent components.
|
||||||
*/
|
*/
|
||||||
void cleanup() {
|
void cleanup() {
|
||||||
if (s_mainLoop) {
|
if (s_network != nullptr) {
|
||||||
|
delete s_network;
|
||||||
|
s_network = nullptr;
|
||||||
|
}
|
||||||
|
if (s_mainLoop != nullptr) {
|
||||||
delete s_mainLoop;
|
delete s_mainLoop;
|
||||||
s_mainLoop = nullptr;
|
s_mainLoop = nullptr;
|
||||||
}
|
}
|
||||||
if (s_messageMap) {
|
if (s_requestQueue != nullptr) {
|
||||||
|
Request* msg;
|
||||||
|
while ((msg = s_requestQueue->pop()) != nullptr) {
|
||||||
|
delete msg;
|
||||||
|
}
|
||||||
|
delete s_requestQueue;
|
||||||
|
s_requestQueue = nullptr;
|
||||||
|
}
|
||||||
|
if (s_messageMap != nullptr) {
|
||||||
delete s_messageMap;
|
delete s_messageMap;
|
||||||
s_messageMap = nullptr;
|
s_messageMap = nullptr;
|
||||||
}
|
}
|
||||||
if (s_scanHelper) {
|
if (s_scanHelper != nullptr) {
|
||||||
delete s_scanHelper;
|
delete s_scanHelper;
|
||||||
s_scanHelper = nullptr;
|
s_scanHelper = nullptr;
|
||||||
}
|
}
|
||||||
@@ -1030,8 +1049,10 @@ int main(int argc, char* argv[], char* envp[]) {
|
|||||||
// load configuration files
|
// load configuration files
|
||||||
s_scanHelper->loadConfigFiles(!s_opt.scanConfig);
|
s_scanHelper->loadConfigFiles(!s_opt.scanConfig);
|
||||||
|
|
||||||
|
s_requestQueue = new Queue<Request*>();
|
||||||
|
|
||||||
// create the MainLoop and start it
|
// create the MainLoop and start it
|
||||||
s_mainLoop = new MainLoop(s_opt, device, s_messageMap, s_scanHelper);
|
s_mainLoop = new MainLoop(s_opt, device, s_messageMap, s_scanHelper, s_requestQueue);
|
||||||
if (s_opt.injectMessages) {
|
if (s_opt.injectMessages) {
|
||||||
BusHandler* busHandler = s_mainLoop->getBusHandler();
|
BusHandler* busHandler = s_mainLoop->getBusHandler();
|
||||||
int scanAdrCount = 0;
|
int scanAdrCount = 0;
|
||||||
@@ -1064,6 +1085,9 @@ int main(int argc, char* argv[], char* envp[]) {
|
|||||||
}
|
}
|
||||||
s_mainLoop->start("mainloop");
|
s_mainLoop->start("mainloop");
|
||||||
|
|
||||||
|
s_network = new Network(s_opt.localOnly, s_opt.port, s_opt.httpPort, s_requestQueue);
|
||||||
|
s_network->start("network");
|
||||||
|
|
||||||
// wait for end of MainLoop
|
// wait for end of MainLoop
|
||||||
s_mainLoop->join();
|
s_mainLoop->join();
|
||||||
|
|
||||||
|
|||||||
+46
-91
@@ -105,12 +105,13 @@ result_t UserList::addFromFile(const string& filename, unsigned int lineNo, map<
|
|||||||
#define VERBOSITY_4 (VERBOSITY_3 | OF_ALL_ATTRS)
|
#define VERBOSITY_4 (VERBOSITY_3 | OF_ALL_ATTRS)
|
||||||
|
|
||||||
|
|
||||||
MainLoop::MainLoop(const struct options& opt, Device *device, MessageMap* messages, ScanHelper* scanHelper)
|
MainLoop::MainLoop(const struct options& opt, Device *device, MessageMap* messages, ScanHelper* scanHelper,
|
||||||
|
Queue<Request*>* requestQueue)
|
||||||
: Thread(), m_device(device), m_reconnectCount(0), m_userList(opt.accessLevel), m_messages(messages),
|
: Thread(), m_device(device), m_reconnectCount(0), m_userList(opt.accessLevel), m_messages(messages),
|
||||||
m_scanHelper(scanHelper), m_address(opt.address), m_scanConfig(opt.scanConfig),
|
m_scanHelper(scanHelper), m_address(opt.address), m_scanConfig(opt.scanConfig),
|
||||||
m_initialScan(opt.readOnly ? ESC : opt.initialScan), m_scanStatus(SCAN_STATUS_NONE),
|
m_initialScan(opt.readOnly ? ESC : opt.initialScan), m_scanStatus(SCAN_STATUS_NONE),
|
||||||
m_polling(opt.pollInterval > 0), m_enableHex(opt.enableHex),
|
m_polling(opt.pollInterval > 0), m_enableHex(opt.enableHex),
|
||||||
m_shutdown(false), m_runUpdateCheck(opt.updateCheck), m_httpClient() {
|
m_shutdown(false), m_runUpdateCheck(opt.updateCheck), m_httpClient(), m_requestQueue(requestQueue) {
|
||||||
m_device->setListener(this);
|
m_device->setListener(this);
|
||||||
// open Device
|
// open Device
|
||||||
result_t result = m_device->open();
|
result_t result = m_device->open();
|
||||||
@@ -160,8 +161,6 @@ MainLoop::MainLoop(const struct options& opt, Device *device, MessageMap* messag
|
|||||||
|
|
||||||
// create network
|
// create network
|
||||||
m_htmlPath = opt.htmlPath;
|
m_htmlPath = opt.htmlPath;
|
||||||
m_network = new Network(opt.localOnly, opt.port, opt.httpPort, &m_netQueue);
|
|
||||||
m_network->start("network");
|
|
||||||
logInfo(lf_main, "registering data handlers");
|
logInfo(lf_main, "registering data handlers");
|
||||||
if (datahandler_register(&m_userList, m_busHandler, messages, &m_dataHandlers)) {
|
if (datahandler_register(&m_userList, m_busHandler, messages, &m_dataHandlers)) {
|
||||||
logInfo(lf_main, "registered data handlers");
|
logInfo(lf_main, "registered data handlers");
|
||||||
@@ -192,10 +191,6 @@ MainLoop::~MainLoop() {
|
|||||||
delete m_logRawFile;
|
delete m_logRawFile;
|
||||||
m_logRawFile = nullptr;
|
m_logRawFile = nullptr;
|
||||||
}
|
}
|
||||||
if (m_network != nullptr) {
|
|
||||||
delete m_network;
|
|
||||||
m_network = nullptr;
|
|
||||||
}
|
|
||||||
if (m_busHandler != nullptr) {
|
if (m_busHandler != nullptr) {
|
||||||
delete m_busHandler;
|
delete m_busHandler;
|
||||||
m_busHandler = nullptr;
|
m_busHandler = nullptr;
|
||||||
@@ -204,10 +199,6 @@ MainLoop::~MainLoop() {
|
|||||||
delete m_device;
|
delete m_device;
|
||||||
m_device = nullptr;
|
m_device = nullptr;
|
||||||
}
|
}
|
||||||
NetMessage* msg;
|
|
||||||
while ((msg = m_netQueue.pop()) != nullptr) {
|
|
||||||
delete msg;
|
|
||||||
}
|
|
||||||
if (m_newlyDefinedMessages) {
|
if (m_newlyDefinedMessages) {
|
||||||
delete m_newlyDefinedMessages;
|
delete m_newlyDefinedMessages;
|
||||||
m_newlyDefinedMessages = nullptr;
|
m_newlyDefinedMessages = nullptr;
|
||||||
@@ -245,8 +236,8 @@ void MainLoop::run() {
|
|||||||
dataHandler->startHandler();
|
dataHandler->startHandler();
|
||||||
}
|
}
|
||||||
while (!m_shutdown) {
|
while (!m_shutdown) {
|
||||||
// pick the next message to handle
|
// pick the next request to handle
|
||||||
NetMessage* netMessage = m_netQueue.pop(taskDelay);
|
Request* req = m_requestQueue->pop(taskDelay);
|
||||||
time(&now);
|
time(&now);
|
||||||
if (now < lastTaskRun) {
|
if (now < lastTaskRun) {
|
||||||
// clock skew
|
// clock skew
|
||||||
@@ -410,66 +401,62 @@ void MainLoop::run() {
|
|||||||
m_messages->unlock();
|
m_messages->unlock();
|
||||||
sinkSince = now;
|
sinkSince = now;
|
||||||
}
|
}
|
||||||
if (netMessage == nullptr) {
|
if (req == nullptr) {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
if (m_shutdown) {
|
if (m_shutdown) {
|
||||||
netMessage->setResult("ERR: shutdown", "", nullptr, now, true);
|
req->setResult("ERR: shutdown", "", nullptr, now, true);
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
string request = netMessage->getRequest();
|
string user = req->getUser();
|
||||||
string user = netMessage->getUser();
|
RequestMode reqMode = req->getMode(&since);
|
||||||
ClientSettings settings = netMessage->getSettings(&since);
|
if (reqMode.listenMode == lm_none) {
|
||||||
if (!netMessage->isListeningMode()) {
|
|
||||||
since = now;
|
since = now;
|
||||||
}
|
}
|
||||||
ostringstream ostream;
|
ostringstream ostream;
|
||||||
bool connected = true;
|
bool connected = true;
|
||||||
if (request.length() > 0) {
|
if (!req->empty()) {
|
||||||
logDebug(lf_main, ">>> %s", request.c_str());
|
req->log();
|
||||||
result_t result = decodeMessage(request, netMessage->isHttp(), &connected, &settings, &user, &reload, &ostream);
|
result_t result = decodeRequest(req, &connected, &reqMode, &user, &reload, &ostream);
|
||||||
if (!netMessage->isHttp() && (ostream.tellp() == 0 || result != RESULT_OK)) {
|
if (!req->isHttp() && (ostream.tellp() == 0 || result != RESULT_OK)) {
|
||||||
if (settings.mode != cm_direct) {
|
if (reqMode.listenMode != lm_direct) {
|
||||||
ostream.str("");
|
ostream.str("");
|
||||||
}
|
}
|
||||||
ostream << getResultCode(result);
|
ostream << getResultCode(result);
|
||||||
}
|
}
|
||||||
if (ostream.tellp() > 100) {
|
const auto resp = ostream.str();
|
||||||
logDebug(lf_main, "<<< %s ...", ostream.str().substr(0, 100).c_str());
|
req->log(&resp);
|
||||||
} else {
|
|
||||||
logDebug(lf_main, "<<< %s", ostream.str().c_str());
|
|
||||||
}
|
|
||||||
if (ostream.tellp() == 0) {
|
if (ostream.tellp() == 0) {
|
||||||
ostream << "\n"; // only for HTTP
|
ostream << "\n"; // only for HTTP
|
||||||
} else if (!netMessage->isHttp()) {
|
} else if (!req->isHttp()) {
|
||||||
ostream << (settings.mode == cm_direct ? "\n" : "\n\n");
|
ostream << (reqMode.listenMode == lm_direct ? "\n" : "\n\n");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if (settings.mode == cm_listen) {
|
if (reqMode.listenMode == lm_listen) {
|
||||||
if (!settings.listenOnlyUnknown) {
|
if (!reqMode.listenOnlyUnknown) {
|
||||||
string levels = getUserLevels(user);
|
string levels = getUserLevels(user);
|
||||||
messages.clear();
|
messages.clear();
|
||||||
m_messages->findAll("", "", levels, false, true, true, true, true, true, since, now, true, &messages);
|
m_messages->findAll("", "", levels, false, true, true, true, true, true, since, now, true, &messages);
|
||||||
for (const auto message : messages) {
|
for (const auto message : messages) {
|
||||||
ostream << message->getCircuit() << " " << message->getName() << " = " << dec;
|
ostream << message->getCircuit() << " " << message->getName() << " = " << dec;
|
||||||
message->decodeLastData(false, nullptr, -1, settings.format, &ostream);
|
message->decodeLastData(false, nullptr, -1, reqMode.format, &ostream);
|
||||||
ostream << endl;
|
ostream << endl;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if (settings.listenWithUnknown || settings.listenOnlyUnknown) {
|
if (reqMode.listenWithUnknown || reqMode.listenOnlyUnknown) {
|
||||||
if (m_busHandler->isGrabEnabled()) {
|
if (m_busHandler->isGrabEnabled()) {
|
||||||
m_busHandler->formatGrabResult(true, OF_NONE, &ostream, true, since, now);
|
m_busHandler->formatGrabResult(true, OF_NONE, &ostream, true, since, now);
|
||||||
} else {
|
} else {
|
||||||
m_busHandler->enableGrab(true); // needed for listening to all messages
|
m_busHandler->enableGrab(true); // needed for listening to all messages
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
} else if (settings.mode == cm_direct) {
|
} else if (reqMode.listenMode == lm_direct) {
|
||||||
if (m_busHandler->isGrabEnabled()) {
|
if (m_busHandler->isGrabEnabled()) {
|
||||||
m_busHandler->formatGrabResult(false, OF_NONE, &ostream, true, since, now);
|
m_busHandler->formatGrabResult(false, OF_NONE, &ostream, true, since, now);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
// send result to client
|
// send result to client
|
||||||
netMessage->setResult(ostream.str(), user, &settings, now, !connected);
|
req->setResult(ostream.str(), user, &reqMode, now, !connected);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -529,49 +516,18 @@ void MainLoop::notifyStatus(bool error, const char* message) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
result_t MainLoop::decodeMessage(const string &data, bool isHttp, bool* connected, ClientSettings* settings,
|
result_t MainLoop::decodeRequest(Request* req, bool* connected, RequestMode* reqMode,
|
||||||
string* user, bool* reload, ostringstream* ostream) {
|
string* user, bool* reload, ostringstream* ostream) {
|
||||||
string token, previous;
|
|
||||||
istringstream stream(data);
|
|
||||||
vector<string> args;
|
vector<string> args;
|
||||||
char escaped = 0;
|
req->split(&args);
|
||||||
|
string cmd = args.size() > 0 ? args[0] : "";
|
||||||
char delim = ' ';
|
if (req->isHttp()) {
|
||||||
while (getline(stream, token, delim)) {
|
|
||||||
if (!isHttp) {
|
|
||||||
if (escaped) {
|
|
||||||
args.pop_back();
|
|
||||||
if (token.length() > 0 && token[token.length()-1] == escaped) {
|
|
||||||
token.erase(token.length() - 1, 1);
|
|
||||||
escaped = 0;
|
|
||||||
}
|
|
||||||
token = previous + " " + token;
|
|
||||||
} else if (token.length() == 0) { // allow multiple space chars for a single delimiter
|
|
||||||
continue;
|
|
||||||
} else if (token[0] == '"' || token[0] == '\'') {
|
|
||||||
escaped = token[0];
|
|
||||||
token.erase(0, 1);
|
|
||||||
if (token.length() > 0 && token[token.length()-1] == escaped) {
|
|
||||||
token.erase(token.length() - 1, 1);
|
|
||||||
escaped = 0;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
args.push_back(token);
|
|
||||||
previous = token;
|
|
||||||
if (isHttp) {
|
|
||||||
delim = (args.size() == 1) ? '?' : '\n';
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
if (isHttp) {
|
|
||||||
if (args.size() < 2) {
|
if (args.size() < 2) {
|
||||||
*connected = false;
|
*connected = false;
|
||||||
*ostream << "HTTP/1.0 400 Bad Request\r\n\r\n";
|
*ostream << "HTTP/1.0 400 Bad Request\r\n\r\n";
|
||||||
return RESULT_OK;
|
return RESULT_OK;
|
||||||
}
|
}
|
||||||
const char* str = args.size() > 0 ? args[0].c_str() : "";
|
if (cmd == "GET") {
|
||||||
if (strcmp(str, "GET") == 0) {
|
|
||||||
return executeGet(args, connected, ostream);
|
return executeGet(args, connected, ostream);
|
||||||
}
|
}
|
||||||
*connected = false;
|
*connected = false;
|
||||||
@@ -579,7 +535,6 @@ result_t MainLoop::decodeMessage(const string &data, bool isHttp, bool* connecte
|
|||||||
return RESULT_OK;
|
return RESULT_OK;
|
||||||
}
|
}
|
||||||
|
|
||||||
string cmd = args.size() > 0 ? args[0] : "";
|
|
||||||
transform(cmd.begin(), cmd.end(), cmd.begin(), ::toupper);
|
transform(cmd.begin(), cmd.end(), cmd.begin(), ::toupper);
|
||||||
if (cmd == "?" || cmd == "H" || cmd == "HELP") {
|
if (cmd == "?" || cmd == "H" || cmd == "HELP") {
|
||||||
// found "HELP CMD"
|
// found "HELP CMD"
|
||||||
@@ -587,8 +542,8 @@ result_t MainLoop::decodeMessage(const string &data, bool isHttp, bool* connecte
|
|||||||
transform(cmd.begin(), cmd.end(), cmd.begin(), ::toupper);
|
transform(cmd.begin(), cmd.end(), cmd.begin(), ::toupper);
|
||||||
args.clear(); // empty args is used as command help indicator
|
args.clear(); // empty args is used as command help indicator
|
||||||
}
|
}
|
||||||
if (settings->mode == cm_direct) {
|
if (reqMode->listenMode == lm_direct) {
|
||||||
return executeDirect(args, &settings->mode, ostream);
|
return executeDirect(args, reqMode, ostream);
|
||||||
}
|
}
|
||||||
if (cmd.empty() && args.size() == 0) {
|
if (cmd.empty() && args.size() == 0) {
|
||||||
return executeHelp(ostream);
|
return executeHelp(ostream);
|
||||||
@@ -620,10 +575,10 @@ result_t MainLoop::decodeMessage(const string &data, bool isHttp, bool* connecte
|
|||||||
return executeFind(args, getUserLevels(*user), ostream);
|
return executeFind(args, getUserLevels(*user), ostream);
|
||||||
}
|
}
|
||||||
if (cmd == "L" || cmd == "LISTEN") {
|
if (cmd == "L" || cmd == "LISTEN") {
|
||||||
return executeListen(args, settings, ostream);
|
return executeListen(args, reqMode, ostream);
|
||||||
}
|
}
|
||||||
if (cmd == "DIRECT") {
|
if (cmd == "DIRECT") {
|
||||||
return executeDirect(args, &settings->mode, ostream);
|
return executeDirect(args, reqMode, ostream);
|
||||||
}
|
}
|
||||||
if (cmd == "S" || cmd == "STATE") {
|
if (cmd == "S" || cmd == "STATE") {
|
||||||
return executeState(args, ostream);
|
return executeState(args, ostream);
|
||||||
@@ -1293,10 +1248,10 @@ result_t MainLoop::executeHex(const vector<string>& args, ostringstream* ostream
|
|||||||
return RESULT_OK;
|
return RESULT_OK;
|
||||||
}
|
}
|
||||||
|
|
||||||
result_t MainLoop::executeDirect(const vector<string>& args, ClientMode* mode, ostringstream* ostream) {
|
result_t MainLoop::executeDirect(const vector<string>& args, RequestMode* reqMode, ostringstream* ostream) {
|
||||||
if (*mode != cm_direct) {
|
if (reqMode->listenMode != lm_direct) {
|
||||||
if (args.size() == 1) {
|
if (args.size() == 1) {
|
||||||
*mode = cm_direct;
|
reqMode->listenMode = lm_direct;
|
||||||
m_busHandler->enableGrab(true); // needed for listening to all messages
|
m_busHandler->enableGrab(true); // needed for listening to all messages
|
||||||
*ostream << "direct mode started";
|
*ostream << "direct mode started";
|
||||||
return RESULT_OK;
|
return RESULT_OK;
|
||||||
@@ -1308,7 +1263,7 @@ result_t MainLoop::executeDirect(const vector<string>& args, ClientMode* mode, o
|
|||||||
if (args.size() > 0) {
|
if (args.size() > 0) {
|
||||||
string firstArg = args[0];
|
string firstArg = args[0];
|
||||||
if (firstArg == "stop") {
|
if (firstArg == "stop") {
|
||||||
*mode = cm_normal;
|
reqMode->listenMode = lm_none;
|
||||||
*ostream << "direct mode stopped";
|
*ostream << "direct mode stopped";
|
||||||
return RESULT_OK;
|
return RESULT_OK;
|
||||||
}
|
}
|
||||||
@@ -1321,7 +1276,7 @@ result_t MainLoop::executeDirect(const vector<string>& args, ClientMode* mode, o
|
|||||||
}
|
}
|
||||||
*ostream << ":";
|
*ostream << ":";
|
||||||
if (!m_enableHex) {
|
if (!m_enableHex) {
|
||||||
*ostream << "ERR: command not enabled";
|
*ostream << "ERR: hex command not enabled";
|
||||||
return RESULT_OK;
|
return RESULT_OK;
|
||||||
}
|
}
|
||||||
size_t argPos = 0;
|
size_t argPos = 0;
|
||||||
@@ -1560,7 +1515,7 @@ result_t MainLoop::executeFind(const vector<string>& args, const string& levels,
|
|||||||
return RESULT_OK;
|
return RESULT_OK;
|
||||||
}
|
}
|
||||||
|
|
||||||
result_t MainLoop::executeListen(const vector<string>& args, ClientSettings* settings, ostringstream* ostream) {
|
result_t MainLoop::executeListen(const vector<string>& args, RequestMode* reqMode, ostringstream* ostream) {
|
||||||
size_t argPos = 1;
|
size_t argPos = 1;
|
||||||
OutputFormat verbosity = OF_NONE;
|
OutputFormat verbosity = OF_NONE;
|
||||||
bool listenWithUnknown = false;
|
bool listenWithUnknown = false;
|
||||||
@@ -1594,17 +1549,17 @@ result_t MainLoop::executeListen(const vector<string>& args, ClientSettings* set
|
|||||||
argPos++;
|
argPos++;
|
||||||
}
|
}
|
||||||
if (argPos > 0 && args.size() == argPos) {
|
if (argPos > 0 && args.size() == argPos) {
|
||||||
settings->format = verbosity;
|
reqMode->format = verbosity;
|
||||||
settings->listenWithUnknown = listenWithUnknown;
|
reqMode->listenWithUnknown = listenWithUnknown;
|
||||||
settings->listenOnlyUnknown = listenOnlyUnknown;
|
reqMode->listenOnlyUnknown = listenOnlyUnknown;
|
||||||
if (listenWithUnknown || listenOnlyUnknown) {
|
if (listenWithUnknown || listenOnlyUnknown) {
|
||||||
m_busHandler->enableGrab(true); // needed for listening to all messages
|
m_busHandler->enableGrab(true); // needed for listening to all messages
|
||||||
}
|
}
|
||||||
if (settings->mode == cm_listen) {
|
if (reqMode->listenMode == lm_listen) {
|
||||||
*ostream << "listen continued";
|
*ostream << "listen continued";
|
||||||
return RESULT_OK;
|
return RESULT_OK;
|
||||||
}
|
}
|
||||||
settings->mode = cm_listen;
|
reqMode->listenMode = lm_listen;
|
||||||
*ostream << "listen started";
|
*ostream << "listen started";
|
||||||
return RESULT_OK;
|
return RESULT_OK;
|
||||||
}
|
}
|
||||||
@@ -1620,7 +1575,7 @@ result_t MainLoop::executeListen(const vector<string>& args, ClientSettings* set
|
|||||||
" -U only show unknown messages";
|
" -U only show unknown messages";
|
||||||
return RESULT_OK;
|
return RESULT_OK;
|
||||||
}
|
}
|
||||||
settings->mode = cm_normal;
|
reqMode->listenMode = lm_none;
|
||||||
*ostream << "listen stopped";
|
*ostream << "listen stopped";
|
||||||
return RESULT_OK;
|
return RESULT_OK;
|
||||||
}
|
}
|
||||||
|
|||||||
+15
-23
@@ -26,7 +26,7 @@
|
|||||||
#include <algorithm>
|
#include <algorithm>
|
||||||
#include "ebusd/bushandler.h"
|
#include "ebusd/bushandler.h"
|
||||||
#include "ebusd/datahandler.h"
|
#include "ebusd/datahandler.h"
|
||||||
#include "ebusd/network.h"
|
#include "ebusd/request.h"
|
||||||
#include "ebusd/scan.h"
|
#include "ebusd/scan.h"
|
||||||
#include "lib/ebus/filereader.h"
|
#include "lib/ebus/filereader.h"
|
||||||
#include "lib/ebus/message.h"
|
#include "lib/ebus/message.h"
|
||||||
@@ -102,13 +102,15 @@ class UserList : public UserInfo, public MappedFileReader {
|
|||||||
class MainLoop : public Thread, DeviceListener {
|
class MainLoop : public Thread, DeviceListener {
|
||||||
public:
|
public:
|
||||||
/**
|
/**
|
||||||
* Construct the main loop and create network and bus handling components.
|
* Construct the main loop and create bus handling components.
|
||||||
* @param opt the program options.
|
* @param opt the program options.
|
||||||
* @param device the @a Device instance.
|
* @param device the @a Device instance.
|
||||||
* @param messages the @a MessageMap instance.
|
* @param messages the @a MessageMap instance.
|
||||||
* @param scanHelper the @a ScanHelper instance.
|
* @param scanHelper the @a ScanHelper instance.
|
||||||
|
* @param requestQueue the reference to the @a Request @a Queue.
|
||||||
*/
|
*/
|
||||||
MainLoop(const struct options& opt, Device *device, MessageMap* messages, ScanHelper* scanHelper);
|
MainLoop(const struct options& opt, Device *device, MessageMap* messages, ScanHelper* scanHelper,
|
||||||
|
Queue<Request*>* requestQueue);
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Destructor.
|
* Destructor.
|
||||||
@@ -126,12 +128,6 @@ class MainLoop : public Thread, DeviceListener {
|
|||||||
*/
|
*/
|
||||||
BusHandler* getBusHandler() { return m_busHandler; }
|
BusHandler* getBusHandler() { return m_busHandler; }
|
||||||
|
|
||||||
/**
|
|
||||||
* Add a client @a NetMessage to the queue.
|
|
||||||
* @param message the client @a NetMessage to handle.
|
|
||||||
*/
|
|
||||||
void addMessage(NetMessage* message) { m_netQueue.push(message); }
|
|
||||||
|
|
||||||
// @copydoc
|
// @copydoc
|
||||||
void notifyDeviceData(symbol_t symbol, bool received) override;
|
void notifyDeviceData(symbol_t symbol, bool received) override;
|
||||||
|
|
||||||
@@ -146,17 +142,16 @@ class MainLoop : public Thread, DeviceListener {
|
|||||||
|
|
||||||
private:
|
private:
|
||||||
/**
|
/**
|
||||||
* Decode and execute client message.
|
* Decode and execute client request.
|
||||||
* @param data the data string to decode (may be empty).
|
* @param req the @a Request to decode.
|
||||||
* @param connected set to false when the client connection shall be closed.
|
* @param connected set to false when the client connection shall be closed.
|
||||||
* @param isHttp true for HTTP message.
|
* @param reqMode the @a RequestMode to use and update.
|
||||||
* @param settings set to the new client settings.
|
|
||||||
* @param user set to the new user name when changed by authentication.
|
* @param user set to the new user name when changed by authentication.
|
||||||
* @param reload set to true when the configuration files were reloaded.
|
* @param reload set to true when the configuration files were reloaded.
|
||||||
* @param ostream the @a ostringstream to format the result string to.
|
* @param ostream the @a ostringstream to format the result string to.
|
||||||
* @return the result code.
|
* @return the result code.
|
||||||
*/
|
*/
|
||||||
result_t decodeMessage(const string& data, bool isHttp, bool* connected, ClientSettings* settings,
|
result_t decodeRequest(Request* req, bool* connected, RequestMode* reqMode,
|
||||||
string* user, bool* reload, ostringstream* ostream);
|
string* user, bool* reload, ostringstream* ostream);
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -227,11 +222,11 @@ class MainLoop : public Thread, DeviceListener {
|
|||||||
/**
|
/**
|
||||||
* Execute the direct command.
|
* Execute the direct command.
|
||||||
* @param args the arguments passed to the command (starting with the command itself), or empty for help.
|
* @param args the arguments passed to the command (starting with the command itself), or empty for help.
|
||||||
* @param mode set to the new client mode.
|
* @param reqMode the @a RequestMode to use and update.
|
||||||
* @param ostream the @a ostringstream to format the result string to.
|
* @param ostream the @a ostringstream to format the result string to.
|
||||||
* @return the result code.
|
* @return the result code.
|
||||||
*/
|
*/
|
||||||
result_t executeDirect(const vector<string>& args, ClientMode* mode, ostringstream* ostream);
|
result_t executeDirect(const vector<string>& args, RequestMode* reqMode, ostringstream* ostream);
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Execute the find command.
|
* Execute the find command.
|
||||||
@@ -245,11 +240,11 @@ class MainLoop : public Thread, DeviceListener {
|
|||||||
/**
|
/**
|
||||||
* Execute the listen command.
|
* Execute the listen command.
|
||||||
* @param args the arguments passed to the command (starting with the command itself), or empty for help.
|
* @param args the arguments passed to the command (starting with the command itself), or empty for help.
|
||||||
* @param settings set to the new client settings.
|
* @param reqMode the @a RequestMode to use and update.
|
||||||
* @param ostream the @a ostringstream to format the result string to.
|
* @param ostream the @a ostringstream to format the result string to.
|
||||||
* @return the result code.
|
* @return the result code.
|
||||||
*/
|
*/
|
||||||
result_t executeListen(const vector<string>& args, ClientSettings* settings, ostringstream* ostream);
|
result_t executeListen(const vector<string>& args, RequestMode* reqMode, ostringstream* ostream);
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Execute the state command.
|
* Execute the state command.
|
||||||
@@ -445,11 +440,8 @@ class MainLoop : public Thread, DeviceListener {
|
|||||||
/** the created @a BusHandler instance. */
|
/** the created @a BusHandler instance. */
|
||||||
BusHandler* m_busHandler;
|
BusHandler* m_busHandler;
|
||||||
|
|
||||||
/** the created @a Network instance. */
|
/** the reference to the @a Request @a Queue. */
|
||||||
Network* m_network;
|
Queue<Request*>* m_requestQueue;
|
||||||
|
|
||||||
/** the @a NetMessage @a Queue. */
|
|
||||||
Queue<NetMessage*> m_netQueue;
|
|
||||||
|
|
||||||
/** the path for HTML files served by the HTTP port. */
|
/** the path for HTML files served by the HTTP port. */
|
||||||
string m_htmlPath;
|
string m_htmlPath;
|
||||||
|
|||||||
+15
-45
@@ -37,39 +37,6 @@ int Connection::m_ids = 0;
|
|||||||
#define POLLRDHUP 0
|
#define POLLRDHUP 0
|
||||||
#endif
|
#endif
|
||||||
|
|
||||||
bool NetMessage::add(const char* request) {
|
|
||||||
if (request && request[0]) {
|
|
||||||
string add = request;
|
|
||||||
add.erase(remove(add.begin(), add.end(), '\r'), add.end());
|
|
||||||
m_request.append(add);
|
|
||||||
}
|
|
||||||
size_t pos = m_request.find(m_isHttp ? "\n\n" : "\n");
|
|
||||||
if (pos != string::npos) {
|
|
||||||
if (m_isHttp) {
|
|
||||||
pos = m_request.find("\n");
|
|
||||||
m_request.resize(pos); // reduce to first line
|
|
||||||
// typical first line: GET /ehp/outsidetemp HTTP/1.1
|
|
||||||
pos = m_request.rfind(" HTTP/");
|
|
||||||
if (pos != string::npos) {
|
|
||||||
m_request.resize(pos); // remove "HTTP/x.x" suffix
|
|
||||||
}
|
|
||||||
pos = 0;
|
|
||||||
while ((pos=m_request.find('%', pos)) != string::npos && pos+2 <= m_request.length()) {
|
|
||||||
unsigned int value1, value2;
|
|
||||||
if (sscanf("%1x%1x", m_request.c_str()+pos+1, &value1, &value2) < 2) {
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
m_request[pos] = static_cast<char>(((value1&0x0f) << 4) | (value2&0x0f));
|
|
||||||
m_request.erase(pos+1, 2);
|
|
||||||
}
|
|
||||||
} else if (pos+1 == m_request.length()) {
|
|
||||||
m_request.resize(pos); // reduce to complete lines
|
|
||||||
}
|
|
||||||
return true;
|
|
||||||
}
|
|
||||||
return m_request.length() == 0 && isListeningMode();
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
void Connection::run() {
|
void Connection::run() {
|
||||||
int ret;
|
int ret;
|
||||||
@@ -108,7 +75,7 @@ void Connection::run() {
|
|||||||
#endif
|
#endif
|
||||||
|
|
||||||
bool closed = false;
|
bool closed = false;
|
||||||
NetMessage message(m_isHttp);
|
RequestImpl req(m_isHttp);
|
||||||
|
|
||||||
while (!closed) {
|
while (!closed) {
|
||||||
#ifdef HAVE_PPOLL
|
#ifdef HAVE_PPOLL
|
||||||
@@ -146,7 +113,7 @@ void Connection::run() {
|
|||||||
#endif
|
#endif
|
||||||
}
|
}
|
||||||
|
|
||||||
if (newData || message.isListeningMode()) {
|
if (newData || req.getMode().listenMode != lm_none) {
|
||||||
char data[256];
|
char data[256];
|
||||||
|
|
||||||
if (!m_socket->isValid()) {
|
if (!m_socket->isValid()) {
|
||||||
@@ -165,21 +132,24 @@ void Connection::run() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// decode client data
|
// decode client data
|
||||||
if (message.add(data)) {
|
if (req.add(data)) {
|
||||||
m_netQueue->push(&message);
|
m_requestQueue->push(&req);
|
||||||
|
|
||||||
// wait for result
|
// wait for result
|
||||||
logDebug(lf_network, "[%05d] wait for result", getID());
|
logDebug(lf_network, "[%05d] wait for result", getID());
|
||||||
string result;
|
string result;
|
||||||
message.getResult(&result);
|
bool disconnect = req.waitResponse(&result);
|
||||||
|
|
||||||
if (!m_socket->isValid()) {
|
if (!m_socket->isValid()) {
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
m_socket->send(result.c_str(), result.size());
|
m_socket->send(result.c_str(), result.size());
|
||||||
|
if (disconnect) {
|
||||||
|
break;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if (message.isDisconnect() || !m_socket->isValid()) {
|
if (!m_socket->isValid()) {
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -193,8 +163,8 @@ void Connection::run() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
Network::Network(const bool local, const uint16_t port, const uint16_t httpPort, Queue<NetMessage*>* netQueue)
|
Network::Network(const bool local, const uint16_t port, const uint16_t httpPort, Queue<Request*>* requestQueue)
|
||||||
: Thread(), m_netQueue(netQueue), m_listening(false) {
|
: Thread(), m_requestQueue(requestQueue), m_listening(false) {
|
||||||
m_tcpServer = new TCPServer(port, local ? "127.0.0.1" : "0.0.0.0");
|
m_tcpServer = new TCPServer(port, local ? "127.0.0.1" : "0.0.0.0");
|
||||||
|
|
||||||
if (m_tcpServer != nullptr && m_tcpServer->start() == 0) {
|
if (m_tcpServer != nullptr && m_tcpServer->start() == 0) {
|
||||||
@@ -214,9 +184,9 @@ Network::Network(const bool local, const uint16_t port, const uint16_t httpPort,
|
|||||||
|
|
||||||
Network::~Network() {
|
Network::~Network() {
|
||||||
stop();
|
stop();
|
||||||
NetMessage* netMsg;
|
Request* req;
|
||||||
while ((netMsg = m_netQueue->pop()) != nullptr) {
|
while ((req = m_requestQueue->pop()) != nullptr) {
|
||||||
netMsg->setResult("ERR: shutdown", "", nullptr, 0, true);
|
req->setResult("ERR: shutdown", "", nullptr, 0, true);
|
||||||
}
|
}
|
||||||
while (!m_connections.empty()) {
|
while (!m_connections.empty()) {
|
||||||
Connection* connection = m_connections.back();
|
Connection* connection = m_connections.back();
|
||||||
@@ -333,7 +303,7 @@ void Network::run() {
|
|||||||
if (socket == nullptr) {
|
if (socket == nullptr) {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
Connection* connection = new Connection(socket, isHttp, m_netQueue);
|
Connection* connection = new Connection(socket, isHttp, m_requestQueue);
|
||||||
string ip = socket->getIP();
|
string ip = socket->getIP();
|
||||||
connection->start("connection");
|
connection->start("connection");
|
||||||
m_connections.push_back(connection);
|
m_connections.push_back(connection);
|
||||||
|
|||||||
+12
-190
@@ -23,6 +23,7 @@
|
|||||||
#include <cstdio>
|
#include <cstdio>
|
||||||
#include <algorithm>
|
#include <algorithm>
|
||||||
#include <list>
|
#include <list>
|
||||||
|
#include "ebusd/request.h"
|
||||||
#include "lib/ebus/datatype.h"
|
#include "lib/ebus/datatype.h"
|
||||||
#include "lib/utils/tcpsocket.h"
|
#include "lib/utils/tcpsocket.h"
|
||||||
#include "lib/utils/queue.h"
|
#include "lib/utils/queue.h"
|
||||||
@@ -35,188 +36,9 @@ namespace ebusd {
|
|||||||
* The TCP and HTTP client request handling.
|
* The TCP and HTTP client request handling.
|
||||||
*/
|
*/
|
||||||
|
|
||||||
/** Forward declaration for @a Connection. */
|
|
||||||
class Connection;
|
|
||||||
|
|
||||||
/** the possible client modes. */
|
|
||||||
enum ClientMode {
|
|
||||||
cm_normal, //!< normal mode
|
|
||||||
cm_listen, //!< listening mode
|
|
||||||
cm_direct, //!< direct mode
|
|
||||||
};
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Combination of client settings.
|
* Instance of a connected client, either TCP or HTTP.
|
||||||
*/
|
|
||||||
struct ClientSettings {
|
|
||||||
ClientMode mode; //!< the current client mode
|
|
||||||
OutputFormat format; //!< the output format settings for listen mode
|
|
||||||
bool listenWithUnknown; //!< include unknown messages in listen mode
|
|
||||||
bool listenOnlyUnknown; //!< only print unknown messages in listen mode
|
|
||||||
};
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Class for data/message transfer between @a Connection and @a MainLoop.
|
|
||||||
*/
|
|
||||||
class NetMessage {
|
|
||||||
public:
|
|
||||||
/**
|
|
||||||
* Constructor.
|
|
||||||
* @param isHttp whether this is a HTTP message.
|
|
||||||
*/
|
|
||||||
explicit NetMessage(bool isHttp)
|
|
||||||
: m_isHttp(isHttp), m_resultSet(false), m_disconnect(false), m_listenSince(0) {
|
|
||||||
m_settings.mode = cm_normal;
|
|
||||||
m_settings.format = OF_NONE;
|
|
||||||
m_settings.listenWithUnknown = false;
|
|
||||||
m_settings.listenOnlyUnknown = false;
|
|
||||||
pthread_mutex_init(&m_mutex, nullptr);
|
|
||||||
pthread_cond_init(&m_cond, nullptr);
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Destructor.
|
|
||||||
*/
|
|
||||||
~NetMessage() {
|
|
||||||
m_resultSet = true;
|
|
||||||
pthread_mutex_destroy(&m_mutex);
|
|
||||||
pthread_cond_destroy(&m_cond);
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
private:
|
|
||||||
/**
|
|
||||||
* Hidden copy constructor.
|
|
||||||
* @param src the object to copy from.
|
|
||||||
*/
|
|
||||||
NetMessage(const NetMessage& src);
|
|
||||||
|
|
||||||
|
|
||||||
public:
|
|
||||||
/**
|
|
||||||
* Add request data received from the client.
|
|
||||||
* @param request the request data from the client.
|
|
||||||
* @return true when the request is complete and the response shall be prepared.
|
|
||||||
*/
|
|
||||||
bool add(const char* request);
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Return whether this is a HTTP message.
|
|
||||||
* @return whether this is a HTTP message.
|
|
||||||
*/
|
|
||||||
bool isHttp() const { return m_isHttp; }
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Return the request string.
|
|
||||||
* @return the request string.
|
|
||||||
*/
|
|
||||||
const string& getRequest() const { return m_request; }
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Return the current user name.
|
|
||||||
* @return the current user name.
|
|
||||||
*/
|
|
||||||
const string& getUser() const { return m_user; }
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Wait for the result being set and return the result string.
|
|
||||||
* @param result the variable in which to store the result string.
|
|
||||||
*/
|
|
||||||
void getResult(string* result) {
|
|
||||||
pthread_mutex_lock(&m_mutex);
|
|
||||||
|
|
||||||
if (!m_resultSet) {
|
|
||||||
pthread_cond_wait(&m_cond, &m_mutex);
|
|
||||||
}
|
|
||||||
m_request.clear();
|
|
||||||
*result = m_result;
|
|
||||||
m_result.clear();
|
|
||||||
m_resultSet = false;
|
|
||||||
pthread_mutex_unlock(&m_mutex);
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Set the result string and notify the waiting thread.
|
|
||||||
* @param result the result string.
|
|
||||||
* @param user the new user name.
|
|
||||||
* @param settings the new client settings.
|
|
||||||
* @param listenUntil the end time to which to updates were added (exclusive).
|
|
||||||
* @param disconnect true when the client shall be disconnected.
|
|
||||||
*/
|
|
||||||
void setResult(const string& result, const string& user, ClientSettings* settings, time_t listenUntil,
|
|
||||||
bool disconnect) {
|
|
||||||
pthread_mutex_lock(&m_mutex);
|
|
||||||
m_result = result;
|
|
||||||
m_user = user;
|
|
||||||
m_disconnect = disconnect;
|
|
||||||
if (settings) {
|
|
||||||
m_settings = *settings;
|
|
||||||
}
|
|
||||||
m_listenSince = listenUntil;
|
|
||||||
m_resultSet = true;
|
|
||||||
pthread_cond_signal(&m_cond);
|
|
||||||
pthread_mutex_unlock(&m_mutex);
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Return the client settings.
|
|
||||||
* @param listenSince set listening to the specified start time from which to add updates (inclusive).
|
|
||||||
* @return the client settings.
|
|
||||||
*/
|
|
||||||
ClientSettings getSettings(time_t* listenSince = nullptr) {
|
|
||||||
if (listenSince) {
|
|
||||||
*listenSince = m_listenSince;
|
|
||||||
}
|
|
||||||
return m_settings;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Return whether this instance is in one of the listening modes.
|
|
||||||
* @return whether this instance is in one of the listening modes.
|
|
||||||
*/
|
|
||||||
bool isListeningMode() { return m_settings.mode == cm_listen || m_settings.mode == cm_direct; }
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Return whether the client shall be disconnected.
|
|
||||||
* @return true when the client shall be disconnected.
|
|
||||||
*/
|
|
||||||
bool isDisconnect() { return m_disconnect; }
|
|
||||||
|
|
||||||
|
|
||||||
private:
|
|
||||||
/** whether this is a HTTP message. */
|
|
||||||
const bool m_isHttp;
|
|
||||||
|
|
||||||
/** the request string. */
|
|
||||||
string m_request;
|
|
||||||
|
|
||||||
/** the current user name. */
|
|
||||||
string m_user;
|
|
||||||
|
|
||||||
/** whether the result was already set. */
|
|
||||||
bool m_resultSet;
|
|
||||||
|
|
||||||
/** the result string. */
|
|
||||||
string m_result;
|
|
||||||
|
|
||||||
/** set to true when the client shall be disconnected. */
|
|
||||||
bool m_disconnect;
|
|
||||||
|
|
||||||
/** mutex variable for exclusive lock. */
|
|
||||||
pthread_mutex_t m_mutex;
|
|
||||||
|
|
||||||
/** condition variable for exclusive lock. */
|
|
||||||
pthread_cond_t m_cond;
|
|
||||||
|
|
||||||
/** the client settings. */
|
|
||||||
ClientSettings m_settings;
|
|
||||||
|
|
||||||
/** start timestamp of listening update. */
|
|
||||||
time_t m_listenSince;
|
|
||||||
};
|
|
||||||
|
|
||||||
/**
|
|
||||||
* class connection which handle client and baseloop communication.
|
|
||||||
*/
|
*/
|
||||||
class Connection : public Thread {
|
class Connection : public Thread {
|
||||||
public:
|
public:
|
||||||
@@ -224,10 +46,10 @@ class Connection : public Thread {
|
|||||||
* Constructor.
|
* Constructor.
|
||||||
* @param socket the @a TCPSocket for communication.
|
* @param socket the @a TCPSocket for communication.
|
||||||
* @param isHttp whether this is a HTTP message.
|
* @param isHttp whether this is a HTTP message.
|
||||||
* @param netQueue the reference to the @a NetMessage @a Queue.
|
* @param requestQueue the reference to the @a Request @a Queue.
|
||||||
*/
|
*/
|
||||||
Connection(TCPSocket* socket, const bool isHttp, Queue<NetMessage*>* netQueue)
|
Connection(TCPSocket* socket, const bool isHttp, Queue<Request*>* requestQueue)
|
||||||
: Thread(), m_isHttp(isHttp), m_socket(socket), m_netQueue(netQueue), m_endedAt(0) {
|
: Thread(), m_isHttp(isHttp), m_socket(socket), m_requestQueue(requestQueue), m_endedAt(0) {
|
||||||
m_id = ++m_ids;
|
m_id = ++m_ids;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -267,8 +89,8 @@ class Connection : public Thread {
|
|||||||
/** the @a TCPSocket for communication. */
|
/** the @a TCPSocket for communication. */
|
||||||
TCPSocket* m_socket;
|
TCPSocket* m_socket;
|
||||||
|
|
||||||
/** the reference to the @a NetMessage @a Queue. */
|
/** the reference to the @a Request @a Queue. */
|
||||||
Queue<NetMessage*>* m_netQueue;
|
Queue<Request*>* m_requestQueue;
|
||||||
|
|
||||||
/** notification object for shutdown procedure. */
|
/** notification object for shutdown procedure. */
|
||||||
Notify m_notify;
|
Notify m_notify;
|
||||||
@@ -284,7 +106,7 @@ class Connection : public Thread {
|
|||||||
};
|
};
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* class network which listening on tcp socket for incoming connections.
|
* Handler for all TCP and HTTP client connections and registry of active connections.
|
||||||
*/
|
*/
|
||||||
class Network : public Thread {
|
class Network : public Thread {
|
||||||
public:
|
public:
|
||||||
@@ -293,9 +115,9 @@ class Network : public Thread {
|
|||||||
* @param local true to accept connections only for local host.
|
* @param local true to accept connections only for local host.
|
||||||
* @param port the port to listen for command line connections.
|
* @param port the port to listen for command line connections.
|
||||||
* @param httpPort the port to listen for HTTP connections, or 0.
|
* @param httpPort the port to listen for HTTP connections, or 0.
|
||||||
* @param netQueue the reference to the @a NetMessage @a Queue.
|
* @param requestQueue the reference to the @a Request @a Queue.
|
||||||
*/
|
*/
|
||||||
Network(const bool local, const uint16_t port, const uint16_t httpPort, Queue<NetMessage*>* netQueue);
|
Network(const bool local, const uint16_t port, const uint16_t httpPort, Queue<Request*>* requestQueue);
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* destructor.
|
* destructor.
|
||||||
@@ -317,8 +139,8 @@ class Network : public Thread {
|
|||||||
/** the list of active @a Connection instances. */
|
/** the list of active @a Connection instances. */
|
||||||
list<Connection*> m_connections;
|
list<Connection*> m_connections;
|
||||||
|
|
||||||
/** the reference to the @a NetMessage @a Queue. */
|
/** the reference to the @a Request @a Queue. */
|
||||||
Queue<NetMessage*>* m_netQueue;
|
Queue<Request*>* m_requestQueue;
|
||||||
|
|
||||||
/** the command line @a TCPServer instance. */
|
/** the command line @a TCPServer instance. */
|
||||||
TCPServer* m_tcpServer;
|
TCPServer* m_tcpServer;
|
||||||
|
|||||||
@@ -0,0 +1,143 @@
|
|||||||
|
/*
|
||||||
|
* ebusd - daemon for communication with eBUS heating systems.
|
||||||
|
* Copyright (C) 2023 John Baier <ebusd@ebusd.eu>
|
||||||
|
*
|
||||||
|
* 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 3 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, see <http://www.gnu.org/licenses/>.
|
||||||
|
*/
|
||||||
|
|
||||||
|
#ifdef HAVE_CONFIG_H
|
||||||
|
# include <config.h>
|
||||||
|
#endif
|
||||||
|
|
||||||
|
#include "ebusd/request.h"
|
||||||
|
#include <sys/types.h>
|
||||||
|
#include <sys/socket.h>
|
||||||
|
#include <cstring>
|
||||||
|
#include "lib/utils/log.h"
|
||||||
|
|
||||||
|
namespace ebusd {
|
||||||
|
|
||||||
|
RequestImpl::RequestImpl(bool isHttp)
|
||||||
|
: Request(), m_isHttp(isHttp), m_resultSet(false), m_disconnect(false), m_listenSince(0) {
|
||||||
|
m_mode.listenMode = lm_none;
|
||||||
|
m_mode.format = OF_NONE;
|
||||||
|
m_mode.listenWithUnknown = false;
|
||||||
|
m_mode.listenOnlyUnknown = false;
|
||||||
|
pthread_mutex_init(&m_mutex, nullptr);
|
||||||
|
pthread_cond_init(&m_cond, nullptr);
|
||||||
|
}
|
||||||
|
|
||||||
|
RequestImpl::~RequestImpl() {
|
||||||
|
m_resultSet = true;
|
||||||
|
pthread_mutex_destroy(&m_mutex);
|
||||||
|
pthread_cond_destroy(&m_cond);
|
||||||
|
}
|
||||||
|
|
||||||
|
bool RequestImpl::add(const char* request) {
|
||||||
|
if (request && request[0]) {
|
||||||
|
string add = request;
|
||||||
|
add.erase(remove(add.begin(), add.end(), '\r'), add.end());
|
||||||
|
m_request.append(add);
|
||||||
|
}
|
||||||
|
size_t pos = m_request.find(m_isHttp ? "\n\n" : "\n");
|
||||||
|
if (pos != string::npos) {
|
||||||
|
if (m_isHttp) {
|
||||||
|
pos = m_request.find("\n");
|
||||||
|
m_request.resize(pos); // reduce to first line
|
||||||
|
// typical first line: GET /ehp/outsidetemp HTTP/1.1
|
||||||
|
pos = m_request.rfind(" HTTP/");
|
||||||
|
if (pos != string::npos) {
|
||||||
|
m_request.resize(pos); // remove "HTTP/x.x" suffix
|
||||||
|
}
|
||||||
|
pos = 0;
|
||||||
|
while ((pos=m_request.find('%', pos)) != string::npos && pos+2 <= m_request.length()) {
|
||||||
|
unsigned int value1, value2;
|
||||||
|
if (sscanf("%1x%1x", m_request.c_str()+pos+1, &value1, &value2) < 2) {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
m_request[pos] = static_cast<char>(((value1&0x0f) << 4) | (value2&0x0f));
|
||||||
|
m_request.erase(pos+1, 2);
|
||||||
|
}
|
||||||
|
} else if (pos+1 == m_request.length()) {
|
||||||
|
m_request.resize(pos); // reduce to complete lines
|
||||||
|
}
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
return m_request.length() == 0 && m_mode.listenMode != lm_none;
|
||||||
|
}
|
||||||
|
|
||||||
|
void RequestImpl::split(vector<string>* args) {
|
||||||
|
string token, previous;
|
||||||
|
istringstream stream(m_request);
|
||||||
|
char escaped = 0;
|
||||||
|
|
||||||
|
char delim = ' ';
|
||||||
|
while (getline(stream, token, delim)) {
|
||||||
|
if (!m_isHttp) {
|
||||||
|
if (escaped) {
|
||||||
|
args->pop_back();
|
||||||
|
if (token.length() > 0 && token[token.length()-1] == escaped) {
|
||||||
|
token.erase(token.length() - 1, 1);
|
||||||
|
escaped = 0;
|
||||||
|
}
|
||||||
|
token = previous + " " + token;
|
||||||
|
} else if (token.length() == 0) { // allow multiple space chars for a single delimiter
|
||||||
|
continue;
|
||||||
|
} else if (token[0] == '"' || token[0] == '\'') {
|
||||||
|
escaped = token[0];
|
||||||
|
token.erase(0, 1);
|
||||||
|
if (token.length() > 0 && token[token.length()-1] == escaped) {
|
||||||
|
token.erase(token.length() - 1, 1);
|
||||||
|
escaped = 0;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
args->push_back(token);
|
||||||
|
previous = token;
|
||||||
|
if (m_isHttp) {
|
||||||
|
delim = (args->size() == 1) ? '?' : '\n';
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
bool RequestImpl::waitResponse(string* result) {
|
||||||
|
pthread_mutex_lock(&m_mutex);
|
||||||
|
|
||||||
|
if (!m_resultSet) {
|
||||||
|
pthread_cond_wait(&m_cond, &m_mutex);
|
||||||
|
}
|
||||||
|
m_request.clear();
|
||||||
|
*result = m_result;
|
||||||
|
m_result.clear();
|
||||||
|
m_resultSet = false;
|
||||||
|
pthread_mutex_unlock(&m_mutex);
|
||||||
|
return m_disconnect;
|
||||||
|
}
|
||||||
|
|
||||||
|
void RequestImpl::setResult(const string& result, const string& user, RequestMode* mode, time_t listenUntil,
|
||||||
|
bool disconnect) {
|
||||||
|
pthread_mutex_lock(&m_mutex);
|
||||||
|
m_result = result;
|
||||||
|
m_user = user;
|
||||||
|
m_disconnect = disconnect;
|
||||||
|
if (mode) {
|
||||||
|
m_mode = *mode;
|
||||||
|
}
|
||||||
|
m_listenSince = listenUntil;
|
||||||
|
m_resultSet = true;
|
||||||
|
pthread_cond_signal(&m_cond);
|
||||||
|
pthread_mutex_unlock(&m_mutex);
|
||||||
|
}
|
||||||
|
|
||||||
|
} // namespace ebusd
|
||||||
@@ -0,0 +1,230 @@
|
|||||||
|
/*
|
||||||
|
* ebusd - daemon for communication with eBUS heating systems.
|
||||||
|
* Copyright (C) 2023 John Baier <ebusd@ebusd.eu>
|
||||||
|
*
|
||||||
|
* 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 3 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, see <http://www.gnu.org/licenses/>.
|
||||||
|
*/
|
||||||
|
|
||||||
|
#ifndef EBUSD_REQUEST_H_
|
||||||
|
#define EBUSD_REQUEST_H_
|
||||||
|
|
||||||
|
#include <string>
|
||||||
|
#include <cstdio>
|
||||||
|
#include <list>
|
||||||
|
#include <vector>
|
||||||
|
#include "lib/ebus/datatype.h"
|
||||||
|
#include "lib/utils/queue.h"
|
||||||
|
#include "lib/utils/notify.h"
|
||||||
|
#include "lib/utils/thread.h"
|
||||||
|
#include "lib/utils/log.h"
|
||||||
|
|
||||||
|
namespace ebusd {
|
||||||
|
|
||||||
|
/** \file ebusd/request.h
|
||||||
|
* Abstraction of ebusd client requests.
|
||||||
|
*/
|
||||||
|
|
||||||
|
/** the request listen mode. */
|
||||||
|
enum ListenMode {
|
||||||
|
lm_none, //!< normal mode (no listening)
|
||||||
|
lm_listen, //!< listening mode
|
||||||
|
lm_direct, //!< direct mode
|
||||||
|
};
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Request mode info.
|
||||||
|
*/
|
||||||
|
struct RequestMode {
|
||||||
|
ListenMode listenMode; //!< whether in listening or direct mode
|
||||||
|
OutputFormat format; //!< the output format settings for listen mode
|
||||||
|
bool listenWithUnknown; //!< include unknown messages in listen/direct mode
|
||||||
|
bool listenOnlyUnknown; //!< only print unknown messages in listen/direct mode
|
||||||
|
};
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Abstract class for request/response.
|
||||||
|
*/
|
||||||
|
class Request {
|
||||||
|
public:
|
||||||
|
/**
|
||||||
|
* Destructor.
|
||||||
|
*/
|
||||||
|
virtual ~Request() { }
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Add request data from the client.
|
||||||
|
* @param request the request data from the client.
|
||||||
|
* @return true when the request is complete and the response shall be prepared.
|
||||||
|
*/
|
||||||
|
virtual bool add(const char* request) = 0;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* @return whether the request is still empty.
|
||||||
|
*/
|
||||||
|
virtual bool empty() const = 0;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Split the request into arguments.
|
||||||
|
* @param args the @a vector to push the arguments to.
|
||||||
|
*/
|
||||||
|
virtual void split(vector<string>* args) = 0;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Return whether this is a HTTP request.
|
||||||
|
* @return whether this is a HTTP request.
|
||||||
|
*/
|
||||||
|
virtual bool isHttp() const = 0;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Log the request or the given response in debug level.
|
||||||
|
*/
|
||||||
|
virtual void log(const string* response = nullptr) const = 0;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Return the current user name.
|
||||||
|
* @return the current user name.
|
||||||
|
*/
|
||||||
|
virtual const string& getUser() const = 0;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Wait for the response being set and return the result string.
|
||||||
|
* @param result the variable in which to store the result string.
|
||||||
|
* @return true when the client shall be disconnected.
|
||||||
|
*/
|
||||||
|
virtual bool waitResponse(string* result) = 0;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Set the result string and notify a waiting thread.
|
||||||
|
* @param result the result string.
|
||||||
|
* @param user the new user name.
|
||||||
|
* @param newMode the new @a RequestMode.
|
||||||
|
* @param listenUntil the end time to which to updates were added (exclusive).
|
||||||
|
* @param disconnect true when the client shall be disconnected.
|
||||||
|
*/
|
||||||
|
virtual void setResult(const string& result, const string& user, RequestMode* newMode, time_t listenUntil,
|
||||||
|
bool disconnect) = 0;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Return the @a RequestMode.
|
||||||
|
* @param listenSince set listening to the specified start time from which to add updates (inclusive).
|
||||||
|
* @return the @a RequestMode.
|
||||||
|
*/
|
||||||
|
virtual RequestMode getMode(time_t* listenSince = nullptr) = 0;
|
||||||
|
};
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Default @a Request implementation.
|
||||||
|
*/
|
||||||
|
class RequestImpl : public Request {
|
||||||
|
public:
|
||||||
|
/**
|
||||||
|
* Constructor.
|
||||||
|
* @param isHttp whether this is a HTTP request.
|
||||||
|
*/
|
||||||
|
explicit RequestImpl(bool isHttp);
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Destructor.
|
||||||
|
*/
|
||||||
|
virtual ~RequestImpl();
|
||||||
|
|
||||||
|
|
||||||
|
private:
|
||||||
|
/**
|
||||||
|
* Hidden copy constructor.
|
||||||
|
* @param src the object to copy from.
|
||||||
|
*/
|
||||||
|
RequestImpl(const RequestImpl& src);
|
||||||
|
|
||||||
|
|
||||||
|
public:
|
||||||
|
// @copydoc
|
||||||
|
bool add(const char* request) override;
|
||||||
|
|
||||||
|
// @copydoc
|
||||||
|
bool empty() const override { return m_request.empty(); }
|
||||||
|
|
||||||
|
// @copydoc
|
||||||
|
void split(vector<string>* args) override;
|
||||||
|
|
||||||
|
// @copydoc
|
||||||
|
bool isHttp() const override { return m_isHttp; }
|
||||||
|
|
||||||
|
// @copydoc
|
||||||
|
void log(const string* response = nullptr) const override {
|
||||||
|
if (response) {
|
||||||
|
if (response->length() > 100) {
|
||||||
|
logDebug(lf_main, "<<< %s ...", response->substr(0, 100).c_str());
|
||||||
|
} else {
|
||||||
|
logDebug(lf_main, "<<< %s", response->c_str());
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
logDebug(lf_main, ">>> %s", m_request.c_str());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// @copydoc
|
||||||
|
const string& getUser() const override { return m_user; }
|
||||||
|
|
||||||
|
// @copydoc
|
||||||
|
bool waitResponse(string* result) override;
|
||||||
|
|
||||||
|
// @copydoc
|
||||||
|
void setResult(const string& result, const string& user, RequestMode* mode, time_t listenUntil,
|
||||||
|
bool disconnect) override;
|
||||||
|
|
||||||
|
// @copydoc
|
||||||
|
RequestMode getMode(time_t* listenSince = nullptr) override {
|
||||||
|
if (listenSince) {
|
||||||
|
*listenSince = m_listenSince;
|
||||||
|
}
|
||||||
|
return m_mode;
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
private:
|
||||||
|
/** whether this is a HTTP message. */
|
||||||
|
const bool m_isHttp;
|
||||||
|
|
||||||
|
/** the request string. */
|
||||||
|
string m_request;
|
||||||
|
|
||||||
|
/** the current user name. */
|
||||||
|
string m_user;
|
||||||
|
|
||||||
|
/** whether the result was already set. */
|
||||||
|
bool m_resultSet;
|
||||||
|
|
||||||
|
/** the result string. */
|
||||||
|
string m_result;
|
||||||
|
|
||||||
|
/** set to true when the client shall be disconnected. */
|
||||||
|
bool m_disconnect;
|
||||||
|
|
||||||
|
/** mutex variable for exclusive lock. */
|
||||||
|
pthread_mutex_t m_mutex;
|
||||||
|
|
||||||
|
/** condition variable for exclusive lock. */
|
||||||
|
pthread_cond_t m_cond;
|
||||||
|
|
||||||
|
/** the @a RequestMode. */
|
||||||
|
RequestMode m_mode;
|
||||||
|
|
||||||
|
/** start timestamp of listening update. */
|
||||||
|
time_t m_listenSince;
|
||||||
|
};
|
||||||
|
|
||||||
|
} // namespace ebusd
|
||||||
|
|
||||||
|
#endif // EBUSD_REQUEST_H_
|
||||||
Reference in New Issue
Block a user