From d1c56426685109b07b2b8900faae2940ec9a90d1 Mon Sep 17 00:00:00 2001 From: canpadawan Date: Thu, 16 Jun 2016 13:25:56 +0200 Subject: [PATCH] more thread abstraction --- connections/canconnection.cpp | 148 +++++++++++++++++++++++- connections/canconnection.h | 88 +++++++++++---- connections/gvretserial.cpp | 207 ++++++++++++++-------------------- connections/gvretserial.h | 24 ++-- connections/socketcan.cpp | 73 +++--------- connections/socketcan.h | 25 ++-- 6 files changed, 334 insertions(+), 231 deletions(-) diff --git a/connections/canconnection.cpp b/connections/canconnection.cpp index b3059b8..7ec3ab7 100644 --- a/connections/canconnection.cpp +++ b/connections/canconnection.cpp @@ -1,35 +1,179 @@ +#include #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"); + qRegisterMetaType("CANFrame"); qRegisterMetaType("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 ; iquit(); + 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& 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&, pFrames)); + return; + } + + return piSendFrameBatch(pFrames); } diff --git a/connections/canconnection.h b/connections/canconnection.h index 5c9c68f..eb02f60 100644 --- a/connections/canconnection.h +++ b/connections/canconnection.h @@ -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& 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 *) = 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&); 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&) = 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 diff --git a/connections/gvretserial.cpp b/connections/gvretserial.cpp index 96666d5..1d34f38 100644 --- a/connections/gvretserial.cpp +++ b/connections/gvretserial.cpp @@ -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 *){} +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&){} + + +/****************************************************************/ 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); -} diff --git a/connections/gvretserial.h b/connections/gvretserial.h index b223aee..350277f 100644 --- a/connections/gvretserial.h +++ b/connections/gvretserial.h @@ -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 *); - - 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&); + void disconnectDevice(); private slots: diff --git a/connections/socketcan.cpp b/connections/socketcan.cpp index 76b6fe3..2b41f2e 100644 --- a/connections/socketcan.cpp +++ b/connections/socketcan.cpp @@ -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 *){} +void SocketCanConnection::piSendFrame(const CANFrame&) {} +void SocketCanConnection::piSendFrameBatch(const QList&){} /***********************************/ @@ -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(); diff --git a/connections/socketcan.h b/connections/socketcan.h index eea37a9..315aad1 100644 --- a/connections/socketcan.h +++ b/connections/socketcan.h @@ -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 *); - - 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&); + void disconnectDevice(); private slots: @@ -57,7 +51,6 @@ private slots: protected: QCanBusDevice* mDev_p; QTimer mTimer; - QThread mThread; };