diff --git a/src/ebusd/bushandler.cpp b/src/ebusd/bushandler.cpp index 45c664f2..565c2ead 100755 --- a/src/ebusd/bushandler.cpp +++ b/src/ebusd/bushandler.cpp @@ -276,7 +276,7 @@ bool GrabbedMessage::dump(bool unknown, MessageMap* messages, bool first, bool d if (remain == 0) { return true; } - for (const auto it : *types) { + for (const auto& it : *types) { const DataType* baseType = it.second; if ((baseType->getBitCount() % 8) != 0 || baseType->isIgnored()) { // skip bit and ignored types continue; @@ -422,7 +422,6 @@ result_t BusHandler::handleSymbol() { unsigned int timeout = SYN_TIMEOUT; symbol_t sendSymbol = ESC; bool sending = false; - BusRequest* startRequest = nullptr; // check if another symbol has to be sent and determine timeout for receive switch (m_state) { @@ -432,13 +431,8 @@ result_t BusHandler::handleSymbol() { case bs_skip: timeout = SYN_TIMEOUT; - break; - - case bs_ready: - if (m_currentRequest != nullptr) { - setState(bs_ready, RESULT_ERR_TIMEOUT); // just to be sure an old BusRequest is cleaned up - } else if (m_remainLockCount == 0) { - startRequest = m_nextRequests.peek(); + if (!m_device->isArbitrating() && m_currentRequest == nullptr && m_remainLockCount == 0) { + BusRequest* startRequest = m_nextRequests.peek(); if (startRequest == nullptr && m_pollInterval > 0) { // check for poll/scan time_t now; time(&now); @@ -446,7 +440,7 @@ result_t BusHandler::handleSymbol() { Message* message = m_messages->getNextPoll(); if (message != nullptr) { m_lastPoll = now; - PollRequest* request = new PollRequest(message); + auto request = new PollRequest(message); result_t ret = request->prepare(m_ownMasterAddress); if (ret != RESULT_OK) { logError(lf_bus, "prepare poll message: %s", getResultCode(ret)); @@ -459,12 +453,23 @@ result_t BusHandler::handleSymbol() { } } if (startRequest != nullptr) { // initiate arbitration - sendSymbol = startRequest->m_master[0]; - sending = true; + result_t ret = m_device->startArbitration(startRequest->m_master[0]); + if (ret != RESULT_OK) { + logError(lf_bus, "arbitration start: %s", getResultCode(ret)); + m_currentRequest = startRequest; + setState(bs_ready, ret); // force the failed request to be notified + startRequest = nullptr; + } } } break; + case bs_ready: + if (m_currentRequest != nullptr) { + setState(bs_ready, RESULT_ERR_TIMEOUT); // just to be sure an old BusRequest is cleaned up + } + break; + case bs_recvCmd: case bs_recvCmdCrc: timeout = m_slaveRecvTimeout; @@ -537,11 +542,11 @@ result_t BusHandler::handleSymbol() { // send symbol if necessary result_t result; - struct timespec sentTime, recvTime; + struct timespec sentTime = {}, recvTime = {}; if (sending) { if (m_state != bs_sendSyn && (sendSymbol == ESC || sendSymbol == SYN)) { if (m_escape) { - sendSymbol = sendSymbol == ESC ? 0x00 : 0x01; + sendSymbol = (symbol_t)(sendSymbol == ESC ? 0x00 : 0x01); } else { m_escape = sendSymbol; sendSymbol = ESC; @@ -558,55 +563,76 @@ result_t BusHandler::handleSymbol() { } else { sending = false; timeout = SYN_TIMEOUT; - if (startRequest != nullptr && m_nextRequests.remove(startRequest)) { - m_currentRequest = startRequest; // force the failed request to be notified - } setState(bs_skip, result); } } // receive next symbol (optionally check reception of sent symbol) symbol_t recvSymbol; - result = m_device->recv(timeout+m_transferLatency, &recvSymbol); + ArbitrationState arbitrationState = as_none; + result = m_device->recv(timeout+m_transferLatency, &recvSymbol, &arbitrationState); if (sending) { clockGettime(&recvTime); } + bool sentAutoSyn = false; if (!sending && result == RESULT_ERR_TIMEOUT && m_generateSynInterval > 0 - && timeout >= m_generateSynInterval && (m_state == bs_noSignal || m_state == bs_skip)) { + && timeout >= m_generateSynInterval && (m_state == bs_noSignal || m_state == bs_skip)) { // check if acting as AUTO-SYN generator is required result = m_device->send(SYN); - if (result == RESULT_OK) { - clockGettime(&sentTime); - recvSymbol = ESC; - result = m_device->recv(SEND_TIMEOUT+m_transferLatency, &recvSymbol); - clockGettime(&recvTime); - if (result == RESULT_ERR_TIMEOUT) { - return setState(bs_noSignal, result); - } - if (result != RESULT_OK) { - logError(lf_bus, "unable to receive sent AUTO-SYN symbol: %s", getResultCode(result)); - } else if (recvSymbol != SYN) { - logError(lf_bus, "received %2.2x instead of AUTO-SYN symbol", recvSymbol); - } else { - measureLatency(&sentTime, &recvTime); - if (m_generateSynInterval != SYN_TIMEOUT) { - // received own AUTO-SYN symbol back again: act as AUTO-SYN generator now - m_generateSynInterval = SYN_TIMEOUT; - logNotice(lf_bus, "acting as AUTO-SYN generator"); - } - m_remainLockCount = 0; - m_lastSynReceiveTime = recvTime; - return setState(bs_ready, result); + if (result != RESULT_OK) { + return setState(bs_skip, result); + } + clockGettime(&sentTime); + recvSymbol = ESC; + result = m_device->recv(SEND_TIMEOUT+m_transferLatency, &recvSymbol, &arbitrationState); + clockGettime(&recvTime); + if (result != RESULT_OK) { + logError(lf_bus, "unable to receive sent AUTO-SYN symbol: %s", getResultCode(result)); + return setState(bs_noSignal, result); + } + if (recvSymbol != SYN) { + logError(lf_bus, "received %2.2x instead of AUTO-SYN symbol", recvSymbol); + return setState(bs_noSignal, result); + } + measureLatency(&sentTime, &recvTime); + if (m_generateSynInterval != SYN_TIMEOUT) { + // received own AUTO-SYN symbol back again: act as AUTO-SYN generator now + m_generateSynInterval = SYN_TIMEOUT; + logNotice(lf_bus, "acting as AUTO-SYN generator"); + } + m_remainLockCount = 0; + m_lastSynReceiveTime = recvTime; + sentAutoSyn = true; + } + if (arbitrationState == as_lost) { + if (m_currentRequest == nullptr) { + BusRequest *startRequest = m_nextRequests.peek(); + if (startRequest != nullptr && m_nextRequests.remove(startRequest)) { + m_currentRequest = startRequest; // force the failed request to be notified } } - return setState(bs_skip, result); + setState(m_state, RESULT_ERR_BUS_LOST); + } else if (arbitrationState == as_won) { // implies RESULT_OK + if (m_currentRequest == nullptr) { + m_currentRequest = m_nextRequests.peek(); + } + if (m_currentRequest == nullptr) { + logDebug(lf_bus, "arbitration won without request"); + } else if (m_state == bs_ready) { + if (!m_nextRequests.remove(m_currentRequest)) { + // request already removed (e.g. due to timeout) + return setState(bs_skip, RESULT_ERR_TIMEOUT); + } + sendSymbol = m_currentRequest->m_master[0]; + sending = true; + } + } + if (sentAutoSyn) { + return setState(bs_ready, RESULT_OK); } time_t now; time(&now); if (result != RESULT_OK) { - if (sending && startRequest != nullptr && m_nextRequests.remove(startRequest)) { - m_currentRequest = startRequest; // force the failed request to be notified - } if ((m_generateSynInterval != SYN_TIMEOUT && difftime(now, m_lastReceive) > 1) // at least one full second has passed since last received symbol || m_state == bs_noSignal) { @@ -672,19 +698,14 @@ result_t BusHandler::handleSymbol() { return RESULT_OK; case bs_ready: - if (startRequest != nullptr && sending) { - if (!m_nextRequests.remove(startRequest)) { - // request already removed (e.g. due to timeout) - return setState(bs_skip, RESULT_ERR_TIMEOUT); - } - m_currentRequest = startRequest; + if (m_currentRequest != nullptr && sending) { // check arbitration if (recvSymbol == sendSymbol) { // arbitration successful // measure arbitration delay long long latencyLong = (sentTime.tv_sec*1000000000 + sentTime.tv_nsec - m_lastSynReceiveTime.tv_sec*1000000000 - m_lastSynReceiveTime.tv_nsec)/1000; if (latencyLong >= 0 && latencyLong <= 10000) { // skip clock skew or out of reasonable range - int latency = static_cast(latencyLong); + auto latency = static_cast(latencyLong); logDebug(lf_bus, "arbitration delay %d micros", latency); if (m_arbitrationDelayMin < 0 || (latency < m_arbitrationDelayMin || latency > m_arbitrationDelayMax)) { if (m_arbitrationDelayMin == -1 || latency < m_arbitrationDelayMin) { @@ -1012,7 +1033,7 @@ void BusHandler::measureLatency(struct timespec* sentTime, struct timespec* recv if (latencyLong < 0 || latencyLong > 1000) { return; // clock skew or out of reasonable range } - int latency = static_cast(latencyLong); + auto latency = static_cast(latencyLong); logDebug(lf_bus, "send/receive symbol latency %d ms", latency); if (m_symbolLatencyMin >= 0 && (latency >= m_symbolLatencyMin && latency <= m_symbolLatencyMax)) { return; @@ -1294,7 +1315,7 @@ bool BusHandler::formatScanResult(symbol_t slave, bool leadingNewline, ostringst *output << endl; } *output << hex << setw(2) << setfill('0') << static_cast(slave); - for (const auto result : it->second) { + for (const auto &result : it->second) { *output << result; } return true; @@ -1423,7 +1444,7 @@ void BusHandler::formatUpdateInfo(ostringstream* output) const { const auto it = m_scanResults.find(address); if (it != m_scanResults.end()) { *output << ",\"s\":\""; - for (const auto result : it->second) { + for (const auto& result : it->second) { *output << result; } *output << "\""; @@ -1439,7 +1460,7 @@ void BusHandler::formatUpdateInfo(ostringstream* output) const { if (!loadedFiles.empty()) { *output << ",\"f\":["; bool first = true; - for (const auto loadedFile : loadedFiles) { + for (const auto& loadedFile : loadedFiles) { if (first) { first = false; } else { diff --git a/src/lib/ebus/device.cpp b/src/lib/ebus/device.cpp index b9f541d9..825b1960 100755 --- a/src/lib/ebus/device.cpp +++ b/src/lib/ebus/device.cpp @@ -120,7 +120,7 @@ result_t Device::send(symbol_t value) { return RESULT_OK; } -result_t Device::recv(unsigned int timeout, symbol_t* value) { +result_t Device::recv(unsigned int timeout, symbol_t* value, ArbitrationState* arbitrationState) { if (!isValid()) { return RESULT_ERR_DEVICE; } @@ -178,9 +178,36 @@ result_t Device::recv(unsigned int timeout, symbol_t* value) { close(); return RESULT_ERR_DEVICE; } + if (*value != SYN || m_arbitrationMaster == SYN) { + if (m_listener != nullptr) { + m_listener->notifyDeviceData(*value, true); + } + if (m_arbitrationMaster != SYN) { + if (m_arbitrationCheck) { + *arbitrationState = *value == m_arbitrationMaster ? as_won : as_lost; + m_arbitrationMaster = SYN; + m_arbitrationCheck = false; + } else { + *arbitrationState = m_arbitrationMaster == SYN ? as_none : as_start; + } + } + return RESULT_OK; + } + ssize_t wcnt = write(m_arbitrationMaster); // send as fast as possible if (m_listener != nullptr) { m_listener->notifyDeviceData(*value, true); } + if (wcnt != 1) { + *arbitrationState = as_error; + m_arbitrationMaster = SYN; + m_arbitrationCheck = false; + return RESULT_OK; + } + if (m_listener != NULL) { + m_listener->notifyDeviceData(m_arbitrationMaster, false); + } + m_arbitrationCheck = true; + *arbitrationState = as_running; return RESULT_OK; } diff --git a/src/lib/ebus/device.h b/src/lib/ebus/device.h index 7c555f8f..c8b5ffbc 100755 --- a/src/lib/ebus/device.h +++ b/src/lib/ebus/device.h @@ -39,6 +39,16 @@ namespace ebusd { * to a file and/or forwarding it to a logging function. */ +/** the arbitration state handled by @a Device. */ +enum ArbitrationState { + as_none, //!< no arbitration in process + as_start, //!< arbitration start requested + as_error, //!< error while sending master address + as_running, //!< arbitration currently running (master address sent, waiting for reception) + as_lost, //!< arbitration lost + as_won, //!< arbitration won +}; + /** * Interface for listening to data received on/sent to a device. */ @@ -72,7 +82,7 @@ class Device { */ Device(const char* name, bool checkDevice, bool readOnly, bool initialSend) : m_name(name), m_checkDevice(checkDevice), m_readOnly(readOnly), m_initialSend(initialSend), m_fd(-1), - m_listener(nullptr) {} + m_listener(nullptr), m_arbitrationMaster(SYN), m_arbitrationCheck(false) {} /** * Destructor. @@ -119,9 +129,31 @@ class Device { * Read a single byte from the device. * @param timeout maximum time to wait for the byte in microseconds, or 0 for infinite. * @param value the reference in which the received byte value is stored. + * @param arbitrationState the reference in which the current @a ArbitrationState is stored on success. When set to + * @a as_won, the received byte is the master address that was successfully arbitrated with. * @return the result_t code. */ - result_t recv(unsigned int timeout, symbol_t* value); + result_t recv(unsigned int timeout, symbol_t* value, ArbitrationState* arbitrationState); + + /** + * Start the arbitration with the specified master address. A subsequent request while an arbitration is currently in + * checking state will always result in @a RESULT_ERR_DUPLICATE. + * @param masterAddress the master address, or @a SYN to cancel a previous arbitration request. + * @return the result_t code. + */ + result_t startArbitration(symbol_t masterAddress) { + if (m_arbitrationCheck) { + return RESULT_ERR_DUPLICATE; + } + if (m_readOnly) { + return RESULT_ERR_SEND; + } + m_arbitrationCheck = false; + m_arbitrationMaster = masterAddress; + return RESULT_OK; + } + + bool isArbitrating() const { return m_arbitrationMaster != SYN; }; /** * Return the device name. @@ -193,6 +225,12 @@ class Device { private: /** the @a DeviceListener, or nullptr. */ DeviceListener* m_listener; + + /** the arbitration master address to send when in arbitration, or @a SYN. */ + symbol_t m_arbitrationMaster; + + /** true when in arbitration and the next received symbol needs to be checked against the sent master address. */ + bool m_arbitrationCheck; }; /**