| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296 |
- #include "UDPServer.h"
- UDPServer::~UDPServer() {
- socket_->close();
- }
- void UDPServer::init_socket()
- {
- static std::mutex init_mutex;
- std::lock_guard<std::mutex> 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<udp::socket>(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<udp::resolver>(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<udp::resolver>(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<std::mutex> lock(targets_mutex_);
- for (auto& target : multi_targets_) {
- if (!target.resolved) {
- resolveTarget(target);
- }
- }
- }
- // 新增接口:设置多个目标
- void UDPServer::setMultiSendParams(const std::vector<UDPTarget>& targets) {
- std::lock_guard<std::mutex> 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<std::mutex> 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<std::mutex> 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<std::mutex> 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<std::mutex> 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<std::mutex> lock(targets_mutex_);
- if (multi_targets_.empty()) {
- LOG_WARNING("No multi-targets set");
- return 0;
- }
- // 复制数据到新缓冲区,确保在异步发送期间有效
- auto send_buffer = std::make_shared<std::vector<char>>(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;
- }
|