From e742e64fb044fb1c8d919886c466c4a49188aec8 Mon Sep 17 00:00:00 2001 From: Keith Burzinski Date: Wed, 30 Sep 2026 12:33:17 -0500 Subject: [PATCH] [serial_proxy] Never let a client write stall the main loop past a fixed budget (#19923) --- .../components/serial_proxy/serial_proxy.cpp | 40 +++++++++++++++++++ .../components/serial_proxy/serial_proxy.h | 17 ++++++++ esphome/components/uart/__init__.py | 6 +++ esphome/components/uart/uart_component.h | 5 +++ .../uart/uart_component_esp8266.cpp | 6 +++ .../components/uart/uart_component_esp8266.h | 1 + .../uart/uart_component_esp_idf.cpp | 26 +++++++++--- .../components/uart/uart_component_esp_idf.h | 6 +++ tests/components/uart/test.esp32-idf.yaml | 1 + 9 files changed, 103 insertions(+), 5 deletions(-) diff --git a/esphome/components/serial_proxy/serial_proxy.cpp b/esphome/components/serial_proxy/serial_proxy.cpp index 129745c1c9..729ec40256 100644 --- a/esphome/components/serial_proxy/serial_proxy.cpp +++ b/esphome/components/serial_proxy/serial_proxy.cpp @@ -2,8 +2,10 @@ #ifdef USE_SERIAL_PROXY +#include "esphome/core/application.h" #include "esphome/core/log.h" +#include #include #include "esphome/core/util.h" @@ -16,6 +18,9 @@ namespace esphome::serial_proxy { static const char *const TAG = "serial_proxy"; +uint32_t SerialProxy::stall_loop_time = 0; +uint32_t SerialProxy::stall_spent_ms = 0; + void SerialProxy::setup() { // Set up modem control pins if configured if (this->rts_pin_ != nullptr) { @@ -293,6 +298,35 @@ void SerialProxy::write_from_client(api::APIConnection *api_connection, const ui #endif if (data == nullptr || len == 0) return; + // Whatever the driver cannot buffer stalls the main loop for its wire time. At high baud + // rates that is brief and losing nothing is worth it; at low ones it would trip the + // watchdog, so cap the stall and drop the rest. The cap covers the whole loop pass: + // several writes can arrive in one, and each alone might stay under it. + const size_t free = this->parent_->available_for_write(); + bool trimmed = false; + if (len > free) { + const uint32_t loop_time = App.get_loop_component_start_time(); + if (loop_time != stall_loop_time) { + stall_loop_time = loop_time; + stall_spent_ms = 0; + } + const uint32_t stall_ms = this->wire_time_ms_(len - free); + trimmed = stall_spent_ms + stall_ms > SERIAL_PROXY_MAX_WRITE_STALL_MS; + if (trimmed && !this->trim_warned_) { + ESP_LOGW(TAG, + "TX buffer full on serial proxy [%" PRIu32 "]: dropping %zu of %zu bytes (would stall %" PRIu32 + " ms at %" PRIu32 " baud); raise the UART tx_buffer_size or pace writes", + this->instance_index_, len - free, len, stall_ms, this->parent_->get_baud_rate()); + } + if (trimmed) { + len = free; + } else { + stall_spent_ms += stall_ms; + } + } + this->trim_warned_ = trimmed; + if (len == 0) + return; this->write_array(data, len); #ifdef USE_SERIAL_PROXY_TAP @@ -303,6 +337,12 @@ void SerialProxy::write_from_client(api::APIConnection *api_connection, const ui #endif } +uint32_t SerialProxy::wire_time_ms_(size_t bytes) const { + const uint32_t bits_per_byte = 1 + this->parent_->get_data_bits() + this->parent_->get_stop_bits() + + (this->parent_->get_parity() != uart::UART_CONFIG_PARITY_NONE ? 1 : 0); + return static_cast(bytes) * bits_per_byte * 1000 / std::max(this->parent_->get_baud_rate(), 1); +} + SerialProxyResult SerialProxy::set_modem_pins(api::APIConnection *api_connection, uint32_t line_states) { #ifdef USE_API if (!this->is_subscriber_(api_connection)) { diff --git a/esphome/components/serial_proxy/serial_proxy.h b/esphome/components/serial_proxy/serial_proxy.h index e3f4264cfa..df775d4f63 100644 --- a/esphome/components/serial_proxy/serial_proxy.h +++ b/esphome/components/serial_proxy/serial_proxy.h @@ -53,6 +53,11 @@ enum class SerialProxyResult : uint8_t { /// Maximum bytes to read from UART in a single loop iteration inline constexpr size_t SERIAL_PROXY_MAX_READ_SIZE = 256; +/// Longest main-loop stall client writes may cause per loop pass, shared by every instance; +/// bytes the UART cannot buffer within it are dropped. Well under the shortest watchdog +/// timeout, since the API hands the proxies up to ten writes in one pass. +inline constexpr uint32_t SERIAL_PROXY_MAX_WRITE_STALL_MS = 1000; + #ifdef USE_SERIAL_PROXY_TAP /// Observes a port's traffic without owning it, and may inject bytes of its own. /// @@ -204,6 +209,9 @@ class SerialProxy final : public uart::UARTDevice, public Component { bool is_subscriber_(api::APIConnection *api_connection) const { return this->api_connection_ == api_connection; } #endif + /// Time the wire needs for the given number of bytes at the current framing + uint32_t wire_time_ms_(size_t bytes) const; + #ifdef USE_SERIAL_PROXY_TAP /// Return the port to RAW when a subscriber goes away, so the mode never outlives it void reset_mode_(); @@ -221,6 +229,12 @@ class SerialProxy final : public uart::UARTDevice, public Component { /// Instance index for identifying this proxy in API messages uint32_t instance_index_{0}; + /// Stall spent by writes in the current loop pass, keyed by the pass's cached start time. + /// Static on purpose: there is one main loop, and writes to different ports that arrive in + /// the same pass all stall it, so the budget is one per device rather than one per port + static uint32_t stall_loop_time; + static uint32_t stall_spent_ms; + /// Subscribed API client (only one allowed at a time) api::APIConnection *api_connection_{nullptr}; @@ -248,6 +262,9 @@ class SerialProxy final : public uart::UARTDevice, public Component { bool rts_state_{false}; bool dtr_state_{false}; + /// Set while writes are being trimmed, so a client streaming into a slow port warns once + bool trim_warned_{false}; + #ifdef USE_SERIAL_PROXY_TAP SerialProxyTap *tap_{nullptr}; #endif diff --git a/esphome/components/uart/__init__.py b/esphome/components/uart/__init__.py index 419659598c..598e3df168 100644 --- a/esphome/components/uart/__init__.py +++ b/esphome/components/uart/__init__.py @@ -46,6 +46,7 @@ from esphome.const import ( CONF_SEQUENCE, CONF_TIMEOUT, CONF_TRIGGER_ID, + CONF_TX_BUFFER_SIZE, CONF_TX_PIN, CONF_UART_ID, PLATFORM_HOST, @@ -297,6 +298,9 @@ CONFIG_SCHEMA = cv.All( ), cv.Optional(CONF_PORT): cv.All(validate_port, cv.only_on(PLATFORM_HOST)), cv.Optional(CONF_RX_BUFFER_SIZE, default=256): cv.validate_bytes, + cv.Optional(CONF_TX_BUFFER_SIZE): cv.All( + cv.only_on_esp32, cv.validate_bytes, cv.int_range(min=129) + ), cv.Optional(CONF_RX_FULL_THRESHOLD): cv.All( cv.only_on_esp32, cv.validate_bytes, cv.int_range(min=1, max=120) ), @@ -391,6 +395,8 @@ async def to_code(config): cg.add(var.set_rx_timeout(config[CONF_RX_TIMEOUT])) if CONF_FLUSH_TIMEOUT in config: cg.add(var.set_flush_timeout(config[CONF_FLUSH_TIMEOUT])) + if (tx_buffer_size := config.get(CONF_TX_BUFFER_SIZE)) is not None: + cg.add(var.set_tx_buffer_size(tx_buffer_size)) # The member already defaults to UART_SCLK_DEFAULT, so only emit a real choice if (clock_source := config.get(CONF_CLOCK_SOURCE, "DEFAULT")) != "DEFAULT": cg.add(var.set_clock_source(UART_CLOCK_SOURCES[clock_source])) diff --git a/esphome/components/uart/uart_component.h b/esphome/components/uart/uart_component.h index 3e52531791..4269ef0b16 100644 --- a/esphome/components/uart/uart_component.h +++ b/esphome/components/uart/uart_component.h @@ -1,6 +1,7 @@ #pragma once #include +#include #include #include "esphome/core/defines.h" #include "esphome/core/component.h" @@ -81,6 +82,10 @@ class UARTComponent { // @return Number of available bytes. virtual size_t available() = 0; + // Returns how many bytes write_array() accepts right now without blocking. + // Platforms that cannot tell return SIZE_MAX: write_array() takes everything and may block. + virtual size_t available_for_write() { return SIZE_MAX; } + // Pure virtual method to block until all bytes have been written to the UART bus. // @return UARTFlushResult indicating whether the flush was confirmed, timed out, failed, or assumed successful. virtual UARTFlushResult flush() = 0; diff --git a/esphome/components/uart/uart_component_esp8266.cpp b/esphome/components/uart/uart_component_esp8266.cpp index 1d492fce47..0f6d7f6577 100644 --- a/esphome/components/uart/uart_component_esp8266.cpp +++ b/esphome/components/uart/uart_component_esp8266.cpp @@ -221,6 +221,12 @@ size_t ESP8266UartComponent::available() { return this->sw_serial_->available(); } } +size_t ESP8266UartComponent::available_for_write() { + if (this->hw_serial_ != nullptr) { + return this->hw_serial_->availableForWrite(); + } + return SIZE_MAX; // software serial bit-bangs each byte synchronously; there is no buffer to fill +} UARTFlushResult ESP8266UartComponent::flush() { ESP_LOGVV(TAG, " Flushing"); if (this->hw_serial_ != nullptr) { diff --git a/esphome/components/uart/uart_component_esp8266.h b/esphome/components/uart/uart_component_esp8266.h index 54bcad3993..0098e752b2 100644 --- a/esphome/components/uart/uart_component_esp8266.h +++ b/esphome/components/uart/uart_component_esp8266.h @@ -58,6 +58,7 @@ class ESP8266UartComponent final : public UARTComponent, public Component { bool read_array(uint8_t *data, size_t len) override; size_t available() override; + size_t available_for_write() override; UARTFlushResult flush() override; uint32_t get_config(); diff --git a/esphome/components/uart/uart_component_esp_idf.cpp b/esphome/components/uart/uart_component_esp_idf.cpp index a052c6015e..5d90fc1c8a 100644 --- a/esphome/components/uart/uart_component_esp_idf.cpp +++ b/esphome/components/uart/uart_component_esp_idf.cpp @@ -7,6 +7,7 @@ #include "esphome/core/log.h" #include "esphome/core/gpio.h" #include "driver/gpio.h" +#include "hal/uart_ll.h" #include "esp_private/gpio.h" #include "soc/gpio_num.h" #include "soc/soc_caps.h" @@ -162,11 +163,11 @@ void IDFUARTComponent::load_settings(bool dump_config) { } err = uart_driver_install(this->uart_num_, // UART number this->rx_buffer_size_, // RX ring buffer size - 0, // TX ring buffer size. If zero, driver will not use a TX buffer and TX function will - // block task until all data has been sent out - 0, // event queue size/depth - nullptr, // event queue - 0 // Flags used to allocate the interrupt + this->tx_buffer_size_, // TX ring buffer size; 0 makes uart_write_bytes() block until + // the FIFO has taken everything + 0, // event queue size/depth + nullptr, // event queue + 0 // Flags used to allocate the interrupt ); if (err != ESP_OK) { ESP_LOGW(TAG, "uart_driver_install failed: %s", esp_err_to_name(err)); @@ -356,6 +357,9 @@ void IDFUARTComponent::dump_config() { " RX Timeout: %u", this->rx_buffer_size_, this->rx_full_threshold_, this->rx_timeout_); } + if (this->tx_buffer_size_ > 0) { + ESP_LOGCONFIG(TAG, " TX Buffer Size: %zu", this->tx_buffer_size_); + } if (this->flush_timeout_ms_ > 0) { ESP_LOGCONFIG(TAG, " Flush Timeout: %" PRIu32 " ms", this->flush_timeout_ms_); } @@ -396,6 +400,18 @@ void IDFUARTComponent::set_rx_timeout(size_t rx_timeout) { this->rx_timeout_ = rx_timeout; } +size_t IDFUARTComponent::available_for_write() { + if (this->uart_num_ == UART_NUM_MAX || !uart_is_driver_installed(this->uart_num_)) + return 0; + if (this->tx_buffer_size_ == 0) { + return uart_ll_get_txfifo_len(UART_LL_GET_HW(this->uart_num_)); + } + // The driver's figure already deducts its ring item headers, so a write of this size fits + size_t free = 0; + uart_get_tx_buffer_free_size(this->uart_num_, &free); + return free; +} + void IDFUARTComponent::write_array(const uint8_t *data, size_t len) { int32_t write_len = uart_write_bytes(this->uart_num_, data, len); if (write_len != (int32_t) len) { diff --git a/esphome/components/uart/uart_component_esp_idf.h b/esphome/components/uart/uart_component_esp_idf.h index 8684937b03..7c93b74fbd 100644 --- a/esphome/components/uart/uart_component_esp_idf.h +++ b/esphome/components/uart/uart_component_esp_idf.h @@ -33,10 +33,15 @@ class IDFUARTComponent final : public UARTComponent, public Component { bool read_array(uint8_t *data, size_t len) override; size_t available() override; + size_t available_for_write() override; UARTFlushResult flush() override; void set_flush_timeout(uint32_t flush_timeout_ms) override { this->flush_timeout_ms_ = flush_timeout_ms; } + /// TX ring buffer size for the driver; 0 leaves TX unbuffered so write_array() blocks until + /// the FIFO has taken everything. + void set_tx_buffer_size(size_t tx_buffer_size) { this->tx_buffer_size_ = tx_buffer_size; } + void set_clock_source(uart_sclk_t clock_source) { this->clock_source_ = static_cast(clock_source); } uint8_t get_hw_serial_number() { return this->uart_num_; } @@ -109,6 +114,7 @@ class IDFUARTComponent final : public UARTComponent, public Component { uint8_t peek_byte_{0}; uint8_t clock_source_{UART_SCLK_DEFAULT}; ///< uart_sclk_t stored in a byte; the IDF values are all small. uint32_t flush_timeout_ms_{0}; ///< 0 means wait indefinitely (portMAX_DELAY). + size_t tx_buffer_size_{0}; #ifdef USE_UART_WAKE_LOOP_ON_RX // ISR callback for UART RX data notification — wakes the main loop directly. diff --git a/tests/components/uart/test.esp32-idf.yaml b/tests/components/uart/test.esp32-idf.yaml index 333576e1c8..9550de911b 100644 --- a/tests/components/uart/test.esp32-idf.yaml +++ b/tests/components/uart/test.esp32-idf.yaml @@ -20,6 +20,7 @@ uart: baud_rate: 9600 data_bits: 8 rx_buffer_size: 512 + tx_buffer_size: 512 rx_full_threshold: 10 rx_timeout: 1 parity: EVEN