diff --git a/CHANGELOG.md b/CHANGELOG.md index 4cdcd8c7..6e461ef8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,7 @@ ### Firmware +- A weak or lost Wi-Fi link no longer freezes the device until a hard reset. Closing the WebSocket used esp_websocket_client_close(), which sends its close frame with no time limit and blocked the app task, and with it the screen, touch and Wi-Fi retries, on a dead link. Sends now run on their own task with a 2 s limit, so a slow link no longer stalls the screen either, and a lost connection is noticed even when its event can't be queued. Contributed by Antonio Lourenco (@tozes). - A console command that times out no longer leaves the app task writing its reply into freed memory. The request slot is shared by both tasks and freed by whichever finishes last, so a late reply is dropped instead of corrupting the stack. The slot is host-tested. - The device refuses an over-the-air image built for another board, by reading the `HGBOARD=` tag as the image streams in, so a wrong image can no longer reach Wi-Fi and Hermes and pass the rollback check. The plugin's own check stays as the first line. - Builds from `main` with that board check no longer refuse every over-the-air update with "the image is for , not this board". The check's search text put a nameless copy of the `HGBOARD=` tag in every image, ahead of the real one, and the device stopped at it. Images now carry the tag once, and the device, the plugin and the release packager skip a nameless tag in images built before this fix. Contributed by Antonio Lourenco (@tozes). diff --git a/firmware/CMakeLists.txt b/firmware/CMakeLists.txt index d9148ab9..018f7fd3 100644 --- a/firmware/CMakeLists.txt +++ b/firmware/CMakeLists.txt @@ -25,7 +25,7 @@ endif() option(HG_BUILD_TESTS "Build the core unit tests" ON) if(HG_BUILD_TESTS) enable_testing() - add_executable(hg_core_tests tests/main.cpp tests/test_basics.cpp tests/test_app.cpp tests/test_qr.cpp tests/test_power.cpp tests/test_ws185.cpp tests/test_speaker_pa.cpp tests/test_speaker_queue.cpp tests/test_band_flush.cpp tests/test_shared_reply.cpp tests/test_tag_scanner.cpp drivers/axp2101.cpp drivers/cores3.cpp) + add_executable(hg_core_tests tests/main.cpp tests/test_basics.cpp tests/test_app.cpp tests/test_qr.cpp tests/test_power.cpp tests/test_ws185.cpp tests/test_speaker_pa.cpp tests/test_speaker_queue.cpp tests/test_band_flush.cpp tests/test_shared_reply.cpp tests/test_tag_scanner.cpp tests/test_ws_link.cpp drivers/axp2101.cpp drivers/cores3.cpp) target_include_directories(hg_core_tests PRIVATE drivers) target_link_libraries(hg_core_tests PRIVATE hg_core) if(MSVC) diff --git a/firmware/drivers/ws_link.hpp b/firmware/drivers/ws_link.hpp new file mode 100644 index 00000000..b2d7fac6 --- /dev/null +++ b/firmware/drivers/ws_link.hpp @@ -0,0 +1,77 @@ +// The WebSocket transport's bookkeeping, apart from ESP-IDF so it can be tested +// on the host: which connection is current, whether it was lost, and what the +// send task does with each queued item (firmware/esp32/main/port_ws.cpp). +#pragma once +#include +#include +#include + +namespace hg::ws { + +// Each connection gets a generation. Leaving a connection moves to the next +// one, so whatever the old client still reports can be told apart and ignored. +class Generations { + public: + uint32_t current() const { return current_.load(); } + bool is_current(uint32_t generation) const { return generation == current_.load(); } + // The app leaves the current connection (before connecting again, or for good). + void leave() { ++current_; } + // A client reports its connection lost (any task). + void lost(uint32_t generation) { lost_.store(generation); } + // True once per loss of the current connection; losses of connections the app + // has already left are not news. + bool take_lost() { + uint32_t g = lost_.exchange(0); + return g && g == current_.load(); + } + + private: + std::atomic current_{0}; + std::atomic lost_{0}; +}; + +// The send queue holds messages plus a few slots only "close this client" may +// use, so leaving a connection never waits for a slow link to drain. +constexpr size_t kTxMessages = 64; +constexpr size_t kTxReserved = 4; +inline bool may_queue_message(size_t free_slots) { return free_slots > kTxReserved; } + +constexpr uint8_t kText = 0x1; +constexpr uint8_t kBinary = 0x2; +constexpr uint8_t kClose = 0xFF; // not a WebSocket opcode: close and destroy the client + +// Every network wait has a limit. The library's own esp_websocket_client_close() +// waits for its close frame without one and froze the app on a dead link. +constexpr uint32_t kSendTimeoutMs = 2000; +constexpr uint32_t kCloseTimeoutMs = 1000; + +struct TxItem { + void* client; + uint32_t generation; + uint8_t opcode; // kText, kBinary or kClose + const uint8_t* data; + size_t len; +}; + +enum class Outcome { Sent, Failed, Dropped, Destroyed }; + +// What the send task does with one item. Ops supplies the client calls: +// bool is_connected(void* client) +// int send(void* client, uint8_t opcode, const uint8_t* data, size_t len, uint32_t timeout_ms) +// void send_close(void* client, uint32_t timeout_ms) +// void destroy(void* client) +// A failed send needs no handling here: the client aborts the connection and +// reports it, which ends in Generations::lost(). +template +Outcome process(const TxItem& item, uint32_t current_generation, Ops& ops) { + if (item.opcode == kClose) { + if (ops.is_connected(item.client)) ops.send_close(item.client, kCloseTimeoutMs); + ops.destroy(item.client); + return Outcome::Destroyed; + } + if (item.generation != current_generation || !ops.is_connected(item.client)) return Outcome::Dropped; + int sent = ops.send(item.client, item.opcode, item.data, item.len, kSendTimeoutMs); + return sent == static_cast(item.len) ? Outcome::Sent : Outcome::Failed; +} + +} // namespace hg::ws diff --git a/firmware/esp32/main/main.cpp b/firmware/esp32/main/main.cpp index c1717f3c..9ac5a840 100644 --- a/firmware/esp32/main/main.cpp +++ b/firmware/esp32/main/main.cpp @@ -274,6 +274,8 @@ extern "C" void app_main(void) { hgp::events::release(ev); } while (hgp::events::receive(ev, 0)); } + // A loss whose WsClosed event was dropped; harmless if it was delivered too. + if (g_transport.take_lost()) app.on_transport_closed("disconnected"); g_buttons.poll(app); g_wifi.tick(app, g_system.now_ms()); if (g_gestures) g_gestures->tick(g_system.now_ms()); diff --git a/firmware/esp32/main/port.hpp b/firmware/esp32/main/port.hpp index e2fde559..7452dac5 100644 --- a/firmware/esp32/main/port.hpp +++ b/firmware/esp32/main/port.hpp @@ -22,6 +22,7 @@ #include "shared_reply.hpp" #include "speaker_pa.hpp" #include "tag_scanner.hpp" +#include "ws_link.hpp" #include "driver/i2c_master.h" #include "driver/i2s_std.h" #include "esp_codec_dev.h" @@ -121,14 +122,25 @@ class WsTransport final : public hg::Transport { bool send_text(std::string_view text) override; bool send_binary(const uint8_t* data, size_t len) override; void close() override; - uint32_t generation() const { return generation_.load(); } + uint32_t generation() const { return gens_.current(); } + // True once if the current connection was lost. Backs up the WsClosed event, + // which the client may post from the app task itself while the queue is full. + bool take_lost() { return gens_.take_lost(); } private: static void on_event(void* arg, const char* base, int32_t id, void* data); + // Sends run on their own task, so a slow or dead link never stalls the app + // task (and with it the screen): send_*() only queue a copy. + static void tx_task(void* arg); + bool enqueue(uint8_t opcode, const void* data, size_t len); esp_websocket_client_handle_t client_ = nullptr; - std::atomic generation_{0}; + QueueHandle_t tx_ = nullptr; + hg::ws::Generations gens_; std::string url_, subprotocol_; - std::string rx_; // fragment reassembly (WebSocket task only) + // Fragment reassembly, WebSocket task only: an old client's task can still be + // delivering a frame after close(), so connect() leaves this alone. Each + // message's first frame clears it. + std::string rx_; uint8_t rx_opcode_ = 0; }; diff --git a/firmware/esp32/main/port_ws.cpp b/firmware/esp32/main/port_ws.cpp index 1953527a..f7cf9280 100644 --- a/firmware/esp32/main/port_ws.cpp +++ b/firmware/esp32/main/port_ws.cpp @@ -1,10 +1,13 @@ -// WebSocket transport over esp_websocket_client. Runs on the client's own -// task; complete messages are posted to the app task as events. +// WebSocket transport over esp_websocket_client. Receiving runs on the client's +// own task and sending on hg-ws-tx; complete messages are posted to the app +// task as events. #include "port.hpp" // first: pulls in FreeRTOS.h ahead of task.h/queue.h +#include #include #include "esp_crt_bundle.h" +#include "esp_heap_caps.h" #include "esp_log.h" namespace hgp { @@ -14,10 +17,42 @@ const char* TAG = "hg.ws"; constexpr int kBufferSize = 4096; constexpr size_t kMaxMessage = 512 * 1024; // images arrive in 4 KB chunks; this only bounds JSON +using hg::ws::TxItem; + +// The client calls hg::ws::process() makes, each with a time limit. +struct EspOps { + static esp_websocket_client_handle_t h(void* c) { return static_cast(c); } + bool is_connected(void* c) { return esp_websocket_client_is_connected(h(c)); } + int send(void* c, uint8_t opcode, const uint8_t* data, size_t len, uint32_t timeout_ms) { + const auto* p = reinterpret_cast(data); + return opcode == hg::ws::kBinary ? esp_websocket_client_send_bin(h(c), p, static_cast(len), pdMS_TO_TICKS(timeout_ms)) + : esp_websocket_client_send_text(h(c), p, static_cast(len), pdMS_TO_TICKS(timeout_ms)); + } + // Not esp_websocket_client_close(): it ignores its timeout for the close frame. + void send_close(void* c, uint32_t timeout_ms) { + esp_websocket_client_send_with_opcode(h(c), WS_TRANSPORT_OPCODES_CLOSE, nullptr, 0, pdMS_TO_TICKS(timeout_ms)); + } + void destroy(void* c) { esp_websocket_client_destroy(h(c)); } +}; + +uint32_t generation_of(const esp_websocket_event_data_t* d) { + return static_cast(reinterpret_cast(d->user_context)); +} + } // namespace void WsTransport::connect(const std::string& url, const std::string& subprotocol) { close(); + if (!tx_) { + tx_ = xQueueCreate(hg::ws::kTxMessages + hg::ws::kTxReserved, sizeof(TxItem)); + if (!tx_ || xTaskCreate(&WsTransport::tx_task, "hg-ws-tx", 6144, this, 5, nullptr) != pdPASS) { + ESP_LOGE(TAG, "can't start the send task"); + if (tx_) vQueueDelete(tx_); + tx_ = nullptr; + events::post(EventType::WsClosed, "client init failed", 18, gens_.current()); + return; + } + } url_ = url; subprotocol_ = subprotocol; esp_websocket_client_config_t cfg = {}; @@ -28,53 +63,77 @@ void WsTransport::connect(const std::string& url, const std::string& subprotocol cfg.disable_auto_reconnect = true; // hg::App owns retry and backoff cfg.network_timeout_ms = 10000; cfg.ping_interval_sec = 0; // the protocol has its own heartbeat + // Every event names the connection it belongs to: a client still being torn + // down on the send task must not be mistaken for the current one. + cfg.user_context = reinterpret_cast(static_cast(gens_.current())); if (url_.rfind("wss://", 0) == 0) cfg.crt_bundle_attach = esp_crt_bundle_attach; client_ = esp_websocket_client_init(&cfg); if (!client_) { - events::post(EventType::WsClosed, "client init failed", 18, generation_.load()); + events::post(EventType::WsClosed, "client init failed", 18, gens_.current()); return; } esp_websocket_register_events(client_, WEBSOCKET_EVENT_ANY, &WsTransport::on_event, this); - rx_.clear(); if (esp_websocket_client_start(client_) != ESP_OK) { - events::post(EventType::WsClosed, "client start failed", 19, generation_.load()); + events::post(EventType::WsClosed, "client start failed", 19, gens_.current()); } } void WsTransport::close() { - ++generation_; // anything still queued from the old connection is now stale + gens_.leave(); // anything still queued from the old connection is now stale if (!client_) return; - if (esp_websocket_client_is_connected(client_)) esp_websocket_client_close(client_, pdMS_TO_TICKS(1000)); - esp_websocket_client_destroy(client_); + // The send task closes and destroys the client after whatever it is sending: + // on a dead link that can take seconds, and the app task must not wait. + TxItem item{client_, 0, hg::ws::kClose, nullptr, 0}; client_ = nullptr; + xQueueSend(tx_, &item, portMAX_DELAY); // the reserved slots keep room for this } -bool WsTransport::send_text(std::string_view text) { +bool WsTransport::enqueue(uint8_t opcode, const void* data, size_t len) { if (!client_ || !esp_websocket_client_is_connected(client_)) return false; - int n = esp_websocket_client_send_text(client_, text.data(), static_cast(text.size()), pdMS_TO_TICKS(2000)); - return n == static_cast(text.size()); + if (!hg::ws::may_queue_message(uxQueueSpacesAvailable(tx_))) return false; // the link isn't keeping up + auto* copy = static_cast(heap_caps_malloc(len ? len : 1, MALLOC_CAP_SPIRAM | MALLOC_CAP_8BIT)); + if (!copy) copy = static_cast(malloc(len ? len : 1)); + if (!copy) return false; + if (len) std::memcpy(copy, data, len); + TxItem item{client_, gens_.current(), opcode, copy, len}; + if (xQueueSend(tx_, &item, 0) != pdTRUE) { + free(copy); + return false; + } + return true; } -bool WsTransport::send_binary(const uint8_t* data, size_t len) { - if (!client_ || !esp_websocket_client_is_connected(client_)) return false; - int n = esp_websocket_client_send_bin(client_, reinterpret_cast(data), static_cast(len), - pdMS_TO_TICKS(2000)); - return n == static_cast(len); +bool WsTransport::send_text(std::string_view text) { return enqueue(hg::ws::kText, text.data(), text.size()); } + +bool WsTransport::send_binary(const uint8_t* data, size_t len) { return enqueue(hg::ws::kBinary, data, len); } + +void WsTransport::tx_task(void* arg) { + auto* self = static_cast(arg); + EspOps ops; + TxItem item; + for (;;) { + if (xQueueReceive(self->tx_, &item, portMAX_DELAY) != pdTRUE) continue; + if (hg::ws::process(item, self->gens_.current(), ops) == hg::ws::Outcome::Failed) ESP_LOGW(TAG, "send failed"); + free(const_cast(item.data)); + } } void WsTransport::on_event(void* arg, const char*, int32_t id, void* event_data) { auto* self = static_cast(arg); auto* d = static_cast(event_data); - const uint32_t gen = self->generation_.load(); + const uint32_t gen = generation_of(d); + if (!self->gens_.is_current(gen)) return; // a client being torn down switch (id) { case WEBSOCKET_EVENT_CONNECTED: events::post(EventType::WsOpen, nullptr, 0, gen); break; case WEBSOCKET_EVENT_DISCONNECTED: case WEBSOCKET_EVENT_CLOSED: + self->gens_.lost(gen); events::post(EventType::WsClosed, "disconnected", 12, gen); break; case WEBSOCKET_EVENT_ERROR: + self->gens_.lost(gen); events::post(EventType::WsClosed, "connection error", 16, gen); break; case WEBSOCKET_EVENT_DATA: { diff --git a/firmware/tests/test_ws_link.cpp b/firmware/tests/test_ws_link.cpp new file mode 100644 index 00000000..c01c9cf6 --- /dev/null +++ b/firmware/tests/test_ws_link.cpp @@ -0,0 +1,181 @@ +#include +#include +#include +#include + +#include "check.hpp" +#include "ws_link.hpp" + +namespace { + +using hg::ws::Generations; +using hg::ws::Outcome; +using hg::ws::TxItem; + +// Records the client calls hg::ws::process() makes. +struct FakeOps { + std::vector calls; + std::vector connected; + int short_by = 0; // a send that writes this many bytes too few + + static std::string name(void* c) { return std::string(1, *static_cast(c)); } + bool is_connected(void* c) { + for (void* x : connected) + if (x == c) return true; + return false; + } + int send(void* c, uint8_t opcode, const uint8_t*, size_t len, uint32_t timeout_ms) { + calls.push_back(std::string(opcode == hg::ws::kBinary ? "binary " : "text ") + name(c) + " " + + std::to_string(len) + " " + std::to_string(timeout_ms) + "ms"); + return static_cast(len) - short_by; + } + void send_close(void* c, uint32_t timeout_ms) { + calls.push_back("close frame " + name(c) + " " + std::to_string(timeout_ms) + "ms"); + } + void destroy(void* c) { calls.push_back("destroy " + name(c)); } +}; + +char client_a = 'A', client_b = 'B'; +const uint8_t kHello[] = {'h', 'e', 'l', 'l', 'o'}; + +TxItem message(void* client, uint32_t generation, uint8_t opcode = hg::ws::kText) { + return TxItem{client, generation, opcode, kHello, sizeof(kHello)}; +} +TxItem close_item(void* client) { return TxItem{client, 0, hg::ws::kClose, nullptr, 0}; } + +} // namespace + +// --- Which connection is current, and whether it was lost ------------------- + +TEST("ws: nothing is lost before the first connection") { + Generations g; + CHECK(!g.take_lost()); +} + +TEST("ws: events from a connection the app has left are ignored") { + Generations g; + g.leave(); + const uint32_t first = g.current(); + g.leave(); + CHECK(!g.is_current(first)); + CHECK(g.is_current(g.current())); +} + +TEST("ws: a lost connection is reported once") { + Generations g; + g.leave(); + g.lost(g.current()); + CHECK(g.take_lost()); + CHECK(!g.take_lost()); +} + +TEST("ws: a loss the app already acted on is not reported again") { + Generations g; + g.leave(); + g.lost(g.current()); + g.leave(); // the app closed that connection before the main loop looked + CHECK(!g.take_lost()); +} + +TEST("ws: the old client's loss doesn't end the new connection") { + Generations g; + g.leave(); + const uint32_t old_conn = g.current(); + g.leave(); // reconnecting: the old client is torn down in the background + g.lost(old_conn); + CHECK(!g.take_lost()); + g.lost(g.current()); + CHECK(g.take_lost()); +} + +// --- The send queue ---------------------------------------------------------- + +TEST("ws: messages leave the reserved slots free") { + CHECK(hg::ws::may_queue_message(hg::ws::kTxReserved + 1)); + CHECK(!hg::ws::may_queue_message(hg::ws::kTxReserved)); + CHECK(!hg::ws::may_queue_message(0)); +} + +TEST("ws: a full queue of messages still has room to close") { + size_t free_slots = hg::ws::kTxMessages + hg::ws::kTxReserved, queued = 0; + while (hg::ws::may_queue_message(free_slots)) { + --free_slots; + ++queued; + } + CHECK_EQ(queued, hg::ws::kTxMessages); + CHECK(free_slots >= 1u); // close() waits for a slot; this one is always there +} + +TEST("ws: a message for the current connection is sent with a time limit") { + FakeOps ops; + ops.connected = {&client_a}; + CHECK(hg::ws::process(message(&client_a, 1), 1, ops) == Outcome::Sent); + CHECK(hg::ws::process(message(&client_a, 1, hg::ws::kBinary), 1, ops) == Outcome::Sent); + CHECK_EQ(ops.calls.size(), 2u); + CHECK_EQ(ops.calls[0], std::string("text A 5 2000ms")); + CHECK_EQ(ops.calls[1], std::string("binary A 5 2000ms")); +} + +TEST("ws: a short send is reported as failed") { + FakeOps ops; + ops.connected = {&client_a}; + ops.short_by = 2; + CHECK(hg::ws::process(message(&client_a, 1), 1, ops) == Outcome::Failed); +} + +TEST("ws: messages for a connection the app has left are dropped unsent") { + FakeOps ops; + ops.connected = {&client_a}; + CHECK(hg::ws::process(message(&client_a, 1), 2, ops) == Outcome::Dropped); + CHECK(ops.calls.empty()); +} + +TEST("ws: messages for a client that lost its link are dropped unsent") { + FakeOps ops; // client A is no longer connected + CHECK(hg::ws::process(message(&client_a, 1), 1, ops) == Outcome::Dropped); + CHECK(ops.calls.empty()); +} + +TEST("ws: the old client is closed after its queued messages, and the new one carries on") { + FakeOps ops; + ops.connected = {&client_a, &client_b}; + Generations g; + g.leave(); + std::deque queue; + queue.push_back(message(&client_a, g.current())); + g.leave(); // the app reconnects: close(A), then messages for B + queue.push_back(close_item(&client_a)); + queue.push_back(message(&client_b, g.current())); + std::vector outcomes; + for (const TxItem& item : queue) outcomes.push_back(hg::ws::process(item, g.current(), ops)); + CHECK(outcomes[0] == Outcome::Dropped); // A's message is stale by the time it is sent + CHECK(outcomes[1] == Outcome::Destroyed); + CHECK(outcomes[2] == Outcome::Sent); + CHECK_EQ(ops.calls.size(), 3u); + CHECK_EQ(ops.calls[0], std::string("close frame A 1000ms")); + CHECK_EQ(ops.calls[1], std::string("destroy A")); + CHECK_EQ(ops.calls[2], std::string("text B 5 2000ms")); +} + +// --- Closing never waits without a limit (the Wi-Fi-loss freeze) ------------- + +TEST("ws: closing sends the close frame with a time limit, then destroys the client") { + FakeOps ops; + ops.connected = {&client_a}; + CHECK(hg::ws::process(close_item(&client_a), 7, ops) == Outcome::Destroyed); + CHECK_EQ(ops.calls.size(), 2u); + CHECK_EQ(ops.calls[0], std::string("close frame A 1000ms")); + CHECK_EQ(ops.calls[1], std::string("destroy A")); +} + +TEST("ws: closing a client whose link is gone skips the close frame") { + FakeOps ops; // Wi-Fi lost: nothing could be written anyway + CHECK(hg::ws::process(close_item(&client_a), 7, ops) == Outcome::Destroyed); + CHECK_EQ(ops.calls.size(), 1u); + CHECK_EQ(ops.calls[0], std::string("destroy A")); +} + +TEST("ws: every network wait has a limit of a few seconds at most") { + CHECK(hg::ws::kSendTimeoutMs > 0 && hg::ws::kSendTimeoutMs <= 5000); + CHECK(hg::ws::kCloseTimeoutMs > 0 && hg::ws::kCloseTimeoutMs <= 5000); +} diff --git a/tests/test_firmware_source.py b/tests/test_firmware_source.py new file mode 100644 index 00000000..5d8a6956 --- /dev/null +++ b/tests/test_firmware_source.py @@ -0,0 +1,32 @@ +"""The firmware's source: patterns that broke the device and must not come back.""" + +from __future__ import annotations + +import re + +from conftest import REPO + +FIRMWARE = REPO / "firmware" + + +def _code_lines(path): + """Each line without its // comment, with its number.""" + for number, line in enumerate(path.read_text(encoding="utf-8").splitlines(), start=1): + yield number, line.split("//", 1)[0] + + +def test_firmware_never_calls_the_websocket_close_that_ignores_its_timeout(): + # esp_websocket_client_close() sends its close frame with no time limit (esp_websocket_client + # 1.8.0), whatever timeout it is given. With Wi-Fi gone that write never finishes, and the app task + # that called it stopped drawing, reading touch and retrying Wi-Fi until a hard reset. Closing goes + # through hg::ws::process() (drivers/ws_link.hpp), which limits every wait. + call = re.compile(r"\besp_websocket_client_close(_with_[a-z_]+)?\s*\(") + found = [] + for path in sorted(FIRMWARE.rglob("*")): + if path.suffix not in (".c", ".cpp", ".h", ".hpp") or "managed_components" in path.parts \ + or ".pio" in path.parts: + continue + for number, code in _code_lines(path): + if call.search(code): + found.append(f"{path.relative_to(REPO)}:{number}: {code.strip()}") + assert not found, "\n".join(found)