mirror of
https://github.com/esphome/esphome.git
synced 2026-10-07 03:16:37 +00:00
[uart_tcp] Add a hardware UART copied to one TCP socket (#19887)
Co-authored-by: pre-commit-ci-lite[bot] <117423508+pre-commit-ci-lite[bot]@users.noreply.github.com> Co-authored-by: J. Nick Koston <nick@home-assistant.io>
This commit is contained in:
co-authored by
pre-commit-ci-lite[bot]
J. Nick Koston
parent
6491eaeaae
commit
26398747a1
@@ -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
|
||||
|
||||
@@ -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)
|
||||
@@ -0,0 +1,170 @@
|
||||
#include "uart_tcp.h"
|
||||
|
||||
#include "esphome/core/log.h"
|
||||
|
||||
#include <algorithm>
|
||||
#include <cerrno>
|
||||
#include <cinttypes>
|
||||
|
||||
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<struct sockaddr *>(&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<struct sockaddr *>(&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<size_t>(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<size_t>(count) == want;
|
||||
this->write_array(tmp, static_cast<size_t>(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<size_t>(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
|
||||
@@ -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 <cstdint>
|
||||
#include <memory>
|
||||
|
||||
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<socket::ListenSocket> 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
|
||||
@@ -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
|
||||
@@ -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
|
||||
@@ -0,0 +1,3 @@
|
||||
packages:
|
||||
uart: !include ../../test_build_components/common/uart/bk72xx-ard.yaml
|
||||
uart_tcp: !include common.yaml
|
||||
@@ -0,0 +1,3 @@
|
||||
packages:
|
||||
uart: !include ../../test_build_components/common/uart/esp32-idf.yaml
|
||||
uart_tcp: !include common.yaml
|
||||
@@ -0,0 +1,3 @@
|
||||
packages:
|
||||
uart: !include ../../test_build_components/common/uart/esp8266-ard.yaml
|
||||
uart_tcp: !include common.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
|
||||
@@ -0,0 +1,3 @@
|
||||
packages:
|
||||
uart: !include ../../test_build_components/common/uart/rp2040-ard.yaml
|
||||
uart_tcp: !include common.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
|
||||
@@ -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()
|
||||
Reference in New Issue
Block a user