You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Boost.Asio客户端中async_read_until与async_write的同时使用方案及即时通讯异常排查

Let's break down your problem and walk through the key issues in your code, along with actionable fixes:

1. Critical Client-side Issue: io_context Isn't Running

The missing '{}' in message_wait_loop logs tell us your async read callbacks are never firing. This is almost certainly because your AsyncTCPClient::m_ios (the asio io_context) isn't running an active event loop. Asio requires this loop to execute any asynchronous operations like async_read_until or async_write—without it, all async work sits idle.

Fix:
Add a dedicated thread to run the io_context in your AsyncTCPClient constructor:

class AsyncTCPClient: public asio::noncopyable {
public:
    AsyncTCPClient() {
        // Start the io_context event loop in a background thread
        io_thread = std::make_unique<std::thread>([this]() {
            m_ios.run();
        });
        io_thread->detach(); // Adjust to join if you need explicit cleanup on shutdown
    }

    ~AsyncTCPClient() {
        // Cleanup: stop the event loop and join the thread if possible
        m_ios.stop();
        if (io_thread && io_thread->joinable()) {
            io_thread->join();
        }
    }

private:
    asio::io_context m_ios;
    std::unique_ptr<std::thread> io_thread; // Add this member to track the loop thread
    // ... existing members ...
};

2. Broken Server-side Message Routing Logic

In Service::onMessageReceived, your check for forwarding messages is completely inverted. Per your architecture, current_sessions maps user IDs to their Service IDs—but you're comparing the recipient's Service ID to the sender's user ID, which will never match. You also aren't checking the open chat sessions you mentioned in your Tracker design.

Fix:
First, add the missing open chat tracking to your Tracker struct:

struct Tracker {
    static std::mutex current_sessions_guard;
    static std::map<long, long> current_sessions; // user_id -> service_id
    static std::map<long, int> client_to_service_id;
    // Track open bidirectional chats (store unordered pairs to avoid duplication)
    static std::set<std::pair<long, long>> open_chats;
};

Then fix the routing logic in onMessageReceived:

void Service::onMessageReceived() {
    std::istream istrm(m_message.get());
    std::string new_message;
    std::getline(istrm, new_message);
    m_message.reset(new asio::streambuf);

    std::unique_lock<std::mutex> tracker_lock(Tracker::current_sessions_guard);

    // Check if the recipient is online
    auto recipient_session_it = Tracker::current_sessions.find(another_party_id);
    if (recipient_session_it != Tracker::current_sessions.end()) {
        // Check if the recipient has an open chat with the sender
        bool chat_is_open = false;
        // Check both orders since chats are bidirectional
        if (Tracker::open_chats.count({client_id, another_party_id}) || 
            Tracker::open_chats.count({another_party_id, client_id})) {
            chat_is_open = true;
        }

        if (chat_is_open) {
            // Get the recipient's Service ID and forward the message
            int recipient_service_id = Tracker::client_to_service_id[another_party_id];
            std::string formatted_msg = _form_message_str(login, new_message);
            spdlog::info("[{}] Forwarding message to service {}: '{}'", 
                         service_id, recipient_service_id, new_message);
            
            auto recipient_service_it = Server::launched_services.find(recipient_service_id);
            if (recipient_service_it != Server::launched_services.end()) {
                recipient_service_it->second->send_to_chat(formatted_msg);
            }
        }
    }

    tracker_lock.unlock();
    receive_message();
}

3. Missing Error Handling for Async Operations

Your code doesn't handle errors from asio's async calls, which can cause silent failures (like broken sockets stopping message flow without any logs).

Fixes:

  • Update Service::send_to_chat to handle write errors:
void send_to_chat(const std::string& new_message) {
    asio::async_write(*m_sock.get(), asio::buffer(new_message), 
        [this](const std::error_code& ec, std::size_t bytes_sent) {
            if (ec) {
                spdlog::error("[{}] Failed to send message: {}", service_id, ec.message());
                // Optional: Clean up the session (remove from Tracker, close socket)
                return;
            }
            onAnotherPartyMessageSent();
        });
}
  • Update AsyncTCPClient::message_wait_loop to handle read errors:
void message_wait_loop(std::shared_ptr<Session> session) {
    logger->info("'{}' in message_wait_loop", session->login);
    std::istream istrm(session->m_received_message.get());
    std::string received_message;
    std::getline(istrm, received_message);
    session->m_received_message.reset(new asio::streambuf);

    std::unique_lock<std::mutex> lock_std_out(std_out_wait_guard);
    session->current_chat.append(received_message + "\n");
    lock_std_out.unlock();

    clear_console();
    m_console.write(session->current_chat);
    m_console.write("Write your message: ");

    // Add error handling to the async read
    asio::async_read_until(session->m_sock, *(session->m_received_message.get()), "\n", 
        [this, session] (const std::error_code& ec, std::size_t bytes_read) {
            if (ec) {
                logger->error("[{}] Failed to receive message: {}", session->login, ec.message());
                return;
            }
            message_wait_loop(session);
        });
}

4. Client-side Thread Management Improvements

Using detached threads for msg_wait_thread and msg_thread can lead to resource leaks or unexpected behavior. Instead, use asio's post to schedule operations on the io_context thread, keeping all network work on a single thread to avoid thread safety issues.

Alternative Fix:
Replace the detached threads in onSentReady with asio post:

void onSentReady(std::shared_ptr<Session> session) {
    // Schedule wait loop on the io_context thread
    asio::post(m_ios, [this, session]() {
        message_wait_loop(session);
    });
    // Schedule send loop on the io_context thread
    asio::post(m_ios, [this, session]() {
        message_send_loop(session);
    });
}

内容的提问来源于stack exchange,提问作者Ivan Ivan

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.29 05:28:17