From 121ad369e8f8bda0cc238a8e48b69ae9dd5a435f Mon Sep 17 00:00:00 2001 From: john30 Date: Sat, 27 Dec 2014 13:37:42 +0100 Subject: [PATCH] solved TODO and fixed issue with concurrent access to BusRequest, added support for determining signal timeout --- src/ebusd/bushandler.cpp | 181 +++++++++++++++++++-------------------- src/ebusd/bushandler.h | 38 ++++---- 2 files changed, 108 insertions(+), 111 deletions(-) diff --git a/src/ebusd/bushandler.cpp b/src/ebusd/bushandler.cpp index 393c08dd..fed9cd22 100644 --- a/src/ebusd/bushandler.cpp +++ b/src/ebusd/bushandler.cpp @@ -43,6 +43,7 @@ extern Logger& L; const char* getStateCode(BusState state) { switch (state) { + case bs_noSignal: return "no signal"; case bs_skip: return "skip"; case bs_ready: return "ready"; case bs_sendCmd: return "send command"; @@ -131,54 +132,15 @@ bool ScanRequest::notify(result_t result, SymbolString& slave) } -ActiveBusRequest::ActiveBusRequest(SymbolString& master, SymbolString& slave) - : BusRequest(master, false), m_finished(false), m_result(RESULT_SYN), m_slave(slave) -{ - pthread_mutex_init(&m_mutex, NULL); - pthread_cond_init(&m_cond, NULL); -} - -ActiveBusRequest::~ActiveBusRequest() -{ - pthread_mutex_destroy(&m_mutex); - pthread_cond_destroy(&m_cond); -} - -bool ActiveBusRequest::wait(int timeout) -{ - m_finished = false; - m_result = RESULT_SYN; - struct timespec t; - clock_gettime(CLOCK_REALTIME, &t); - t.tv_sec += timeout; - int result = 0; - - pthread_mutex_lock(&m_mutex); - - while (m_finished == false && result == 0) - result = pthread_cond_timedwait(&m_cond, &m_mutex, &t); - - if (result == 0 && m_finished == false) - result = 1; - - pthread_mutex_unlock(&m_mutex); - - return result == 0; -} - bool ActiveBusRequest::notify(result_t result, SymbolString& slave) { if (result == RESULT_OK) L.log(bus, event, "read res: %s", slave.getDataStr().c_str()); - pthread_mutex_lock(&m_mutex); - m_result = result; m_slave = SymbolString(slave, false, false); m_finished = true; - pthread_cond_signal(&m_cond); - pthread_mutex_unlock(&m_mutex); return false; } @@ -189,20 +151,24 @@ result_t BusHandler::sendAndWait(SymbolString& master, SymbolString& slave) ActiveBusRequest* request = new ActiveBusRequest(master, slave); for (int sendRetries=m_failedSendRetries+1; sendRetries>=0; sendRetries--) { - m_requests.add(request); - bool success = request->wait(1); // 1 second is still 3 times the theoretical worst-case request duration - if (success == false) - m_requests.remove(request); + m_nextRequests.add(request); + bool success = m_finishedRequests.waitRemove(request); result = success == true ? request->m_result : RESULT_ERR_TIMEOUT; if (result == RESULT_OK) break; - L.log(bus, error, "%s, %s", getResultCode(result), sendRetries>0 ? "retry send" : "give up"); + if (success == false || result == RESULT_ERR_NO_SIGNAL) { + L.log(bus, error, "%s, give up", getResultCode(result)); + break; + } + L.log(bus, error, "%s, %s", getResultCode(result), sendRetries>0 ? "retry send" : ""); + request->m_busLostRetries = 0; + request->m_finished = false; } - delete request; // TODO may be unsave while run() is using the request + delete request; return result; } @@ -230,20 +196,25 @@ result_t BusHandler::handleSymbol() long timeout = SYN_TIMEOUT; unsigned char sendSymbol = ESC; bool sending = false; + BusRequest* startRequest = NULL; // check if another symbol has to be sent and determine timeout for receive switch (m_state) { + case bs_noSignal: + timeout = SIGNAL_TIMEOUT; + break; + case bs_skip: - timeout = 0; // endless + timeout = SYN_TIMEOUT; break; case bs_ready: - if (m_request != NULL) + if (m_currentRequest != NULL) setState(bs_ready, RESULT_ERR_TIMEOUT); // just to be sure an old BusRequest is cleaned up - if (m_remainLockCount == 0) { - m_request = m_requests.next(false); - if (m_request == NULL && m_pollInterval > 0) { // check for poll/scan + if (m_remainLockCount == 0 && m_currentRequest == NULL) { + startRequest = m_nextRequests.next(false); + if (startRequest == NULL && m_pollInterval > 0) { // check for poll/scan time_t now; time(&now); if (m_lastPoll == 0 || difftime(now, m_lastPoll) > m_pollInterval) { @@ -257,14 +228,14 @@ result_t BusHandler::handleSymbol() delete request; } else { - m_request = request; - m_requests.add(request); + startRequest = request; + m_nextRequests.add(request); } } } } - if (m_request != NULL) { // initiate arbitration - sendSymbol = m_request->m_master[0]; + if (startRequest != NULL) { // initiate arbitration + sendSymbol = m_ownMasterAddress; sending = true; } } @@ -278,28 +249,28 @@ result_t BusHandler::handleSymbol() break; case bs_sendCmd: - if (m_request != NULL) { - sendSymbol = m_request->m_master[m_nextSendPos]; + if (m_currentRequest != NULL) { + sendSymbol = m_currentRequest->m_master[m_nextSendPos]; sending = true; } break; case bs_sendResAck: - if (m_request != NULL) { + if (m_currentRequest != NULL) { sendSymbol = m_responseCrcValid ? ACK : NAK; sending = true; } break; case bs_sendCmdAck: - if (m_request != NULL) { + if (m_currentRequest != NULL) { sendSymbol = m_commandCrcValid ? ACK : NAK; sending = true; } break; case bs_sendRes: - if (m_request != NULL) { + if (m_currentRequest != NULL) { sendSymbol = m_response[m_nextSendPos]; sending = true; } @@ -331,9 +302,16 @@ result_t BusHandler::handleSymbol() unsigned char recvSymbol; result = m_port->recv(timeout, recvSymbol); - if (result != RESULT_OK) - return setState(bs_skip, result); // TODO keep "no signal" within auto-syn state + time_t now; + time(&now); + if (result != RESULT_OK) { + if (difftime(now, m_lastReceive) > 1) // at least one full second has passed since last received symbol + return setState(bs_noSignal, result); + return setState(timeout == SIGNAL_TIMEOUT ? bs_noSignal : bs_skip, result); + } + + m_lastReceive = now; if (recvSymbol == SYN) { if (sending == false && m_remainLockCount > 0 && m_command.size() != 1) m_remainLockCount--; @@ -346,15 +324,19 @@ result_t BusHandler::handleSymbol() switch (m_state) { + case bs_noSignal: + return setState(bs_skip, RESULT_OK); + case bs_skip: return RESULT_OK; case bs_ready: - if (m_request != NULL && sending == true) { - if (m_requests.remove(m_request) == false) { - // request already timed out + if (startRequest != NULL && sending == true) { + if (m_nextRequests.remove(startRequest) == false) { + // request already removed (e.g. due to timeout) return setState(bs_skip, RESULT_ERR_TIMEOUT); } + m_currentRequest = startRequest; // check arbitration if (recvSymbol == sendSymbol) { // arbitration successful m_nextSendPos = 1; @@ -415,8 +397,8 @@ result_t BusHandler::handleSymbol() if (m_commandCrcValid == false) return setState(bs_skip, RESULT_ERR_ACK); - if (m_request != NULL) { - if (isMaster(m_request->m_master[1]) == true) { + if (m_currentRequest != NULL) { + if (isMaster(m_currentRequest->m_master[1]) == true) { return setState(bs_sendSyn, RESULT_OK); } } else if (isMaster(m_command[1]) == true) { @@ -432,17 +414,17 @@ result_t BusHandler::handleSymbol() m_repeat = true; m_nextSendPos = 0; m_command.clear(); - if (m_request != NULL) + if (m_currentRequest != NULL) return setState(bs_sendCmd, RESULT_ERR_NAK, true); return setState(bs_recvCmd, RESULT_ERR_NAK); } - if (m_request != NULL) + if (m_currentRequest != NULL) return setState(bs_skip, RESULT_ERR_NAK); return setState(bs_skip, RESULT_ERR_NAK); } - if (m_request != NULL) + if (m_currentRequest != NULL) return setState(bs_skip, RESULT_ERR_ACK); return setState(bs_skip, RESULT_ERR_ACK); @@ -452,7 +434,7 @@ result_t BusHandler::handleSymbol() crcPos = m_response.size() > headerLen ? headerLen + 1 + m_response[headerLen] : 0xff; result = m_response.push_back(recvSymbol, true, m_response.size() < crcPos); if (result < RESULT_OK) { - if (m_request != NULL) + if (m_currentRequest != NULL) return setState(bs_skip, result); return setState(bs_skip, result); @@ -460,18 +442,18 @@ result_t BusHandler::handleSymbol() if (result == RESULT_OK && crcPos != 0xff && m_response.size() == crcPos + 1) { // CRC received m_responseCrcValid = m_response[headerLen + 1 + m_response[headerLen]] == m_response.getCRC(); if (m_responseCrcValid) { - if (m_request != NULL) + if (m_currentRequest != NULL) return setState(bs_sendResAck, RESULT_OK); return setState(bs_recvResAck, RESULT_OK); } if (m_repeat == true) { - if (m_request != NULL) + if (m_currentRequest != NULL) return setState(bs_sendSyn, RESULT_ERR_CRC); return setState(bs_skip, RESULT_ERR_CRC); } - if (m_request != NULL) + if (m_currentRequest != NULL) return setState(bs_sendResAck, RESULT_ERR_CRC); return setState(bs_recvResAck, RESULT_ERR_CRC); @@ -497,13 +479,13 @@ result_t BusHandler::handleSymbol() return setState(bs_skip, RESULT_ERR_ACK); case bs_sendCmd: - if (m_request != NULL && sending == true) { + if (m_currentRequest != NULL && sending == true) { if (recvSymbol == sendSymbol) { // successfully sent m_nextSendPos++; - if (m_nextSendPos >= m_request->m_master.size()) { + if (m_nextSendPos >= m_currentRequest->m_master.size()) { // master data completely sent - if (m_request->m_master[1] == BROADCAST) + if (m_currentRequest->m_master[1] == BROADCAST) return setState(bs_sendSyn, RESULT_OK); m_commandCrcValid = true; @@ -515,7 +497,7 @@ result_t BusHandler::handleSymbol() return setState(bs_skip, RESULT_ERR_INVALID_ARG); case bs_sendResAck: - if (m_request != NULL && sending == true) { + if (m_currentRequest != NULL && sending == true) { if (recvSymbol == sendSymbol) { // successfully sent if (m_responseCrcValid == false) { @@ -592,26 +574,43 @@ result_t BusHandler::handleSymbol() result_t BusHandler::setState(BusState state, result_t result, bool firstRepetition) { - if (m_request != NULL) { - if (result == RESULT_ERR_BUS_LOST && m_request->m_busLostRetries < m_busLostRetries) { + if (m_currentRequest != NULL) { + if (result == RESULT_ERR_BUS_LOST && m_currentRequest->m_busLostRetries < m_busLostRetries) { L.log(bus, error, "%s, retry", getResultCode(result)); - m_request->m_busLostRetries++; - m_requests.add(m_request); // repeat - m_request = NULL; + m_currentRequest->m_busLostRetries++; + m_nextRequests.add(m_currentRequest); // repeat + m_currentRequest = NULL; } else if (state == bs_sendSyn || (result != RESULT_OK && firstRepetition == false)) { L.log(bus, debug, "notify request: %s", getResultCode(result)); - bool restart = m_request->notify(result, m_response); - unsigned char dstAddress = m_request->m_master[1]; + unsigned char dstAddress = m_currentRequest->m_master[1]; if (result == RESULT_OK && isValidAddress(dstAddress, false) == true) m_seenAddresses[dstAddress] = true; + bool restart = m_currentRequest->notify(result, m_response); if (restart == true) { - m_request->m_busLostRetries = 0; - m_requests.add(m_request); + m_currentRequest->m_busLostRetries = 0; + m_nextRequests.add(m_currentRequest); } - else if (m_request->m_deleteOnFinish == true) { - delete m_request; + else if (m_currentRequest->m_deleteOnFinish == true) + delete m_currentRequest; + else + m_finishedRequests.add(m_currentRequest); + + m_currentRequest = NULL; + } + } + + if (state == bs_noSignal) { // notify all requests + m_response.clear(); + while ((m_currentRequest = m_nextRequests.remove(false)) != NULL) { + bool restart = m_currentRequest->notify(RESULT_ERR_NO_SIGNAL, m_response); + if (restart == true) { // should not occur with no signal + m_currentRequest->m_busLostRetries = 0; + m_nextRequests.add(m_currentRequest); } - m_request = NULL; + else if (m_currentRequest->m_deleteOnFinish == true) + delete m_currentRequest; + else + m_finishedRequests.add(m_currentRequest); } } @@ -620,7 +619,7 @@ result_t BusHandler::setState(BusState state, result_t result, bool firstRepetit if (result < RESULT_OK || (result != RESULT_OK && state == bs_skip)) L.log(bus, debug, "%s during %s, switching to %s", getResultCode(result), getStateCode(m_state), getStateCode(state)); - else if (m_request != NULL || state == bs_sendCmd || state==bs_sendResAck || state==bs_sendSyn) + else if (m_currentRequest != NULL || state == bs_sendCmd || state==bs_sendResAck || state==bs_sendSyn) L.log(bus, debug, "switching from %s to %s", getStateCode(m_state), getStateCode(state)); m_state = state; @@ -704,7 +703,7 @@ result_t BusHandler::startScan(bool full) delete request; return result; } - m_requests.add(request); + m_nextRequests.add(request); } return RESULT_OK; } diff --git a/src/ebusd/bushandler.h b/src/ebusd/bushandler.h index 1535d044..47f70008 100644 --- a/src/ebusd/bushandler.h +++ b/src/ebusd/bushandler.h @@ -42,6 +42,9 @@ using namespace std; /** @brief the maximum allowed time [us] for retrieving the AUTO-SYN symbol (45ms + 2*1,2% + 1 Symbol). */ #define SYN_TIMEOUT 50800 +/** @brief the time [us] for determining bus signal availability (AUTO-SYN timeout * 5). */ +#define SIGNAL_TIMEOUT 250000 + /** @brief the maximum duration [us] of a single symbol (Start+8Bit+Stop+Extra @ 2400Bd-2*1,2%). */ #define SYMBOL_DURATION 4700 @@ -50,6 +53,7 @@ using namespace std; /** @brief the possible bus states. */ enum BusState { + bs_noSignal, //!< no signal on the bus bs_skip, //!< skip all symbols until next @a SYN bs_ready, //!< ready for next master (after @a SYN symbol, send/receive QQ) bs_recvCmd, //!< receive command (ZZ, PBSB, master data) [passive set] @@ -215,19 +219,13 @@ public: * @param master reference to the master data @a SymbolString to send. * @param slave reference to @a SymbolString for filling in the received slave data. */ - ActiveBusRequest(SymbolString& master, SymbolString& slave); + ActiveBusRequest(SymbolString& master, SymbolString& slave) + : BusRequest(master, false), m_finished(false), m_result(RESULT_SYN), m_slave(slave) {} /** * @brief Destructor. */ - virtual ~ActiveBusRequest(); - - /** - * @brief Wait for notification. - * @param timeout the maximum time to wait in seconds. - * @return the result code. - */ - bool wait(int timeout); + virtual ~ActiveBusRequest() {} // @copydoc virtual bool notify(result_t result, SymbolString& slave); @@ -243,12 +241,6 @@ private: /** reference to @a SymbolString for filling in the received slave data. */ SymbolString& m_slave; - /** a mutex for wait/notify. */ - pthread_mutex_t m_mutex; - - /** a mutex condition for wait/notify. */ - pthread_cond_t m_cond; - }; @@ -282,9 +274,9 @@ public: m_busLostRetries(busLostRetries), m_failedSendRetries(failedSendRetries), m_busAcquireTimeout(busAcquireTimeout), m_slaveRecvTimeout(slaveRecvTimeout), m_lockCount(lockCount), m_remainLockCount(lockCount), - m_pollInterval(pollInterval), m_lastPoll(0), - m_request(NULL), m_nextSendPos(0), - m_state(bs_skip), m_repeat(false), + m_pollInterval(pollInterval), m_lastReceive(0), m_lastPoll(0), + m_currentRequest(NULL), m_nextSendPos(0), + m_state(bs_noSignal), m_repeat(false), m_commandCrcValid(false), m_responseCrcValid(false), m_scanMessage(NULL) { memset(m_seenAddresses, 0, sizeof(m_seenAddresses)); @@ -389,14 +381,20 @@ private: /** the interval in seconds in which poll messages are cycled, or 0 if disabled. */ const unsigned int m_pollInterval; + /** the time of the last received symbol, or 0 for never. */ + time_t m_lastReceive; + /** the time of the last poll, or 0 for never. */ time_t m_lastPoll; /** the queue of @a BusRequests that shall be handled. */ - WQueue m_requests; + WQueue m_nextRequests; /** the currently handled BusRequest, or NULL. */ - BusRequest* m_request; + BusRequest* m_currentRequest; + + /** the queue of @a BusRequests that are already finished. */ + WQueue m_finishedRequests; /** the offset of the next symbol that needs to be sent from the command or response, * (only relevant if m_request is set and state is bs_command or bs_response). */