class WQueue: method 'T remove(const long delay)' added.
This commit is contained in:
+32
-1
@@ -20,8 +20,11 @@
|
|||||||
#ifndef WQUEUE_H_
|
#ifndef WQUEUE_H_
|
||||||
#define WQUEUE_H_
|
#define WQUEUE_H_
|
||||||
|
|
||||||
#include <pthread.h>
|
|
||||||
#include <list>
|
#include <list>
|
||||||
|
#include <ctime>
|
||||||
|
#include <pthread.h>
|
||||||
|
#include <sys/time.h>
|
||||||
|
#include <errno.h>
|
||||||
|
|
||||||
template <typename T> class WQueue
|
template <typename T> class WQueue
|
||||||
{
|
{
|
||||||
@@ -64,6 +67,34 @@ public:
|
|||||||
return item;
|
return item;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
T remove(const long delay)
|
||||||
|
{
|
||||||
|
struct timeval tv;
|
||||||
|
struct timezone tz;
|
||||||
|
struct timespec timeout;
|
||||||
|
|
||||||
|
int ret = 0;
|
||||||
|
|
||||||
|
pthread_mutex_lock(&m_mutex);
|
||||||
|
|
||||||
|
gettimeofday(&tv, &tz);
|
||||||
|
timeout.tv_sec = tv.tv_sec + delay;
|
||||||
|
timeout.tv_nsec = tv.tv_usec * 1000;
|
||||||
|
|
||||||
|
while (m_queue.size() == 0 && ret != ETIMEDOUT)
|
||||||
|
ret = pthread_cond_timedwait(&m_condv, &m_mutex, &timeout);
|
||||||
|
|
||||||
|
if (ret == ETIMEDOUT)
|
||||||
|
return NULL;
|
||||||
|
|
||||||
|
T item = m_queue.front();
|
||||||
|
m_queue.pop_front();
|
||||||
|
|
||||||
|
pthread_mutex_unlock(&m_mutex);
|
||||||
|
|
||||||
|
return item;
|
||||||
|
}
|
||||||
|
|
||||||
int size()
|
int size()
|
||||||
{
|
{
|
||||||
pthread_mutex_lock(&m_mutex);
|
pthread_mutex_lock(&m_mutex);
|
||||||
|
|||||||
+41
-24
@@ -96,22 +96,27 @@ std::string BaseLoop::decodeMessage(const std::string& data)
|
|||||||
m_ebusloop->addBusCommand(new BusCommand(type, ebusCommand));
|
m_ebusloop->addBusCommand(new BusCommand(type, ebusCommand));
|
||||||
BusCommand* busCommand = m_ebusloop->getBusCommand();
|
BusCommand* busCommand = m_ebusloop->getBusCommand();
|
||||||
|
|
||||||
if (busCommand->getResult().c_str()[0] != '-') {
|
if (busCommand != NULL) {
|
||||||
// decode data
|
if (busCommand->getResult().c_str()[0] != '-') {
|
||||||
Command* command = new Command(index, (*m_commands)[index], busCommand->getResult());
|
// decode data
|
||||||
|
Command* command = new Command(index, (*m_commands)[index], busCommand->getResult());
|
||||||
|
|
||||||
// return result
|
// return result
|
||||||
result << command->calcResult(cmd);
|
result << command->calcResult(cmd);
|
||||||
|
|
||||||
delete command;
|
delete command;
|
||||||
|
} else {
|
||||||
|
L.log(bas, error, " %s", busCommand->getResult().c_str());
|
||||||
|
result << busCommand->getResult();
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
delete busCommand;
|
||||||
} else {
|
} else {
|
||||||
L.log(bas, error, " %s", busCommand->getResult().c_str());
|
L.log(bas, error, " busCommand timeout reached");
|
||||||
result << busCommand->getResult();
|
result << "busCommand timeout reached";
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
delete busCommand;
|
|
||||||
|
|
||||||
} else {
|
} else {
|
||||||
result << "ebus command not found";
|
result << "ebus command not found";
|
||||||
}
|
}
|
||||||
@@ -150,19 +155,25 @@ std::string BaseLoop::decodeMessage(const std::string& data)
|
|||||||
m_ebusloop->addBusCommand(new BusCommand(type, ebusCommand));
|
m_ebusloop->addBusCommand(new BusCommand(type, ebusCommand));
|
||||||
BusCommand* busCommand = m_ebusloop->getBusCommand();
|
BusCommand* busCommand = m_ebusloop->getBusCommand();
|
||||||
|
|
||||||
if (busCommand->getResult().c_str()[0] != '-') {
|
if (busCommand != NULL) {
|
||||||
// decode result
|
if (busCommand->getResult().c_str()[0] != '-') {
|
||||||
if (busCommand->getResult().substr(busCommand->getResult().length()-8) == "00000000")
|
// decode result
|
||||||
result << "done";
|
if (busCommand->getResult().substr(busCommand->getResult().length()-8) == "00000000")
|
||||||
else
|
result << "done";
|
||||||
result << "error";
|
else
|
||||||
|
result << "error";
|
||||||
|
|
||||||
|
} else {
|
||||||
|
L.log(bas, error, " %s", busCommand->getResult().c_str());
|
||||||
|
result << busCommand->getResult();
|
||||||
|
}
|
||||||
|
|
||||||
|
delete busCommand;
|
||||||
} else {
|
} else {
|
||||||
L.log(bas, error, " %s", busCommand->getResult().c_str());
|
L.log(bas, error, " busCommand timeout reached");
|
||||||
result << busCommand->getResult();
|
result << "busCommand timeout reached";
|
||||||
}
|
}
|
||||||
|
|
||||||
delete busCommand;
|
|
||||||
delete command;
|
delete command;
|
||||||
|
|
||||||
} else {
|
} else {
|
||||||
@@ -220,12 +231,18 @@ std::string BaseLoop::decodeMessage(const std::string& data)
|
|||||||
m_ebusloop->addBusCommand(new BusCommand(type, ebusCommand));
|
m_ebusloop->addBusCommand(new BusCommand(type, ebusCommand));
|
||||||
BusCommand* busCommand = m_ebusloop->getBusCommand();
|
BusCommand* busCommand = m_ebusloop->getBusCommand();
|
||||||
|
|
||||||
if (busCommand->getResult().c_str()[0] == '-')
|
if (busCommand != NULL) {
|
||||||
L.log(bas, error, " %s", busCommand->getResult().c_str());
|
if (busCommand->getResult().c_str()[0] == '-')
|
||||||
|
L.log(bas, error, " %s", busCommand->getResult().c_str());
|
||||||
|
|
||||||
result << busCommand->getResult();
|
result << busCommand->getResult();
|
||||||
|
|
||||||
|
delete busCommand;
|
||||||
|
} else {
|
||||||
|
L.log(bas, error, " busCommand timeout reached");
|
||||||
|
result << "busCommand timeout reached";
|
||||||
|
}
|
||||||
|
|
||||||
delete busCommand;
|
|
||||||
} else {
|
} else {
|
||||||
result << "specified message type is incorrect";
|
result << "specified message type is incorrect";
|
||||||
}
|
}
|
||||||
|
|||||||
+1
-1
@@ -40,7 +40,7 @@ public:
|
|||||||
std::string getData() { return m_cycBuffer.remove(); }
|
std::string getData() { return m_cycBuffer.remove(); }
|
||||||
|
|
||||||
void addBusCommand(BusCommand* busCommand) { m_sendBuffer.add(busCommand); }
|
void addBusCommand(BusCommand* busCommand) { m_sendBuffer.add(busCommand); }
|
||||||
BusCommand* getBusCommand() { return m_recvBuffer.remove(); }
|
BusCommand* getBusCommand() { return m_recvBuffer.remove(10); }
|
||||||
|
|
||||||
void dump(const bool dumpState) { m_bus->setDumpState(dumpState); }
|
void dump(const bool dumpState) { m_bus->setDumpState(dumpState); }
|
||||||
|
|
||||||
|
|||||||
+1
-1
@@ -122,7 +122,7 @@ Network::Network(const bool localhost) : m_listening(false), m_running(false)
|
|||||||
else
|
else
|
||||||
m_Server = new TCPServer(A.getParam<int>("p_port"), "0.0.0.0");
|
m_Server = new TCPServer(A.getParam<int>("p_port"), "0.0.0.0");
|
||||||
|
|
||||||
if (m_Server && m_Server->start() == 0)
|
if (m_Server != NULL && m_Server->start() == 0)
|
||||||
m_listening = true;
|
m_listening = true;
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user