Classify UART traffic for arbitrary protocol passthrough

This commit is contained in:
puddly
2026-08-05 14:40:38 -04:00
parent 997e218376
commit 7a05749231
9 changed files with 692 additions and 463 deletions
@@ -36,20 +36,27 @@ void SerialProxy::setup() {
void SerialProxy::loop() {
#ifdef USE_API
// Safety check — loop should only run when subscribed, but guard against races
if (this->api_connection_ == nullptr) [[unlikely]] {
this->disable_loop();
return;
}
// Detect subscriber disconnect
if (this->api_connection_->is_marked_for_removal() || !this->api_connection_->is_connection_setup() ||
!api_is_connected()) {
if (this->api_connection_ != nullptr &&
(this->api_connection_->is_marked_for_removal() || !this->api_connection_->is_connection_setup() ||
!api_is_connected())) {
ESP_LOGW(TAG, "Subscriber disconnected");
this->api_connection_ = nullptr;
this->parent_->release(this);
}
// With no subscriber there is normally nothing to do, but a tap may still need the port
// read -- it does its protocol work precisely while nobody else is listening.
if (this->api_connection_ == nullptr) [[unlikely]] {
#ifdef USE_SERIAL_PROXY_TAP
if (this->tap_ == nullptr || !this->tap_->tap_needs_port()) {
this->disable_loop();
return;
}
#else
this->disable_loop();
return;
#endif
}
// Read available data from UART and forward to subscribed client
@@ -70,11 +77,33 @@ void __attribute__((noinline)) SerialProxy::read_and_send_(size_t available) {
if (!this->read_array(buffer, to_read))
return;
#ifdef USE_SERIAL_PROXY_TAP
// Before forwarding, so a tap that answers the device (an acknowledgement, say) is not
// waiting on the network round trip to a subscriber that may not even exist.
if (this->tap_ != nullptr) {
this->tap_->on_device_rx(buffer, to_read);
}
#endif
if (this->api_connection_ == nullptr) {
return;
}
this->outgoing_msg_.set_data(buffer, to_read);
this->api_connection_->send_serial_proxy_data(this->outgoing_msg_);
}
#endif
#ifdef USE_SERIAL_PROXY_TAP
void SerialProxy::tap_pump() {
#ifdef USE_API
const size_t available = this->available();
if (available > 0) {
this->read_and_send_(available);
}
#endif
}
#endif
void SerialProxy::dump_config() {
ESP_LOGCONFIG(TAG,
"Serial Proxy [%" PRIu32 "]:\n"
@@ -171,6 +200,13 @@ void SerialProxy::write_from_client(api::APIConnection *api_connection, const ui
if (data == nullptr || len == 0)
return;
this->write_array(data, len);
#ifdef USE_SERIAL_PROXY_TAP
// After the write, so the tap observes the same ordering the device does
if (this->tap_ != nullptr) {
this->tap_->on_client_tx(data, len);
}
#endif
}
void SerialProxy::set_modem_pins(api::APIConnection *api_connection, uint32_t line_states) {
+48 -1
View File
@@ -41,6 +41,29 @@ enum SerialProxyLineStateFlag : uint32_t {
/// Maximum bytes to read from UART in a single loop iteration
inline constexpr size_t SERIAL_PROXY_MAX_READ_SIZE = 256;
#ifdef USE_SERIAL_PROXY_TAP
/// Observes a port's traffic without owning it, and may inject bytes of its own.
///
/// This exists so protocol-aware behaviour can be layered onto a plain byte pipe without
/// the pipe knowing anything about the protocol: the tap is compiled in only when some
/// component asks for one, so a proxy carrying an RS485 meter pays nothing for it.
///
/// A tap is an observer, never a gatekeeper -- it cannot suppress or alter the bytes
/// flowing in either direction, so a misbehaving tap cannot corrupt the stream.
class SerialProxyTap {
public:
/// Bytes read from the device, before they are forwarded to any subscriber.
virtual void on_device_rx(const uint8_t *data, size_t len) = 0;
/// Bytes a subscriber sent towards the device, after they have been written.
virtual void on_client_tx(const uint8_t *data, size_t len) = 0;
/// True when the port must keep reading even with no subscriber attached, so a tap can
/// do its own protocol work while nobody is listening.
virtual bool tap_needs_port() const = 0;
};
#endif
class SerialProxy final : public uart::UARTDevice, public Component {
public:
void setup() override;
@@ -103,9 +126,29 @@ class SerialProxy final : public uart::UARTDevice, public Component {
/// Set the DTR GPIO pin (from YAML configuration)
void set_dtr_pin(GPIOPin *pin) { this->dtr_pin_ = pin; }
#ifdef USE_SERIAL_PROXY_TAP
/// Attach a traffic observer. At most one, set once at setup time.
void set_tap(SerialProxyTap *tap) { this->tap_ = tap; }
/// Write bytes originating from the tap rather than from a client. Bypasses the
/// subscriber ownership check, since the tap is part of the device, not a client of it.
void write_from_tap(const uint8_t *data, size_t len) { this->write_array(data, len); }
/// Resume reading after a tap's needs change. loop() disables itself when there is
/// neither a subscriber nor a tap that wants the port, so a tap starting fresh work
/// must ask for it back.
void tap_request_port() { this->enable_loop(); }
/// Run one read-and-dispatch cycle immediately. Lets a tap make progress before the
/// main loop is running -- during setup, for instance, while a component is still
/// blocking on can_proceed().
void tap_pump();
#endif
protected:
#ifdef USE_API
/// Read from UART and send to API client (slow path with 256-byte stack buffer)
/// Read from UART, hand the bytes to any tap, and forward them to a subscriber
/// (slow path with a 256-byte stack buffer)
void read_and_send_(size_t available);
/// True when a live subscriber other than the given connection holds the port
@@ -136,6 +179,10 @@ class SerialProxy final : public uart::UARTDevice, public Component {
/// Current modem pin states
bool rts_state_{false};
bool dtr_state_{false};
#ifdef USE_SERIAL_PROXY_TAP
SerialProxyTap *tap_{nullptr};
#endif
};
} // namespace esphome::serial_proxy
@@ -0,0 +1,256 @@
#include "ash_detector.h"
#ifdef USE_ZIGBEE_PROXY
namespace esphome::zigbee_proxy {
// Control byte of an RSTACK, and the only ASH version byte that can follow it
static constexpr uint8_t ASH_RSTACK_CONTROL = 0xC1;
static constexpr uint8_t ASH_PROTOCOL_VERSION = 0x02;
static constexpr size_t ASH_RSTACK_BODY_SIZE = 3; // control, version, reset code
static constexpr size_t ASH_CRC_SIZE = 2;
// Smallest legal frame on the wire: a bare control byte plus its CRC
static constexpr size_t ASH_MIN_FRAME_SIZE = 1 + ASH_CRC_SIZE;
// The opening EZSP version command is a constant: control 0x00 (frmNum 0, ackNum 0)
// followed by [seq=0][frameControl=0][frameId=0] randomized by 0x42 0x21 0xA8. Only the
// requested version varies, as version ^ 0x54, so it can be recovered for free.
static constexpr uint8_t EZSP_VERSION_CMD_PREFIX[] = {0x00, 0x42, 0x21, 0xA8};
static constexpr size_t EZSP_VERSION_CMD_SIZE = 5;
static constexpr uint8_t EZSP_VERSION_RANDOM_MASK = 0x54;
// Consecutive frames we could not accept, with neither a good frame nor a retransmission
// in between, before concluding the peer is no longer speaking ASH. A real ASH peer must
// retransmit an unacknowledged frame, so the absence of one is the positive evidence
// here -- garbage on the line is not, since noise proves nothing either way.
static constexpr uint8_t MAX_UNCONFIRMED_REJECTS = 4;
bool ash_reset_code_is_known(uint8_t code) {
switch (code) {
case 0x00: // RESET_UNKNOWN
case 0x01: // RESET_EXTERNAL
case 0x02: // RESET_POWER_ON
case 0x03: // RESET_WATCHDOG
case 0x06: // RESET_ASSERT
case 0x09: // RESET_BOOTLOADER
case 0x0B: // RESET_SOFTWARE
case 0x51: // ERROR_EXCEEDED_MAXIMUM_ACK_TIMEOUT_COUNT
case 0x80: // ERROR_CHIP_SPECIFIC
case 0x81: // RESET_CHIP_SPECIFIC
return true;
default:
return false;
}
}
void AshFrameScanner::begin_frame_() {
this->index_ = 0;
this->crc_ = ASH_CRC_INIT;
this->escaped_ = false;
this->poisoned_ = false;
}
void AshFrameScanner::reset() {
this->begin_frame_();
this->frame_length_ = 0;
this->discarding_ = false;
}
ScanResult AshFrameScanner::feed(uint8_t byte) {
if (byte == ASH_FLAG_BYTE) {
// Snapshot everything the verdict depends on: begin_frame_() clears all of it.
const bool discarding = this->discarding_;
const bool poisoned = this->poisoned_;
const bool escaped = this->escaped_;
const size_t index = this->index_;
const uint16_t crc = this->crc_;
// A FLAG always starts the next frame afresh, whatever preceded it
this->begin_frame_();
this->discarding_ = false;
if (discarding || index == 0) {
// Consecutive delimiters carry no frame at all, so there is nothing to judge
this->frame_length_ = 0;
return ScanResult::NONE;
}
// Running the CRC over the body *and* its trailing CRC bytes leaves zero when
// correct, so validity needs no second pass over the frame.
if (poisoned || escaped || index < ASH_MIN_FRAME_SIZE || crc != 0) {
this->frame_length_ = 0;
return ScanResult::INVALID;
}
this->frame_length_ = index - ASH_CRC_SIZE;
return ScanResult::FRAME;
}
if (this->discarding_) {
return ScanResult::NONE;
}
switch (byte) {
case ASH_CANCEL_BYTE:
// Everything received since the last FLAG is to be ignored
this->begin_frame_();
return ScanResult::NONE;
case ASH_SUBSTITUTE_BYTE:
// A low-level error was flagged; ignore everything up to the next FLAG
this->discarding_ = true;
return ScanResult::NONE;
case ASH_XON_BYTE:
case ASH_XOFF_BYTE:
// Transport flow control, not frame content: skip it without disturbing the frame
return ScanResult::NONE;
case ASH_ESCAPE_BYTE:
this->escaped_ = true;
return ScanResult::NONE;
default:
break;
}
uint8_t value = byte;
if (this->escaped_) {
this->escaped_ = false;
value = byte ^ ASH_XOR_BYTE;
// An escape must decode to a reserved byte; anything else is not ASH framing at all
if (!ash_is_reserved(value)) {
this->poisoned_ = true;
return ScanResult::NONE;
}
}
if (this->index_ >= sizeof(this->buffer_)) {
this->poisoned_ = true;
return ScanResult::NONE;
}
this->buffer_[this->index_++] = value;
this->crc_ = ash_crc16(&value, 1, this->crc_);
return ScanResult::NONE;
}
void AshDetector::reset() {
this->ncp_scanner_.reset();
this->host_scanner_.reset();
this->state_ = AshDetectState::IDLE;
this->rx_sequence_ = 0;
this->ack_owed_ = false;
this->data_frame_ready_ = false;
this->unconfirmed_rejects_ = 0;
this->negotiated_version_ = 0;
}
void AshDetector::from_ncp(uint8_t byte) {
this->data_frame_ready_ = false;
switch (this->ncp_scanner_.feed(byte)) {
case ScanResult::FRAME:
this->handle_ncp_frame_();
break;
case ScanResult::INVALID:
// A delimited chunk that is not a frame. While armed this may be a corrupted ASH
// frame, which the peer will retransmit, or a sign the peer stopped speaking ASH.
// reject_() distinguishes the two by whether a retransmission ever arrives.
this->reject_();
break;
case ScanResult::NONE:
break;
}
}
void AshDetector::handle_ncp_frame_() {
const uint8_t *body = this->ncp_scanner_.frame();
const size_t length = this->ncp_scanner_.length();
const uint8_t control = body[0];
// RSTACK is the only way into the handshake, and the only way back after a firmware
// swap: a Spinel or bootloader NCP never emits one, so those stay unarmed forever.
if (control == ASH_RSTACK_CONTROL) {
if (length == ASH_RSTACK_BODY_SIZE && body[1] == ASH_PROTOCOL_VERSION && ash_reset_code_is_known(body[2])) {
this->state_ = AshDetectState::SAW_RSTACK;
this->rx_sequence_ = 0;
this->ack_owed_ = false;
this->unconfirmed_rejects_ = 0;
}
return;
}
if (this->state_ != AshDetectState::ARMED) {
return;
}
if ((control & 0x80) != 0) {
return; // ACK/NAK/RST/ERROR: nothing is owed for these
}
const uint8_t frame_num = (control >> 4) & ASH_MAX_SEQUENCE;
const bool re_tx = (control & 0x08) != 0;
if (frame_num != this->rx_sequence_) {
// A retransmission still proves the peer is speaking ASH even though we cannot use
// this copy, so it clears the suspicion without being acknowledged.
if (re_tx) {
this->unconfirmed_rejects_ = 0;
} else {
this->reject_();
}
return;
}
this->rx_sequence_ = (this->rx_sequence_ + 1) & ASH_MAX_SEQUENCE;
this->pending_ack_ = this->rx_sequence_;
this->ack_owed_ = true;
this->data_frame_ready_ = true;
this->unconfirmed_rejects_ = 0;
}
void AshDetector::reject_() {
if (this->state_ != AshDetectState::ARMED) {
return;
}
if (++this->unconfirmed_rejects_ >= MAX_UNCONFIRMED_REJECTS) {
this->state_ = AshDetectState::IDLE;
this->unconfirmed_rejects_ = 0;
}
}
void AshDetector::from_host(uint8_t byte) {
if (this->host_scanner_.feed(byte) != ScanResult::FRAME) {
return;
}
if (this->state_ != AshDetectState::SAW_RSTACK) {
return;
}
const uint8_t *body = this->host_scanner_.frame();
if (this->host_scanner_.length() != EZSP_VERSION_CMD_SIZE) {
return;
}
for (size_t i = 0; i < sizeof(EZSP_VERSION_CMD_PREFIX); i++) {
if (body[i] != EZSP_VERSION_CMD_PREFIX[i]) {
return;
}
}
this->negotiated_version_ = body[4] ^ EZSP_VERSION_RANDOM_MASK;
this->state_ = AshDetectState::ARMED;
this->rx_sequence_ = 0;
this->ack_owed_ = false;
this->unconfirmed_rejects_ = 0;
}
bool AshDetector::take_pending_ack(uint8_t &ack_num) {
if (!this->ack_owed_) {
return false;
}
this->ack_owed_ = false;
ack_num = this->pending_ack_;
return true;
}
} // namespace esphome::zigbee_proxy
#endif // USE_ZIGBEE_PROXY
@@ -0,0 +1,118 @@
#pragma once
#include "esphome/core/defines.h"
#ifdef USE_ZIGBEE_PROXY
#include "ash_protocol.h"
#include <cstddef>
#include <cstdint>
namespace esphome::zigbee_proxy {
// Decides when it is safe to acknowledge NCP frames on a client's behalf.
//
// The client suppresses its own ACKs, so nobody else will send them, and injecting ASH
// bytes into a stream that is not ASH would corrupt it. Detection is therefore one-sided:
// arm only on the session handshake, which is a fixed byte string, and never on frame
// validity, which non-ASH traffic can satisfy by luck.
//
// RSTACK (NCP -> host) c1 02 <reset_code> <crc> 7e
// version (host -> NCP) 00 42 21 a8 <version^0x54> <crc> 7e
//
// Requiring both, in that order, in opposite directions cannot be satisfied by a
// unidirectional byte stream whatever it contains -- which is exactly the situation
// during a firmware upload. Verified against real .gbl images and real Spinel traffic:
// zero false arms, and neither pattern occurs even as a substring.
//
// Getting it wrong in the other direction is cheap: a frame we decline to acknowledge is
// retransmitted by the NCP, so we see a clean copy and lose only the ack timeout. That
// asymmetry is why this errs towards silence everywhere.
enum class AshDetectState : uint8_t {
IDLE, // Not ASH, or not yet proven to be
SAW_RSTACK, // Handshake half-complete; watching for the version command
ARMED, // Session confirmed; acknowledging on the client's behalf
};
enum class ScanResult : uint8_t {
NONE, // Mid-frame, or a delimiter that carried nothing
FRAME, // frame()/length() hold a complete body with a verified CRC
INVALID, // A delimited chunk arrived but was not a well-formed ASH frame
};
// Reassembles one direction of the byte stream into unstuffed, CRC-checked frames.
// Mirrors bellows' AshProtocol.data_received: FLAG ends a frame, CANCEL discards what
// precedes it, SUBSTITUTE poisons everything up to the next FLAG, and XON/XOFF are
// transport flow control removed without disturbing the frame around them.
class AshFrameScanner {
public:
ScanResult feed(uint8_t byte);
void reset();
// Valid only until the next feed() call, which begins overwriting the buffer.
const uint8_t *frame() const { return this->buffer_; }
size_t length() const { return this->frame_length_; }
private:
void begin_frame_();
// Frames are bounded by the ASH maximum, so a stream carrying no delimiters cannot
// grow the buffer without limit; it just keeps failing.
uint8_t buffer_[MAX_ASH_FRAME_SIZE];
size_t index_{0}; // accumulation position for the frame being read
size_t frame_length_{0}; // body length of the last completed frame
uint16_t crc_{ASH_CRC_INIT};
bool escaped_{false};
bool discarding_{false};
bool poisoned_{false};
};
class AshDetector {
public:
void reset();
// Feed observed traffic. Neither call gates forwarding: the detector only watches.
void from_ncp(uint8_t byte);
void from_host(uint8_t byte);
bool armed() const { return this->state_ == AshDetectState::ARMED; }
// True only while the host direction can affect the state machine, i.e. while waiting
// for the version command. Lets the caller skip scanning that direction entirely the
// rest of the time -- it is the one carrying firmware uploads.
bool needs_host_scan() const { return this->state_ == AshDetectState::SAW_RSTACK; }
// An acknowledgement became owed after the last from_ncp() call. Clears the flag.
bool take_pending_ack(uint8_t &ack_num);
// The EZSP frame carried by the DATA frame just accepted, for metadata sniffing. The
// ASH control byte is skipped, so offset 0 is the EZSP sequence number. Still
// randomized, and valid only until the next from_ncp() call.
const uint8_t *last_ezsp_frame() const { return this->ncp_scanner_.frame() + 1; }
size_t last_ezsp_frame_length() const {
const size_t length = this->ncp_scanner_.length();
return length > 0 ? length - 1 : 0;
}
AshDetectState state() const { return this->state_; }
uint8_t negotiated_version() const { return this->negotiated_version_; }
protected:
void handle_ncp_frame_();
void reject_();
AshFrameScanner ncp_scanner_;
AshFrameScanner host_scanner_;
AshDetectState state_{AshDetectState::IDLE};
uint8_t rx_sequence_{0};
uint8_t pending_ack_{0};
bool ack_owed_{false};
bool data_frame_ready_{false};
uint8_t unconfirmed_rejects_{0};
uint8_t negotiated_version_{0};
};
} // namespace esphome::zigbee_proxy
#endif // USE_ZIGBEE_PROXY
@@ -41,7 +41,7 @@ void ash_randomize(uint8_t *data, size_t length) {
}
}
uint16_t ZigbeeProxy::calculate_crc_(const uint8_t *data, size_t length, uint16_t init) {
uint16_t ash_crc16(const uint8_t *data, size_t length, uint16_t init) {
uint16_t crc = init;
for (size_t i = 0; i < length; i++) {
crc = (crc << 8) ^ CRC_TABLE[(crc >> 8) ^ data[i]];
@@ -49,6 +49,10 @@ uint16_t ZigbeeProxy::calculate_crc_(const uint8_t *data, size_t length, uint16_
return crc;
}
uint16_t ZigbeeProxy::calculate_crc_(const uint8_t *data, size_t length, uint16_t init) {
return ash_crc16(data, length, init);
}
bool ZigbeeProxy::validate_frame_crc_() {
// CRC is calculated over control byte + data
// rx_buffer_[0] contains control byte, rx_buffer_[1..rx_buffer_index_-3] contains data
@@ -85,7 +89,6 @@ bool ZigbeeProxy::handle_ack_num_(uint8_t ack_num) {
this->update_adaptive_timeout_(rtt);
ESP_LOGV(TAG, "Frame %d acknowledged, RTT: %u ms", this->tx_pending_frame_num_, rtt);
this->clear_tx_buffer_();
this->drain_ncp_tx_queue_();
return true;
}
@@ -157,13 +160,6 @@ void ZigbeeProxy::parse_control_byte_(uint8_t control) {
return;
}
// Backpressure: a client-bound frame is still waiting on API TX buffer space, so we
// cannot forward another one. Do not ACK or consume - the NCP will retransmit.
if (!this->boot_sequence_active_ && this->api_connection_ != nullptr && this->client_tx_pending_length_ > 0) {
ESP_LOGV(TAG, "Client TX pending, deferring DATA frame %d to NCP retransmit", frame_num);
return;
}
// Increment RX sequence and send ACK (ack_num = next expected frame)
this->increment_rx_sequence_();
this->send_ack_frame_(this->rx_sequence_);
@@ -172,15 +168,12 @@ void ZigbeeProxy::parse_control_byte_(uint8_t control) {
size_t payload_length = this->rx_buffer_index_ > 3 ? this->rx_buffer_index_ - 3 : 0;
const uint8_t *payload = this->rx_buffer_.data() + 1;
// During boot sequence, route to boot handler. Frames consumed locally must
// be derandomized first; proxied frames below are passed through untouched
// so the client's own randomization survives end to end.
if (this->boot_sequence_active_ && payload_length > 0) {
// This path only runs during the boot harvest, where this component is the ASH
// endpoint and consumes frames itself, so they must be derandomized. A subscribed
// client is served by the transparent relay instead, which never reaches here.
if (payload_length > 0) {
ash_randomize(this->rx_buffer_.data() + 1, payload_length);
this->handle_boot_data_frame_(payload, payload_length);
} else if (this->api_connection_ != nullptr && payload_length > 0) {
// Forward EZSP payload to client via client-side ASH DATA frame
this->forward_ncp_data_to_client_(payload, payload_length);
}
break;
}
+17 -6
View File
@@ -11,8 +11,23 @@ static constexpr uint8_t ASH_ESCAPE_BYTE = 0x7D; // Escape/substitution byt
static constexpr uint8_t ASH_XOR_BYTE = 0x20; // XOR mask for escaped bytes
static constexpr uint8_t ASH_SUBSTITUTE_BYTE = 0x18; // Substitution for invalid bytes
// Reserved bytes that must be escaped
static constexpr uint8_t ASH_RESERVED_BYTES[] = {0x7E, 0x7D, 0x11, 0x13, 0x93, 0xA3};
static constexpr uint8_t ASH_XON_BYTE = 0x11; // Resume transmission
static constexpr uint8_t ASH_XOFF_BYTE = 0x13; // Pause transmission
static constexpr uint8_t ASH_CANCEL_BYTE = 0x1A; // Discards the partial frame before it
// A reserved byte can never appear literally inside a frame; it is escaped as
// ESCAPE followed by the byte XOR 0x20. Rejecting frames that contain one is what
// eliminates most non-ASH traffic before its CRC is ever computed: real firmware
// images and Spinel payloads are dense in 0x11/0x13/0x18/0x1A.
inline bool ash_is_reserved(uint8_t byte) {
return byte == ASH_FLAG_BYTE || byte == ASH_ESCAPE_BYTE || byte == ASH_XON_BYTE || byte == ASH_XOFF_BYTE ||
byte == ASH_SUBSTITUTE_BYTE || byte == ASH_CANCEL_BYTE;
}
// CRC-CCITT (init 0xFFFF, polynomial 0x1021, transmitted big-endian). Note this is a
// different variant from the Kermit FCS that Spinel/HDLC-lite uses over the same
// 0x7E framing, so Spinel frames systematically fail this check.
uint16_t ash_crc16(const uint8_t *data, size_t length, uint16_t init = 0xFFFF);
// Buffer size configuration
#ifdef ZIGBEE_PROXY_BUFFER_SIZE
@@ -32,10 +47,6 @@ static constexpr uint8_t ASH_MAX_RETRIES = 5; // Maximum retransmission a
static constexpr uint16_t ASH_CRC_INIT = 0xFFFF; // CRC-CCITT initial value
static constexpr uint32_t ASH_RESET_TIMEOUT = 3000; // RST/RSTACK timeout in milliseconds
// Client -> NCP queue depth: frames accepted while the single ASH TX window is occupied.
// Overflow beyond this is NAKed to the client, which retransmits.
static constexpr uint8_t NCP_TX_QUEUE_SIZE = 2;
// IEEE address size
static constexpr size_t ZIGBEE_IEEE_ADDR_SIZE = 8; // 64-bit IEEE address
@@ -27,10 +27,17 @@ static constexpr uint8_t EZSP_FRAME_CONTROL_EXTENDED = 0x01;
// NCP starts in legacy mode and has not yet learned the negotiated version.
// Everything after that is extended, with no per-NCP exceptions.
// EZSP Frame IDs - Callbacks (NCP to host, async)
static constexpr uint16_t EZSP_STACK_STATUS_HANDLER = 0x0019; // Stack up/down notification
// EZSP Frame IDs - Commands (host to NCP)
static constexpr uint16_t EZSP_VERSION = 0x0000; // Version negotiation
static constexpr uint16_t EZSP_GET_EUI64 = 0x0026; // Get IEEE address
static constexpr uint16_t EZSP_GET_TOKEN_DATA = 0x0102; // Read an NVM3 token
static constexpr uint16_t EZSP_VERSION = 0x0000; // Version negotiation
static constexpr uint16_t EZSP_GET_EUI64 = 0x0026; // Get IEEE address
static constexpr uint16_t EZSP_GET_NETWORK_PARAMETERS = 0x0028; // Get network parameters
static constexpr uint16_t EZSP_GET_TOKEN_DATA = 0x0102; // Read an NVM3 token
// Extended EZSP header: [sequence] [frame_control_lo] [frame_control_hi] [id_lo] [id_hi]
static constexpr size_t EZSP_EXTENDED_HEADER_SIZE = 5;
// Network metadata comes straight out of NVM3 instead of from a running stack.
// NVM3KEY_STACK_NODE_DATA holds the PAN ID, channel, extended PAN ID and node type of
@@ -44,7 +51,6 @@ static constexpr uint16_t EZSP_GET_TOKEN_DATA = 0x0102; // Read an NVM3 token
static constexpr uint32_t NVM3KEY_STACK_NODE_DATA = 0x0001EE64;
// getTokenData response: [status (4)] [length (4)] [value (length)]
static constexpr size_t TOKEN_DATA_LENGTH_OFFSET = 4;
static constexpr size_t TOKEN_DATA_VALUE_OFFSET = 8;
// NV3StackNodeData value layout (16 bytes, little-endian):
@@ -65,10 +71,18 @@ static constexpr uint8_t NV3_NODE_TYPE_UNKNOWN_DEVICE = 0x00;
// legacy 8-bit EmberStatus.
enum class SlStatus : uint8_t {
OK = 0x00,
NETWORK_UP = 0x15,
NETWORK_DOWN = 0x16,
};
// sl_status_t is 32-bit little-endian on the wire, so a status field occupies
// four bytes even though every code used here fits in the first one.
static constexpr size_t SL_STATUS_SIZE = 4;
// getNetworkParameters response layout, used when sniffing a client's own traffic. This
// is a different shape from the NV3 token the boot harvest reads: 25 bytes of
// [status (4)] [nodeType (1)] [extendedPanId (8)] [panId (2)] [radioTxPower (1)]
// [radioChannel (1)] [joinMethod (1)] [nwkManagerId (2)] [nwkUpdateId (1)] [channels (4)]
static constexpr size_t NETWORK_PARAMS_RESPONSE_SIZE = 25;
static constexpr size_t NETWORK_PARAMS_STATUS_OFFSET = 0;
static constexpr size_t NETWORK_PARAMS_EXT_PAN_ID_OFFSET = 5;
static constexpr size_t NETWORK_PARAMS_PAN_ID_OFFSET = 13;
static constexpr size_t NETWORK_PARAMS_CHANNEL_OFFSET = 16;
} // namespace esphome::zigbee_proxy
+159 -384
View File
@@ -98,8 +98,8 @@ void ZigbeeProxy::loop() {
if (this->boot_sequence_active_) {
this->check_boot_timeouts_();
} else if (this->ash_state_ == AshState::CONNECTING && millis() - this->setup_time_ > ASH_RESET_TIMEOUT) {
// Stuck in CONNECTING state (client-triggered RST, not the boot sequence)
} else if (this->api_connection_ == nullptr && this->ash_state_ == AshState::CONNECTING &&
millis() - this->setup_time_ > ASH_RESET_TIMEOUT) {
ESP_LOGE(TAG, "RSTACK timeout, NCP not responding");
this->ash_state_ = AshState::FAILED;
}
@@ -110,27 +110,11 @@ void ZigbeeProxy::loop() {
this->unsubscribe_api_connection(this->api_connection_);
}
// Retry any client-bound frame that hit API TX buffer backpressure
this->try_send_pending_client_frame_();
// Send any queued client frames if the ASH TX window opened up
this->drain_ncp_tx_queue_();
// Autonomous recovery: periodically retry a failed NCP link so a late-powered NCP does
// not require a client RST. Only while subscribed -- with nobody listening there is
// nothing to recover for, and resetting could disturb another device's use of the bus.
if (this->api_connection_ != nullptr &&
(this->ash_state_ == AshState::FAILED ||
(this->boot_state_ == BootState::FAILED && !this->boot_sequence_active_)) &&
millis() - this->last_recovery_attempt_ > RECOVERY_RETRY_INTERVAL_MS) {
ESP_LOGI(TAG, "Attempting NCP recovery");
this->last_recovery_attempt_ = millis();
// Reset the link only. Re-running the harvest would renegotiate the NCP's EZSP
// version underneath the subscriber, which then keeps addressing a version the NCP
// no longer speaks -- and because the harvest consumes its own RSTACKs, it would
// never find out. Relaying the RSTACK instead lets it resynchronize deliberately.
this->reset_ncp_link_();
}
// No autonomous recovery while a client is subscribed. A subscriber owns the link: it
// opens with its own RST and resets whenever it decides it needs to. Resetting on its
// behalf relays an RSTACK it never asked for, which bellows treats as fatal -- and if it
// happens to be driving a bootloader over this interface, injecting ASH into the
// transfer is worse still. A broken link is the client's to notice and repair.
}
void ZigbeeProxy::check_boot_timeouts_() {
@@ -227,7 +211,15 @@ void ZigbeeProxy::zigbee_proxy_request(api::APIConnection *api_connection, const
}
ESP_LOGD(TAG, "Client subscribed");
this->api_connection_ = api_connection;
this->client_reset_session_();
// A subscriber owns the link from here on, so abandon any harvest in flight rather
// than interleaving our own EZSP commands with the client's session. Metadata for
// this session comes from watching the client's own traffic instead.
if (this->boot_sequence_active_) {
ESP_LOGD(TAG, "Abandoning boot harvest, client owns the link");
this->boot_sequence_active_ = false;
this->boot_state_ = BootState::IDLE;
}
this->detector_.reset();
break;
case api::enums::ZIGBEE_PROXY_REQUEST_TYPE_UNSUBSCRIBE:
@@ -238,7 +230,7 @@ void ZigbeeProxy::zigbee_proxy_request(api::APIConnection *api_connection, const
break;
case api::enums::ZIGBEE_PROXY_REQUEST_TYPE_NETWORK_INFO:
this->send_network_info_response_(api_connection);
this->send_network_info_changed_msg_(api_connection);
break;
default:
@@ -252,10 +244,9 @@ void ZigbeeProxy::unsubscribe_api_connection(api::APIConnection *conn) {
return;
}
this->api_connection_ = nullptr;
// Frames belonging to the departed client's session must not linger
this->client_tx_pending_length_ = 0;
this->ncp_tx_queue_count_ = 0;
this->client_reset_session_();
// Anything already buffered belongs to the departed session
this->relay_length_ = 0;
this->detector_.reset();
}
void ZigbeeProxy::zigbee_proxy_frame(api::APIConnection *api_connection, const api::ZigbeeProxyFrame &msg) {
@@ -264,9 +255,17 @@ void ZigbeeProxy::zigbee_proxy_frame(api::APIConnection *api_connection, const a
return;
}
// Feed raw bytes into the client-side ASH parser
for (size_t i = 0; i < msg.data_len; i++) {
this->client_parse_byte_(msg.data[i]);
// Transparent relay: the client's ASH bytes reach the NCP untouched, so the two share
// one sequence space and nothing here can desynchronize it.
this->write_array(msg.data, msg.data_len);
// Scanning this direction only matters while waiting for the version command that
// completes the handshake. Outside that window it is a pure passthrough -- which is
// what makes a firmware upload, all of which flows this way, essentially free.
if (this->detector_.needs_host_scan()) {
for (size_t i = 0; i < msg.data_len; i++) {
this->detector_.from_host(msg.data[i]);
}
}
}
@@ -278,40 +277,6 @@ uint64_t ZigbeeProxy::get_ieee_address() const {
return addr;
}
bool ZigbeeProxy::send_frame(const uint8_t *data, size_t length) {
if (this->ash_state_ != AshState::CONNECTED) {
ESP_LOGW(TAG, "Cannot send frame, not connected");
return false;
}
// Transmit directly when the ASH window is free and nothing is queued ahead
if (!this->tx_buffer_pending_ && this->ncp_tx_queue_count_ == 0) {
return this->send_data_frame_(data, length, false);
}
// Window occupied - queue for transmission when the pending frame is ACKed
if (this->ncp_tx_queue_count_ >= NCP_TX_QUEUE_SIZE || length > MAX_ASH_FRAME_SIZE) {
return false; // Caller NAKs the client, which retransmits
}
QueuedTxFrame &entry =
this->ncp_tx_queue_[(this->ncp_tx_queue_head_ + this->ncp_tx_queue_count_) % NCP_TX_QUEUE_SIZE];
memcpy(entry.data.data(), data, length);
entry.length = length;
this->ncp_tx_queue_count_++;
ESP_LOGV(TAG, "Queued frame (%u bytes, %u queued)", length, this->ncp_tx_queue_count_);
return true;
}
void ZigbeeProxy::drain_ncp_tx_queue_() {
if (this->tx_buffer_pending_ || this->ncp_tx_queue_count_ == 0 || this->ash_state_ != AshState::CONNECTED) {
return;
}
QueuedTxFrame &entry = this->ncp_tx_queue_[this->ncp_tx_queue_head_];
this->ncp_tx_queue_head_ = (this->ncp_tx_queue_head_ + 1) % NCP_TX_QUEUE_SIZE;
this->ncp_tx_queue_count_--;
this->send_data_frame_(entry.data.data(), entry.length, false);
}
void ZigbeeProxy::set_timeout_config(uint32_t initial_ms, uint32_t min_ms, uint32_t max_ms) {
this->timeout_config_.initial_timeout_ms = initial_ms;
this->timeout_config_.min_timeout_ms = min_ms;
@@ -335,7 +300,6 @@ void ZigbeeProxy::reset_ash_protocol_() {
this->rx_sequence_ = 0;
this->tx_buffer_pending_ = false;
this->tx_retry_count_ = 0;
this->ncp_tx_queue_count_ = 0;
this->parsing_state_ = ParsingState::WAIT_FLAG_START;
this->setup_time_ = millis();
this->boot_start_time_ = this->setup_time_;
@@ -364,9 +328,9 @@ void ZigbeeProxy::reset_ncp_link_() {
this->rx_sequence_ = 0;
this->tx_buffer_pending_ = false;
this->tx_retry_count_ = 0;
this->ncp_tx_queue_count_ = 0;
this->parsing_state_ = ParsingState::WAIT_FLAG_START;
this->client_reset_session_();
this->relay_length_ = 0;
this->detector_.reset();
this->send_rst_frame_(false);
}
@@ -445,10 +409,6 @@ void ZigbeeProxy::handle_rstack_frame_(const uint8_t *data, size_t length) {
// RSTACK during connecting (triggered by client RST forwarding)
ESP_LOGV(TAG, "Received RSTACK, NCP ready");
this->ash_state_ = AshState::CONNECTED;
// Forward to client if subscribed
if (this->api_connection_ != nullptr) {
this->forward_ncp_rstack_to_client_(this->rx_buffer_.data() + 1, this->rx_buffer_index_ - 3);
}
} else if (solicited_by_us) {
// Surplus RSTACK from one of our own resets, most often an RST retry racing a reply
// that was merely slow. The client never asked for it, so swallow it.
@@ -459,7 +419,6 @@ void ZigbeeProxy::handle_rstack_frame_(const uint8_t *data, size_t length) {
// client's session, so it has to hear about it.
ESP_LOGW(TAG, "NCP reset unexpectedly, notifying client");
this->ash_state_ = AshState::CONNECTED;
this->forward_ncp_rstack_to_client_(this->rx_buffer_.data() + 1, this->rx_buffer_index_ - 3);
} else {
ESP_LOGW(TAG, "Unexpected RSTACK received (boot_state=%d)", static_cast<int>(this->boot_state_));
}
@@ -508,7 +467,6 @@ void ZigbeeProxy::handle_error_frame_(const uint8_t *data, size_t length) {
if (this->api_connection_ != nullptr) {
// Forward error to client
this->forward_ncp_error_to_client_(data, length);
} else {
// No client, attempt recovery ourselves
ESP_LOGV(TAG, "Attempting recovery");
@@ -984,8 +942,6 @@ void ZigbeeProxy::send_network_info_changed_msg_(api::APIConnection *conn) {
}
}
void ZigbeeProxy::send_network_info_response_(api::APIConnection *conn) { this->send_network_info_changed_msg_(conn); }
// WiFi/Zigbee channel conflict detection
void ZigbeeProxy::check_wifi_zigbee_conflict_() {
#ifdef USE_WIFI
@@ -1087,326 +1043,145 @@ void ZigbeeProxy::process_uart_slow_() {
this->bootloader_state_ = BootloaderState::NORMAL;
}
this->parse_byte_(byte);
} while (this->available());
}
// ==================== Client-side ASH session ====================
void ZigbeeProxy::client_reset_session_() {
this->client_tx_sequence_ = 0;
this->client_rx_sequence_ = 0;
this->client_rx_buffer_index_ = 0;
this->client_escape_next_byte_ = false;
this->client_tx_pending_length_ = 0;
this->client_ash_state_ = AshState::DISCONNECTED;
this->client_parsing_state_ = ParsingState::WAIT_FLAG_START;
ESP_LOGV(TAG, "Client ASH session reset");
}
bool ZigbeeProxy::send_to_client_(const uint8_t *data, size_t length) {
if (this->api_connection_ == nullptr) {
return false;
}
this->outgoing_proto_msg_.data = data;
this->outgoing_proto_msg_.data_len = length;
return this->api_connection_->send_zigbee_proxy_frame(this->outgoing_proto_msg_);
}
void ZigbeeProxy::client_send_raw_frame_(const uint8_t *frame, size_t length) {
// Failure here means API TX buffer backpressure. Losing a control frame is recoverable:
// an unsent ACK triggers a client retransmit, which the duplicate-frame path re-ACKs.
if (!this->send_to_client_(frame, length)) {
ESP_LOGV(TAG, "Dropped %u byte control frame to client (API TX buffer full)", length);
}
}
void ZigbeeProxy::client_send_ack_frame_(uint8_t ack_num) {
uint8_t frame[8];
size_t length = this->build_frame_(frame, sizeof(frame), nullptr, 0, AshFrameType::ACK, 0, ack_num);
this->client_send_raw_frame_(frame, length);
ESP_LOGV(TAG, "Sent client ACK for frame %d", ack_num);
}
void ZigbeeProxy::client_send_nak_frame_(uint8_t ack_num) {
uint8_t frame[8];
size_t length = this->build_frame_(frame, sizeof(frame), nullptr, 0, AshFrameType::NAK, 0, ack_num);
this->client_send_raw_frame_(frame, length);
ESP_LOGV(TAG, "Sent client NAK for frame %d", ack_num);
}
void ZigbeeProxy::client_send_rstack_frame_(uint8_t reset_code) {
// RSTACK payload: [version] [reset_code]
uint8_t payload[] = {0x02, reset_code};
uint8_t frame[16];
size_t length = this->build_frame_(frame, sizeof(frame), payload, sizeof(payload), AshFrameType::RSTACK);
this->client_send_raw_frame_(frame, length);
ESP_LOGV(TAG, "Sent client RSTACK (code=0x%02X)", reset_code);
}
void ZigbeeProxy::client_send_data_frame_(const uint8_t *data, size_t length) {
size_t frame_length = this->build_frame_(this->client_tx_buffer_.data(), this->client_tx_buffer_.size(), data, length,
AshFrameType::DATA, this->client_tx_sequence_, this->client_rx_sequence_);
if (frame_length == 0) {
return; // build_frame_ logged the error; payload cannot be represented in the buffer
}
this->client_tx_sequence_ = (this->client_tx_sequence_ + 1) & ASH_MAX_SEQUENCE;
if (!this->send_to_client_(this->client_tx_buffer_.data(), frame_length)) {
// API TX buffer backpressure: keep the frame and retry from loop(). While a frame is
// pending, incoming NCP DATA frames are left unACKed so the NCP provides flow control.
ESP_LOGV(TAG, "Client DATA frame deferred (API TX buffer full)");
this->client_tx_pending_length_ = frame_length;
this->client_tx_pending_since_ = millis();
return;
}
ESP_LOGV(TAG, "Sent client DATA frame, payload %u bytes", length);
}
void ZigbeeProxy::try_send_pending_client_frame_() {
if (this->client_tx_pending_length_ == 0) {
return;
}
if (this->send_to_client_(this->client_tx_buffer_.data(), this->client_tx_pending_length_)) {
ESP_LOGV(TAG, "Sent deferred client DATA frame");
this->client_tx_pending_length_ = 0;
return;
}
if (millis() - this->client_tx_pending_since_ > CLIENT_TX_RETRY_TIMEOUT_MS) {
// The API connection is not draining; abandon the frame and force the client to
// re-establish a clean ASH session (best effort - the ERROR frame may also fail)
ESP_LOGE(TAG, "Client TX stalled for %u ms, resetting client session", CLIENT_TX_RETRY_TIMEOUT_MS);
this->client_tx_pending_length_ = 0;
this->client_send_error_frame_(static_cast<uint8_t>(EzspError::EXCEEDED_MAXIMUM_ACK_TIMEOUT_COUNT));
this->client_reset_session_();
}
}
void ZigbeeProxy::client_send_error_frame_(uint8_t error_code) {
uint8_t payload[] = {0x02, error_code};
uint8_t frame[16];
size_t length = this->build_frame_(frame, sizeof(frame), payload, sizeof(payload), AshFrameType::ERROR);
this->client_send_raw_frame_(frame, length);
ESP_LOGV(TAG, "Sent client ERROR (code=0x%02X)", error_code);
}
void ZigbeeProxy::forward_ncp_data_to_client_(const uint8_t *payload, size_t length) {
this->client_send_data_frame_(payload, length);
}
void ZigbeeProxy::forward_ncp_rstack_to_client_(const uint8_t *data, size_t length) {
// RSTACK payload is [version] [reset_code] per spec; cap defensively so a malformed
// NCP frame cannot overflow the stack buffer
length = std::min(length, static_cast<size_t>(2));
uint8_t frame[16];
size_t frame_length = this->build_frame_(frame, sizeof(frame), data, length, AshFrameType::RSTACK);
this->client_send_raw_frame_(frame, frame_length);
// Reset client-side sequence numbers since RSTACK means new session
this->client_tx_sequence_ = 0;
this->client_rx_sequence_ = 0;
this->client_tx_pending_length_ = 0;
this->client_ash_state_ = AshState::CONNECTED;
ESP_LOGV(TAG, "Forwarded RSTACK to client");
}
void ZigbeeProxy::forward_ncp_error_to_client_(const uint8_t *data, size_t length) {
// ERROR payload is [version] [error_code] per spec; cap defensively (see RSTACK above)
length = std::min(length, static_cast<size_t>(2));
uint8_t frame[16];
size_t frame_length = this->build_frame_(frame, sizeof(frame), data, length, AshFrameType::ERROR);
this->client_send_raw_frame_(frame, frame_length);
ESP_LOGV(TAG, "Forwarded ERROR to client");
}
void ZigbeeProxy::client_parse_byte_(uint8_t byte) {
static constexpr uint8_t ASH_CAN_BYTE = 0x1A;
static constexpr uint8_t ASH_XON_BYTE = 0x11;
static constexpr uint8_t ASH_XOFF_BYTE = 0x13;
// Reserved bytes count only when bare, never as the second half of an escape
// sequence (see the matching comment in parse_byte_).
if (!this->client_escape_next_byte_) {
if (byte == ASH_CAN_BYTE) {
// Cancel: discard any partial frame
this->client_rx_buffer_index_ = 0;
this->client_parsing_state_ = ParsingState::WAIT_FLAG_START;
return;
}
if (byte == ASH_XON_BYTE || byte == ASH_XOFF_BYTE) {
// Flow control: not part of any frame
return;
}
}
switch (this->client_parsing_state_) {
case ParsingState::WAIT_FLAG_START:
if (byte == ASH_ESCAPE_BYTE) {
this->client_escape_next_byte_ = true;
return;
}
if (this->client_escape_next_byte_) {
byte ^= ASH_XOR_BYTE;
this->client_escape_next_byte_ = false;
}
if (byte == ASH_FLAG_BYTE) {
this->client_rx_buffer_index_ = 0;
this->client_escape_next_byte_ = false;
this->client_parsing_state_ = ParsingState::WAIT_CONTROL;
} else if ((byte & 0x80) != 0 || this->client_ash_state_ == AshState::CONNECTED) {
// Accept control byte without leading FLAG
this->client_rx_buffer_index_ = 0;
this->client_rx_buffer_[this->client_rx_buffer_index_++] = byte;
this->client_parsing_state_ = ParsingState::WAIT_DATA;
}
break;
case ParsingState::WAIT_CONTROL:
if (byte == ASH_FLAG_BYTE) {
// Empty frame or repeated FLAG
this->client_rx_buffer_index_ = 0;
return;
}
if (byte == ASH_ESCAPE_BYTE) {
this->client_escape_next_byte_ = true;
return;
}
if (this->client_escape_next_byte_) {
byte ^= ASH_XOR_BYTE;
this->client_escape_next_byte_ = false;
}
this->client_rx_buffer_[this->client_rx_buffer_index_++] = byte;
this->client_parsing_state_ = ParsingState::WAIT_DATA;
break;
case ParsingState::WAIT_DATA:
if (byte == ASH_FLAG_BYTE) {
// End of frame - validate and process
if (this->client_validate_frame_crc_()) {
this->client_parse_control_byte_(this->client_rx_buffer_[0]);
} else {
ESP_LOGW(TAG, "Client frame CRC failed (%u bytes)", this->client_rx_buffer_index_);
}
this->client_parsing_state_ = ParsingState::WAIT_FLAG_START;
return;
}
if (byte == ASH_ESCAPE_BYTE) {
this->client_escape_next_byte_ = true;
return;
}
if (this->client_escape_next_byte_) {
byte ^= ASH_XOR_BYTE;
this->client_escape_next_byte_ = false;
}
if (this->client_rx_buffer_index_ >= MAX_ASH_FRAME_SIZE) {
ESP_LOGE(TAG, "Client RX buffer overflow");
this->client_parsing_state_ = ParsingState::WAIT_FLAG_START;
return;
}
this->client_rx_buffer_[this->client_rx_buffer_index_++] = byte;
break;
default:
this->client_parsing_state_ = ParsingState::WAIT_FLAG_START;
break;
}
}
bool ZigbeeProxy::client_validate_frame_crc_() {
if (this->client_rx_buffer_index_ < 3) {
return false;
}
uint16_t calculated = this->calculate_crc_(this->client_rx_buffer_.data(), this->client_rx_buffer_index_ - 2);
uint16_t received = (static_cast<uint16_t>(this->client_rx_buffer_[this->client_rx_buffer_index_ - 2]) << 8) |
this->client_rx_buffer_[this->client_rx_buffer_index_ - 1];
return calculated == received;
}
void ZigbeeProxy::client_parse_control_byte_(uint8_t control) {
AshFrameType frame_type;
if ((control & 0x80) == 0) {
frame_type = AshFrameType::DATA;
} else if ((control & 0xC0) == 0x80) {
frame_type = ((control & 0x20) == 0) ? AshFrameType::ACK : AshFrameType::NAK;
} else {
uint8_t control_bits = control & 0x07;
if (control_bits == 0x00) {
frame_type = AshFrameType::RST;
} else if (control_bits == 0x01) {
frame_type = AshFrameType::RSTACK;
} else if (control_bits == 0x02) {
frame_type = AshFrameType::ERROR;
if (this->boot_sequence_active_) {
// Harvest: this component is the ASH endpoint and consumes the frames itself
this->parse_byte_(byte);
} else {
ESP_LOGW(TAG, "Client: unknown control frame type: 0x%02X", control);
return;
this->relay_ncp_byte_(byte);
}
} while (this->available());
this->relay_flush_();
}
// ==================== Transparent relay ====================
void ZigbeeProxy::relay_ncp_byte_(uint8_t byte) {
if (this->relay_length_ >= sizeof(this->relay_buffer_)) {
this->relay_flush_();
}
this->relay_buffer_[this->relay_length_++] = byte;
// Observation only: the detector never gates forwarding, so it adds no latency and a
// frame it cannot parse still reaches the client, which judges it for itself.
this->detector_.from_ncp(byte);
// In relay mode the detector is the only thing watching the link, so its progress is
// what tells us the NCP is alive. Without this ash_state_ sits at CONNECTING, times out
// into FAILED, and autonomous recovery resets the NCP underneath a working session --
// relaying an RSTACK the client never asked for, which kills it outright.
if (this->detector_.state() != AshDetectState::IDLE) {
this->ash_state_ = AshState::CONNECTED;
}
uint8_t frame_num = (control >> 4) & 0x07;
uint8_t ack_num = control & 0x07;
uint8_t ack_num;
if (this->detector_.take_pending_ack(ack_num)) {
// The client suppresses its own ACKs, so this is the only acknowledgement the NCP
// will see. Only ever sent for a frame that passed CRC and arrived in sequence.
this->send_ack_frame_(ack_num);
const uint8_t *ezsp = this->detector_.last_ezsp_frame();
const size_t ezsp_length = this->detector_.last_ezsp_frame_length();
this->sniff_network_info_(ezsp, ezsp_length);
this->sniff_stack_status_(ezsp, ezsp_length);
}
}
switch (frame_type) {
case AshFrameType::DATA: {
// Verify sequence number
if (frame_num != this->client_rx_sequence_) {
uint8_t retx_bit = this->client_rx_buffer_[0] & 0x08;
if (retx_bit != 0 && frame_num == ((this->client_rx_sequence_ - 1) & ASH_MAX_SEQUENCE)) {
// Retransmission of a frame we already accepted (our ACK was lost) - re-ACK and discard
ESP_LOGV(TAG, "Client: duplicate DATA frame %d, re-sending ACK", frame_num);
this->client_send_ack_frame_(this->client_rx_sequence_);
} else {
ESP_LOGW(TAG, "Client: out of sequence DATA frame: expected %d, got %d", this->client_rx_sequence_,
frame_num);
this->client_send_nak_frame_(this->client_rx_sequence_);
}
return;
}
// A proxied getNetworkParameters response is the only authoritative view of the network
// available while a client owns the link, so metadata is refreshed from the client's own
// traffic rather than by injecting commands. Read-only: a frame that fails any check
// simply leaves the previous values in place.
void ZigbeeProxy::sniff_network_info_(const uint8_t *frame, size_t length) {
// Every frame after the version handshake uses extended framing, so the header size is
// fixed and needs no knowledge of the negotiated version.
if (length < EZSP_EXTENDED_HEADER_SIZE + NETWORK_PARAMS_RESPONSE_SIZE) {
return;
}
// Extract EZSP payload (skip control byte, exclude CRC)
size_t payload_length = this->client_rx_buffer_index_ > 3 ? this->client_rx_buffer_index_ - 3 : 0;
const uint8_t *payload = this->client_rx_buffer_.data() + 1;
// The relay never derandomizes, so work on a copy: the original bytes are already on
// their way to the client and must stay untouched.
uint8_t decoded[EZSP_EXTENDED_HEADER_SIZE + NETWORK_PARAMS_RESPONSE_SIZE];
memcpy(decoded, frame, sizeof(decoded));
ash_randomize(decoded, sizeof(decoded));
if (payload_length > 0) {
// Forward EZSP payload to NCP via right-side ASH; NAK without consuming if the
// NCP link is down or the TX queue is full so the client retransmits
ESP_LOGV(TAG, "Client DATA → NCP, EZSP payload %u bytes", payload_length);
if (!this->send_frame(payload, payload_length)) {
this->client_send_nak_frame_(this->client_rx_sequence_);
return;
}
}
const uint16_t frame_id = decoded[3] | (static_cast<uint16_t>(decoded[4]) << 8);
if (frame_id != EZSP_GET_NETWORK_PARAMETERS || (decoded[1] & EZSP_FRAME_CONTROL_RESPONSE) == 0) {
return;
}
// Accepted: advance the sequence and ACK immediately rather than relying on the
// piggybacked ACK of an eventual NCP response
this->client_rx_sequence_ = (this->client_rx_sequence_ + 1) & ASH_MAX_SEQUENCE;
this->client_send_ack_frame_(this->client_rx_sequence_);
break;
}
const uint8_t *params = decoded + EZSP_EXTENDED_HEADER_SIZE;
if (params[NETWORK_PARAMS_STATUS_OFFSET] != static_cast<uint8_t>(SlStatus::OK)) {
return;
}
case AshFrameType::ACK:
// Client ACKed our data - nothing to retransmit on client side for now
ESP_LOGV(TAG, "Client ACK received for frame %d", ack_num);
break;
const uint16_t pan_id =
params[NETWORK_PARAMS_PAN_ID_OFFSET] | (static_cast<uint16_t>(params[NETWORK_PARAMS_PAN_ID_OFFSET + 1]) << 8);
const uint8_t channel = params[NETWORK_PARAMS_CHANNEL_OFFSET];
const bool changed =
pan_id != this->network_info_.pan_id || channel != this->network_info_.channel ||
memcmp(this->network_info_.extended_pan_id.data(), params + NETWORK_PARAMS_EXT_PAN_ID_OFFSET, 8) != 0;
if (!changed) {
return;
}
case AshFrameType::NAK:
ESP_LOGW(TAG, "Client NAK received for frame %d", ack_num);
break;
memcpy(this->network_info_.extended_pan_id.data(), params + NETWORK_PARAMS_EXT_PAN_ID_OFFSET, 8);
this->network_info_.pan_id = pan_id;
this->network_info_.channel = channel;
this->network_info_.valid = true;
ESP_LOGD(TAG, "Network info from proxied traffic: PAN 0x%04X, channel %u", pan_id, channel);
this->send_network_info_changed_msg_();
}
case AshFrameType::RST:
// Client wants to reset - forward RST to NCP
// Don't use reset_ash_protocol_() as that enters boot sequence
ESP_LOGV(TAG, "Client RST → forwarding to NCP");
this->reset_ncp_link_();
break;
void ZigbeeProxy::sniff_stack_status_(const uint8_t *frame, size_t length) {
// stackStatusHandler: [status (4)]. It carries no network parameters, so it can only
// invalidate, never refresh. Worth acting on anyway: a client leaving a network need not
// read parameters afterwards, and without this the old PAN would be reported forever.
if (length < EZSP_EXTENDED_HEADER_SIZE + 1) {
return;
}
case AshFrameType::RSTACK:
// Client shouldn't send RSTACK, ignore
ESP_LOGW(TAG, "Client sent unexpected RSTACK");
break;
uint8_t decoded[EZSP_EXTENDED_HEADER_SIZE + 1];
memcpy(decoded, frame, sizeof(decoded));
ash_randomize(decoded, sizeof(decoded));
case AshFrameType::ERROR:
// Client sent error, log it
ESP_LOGW(TAG, "Client sent ERROR frame");
break;
const uint16_t frame_id = decoded[3] | (static_cast<uint16_t>(decoded[4]) << 8);
// A callback sets both the response direction and the callback bit, so match on both
if (frame_id != EZSP_STACK_STATUS_HANDLER ||
(decoded[1] & EZSP_FRAME_CONTROL_CALLBACK) != EZSP_FRAME_CONTROL_CALLBACK) {
return;
}
if (decoded[EZSP_EXTENDED_HEADER_SIZE] != static_cast<uint8_t>(SlStatus::NETWORK_DOWN)) {
return; // NETWORK_UP and everything else leaves what we have standing
}
if (this->network_info_.pan_id == 0 && this->network_info_.channel == 0) {
return; // Already reporting no network
}
ESP_LOGD(TAG, "Stack reported NETWORK_DOWN, dropping stale network info");
this->network_info_.pan_id = 0;
this->network_info_.channel = 0;
this->network_info_.extended_pan_id.fill(0);
// The IEEE address is a hardware property and survives leaving a network, so it stays.
this->send_network_info_changed_msg_();
}
void ZigbeeProxy::relay_flush_() {
if (this->relay_length_ == 0) {
return;
}
const size_t length = this->relay_length_;
this->relay_length_ = 0;
if (this->api_connection_ == nullptr) {
return;
}
this->outgoing_proto_msg_.data = this->relay_buffer_;
this->outgoing_proto_msg_.data_len = length;
if (!this->api_connection_->send_zigbee_proxy_frame(this->outgoing_proto_msg_)) {
// API TX backpressure. Dropping bytes is recoverable: the client sees a truncated
// frame, fails its CRC and NAKs, and the NCP retransmits. Withholding our ACK would
// achieve the same thing more slowly, and buffering risks unbounded growth.
ESP_LOGW(TAG, "Dropped %u relayed bytes (API TX buffer full)", length);
}
}
+20 -41
View File
@@ -9,6 +9,7 @@
#include "esphome/core/helpers.h"
#include "esphome/components/uart/uart.h"
#include "ash_protocol.h"
#include "ash_detector.h"
#include <array>
@@ -86,11 +87,6 @@ class ZigbeeProxy : public uart::UARTDevice, public Component {
const NetworkInfo &get_network_info() const { return this->network_info_; }
uint64_t get_ieee_address() const;
// Send an EZSP payload to the NCP: transmits immediately when the ASH link is idle,
// otherwise queues it. Returns false if the link is down or the queue is full (caller
// should NAK the client so it retransmits).
bool send_frame(const uint8_t *data, size_t length);
// Timeout configuration (callable from Python/API)
void set_timeout_config(uint32_t initial_ms, uint32_t min_ms, uint32_t max_ms);
void set_initial_timeout(uint32_t timeout_ms) { this->timeout_config_.initial_timeout_ms = timeout_ms; }
@@ -151,9 +147,6 @@ class ZigbeeProxy : public uart::UARTDevice, public Component {
this->tx_retry_count_ = 0;
}
// Client -> NCP pending frame queue (absorbs frames arriving while a TX is unacknowledged)
void drain_ncp_tx_queue_();
// Boot-time NCP initialization
void advance_boot_state_();
void check_boot_timeouts_();
@@ -211,21 +204,15 @@ class ZigbeeProxy : public uart::UARTDevice, public Component {
void client_send_data_frame_(const uint8_t *data, size_t length);
void client_send_error_frame_(uint8_t error_code);
void client_send_raw_frame_(const uint8_t *frame, size_t length);
void client_reset_session_();
// Retry a client-bound DATA frame that failed to send (API TX buffer backpressure)
void try_send_pending_client_frame_();
// Send raw bytes to API client; returns false if the API TX buffer rejected the message
bool send_to_client_(const uint8_t *data, size_t length);
// Network info request/push (packed: ieee[8] + extended_pan[8] + pan_id[2] + channel[1], little-endian)
void send_network_info_response_(api::APIConnection *conn);
// Forward NCP frames to client
void forward_ncp_data_to_client_(const uint8_t *payload, size_t length);
void forward_ncp_rstack_to_client_(const uint8_t *data, size_t length);
void forward_ncp_error_to_client_(const uint8_t *data, size_t length);
// Transparent relay. NCP bytes are forwarded to the client verbatim and in bulk; the
// detector only observes them, so it never gates or delays forwarding.
void relay_ncp_byte_(uint8_t byte);
void relay_flush_();
// Reads network metadata out of a proxied getNetworkParameters response. Read-only, so a
// misparse costs a missed update rather than corrupting anything.
void sniff_network_info_(const uint8_t *frame, size_t length);
// Invalidates network metadata when the stack reports it has left the network.
void sniff_stack_status_(const uint8_t *frame, size_t length);
// Pre-allocated message - always ready to send
api::ZigbeeProxyFrame outgoing_proto_msg_;
@@ -236,16 +223,8 @@ class ZigbeeProxy : public uart::UARTDevice, public Component {
std::array<uint8_t, MAX_ASH_FRAME_SIZE> tx_pending_buffer_; // For retransmission
// Client-side (left) ASH buffers
std::array<uint8_t, MAX_ASH_FRAME_SIZE> client_rx_buffer_;
std::array<uint8_t, MAX_ASH_FRAME_SIZE> client_tx_buffer_;
// Client -> NCP queue: EZSP payloads accepted while the ASH TX window is occupied
struct QueuedTxFrame {
uint16_t length;
std::array<uint8_t, MAX_ASH_FRAME_SIZE> data;
};
std::array<QueuedTxFrame, NCP_TX_QUEUE_SIZE> ncp_tx_queue_;
// Network information
NetworkInfo network_info_;
@@ -263,7 +242,6 @@ class ZigbeeProxy : public uart::UARTDevice, public Component {
uint32_t last_recovery_attempt_{0}; // Time of last automatic reset attempt from FAILED
// Client-side (left) 32-bit values
uint32_t client_tx_pending_since_{0}; // Time the pending client frame first failed to send
// NCP-side (right) 16-bit values
uint16_t rx_buffer_index_{0}; // Index for populating rx_buffer_
@@ -271,8 +249,6 @@ class ZigbeeProxy : public uart::UARTDevice, public Component {
uint16_t calculated_crc_{0}; // CRC calculated during frame reception
// Client-side (left) 16-bit values
uint16_t client_rx_buffer_index_{0};
uint16_t client_tx_pending_length_{0}; // >0: client_tx_buffer_ holds an unsent DATA frame
// NCP-side (right) 8-bit values
uint8_t tx_sequence_{0}; // TX sequence number (0-7)
@@ -280,13 +256,9 @@ class ZigbeeProxy : public uart::UARTDevice, public Component {
uint8_t tx_retry_count_{0}; // Number of retransmission attempts
uint8_t tx_pending_frame_num_{0}; // Frame number of pending TX frame
uint8_t last_ack_sent_{0}; // Last ACK number sent
uint8_t ncp_tx_queue_head_{0}; // Oldest entry in ncp_tx_queue_
uint8_t ncp_tx_queue_count_{0}; // Number of queued entries
uint8_t last_rx_byte_{0}; // Previous raw RX byte (bootloader detection)
// Client-side (left) 8-bit values
uint8_t client_tx_sequence_{0}; // Client-facing TX sequence (proxy → client)
uint8_t client_rx_sequence_{0}; // Client-facing RX sequence (client → proxy)
// NCP-side enums and booleans
AshState ash_state_{AshState::DISCONNECTED};
@@ -295,8 +267,6 @@ class ZigbeeProxy : public uart::UARTDevice, public Component {
BootState boot_state_{BootState::IDLE};
// Client-side enums and booleans
AshState client_ash_state_{AshState::DISCONNECTED};
ParsingState client_parsing_state_{ParsingState::WAIT_FLAG_START};
uint8_t ezsp_version_{0}; // NCP's EZSP protocol version
uint8_t ezsp_sequence_{0}; // EZSP frame sequence number
@@ -308,13 +278,22 @@ class ZigbeeProxy : public uart::UARTDevice, public Component {
bool tx_buffer_pending_{false}; // True if waiting for ACK from NCP
bool escape_next_byte_{false}; // True if next NCP byte should be unescaped
bool client_escape_next_byte_{false}; // True if next client byte should be unescaped
bool network_info_ready_{false}; // True when network info retrieved
bool owns_uart_{false}; // True while this component drives the UART
uint32_t configured_baud_rate_{0}; // Line rate to restore after another device
bool boot_sequence_active_{false}; // True during boot-time init
// Decides when acknowledging on the client's behalf is safe. Armed only by the ASH
// session handshake, so a bootloader or Thread NCP never triggers it.
AshDetector detector_;
// Bytes staged for the client. Forwarding in bulk once per UART drain avoids an API
// message per byte; the size only bounds latency, not correctness.
static constexpr size_t RELAY_BUFFER_SIZE = 256;
uint8_t relay_buffer_[RELAY_BUFFER_SIZE];
size_t relay_length_{0};
// RSTACKs still owed to us for resets we sent ourselves. A retry can put two RSTs on
// the wire when the first RSTACK was only slow rather than lost, so the NCP answers
// with more RSTACKs than we asked for; the surplus must not reach a client.