#include "UDPServer.h" UDPServer::~UDPServer() { socket_->close(); } void UDPServer::init_socket() { static std::mutex init_mutex; std::lock_guard lock(init_mutex); // 确保 io_context 没有在运行 if (!io_context_.stopped()) { // 暂停 io_context io_context_.stop(); std::this_thread::sleep_for(std::chrono::milliseconds(10)); } // 创建 socket socket_ = std::make_unique(io_context_); // 绑定端口 asio::error_code ec; socket_->open(udp::v4(), ec); if (!ec) { socket_->bind(udp::endpoint(udp::v4(), nPort_), ec); } if (ec) { LOG_ERROR("Failed to bind socket: " + ec.message()); socket_.reset(); return; } // 重启 io_context io_context_.restart(); // 开始接收 start_receive(); } void UDPServer::setCMDQueue(HDDSCMDQueue* pQueue) { m_pCMDQueue = pQueue; } void UDPServer::setSendParam(const char* hostIP, int nPort) { // 延迟创建resolver if (!resolver_) { resolver_ = std::make_unique(io_context_); } char cPort[30]; snprintf(cPort, sizeof(cPort), "%d", nPort); endpoint_ = resolver_->resolve(udp::v4(), hostIP, cPort); } void UDPServer::start_receive() { socket_->async_receive_from( asio::buffer(data, max_length), remote_endpoint_, [this](std::error_code ec, std::size_t bytes_recvd) { this->handle_receive(ec, bytes_recvd); } ); } void UDPServer::handle_receive(std::error_code ec, std::size_t bytes_recvd) { if (!ec && bytes_recvd > 0) { HDDSCMD pCMD; if (bytes_recvd <= sizeof(HDDSCMD)) { memcpy(&pCMD, data, sizeof(HDDSCMD)); if (m_pCMDQueue) { m_pCMDQueue->push(pCMD); } } else { LOG_WARNING("Received data too small for HDDSCMD structure"); } } else if (ec) { LOG_ERROR("Receive error: " + ec.message()); } start_receive(); } int UDPServer::blockingSendData(LPHDDSCMD pData, long DataLen) { if (DataLen > max_length + 12) { LOG_WARNING("Data too large to send"); return -1; } memcpy(sendData, pData, DataLen); try { return socket_->send_to(asio::buffer(sendData, DataLen), *endpoint_.begin()); } catch (std::exception& e) { LOG_ERROR("Send error: " + std::string(e.what())); return -1; } } int UDPServer::asyncSendData(LPHDDSCMD pData, long DataLen) { // 异步发送实现 return 1; } // 解析单个目标地址 bool UDPServer::resolveTarget(UDPTarget& target) { if (!resolver_) { resolver_ = std::make_unique(io_context_); } try { char cPort[30]; snprintf(cPort, sizeof(cPort), "%d", target.port); auto endpoints = resolver_->resolve(udp::v4(), target.ip, cPort); if (endpoints.begin() != endpoints.end()) { target.endpoint = *endpoints.begin(); target.resolved = true; return true; } } catch (const std::exception& e) { LOG_ERROR("Failed to resolve target " + target.ip + ":" + std::to_string(target.port) + " - " + std::string(e.what())); } return false; } // 解析所有未解析的目标 void UDPServer::resolveAllTargets() { std::lock_guard lock(targets_mutex_); for (auto& target : multi_targets_) { if (!target.resolved) { resolveTarget(target); } } } // 新增接口:设置多个目标 void UDPServer::setMultiSendParams(const std::vector& targets) { std::lock_guard lock(targets_mutex_); multi_targets_.clear(); multi_targets_ = targets; // 解析所有目标地址 resolveAllTargets(); LOG_INFO("Set " + std::to_string(multi_targets_.size()) + " multi-targets"); } void UDPServer::addSendTarget(const char* hostIP, int nPort) { std::lock_guard lock(targets_mutex_); // 检查是否已存在 for (const auto& target : multi_targets_) { if (target.ip == hostIP && target.port == nPort) { LOG_WARNING("Target already exists: " + std::string(hostIP) + ":" + std::to_string(nPort)); return; } } // 添加新目标 UDPTarget new_target; new_target.ip = hostIP; new_target.port = nPort; new_target.resolved = false; // 立即解析 if (resolveTarget(new_target)) { multi_targets_.push_back(new_target); LOG_INFO("Added target: " + std::string(hostIP) + ":" + std::to_string(nPort)); } else { LOG_ERROR("Failed to add target: " + std::string(hostIP) + ":" + std::to_string(nPort)); } } void UDPServer::removeSendTarget(const char* hostIP, int nPort) { std::lock_guard lock(targets_mutex_); auto it = std::remove_if(multi_targets_.begin(), multi_targets_.end(), [hostIP, nPort](const UDPTarget& target) { return target.ip == hostIP && target.port == nPort; }); if (it != multi_targets_.end()) { multi_targets_.erase(it, multi_targets_.end()); LOG_INFO("Removed target: " + std::string(hostIP) + ":" + std::to_string(nPort)); } else { LOG_WARNING("Target not found: " + std::string(hostIP) + ":" + std::to_string(nPort)); } } void UDPServer::clearSendTargets() { std::lock_guard lock(targets_mutex_); multi_targets_.clear(); LOG_INFO("Cleared all send targets"); } int UDPServer::blockingSendDataToMulti(LPHDDSCMD pData, long DataLen) { if (DataLen > max_length + 12) { LOG_WARNING("Data too large to send"); return -1; } std::lock_guard lock(targets_mutex_); if (multi_targets_.empty()) { LOG_WARNING("No multi-targets set"); return 0; } memcpy(sendData, pData, DataLen); int total_sent = 0; int success_count = 0; for (const auto& target : multi_targets_) { if (target.resolved) { try { int sent = socket_->send_to(asio::buffer(sendData, DataLen), target.endpoint); total_sent += sent; success_count++; } catch (std::exception& e) { LOG_ERROR("Send to " + target.ip + ":" + std::to_string(target.port) + " error: " + std::string(e.what())); } } else { LOG_WARNING("Target not resolved: " + target.ip + ":" + std::to_string(target.port)); } } LOG_INFO("Sent to " + std::to_string(success_count) + "/" + std::to_string(multi_targets_.size()) + " targets"); return total_sent; } // 新增接口:异步发送到多个目标 int UDPServer::asyncSendDataToMulti(LPHDDSCMD pData, long DataLen) { if (DataLen > max_length + 12) { LOG_WARNING("Data too large to send"); return -1; } std::lock_guard lock(targets_mutex_); if (multi_targets_.empty()) { LOG_WARNING("No multi-targets set"); return 0; } // 复制数据到新缓冲区,确保在异步发送期间有效 auto send_buffer = std::make_shared>(DataLen); memcpy(send_buffer->data(), pData, DataLen); int success_count = 0; for (const auto& target : multi_targets_) { if (target.resolved) { try { socket_->async_send_to( asio::buffer(*send_buffer), target.endpoint, [this, send_buffer, target](std::error_code ec, std::size_t bytes_sent) { if (!ec) { LOG_INFO("Async sent " + std::to_string(bytes_sent) + " bytes to " + target.ip + ":" + std::to_string(target.port)); } else { LOG_ERROR("Async send to " + target.ip + ":" + std::to_string(target.port) + " error: " + ec.message()); } } ); success_count++; } catch (std::exception& e) { LOG_ERROR("Async send to " + target.ip + ":" + std::to_string(target.port) + " error: " + std::string(e.what())); } } else { LOG_WARNING("Target not resolved: " + target.ip + ":" + std::to_string(target.port)); } } LOG_INFO("Async sending to " + std::to_string(success_count) + "/" + std::to_string(multi_targets_.size()) + " targets"); return success_count; }