From 3885c67c8eaf46ce66e28252338df783ca178a95 Mon Sep 17 00:00:00 2001 From: Alexander Hoffer Date: Mon, 20 Jul 2026 15:47:03 +0100 Subject: [PATCH] fix: synchronize BLE receive queue --- src/helpers/BoundedQueue.h | 39 +++++++++++++ src/helpers/esp32/SerialBLEInterface.cpp | 41 ++++++++----- src/helpers/esp32/SerialBLEInterface.h | 10 ++-- .../test_bounded_queue/test_bounded_queue.cpp | 58 +++++++++++++++++++ 4 files changed, 129 insertions(+), 19 deletions(-) create mode 100644 src/helpers/BoundedQueue.h create mode 100644 test/test_bounded_queue/test_bounded_queue.cpp diff --git a/src/helpers/BoundedQueue.h b/src/helpers/BoundedQueue.h new file mode 100644 index 00000000..07ecb4ff --- /dev/null +++ b/src/helpers/BoundedQueue.h @@ -0,0 +1,39 @@ +#pragma once + +#include + +namespace mesh { + +template +class BoundedQueue { + static_assert(Capacity > 0, "BoundedQueue capacity must be greater than zero"); + + T _items[Capacity]; + size_t _head = 0; + size_t _size = 0; + +public: + bool push(const T& item) { + if (_size == Capacity) return false; + _items[(_head + _size) % Capacity] = item; + _size++; + return true; + } + + bool pop(T& item) { + if (_size == 0) return false; + item = _items[_head]; + _head = (_head + 1) % Capacity; + _size--; + return true; + } + + void clear() { + _head = 0; + _size = 0; + } + + size_t size() const { return _size; } +}; + +} // namespace mesh diff --git a/src/helpers/esp32/SerialBLEInterface.cpp b/src/helpers/esp32/SerialBLEInterface.cpp index dcfa0e1e..9c81dc7b 100644 --- a/src/helpers/esp32/SerialBLEInterface.cpp +++ b/src/helpers/esp32/SerialBLEInterface.cpp @@ -118,17 +118,30 @@ void SerialBLEInterface::onWrite(BLECharacteristic* pCharacteristic, esp_ble_gat if (len > MAX_FRAME_SIZE) { BLE_DEBUG_PRINTLN("ERROR: onWrite(), frame too big, len=%d", len); - } else if (recv_queue_len >= FRAME_QUEUE_SIZE) { - BLE_DEBUG_PRINTLN("ERROR: onWrite(), recv_queue is full!"); } else { - recv_queue[recv_queue_len].len = len; - memcpy(recv_queue[recv_queue_len].buf, rxValue, len); - recv_queue_len++; + Frame frame = {}; + frame.len = len; + memcpy(frame.buf, rxValue, len); + + portENTER_CRITICAL(&recv_queue_mux); + bool queued = recv_queue.push(frame); + portEXIT_CRITICAL(&recv_queue_mux); + + if (!queued) { + BLE_DEBUG_PRINTLN("ERROR: onWrite(), recv_queue is full!"); + } } } // ---------- public methods +void SerialBLEInterface::clearBuffers() { + portENTER_CRITICAL(&recv_queue_mux); + recv_queue.clear(); + portEXIT_CRITICAL(&recv_queue_mux); + send_queue_len = 0; +} + void SerialBLEInterface::enable() { if (_isEnabled) return; @@ -202,17 +215,15 @@ size_t SerialBLEInterface::checkRecvFrame(uint8_t dest[]) { } } - if (recv_queue_len > 0) { // check recv queue - size_t len = recv_queue[0].len; // take from top of queue - memcpy(dest, recv_queue[0].buf, len); + Frame frame = {}; + portENTER_CRITICAL(&recv_queue_mux); + bool received = recv_queue.pop(frame); + portEXIT_CRITICAL(&recv_queue_mux); - BLE_DEBUG_PRINTLN("readBytes: sz=%d, hdr=%d", len, (uint32_t) dest[0]); - - recv_queue_len--; - for (int i = 0; i < recv_queue_len; i++) { // delete top item from queue - recv_queue[i] = recv_queue[i + 1]; - } - return len; + if (received) { + memcpy(dest, frame.buf, frame.len); + BLE_DEBUG_PRINTLN("readBytes: sz=%d, hdr=%d", (uint32_t) frame.len, (uint32_t) dest[0]); + return frame.len; } if (pServer->getConnectedCount() == 0) deviceConnected = false; diff --git a/src/helpers/esp32/SerialBLEInterface.h b/src/helpers/esp32/SerialBLEInterface.h index 965e90fd..e2d600a0 100644 --- a/src/helpers/esp32/SerialBLEInterface.h +++ b/src/helpers/esp32/SerialBLEInterface.h @@ -1,10 +1,12 @@ #pragma once #include "../BaseSerialInterface.h" +#include "../BoundedQueue.h" #include #include #include #include +#include class SerialBLEInterface : public BaseSerialInterface, BLESecurityCallbacks, BLEServerCallbacks, BLECharacteristicCallbacks { BLEServer *pServer; @@ -24,12 +26,12 @@ class SerialBLEInterface : public BaseSerialInterface, BLESecurityCallbacks, BLE }; #define FRAME_QUEUE_SIZE 4 - int recv_queue_len; - Frame recv_queue[FRAME_QUEUE_SIZE]; + portMUX_TYPE recv_queue_mux = portMUX_INITIALIZER_UNLOCKED; + mesh::BoundedQueue recv_queue; int send_queue_len; Frame send_queue[FRAME_QUEUE_SIZE]; - void clearBuffers() { recv_queue_len = 0; send_queue_len = 0; } + void clearBuffers(); protected: // BLESecurityCallbacks methods @@ -58,7 +60,7 @@ public: _isEnabled = false; _last_write = 0; last_conn_id = 0; - send_queue_len = recv_queue_len = 0; + send_queue_len = 0; } /** diff --git a/test/test_bounded_queue/test_bounded_queue.cpp b/test/test_bounded_queue/test_bounded_queue.cpp new file mode 100644 index 00000000..3e2e2de7 --- /dev/null +++ b/test/test_bounded_queue/test_bounded_queue.cpp @@ -0,0 +1,58 @@ +#include + +#include + +TEST(BoundedQueue, PreservesOrderAcrossWraparound) { + mesh::BoundedQueue queue; + int value = 0; + + for (int i = 0; i < 4; i++) EXPECT_TRUE(queue.push(i)); + EXPECT_FALSE(queue.push(4)); + + EXPECT_TRUE(queue.pop(value)); + EXPECT_EQ(0, value); + EXPECT_TRUE(queue.pop(value)); + EXPECT_EQ(1, value); + + EXPECT_TRUE(queue.push(4)); + EXPECT_TRUE(queue.push(5)); + + for (int expected = 2; expected < 6; expected++) { + EXPECT_TRUE(queue.pop(value)); + EXPECT_EQ(expected, value); + } + EXPECT_FALSE(queue.pop(value)); +} + +TEST(BoundedQueue, SurvivesRepeatedFillAndDrainCycles) { + mesh::BoundedQueue queue; + int value = 0; + + for (int cycle = 0; cycle < 1000; cycle++) { + for (int i = 0; i < 4; i++) EXPECT_TRUE(queue.push(cycle * 4 + i)); + for (int i = 0; i < 4; i++) { + EXPECT_TRUE(queue.pop(value)); + EXPECT_EQ(cycle * 4 + i, value); + } + } +} + +TEST(BoundedQueue, ClearDropsPendingItems) { + mesh::BoundedQueue queue; + int value = 0; + + EXPECT_TRUE(queue.push(1)); + EXPECT_TRUE(queue.push(2)); + queue.clear(); + + EXPECT_EQ(0u, queue.size()); + EXPECT_FALSE(queue.pop(value)); + EXPECT_TRUE(queue.push(3)); + EXPECT_TRUE(queue.pop(value)); + EXPECT_EQ(3, value); +} + +int main(int argc, char **argv) { + ::testing::InitGoogleTest(&argc, argv); + return RUN_ALL_TESTS(); +}