Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 18 additions & 7 deletions examples/companion_radio/MyMesh.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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);
Expand Down
3 changes: 2 additions & 1 deletion examples/companion_radio/MyMesh.h
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand Down
1 change: 1 addition & 0 deletions platformio.ini
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
9 changes: 7 additions & 2 deletions src/helpers/ArduinoSerialInterface.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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[]) {
Expand Down
5 changes: 5 additions & 0 deletions src/helpers/BaseSerialInterface.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;
};
17 changes: 11 additions & 6 deletions src/helpers/MultiSerialInterface.h
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,80 @@
#include <gtest/gtest.h>

#include <cstdint>
#include <vector>

#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<uint8_t> wire;

private:
size_t _capacity;
};

const uint8_t PAYLOAD[] = { 0x10, 0x20, 0x30, 0x40 };
const size_t LEN = sizeof(PAYLOAD);
const std::vector<uint8_t> 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();
}
97 changes: 97 additions & 0 deletions test/test_multi_serial_interface/test_multi_serial_interface.cpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,97 @@
#include <gtest/gtest.h>

#include <cstdint>
#include <vector>

#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<uint8_t> last;

private:
bool _accepts;
bool _enabled = false;
};

const uint8_t FRAME[] = { 0x10, 0x20, 0x30 };
const std::vector<uint8_t> 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();
}