diff --git a/src/lib/ebus/device.cpp b/src/lib/ebus/device.cpp index 41b14410..ce324c75 100755 --- a/src/lib/ebus/device.cpp +++ b/src/lib/ebus/device.cpp @@ -57,9 +57,30 @@ namespace ebusd { #define ENH_STARTED ((symbol_t)0x02) #define ENH_FAILED ((symbol_t)0x82) +/** + * Construct a new instance. + * @param name the device name (e.g. "/dev/ttyUSB0" for serial, "127.0.0.1:1234" for network). + * @param address the socket address of the device. + * @param readOnly whether to allow read access to the device only. + * @param initialSend whether to send an initial @a ESC symbol in @a open(). + * @param udp true for UDP, false to TCP. + * @param enhancedProto whether the device supports the ebusd enhanced protocol. + */ +Device::Device(const char* name, bool checkDevice, bool readOnly, bool initialSend, bool enhancedProto) + : m_name(name), m_checkDevice(checkDevice), m_readOnly(readOnly), m_initialSend(initialSend), + m_enhancedProto(enhancedProto), m_fd(-1), m_listener(nullptr), m_arbitrationMaster(SYN), + m_arbitrationCheck(false), m_bufSize(((MAX_LEN+1+3)/4)*4), m_bufLen(0), m_bufPos(0) { + m_buffer = reinterpret_cast(malloc(m_bufSize)); + if (!m_buffer) { + m_bufSize = 0; + } +} Device::~Device() { close(); + if (m_buffer) { + free(m_buffer); + } } Device* Device::create(const char* name, bool checkDevice, bool readOnly, bool initialSend) { @@ -72,6 +93,10 @@ Device* Device::create(const char* name, bool checkDevice, bool readOnly, bool i if (portpos >= addrpos+3 && strncmp(addrpos, "enh", 3) == 0) { enhanced = true; addrpos += 3; + if (portpos == addrpos) { + addrpos++; + portpos = strchr(addrpos, ':'); + } } if (portpos == addrpos+3 && (strncmp(addrpos, "tcp", 3) == 0 || (udp=(strncmp(addrpos, "udp", 3) == 0)))) { addrpos += 4; @@ -106,11 +131,17 @@ Device* Device::create(const char* name, bool checkDevice, bool readOnly, bool i return new SerialDevice(name, checkDevice, readOnly, initialSend); } +result_t Device::open() { + close(); + return m_bufSize == 0 ? RESULT_ERR_DEVICE : RESULT_OK; +} + void Device::close() { if (m_fd != -1) { ::close(m_fd); m_fd = -1; } + m_bufLen = 0; // flush read buffer } bool Device::isValid() { @@ -140,7 +171,8 @@ result_t Device::recv(unsigned int timeout, symbol_t* value, ArbitrationState* a if (!isValid()) { return RESULT_ERR_DEVICE; } - if (!available() && timeout > 0) { + bool isAvailable = available(); + if (!isAvailable && timeout > 0) { int ret; struct timespec tdiff; @@ -185,18 +217,18 @@ result_t Device::recv(unsigned int timeout, symbol_t* value, ArbitrationState* a } } + ArbitrationState prevState = *arbitrationState; // directly read byte from device - if (!read(value, m_enhancedProto ? arbitrationState : nullptr)) { + if (!read(value, isAvailable, arbitrationState)) { close(); return RESULT_ERR_DEVICE; } - ArbitrationState prevState = m_enhancedProto && arbitrationState ? *arbitrationState : as_none; - if (*value != SYN || m_arbitrationMaster == SYN || m_enhancedProto) { + if (m_enhancedProto || *value != SYN || m_arbitrationMaster == SYN) { if (m_listener != nullptr) { m_listener->notifyDeviceData(*value, true); } if (m_enhancedProto) { - if (arbitrationState && *arbitrationState != prevState) { + if (*arbitrationState != prevState) { m_arbitrationMaster = SYN; m_arbitrationCheck = false; } @@ -236,32 +268,116 @@ result_t Device::startArbitration(symbol_t masterAddress) { if (m_readOnly) { return RESULT_ERR_SEND; } - m_arbitrationCheck = false; m_arbitrationMaster = masterAddress; - if (m_enhancedProto) { + m_arbitrationCheck = false; + if (m_enhancedProto && masterAddress != SYN) { if (!write(masterAddress, true)) { m_arbitrationMaster = SYN; - m_arbitrationCheck = false; return RESULT_ERR_SEND; } m_arbitrationCheck = true; - return RESULT_OK; } return RESULT_OK; } bool Device::write(symbol_t value, bool startArbitration) { + if (m_enhancedProto) { + symbol_t buf[2] = {startArbitration ? ENH_START : ENH_SEND, value}; + return ::write(m_fd, buf, 2) == 2; + } return ::write(m_fd, &value, 1) == 1; } -bool Device::read(symbol_t* value, ArbitrationState* arbitrationState) { - return ::read(m_fd, value, 1) == 1; +bool Device::available() { + if (m_bufLen <= 0) { + return false; + } + if (!m_enhancedProto) { + return true; + } + // peek into the received enhanced proto bytes to determine symbol availability + for (size_t pos = 0; pos < m_bufLen; pos++) { + symbol_t ch = m_buffer[(pos+m_bufPos)%m_bufSize]; + if (ch == ENH_RECEIVED) { + return pos+1 < m_bufLen; + } + } + return false; +} + +bool Device::read(symbol_t* value, bool isAvailable, ArbitrationState* arbitrationState) { + if (!isAvailable) { + if (m_bufLen > 0 && m_bufPos != 0) { + if (m_bufLen > m_bufSize / 2) { + m_bufLen = 0; // TODO report error + } else { + size_t tail; + if (m_bufPos+m_bufLen > m_bufSize) { + // move wrapped tail away + tail = (m_bufPos+m_bufLen) % m_bufSize; + size_t head = m_bufLen-tail; + memmove(m_buffer+head, m_buffer, tail); + } else { + tail = 0; + } + // move head to first position + memmove(m_buffer, m_buffer + m_bufPos, m_bufLen - tail); + } + } + m_bufPos = 0; + // fill up the buffer + ssize_t size = ::read(m_fd, m_buffer + m_bufLen, m_bufSize - m_bufLen); + if (size <= 0) { + return false; + } + m_bufLen += size; + } + if (!available()) { + return false; + } + if (!m_enhancedProto) { + *value = m_buffer[m_bufPos]; + m_bufPos = (m_bufPos+1)%m_bufSize; + m_bufLen--; + return true; + } + while (m_bufLen > 0) { + symbol_t ch = m_buffer[m_bufPos]; + m_bufPos = (m_bufPos+1)%m_bufSize; + m_bufLen--; + switch (ch) { + case ENH_STARTED: + *arbitrationState = as_won; + m_arbitrationMaster = SYN; + break; + case ENH_FAILED: + *arbitrationState = as_error; + m_arbitrationMaster = SYN; + break; + case ENH_RECEIVED: + if (m_bufLen <= 0) { + return false; + } + *value = m_buffer[m_bufPos]; + m_bufPos = (m_bufPos+1)%m_bufSize; + m_bufLen--; + return true; + case ENH_RESETTED: // TODO + *arbitrationState = as_error; + break; + default: + // TODO proto error + return false; + } + } + return false; } result_t SerialDevice::open() { - if (m_fd != -1) { - close(); + result_t result = Device::open(); + if (result != RESULT_OK) { + return result; } struct termios newSettings; @@ -339,8 +455,9 @@ void SerialDevice::checkDevice() { result_t NetworkDevice::open() { - if (m_fd != -1) { - close(); + result_t result = Device::open(); + if (result != RESULT_OK) { + return result; } m_fd = socket(AF_INET, m_udp ? SOCK_DGRAM : SOCK_STREAM, 0); if (m_fd < 0) { @@ -376,13 +493,6 @@ result_t NetworkDevice::open() { break; } } - if (m_bufSize == 0) { - m_bufSize = MAX_LEN+1; - m_buffer = reinterpret_cast(malloc(m_bufSize)); - if (!m_buffer) { - m_bufSize = 0; - } - } m_bufLen = 0; if (m_initialSend && !write(ESC)) { return RESULT_ERR_SEND; @@ -390,11 +500,6 @@ result_t NetworkDevice::open() { return RESULT_OK; } -void NetworkDevice::close() { - m_bufLen = 0; // flush read buffer - Device::close(); -} - void NetworkDevice::checkDevice() { int cnt; if (ioctl(m_fd, FIONREAD, &cnt) < 0) { @@ -402,74 +507,4 @@ void NetworkDevice::checkDevice() { } } -bool NetworkDevice::available() { - return m_buffer && m_bufLen > 0; -} - -bool NetworkDevice::write(symbol_t value, bool startArbitration) { - m_bufLen = 0; // flush read buffer - if (m_bufSize > 0 && m_enhancedProto) { - m_buffer[0] = startArbitration ? ENH_START : ENH_SEND; - m_buffer[1] = value; - return ::write(m_fd, m_buffer, 2) == 2; - } - return Device::write(value); -} - -bool NetworkDevice::read(symbol_t* value, ArbitrationState* arbitrationState) { - if (available()) { - *value = m_buffer[m_bufPos]; - m_bufPos = (m_bufPos+1)%m_bufSize; - m_bufLen--; - return true; - } - if (m_bufSize > 0) { - ssize_t size = ::read(m_fd, m_buffer, m_bufSize); - if (size <= 0) { - return false; - } - if (m_enhancedProto) { - symbol_t* buf = m_buffer; - m_bufPos = 0; - m_bufLen = 0; - while (size > 0) { - buf++; - size--; - switch (*buf) { - case ENH_STARTED: - if (arbitrationState) { - *arbitrationState = as_won; - } - break; - case ENH_FAILED: - if (arbitrationState) { - *arbitrationState = as_error; - } - break; - case ENH_RECEIVED: - m_buffer[m_bufPos++] = *buf; - m_bufLen++; - break; - case ENH_RESETTED: // TODO - break; - default: - return false; - } - } - if (m_bufLen > 0) { - *value = m_buffer[0]; - m_bufPos = 1; - m_bufLen--; - return true; - } - return false; - } - *value = m_buffer[0]; - m_bufPos = 1; - m_bufLen = size-1; - return true; - } - return Device::read(value, arbitrationState); -} - } // namespace ebusd diff --git a/src/lib/ebus/device.h b/src/lib/ebus/device.h index 92768e33..84576291 100755 --- a/src/lib/ebus/device.h +++ b/src/lib/ebus/device.h @@ -81,10 +81,7 @@ class Device { * @param initialSend whether to send an initial @a ESC symbol in @a open(). * @param enhancedProto whether the device supports the ebusd enhanced protocol. */ - Device(const char* name, bool checkDevice, bool readOnly, bool initialSend, bool enhancedProto=false) - : m_name(name), m_checkDevice(checkDevice), m_readOnly(readOnly), m_initialSend(initialSend), - m_enhancedProto(enhancedProto), m_fd(-1), m_listener(nullptr), m_arbitrationMaster(SYN), - m_arbitrationCheck(false) {} + Device(const char* name, bool checkDevice, bool readOnly, bool initialSend, bool enhancedProto=false); /** * Destructor. @@ -113,7 +110,7 @@ class Device { * Open the file descriptor. * @return the @a result_t code. */ - virtual result_t open() = 0; // abstract + virtual result_t open(); /** * Close the file descriptor if opened. @@ -178,12 +175,6 @@ class Device { */ virtual void checkDevice() = 0; // abstract - /** - * Check whether a byte is available immediately (without waiting). - * @return true when a a byte is available immediately. - */ - virtual bool available() { return false; } - /** * Write a single byte. * @param value the byte value to write. @@ -191,12 +182,20 @@ class Device { */ virtual bool write(symbol_t value, bool startArbitration=false); + /** + * Check whether a symbol is available for reading immediately (without waiting). + * @return true when a symbol is available for reading immediately. + */ + virtual bool available(); + /** * Read a single byte. * @param value the reference in which the read byte value is stored. + * @param isAvailable the result of the immediately preceding call to @a available(). + * @param arbitrationState the variable in which to store the received arbitration state (mandatory for enhanced proto). * @return true on success, false on error. */ - virtual bool read(symbol_t* value, ArbitrationState* arbitrationState=nullptr); + virtual bool read(symbol_t* value, bool isAvailable, ArbitrationState* arbitrationState=nullptr); /** the device name (e.g. "/dev/ttyUSB0" for serial, "127.0.0.1:1234" for network). */ const char* m_name; @@ -219,15 +218,28 @@ class Device { private: /** the @a DeviceListener, or nullptr. */ - DeviceListener* m_listener; +public: 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; + + /** the read buffer. */ + symbol_t* m_buffer; + + /** the read buffer size (multiple of 4). */ + size_t m_bufSize; + + /** the read buffer fill length. */ + size_t m_bufLen; + + /** the read buffer read position. */ + size_t m_bufPos; }; + /** * The @a Device for directly connected serial interfaces (tty). */ @@ -276,17 +288,12 @@ class NetworkDevice : public Device { */ NetworkDevice(const char* name, const struct sockaddr_in& address, bool readOnly, bool initialSend, bool udp, bool enhancedProto=false) - : Device(name, true, readOnly, initialSend, enhancedProto), m_address(address), m_udp(udp), - m_buffer(nullptr), m_bufSize(0), m_bufLen(0), m_bufPos(0) {} + : Device(name, true, readOnly, initialSend, enhancedProto), m_address(address), m_udp(udp) {} /** * Destructor. */ - virtual ~NetworkDevice() { - if (m_buffer) { - free(m_buffer); - } - } + ~NetworkDevice() override {} // @copydoc unsigned int getLatency() const override { return 10000; } @@ -294,22 +301,11 @@ class NetworkDevice : public Device { // @copydoc result_t open() override; - // @copydoc - void close() override; protected: // @copydoc void checkDevice() override; - // @copydoc - bool available() override; - - // @copydoc - bool write(symbol_t value, bool startArbitration=false) override; - - // @copydoc - bool read(symbol_t* value, ArbitrationState* arbitrationState=nullptr) override; - private: /** the socket address of the device. */ @@ -317,18 +313,6 @@ class NetworkDevice : public Device { /** true for UDP, false to TCP. */ const bool m_udp; - - /** the buffer memory, or nullptr. */ - symbol_t* m_buffer; - - /** the buffer size. */ - size_t m_bufSize; - - /** the buffer fill length. */ - size_t m_bufLen; - - /** the buffer read position. */ - size_t m_bufPos; }; } // namespace ebusd