From bf42395a540827e759cf7f8b45338a36c0ae65a0 Mon Sep 17 00:00:00 2001 From: notbucki <103531753+notbucki@users.noreply.github.com> Date: Sat, 5 Sep 2026 21:14:58 +0200 Subject: [PATCH] Companion: keep a queued message until the sync reply is actually written CMD_SYNC_NEXT_MESSAGE took the oldest frame out of the offline queue and then ignored the return value of writeFrame(). BLE and WiFi refuse a frame while their send queue is full, and a serial write can come back short, so in each of those cases the message was gone for good while the client saw nothing and asked for the next one. The frame now stays queued until the transport reports that it took it whole; on failure the client simply repeats the request. MultiSerialInterface reported success only when every enabled interface accepted the frame. An interface that is enabled but has nobody attached (BLE advertising, an idle port) fails each write, which would have made the retry above resend the same message forever on any build with more than one interface. It now counts a frame as delivered once any enabled interface took it. Covered by a native googletest. For that to be safe, ArduinoSerialInterface::writeFrame() now reports 0 only when nothing went out. A torn frame (short write on a CDC port with a stalled host) counts as taken: a retry cannot mend it, the receiver would swallow the next header as the missing payload, so the message is consumed as it was before. The contract is spelled out on BaseSerialInterface; BLE, WiFi and Ethernet already behaved that way. Covered by a second native googletest. Co-Authored-By: Claude Fable 5.1 --- examples/companion_radio/MyMesh.cpp | 25 +++-- examples/companion_radio/MyMesh.h | 3 +- platformio.ini | 1 + src/helpers/ArduinoSerialInterface.cpp | 9 +- src/helpers/BaseSerialInterface.h | 5 + src/helpers/MultiSerialInterface.h | 17 ++-- .../test_arduino_serial_interface.cpp | 80 +++++++++++++++ .../test_multi_serial_interface.cpp | 97 +++++++++++++++++++ 8 files changed, 221 insertions(+), 16 deletions(-) create mode 100644 test/test_arduino_serial_interface/test_arduino_serial_interface.cpp create mode 100644 test/test_multi_serial_interface/test_multi_serial_interface.cpp 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(); +}