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

如何使用Boost async_read处理带长度前缀无终止符的Protobuf消息?

异步SSL客户端读取Protobuf消息的问题与解决方案

问题背景

你有Java和C#开发经验,现在尝试用Boost.Asio实现一个基于异步Socket的SSL客户端,用来和服务器收发Protobuf消息。消息格式是前4字节表示后续消息的长度X,剩余X字节是Protobuf消息内容。

你遇到的核心问题是:Boost的async_read()需要指定精确的读取长度,你知道可以先读4字节的长度头,但不知道后续怎么正确读取剩下的消息内容;而async_read_until()需要终止符,显然这里不适用。你的源代码如下(可以搜索TODO定位读取逻辑相关部分):

#include <boost/asio.hpp>
#include <boost/asio/ssl.hpp>
#include <boost/bind.hpp>
#include <iostream>
#include <istream>
#include <ostream>
#include <string>
#include "client/connection/authentication.pb.h"
#include "client/connection/authentication.pb.cc"
#include "client/msg.pb.h"
#include "client/msg.pb.cc"
#include "client/common.pb.h"
#include "client/common.pb.cc"
class client {
public:
client(boost::asio::io_service& io_service, boost::asio::ssl::context& context, boost::asio::ip::tcp::resolver::iterator endpoint_iterator)
: socket_(io_service, context) {
socket_.set_verify_mode(boost::asio::ssl::context::verify_none);
socket_.set_verify_callback(boost::bind(&client::verify_certificate, this, _1, _2));
boost::asio::async_connect(socket_.lowest_layer(), endpoint_iterator,
boost::bind(&client::handle_connect, this, boost::asio::placeholders::error));
}
bool verify_certificate(bool preverified, boost::asio::ssl::verify_context& ctx) {
char subject_name[256];
X509* cert = X509_STORE_CTX_get_current_cert(ctx.native_handle());
X509_NAME_oneline(X509_get_subject_name(cert), subject_name, 256);
std::cout << "Verifying:\n" << subject_name << std::endl;
return preverified;
}
void handle_connect(const boost::system::error_code& error) {
if(!error){
std::cout << "Connection OK!" << std::endl;
socket_.async_handshake(boost::asio::ssl::stream_base::client,
boost::bind(&client::handle_handshake, this, boost::asio::placeholders::error));
}else{
std::cout << "Connect failed: " << error.message() << std::endl;
}
}
void handle_handshake(const boost::system::error_code& error) {
if(!error){
std::cout << "Sending request: " << std::endl;
// std::stringstream request_;
// request_ << "GET /api/0/data/ticker.php HTTP 1.1\r\n";
// request_ << "Host: mtgox.com\r\n";
// request_ << "Accept-Encoding: *\r\n";
// request_ << "\r\n";
protobuf::Message msg;
char *data = new char[dlen];
bool ok = msg.SerializeToArray(data,dlen);
// uint32_t n = htonl(dlen);
uint32_t n = dlen;
// char bytes[4];
// bytes[0] = (n >> 24) & 0xFF;
// bytes[1] = (n >> 16) & 0xFF;
// bytes[2] = (n >> 8) & 0xFF;
// bytes[3] = n & 0xFF;
int sizeOfPacket = 4 + dlen;
char* rq = new char[sizeOfPacket];
// strncpy(rq, bytes, 4);
rq[0] = (n >> 24) & 0xFF;
rq[1] = (n >> 16) & 0xFF;
rq[2] = (n >> 8) & 0xFF;
rq[3] = n & 0xFF;
strncpy(rq + 4,data, dlen);
// request_<< rq;
// std::cerr << request_.str() << std::endl;
boost::asio::async_write(socket_, boost::asio::buffer(rq,sizeOfPacket),
boost::bind(&client::handle_write, this, boost::asio::placeholders::error,
boost::asio::placeholders::bytes_transferred));
}else{
std::cout << "Handshake failed: " << error.message() << std::endl;
}
}
void handle_write(const boost::system::error_code& error, size_t bytes_transferred) {
if (!error){
std::cout << "Sending request OK!" << std::endl;
char respond[4] = "";
boost::asio::async_read(socket_, boost::asio::buffer(respond,4),
boost::bind(&client::handle_read, this, boost::asio::placeholders::error,
boost::asio::placeholders::bytes_transferred));
// std::cerr << "respond is " << respond;
//TODO
}else{
std::cout << "Write failed: " << error.message() << std::endl;
}
}
void handle_read(const boost::system::error_code& error, size_t bytes_transferred) {
if (!error){
std::cout << "Reply: ";
std::cout.write(reply_, bytes_transferred);
std::cout << "\n";
}else{
std::cout << "Read failed: " << error.message() << std::endl;
}
}
private:
boost::asio::ssl::stream<boost::asio::ip::tcp::socket> socket_;
char reply_[0x1 << 16];
};
int main(int argc, char* argv[]) {
try{
boost::asio::io_service io_service;
boost::asio::ip::tcp::resolver resolver(io_service);
boost::asio::ip::tcp::resolver::query query("192.168.2.32", "443");
boost::asio::ip::tcp::resolver::iterator iterator = resolver.resolve(query);
boost::asio::ssl::context context(boost::asio::ssl::context::sslv23);
// context.load_verify_file("key.pem");
client c(io_service, context, iterator);
io_service.run();
}catch (std::exception& e){
std::cerr << "Exception: " << e.what() << "\n";
}
std::cin.get();
return 0;
}

你的尝试(存在的问题)

