mirror of
https://github.com/esphome/esphome.git
synced 2026-10-07 11:26:39 +00:00
[socket] Move the outgoing buffer into the TCP client link (#20010)
This commit is contained in:
@@ -7,6 +7,7 @@
|
||||
|
||||
#include <algorithm>
|
||||
#include <cerrno>
|
||||
#include <cstring>
|
||||
|
||||
namespace esphome::socket {
|
||||
|
||||
@@ -114,7 +115,7 @@ ssize_t TcpClientLink::read(uint8_t *buf, size_t len) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
ssize_t TcpClientLink::write(const uint8_t *buf, size_t len) {
|
||||
ssize_t TcpClientLink::write_(const uint8_t *buf, size_t len) {
|
||||
if (!this->connected_ || len == 0) {
|
||||
return 0;
|
||||
}
|
||||
@@ -129,6 +130,26 @@ ssize_t TcpClientLink::write(const uint8_t *buf, size_t len) {
|
||||
return -1;
|
||||
}
|
||||
|
||||
size_t TcpClientLink::queue(const uint8_t *data, size_t len) {
|
||||
size_t room = this->tx_free();
|
||||
if (len > room) {
|
||||
len = room;
|
||||
}
|
||||
std::memcpy(this->tx_ + this->tx_len_, data, len);
|
||||
this->tx_len_ += static_cast<uint16_t>(len);
|
||||
return len;
|
||||
}
|
||||
|
||||
void TcpClientLink::flush_tx_slow_() {
|
||||
ssize_t sent = this->write_(this->tx_, this->tx_len_);
|
||||
if (sent > 0) {
|
||||
this->tx_len_ -= static_cast<uint16_t>(sent);
|
||||
if (this->tx_len_ != 0) {
|
||||
std::memmove(this->tx_, this->tx_ + sent, this->tx_len_);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void TcpClientLink::close() {
|
||||
if (this->sock_ != nullptr) {
|
||||
this->sock_->shutdown(SHUT_RDWR);
|
||||
@@ -136,6 +157,7 @@ void TcpClientLink::close() {
|
||||
this->sock_.reset();
|
||||
}
|
||||
this->connected_ = false;
|
||||
this->tx_len_ = 0;
|
||||
this->resolved_.forget();
|
||||
}
|
||||
|
||||
|
||||
@@ -16,8 +16,9 @@
|
||||
namespace esphome::socket {
|
||||
|
||||
/// A reconnecting TCP stream driven from loop(). Owns the socket, the DNS
|
||||
/// lookup and the retry backoff. A fatal read/write error closes the link
|
||||
/// and schedules the next attempt; the caller sees the edge via connected().
|
||||
/// lookup, the retry backoff and the outgoing buffer. A fatal read/write
|
||||
/// error closes the link and schedules the next attempt; the caller sees
|
||||
/// the edge via connected().
|
||||
class TcpClientLink {
|
||||
public:
|
||||
void set_host(const char *host) { this->host_ = StringRef(host); }
|
||||
@@ -41,7 +42,21 @@ class TcpClientLink {
|
||||
void adopt(std::unique_ptr<Socket> sock);
|
||||
/// Returns bytes moved, 0 when nothing can move now, -1 when the link dropped.
|
||||
ssize_t read(uint8_t *buf, size_t len);
|
||||
ssize_t write(const uint8_t *buf, size_t len);
|
||||
/// Copy into the outgoing buffer; returns how many bytes fit.
|
||||
size_t queue(const uint8_t *data, size_t len);
|
||||
/// Direct access to the buffer's free tail. Fill at most tx_free() bytes,
|
||||
/// then tx_commit() the count; neither is bounds checked.
|
||||
uint8_t *tx_tail() { return this->tx_ + this->tx_len_; }
|
||||
void tx_commit(size_t len) { this->tx_len_ += static_cast<uint16_t>(len); }
|
||||
size_t tx_free() const { return this->connected_ ? TX_BUFFER_SIZE - this->tx_len_ : 0; }
|
||||
/// Send the front of the buffer; true once it is empty.
|
||||
/// A partial write keeps the rest; inline no-op while nothing is queued.
|
||||
bool flush_tx() {
|
||||
if (this->tx_len_ != 0) {
|
||||
this->flush_tx_slow_();
|
||||
}
|
||||
return this->tx_len_ == 0;
|
||||
}
|
||||
/// Close without scheduling a reconnect (shutdown).
|
||||
void close();
|
||||
|
||||
@@ -54,6 +69,11 @@ class TcpClientLink {
|
||||
}
|
||||
|
||||
protected:
|
||||
static constexpr size_t TX_BUFFER_SIZE = 1024;
|
||||
|
||||
/// The raw stream write behind flush_tx(); drops the link on a fatal error.
|
||||
ssize_t write_(const uint8_t *buf, size_t len);
|
||||
void flush_tx_slow_();
|
||||
void poll_slow_();
|
||||
void try_connect_();
|
||||
/// Close after a failure, log what and errno, schedule the next attempt.
|
||||
@@ -66,7 +86,9 @@ class TcpClientLink {
|
||||
uint32_t reconnect_interval_ms_{5000};
|
||||
Ipv4Resolve resolved_;
|
||||
uint16_t port_{0};
|
||||
uint16_t tx_len_{0};
|
||||
bool connected_{false};
|
||||
uint8_t tx_[TX_BUFFER_SIZE]{};
|
||||
};
|
||||
|
||||
} // namespace esphome::socket
|
||||
|
||||
@@ -33,7 +33,6 @@ void TcpUart::sync_link_() {
|
||||
this->link_was_up_ = up;
|
||||
if (!up) {
|
||||
this->rx_start_ = this->rx_end_ = 0;
|
||||
this->tx_len_ = 0;
|
||||
}
|
||||
if (this->connected_sensor_ != nullptr) {
|
||||
this->connected_sensor_->publish_state(up);
|
||||
@@ -63,14 +62,6 @@ void TcpUart::read_socket_() {
|
||||
this->rx_pending_ = static_cast<size_t>(count) == room;
|
||||
}
|
||||
|
||||
void TcpUart::flush_tx_() {
|
||||
ssize_t sent = this->link_.write(this->tx_, this->tx_len_);
|
||||
if (sent > 0) {
|
||||
this->tx_len_ -= static_cast<uint16_t>(sent);
|
||||
std::memmove(this->tx_, this->tx_ + sent, this->tx_len_);
|
||||
}
|
||||
}
|
||||
|
||||
void TcpUart::loop() {
|
||||
this->link_.poll();
|
||||
if (this->link_.connected() != this->link_was_up_) {
|
||||
@@ -82,25 +73,20 @@ void TcpUart::loop() {
|
||||
if (this->rx_pending_ || this->link_.ready()) {
|
||||
this->read_socket_();
|
||||
}
|
||||
if (this->tx_len_ != 0) {
|
||||
this->flush_tx_();
|
||||
}
|
||||
this->link_.flush_tx();
|
||||
}
|
||||
|
||||
void TcpUart::write_array(const uint8_t *data, size_t len) {
|
||||
size_t room = this->link_.connected() ? sizeof(this->tx_) - this->tx_len_ : 0;
|
||||
if (len > room) {
|
||||
size_t queued = this->link_.queue(data, len);
|
||||
if (queued < len) {
|
||||
uint32_t now = App.get_loop_component_start_time();
|
||||
if (this->last_drop_log_ms_ == 0 || now - this->last_drop_log_ms_ >= DROP_LOG_INTERVAL_MS) {
|
||||
ESP_LOGW(TAG, "%s, dropped %u bytes",
|
||||
this->link_.connected() ? LOG_STR_LITERAL("TX buffer full") : LOG_STR_LITERAL("Not connected"),
|
||||
static_cast<unsigned>(len - room));
|
||||
static_cast<unsigned>(len - queued));
|
||||
this->last_drop_log_ms_ = now;
|
||||
}
|
||||
len = room;
|
||||
}
|
||||
std::memcpy(this->tx_ + this->tx_len_, data, len);
|
||||
this->tx_len_ += static_cast<uint16_t>(len);
|
||||
}
|
||||
|
||||
bool TcpUart::peek_byte(uint8_t *data) {
|
||||
@@ -121,11 +107,8 @@ bool TcpUart::read_array(uint8_t *data, size_t len) {
|
||||
}
|
||||
|
||||
uart::UARTFlushResult TcpUart::flush() {
|
||||
this->flush_tx_();
|
||||
if (this->tx_len_ == 0) {
|
||||
return uart::UARTFlushResult::UART_FLUSH_RESULT_SUCCESS;
|
||||
}
|
||||
return uart::UARTFlushResult::UART_FLUSH_RESULT_TIMEOUT;
|
||||
return this->link_.flush_tx() ? uart::UARTFlushResult::UART_FLUSH_RESULT_SUCCESS
|
||||
: uart::UARTFlushResult::UART_FLUSH_RESULT_TIMEOUT;
|
||||
}
|
||||
|
||||
} // namespace esphome::tcp_uart
|
||||
|
||||
@@ -32,7 +32,7 @@ class TcpUart : public uart::UARTComponent, public Component {
|
||||
bool read_array(uint8_t *data, size_t len) override;
|
||||
size_t available() override { return static_cast<size_t>(this->rx_end_ - this->rx_start_); }
|
||||
// Same room write_array() grants, so consumers can apply backpressure.
|
||||
size_t available_for_write() override { return this->link_.connected() ? sizeof(this->tx_) - this->tx_len_ : 0; }
|
||||
size_t available_for_write() override { return this->link_.tx_free(); }
|
||||
uart::UARTFlushResult flush() override;
|
||||
bool is_connected() override { return this->link_.connected(); }
|
||||
#if defined(USE_ESP8266) || defined(USE_ESP32)
|
||||
@@ -43,24 +43,20 @@ class TcpUart : public uart::UARTComponent, public Component {
|
||||
void check_logger_conflict() override {}
|
||||
void sync_link_();
|
||||
void read_socket_();
|
||||
void flush_tx_();
|
||||
|
||||
static constexpr size_t RX_BUFFER_SIZE = 1024;
|
||||
static constexpr size_t TX_BUFFER_SIZE = 1024;
|
||||
|
||||
socket::TcpClientLink link_;
|
||||
binary_sensor::BinarySensor *connected_sensor_{nullptr};
|
||||
uint32_t last_drop_log_ms_{0};
|
||||
uint16_t tx_len_{0};
|
||||
// rx_[rx_start_, rx_end_) holds unread bytes; read_socket_() compacts to the front.
|
||||
uint16_t rx_start_{0};
|
||||
uint16_t rx_end_{0};
|
||||
// The link state loop() saw last; edges clear the buffers and publish the sensor.
|
||||
// The link state loop() saw last; edges clear rx_ 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};
|
||||
uint8_t rx_[RX_BUFFER_SIZE]{};
|
||||
uint8_t tx_[TX_BUFFER_SIZE]{};
|
||||
};
|
||||
|
||||
} // namespace esphome::tcp_uart
|
||||
|
||||
+8
-2
@@ -14,14 +14,20 @@ void TcpClientLinkTestComponent::loop() {
|
||||
this->was_up_ = up;
|
||||
ESP_LOGI(TAG, "Link %s", up ? LOG_STR_LITERAL("up") : LOG_STR_LITERAL("down"));
|
||||
}
|
||||
if (!up || !this->link_.ready()) {
|
||||
if (!up) {
|
||||
return;
|
||||
}
|
||||
// Retries a partial echo; inline no-op when nothing is queued.
|
||||
this->link_.flush_tx();
|
||||
if (!this->link_.ready()) {
|
||||
return;
|
||||
}
|
||||
uint8_t buf[64];
|
||||
ssize_t count = this->link_.read(buf, sizeof(buf));
|
||||
if (count > 0) {
|
||||
ESP_LOGI(TAG, "Echoing %d bytes", static_cast<int>(count));
|
||||
this->link_.write(buf, static_cast<size_t>(count));
|
||||
this->link_.queue(buf, static_cast<size_t>(count));
|
||||
this->link_.flush_tx();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user