mirror of
https://github.com/esphome/esphome.git
synced 2026-10-06 02:51:29 +00:00
[modbus] Rework the client queue as a per-frame state machine (#17922)
Co-authored-by: Claude Fable 5 <noreply@anthropic.com> Co-authored-by: J. Nick Koston <nick@koston.org>
This commit is contained in:
co-authored by
Claude Fable 5
J. Nick Koston
parent
5821915aad
commit
d38a458de9
@@ -45,26 +45,46 @@ void Modbus::loop() {
|
||||
}
|
||||
|
||||
void ModbusClientHub::loop() {
|
||||
// Call base class to receive bytes and parse frames
|
||||
this->Modbus::loop();
|
||||
// Drain anything owed since the last loop (e.g. an external clear) before the watchdog runs, so it
|
||||
// never times out an entry whose pending count has not been drained. No-op when nothing is owed.
|
||||
this->sweep_();
|
||||
|
||||
// If we're past the send_wait_time timeout and response buffer doesn't have the start of the expected response
|
||||
if (this->waiting_for_response_.has_value()) {
|
||||
ModbusDeviceCommand &wfr = this->waiting_for_response_.value();
|
||||
uint8_t expected_address = wfr.frame.address();
|
||||
if (this->last_receive_check_ - this->last_send_ > this->last_send_tx_offset_ + this->send_wait_time_ &&
|
||||
(this->rx_buffer_.empty() || this->rx_buffer_[0] != expected_address)) {
|
||||
ESP_LOGW(TAG, "Stop waiting for response from %" PRIu8 " %" PRIu32 "ms after last send", expected_address,
|
||||
this->last_receive_check_ - this->last_send_);
|
||||
this->notify_no_response_(wfr);
|
||||
this->waiting_for_response_.reset();
|
||||
}
|
||||
this->Modbus::loop(); // receive bytes and parse frames
|
||||
|
||||
// Send-wait watchdog: only the cheap time check runs at loop rate; expire_waiting_() looks the
|
||||
// entry up and holds off if the response has started arriving.
|
||||
if (this->waiting_for_response_ &&
|
||||
this->last_receive_check_ - this->last_send_ > this->last_send_tx_offset_ + this->send_wait_time_) {
|
||||
this->expire_waiting_();
|
||||
}
|
||||
|
||||
// If there's no response pending and there's commands in the buffer
|
||||
this->sweep_(); // deliver owed callbacks with the hub quiescent
|
||||
this->send_next_frame_();
|
||||
}
|
||||
|
||||
void ModbusClientHub::expire_waiting_() {
|
||||
ModbusDeviceCommand *cmd = this->find_waiting_();
|
||||
if (cmd == nullptr) {
|
||||
this->waiting_for_response_ = false;
|
||||
return;
|
||||
}
|
||||
if (!this->rx_buffer_.empty() && this->rx_buffer_[0] == cmd->frame.address()) {
|
||||
// The start of the response is in the buffer: let the frame finish arriving.
|
||||
return;
|
||||
}
|
||||
// Only a genuine WAITING entry warrants the log (a cleared or interrupted shell timing out is expected).
|
||||
if (cmd->state == FrameState::WAITING) {
|
||||
ESP_LOGW(TAG, "Stop waiting for response from %" PRIu8 " %" PRIu32 "ms after last send", cmd->frame.address(),
|
||||
this->last_receive_check_ - this->last_send_);
|
||||
}
|
||||
// Deliver on_no_response directly, the way the parse path delivers response()/error(): the entry
|
||||
// lands in TIMED_OUT and the following sweep reschedules a retry or erases it. Free the
|
||||
// wire first so a resend from inside the callback sees it available.
|
||||
this->waiting_for_response_ = false;
|
||||
this->sweep_needed_ = true;
|
||||
cmd->timed_out();
|
||||
}
|
||||
|
||||
bool Modbus::timeout_() {
|
||||
// If the response frame is finished (including interframe delay) - we timeout.
|
||||
// The long_rx_buffer_delay accounts for long responses (larger than the UART rx_full_threshold) to avoid timeouts
|
||||
@@ -110,12 +130,21 @@ bool Modbus::tx_blocked() {
|
||||
|
||||
bool ModbusClientHub::tx_blocked() {
|
||||
// We block transmission in any of these case:
|
||||
// 1. We're waiting for a response
|
||||
// 1. We're waiting for a response (a waiting entry: WAITING/INTERRUPTED/WAITING_RETIRED/INTERRUPTED_RETIRED)
|
||||
// 2. Any of the base class tx_blocked conditions
|
||||
return (this->waiting_for_response_.has_value()) || this->Modbus::tx_blocked();
|
||||
return this->waiting_for_response_ || this->Modbus::tx_blocked();
|
||||
}
|
||||
|
||||
bool ModbusClientHub::tx_buffer_empty() { return this->tx_buffer_.empty(); }
|
||||
bool ModbusClientHub::tx_buffer_empty() {
|
||||
// "Empty" for ready_for_immediate_send(): no one-shot is queued ahead of the caller. Entries in
|
||||
// other states are mid-transaction or owed bookkeeping, not queued sends - and a READY continuous
|
||||
// poll does not count either, since it ranks below every one-shot, so a new send goes out first.
|
||||
for (const auto &cmd : this->tx_buffer_) {
|
||||
if (cmd.state == FrameState::READY && !cmd.continuous)
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
void Modbus::receive_bytes_() {
|
||||
this->last_receive_check_ = millis();
|
||||
@@ -257,69 +286,57 @@ bool ModbusServerHub::parse_modbus_client_frame_() {
|
||||
return true;
|
||||
}
|
||||
|
||||
// Bounds contract, enforced by the parser (parse_modbus_server_frame_) rather than locally:
|
||||
// - pdu is never empty: helpers::server_pdu_length() returns at least MIN_PDU_SIZE (1) on every
|
||||
// branch, and find_custom_frame_end_() only ever lengthens the frame, so the PDU always holds
|
||||
// at least the function code.
|
||||
// - When the exception bit is set, pdu has at least 2 bytes: server_pdu_length() checks the
|
||||
// exception bit before anything else and pins those PDUs to 2 bytes, so the exception code
|
||||
// read below is always present.
|
||||
// Keep those guarantees in mind when changing server_pdu_length() or adding callers.
|
||||
// The parser (parse_modbus_server_frame_) guarantees the bounds relied on here: pdu is never empty,
|
||||
// and an exception-flagged pdu is at least 2 bytes. Keep that in mind when changing server_pdu_length().
|
||||
void ModbusClientHub::process_modbus_server_frame(uint8_t address, std::span<const uint8_t> pdu) {
|
||||
const uint8_t function_code = pdu[0];
|
||||
if (!this->waiting_for_response_.has_value()) {
|
||||
ModbusDeviceCommand *cmd = this->waiting_for_response_ ? this->find_waiting_() : nullptr;
|
||||
if (cmd == nullptr) {
|
||||
ESP_LOGW(TAG,
|
||||
"Received unexpected frame from address %" PRIu8 ", function code 0x%X, %" PRIu32 "ms after last send",
|
||||
address, function_code, this->last_modbus_byte_ - this->last_send_);
|
||||
return;
|
||||
} else { // We are waiting for a response
|
||||
// Check if the response matches the expected address and function code
|
||||
}
|
||||
|
||||
ModbusDeviceCommand &wfr = this->waiting_for_response_.value();
|
||||
uint8_t expected_address = wfr.frame.address();
|
||||
uint8_t expected_function_code = wfr.frame.pdu()[0];
|
||||
if (expected_address != address || expected_function_code != (function_code & FUNCTION_CODE_MASK)) {
|
||||
ESP_LOGW(TAG,
|
||||
"Received incorrect frame address %" PRIu8 " <> %" PRIu8 " or function code 0x%X <> 0x%X, %" PRIu32
|
||||
"ms after last send",
|
||||
address, expected_address, (function_code & FUNCTION_CODE_MASK), expected_function_code,
|
||||
this->last_modbus_byte_ - this->last_send_);
|
||||
// Invalidate the device; the entry survives as an interrupted shell so the late response is ignored.
|
||||
// A retry requested here stays queued behind the shell until the send-wait timeout clears it.
|
||||
this->notify_no_response_(wfr);
|
||||
wfr.interrupted = true;
|
||||
return;
|
||||
}
|
||||
// Check if the response matches the expected address and function code
|
||||
const uint8_t expected_address = cmd->frame.address();
|
||||
const uint8_t expected_function_code = cmd->frame.pdu()[0];
|
||||
if (expected_address != address || expected_function_code != (function_code & FUNCTION_CODE_MASK)) {
|
||||
ESP_LOGW(TAG,
|
||||
"Received incorrect frame address %" PRIu8 " <> %" PRIu8 " or function code 0x%X <> 0x%X, %" PRIu32
|
||||
"ms after last send",
|
||||
address, expected_address, (function_code & FUNCTION_CODE_MASK), expected_function_code,
|
||||
this->last_modbus_byte_ - this->last_send_);
|
||||
// Unexpected frame: flip a WAITING entry to an INTERRUPTED shell that ignores the rest of this
|
||||
// transaction and blocks tx until the send-wait timeout, where it gets its on_no_response.
|
||||
cmd->interrupt();
|
||||
return;
|
||||
}
|
||||
|
||||
if (wfr.interrupted) {
|
||||
ESP_LOGW(TAG,
|
||||
"Ignoring response from %" PRIu8 " - transmission interrupted by previous unexpected response, %" PRIu32
|
||||
"ms after last send",
|
||||
address, this->last_modbus_byte_ - this->last_send_);
|
||||
return;
|
||||
} else { // We have a valid device waiting for this response
|
||||
if (cmd->state == FrameState::INTERRUPTED || cmd->state == FrameState::INTERRUPTED_RETIRED) {
|
||||
// An interrupted shell keeps blocking until the send-wait timeout; a late response for it is
|
||||
// ignored and does NOT free the wire. The distrust survives a clear (INTERRUPTED_RETIRED), so a
|
||||
// cleared-interrupted frame still ends in on_no_response rather than delivering a late response.
|
||||
ESP_LOGW(TAG,
|
||||
"Ignoring response from %" PRIu8 " - transmission interrupted by previous unexpected response, %" PRIu32
|
||||
"ms after last send",
|
||||
address, this->last_modbus_byte_ - this->last_send_);
|
||||
return;
|
||||
}
|
||||
|
||||
// Move the command out of the waiting slot so the request PDU stays alive for the callback.
|
||||
ModbusDeviceCommand command = std::move(this->waiting_for_response_.value());
|
||||
this->waiting_for_response_.reset();
|
||||
ModbusClientDevice *device = command.device;
|
||||
// The request PDU is the sent frame without the leading address and the trailing CRC.
|
||||
std::span<const uint8_t> request_pdu = command.frame.pdu();
|
||||
// Is it an error response?
|
||||
if (helpers::is_function_code_exception(function_code)) {
|
||||
uint8_t exception = pdu[1]; // exception frames are fixed-length, so the code is always present
|
||||
ESP_LOGW(TAG,
|
||||
"Error function code: 0x%X exception: %" PRIu8 ", address: %" PRIu8 ", %" PRIu32 "ms after last send",
|
||||
function_code, exception, address, this->last_modbus_byte_ - this->last_send_);
|
||||
if (device)
|
||||
device->on_error(request_pdu, static_cast<ExceptionCode>(exception));
|
||||
} else if (device) { // Not an error response
|
||||
device->on_response(request_pdu, pdu);
|
||||
} else { // Not an error response, but no device to respond to
|
||||
ESP_LOGV(TAG, "Ignoring response from %" PRIu8 " - no callback device set, %" PRIu32 "ms after last send",
|
||||
address, this->last_modbus_byte_ - this->last_send_);
|
||||
}
|
||||
}
|
||||
// Deliver at parse time so the response span can point into the rx buffer (zero copy). error()/
|
||||
// response() set the state and consume the request BEFORE the callback, so a clear from inside it
|
||||
// ("stop polling now") wins. A device-less shell runs no callback and the sweep erases it.
|
||||
this->waiting_for_response_ = false;
|
||||
this->sweep_needed_ = true;
|
||||
if (helpers::is_function_code_exception(function_code)) {
|
||||
uint8_t exception = pdu[1]; // exception frames are fixed-length, so the code is always present
|
||||
ESP_LOGW(TAG, "Error function code: 0x%X exception: %" PRIu8 ", address: %" PRIu8 ", %" PRIu32 "ms after last send",
|
||||
function_code, exception, address, this->last_modbus_byte_ - this->last_send_);
|
||||
cmd->error(static_cast<ExceptionCode>(exception));
|
||||
} else if (!cmd->response(pdu)) {
|
||||
ESP_LOGV(TAG, "Ignoring response from %" PRIu8 " - no callback device set, %" PRIu32 "ms after last send", address,
|
||||
this->last_modbus_byte_ - this->last_send_);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -460,21 +477,20 @@ void ModbusServerHub::process_modbus_client_frame_(uint8_t address, uint8_t func
|
||||
}
|
||||
}
|
||||
|
||||
// Callers gate on tx_blocked() first, but the pre-send delay below can span several ms, so re-check
|
||||
// after it and refuse (return false) if a byte arrived in that window rather than transmit over it.
|
||||
bool Modbus::send_frame_(const ModbusFrame &frame) {
|
||||
if (this->tx_blocked()) {
|
||||
ESP_LOGE(TAG, "Attempted to send while transmission blocked");
|
||||
return false;
|
||||
}
|
||||
if (frame.size() > MAX_FRAME_SIZE) {
|
||||
ESP_LOGE(TAG, "Attempted to send frame larger than max frame size of %" PRIu16 " bytes", MAX_FRAME_SIZE);
|
||||
return false;
|
||||
}
|
||||
|
||||
const int32_t tx_delay_remaining = this->tx_delay_remaining();
|
||||
if (tx_delay_remaining > 0) {
|
||||
delay(tx_delay_remaining);
|
||||
}
|
||||
|
||||
// The delay above can span several ms; a byte arriving in that window blocks transmission after the
|
||||
// caller's gate already passed. Don't collide with the incoming frame - leave the entry to retry.
|
||||
if (this->tx_blocked()) {
|
||||
return false;
|
||||
}
|
||||
|
||||
if (this->flow_control_pin_ != nullptr) {
|
||||
this->flow_control_pin_->digital_write(true);
|
||||
this->write_array(frame.data.data(), frame.size());
|
||||
@@ -498,37 +514,21 @@ bool Modbus::send_frame_(const ModbusFrame &frame) {
|
||||
}
|
||||
|
||||
void ModbusClientHub::send_next_frame_() {
|
||||
if (this->tx_buffer_.empty()) {
|
||||
if (this->tx_blocked())
|
||||
return;
|
||||
|
||||
ModbusDeviceCommand *cmd = this->select_next_ready_();
|
||||
if (cmd == nullptr)
|
||||
return;
|
||||
|
||||
if (!this->send_frame_(cmd->frame)) {
|
||||
ESP_LOGV(TAG, "Send deferred for %" PRIu8 ": a frame arrived during the send delay, will retry",
|
||||
cmd->frame.address());
|
||||
return;
|
||||
}
|
||||
|
||||
if (this->tx_blocked()) {
|
||||
return;
|
||||
}
|
||||
|
||||
// Move the command out and pop BEFORE attempting the send: no callback may run while the frame still
|
||||
// sits in the queue (the same principle as the clear sweep). A failure callback that sends would
|
||||
// otherwise queue a new frame and pop_front() could discard the wrong one - and the deque
|
||||
// reference / PDU span could be invalidated mid-callback.
|
||||
ModbusDeviceCommand command = std::move(this->tx_buffer_.front());
|
||||
this->tx_buffer_.pop_front();
|
||||
ModbusClientDevice *device = command.device;
|
||||
const bool sent = this->send_frame_(command.frame);
|
||||
|
||||
if (sent) {
|
||||
// The frame now lives in the waiting slot; its PDU is the frame without the leading address and
|
||||
// trailing CRC.
|
||||
ModbusDeviceCommand &wfr = this->waiting_for_response_.emplace(std::move(command));
|
||||
if (device != nullptr)
|
||||
device->on_sent(wfr.frame.pdu());
|
||||
} else {
|
||||
if (device != nullptr)
|
||||
device->trigger_not_sent(command.frame.pdu());
|
||||
}
|
||||
|
||||
if (!this->tx_buffer_.empty()) {
|
||||
ESP_LOGV(TAG, "Write queue contains %zu items.", this->tx_buffer_.size());
|
||||
}
|
||||
cmd->sent();
|
||||
this->waiting_for_response_ = true;
|
||||
}
|
||||
|
||||
void ModbusClientHub::dump_config() {
|
||||
@@ -579,118 +579,258 @@ void ModbusServerHub::send_exception_(uint8_t address, uint8_t function_code, Ex
|
||||
this->send_raw_(raw_frame, 3);
|
||||
}
|
||||
|
||||
void ModbusClientHub::notify_no_response_(ModbusDeviceCommand &wfr) {
|
||||
if (wfr.device == nullptr)
|
||||
return;
|
||||
const bool retry = wfr.device->on_no_response(wfr.frame.pdu());
|
||||
// The callback may have detached the device (e.g. clear_tx_queue_for_device()); honor the detach
|
||||
// over the retry request rather than re-queueing a frame that can no longer be routed.
|
||||
if (retry && wfr.device != nullptr)
|
||||
this->requeue_waiting_frame_(wfr);
|
||||
// The old transaction is over either way; never deliver anything else to the device through it.
|
||||
wfr.device = nullptr;
|
||||
ModbusDeviceCommand *ModbusClientHub::find_waiting_() {
|
||||
for (auto &cmd : this->tx_buffer_) {
|
||||
if (cmd.waiting_state())
|
||||
return &cmd;
|
||||
}
|
||||
return nullptr;
|
||||
}
|
||||
|
||||
void ModbusClientHub::requeue_waiting_frame_(ModbusDeviceCommand &wfr) {
|
||||
const ModbusFrame &frame = wfr.frame;
|
||||
if (this->tx_buffer_.size() >= MODBUS_TX_BUFFER_SIZE) {
|
||||
ESP_LOGE(TAG, "Write buffer full, dropped retry for address %" PRIu8, frame.address());
|
||||
if (wfr.device != nullptr)
|
||||
wfr.device->trigger_not_sent(frame.pdu());
|
||||
ModbusDeviceCommand *ModbusClientHub::select_next_ready_() {
|
||||
// Class first (WRITE, then one-shot READ, then CONTINUOUS), oldest within a class. seq is a
|
||||
// free-running counter, so compare each entry's AGE against it (correct across the full range).
|
||||
const uint16_t now = this->next_seq_;
|
||||
const auto age = [now](const ModbusDeviceCommand &cmd) -> uint16_t { return now - cmd.seq; };
|
||||
const auto older = [&age](const ModbusDeviceCommand &a, const ModbusDeviceCommand &b) { return age(a) > age(b); };
|
||||
ModbusDeviceCommand *best = nullptr;
|
||||
for (auto &cmd : this->tx_buffer_) {
|
||||
if (cmd.state != FrameState::READY)
|
||||
continue;
|
||||
if (best == nullptr || cmd.priority() > best->priority() ||
|
||||
(cmd.priority() == best->priority() && older(cmd, *best))) {
|
||||
best = &cmd;
|
||||
}
|
||||
}
|
||||
return best;
|
||||
}
|
||||
|
||||
bool ModbusDeviceCommand::sent() {
|
||||
this->state = FrameState::WAITING;
|
||||
// on_sent() is not a terminal, so nothing is consumed.
|
||||
if (this->device == nullptr)
|
||||
return false;
|
||||
this->device->on_sent(this->frame.pdu());
|
||||
return true;
|
||||
}
|
||||
|
||||
bool ModbusDeviceCommand::notify_retired() {
|
||||
if (!this->decrement_pending())
|
||||
return false; // nothing owed - stop the sweep draining this entry
|
||||
if (this->device != nullptr)
|
||||
this->device->on_not_sent(this->frame.pdu());
|
||||
return true; // consumed one debt (delivered, or silent when device-less) - keep draining to zero
|
||||
}
|
||||
|
||||
bool ModbusDeviceCommand::response(std::span<const uint8_t> response_pdu) {
|
||||
this->state = this->state == FrameState::WAITING_RETIRED ? FrameState::RETIRED : FrameState::RECEIVED_RESPONSE;
|
||||
// A continuous poll is never consumed by its own response; a one-shot consumes one request here.
|
||||
if (!this->continuous)
|
||||
this->decrement_pending();
|
||||
if (this->device == nullptr)
|
||||
return false;
|
||||
this->device->on_response(this->frame.pdu(), response_pdu);
|
||||
return true;
|
||||
}
|
||||
|
||||
bool ModbusDeviceCommand::error(ExceptionCode exception_code) {
|
||||
this->state = this->state == FrameState::WAITING_RETIRED ? FrameState::RETIRED : FrameState::RECEIVED_EXCEPTION;
|
||||
// An exception ends a continuous poll too, so decrement unconditionally.
|
||||
this->decrement_pending();
|
||||
if (this->device == nullptr)
|
||||
return false;
|
||||
this->device->on_error(this->frame.pdu(), exception_code);
|
||||
return true;
|
||||
}
|
||||
|
||||
bool ModbusDeviceCommand::interrupt() {
|
||||
// An unexpected frame distrusts the transaction. A cleared-but-still-waiting shell distrusts too, so
|
||||
// the interrupt survives the clear in either order (WAITING_RETIRED -> INTERRUPTED_RETIRED).
|
||||
if (this->state == FrameState::WAITING) {
|
||||
this->state = FrameState::INTERRUPTED;
|
||||
return true;
|
||||
}
|
||||
if (this->state == FrameState::WAITING_RETIRED) {
|
||||
this->state = FrameState::INTERRUPTED_RETIRED;
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
bool ModbusDeviceCommand::timed_out() {
|
||||
this->state = FrameState::TIMED_OUT; // advance BEFORE the callback so a clear from inside it wins
|
||||
this->decrement_pending(); // resolve this request (WAITING-origin, so pending >= 1)
|
||||
if (this->device == nullptr)
|
||||
return false; // resolved, no one to tell
|
||||
if (this->device->on_no_response(this->frame.pdu()))
|
||||
this->increment_pending(); // granted retry = re-request (capped)
|
||||
return true;
|
||||
}
|
||||
|
||||
void ModbusClientHub::sweep_() {
|
||||
if (!this->sweep_needed_)
|
||||
return;
|
||||
this->sweep_needed_ = false;
|
||||
// Serve only the entries present now: a callback may append (a re-send), but those sit beyond
|
||||
// work_set and are left for the next sweep, which bounds the work and is the termination argument.
|
||||
// Entries leave the container only in the erase pass below, so indices/references stay valid.
|
||||
const size_t work_set = this->tx_buffer_.size();
|
||||
// Restart the walk after every callback: a handler may have moved any entry to any state.
|
||||
bool callback_ran = true;
|
||||
while (callback_ran) {
|
||||
callback_ran = false;
|
||||
for (size_t i = 0; i != work_set && !callback_ran; i++) {
|
||||
ModbusDeviceCommand &cmd = this->tx_buffer_[i];
|
||||
switch (cmd.state) {
|
||||
case FrameState::RECEIVED_RESPONSE:
|
||||
case FrameState::RECEIVED_EXCEPTION:
|
||||
case FrameState::TIMED_OUT:
|
||||
// Off the wire, callback already delivered: reschedule what is still pending, else erase.
|
||||
if (cmd.pending)
|
||||
cmd.requeue(this->next_seq_++);
|
||||
break;
|
||||
case FrameState::RETIRED:
|
||||
// Owes one on_not_sent() per accepted request; notify_retired() consumes one and reports
|
||||
// whether a debt remained, so the restart loop drains the entry to zero - even a device-less
|
||||
// shell with pending > 1 (no callback fires, but it still drains rather than stranding).
|
||||
callback_ran = cmd.notify_retired();
|
||||
break;
|
||||
case FrameState::WAITING_RETIRED:
|
||||
case FrameState::INTERRUPTED_RETIRED:
|
||||
// Cleared shell: drain only the un-run duplicates; the request in flight keeps pending 1
|
||||
// and gets its usual callback when it resolves.
|
||||
if (cmd.pending > 1)
|
||||
callback_ran = cmd.notify_retired();
|
||||
break;
|
||||
default: // READY / WAITING / INTERRUPTED: idle or waiting for a response, nothing owed until the timeout
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
// Erase pass: the only place entries leave the container. Storage order carries no meaning, so a
|
||||
// finished entry is swap-and-popped; walking backwards means a moved-down entry is already seen.
|
||||
for (size_t i = this->tx_buffer_.size(); i-- > 0;) {
|
||||
const ModbusDeviceCommand &cmd = this->tx_buffer_[i];
|
||||
// pending == 0 is erasable, but shells still waiting for a response are exempt until it resolves.
|
||||
if (cmd.pending != 0 || cmd.waiting_state())
|
||||
continue;
|
||||
if (i + 1 != this->tx_buffer_.size())
|
||||
this->tx_buffer_[i] = std::move(this->tx_buffer_.back());
|
||||
this->tx_buffer_.pop_back();
|
||||
}
|
||||
// Re-queue a copy (not a move): the waiting entry may have to survive as an interrupted shell.
|
||||
this->tx_buffer_.emplace_back(wfr.device, frame.address(), frame.pdu());
|
||||
}
|
||||
|
||||
// Raw send for client: pushes to tx queue. Everything except the CRC must be contained in payload.
|
||||
void ModbusClientHub::send_pdu(uint8_t address, std::span<const uint8_t> pdu, ModbusClientDevice *device) {
|
||||
bool ModbusClientHub::send_pdu(uint8_t address, std::span<const uint8_t> pdu, ModbusClientDevice *device,
|
||||
CommandOptions options) {
|
||||
// Requests refused here never enter the machine and get no callback - the false return is it.
|
||||
if (pdu.empty()) {
|
||||
if (device != nullptr)
|
||||
device->trigger_not_sent(pdu);
|
||||
return;
|
||||
ESP_LOGW(TAG, "Empty PDU refused for address %" PRIu8, address);
|
||||
return false;
|
||||
}
|
||||
|
||||
// Bound the PDU so the wire frame (address + pdu + CRC) stays within the Modbus RTU 256-byte limit.
|
||||
if (pdu.size() > MAX_PDU_SIZE) {
|
||||
ESP_LOGE(TAG, "Frame too large, dropped: %" PRIu8 ":%zu bytes", address, pdu.size());
|
||||
if (device != nullptr)
|
||||
device->trigger_not_sent(pdu);
|
||||
return;
|
||||
ESP_LOGE(TAG, "Frame too large, refused: %" PRIu8 ":%zu bytes", address, pdu.size());
|
||||
return false;
|
||||
}
|
||||
|
||||
if (this->tx_buffer_.size() < MODBUS_TX_BUFFER_SIZE) {
|
||||
#if ESPHOME_LOG_LEVEL >= ESPHOME_LOG_LEVEL_VERBOSE
|
||||
char hex_buf[format_hex_pretty_size(MODBUS_MAX_LOG_BYTES)];
|
||||
#endif
|
||||
ESP_LOGV(TAG, "Adding frame to tx queue: %" PRIu8 ":%s", address,
|
||||
format_hex_pretty_to(hex_buf, pdu.data(), pdu.size()));
|
||||
this->tx_buffer_.emplace_back(device, address, pdu);
|
||||
} else {
|
||||
// continuous is ignored for every mutating code (re-writing a value forever is never intended).
|
||||
const bool mutates = ModbusDeviceCommand::classify(pdu[0]) == CommandPriority::WRITE;
|
||||
bool continuous = false;
|
||||
if (options.continuous) {
|
||||
if (mutates) {
|
||||
ESP_LOGV(TAG, "continuous is ignored for a mutating function (0x%X, address %" PRIu8 ")", pdu[0], address);
|
||||
} else {
|
||||
continuous = true;
|
||||
}
|
||||
}
|
||||
|
||||
// A duplicate of a live entry with the same owner is not queued twice; it resolves against that
|
||||
// entry: anonymous -> dropped; continuous incoming -> convert the entry to a poll; one-shot onto a
|
||||
// poll -> downgrade the poll to one-shot; both one-shots -> pending++ below the cap, else refused.
|
||||
for (auto &item : this->tx_buffer_) {
|
||||
if (item.state == FrameState::RETIRED || item.state == FrameState::WAITING_RETIRED ||
|
||||
item.state == FrameState::INTERRUPTED_RETIRED)
|
||||
continue; // cleared, on their way out: a new identical send queues fresh, never absorbs
|
||||
if (item.device != device || !item.same_frame(address, pdu))
|
||||
continue;
|
||||
if (device == nullptr) {
|
||||
// A dropped read is routine (DEBUG); a dropped write/custom warns (unobservable without a device).
|
||||
const bool requeueable = !helpers::is_function_code_exception(pdu[0]) && helpers::is_function_code_read(pdu[0]);
|
||||
if (requeueable) {
|
||||
ESP_LOGD(TAG, "Anonymous duplicate of active frame for %" PRIu8 " (function 0x%X), dropped", address, pdu[0]);
|
||||
} else {
|
||||
ESP_LOGW(TAG,
|
||||
"Anonymous duplicate of active frame for %" PRIu8 " (function 0x%X), dropped - register a "
|
||||
"device for delivery accounting",
|
||||
address, pdu[0]);
|
||||
}
|
||||
return false; // dropped: no entry, no callbacks - the refusal is the return value
|
||||
}
|
||||
if (continuous) {
|
||||
item.make_continuous(true);
|
||||
ESP_LOGV(TAG, "Frame already active for %" PRIu8 ", now polled continuously", address);
|
||||
} else if (item.continuous) {
|
||||
// A one-shot duplicate downgrades the poll to a one-shot: it runs one more cycle to serve this
|
||||
// request, then stops (mirrors continuous incoming converting a one-shot the other way).
|
||||
item.make_continuous(false);
|
||||
ESP_LOGV(TAG, "Frame already active for %" PRIu8 ", downgraded from continuous to one-shot", address);
|
||||
} else if (!item.increment_pending()) {
|
||||
// At the servable cap, so refused. (An absorbed duplicate leaves seq alone - the entry keeps
|
||||
// its place in line, held by its oldest outstanding request.)
|
||||
ESP_LOGD(TAG, "Frame already active for %" PRIu8 " with %" PRIu8 " requests pending, refused", address,
|
||||
item.pending);
|
||||
return false;
|
||||
} else {
|
||||
ESP_LOGV(TAG, "Frame already active for %" PRIu8 ", request absorbed (pending %" PRIu8 ")", address,
|
||||
item.pending);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
// Backstop counts every entry; dead ones are gone by the sweep's end, so at worst they cost one
|
||||
// refusal at the very cap for one loop.
|
||||
if (this->tx_buffer_.size() >= MODBUS_TX_BUFFER_SIZE) {
|
||||
#if ESPHOME_LOG_LEVEL >= ESPHOME_LOG_LEVEL_ERROR
|
||||
char hex_buf[format_hex_pretty_size(MODBUS_MAX_LOG_BYTES)];
|
||||
#endif
|
||||
ESP_LOGE(TAG, "Write buffer full, dropped: %" PRIu8 ":%s", address,
|
||||
ESP_LOGE(TAG, "Write buffer full, refused: %" PRIu8 ":%s", address,
|
||||
format_hex_pretty_to(hex_buf, pdu.data(), pdu.size()));
|
||||
if (device != nullptr)
|
||||
device->trigger_not_sent(pdu);
|
||||
return false;
|
||||
}
|
||||
#if ESPHOME_LOG_LEVEL >= ESPHOME_LOG_LEVEL_VERBOSE
|
||||
char hex_buf[format_hex_pretty_size(MODBUS_MAX_LOG_BYTES)];
|
||||
#endif
|
||||
ESP_LOGV(TAG, "Adding frame to tx queue: %" PRIu8 ":%s", address,
|
||||
format_hex_pretty_to(hex_buf, pdu.data(), pdu.size()));
|
||||
this->tx_buffer_.emplace_back(device, address, pdu, continuous, this->next_seq_++);
|
||||
return true;
|
||||
}
|
||||
|
||||
void ModbusClientHub::clear_tx_queue_for_address(uint8_t address, bool clear_sent) {
|
||||
// Drop the queued frames for this address, delivering on_not_sent() to each frame's owner: other
|
||||
// devices talking to the same physical device (e.g. a modbus_client action alongside a controller that
|
||||
// just went offline) must observe the drop, or their command never resolves. Mark first, then sweep
|
||||
// only marked frames: anything a callback re-queues is unmarked, so it
|
||||
// is never swept - or re-notified - by the clear that triggered it. Each marked frame is moved out and erased BEFORE
|
||||
// its callback runs, so handlers see a consistent queue; termination is guaranteed because only the initially-marked
|
||||
// frames are ever swept.
|
||||
void ModbusClientHub::clear_tx_queue_for_address(uint8_t address) {
|
||||
// A clear is a pure state flip; the sweep delivers every owed on_not_sent() from a quiescent hub.
|
||||
for (auto &cmd : this->tx_buffer_) {
|
||||
if (cmd.frame.address() == address)
|
||||
cmd.marked_for_deletion = true;
|
||||
}
|
||||
for (;;) {
|
||||
auto it = std::find_if(this->tx_buffer_.begin(), this->tx_buffer_.end(),
|
||||
[](const ModbusDeviceCommand &cmd) { return cmd.marked_for_deletion; });
|
||||
if (it == this->tx_buffer_.end())
|
||||
break;
|
||||
ModbusDeviceCommand dropped = std::move(*it);
|
||||
this->tx_buffer_.erase(it);
|
||||
// The sweep delivers through the same per-device guard as refusals: a device clearing from inside
|
||||
// its own on_not_sent() gets its remaining frames resolved silently (documented in the lifecycle
|
||||
// contract), other owners are notified normally, and every nested clear stays bounded.
|
||||
if (dropped.device != nullptr)
|
||||
dropped.device->trigger_not_sent(dropped.frame.pdu());
|
||||
}
|
||||
|
||||
if (clear_sent && this->waiting_for_response_.has_value() && this->waiting_for_response_.value().device) {
|
||||
if (this->waiting_for_response_.value().frame.address() == address) {
|
||||
ESP_LOGV(TAG, "Clearing waiting for response for address %" PRIu8, address);
|
||||
// Invalidate the waiting device so it won't process a response.
|
||||
this->waiting_for_response_.value().device = nullptr;
|
||||
}
|
||||
if (cmd.frame.address() != address)
|
||||
continue;
|
||||
cmd.retire();
|
||||
this->sweep_needed_ = true;
|
||||
}
|
||||
}
|
||||
void ModbusClientHub::clear_tx_queue_for_device(ModbusClientDevice *device) {
|
||||
// Remove any pending commands for this address from the tx buffer
|
||||
auto &tx_buffer = this->tx_buffer_;
|
||||
tx_buffer.erase(std::remove_if(tx_buffer.begin(), tx_buffer.end(),
|
||||
[device](const ModbusDeviceCommand &cmd) { return cmd.device == device; }),
|
||||
tx_buffer.end());
|
||||
|
||||
if (this->waiting_for_response_.has_value() && this->waiting_for_response_.value().device) {
|
||||
if (this->waiting_for_response_.value().device == device) {
|
||||
ESP_LOGV(TAG, "Clearing waiting for response");
|
||||
// Invalidate the waiting device so it won't process a response.
|
||||
this->waiting_for_response_.value().device = nullptr;
|
||||
}
|
||||
void ModbusClientHub::clear_tx_queue_for_device(ModbusClientDevice *device) {
|
||||
// Silent teardown (supersede semantics): the caller's own frames vanish without callbacks; see
|
||||
// the lifecycle note on ModbusClientDevice.
|
||||
for (auto &cmd : this->tx_buffer_) {
|
||||
if (cmd.device != device)
|
||||
continue;
|
||||
cmd.silent_retire();
|
||||
this->sweep_needed_ = true;
|
||||
}
|
||||
}
|
||||
|
||||
void ModbusClientHub::send_raw(const std::vector<uint8_t> &payload, ModbusClientDevice *device) {
|
||||
if (payload.size() < 2) {
|
||||
if (device != nullptr)
|
||||
device->trigger_not_sent({}); // too short to contain a PDU
|
||||
ESP_LOGW(TAG, "send_raw() payload too short to contain a PDU, refused");
|
||||
return;
|
||||
}
|
||||
this->send_pdu(payload[0], std::span<const uint8_t>(payload).subspan(1), device);
|
||||
@@ -706,23 +846,26 @@ void ModbusServerHub::send_raw_(const uint8_t *payload, uint16_t len) {
|
||||
return;
|
||||
}
|
||||
|
||||
// In the rare case that the server is blocked (frame delay has not elapsed), we delay the send.
|
||||
// This should only happen at low baud rates with long frame delays.
|
||||
// If blocked now (frame delay not elapsed at low baud, or a frame arriving), defer rather than
|
||||
// busy-waiting the loop; send_frame_ itself re-checks after its delay, so the deferred callback
|
||||
// just reports whatever it returns.
|
||||
if (this->tx_blocked()) {
|
||||
// Stash the raw payload in a single member buffer so the deferred callback can rebuild the frame
|
||||
// without a heap allocation. Only one server reply is ever in flight, and the named timeout ensures
|
||||
// only one deferred send is pending, so a single buffer is sufficient.
|
||||
// without a heap allocation. Only one server reply is ever waiting, so a single buffer suffices.
|
||||
std::memcpy(this->deferred_payload_.data(), payload, len);
|
||||
this->deferred_payload_len_ = len;
|
||||
this->set_timeout("deferred_send", this->tx_delay_remaining(), [this]() {
|
||||
ModbusFrame frame(this->deferred_payload_[0], this->deferred_payload_.data() + 1,
|
||||
this->deferred_payload_len_ - 1);
|
||||
this->send_frame_(frame);
|
||||
if (!this->send_frame_(frame))
|
||||
ESP_LOGE(TAG, "Deferred server reply dropped: transmission still blocked");
|
||||
});
|
||||
} else {
|
||||
ModbusFrame frame(payload[0], payload + 1, len - 1);
|
||||
this->send_frame_(frame);
|
||||
return;
|
||||
}
|
||||
|
||||
ModbusFrame frame(payload[0], payload + 1, len - 1);
|
||||
if (!this->send_frame_(frame))
|
||||
ESP_LOGE(TAG, "Server reply dropped: a frame arrived during the send delay");
|
||||
}
|
||||
|
||||
void Modbus::clear_rx_buffer_(const LogString *reason, bool warn, size_t bytes_to_clear) {
|
||||
@@ -778,7 +921,7 @@ void ModbusClientDevice::dispatch_response_(std::span<const uint8_t> request_pdu
|
||||
const bool bits =
|
||||
function_code == FunctionCode::READ_COILS || function_code == FunctionCode::READ_DISCRETE_INPUTS;
|
||||
const size_t expected_data_size =
|
||||
bits ? (static_cast<size_t>(count_or_value) + 7) / 8 : static_cast<size_t>(count_or_value) * 2;
|
||||
bits ? packed_bit_bytes(count_or_value) : static_cast<size_t>(count_or_value) * 2;
|
||||
if (response_pdu.size() != expected_data_size + 2) {
|
||||
ESP_LOGD(TAG, "Response length %zu does not match request (expected %zu) for function code 0x%X",
|
||||
response_pdu.size(), expected_data_size + 2, static_cast<uint8_t>(function_code));
|
||||
|
||||
Reference in New Issue
Block a user