diff --git a/CODEOWNERS b/CODEOWNERS index 5f89315d23..b453cf4538 100644 --- a/CODEOWNERS +++ b/CODEOWNERS @@ -602,6 +602,7 @@ esphome/components/uart/button/* @ssieb esphome/components/uart/event/* @eoasmxd esphome/components/uart/packet_transport/* @clydebarrow esphome/components/uart_mux/* @kbx81 +esphome/components/uart_tcp/* @Bascht74 esphome/components/udp/* @clydebarrow esphome/components/ufire_ec/* @pvizeli esphome/components/ufire_ise/* @pvizeli diff --git a/esphome/components/uart_tcp/__init__.py b/esphome/components/uart_tcp/__init__.py new file mode 100644 index 0000000000..e7a85a0f2d --- /dev/null +++ b/esphome/components/uart_tcp/__init__.py @@ -0,0 +1,71 @@ +import esphome.codegen as cg +from esphome.components import binary_sensor, socket, uart +from esphome.components.const import CONF_HOST, CONF_RECONNECT_INTERVAL, CONF_ROLE +import esphome.config_validation as cv +from esphome.const import ( + CONF_ID, + CONF_PORT, + CONF_UART_ID, + DEVICE_CLASS_CONNECTIVITY, + ENTITY_CATEGORY_DIAGNOSTIC, +) +from esphome.types import ConfigType + +CODEOWNERS = ["@Bascht74"] +DEPENDENCIES = ["network", "uart"] +AUTO_LOAD = ["binary_sensor", "socket"] +MULTI_CONF = True + +uart_tcp_ns = cg.esphome_ns.namespace("uart_tcp") +UartTcp = uart_tcp_ns.class_("UartTcp", cg.Component, uart.UARTDevice) + +CONF_CONNECTED = "connected" + + +def _consume_sockets(config: ConfigType) -> ConfigType: + if config[CONF_ROLE] == "server": + socket.consume_sockets(1, "uart_tcp", socket.SocketType.TCP_LISTEN)(config) + return socket.consume_sockets(1, "uart_tcp")(config) + + +BASE_SCHEMA = cv.Schema( + { + cv.GenerateID(): cv.declare_id(UartTcp), + cv.Required(CONF_UART_ID): cv.use_id(uart.UARTComponent), + cv.Required(CONF_PORT): cv.port, + cv.Optional( + CONF_RECONNECT_INTERVAL, default="5s" + ): cv.positive_time_period_milliseconds, + cv.Optional(CONF_CONNECTED): binary_sensor.binary_sensor_schema( + device_class=DEVICE_CLASS_CONNECTIVITY, + entity_category=ENTITY_CATEGORY_DIAGNOSTIC, + ), + } +).extend(cv.COMPONENT_SCHEMA) + +CONFIG_SCHEMA = cv.All( + cv.typed_schema( + { + "client": BASE_SCHEMA.extend({cv.Required(CONF_HOST): cv.string}), + "server": BASE_SCHEMA, + }, + key=CONF_ROLE, + default_type="client", + lower=True, + ), + _consume_sockets, +) + + +async def to_code(config: ConfigType) -> None: + socket.require_tcp_client_link() + var = cg.new_Pvariable(config[CONF_ID]) + await cg.register_component(var, config) + await uart.register_uart_device(var, config) + cg.add(var.set_server(config[CONF_ROLE] == "server")) + cg.add(var.set_port(config[CONF_PORT])) + cg.add(var.set_reconnect_interval(config[CONF_RECONNECT_INTERVAL])) + if (host := config.get(CONF_HOST)) is not None: + cg.add(var.set_host(host)) + binary_sensors = binary_sensor.sub_binary_sensors(config) + await binary_sensors(CONF_CONNECTED, var.set_connected_sensor) diff --git a/esphome/components/uart_tcp/uart_tcp.cpp b/esphome/components/uart_tcp/uart_tcp.cpp new file mode 100644 index 0000000000..2531b73b20 --- /dev/null +++ b/esphome/components/uart_tcp/uart_tcp.cpp @@ -0,0 +1,170 @@ +#include "uart_tcp.h" + +#include "esphome/core/log.h" + +#include +#include +#include + +namespace esphome::uart_tcp { + +static const char *const TAG = "uart_tcp"; + +// One client at a time; a second connection waits in the stack until the first drops. +static constexpr int LISTEN_BACKLOG = 1; +// Bytes per 16 ms loop pass at 10 bits per byte: baud / 10 / 62.5. +static constexpr uint32_t BAUD_PACE_DIVISOR = 625; + +void UartTcp::setup() { + this->link_.begin(TAG); + if (this->connected_sensor_ != nullptr) { + this->connected_sensor_->publish_state(false); + } +} + +void UartTcp::dump_config() { + ESP_LOGCONFIG(TAG, + "UART TCP:\n" + " %s: %s:%u\n" + " Reconnect Interval: %" PRIu32 "ms", + this->server_ ? LOG_STR_LITERAL("Listen") : LOG_STR_LITERAL("Host"), + this->server_ ? LOG_STR_LITERAL("*") : this->link_.host(), this->link_.port(), + this->link_.reconnect_interval()); + LOG_BINARY_SENSOR(" ", "Connected", this->connected_sensor_); +} + +void UartTcp::on_shutdown() { + this->link_.close(); + this->listen_.reset(); +} + +void UartTcp::sync_link_() { + bool up = this->link_.connected(); + this->link_was_up_ = up; + if (up) { + // The driver kept whatever arrived while the link was down. + this->discard_uart_(); + } + if (this->connected_sensor_ != nullptr) { + this->connected_sensor_->publish_state(up); + } +} + +void UartTcp::try_listen_() { + this->listen_ = socket::socket_ip_loop_monitored(SOCK_STREAM, IPPROTO_TCP); + int err = errno; + if (this->listen_ != nullptr) { + int yes = 1; + this->listen_->setsockopt(SOL_SOCKET, SO_REUSEADDR, &yes, sizeof(yes)); + struct sockaddr_storage local; + socklen_t local_len = + socket::set_sockaddr_any(reinterpret_cast(&local), sizeof(local), this->link_.port()); + // A blocking listener would stall loop() inside accept(), so its + // setblocking result is part of the success condition. + if (this->listen_->setblocking(false) == 0 && local_len != 0 && + this->listen_->bind(reinterpret_cast(&local), local_len) == 0 && + this->listen_->listen(LISTEN_BACKLOG) == 0) { + ESP_LOGI(TAG, "Listening on %u", this->link_.port()); + return; + } + // Captured before reset(); the close inside can overwrite errno. + err = errno; + this->listen_.reset(); + } + ESP_LOGW(TAG, "Listen on %u failed: %d", this->link_.port(), err); + this->link_.note_attempt(); +} + +void UartTcp::accept_client_() { + auto client = this->listen_->accept_loop_monitored(nullptr, nullptr); + if (client == nullptr) { + // A reset during the handshake or a signal only affects that connection. + if (errno == EAGAIN || errno == EWOULDBLOCK || errno == ECONNABORTED || errno == EINTR) { + return; + } + // Rebuild the listener after the backoff instead of spinning on it. + int err = errno; + this->listen_.reset(); + ESP_LOGW(TAG, "Accept failed: %d", err); + this->link_.note_attempt(); + return; + } + this->link_.adopt(std::move(client)); + ESP_LOGI(TAG, "Client connected"); +} + +void UartTcp::read_socket_() { + // A hardware write blocks until the driver takes every byte. Leave what does + // not fit in the socket, so TCP flow control throttles the peer. + size_t room = this->parent_->available_for_write(); + if (room == SIZE_MAX) { + // Capacity unknown on this platform; pace to one loop pass of UART time + // (16 ms at 10 bits per byte) so a blocking write stays short. + room = std::max(1, this->parent_->get_baud_rate() / BAUD_PACE_DIVISOR); + } + if (room == 0) { + this->rx_pending_ = true; + return; + } + uint8_t tmp[READ_CHUNK]; + size_t want = std::min(room, sizeof(tmp)); + ssize_t count = this->link_.read(tmp, want); + if (count <= 0) { + // A dropped link (-1) is cleaned up by sync_link_() on the next loop. + if (count == 0) { + this->rx_pending_ = false; + } + return; + } + this->rx_pending_ = static_cast(count) == want; + this->write_array(tmp, static_cast(count)); +} + +void UartTcp::discard_uart_() { + // Drain exactly what was buffered while the link was down; later bytes are live. + uint8_t dump[32]; + size_t left = this->available(); + while (left != 0) { + size_t n = std::min(left, sizeof(dump)); + if (!this->read_array(dump, n)) { + return; + } + left -= n; + } +} + +void UartTcp::read_uart_() { + size_t want = std::min(this->available(), this->link_.tx_free()); + if (want != 0 && this->read_array(this->link_.tx_tail(), want)) { + this->link_.tx_commit(want); + } +} + +void UartTcp::loop() { + if (this->server_) { + if (this->listen_ == nullptr && !this->link_.in_backoff()) { + this->try_listen_(); + } + // link_was_up_ holds the accept until the previous drop's edge has run, + // so the sensor and the stale UART discard always see the disconnect. + if (this->listen_ != nullptr && !this->link_.connected() && !this->link_was_up_ && this->listen_->ready()) { + this->accept_client_(); + } + } else { + this->link_.poll(); + } + if (this->link_.connected() != this->link_was_up_) { + this->sync_link_(); + } + if (!this->link_was_up_) { + return; + } + if (this->rx_pending_ || this->link_.ready()) { + this->read_socket_(); + } + // UART bytes picked up here go out in the same pass. + this->read_uart_(); + this->link_.flush_tx(); +} + +} // namespace esphome::uart_tcp diff --git a/esphome/components/uart_tcp/uart_tcp.h b/esphome/components/uart_tcp/uart_tcp.h new file mode 100644 index 0000000000..2a59bf968f --- /dev/null +++ b/esphome/components/uart_tcp/uart_tcp.h @@ -0,0 +1,48 @@ +#pragma once + +#include "esphome/components/binary_sensor/binary_sensor.h" +#include "esphome/components/socket/tcp_client_link.h" +#include "esphome/components/uart/uart.h" +#include "esphome/core/component.h" + +#include +#include + +namespace esphome::uart_tcp { + +/// Copies raw bytes between one hardware UART and one TCP socket. +class UartTcp : public Component, public uart::UARTDevice { + public: + void set_host(const char *host) { this->link_.set_host(host); } + void set_port(uint16_t port) { this->link_.set_port(port); } + void set_server(bool server) { this->server_ = server; } + void set_reconnect_interval(uint32_t ms) { this->link_.set_reconnect_interval(ms); } + void set_connected_sensor(binary_sensor::BinarySensor *sensor) { this->connected_sensor_ = sensor; } + + void setup() override; + void loop() override; + void dump_config() override; + void on_shutdown() override; + float get_setup_priority() const override { return setup_priority::AFTER_WIFI; } + + protected: + void sync_link_(); + void try_listen_(); + void accept_client_(); + void read_socket_(); + void read_uart_(); + void discard_uart_(); + + static constexpr size_t READ_CHUNK = 128; + + socket::TcpClientLink link_; + std::unique_ptr listen_; + binary_sensor::BinarySensor *connected_sensor_{nullptr}; + bool server_{false}; + // The link state loop() saw last; edges clear the buffer and publish the sensor. + bool link_was_up_{false}; + // A read stopped before EAGAIN. ready() stays false until new data arrives. + bool rx_pending_{false}; +}; + +} // namespace esphome::uart_tcp diff --git a/tests/components/uart_tcp/common.yaml b/tests/components/uart_tcp/common.yaml new file mode 100644 index 0000000000..8010878e89 --- /dev/null +++ b/tests/components/uart_tcp/common.yaml @@ -0,0 +1,11 @@ +wifi: + ssid: MySSID + password: password1 + +uart_tcp: + - id: bridge + uart_id: uart_bus + role: server + port: 502 + connected: + name: UART TCP Connected diff --git a/tests/components/uart_tcp/test-client.esp32-idf.yaml b/tests/components/uart_tcp/test-client.esp32-idf.yaml new file mode 100644 index 0000000000..92b0213f74 --- /dev/null +++ b/tests/components/uart_tcp/test-client.esp32-idf.yaml @@ -0,0 +1,16 @@ +packages: + uart: !include ../../test_build_components/common/uart/esp32-idf.yaml + +wifi: + ssid: MySSID + password: password1 + +uart_tcp: + - id: bridge + uart_id: uart_bus + role: client + host: 192.0.2.20 + port: 502 + reconnect_interval: 10s + connected: + name: UART TCP Connected diff --git a/tests/components/uart_tcp/test.bk72xx-ard.yaml b/tests/components/uart_tcp/test.bk72xx-ard.yaml new file mode 100644 index 0000000000..719d654f14 --- /dev/null +++ b/tests/components/uart_tcp/test.bk72xx-ard.yaml @@ -0,0 +1,3 @@ +packages: + uart: !include ../../test_build_components/common/uart/bk72xx-ard.yaml + uart_tcp: !include common.yaml diff --git a/tests/components/uart_tcp/test.esp32-idf.yaml b/tests/components/uart_tcp/test.esp32-idf.yaml new file mode 100644 index 0000000000..8e7422f24b --- /dev/null +++ b/tests/components/uart_tcp/test.esp32-idf.yaml @@ -0,0 +1,3 @@ +packages: + uart: !include ../../test_build_components/common/uart/esp32-idf.yaml + uart_tcp: !include common.yaml diff --git a/tests/components/uart_tcp/test.esp8266-ard.yaml b/tests/components/uart_tcp/test.esp8266-ard.yaml new file mode 100644 index 0000000000..18ae2ef804 --- /dev/null +++ b/tests/components/uart_tcp/test.esp8266-ard.yaml @@ -0,0 +1,3 @@ +packages: + uart: !include ../../test_build_components/common/uart/esp8266-ard.yaml + uart_tcp: !include common.yaml diff --git a/tests/components/uart_tcp/test.host.yaml b/tests/components/uart_tcp/test.host.yaml new file mode 100644 index 0000000000..96a636fa98 --- /dev/null +++ b/tests/components/uart_tcp/test.host.yaml @@ -0,0 +1,12 @@ +uart: + - id: uart_bus + baud_rate: 9600 + port: /dev/ttyS0 + +uart_tcp: + - id: bridge + uart_id: uart_bus + host: 127.0.0.1 + port: 44502 + connected: + name: UART TCP Connected diff --git a/tests/components/uart_tcp/test.rp2040-ard.yaml b/tests/components/uart_tcp/test.rp2040-ard.yaml new file mode 100644 index 0000000000..5ece3bc3f4 --- /dev/null +++ b/tests/components/uart_tcp/test.rp2040-ard.yaml @@ -0,0 +1,3 @@ +packages: + uart: !include ../../test_build_components/common/uart/rp2040-ard.yaml + uart_tcp: !include common.yaml diff --git a/tests/integration/fixtures/uart_tcp_bridge.yaml b/tests/integration/fixtures/uart_tcp_bridge.yaml new file mode 100644 index 0000000000..c61add801f --- /dev/null +++ b/tests/integration/fixtures/uart_tcp_bridge.yaml @@ -0,0 +1,22 @@ +esphome: + name: uart-tcp-bridge-test + +host: + +api: + +logger: + level: INFO + +uart: + - id: uart_bus + baud_rate: 115200 + port: PTY_PATH + +uart_tcp: + - id: bridge + uart_id: uart_bus + role: server + port: 18126 + connected: + name: Bridge Connected diff --git a/tests/integration/test_uart_tcp_bridge.py b/tests/integration/test_uart_tcp_bridge.py new file mode 100644 index 0000000000..b98d3ffac9 --- /dev/null +++ b/tests/integration/test_uart_tcp_bridge.py @@ -0,0 +1,110 @@ +"""Integration test for the uart_tcp bridge on host. + +The UART bus is backed by a pty; pytest holds the controller side and connects +as the TCP client. Covers both transfer directions, the stale-byte discard +at every accept, and the drop plus client replacement path. +""" + +from __future__ import annotations + +import asyncio +import os +import pathlib + +import pytest + +from .log_utils import LineWaiter +from .types import APIClientConnectedFactory, RunCompiledFunction + + +@pytest.mark.asyncio +async def test_uart_tcp_bridge( + yaml_config: str, + run_compiled: RunCompiledFunction, + api_client_connected: APIClientConnectedFactory, + unused_tcp_port_factory, +) -> None: + server_port = unused_tcp_port_factory() + controller_fd, device_fd = os.openpty() + os.set_blocking(controller_fd, False) + # uart's validate_port wants a two segment device path; Linux ptys live at + # /dev/pts/N, so hand the config a /tmp symlink instead. + pty_link = f"/tmp/uart-tcp-pty-{os.getpid()}" + pathlib.Path(pty_link).symlink_to(os.ttyname(device_fd)) + yaml_config = yaml_config.replace("port: 18126", f"port: {server_port}") + yaml_config = yaml_config.replace("PTY_PATH", pty_link) + + lines = LineWaiter() + loop = asyncio.get_running_loop() + uart_rx = bytearray() + uart_rx_event = asyncio.Event() + + def on_controller_readable() -> None: + try: + chunk = os.read(controller_fd, 256) + except BlockingIOError: + return + if chunk: + uart_rx.extend(chunk) + uart_rx_event.set() + + async def read_uart(count: int, timeout: float = 10.0) -> bytes: + while len(uart_rx) < count: + uart_rx_event.clear() + await asyncio.wait_for(uart_rx_event.wait(), timeout) + data = bytes(uart_rx[:count]) + del uart_rx[:count] + return data + + async def wait_log_count(needle: str, count: int, timeout: float = 15.0) -> None: + async with asyncio.timeout(timeout): + while sum(needle in line for line in lines.lines) < count: + await asyncio.sleep(0.05) + + loop.add_reader(controller_fd, on_controller_readable) + try: + async with ( + run_compiled(yaml_config, line_callback=lines.callback), + api_client_connected() as client, + ): + device_info = await client.device_info() + assert device_info is not None + assert device_info.name == "uart-tcp-bridge-test" + await lines.wait_for("Listening on") + + # Bytes written before any client connects must never reach one. + os.write(controller_fd, b"STALE") + await asyncio.sleep(0.2) + + reader, writer = await asyncio.open_connection("127.0.0.1", server_port) + await wait_log_count("Client connected", 1) + os.write(controller_fd, b"live!") + assert await asyncio.wait_for(reader.readexactly(5), 10) == b"live!", ( + "First bytes to the client were not the live payload" + ) + writer.write(b"down1") + await writer.drain() + assert await read_uart(5) == b"down1" + + # Drop the client; bytes while no client is connected are discarded + # when the next one is accepted. + writer.close() + await lines.wait_for("Connection lost") + os.write(controller_fd, b"gap") + await asyncio.sleep(0.2) + + reader, writer = await asyncio.open_connection("127.0.0.1", server_port) + await wait_log_count("Client connected", 2) + os.write(controller_fd, b"live2") + assert await asyncio.wait_for(reader.readexactly(5), 10) == b"live2", ( + "Second client received stale bytes from the gap" + ) + writer.write(b"down2") + await writer.drain() + assert await read_uart(5) == b"down2" + writer.close() + finally: + loop.remove_reader(controller_fd) + os.close(controller_fd) + os.close(device_fd) + pathlib.Path(pty_link).unlink()