【问题标题】:Protobuf, CodedInputStream parsing partial messagesProtobuf, CodedInputStream 解析部分消息
【发布时间】:2014-10-30 14:31:28
【问题描述】:

我正在尝试实现与 java 版本兼容的 protobuf 发送/接收,其中首先包含 varint32-prefix。

我几乎可以让它工作了,但由于某种原因,一些消息变得不完整,导致 assert() 失败。

/receiver.cpp:69: void tcp_connection::handle_read_message(const boost::system::error_code&, size_t): Assertion `line.ParseFromCodedStream(&input)' failed.

semder.pp

boost::asio::streambuf buffer;
std::ostream writer(&buffer);
bool packet_full = false;
uint32_t sent_lines = 0;
{ //new scope for protobuf streams, these flush in dtor
    google::protobuf::io::OstreamOutputStream osostream(&writer);
    google::protobuf::io::CodedOutputStream output(&osostream);
    std::string lines;
    while(std::getline(reader, line)) {
        lines += line + "\n";
        ++sent_lines;
        if(sent_lines > 100) {
            packet_full = true;
            break;
        }
    }
    if(!lines.empty()) {
        msg->set_text(lines);
        const uint32_t size = msg->ByteSize();
        output.WriteVarint32(size);
        uint8_t* buffer = output.GetDirectBufferForNBytesAndAdvance(size);
        if(buffer != 0) {
            msg->SerializeWithCachedSizesToArray(buffer);
        } else {
            msg->SerializeWithCachedSizes(&output);
        }
}
if(sent_lines > 0) {
    sock.send(buffer.data());
    if(!packet_full && !reader.eof()) { //Read ended, and not due to end of file
        std::cout << "An error occured" << std::endl;
        break;
    }
    reader.clear(); //clear EOF flag
}

receiver.cpp

这是一个 boost asio 回调。

成员变量:

boost::asio::ip::tcp::socket socket_;
boost::asio::streambuf buffer_;

代码

void handle_read_message(const boost::system::error_code& error,
                      size_t bytes_transferred) {


   if(!error) {
      buffer_.commit(bytes_transferred);
      std::istream reader(&buffer_);
      google::protobuf::io::IstreamInputStream isistream(&reader);
      google::protobuf::io::CodedInputStream input(&isistream);
      uint32_t size = 0;
      assert(input.ReadVarint32(&size));
      auto limit = input.PushLimit(size);
      msgs::Line line;
      assert(line.ParseFromCodedStream(&input));
      assert(input.ConsumedEntireMessage());
      input.PopLimit(limit);

      start();  
    } else {
      std::cout <<"error during handle_read_message: " << error << std::endl;
    }
}

这个主要是基于https://stackoverflow.com/a/22899712

编辑: 新的接收器版本,reader_ 现在是成员变量:

void handle_read_message(const boost::system::error_code& error,
                          size_t bytes_transferred) {
    std::cout << "handle_read_message(" << bytes_transferred << ")" <<std::endl;
    if(!error) {
      buffer_.commit(bytes_transferred);
      uint32_t size = 0;
      google::protobuf::io::IstreamInputStream isistream_(&reader_);
      {
        google::protobuf::io::CodedInputStream input(&isistream_);
        if(!input.ReadVarint32(&size)) {
          std::cout << "Failed to read size, waiting for more data" << std::endl;
          start();
          return;
        }
      }
      std::size_t varint_size = isistream_.ByteCount();
      std::cout <<"varintsize: " << varint_size << ", size: " << size << ", have bytes: " << buffer_.size() << std::endl;
      if(varint_size + size > buffer_.size()) {
        std::cout << "Not enough data received, waiting for more" << std::endl;
        start();
        return;
      }
      google::protobuf::io::CodedInputStream input(&isistream_);
      auto limit = input.PushLimit(size);
      msgs::Line line;
      assert(line.ParseFromCodedStream(&input));
      std::cout << line.text() << std::endl;
      assert(input.ConsumedEntireMessage());
      input.PopLimit(limit);

      start();  
    } else {
      std::cout <<"error during handle_read_message: " << error << std::endl;
    } 
  }

【问题讨论】:

    标签: c++ boost-asio protocol-buffers


    【解决方案1】:

    如果您在接收端使用异步 I/O,则需要确保在开始解析之前确实收到了整个消息。请记住,TCP 连接是一个流。只要有可用数据,异步回调就会运行——即使它不完整。您可能只收到部分消息,或者您可能收到整条消息加上下一条消息。这就是为什么首先需要 readDelimitedFrom() 的原因:在解析之前准确计算出需要等待多少字节。

    因此,在使用异步 I/O 时,您需要以不同的方式编写代码。你可以使用这样的策略:

    • 维护一个包含您目前收到的所有字节的缓冲区。
    • 每次收到更多字节时,将它们添加到缓冲区中。然后,开始尝试按如下方式解析它们 - 您必须始终从头开始,使用全新的 ZeroCopyInputStream 和 CodedInputStream。
    • 然后,尝试使用ReadVarint32() 读取大小。如果 ReadVarint32 失败,那么您还没有收到整个大小,因此请停止并等待更多字节。
    • 如果ReadVarint32() 成功,则销毁CodedInputStream,然后在底层ZeroCopyInputStream 上调用ByteCount() 以了解varint 消耗了多少字节。
    • 您现在知道了消息的大小和 varint 前缀的大小。把它们加在一起。如果缓冲区中的字节数少于此数量,请停下来等待更多。
    • 您现在拥有消息的所有字节。继续将它们从缓冲区中拉出并解析它们。请注意,如果缓冲区中的字节数超过消息的大小,则应将多余的字节留在缓冲区中,因为它们是下一条消息的一部分。

    (另外:您的 sender.cpp 代码中似乎缺少右括号。如果原始文件有相同的错误,则可能是您在 CodedOutputStream 刷新之前发送数据。但我猜错误不在原件中。)

    【讨论】:

    • 感谢您写得很好的答案,我已经提出了一个新版本,我认为它涵盖了您的观点。你能看透吗?我仍然在 ParseFromCodedStream 上遇到断言崩溃。
    • 嗯,我没有看到明显的问题。我建议在两端转储原始字节并进行比较。 IE。在发送端使用SerializeAsString(),打印它产生的字节的十六进制版本,然后在接收端打印缓冲区内容的十六进制版本,然后再解析它。如果它们不相同,则说明在传输过程中出现问题,您可以开始缩小范围。
    猜你喜欢
    • 1970-01-01
    • 2014-04-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-05-20
    • 1970-01-01
    • 1970-01-01
    • 2012-02-21
    相关资源
    最近更新 更多