你试着修改了handle_write函数,实现思路是先读4字节头,解析出消息长度,再读对应长度的内容,但这段代码有致命的异步逻辑错误:

void handle_write(const boost::system::error_code& error, size_t bytes_transferred) {
if (!error){
std::cout << "Sending request OK!" << std::endl;
char respond[4] = "";
boost::asio::async_read(socket_, boost::asio::buffer(respond,4),
boost::bind(&client::handle_read, this, boost::asio::placeholders::error,
boost::asio::placeholders::bytes_transferred));
int sizeOfMessage = getSizeFromHeader(respond);//need implement
char message[sizeOfMessage] ="";
boost::asio::async_read(socket_, boost::asio::buffer(message,sizeOfMessage),
boost::bind(&client::handle_read, this, boost::asio::placeholders::error,
boost::asio::placeholders::bytes_transferred));
decodeMessage(message);
// std::cerr << "respond is " << respond;
//TODO
}else{
std::cout << "Write failed: " << error.message() << std::endl;
}
}

问题出在:Boost.Asio的异步操作是非阻塞立即返回的,你发起async_read读取长度头后,马上就去调用getSizeFromHeader(respond),这时候长度头的读取还没完成,respond数组里的内容是随机的垃圾值,后续的读取和解析自然全错了。

异步编程的核心是回调驱动——必须等前一个异步操作完成(进入它的回调函数)后,才能执行后续逻辑。

正确实现方案

我们需要把读取流程拆分成两个异步步骤,用回调函数串联起来:

步骤1:修改Client类,添加必要的成员变量

我们需要存储解析后的消息长度和消息缓冲区,避免在回调里传递过多参数:

class client {
public:
    // ... 原有的构造函数和成员函数 ...
private:
    boost::asio::ssl::stream<boost::asio::ip::tcp::socket> socket_;
    char reply_[0x1 << 16];
    uint32_t message_length_; // 存储解析后的消息长度(主机字节序)
    std::vector<char> message_buffer_; // 动态存储消息内容的缓冲区
};

步骤2:修正handle_write函数,仅发起读取长度头的操作

void handle_write(const boost::system::error_code& error, size_t bytes_transferred) {
    if (!error){
        std::cout << "Sending request OK!" << std::endl;
        // 用vector存储长度头,避免栈溢出和生命周期问题
        std::vector<char> header_buffer(4);
        // 发起异步读取4字节长度头,完成后进入handle_read_header回调
        boost::asio::async_read(socket_, boost::asio::buffer(header_buffer),
            boost::bind(&client::handle_read_header, this, 
                boost::asio::placeholders::error,
                boost::asio::placeholders::bytes_transferred,
                header_buffer));
    }else{
        std::cout << "Write failed: " << error.message() << std::endl;
    }
}

步骤3:实现handle_read_header,解析长度并发起读取消息内容

void handle_read_header(const boost::system::error_code& error, size_t bytes_transferred, std::vector<char> header_buffer) {
    if (!error && bytes_transferred == 4) {
        // 把4字节的网络字节序转成主机字节序(和你发送时的大端写法对应)
        uint32_t network_length = *reinterpret_cast<const uint32_t*>(header_buffer.data());
        message_length_ = ntohl(network_length);

        // 初始化消息缓冲区,大小为解析出的消息长度
        message_buffer_.resize(message_length_);
        // 发起异步读取消息内容,完成后进入handle_read_message回调
        boost::asio::async_read(socket_, boost::asio::buffer(message_buffer_),
            boost::bind(&client::handle_read_message, this,
                boost::asio::placeholders::error,
                boost::asio::placeholders::bytes_transferred));
    } else {
        std::cout << "Read header failed: " << error.message() << std::endl;
    }
}

步骤4:实现handle_read_message,解码并处理Protobuf消息

void handle_read_message(const boost::system::error_code& error, size_t bytes_transferred) {
    if (!error && bytes_transferred == message_length_) {
        std::cout << "Received message of size: " << bytes_transferred << std::endl;
        // 替换成你实际使用的Protobuf消息类型,比如AuthenticationMsg或者Msg
        YourProtobufMessage msg;
        if (msg.ParseFromArray(message_buffer_.data(), message_length_)) {
            std::cout << "Decoded message successfully!" << std::endl;
            // 这里添加你的业务逻辑,比如打印消息字段、处理请求等
            // 如果需要继续收发消息,可以在这里再次发起读取长度头的操作
        } else {
            std::cout << "Failed to parse Protobuf message!" << std::endl;
        }
    } else {
        std::cout << "Read message failed: " << error.message() << std::endl;
    }
}

额外需要注意的细节

  1. 字节序问题:你发送时用了大端字节序((n >> 24) & 0xFF),读取时必须用ntohl()把网络字节序转成主机字节序,否则在小端机器上解析出来的长度会完全错误。
  2. 内存安全:不要用栈上的变长数组(char message[sizeOfMessage]),这是非标准C++写法,容易导致栈溢出,用std::vector<char>是更安全的选择。
  3. Protobuf编译:你现在直接包含.cc文件的做法虽然能编译,但建议把Protobuf文件编译成静态库/动态库链接到项目中,避免重复编译和代码冗余。
  4. SSL证书验证:生产环境下不要用verify_none,建议加载CA证书开启验证,避免中间人攻击。
  5. 错误处理:每个异步回调都要处理错误,比如读取长度头时字节数不足、解析Protobuf失败等情况,确保程序能优雅处理异常。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:26:21