[serial_proxy] Never let a client write stall the main loop past a fixed budget (#19923)

This commit is contained in:
Keith Burzinski
2026-09-30 12:33:17 -05:00
committed by GitHub
parent 52079dd3f3
commit e742e64fb0
9 changed files with 103 additions and 5 deletions
@@ -2,8 +2,10 @@
#ifdef USE_SERIAL_PROXY
#include "esphome/core/application.h"
#include "esphome/core/log.h"
#include <algorithm>
#include <cinttypes>
#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<uint64_t>(bytes) * bits_per_byte * 1000 / std::max<uint32_t>(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)) {
@@ -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
+6
View File
@@ -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]))
+5
View File
@@ -1,6 +1,7 @@
#pragma once
#include <vector>
#include <cstdint>
#include <cstring>
#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;
@@ -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) {
@@ -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();
@@ -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) {
@@ -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<uint8_t>(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.
@@ -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