more thread abstraction

This commit is contained in:
canpadawan
2016-06-16 13:25:56 +02:00
parent 71e989ad05
commit d1c5642668
6 changed files with 334 additions and 231 deletions
+146 -2
View File
@@ -1,35 +1,179 @@
#include <QThread>
#include "canconnection.h"
CANConnection::CANConnection(QString pPort, CANCon::type pType, int pNumBuses) :
CANConnection::CANConnection(QString pPort,
CANCon::type pType,
int pNumBuses,
int pQueueLen,
bool pUseThread) :
mQueue(),
mNumBuses(pNumBuses),
mPort(pPort),
mType(pType),
mIsCapSuspended(false),
mStatus(CANCon::NOT_CONNECTED)
mStatus(CANCon::NOT_CONNECTED),
mThread_p(NULL)
{
qDebug() << "CANConnection()";
/* register types */
qRegisterMetaType<CANBus>("CANBus");
qRegisterMetaType<CANFrame>("CANFrame");
qRegisterMetaType<CANCon::status>("CANCon::status");
/* set queue size */
mQueue.setSize(pQueueLen); /*TODO add check on returned value */
/* allocate buses */
mBus = new CANBus[mNumBuses];
mConfigured = new bool[mNumBuses];
for(int i=0 ; i<mNumBuses ; i++)
mConfigured[i] = false;
/* if needed, create a thread and move ourself into it */
if(pUseThread) {
mThread_p = new QThread();
}
}
CANConnection::~CANConnection()
{
qDebug() << "~CANConnection()";
/* stop and delete thread */
if(mThread_p) {
mThread_p->quit();
mThread_p->wait();
delete mThread_p;
mThread_p = NULL;
}
/* delete bus table */
delete[] mBus;
mBus = NULL;
/* configured table */
delete[] mConfigured;
mConfigured = NULL;
/* delete queue table */
mQueue.setSize(0);
}
void CANConnection::start()
{
if( mThread_p && (mThread_p != QThread::currentThread()) )
{
/* move ourself to the thread */
moveToThread(mThread_p); /*TODO handle errors */
/* connect started() */
connect(mThread_p, SIGNAL(started()), this, SLOT(start()));
/* start the thread */
mThread_p->start(QThread::HighPriority);
return;
}
/* in multithread case, this will be called before entering thread event loop */
return piStarted();
}
void CANConnection::suspend(bool pSuspend)
{
/* execute in mThread_p context */
if( mThread_p && (mThread_p != QThread::currentThread()) ) {
QMetaObject::invokeMethod(this, "suspend",
Qt::BlockingQueuedConnection,
Q_ARG(bool, pSuspend));
return;
}
return piSuspend(pSuspend);
}
void CANConnection::stop()
{
/* 1) execute in mThread_p context */
if( mThread_p && (mThread_p != QThread::currentThread()) )
{
/* if thread is finished, it means we call this function for the second time so we can leave */
if( !mThread_p->isFinished() )
{
/* we need to call piStop() */
QMetaObject::invokeMethod(this, "stop",
Qt::BlockingQueuedConnection);
/* 3) stop thread */
mThread_p->quit();
if(!mThread_p->wait()) {
qDebug() << "can't stop thread";
}
}
return;
}
/* 2) call piStop in mThread context */
return piStop();
}
bool CANConnection::getBusSettings(int pBusIdx, CANBus& pBus)
{
/* make sure we execute in mThread context */
if( mThread_p && (mThread_p != QThread::currentThread()) ) {
bool ret;
QMetaObject::invokeMethod(this, "getBusSettings",
Qt::BlockingQueuedConnection,
Q_RETURN_ARG(bool, ret),
Q_ARG(int , pBusIdx),
Q_ARG(CANBus& , pBus));
return ret;
}
return piGetBusSettings(pBusIdx, pBus);
}
void CANConnection::setBusSettings(int pBusIdx, CANBus pBus)
{
/* make sure we execute in mThread context */
if( mThread_p && (mThread_p != QThread::currentThread()) ) {
QMetaObject::invokeMethod(this, "setBusSettings",
Qt::BlockingQueuedConnection,
Q_ARG(int, pBusIdx),
Q_ARG(CANBus, pBus));
return;
}
return piSetBusSettings(pBusIdx, pBus);
}
void CANConnection::sendFrame(const CANFrame& pFrame)
{
/* make sure we execute in mThread context */
if( mThread_p && (mThread_p != QThread::currentThread()) ) {
QMetaObject::invokeMethod(this, "sendFrame",
Qt::BlockingQueuedConnection,
Q_ARG(const CANFrame&, pFrame));
return;
}
return piSendFrame(pFrame);
}
void CANConnection::sendFrameBatch(const QList<CANFrame>& pFrames)
{
/* make sure we execute in mThread context */
if( mThread_p && (mThread_p != QThread::currentThread()) ) {
QMetaObject::invokeMethod(this, "sendFrameBatch",
Qt::BlockingQueuedConnection,
Q_ARG(const QList<CANFrame>&, pFrames));
return;
}
return piSendFrameBatch(pFrames);
}
+65 -23
View File
@@ -20,8 +20,11 @@ public:
* @param pPort: string containing port name
* @param pType: the type of connection @ref CANCon::type
* @param pNumBuses: the number of buses the device has
* @param pQueueLen: the length of the lock free queue to use
* @param pUseThread: if set to true, object will be execute in a dedicated thread
*/
CANConnection(QString pPort, CANCon::type pType, int pNumBuses);
CANConnection(QString pPort, CANCon::type pType, int pNumBuses,
int pQueueLen, bool pUseThread);
/**
* @brief CANConnection destructor
*/
@@ -31,50 +34,33 @@ public:
/**
* @brief getNumBuses
* @return returns the number of buses of the device
* @note multithread safe
*/
int getNumBuses();
/**
* @brief getPort
* @return returns the port name of the device
* @note multithread safe
*/
QString getPort();
/**
* @brief getQueue is call by reader to get a reference on the queue to monitor
* @return the lock free queue of the device
* @note multithread safe
*/
LFQueue<CANFrame>& getQueue();
/**
* @brief getType
* @return the @ref CANCon::type of the device
* @note multithread safe
*/
CANCon::type getType();
/**
* @brief getStatus
* @return the @ref CANCon::status of the device (either connected or not)
* @note multithread safe
*/
CANCon::status getStatus();
/**
* @brief start the device
* @note start a working thread here if needed
*/
virtual void start() = 0;
/**
* @brief stop the device
* @note stop the working thread here if one has been started
*/
virtual void stop() = 0;
signals:
void error(const QString &);
@@ -95,30 +81,45 @@ signals:
public slots:
virtual void sendFrame(const CANFrame *) = 0;
virtual void sendFrameBatch(const QList<CANFrame> *) = 0;
/**
* @brief start the device, this calls piStarted
* @note starts the working thread if required (piStarted in the working thread context)
*/
void start();
/**
* @brief stop the device, this calls piStop
* @note if a working thread is used, piStop is called before exiting the working thread
*/
void stop();
/**
* @brief setBusSettings
* @param pBusIdx: the index of the bus for which settings have to be set
* @param pBus: the settings to set
* @note this calls piSetBusSettings in the working thread context (if one has been started)
*/
virtual void setBusSettings(int pBusIdx, CANBus pBus) = 0;
void setBusSettings(int pBusIdx, CANBus pBus);
/**
* @brief getBusSettings
* @param pBusIdx: the index of the bus for which settings have to be retrieved
* @param pBus: the CANBus struct to fill with information
* @return true if operation succeeds, false if pBusIdx is invalid or bus has not been configured yet
* @note this calls piGetBusSettings in the working thread context (if one has been started)
*/
virtual bool getBusSettings(int pBusIdx, CANBus& pBus) = 0;
bool getBusSettings(int pBusIdx, CANBus& pBus);
/**
* @brief suspends/restarts data capture
* @param pSuspend: suspends capture if true else restarts it
* @note this calls piSuspend in the working thread context (if one has been started)
* @note the caller will not access the queue when capture is suspended, so it is safe for callee to flush the queue
*/
virtual void suspend(bool pSuspend) = 0;
void suspend(bool pSuspend);
void sendFrame(const CANFrame&);
void sendFrameBatch(const QList<CANFrame>&);
protected:
@@ -177,6 +178,46 @@ protected:
*/
void setCapSuspended(bool pIsSuspended);
protected:
/**
* @brief start the device
* @note start a working thread here if needed
*/
virtual void piStarted() = 0;
/**
* @brief stop the device
* @note stop the working thread here if one has been started
*/
virtual void piStop() = 0;
/**
* @brief setBusSettings
* @param pBusIdx: the index of the bus for which settings have to be set
* @param pBus: the settings to set
*/
virtual void piSetBusSettings(int pBusIdx, CANBus pBus) = 0;
/**
* @brief getBusSettings
* @param pBusIdx: the index of the bus for which settings have to be retrieved
* @param pBus: the CANBus struct to fill with information
* @return true if operation succeeds, false if pBusIdx is invalid or bus has not been configured yet
*/
virtual bool piGetBusSettings(int pBusIdx, CANBus& pBus) = 0;
/**
* @brief suspends/restarts data capture
* @param pSuspend: suspends capture if true else restarts it
* @note the caller will not access the queue when capture is suspended, so it is safe for callee to flush the queue
*/
virtual void piSuspend(bool pSuspend) = 0;
virtual void piSendFrame(const CANFrame&) = 0;
virtual void piSendFrameBatch(const QList<CANFrame>&) = 0;
private:
CANBus* mBus;
bool* mConfigured;
@@ -186,6 +227,7 @@ private:
const CANCon::type mType;
bool mIsCapSuspended;
QAtomicInt mStatus;
QThread* mThread_p;
};
#endif // CANCONNECTION_H
+87 -120
View File
@@ -7,13 +7,9 @@
#include "gvretserial.h"
GVRetSerial::GVRetSerial(QString portName) :
CANConnection(portName, CANCon::GVRET_SERIAL, 2),
mThread(NULL)
CANConnection(portName, CANCon::GVRET_SERIAL, 2, 4000, true),
mTimer(this) /*NB: set this as parent of timer to manage it from working thread */
{
getQueue().setSize(2000); /*TODO add check on returned value */
/* move ourself to the thread */
moveToThread(&mThread);
qDebug() << "GVRetSerial()";
serial = NULL;
@@ -35,13 +31,9 @@ GVRetSerial::~GVRetSerial()
}
void GVRetSerial::start()
void GVRetSerial::piStarted()
{
qDebug() << "enter thread";
/* start thread */
mThread.start(QThread::HighPriority);
/* connect device in thread context */
connect(&mThread, SIGNAL(started()), this, SLOT(connectDevice()));
connectDevice();
/* start timer */
connect(&mTimer, SIGNAL(timeout()), this, SLOT(handleTick()));
@@ -51,15 +43,8 @@ void GVRetSerial::start()
}
void GVRetSerial::suspend(bool pSuspend) {
/* make sure we execute in mThread context */
if(QThread::currentThread() != &mThread) {
QMetaObject::invokeMethod(this, "suspend",
Qt::BlockingQueuedConnection,
Q_ARG(bool, pSuspend));
return;
}
void GVRetSerial::piSuspend(bool pSuspend)
{
/* update capSuspended */
setCapSuspended(pSuspend);
@@ -69,36 +54,98 @@ void GVRetSerial::suspend(bool pSuspend) {
}
void GVRetSerial::stop() {
void GVRetSerial::piStop()
{
mTimer.stop();
mThread.quit();
if(!mThread.wait()) {
qDebug() << "can't stop thread";
}
disconnectDevice();
}
bool GVRetSerial::getBusSettings(int pBusIdx, CANBus& pBus) {
/* make sure we execute in mThread context */
if(QThread::currentThread() != &mThread) {
bool ret;
QMetaObject::invokeMethod(this, "getBusSettings",
Qt::BlockingQueuedConnection,
Q_RETURN_ARG(bool, ret),
Q_ARG(int , pBusIdx),
Q_ARG(CANBus& , pBus));
return ret;
}
bool GVRetSerial::piGetBusSettings(int pBusIdx, CANBus& pBus)
{
return getBusConfig(pBusIdx, pBus);
}
void GVRetSerial::sendFrame(const CANFrame *) {}
void GVRetSerial::sendFrameBatch(const QList<CANFrame> *){}
void GVRetSerial::piSetBusSettings(int pBusIdx, CANBus bus)
{
/* sanity checks */
if( (pBusIdx < 0) || pBusIdx >= getNumBuses())
return;
/* copy bus config */
setBusConfig(pBusIdx, bus);
qDebug() << "About to update bus " << pBusIdx << " on GVRET";
if (pBusIdx == 0)
{
can0Baud = bus.getSpeed();
can0Baud |= 0x80000000;
if (bus.isActive())
{
can0Baud |= 0x40000000;
can0Enabled = true;
}
else can0Enabled = false;
if (bus.isListenOnly())
{
can0Baud |= 0x20000000;
can0ListenOnly = true;
}
else can0ListenOnly = false;
}
else if (pBusIdx == 1)
{
can1Baud = bus.getSpeed();
can1Baud |= 0x80000000;
if (bus.isActive())
{
can1Baud |= 0x40000000;
can1Enabled = true;
}
else can1Enabled = false;
if (bus.isListenOnly())
{
can1Baud |= 0x20000000;
can1ListenOnly = true;
}
else can1ListenOnly = false;
if (bus.isSingleWire())
{
can1Baud |= 0x10000000;
deviceSingleWireMode = 1;
}
else deviceSingleWireMode = 0;
}
/* update baud rates */
QByteArray buffer;
qDebug() << "Got signal to update bauds. 1: " << can0Baud <<" 2: " << can1Baud;
buffer[0] = (char)0xF1; //start of a command over serial
buffer[1] = 5; //setup canbus
buffer[2] = (unsigned char)(can0Baud & 0xFF); //four bytes of ID LSB first
buffer[3] = (unsigned char)(can0Baud >> 8);
buffer[4] = (unsigned char)(can0Baud >> 16);
buffer[5] = (unsigned char)(can0Baud >> 24);
buffer[6] = (unsigned char)(can1Baud & 0xFF); //four bytes of ID LSB first
buffer[7] = (unsigned char)(can1Baud >> 8);
buffer[8] = (unsigned char)(can1Baud >> 16);
buffer[9] = (unsigned char)(can1Baud >> 24);
buffer[10] = 0;
if (serial == NULL) return;
if (!serial->isOpen()) return;
serial->write(buffer);
}
void GVRetSerial::piSendFrame(const CANFrame&) {}
void GVRetSerial::piSendFrameBatch(const QList<CANFrame>&){}
/****************************************************************/
void GVRetSerial::readSettings()
{
@@ -528,84 +575,4 @@ void GVRetSerial::sendCommValidation()
}
void GVRetSerial::setBusSettings(int pBusIdx, CANBus bus)
{
/* make sure we execute in mThread context */
if(QThread::currentThread() != &mThread) {
QMetaObject::invokeMethod(this, "setBusSettings",
Qt::BlockingQueuedConnection,
Q_ARG(int, pBusIdx),
Q_ARG(CANBus, bus));
return;
}
/* sanity checks */
if( (pBusIdx < 0) || pBusIdx >= getNumBuses())
return;
/* copy bus config */
setBusConfig(pBusIdx, bus);
qDebug() << "About to update bus " << pBusIdx << " on GVRET";
if (pBusIdx == 0)
{
can0Baud = bus.getSpeed();
can0Baud |= 0x80000000;
if (bus.isActive())
{
can0Baud |= 0x40000000;
can0Enabled = true;
}
else can0Enabled = false;
if (bus.isListenOnly())
{
can0Baud |= 0x20000000;
can0ListenOnly = true;
}
else can0ListenOnly = false;
}
else if (pBusIdx == 1)
{
can1Baud = bus.getSpeed();
can1Baud |= 0x80000000;
if (bus.isActive())
{
can1Baud |= 0x40000000;
can1Enabled = true;
}
else can1Enabled = false;
if (bus.isListenOnly())
{
can1Baud |= 0x20000000;
can1ListenOnly = true;
}
else can1ListenOnly = false;
if (bus.isSingleWire())
{
can1Baud |= 0x10000000;
deviceSingleWireMode = 1;
}
else deviceSingleWireMode = 0;
}
/* update baud rates */
QByteArray buffer;
qDebug() << "Got signal to update bauds. 1: " << can0Baud <<" 2: " << can1Baud;
buffer[0] = (char)0xF1; //start of a command over serial
buffer[1] = 5; //setup canbus
buffer[2] = (unsigned char)(can0Baud & 0xFF); //four bytes of ID LSB first
buffer[3] = (unsigned char)(can0Baud >> 8);
buffer[4] = (unsigned char)(can0Baud >> 16);
buffer[5] = (unsigned char)(can0Baud >> 24);
buffer[6] = (unsigned char)(can1Baud & 0xFF); //four bytes of ID LSB first
buffer[7] = (unsigned char)(can1Baud >> 8);
buffer[8] = (unsigned char)(can1Baud >> 16);
buffer[9] = (unsigned char)(can1Baud >> 24);
buffer[10] = 0;
if (serial == NULL) return;
if (!serial->isOpen()) return;
serial->write(buffer);
}
+9 -15
View File
@@ -44,10 +44,6 @@ public:
GVRetSerial(QString portName);
virtual ~GVRetSerial();
virtual void start();
virtual void stop();
signals:
void error(const QString &);
@@ -60,18 +56,16 @@ signals:
//being passed. Just set for things that really are being updated.
void busStatus(int, int, int);
public slots:
virtual void sendFrame(const CANFrame *);
virtual void sendFrameBatch(const QList<CANFrame> *);
virtual void setBusSettings(int, CANBus);
virtual bool getBusSettings(int pBusIdx, CANBus& pBus);
virtual void suspend(bool);
protected:
virtual void piStarted();
virtual void piStop();
virtual void piSetBusSettings(int pBusIdx, CANBus pBus);
virtual bool piGetBusSettings(int pBusIdx, CANBus& pBus);
virtual void piSuspend(bool pSuspend);
virtual void piSendFrame(const CANFrame&) ;
virtual void piSendFrameBatch(const QList<CANFrame>&);
void disconnectDevice();
private slots:
+18 -55
View File
@@ -11,14 +11,10 @@
/***********************************/
SocketCanConnection::SocketCanConnection(QString portName) :
CANConnection(portName, CANCon::SOCKETCAN, 1),
CANConnection(portName, CANCon::SOCKETCAN, 1, 4000, true),
mDev_p(NULL),
mThread(NULL)
mTimer(this) /*NB: set connection as parent of timer to manage it from working thread */
{
getQueue().setSize(2000); /*TODO add check on returned value */
/* move ourself to the thread */
moveToThread(&mThread);
qDebug() << "SocketCanConnection()";
}
@@ -30,11 +26,8 @@ SocketCanConnection::~SocketCanConnection()
}
void SocketCanConnection::start() {
qDebug() << "enter thread";
mThread.start(QThread::HighPriority);
void SocketCanConnection::piStarted()
{
connect(&mTimer, SIGNAL(timeout()), this, SLOT(testConnection()));
mTimer.setInterval(1000);
mTimer.setSingleShot(false); //keep ticking
@@ -42,15 +35,8 @@ void SocketCanConnection::start() {
}
void SocketCanConnection::suspend(bool pSuspend) {
/* make sure we execute in mThread context */
if(QThread::currentThread() != &mThread) {
QMetaObject::invokeMethod(this, "suspend",
Qt::BlockingQueuedConnection,
Q_ARG(bool, pSuspend));
return;
}
void SocketCanConnection::piSuspend(bool pSuspend)
{
/* update capSuspended */
setCapSuspended(pSuspend);
@@ -60,43 +46,20 @@ void SocketCanConnection::suspend(bool pSuspend) {
}
void SocketCanConnection::stop() {
void SocketCanConnection::piStop() {
mTimer.stop();
mThread.quit();
if(!mThread.wait()) {
qDebug() << "can't stop thread";
}
disconnectDevice();
}
bool SocketCanConnection::getBusSettings(int pBusIdx, CANBus& pBus) {
/* make sure we execute in mThread context */
if(QThread::currentThread() != &mThread) {
bool ret;
QMetaObject::invokeMethod(this, "getBusSettings",
Qt::BlockingQueuedConnection,
Q_RETURN_ARG(bool, ret),
Q_ARG(int , pBusIdx),
Q_ARG(CANBus& , pBus));
return ret;
}
bool SocketCanConnection::piGetBusSettings(int pBusIdx, CANBus& pBus)
{
return getBusConfig(pBusIdx, pBus);
}
void SocketCanConnection::setBusSettings(int pBusIdx, CANBus bus)
void SocketCanConnection::piSetBusSettings(int pBusIdx, CANBus bus)
{
/* make sure we execute in mThread context */
if(QThread::currentThread() != &mThread) {
QMetaObject::invokeMethod(this, "setBusSettings",
Qt::BlockingQueuedConnection,
Q_ARG(int, pBusIdx),
Q_ARG(CANBus, bus));
return;
}
/* sanity checks */
if(0 != pBusIdx)
return;
@@ -139,8 +102,8 @@ void SocketCanConnection::setBusSettings(int pBusIdx, CANBus bus)
}
void SocketCanConnection::sendFrame(const CANFrame *) {}
void SocketCanConnection::sendFrameBatch(const QList<CANFrame> *){}
void SocketCanConnection::piSendFrame(const CANFrame&) {}
void SocketCanConnection::piSendFrameBatch(const QList<CANFrame>&){}
/***********************************/
@@ -239,12 +202,12 @@ void SocketCanConnection::testConnection() {
break;
case CANCon::NOT_CONNECTED:
if (dev_p && dev_p->connectDevice()) {
/* try to reconnect */
CANBus bus;
if(getBusConfig(0, bus))
setBusSettings(0, bus);
if(!mDev_p) {
/* try to reconnect */
CANBus bus;
if(getBusConfig(0, bus))
setBusSettings(0, bus);
}
/* disconnect test instance */
dev_p->disconnectDevice();
+9 -16
View File
@@ -18,10 +18,6 @@ public:
SocketCanConnection(QString portName);
virtual ~SocketCanConnection();
virtual void start();
virtual void stop();
signals:
void error(const QString &);
@@ -34,18 +30,16 @@ signals:
//being passed. Just set for things that really are being updated.
void busStatus(int, int, int);
public slots:
virtual void sendFrame(const CANFrame *);
virtual void sendFrameBatch(const QList<CANFrame> *);
virtual void setBusSettings(int, CANBus);
virtual bool getBusSettings(int pBusIdx, CANBus& pBus);
virtual void suspend(bool);
protected:
virtual void piStarted();
virtual void piStop();
virtual void piSetBusSettings(int pBusIdx, CANBus pBus);
virtual bool piGetBusSettings(int pBusIdx, CANBus& pBus);
virtual void piSuspend(bool pSuspend);
virtual void piSendFrame(const CANFrame&) ;
virtual void piSendFrameBatch(const QList<CANFrame>&);
void disconnectDevice();
private slots:
@@ -57,7 +51,6 @@ private slots:
protected:
QCanBusDevice* mDev_p;
QTimer mTimer;
QThread mThread;
};