diff --git a/examples/companion_radio/MyMesh.cpp b/examples/companion_radio/MyMesh.cpp index ee8114ca96..069509822d 100644 --- a/examples/companion_radio/MyMesh.cpp +++ b/examples/companion_radio/MyMesh.cpp @@ -246,18 +246,22 @@ void MyMesh::addToOfflineQueue(const uint8_t frame[], int len) { } } -int MyMesh::getFromOfflineQueue(uint8_t frame[]) { +int MyMesh::peekOfflineQueue(uint8_t frame[]) { if (offline_queue_len > 0) { // check offline queue - size_t len = offline_queue[0].len; // take from top of queue + size_t len = offline_queue[0].len; // copy from top of queue, but leave it there memcpy(frame, offline_queue[0].buf, len); + return len; + } + return 0; // queue is empty +} +void MyMesh::popOfflineQueue() { + if (offline_queue_len > 0) { offline_queue_len--; for (int i = 0; i < offline_queue_len; i++) { // delete top item from queue offline_queue[i] = offline_queue[i + 1]; } - return len; } - return 0; // queue is empty } float MyMesh::getAirtimeBudgetFactor() const { @@ -1449,11 +1453,18 @@ void MyMesh::handleCmdFrame(size_t len) { } } else if (cmd_frame[0] == CMD_SYNC_NEXT_MESSAGE) { int out_len; - if ((out_len = getFromOfflineQueue(out_frame)) > 0) { - _serial->writeFrame(out_frame, out_len); + if ((out_len = peekOfflineQueue(out_frame)) > 0) { + // The frame stays queued until a transport has taken it. A write that + // puts nothing on the wire (BLE/WiFi send queue full, serial link not + // accepting) used to lose the message for good; now the client simply + // asks again. writeFrame() returns 0 only in that case, see + // BaseSerialInterface, so any other result consumes the frame. + if (_serial->writeFrame(out_frame, out_len) != 0) { + popOfflineQueue(); #ifdef DISPLAY_CLASS - if (_ui) _ui->msgRead(offline_queue_len); + if (_ui) _ui->msgRead(offline_queue_len); #endif + } } else { out_frame[0] = RESP_CODE_NO_MORE_MESSAGES; _serial->writeFrame(out_frame, 1); diff --git a/examples/companion_radio/MyMesh.h b/examples/companion_radio/MyMesh.h index 02b00ad259..3fd0a365e3 100644 --- a/examples/companion_radio/MyMesh.h +++ b/examples/companion_radio/MyMesh.h @@ -210,7 +210,8 @@ class MyMesh : public BaseChatMesh, public DataStoreHost { void writeContactRespFrame(uint8_t code, const ContactInfo &contact); void updateContactFromFrame(ContactInfo &contact, uint32_t& last_mod, const uint8_t *frame, int len); void addToOfflineQueue(const uint8_t frame[], int len); - int getFromOfflineQueue(uint8_t frame[]); + int peekOfflineQueue(uint8_t frame[]); // copies the oldest frame, leaves it queued + void popOfflineQueue(); // discards the oldest frame int getBlobByKey(const uint8_t key[], int key_len, uint8_t dest_buf[]) override { return _store->getBlobByKey(key, key_len, dest_buf); } diff --git a/platformio.ini b/platformio.ini index 2219c97862..f0d4b64ea5 100644 --- a/platformio.ini +++ b/platformio.ini @@ -172,6 +172,7 @@ build_src_filter = +<../src/Packet.cpp> +<../src/helpers/ConfigSerializer.cpp> +<../src/helpers/DynamicConfigSerializer.cpp> + +<../src/helpers/ArduinoSerialInterface.cpp> lib_deps = google/googletest @ 1.17.0 diff --git a/src/helpers/ArduinoSerialInterface.cpp b/src/helpers/ArduinoSerialInterface.cpp index a01fa5866f..16804112c6 100644 --- a/src/helpers/ArduinoSerialInterface.cpp +++ b/src/helpers/ArduinoSerialInterface.cpp @@ -32,8 +32,13 @@ size_t ArduinoSerialInterface::writeFrame(const uint8_t src[], size_t len) { hdr[1] = (len & 0xFF); // LSB hdr[2] = (len >> 8); // MSB - _serial->write(hdr, 3); - return _serial->write(src, len); + // See BaseSerialInterface::writeFrame(): 0 only when nothing went out, so + // the caller may retry; a torn frame counts as taken because a retry would + // only hand the receiver another header as payload. + size_t n = _serial->write(hdr, 3); + if (n == 0) return 0; + if (n == 3) _serial->write(src, len); + return len; } size_t ArduinoSerialInterface::checkRecvFrame(uint8_t dest[]) { diff --git a/src/helpers/BaseSerialInterface.h b/src/helpers/BaseSerialInterface.h index 23933fcb4b..f55ab22635 100644 --- a/src/helpers/BaseSerialInterface.h +++ b/src/helpers/BaseSerialInterface.h @@ -17,6 +17,11 @@ class BaseSerialInterface { virtual void loop() {}; virtual bool isWriteBusy() const = 0; + // Returns len once the transport has taken the frame, and 0 only when + // nothing of it was written, so that a caller may safely offer the same + // frame again. A frame the transport tore (a short write on a serial link) + // counts as taken: a retry cannot mend it, the receiver would swallow the + // next header as the missing payload. virtual size_t writeFrame(const uint8_t src[], size_t len) = 0; virtual size_t checkRecvFrame(uint8_t dest[]) = 0; }; diff --git a/src/helpers/MultiSerialInterface.h b/src/helpers/MultiSerialInterface.h index f7742b24c3..3245eaf9ce 100644 --- a/src/helpers/MultiSerialInterface.h +++ b/src/helpers/MultiSerialInterface.h @@ -163,18 +163,23 @@ class MultiSerialInterface : public BaseSerialInterface { return 0; } - // write frame to all enabled interfaces - bool allSuccessful = true; + // write frame to all enabled interfaces. It counts as delivered once ANY + // of them took it (see BaseSerialInterface::writeFrame): an interface that + // is enabled but has nobody attached (BLE advertising, an idle port) fails + // its write, and must not veto the delivery to the client that asked -- a + // caller that keeps a frame until it is delivered would otherwise resend + // it forever. + bool delivered = false; for(auto iface : _interfaces){ if(iface.instance && iface.instance->isEnabled()){ - if(iface.instance->writeFrame(src, len) != len){ - allSuccessful = false; + if(iface.instance->writeFrame(src, len) == len){ + delivered = true; } } } - // report success if all writes completed successfully - return allSuccessful ? len : 0; + // report success if at least one interface delivered the frame + return delivered ? len : 0; } size_t checkRecvFrame(uint8_t dest[]) override { diff --git a/test/test_arduino_serial_interface/test_arduino_serial_interface.cpp b/test/test_arduino_serial_interface/test_arduino_serial_interface.cpp new file mode 100644 index 0000000000..75fea8bf41 --- /dev/null +++ b/test/test_arduino_serial_interface/test_arduino_serial_interface.cpp @@ -0,0 +1,80 @@ +#include + +#include +#include + +#include "helpers/ArduinoSerialInterface.h" + +namespace { + +// A stream that accepts only `capacity` more bytes, the way a CDC port with a +// stalled host or a nearly full TX ring answers a write with a short count. +class CappedStream : public Stream { +public: + explicit CappedStream(size_t capacity) : _capacity(capacity) {} + + size_t write(uint8_t b) override { return write(&b, 1); } + size_t write(const uint8_t* buffer, size_t size) override { + size_t n = size < _capacity ? size : _capacity; + wire.insert(wire.end(), buffer, buffer + n); + _capacity -= n; + return n; + } + + std::vector wire; + +private: + size_t _capacity; +}; + +const uint8_t PAYLOAD[] = { 0x10, 0x20, 0x30, 0x40 }; +const size_t LEN = sizeof(PAYLOAD); +const std::vector FRAME_ON_WIRE = { '>', LEN, 0, 0x10, 0x20, 0x30, 0x40 }; + +size_t writeTo(CappedStream& stream) { + ArduinoSerialInterface iface; + iface.begin(stream); + iface.enable(); + return iface.writeFrame(PAYLOAD, LEN); +} + +} // namespace + +TEST(ArduinoSerialInterface, WholeFrameIsTakenAndReportedAsLen) { + CappedStream stream(64); + EXPECT_EQ(writeTo(stream), LEN); + EXPECT_EQ(stream.wire, FRAME_ON_WIRE); +} + +TEST(ArduinoSerialInterface, NothingWrittenReportsZeroSoTheCallerMayRetry) { + CappedStream stream(0); + EXPECT_EQ(writeTo(stream), 0u); + EXPECT_TRUE(stream.wire.empty()); +} + +TEST(ArduinoSerialInterface, TornHeaderCountsAsTaken) { + CappedStream stream(2); // header cut short, payload never attempted + EXPECT_EQ(writeTo(stream), LEN); + EXPECT_EQ(stream.wire.size(), 2u); +} + +TEST(ArduinoSerialInterface, TornPayloadCountsAsTaken) { + CappedStream stream(3 + LEN - 1); // header out, payload one byte short + EXPECT_EQ(writeTo(stream), LEN); + EXPECT_EQ(stream.wire.size(), 3 + LEN - 1); +} + +TEST(ArduinoSerialInterface, OversizedFrameIsRefusedWithoutWriting) { + CappedStream stream(1024); + uint8_t big[MAX_FRAME_SIZE + 1] = {0}; + ArduinoSerialInterface iface; + iface.begin(stream); + iface.enable(); + EXPECT_EQ(iface.writeFrame(big, sizeof(big)), 0u); + EXPECT_TRUE(stream.wire.empty()); +} + +int main(int argc, char** argv) { + ::testing::InitGoogleTest(&argc, argv); + return RUN_ALL_TESTS(); +} diff --git a/test/test_multi_serial_interface/test_multi_serial_interface.cpp b/test/test_multi_serial_interface/test_multi_serial_interface.cpp new file mode 100644 index 0000000000..2ff304a8e5 --- /dev/null +++ b/test/test_multi_serial_interface/test_multi_serial_interface.cpp @@ -0,0 +1,97 @@ +#include + +#include +#include + +#include "helpers/MultiSerialInterface.h" + +namespace { + +// A transport that either takes every frame whole or refuses them all, the way +// a BLE/WiFi interface behaves with no client attached or a full send queue. +class FakeInterface : public BaseSerialInterface { +public: + explicit FakeInterface(bool accepts) : _accepts(accepts) {} + + void enable() override { _enabled = true; } + void disable() override { _enabled = false; } + bool isEnabled() const override { return _enabled; } + bool isConnected() const override { return _accepts; } + bool isWriteBusy() const override { return false; } + size_t writeFrame(const uint8_t src[], size_t len) override { + writes++; + if (!_accepts) return 0; + last.assign(src, src + len); + return len; + } + size_t checkRecvFrame(uint8_t dest[]) override { return 0; } + + int writes = 0; + std::vector last; + +private: + bool _accepts; + bool _enabled = false; +}; + +const uint8_t FRAME[] = { 0x10, 0x20, 0x30 }; +const std::vector FRAME_BYTES(FRAME, FRAME + sizeof(FRAME)); + +} // namespace + +TEST(MultiSerialInterface, DeliveredWhenOneOfTwoEnabledInterfacesAcceptsIt) { + FakeInterface ble(false); // advertising, nobody connected + FakeInterface usb(true); + MultiSerialInterface multi; + multi.addInterface(InterfaceType::Bluetooth, &ble); + multi.addInterface(InterfaceType::USB, &usb); + multi.enable(); + + EXPECT_EQ(multi.writeFrame(FRAME, sizeof(FRAME)), sizeof(FRAME)); + EXPECT_EQ(ble.writes, 1); // still offered to every enabled interface + EXPECT_EQ(usb.writes, 1); + EXPECT_EQ(usb.last, FRAME_BYTES); +} + +TEST(MultiSerialInterface, NotDeliveredWhenNoInterfaceAcceptsIt) { + FakeInterface ble(false); + FakeInterface usb(false); + MultiSerialInterface multi; + multi.addInterface(InterfaceType::Bluetooth, &ble); + multi.addInterface(InterfaceType::USB, &usb); + multi.enable(); + + EXPECT_EQ(multi.writeFrame(FRAME, sizeof(FRAME)), 0u); + EXPECT_EQ(ble.writes, 1); + EXPECT_EQ(usb.writes, 1); +} + +TEST(MultiSerialInterface, DisabledInterfaceIsNeitherWrittenNorCounted) { + FakeInterface ble(true); + FakeInterface usb(false); + MultiSerialInterface multi; + multi.addInterface(InterfaceType::Bluetooth, &ble); + multi.addInterface(InterfaceType::USB, &usb); + multi.enable(); + ble.disable(); + + EXPECT_EQ(multi.writeFrame(FRAME, sizeof(FRAME)), 0u); + EXPECT_EQ(ble.writes, 0); + EXPECT_EQ(usb.writes, 1); +} + +TEST(MultiSerialInterface, NothingIsWrittenWhileDisabledOrEmpty) { + FakeInterface usb(true); + MultiSerialInterface multi; + multi.addInterface(InterfaceType::USB, &usb); + + EXPECT_EQ(multi.writeFrame(FRAME, sizeof(FRAME)), 0u); // manager not enabled + multi.enable(); + EXPECT_EQ(multi.writeFrame(FRAME, 0), 0u); // empty frame + EXPECT_EQ(usb.writes, 0); +} + +int main(int argc, char** argv) { + ::testing::InitGoogleTest(&argc, argv); + return RUN_ALL_TESTS(); +}