#include "plc_communication_service.h" #include "plc_communication_error_classifier.h" #include "plc_register_repository.h" #include #include #include #include #include #include #include #include namespace { constexpr int kRecoveryProbeIntervalMs = 2000; std::string toUtf8(const QString &value) { const QByteArray bytes = value.toUtf8(); return std::string(bytes.constData(), static_cast(bytes.size())); } QModbusDataUnit::RegisterType registerType(RegisterArea area) { return area == RegisterArea::M ? QModbusDataUnit::Coils : QModbusDataUnit::HoldingRegisters; } bool isRecoverableTimeout(PlcCommunicationError error) { return error == PlcCommunicationError::PlcNotResponding || error == PlcCommunicationError::CommunicationTimeout; } bool isReadingState(PlcConnectionState state) { return state == PlcConnectionState::Connected || state == PlcConnectionState::Recovering; } bool normalizePollAddresses( const std::vector &addresses, std::vector *normalized, std::string *error) { *normalized = addresses; if (std::any_of( normalized->cbegin(), normalized->cend(), [](const RegisterAddress &address) { return !address.isValid(); })) { *error = "PLC 轮询地址中包含无效的 M/D 地址"; return false; } std::sort( normalized->begin(), normalized->end(), [](const RegisterAddress &left, const RegisterAddress &right) { return left.area() == right.area() ? left.index() < right.index() : left.area() == RegisterArea::M; }); normalized->erase( std::unique(normalized->begin(), normalized->end()), normalized->end()); if (normalized->size() > ProjectLimits::kMaximumPollAddresses) { *error = "PLC 轮询的去重 M/D 地址最多为 1024 个"; return false; } return true; } std::size_t pollBlockCount(const std::vector &addresses) { if (addresses.empty()) { return 2U; } std::size_t blocks = 0U; RegisterArea current_area = RegisterArea::M; int start_address = 0; int count = 0; for (const RegisterAddress &address : addresses) { if (count == 0 || address.area() != current_area || address.index() > start_address + count || count >= ProjectLimits::kMaximumModbusReadCount) { ++blocks; current_area = address.area(); start_address = address.index(); count = 1; } else { count = address.index() - start_address + 1; } } return blocks; } } // namespace PlcCommunicationService::PlcCommunicationService( PlcRegisterRepository &repository, QObject *parent) : QObject(parent), repository_(repository), master_(std::make_unique()) { repository_.setWriteHandlers( [this](const RegisterAddress &address, bool value) { return sendBitWrite(address, value); }, [this](const RegisterAddress &address, std::int16_t value) { return sendWordWrite(address, value); }); connect(&poll_timer_, &QTimer::timeout, this, &PlcCommunicationService::pollNextBlock); recovery_timer_.setSingleShot(true); connect( &recovery_timer_, &QTimer::timeout, this, &PlcCommunicationService::probeRecovery); connect(master_.get(), &QModbusClient::stateChanged, this, [this](QModbusDevice::State device_state) { if (device_state == QModbusDevice::ConnectedState) { serial_session_opened_ = true; setState(PlcConnectionState::Connected); poll_timer_.start(configuration_.pollIntervalMs); pollNextBlock(); } else if (device_state == QModbusDevice::ConnectingState) { serial_session_opened_ = false; setState(PlcConnectionState::Connecting); } else if (device_state == QModbusDevice::UnconnectedState) { if (disconnecting_) { serial_session_opened_ = false; } else if (serial_session_opened_) { handleUnexpectedDisconnect(); } else if (state_ != PlcConnectionState::Faulted) { setState(PlcConnectionState::Disconnected); } } }); connect(master_.get(), &QModbusClient::errorOccurred, this, [this](QModbusDevice::Error error) { if (error != QModbusDevice::NoError) { handleModbusError(error); } }); } PlcCommunicationService::~PlcCommunicationService() = default; PlcCommunicationResult PlcCommunicationService::connectDevice( const PlcSerialConfiguration &configuration) { const PlcCommunicationResult validation = validatePlcSerialConfiguration(configuration); if (!validation.succeeded) { return validation; } if (state_ == PlcConnectionState::Disconnected && master_->state() != QModbusDevice::UnconnectedState) { closeSerialSession(); } if (master_->state() != QModbusDevice::UnconnectedState) { return {false, "PLC 连接已经启动,请先断开当前连接"}; } // 每次连接都递增代次;旧连接的异步回复即使晚到也不能污染新缓存 configuration_ = configuration; ++connection_generation_; disconnecting_ = false; serial_session_opened_ = false; received_valid_response_ = false; last_error_type_ = PlcCommunicationError::None; last_error_.clear(); master_->setConnectionParameter( QModbusDevice::SerialPortNameParameter, QString::fromStdString(configuration.portName)); master_->setConnectionParameter( QModbusDevice::SerialBaudRateParameter, configuration.baudRate); master_->setConnectionParameter( QModbusDevice::SerialDataBitsParameter, static_cast(configuration.dataBits)); master_->setConnectionParameter( QModbusDevice::SerialParityParameter, static_cast(configuration.parity)); master_->setConnectionParameter( QModbusDevice::SerialStopBitsParameter, static_cast(configuration.stopBits)); master_->setTimeout(configuration.responseTimeoutMs); master_->setNumberOfRetries(configuration.retries); // 新连接必须重新完成完整首读,不能沿用上一条连接的运行资格 repository_.invalidate(); updateInitialReadCompleted(false, true); rebuildPollBlocks(); setState(PlcConnectionState::Connecting); if (!master_->connectDevice()) { if (last_error_.empty()) { const QModbusDevice::Error error = master_->error() == QModbusDevice::NoError ? QModbusDevice::ConnectionError : master_->error(); handleModbusError(error); } return {false, last_error_}; } return {true, {}}; } void PlcCommunicationService::disconnectDevice() { closeSerialSession(); repository_.invalidate(); updateInitialReadCompleted(false, true); last_error_type_ = PlcCommunicationError::None; last_error_.clear(); setState(PlcConnectionState::Disconnected); } PlcCommunicationResult PlcCommunicationService::setPollAddresses( const std::vector &addresses) { std::vector normalized; std::string error; if (!normalizePollAddresses(addresses, &normalized, &error)) { return {false, error}; } if (pollBlockCount(normalized) > ProjectLimits::kMaximumPollBlocks) { return {false, "PLC 轮询地址拆分后最多允许 64 个读块"}; } // 当前读请求完成前暂存新集合,避免按半旧半新的地址解析回复 if (pending_reply_ != nullptr) { pending_poll_addresses_ = std::move(normalized); poll_update_pending_ = true; return {true, {}}; } applyPollAddresses(normalized); return {true, {}}; } PlcConnectionState PlcCommunicationService::state() const { return state_; } bool PlcCommunicationService::initialReadCompleted() const { return initial_read_completed_; } PlcCommunicationError PlcCommunicationService::lastErrorType() const { return last_error_type_; } const std::string &PlcCommunicationService::lastError() const { return last_error_; } const PlcSerialConfiguration &PlcCommunicationService::configuration() const { return configuration_; } void PlcCommunicationService::setCallbacks( std::function state_changed, std::function initial_read_changed, std::function cache_updated, std::function error_reported) { state_changed_callback_ = std::move(state_changed); initial_read_changed_callback_ = std::move(initial_read_changed); cache_updated_callback_ = std::move(cache_updated); error_reported_callback_ = std::move(error_reported); } void PlcCommunicationService::rebuildPollBlocks() { std::vector addresses = poll_addresses_; if (addresses.empty()) { addresses = { RegisterAddress{RegisterArea::M, 0}, RegisterAddress{RegisterArea::D, 0}}; } std::sort( addresses.begin(), addresses.end(), [](const RegisterAddress &left, const RegisterAddress &right) { if (left.area() != right.area()) { return left.area() == RegisterArea::M; } return left.index() < right.index(); }); addresses.erase(std::unique(addresses.begin(), addresses.end()), addresses.end()); // 将相邻地址合并为读块,同时遵守 Modbus 单次读取上限 poll_blocks_.clear(); for (const RegisterAddress &address : addresses) { if (!address.isValid()) { continue; } if (poll_blocks_.empty() || poll_blocks_.back().area != address.area() || address.index() > poll_blocks_.back().startAddress + poll_blocks_.back().count || poll_blocks_.back().count >= ProjectLimits::kMaximumModbusReadCount) { poll_blocks_.push_back({address.area(), address.index(), 1}); } else { poll_blocks_.back().count = address.index() - poll_blocks_.back().startAddress + 1; } } next_poll_block_ = 0; if (!initial_read_completed_) { initial_blocks_read_.assign(poll_blocks_.size(), false); } else { initial_blocks_read_.clear(); } } void PlcCommunicationService::applyPollAddresses( const std::vector &addresses) { poll_addresses_ = addresses; rebuildPollBlocks(); poll_update_pending_ = false; pending_poll_addresses_.clear(); if (isReadingState(state_) && pending_reply_ == nullptr) { pollNextBlock(); } } void PlcCommunicationService::pollNextBlock() { if (!isReadingState(state_) || pending_reply_ != nullptr || poll_blocks_.empty()) { return; } // 每次只发一个异步请求,完成回调中再推进到下一个块 const std::size_t block_index = next_poll_block_; const PollBlock block = poll_blocks_.at(block_index); next_poll_block_ = (next_poll_block_ + 1U) % poll_blocks_.size(); QModbusDataUnit request( registerType(block.area), block.startAddress, static_cast(block.count)); QModbusReply *reply = master_->sendReadRequest(request, configuration_.serverAddress); if (reply == nullptr) { handleModbusError(master_->error()); return; } pending_reply_ = reply; const std::uint64_t generation = connection_generation_; connect(reply, &QModbusReply::finished, this, [this, reply, block, block_index, generation] { // 连接已经重建时,直接丢弃旧回复 if (generation != connection_generation_) { reply->deleteLater(); return; } handleReadFinished(reply, block); if (isReadingState(state_) && reply->error() == QModbusDevice::NoError && !poll_update_pending_ && block_index < initial_blocks_read_.size()) { initial_blocks_read_[block_index] = true; const bool completed = std::all_of( initial_blocks_read_.cbegin(), initial_blocks_read_.cend(), [](bool read) { return read; }); // 所有轮询块都成功读过,才授予真机运行资格 if (completed && !initial_read_completed_) { updateInitialReadCompleted(true); if (state_ == PlcConnectionState::Recovering) { setState(PlcConnectionState::Connected); } } } if (pending_reply_ == reply) { pending_reply_ = nullptr; } reply->deleteLater(); if (poll_update_pending_) { const std::vector addresses = pending_poll_addresses_; applyPollAddresses(addresses); } }); } void PlcCommunicationService::probeRecovery() { if (state_ != PlcConnectionState::Faulted || !isRecoverableTimeout(last_error_type_)) { return; } if (master_->state() != QModbusDevice::ConnectedState) { return; } if (pending_reply_ != nullptr) { recovery_timer_.start(kRecoveryProbeIntervalMs); return; } if (poll_blocks_.empty()) { return; } const PollBlock block = poll_blocks_.front(); const QModbusDataUnit request( registerType(block.area), block.startAddress, 1); QModbusReply *reply = master_->sendReadRequest(request, configuration_.serverAddress); if (reply == nullptr) { handleRecoveryProbeFailure(master_->error()); return; } pending_reply_ = reply; const std::uint64_t generation = connection_generation_; connect(reply, &QModbusReply::finished, this, [this, reply, generation] { if (generation != connection_generation_) { reply->deleteLater(); return; } if (pending_reply_ == reply) { pending_reply_ = nullptr; } const QModbusDevice::Error error = reply->error(); reply->deleteLater(); if (poll_update_pending_) { const std::vector addresses = pending_poll_addresses_; applyPollAddresses(addresses); } if (state_ != PlcConnectionState::Faulted || !isRecoverableTimeout(last_error_type_)) { return; } if (error != QModbusDevice::NoError) { handleRecoveryProbeFailure(error); return; } restoreCommunication(); }); } void PlcCommunicationService::handleRecoveryProbeFailure(QModbusDevice::Error error) { if (state_ != PlcConnectionState::Faulted || !isRecoverableTimeout(last_error_type_)) { return; } if (error == QModbusDevice::ConnectionError) { const PlcCommunicationErrorContext context{ QString::fromStdString(configuration_.portName), serial_session_opened_, received_valid_response_}; setError(classifyPlcCommunicationError(error, context)); return; } recovery_timer_.start(kRecoveryProbeIntervalMs); } void PlcCommunicationService::restoreCommunication() { received_valid_response_ = true; last_error_type_ = PlcCommunicationError::None; last_error_.clear(); rebuildPollBlocks(); setState(PlcConnectionState::Recovering); poll_timer_.start(configuration_.pollIntervalMs); pollNextBlock(); } void PlcCommunicationService::handleReadFinished(QModbusReply *reply, PollBlock block) { if (!isReadingState(state_)) { return; } if (reply->error() != QModbusDevice::NoError) { handleModbusError(reply->error()); return; } // 只有成功回复才能更新 PLC 缓存;写请求不会直接改缓存 const QModbusDataUnit result = reply->result(); for (uint index = 0; index < result.valueCount(); ++index) { const int address = block.startAddress + static_cast(index); if (block.area == RegisterArea::M) { repository_.updateBit(address, result.value(index) != 0U); } else { repository_.updateWord(address, static_cast(result.value(index))); } } received_valid_response_ = true; last_error_type_ = PlcCommunicationError::None; last_error_.clear(); emit cacheUpdated(); if (cache_updated_callback_) { cache_updated_callback_(); } } RegisterWriteResult PlcCommunicationService::sendBitWrite( const RegisterAddress &address, bool value) { if (state_ != PlcConnectionState::Connected) { return {false, RegisterError::Unavailable}; } if (pending_write_reply_ != nullptr) { return {false, RegisterError::WriteRejected}; } QModbusDataUnit unit(QModbusDataUnit::Coils, address.index(), 1); unit.setValue(0, value ? 1U : 0U); QModbusReply *reply = master_->sendWriteRequest(unit, configuration_.serverAddress); if (reply == nullptr) { handleModbusError(master_->error()); return {false, RegisterError::WriteRejected}; } pending_write_reply_ = reply; connect(reply, &QModbusReply::finished, this, [this, reply, generation = connection_generation_] { if (pending_write_reply_ == reply) { pending_write_reply_ = nullptr; } if (generation != connection_generation_) { reply->deleteLater(); return; } if (reply->error() != QModbusDevice::NoError) { handleModbusError(reply->error()); } reply->deleteLater(); }); return {true, RegisterError::None}; } RegisterWriteResult PlcCommunicationService::sendWordWrite( const RegisterAddress &address, std::int16_t value) { if (state_ != PlcConnectionState::Connected) { return {false, RegisterError::Unavailable}; } if (pending_write_reply_ != nullptr) { return {false, RegisterError::WriteRejected}; } QModbusDataUnit unit(QModbusDataUnit::HoldingRegisters, address.index(), 1); unit.setValue(0, static_cast(value)); QModbusReply *reply = master_->sendWriteRequest(unit, configuration_.serverAddress); if (reply == nullptr) { handleModbusError(master_->error()); return {false, RegisterError::WriteRejected}; } pending_write_reply_ = reply; connect(reply, &QModbusReply::finished, this, [this, reply, generation = connection_generation_] { if (pending_write_reply_ == reply) { pending_write_reply_ = nullptr; } if (generation != connection_generation_) { reply->deleteLater(); return; } if (reply->error() != QModbusDevice::NoError) { handleModbusError(reply->error()); } reply->deleteLater(); }); return {true, RegisterError::None}; } void PlcCommunicationService::updateInitialReadCompleted( bool completed, bool force_notification) { if (initial_read_completed_ == completed && !force_notification) { return; } initial_read_completed_ = completed; emit initialReadCompletedChanged(completed); if (initial_read_changed_callback_) { initial_read_changed_callback_(completed); } } void PlcCommunicationService::handleUnexpectedDisconnect() { serial_session_opened_ = false; const QString port_name = QString::fromStdString(configuration_.portName).trimmed(); setError({ PlcCommunicationError::SerialConnectionLost, QStringLiteral( "PLC 串口 %1 连接已中断;请检查 USB 转串口是否被拔出或已经失效") .arg(port_name)}); } void PlcCommunicationService::handleModbusError(QModbusDevice::Error error) { if (disconnecting_ || (error == QModbusDevice::ReplyAbortedError && state_ == PlcConnectionState::Disconnected) || (state_ == PlcConnectionState::Disconnected && last_error_type_ == PlcCommunicationError::SerialConnectionLost) || state_ == PlcConnectionState::Faulted) { return; } const PlcCommunicationErrorContext context{ QString::fromStdString(configuration_.portName), serial_session_opened_, received_valid_response_}; setError(classifyPlcCommunicationError(error, context)); } void PlcCommunicationService::closeSerialSession() { ++connection_generation_; disconnecting_ = true; poll_timer_.stop(); recovery_timer_.stop(); pending_reply_ = nullptr; pending_write_reply_ = nullptr; poll_update_pending_ = false; pending_poll_addresses_.clear(); if (master_->state() != QModbusDevice::UnconnectedState) { master_->disconnectDevice(); } serial_session_opened_ = false; received_valid_response_ = false; disconnecting_ = false; } void PlcCommunicationService::setState(PlcConnectionState state) { if (state_ == state) { return; } state_ = state; emit stateChanged(); if (state_changed_callback_) { state_changed_callback_(); } } void PlcCommunicationService::setError(const PlcCommunicationFailure &failure) { last_error_type_ = failure.type; last_error_ = toUtf8(failure.message); poll_timer_.stop(); recovery_timer_.stop(); // 任意通信故障都会撤销首读资格;恢复后必须重新完整读取 updateInitialReadCompleted(false); const bool disconnected = failure.type == PlcCommunicationError::SerialPortOpenFailed || failure.type == PlcCommunicationError::SerialConnectionLost || failure.type == PlcCommunicationError::UsbSerialAdapterRemoved; if (disconnected) { closeSerialSession(); } setState(disconnected ? PlcConnectionState::Disconnected : PlcConnectionState::Faulted); emit communicationError(failure.message); if (error_reported_callback_) { error_reported_callback_(last_error_); } if (isRecoverableTimeout(failure.type) && master_->state() == QModbusDevice::ConnectedState) { recovery_timer_.start(kRecoveryProbeIntervalMs); } }