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

