solved TODO and fixed issue with concurrent access to BusRequest, added support for determining signal timeout

This commit is contained in:
john30
2014-12-27 13:37:42 +01:00
parent ff7caba5ca
commit 121ad369e8
2 changed files with 108 additions and 111 deletions
+90 -91
View File
@@ -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;
}
+18 -20
View File
@@ -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<BusRequest*> m_requests;
WQueue<BusRequest*> m_nextRequests;
/** the currently handled BusRequest, or NULL. */
BusRequest* m_request;
BusRequest* m_currentRequest;
/** the queue of @a BusRequests that are already finished. */
WQueue<BusRequest*> 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). */