From efc85b1afe5a281abdbb36bea8e61861c3c38e16 Mon Sep 17 00:00:00 2001 From: mikecarper Date: Fri, 25 Sep 2026 12:06:18 -0700 Subject: [PATCH] Drain companion queue on explicit OTA entry --- docs/companion_radio_full.md | 5 ++ examples/companion_radio/MyMesh.cpp | 48 +++++++++++++++++-- examples/companion_radio/MyMesh.h | 1 + examples/companion_radio/main.cpp | 9 +++- src/helpers/ota/OtaContext.cpp | 17 +++++-- src/helpers/ota/OtaContext.h | 5 +- .../test_shared_mota_queue.cpp | 42 ++++++++++++++-- .../test_wifi_shared_queue.h | 2 +- test/test_companion_ota_config.py | 13 ++++- 9 files changed, 128 insertions(+), 14 deletions(-) diff --git a/docs/companion_radio_full.md b/docs/companion_radio_full.md index 5a6adbb89..174c1c719 100644 --- a/docs/companion_radio_full.md +++ b/docs/companion_radio_full.md @@ -361,6 +361,11 @@ Every corrected nRF52 Full profile keeps **256 offline frames normally** and temporarily lends 128 slots to mOTA to leave room for Bluetooth tasks, displays and UI allocations. Queue sharing preserves the board's contacts, channels, and USB/Bluetooth mOTA sending. See the [memory correction](releases/1.17.1.5.md#t096-full-companion-bluetooth-and-menu-freeze-report). +If an explicit OTA source, discovery, download, install, or temporary-radio +request needs that workspace while more than 128 unread frames are queued, the +oldest volatile frames are consumed until 128 remain. The command reports the +number removed. Read-only OTA status queries do not remove messages. Sync the +queue to an app before starting OTA when those unread messages matter. The nRF52 target inherits the board's ordinary USB Companion installation format and adds BLE plus the serial mOTA source. It does not enable an SD cache diff --git a/examples/companion_radio/MyMesh.cpp b/examples/companion_radio/MyMesh.cpp index 261def548..5f518459e 100644 --- a/examples/companion_radio/MyMesh.cpp +++ b/examples/companion_radio/MyMesh.cpp @@ -461,6 +461,19 @@ mesh::ota::OtaContext* MyMesh::acquireOfflineQueueForOta(void* owner) { return mesh->offline_queue.acquire(mesh->offline_queue_len, mesh->offline_queue_head); } +uint16_t MyMesh::drainOfflineQueueForOta(void* owner) { + MyMesh* mesh = static_cast(owner); + uint8_t discarded[MAX_FRAME_SIZE]; + uint16_t removed = 0; + // Explicit OTA work borrows the upper half of this volatile queue. Keep the + // newest 128 frames and notify the UI for each oldest frame consumed. + while (mesh->offline_queue_len > 128) { + if (mesh->getFromOfflineQueue(discarded) <= 0) break; + ++removed; + } + return removed; +} + void MyMesh::releaseOfflineQueueFromOta(void* owner) { MyMesh* mesh = static_cast(owner); mesh->offline_queue.release(mesh->offline_queue_head); @@ -469,7 +482,8 @@ void MyMesh::releaseOfflineQueueFromOta(void* owner) { void MyMesh::initializeOfflineQueue() { #if defined(OTA_SHARED_COMPANION_QUEUE) - mesh::ota::ota_set_context_storage(this, acquireOfflineQueueForOta, releaseOfflineQueueFromOta); + mesh::ota::ota_set_context_storage(this, acquireOfflineQueueForOta, + releaseOfflineQueueFromOta, drainOfflineQueueForOta); #elif defined(ESP32_PLATFORM) && defined(BOARD_HAS_PSRAM) if (offline_queue != offline_queue_fallback || OFFLINE_QUEUE_SIZE <= offline_queue_capacity) return; @@ -2287,7 +2301,8 @@ bool MyMesh::scheduleTempRadio(float freq, float bw, uint8_t sf, uint8_t cr, } #if defined(OTA_SHARED_COMPANION_QUEUE) - if (!mesh::ota::ota_acquire_context(reply, reply_size)) return false; + uint16_t drained = 0; + if (!mesh::ota::ota_acquire_context(reply, reply_size, true, &drained)) return false; mesh::ota::ota_ctx().release_when_idle = false; #endif _temp_radio_freq = freq; @@ -2301,6 +2316,13 @@ bool MyMesh::scheduleTempRadio(float freq, float bw, uint8_t sf, uint8_t cr, snprintf(reply, reply_size, "OK - temp params for %lu mins", (unsigned long)timeout_mins); appendRxPowerSavingAdjustmentNote(reply, reply_size, sf, bw); +#if defined(OTA_SHARED_COMPANION_QUEUE) + if (drained) { + size_t const used = strlen(reply); + snprintf(reply + used, reply_size - used, "; OTA cleared %u oldest queued messages", + (unsigned)drained); + } +#endif return true; } @@ -3375,7 +3397,22 @@ bool MyMesh::handleLocalControlCommand(const char* command, char* reply, if (strncmp(command, "ota", 3) == 0 && (command[3] == 0 || command[3] == ' ')) { char ota_reply[160] = {0}; - if (!mesh::ota::ota_acquire_context(ota_reply, sizeof(ota_reply))) { + const char* action = command + 3; + while (*action == ' ') ++action; + // Only an explicit transfer, catalog search, or source attach may evict + // unread offline frames. Read-only OTA queries leave the queue intact. + bool const drain_queue = strcmp(action, "folder on") == 0 + || strcmp(action, "ls") == 0 || strncmp(action, "ls ", 3) == 0 + || strncmp(action, "pull ", 5) == 0 + || strncmp(action, "get ", 4) == 0 + || strncmp(action, "download ", 9) == 0 + || strcmp(action, "install") == 0 + || strcmp(action, "apply") == 0 + || strcmp(action, "applydelta") == 0 + || strncmp(action, "bootloader install ", 19) == 0 + || strncmp(action, "rescue install ", 15) == 0; + uint16_t drained = 0; + if (!mesh::ota::ota_acquire_context(ota_reply, sizeof(ota_reply), drain_queue, &drained)) { snprintf(reply, reply_size, "%s", ota_reply); return true; } @@ -3389,6 +3426,11 @@ bool MyMesh::handleLocalControlCommand(const char* command, char* reply, } context.config_dirty = false; } + if (drained) { + size_t const used = strlen(ota_reply); + snprintf(ota_reply + used, sizeof(ota_reply) - used, + "; cleared %u oldest queued messages for OTA", (unsigned)drained); + } snprintf(reply, reply_size, "%s", ota_reply); return true; } diff --git a/examples/companion_radio/MyMesh.h b/examples/companion_radio/MyMesh.h index 0a82dae87..8649c1697 100644 --- a/examples/companion_radio/MyMesh.h +++ b/examples/companion_radio/MyMesh.h @@ -646,6 +646,7 @@ private: static_assert(OFFLINE_QUEUE_SIZE == 256, "Shared mOTA queue requires 256 normal slots"); mesh::BorrowableFrameBuffer offline_queue; static mesh::ota::OtaContext* acquireOfflineQueueForOta(void* owner); + static uint16_t drainOfflineQueueForOta(void* owner); static void releaseOfflineQueueFromOta(void* owner); #elif defined(ESP32_PLATFORM) && defined(BOARD_HAS_PSRAM) enum { diff --git a/examples/companion_radio/main.cpp b/examples/companion_radio/main.cpp index 633077239..8d9d513a5 100644 --- a/examples/companion_radio/main.cpp +++ b/examples/companion_radio/main.cpp @@ -329,7 +329,8 @@ public: return false; } - if (!mesh::ota::ota_acquire_context(reply, reply_size)) return false; + uint16_t drained = 0; + if (!mesh::ota::ota_acquire_context(reply, reply_size, true, &drained)) return false; mesh::ota::OtaContext& context = mesh::ota::ota_ctx(); if (context.folder_active && context.folderLink() != mesh::ota::OtaContext::FOLDER_LINK_BLE) { @@ -348,6 +349,12 @@ public: return false; } + if (drained) { + size_t const used = strlen(reply); + snprintf(reply + used, reply_size - used, + "; cleared %u oldest queued messages for OTA", (unsigned)drained); + } + context.manager.announce(); mesh::usbLoggingPort().println("Bluetooth mOTA source attached"); return true; diff --git a/src/helpers/ota/OtaContext.cpp b/src/helpers/ota/OtaContext.cpp index a85e7108c..6ae82b703 100644 --- a/src/helpers/ota/OtaContext.cpp +++ b/src/helpers/ota/OtaContext.cpp @@ -17,6 +17,7 @@ OtaContext* active_context = nullptr; void* storage_owner = nullptr; OtaContext* (*acquire_storage)(void*) = nullptr; void (*release_storage)(void*) = nullptr; +uint16_t (*drain_storage)(void*) = nullptr; OtaSend saved_send = nullptr; void* saved_send_ctx = nullptr; uint32_t saved_target = 0; @@ -50,11 +51,12 @@ void releaseHeapContext(void*) { } void ota_set_context_storage(void* owner, OtaContext* (*acquire)(void*), - void (*release)(void*)) { + void (*release)(void*), uint16_t (*drain)(void*)) { assert(!active_context); storage_owner = owner; acquire_storage = acquire; release_storage = release; + drain_storage = drain; } OtaContext& ota_ctx() { @@ -80,7 +82,8 @@ void ota_refresh_seeder_identity(const uint8_t* seeder_id) { if (active_context) active_context->manager.set_seeder_id(seeder_id); } -bool ota_acquire_context(char* reply, size_t cap) { +bool ota_acquire_context(char* reply, size_t cap, bool drain_queue, uint16_t* drained) { + if (drained) *drained = 0; if (active_context) return true; #if defined(OTA_HEAP_CONTEXT) if (!acquire_storage) { // no owner registered one: fall back to the heap @@ -93,6 +96,11 @@ bool ota_acquire_context(char* reply, size_t cap) { return false; } active_context = acquire_storage(storage_owner); + if (!active_context && drain_queue && drain_storage) { + uint16_t const count = drain_storage(storage_owner); + if (drained) *drained = count; + active_context = acquire_storage(storage_owner); + } if (!active_context) { if (reply && cap) snprintf(reply, cap, #if defined(OTA_HEAP_CONTEXT) @@ -153,7 +161,10 @@ OtaContext& ota_ctx() { } OtaContext* ota_context_if_active() { return &ota_ctx(); } -bool ota_acquire_context(char*, size_t) { return true; } +bool ota_acquire_context(char*, size_t, bool, uint16_t* drained) { + if (drained) *drained = 0; + return true; +} void ota_begin_context(uint32_t target, OtaSend send, void* ctx, const char* hw, const uint8_t* seeder_id) { ota_ctx().begin(target, send, ctx, hw); diff --git a/src/helpers/ota/OtaContext.h b/src/helpers/ota/OtaContext.h index 1698ab137..c0cf8a048 100644 --- a/src/helpers/ota/OtaContext.h +++ b/src/helpers/ota/OtaContext.h @@ -744,7 +744,8 @@ OtaContext* ota_context_if_active(); // Optional role-specific policy reload; registration updates idle/active policy // without allocating a context. A failed load leaves the prior policy intact. void ota_set_context_config_loader(bool (*load)(OtaConfigState&)); -bool ota_acquire_context(char* reply, size_t cap); +bool ota_acquire_context(char* reply, size_t cap, bool drain_queue = false, + uint16_t* drained = nullptr); void ota_begin_context(uint32_t target, OtaSend send, void* ctx, const char* hw, const uint8_t* seeder_id); // Identity can become available after Mesh::begin or change through key import. @@ -766,7 +767,7 @@ uint8_t ota_hop_limit(); // ota_ctx() is only valid while storage is held. Callers that can run before a // successful ota_acquire_context() must gate on ota_context_if_active() first. void ota_set_context_storage(void* owner, OtaContext* (*acquire)(void*), - void (*release)(void*)); + void (*release)(void*), uint16_t (*drain)(void*) = nullptr); void ota_release_context_if_idle(bool temporary_radio_active); // Loop helper for roles whose LoRa OTA only runs under the temporary radio diff --git a/test/fixtures/shared_mota_queue/test_shared_mota_queue.cpp b/test/fixtures/shared_mota_queue/test_shared_mota_queue.cpp index 083efb648..75cbb603f 100644 --- a/test/fixtures/shared_mota_queue/test_shared_mota_queue.cpp +++ b/test/fixtures/shared_mota_queue/test_shared_mota_queue.cpp @@ -25,6 +25,16 @@ struct Queue { Queue& q = *static_cast(owner); return q.buffer.acquire(q.count, q.head); } + static uint16_t drain(void* owner) { + Queue& q = *static_cast(owner); + uint16_t removed = 0; + while (q.count > 128) { + q.head = (q.head + 1) % q.buffer.capacity(); + --q.count; + ++removed; + } + return removed; + } static void release(void* owner) { Queue& q = *static_cast(owner); q.buffer.release(q.head); @@ -119,6 +129,7 @@ static void transfer(OtaContext& context, void* reply_route) { int main() { Queue q; + unsigned first_retained = 0; for (int i = 0; i < q.count; ++i) { Frame& frame = q.buffer.at((q.head + i) % q.buffer.capacity()); frame.len = 176; @@ -128,13 +139,13 @@ int main() { for (int i = 0; i < q.count; ++i) { const Frame& frame = q.buffer.at((q.head + i) % q.buffer.capacity()); assert(frame.len == 176); - for (auto byte : frame.buf) assert(byte == i); + for (auto byte : frame.buf) assert(byte == i + first_retained); } }; char reply[160] = {}; assert(!ota_acquire_context(reply, sizeof reply)); assert(strstr(reply, "not ready")); - ota_set_context_storage(&q, Queue::acquire, Queue::release); + ota_set_context_storage(&q, Queue::acquire, Queue::release, Queue::drain); uint8_t identity[4] = {1, 2, 3, 4}; ota_begin_context(0, send, nullptr, "Heltec_t096", identity); assert(!ota_context_if_active() && q.buffer.capacity() == 256); @@ -143,7 +154,11 @@ int main() { assert(strstr(reply, "sync unread messages")); assert(q.head == 191 && q.count == 129 && q.buffer.capacity() == 256); check_messages(); - --q.count; + uint16_t drained = 0; + assert(ota_acquire_context(reply, sizeof reply, true, &drained)); + assert(drained == 1 && q.count == 128 && q.head == 0); + first_retained = 1; + check_messages(); for (unsigned cycle = 0; cycle < 8; ++cycle) { assert(ota_acquire_context(reply, sizeof reply)); @@ -249,4 +264,25 @@ int main() { assert(!ota_context_if_active() && q.buffer.capacity() == 256); check_messages(); #endif + + Queue full; + full.count = 256; + for (int i = 0; i < full.count; ++i) { + Frame& frame = full.buffer.at((full.head + i) % full.buffer.capacity()); + frame.len = 176; + memset(frame.buf, i, sizeof frame.buf); + } + ota_set_context_storage(&full, Queue::acquire, Queue::release, Queue::drain); + ota_begin_context(0, send, nullptr, "Heltec_t096", identity); + assert(!ota_acquire_context(reply, sizeof reply)); + drained = 0; + assert(ota_acquire_context(reply, sizeof reply, true, &drained)); + assert(drained == 128 && full.count == 128 && full.head == 0); + for (int i = 0; i < full.count; ++i) { + const Frame& frame = full.buffer.at(i); + assert(frame.len == 176); + for (auto byte : frame.buf) assert(byte == i + 128); + } + ota_release_context_if_idle(false); + assert(!ota_context_if_active() && full.buffer.capacity() == 256); } diff --git a/test/fixtures/shared_mota_queue/test_wifi_shared_queue.h b/test/fixtures/shared_mota_queue/test_wifi_shared_queue.h index 80c7f14e4..ed2133351 100644 --- a/test/fixtures/shared_mota_queue/test_wifi_shared_queue.h +++ b/test/fixtures/shared_mota_queue/test_wifi_shared_queue.h @@ -67,7 +67,7 @@ static void test_wifi_shared_queue(Queue& q, CheckMessages check_messages) { q.count = 129; Frame& extra = q.buffer.at((q.head + 128) % 256); extra.len = 176; - memset(extra.buf, 128, sizeof extra.buf); + memset(extra.buf, 129, sizeof extra.buf); auto refused = connect_host(); WiFiOtaSeeder::loop(); assert(!refused->connected && refused->requests == 0); diff --git a/test/test_companion_ota_config.py b/test/test_companion_ota_config.py index c5a6fba43..1886de863 100644 --- a/test/test_companion_ota_config.py +++ b/test/test_companion_ota_config.py @@ -47,7 +47,11 @@ struct OtaContext { }; static OtaContext context; static bool acquire_ok=true; -static bool ota_acquire_context(char* reply, size_t cap) { +static bool drain_requested=false; +static bool ota_acquire_context(char* reply, size_t cap, bool drain=false, + uint16_t* drained=nullptr) { + drain_requested=drain; + if (drained) *drained=0; if (!acquire_ok) snprintf(reply, cap, "ERR unavailable"); return acquire_ok; } @@ -62,6 +66,10 @@ static bool config(const char* rest, char* reply, OtaContext& c) { return true; } static bool handle_ota_command(const char* command, char* reply, int) { + if (!strcmp(command, "ota folder on")) { + strcpy(reply, "OK folder attached"); + return true; + } if (!strcmp(command, "ota key add test-key")) { uint8_t key[32]={93}; context.allow.add(key); context.config_dirty=true; @@ -114,7 +122,10 @@ int main() { assert(!context.config_dirty); }; reboot(fs); + command("ota folder on", "OK folder attached"); + assert(drain_requested); command("ota config", "ota config: speed=1x packet=0.5x"); + assert(!drain_requested); for (const auto* setting : {"hops", "checkpoint", "advert"}) { const auto before=OtaConfigState::capture(context); for (const auto* bad : {"", " ", "x", "-1", "1x", "1 2", "1.0", "4294967296", "9999999999999999999999"}) {