【问题标题】:boost socket iostreams echo server with zlib compression sleeps until the connection will be closed使用 zlib 压缩的 boost socket iostreams echo 服务器休眠,直到连接关闭
【发布时间】:2021-01-18 10:21:11
【问题描述】:

我尝试按照this 和this 示例创建带有zlib 压缩的简单回显服务器。

我的想法是现在发送一些字符串,因为我可以在发送之前将 POD 类型转换为字符串 (std::string(reinterpret_cast<const char *>(&pod), sizeof(pod))),这样我才能确定传输层工作正常。

这里有一个问题。客户端压缩数据,发送它并说数据已发送,但服务器在读取数据时被阻止。我不明白为什么会这样。

我尝试使用operator<< 和out.flush(),也尝试使用boost::iostreams::copy()。结果是一样的。代码示例是(我根据参数为服务器和客户端使用相同的源文件):

#include <boost/iostreams/filtering_stream.hpp>
#include <boost/iostreams/filter/zlib.hpp>
#include <boost/iostreams/copy.hpp>
#include <boost/asio.hpp>

#include <iostream>
#include <sstream>

namespace ip = boost::asio::ip;
using ip::tcp;

const unsigned short port = 9999;
const char host[] = "127.0.0.1";

void receive()
{
    boost::asio::io_context ctx;
    tcp::endpoint ep(ip::address::from_string(host), port);
    tcp::acceptor a(ctx, ep);

    tcp::iostream stream;
    a.accept(stream.socket());

    std::stringstream buffer;

    std::cout << "start session" << std::endl;

    try
    {
        for (;;)
        {
            {
                boost::iostreams::filtering_istream in;
                in.push(boost::iostreams::zlib_decompressor());
                in.push(stream);

                std::cout << "start reading" << std::endl;

                // looks like server is blocked here
                boost::iostreams::copy(in, buffer);
            }

            std::cout << "data: " << buffer.str() << std::endl;

            {
                boost::iostreams::filtering_ostream out;
                out.push(boost::iostreams::zlib_compressor());
                out.push(stream);

                boost::iostreams::copy(buffer, out);
            }

            std::cout << "Reply is sended" << std::endl;
        }
    }
    catch(const boost::iostreams::zlib_error &e)
    {
        std::cerr << e.what() << e.error() << '\n';
        stream.close();
    }
}

void send(const std::string &data)
{
    tcp::endpoint ep(ip::address::from_string(host), port);
    
    tcp::iostream stream;
    stream.connect(ep);

    std::stringstream buffer;
    buffer << data;

    if (!stream)
    {
        std::cerr << "Cannot connect to " << host << ":" << port << std::endl;
        return;
    }

    try
    {
        {
            boost::iostreams::filtering_ostream out;
            out.push(boost::iostreams::zlib_compressor());
            out.push(stream);

            out << buffer.str();
            out.flush();
        }

        std::cout << "sended: " << data << std::endl;
        buffer.str("");

        {
            boost::iostreams::filtering_istream in;
            in.push(boost::iostreams::zlib_decompressor());
            in.push(stream);

            // looks like client is blocked here
            boost::iostreams::copy(in, buffer);
        }

        std::cout << "result: " << buffer.str() << std::endl;
    }
    catch(const boost::iostreams::zlib_error &e)
    {
        std::cerr << e.what() << '\n';
    }
}

int main(int argc, const char *argv[])
{
    if (argc > 1 && argv[1] ==  std::string("sender"))
        send("hello world");
    else
        receive();

    return 0;
}

首先我启动服务器,然后我启动客户端。产生以下输出:

服务器

$ ./example
# now it waits while client will be accepted
start session
start reading

客户

$ ./example sender
sended: hello world

程序被上面的输出阻塞。我猜服务器仍在等待来自客户端的数据,并且它不知道客户端发送了它所拥有的所有内容。

如果我使用Ctrl + C 关闭客户端,则输出如下:

$ ./example
# now it waits while client will be accepted
start session
start reading
# now it is blocked until I press Ctrl + C
data: hello world
Reply is sended
start reading
zlib error-5

和

$ ./example sender
sended: hello world
^C

我猜zlib error-5 是因为服务器认为存档不完整。

预期行为没有阻塞。当客户端启动时,该消息必须出现在服务器程序输出中。

为什么程序无法读取?我该如何解决?

【问题讨论】:

  • 接收者应该如何知道消息的长度?您必须对某种未压缩的消息长度进行编码,才能知道何时停止阅读并实际解压缩到目前为止收到的内容。
  • 我认为这应该由zlib_compressor 完成。如果它获取数据,那么它会打包数据(我们得到一个存档对象),tcp::iostream 对象通过网络发送存档对象。并且接收者获取具有大小的存档对象,不是吗?如果我错了,请纠正我

标签: c++ boost boost-asio boost-iostreams


【解决方案1】:

iostreams::copy 就是这样做的:它复制流。

赞美您的代码。它非常可读:) 它让我想起了这个答案Reading and writing files with boost iostream socket。主要区别在于该答案发送一个 单个 压缩 blob 并关闭。

你是“正确的”解压缩器知道一个压缩块何时完成,但它不决定另一个不会跟随。

所以你需要添加框架。传统的方法是在带外传递一个长度。我已经实现了这些更改,同时还通过使用 IO 操纵器减少了代码重复。

template <typename T> struct LengthPrefixed {
    T _wrapped;

    friend std::ostream& operator<<(std::ostream& os, LengthPrefixed lp) ;
    friend std::istream& operator>>(std::istream& is, LengthPrefixed lp) ;
};

和

template <typename T> struct ZLIB {
    T& data;
    ZLIB(T& ref) : data(ref){}

    friend std::ostream& operator<<(std::ostream& os, ZLIB z) ;
    friend std::istream& operator>>(std::istream& is, ZLIB z) ;
};

ZLIB机械手

这个主要封装了你在发送方/接收方之间复制的代码:

template <typename T> struct ZLIB {
    T& data;
    ZLIB(T& ref) : data(ref){}

    friend std::ostream& operator<<(std::ostream& os, ZLIB z) {
        {
            boost::iostreams::filtering_ostream out;
            out.push(boost::iostreams::zlib_compressor());
            out.push(os);
            out << z.data << std::flush;
        }
        return os.flush();
    }

    friend std::istream& operator>>(std::istream& is, ZLIB z) {
        boost::iostreams::filtering_istream in;
        in.push(boost::iostreams::zlib_decompressor());
        in.push(is);

        std::ostringstream oss;
        copy(in, oss);
        z.data = oss.str();

        return is;
    }
};

我制作了T 模板,因此它可以根据需要存储std::string&amp; 或std::string const&amp;。

LengthPrefixed机械手

这个操纵器不关心被序列化的内容,而只是简单地在其前面加上有效的在线长度:

template <typename T> struct LengthPrefixed {
    T _wrapped;

    friend std::ostream& operator<<(std::ostream& os, LengthPrefixed lp) {
        std::ostringstream oss;
        oss << lp._wrapped;
        auto on_the_wire = std::move(oss).str();

        debug << "Writing length " << on_the_wire.length() << std::endl;
        return os << on_the_wire.length() << "\n" << on_the_wire << std::flush;
    }

    friend std::istream& operator>>(std::istream& is, LengthPrefixed lp) {
        size_t len;
        if (is >> std::noskipws >> len && is.ignore(1, '\n')) {
            debug << "Reading length " << len << std::endl;

            std::string on_the_wire(len, '\0');
            if (is.read(on_the_wire.data(), on_the_wire.size())) {
                std::istringstream iss(on_the_wire);
                iss >> lp._wrapped;
            }
        }
        return is;
    }
};

我们添加了一个微妙之处:通过根据我们的构造来存储引用或值,我们还可以接受临时变量(如 ZLIB 操纵器):

template <typename T> LengthPrefixed(T&&) -> LengthPrefixed<T>;
template <typename T> LengthPrefixed(T&) -> LengthPrefixed<T&>;

我没想过要让ZLIB 操纵器同样通用。所以我把它作为一个驱魔留给读者

演示程序

结合这两者,您可以将发送者/接收者简单地写为:

void server() {
    boost::asio::io_context ctx;
    tcp::endpoint ep(ip::address::from_string(host), port);
    tcp::acceptor a(ctx, ep);

    tcp::iostream stream;
    a.accept(stream.socket());

    std::cout << "start session" << std::endl;

    for (std::string data; stream >> LengthPrefixed{ZLIB{data}};) {
        std::cout << "data: " << std::quoted(data) << std::endl;
        stream << LengthPrefixed{ZLIB{data}} << std::flush;
    }
}

void client(std::string data) {
    tcp::endpoint ep(ip::address::from_string(host), port);
    tcp::iostream stream(ep);

    stream << LengthPrefixed{ZLIB{data}} << std::flush;
    std::cout << "sent: " << std::quoted(data) << std::endl;

    stream >> LengthPrefixed{ZLIB{data}};
    std::cout << "result: " << std::quoted(data) << std::endl;
}

确实,它会打印:

reader: start session
sender: Writing length 19
reader: Reading length 19
sender: sent: "hello world"
reader: data: "hello world"
reader: Writing length 19
sender: Reading length 19
sender: result: "hello world"

完整清单

#include <boost/iostreams/filtering_stream.hpp>
#include <boost/iostreams/filter/zlib.hpp>
#include <boost/iostreams/copy.hpp>
#include <boost/asio.hpp>

#include <iostream>
#include <iomanip>
#include <sstream>

namespace ip = boost::asio::ip;
using ip::tcp;

const unsigned short port = 9999;
const char host[] = "127.0.0.1";

#ifdef DEBUG
    std::ostream debug(std::cerr.rdbuf());
#else
    std::ostream debug(nullptr);
#endif

template <typename T> struct LengthPrefixed {
    T _wrapped;

    friend std::ostream& operator<<(std::ostream& os, LengthPrefixed lp) {
        std::ostringstream oss;
        oss << lp._wrapped;
        auto on_the_wire = std::move(oss).str();

        debug << "Writing length " << on_the_wire.length() << std::endl;
        return os << on_the_wire.length() << "\n" << on_the_wire << std::flush;
    }

    friend std::istream& operator>>(std::istream& is, LengthPrefixed lp) {
        size_t len;
        if (is >> std::noskipws >> len && is.ignore(1, '\n')) {
            debug << "Reading length " << len << std::endl;

            std::string on_the_wire(len, '\0');
            if (is.read(on_the_wire.data(), on_the_wire.size())) {
                std::istringstream iss(on_the_wire);
                iss >> lp._wrapped;
            }
        }
        return is;
    }
};

template <typename T> LengthPrefixed(T&&) -> LengthPrefixed<T>;
template <typename T> LengthPrefixed(T&) -> LengthPrefixed<T&>;

template <typename T> struct ZLIB {
    T& data;
    ZLIB(T& ref) : data(ref){}

    friend std::ostream& operator<<(std::ostream& os, ZLIB z) {
        {
            boost::iostreams::filtering_ostream out;
            out.push(boost::iostreams::zlib_compressor());
            out.push(os);
            out << z.data << std::flush;
        }
        return os.flush();
    }

    friend std::istream& operator>>(std::istream& is, ZLIB z) {
        boost::iostreams::filtering_istream in;
        in.push(boost::iostreams::zlib_decompressor());
        in.push(is);

        std::ostringstream oss;
        copy(in, oss);
        z.data = oss.str();

        return is;
    }
};

void server() {
    boost::asio::io_context ctx;
    tcp::endpoint ep(ip::address::from_string(host), port);
    tcp::acceptor a(ctx, ep);

    tcp::iostream stream;
    a.accept(stream.socket());

    std::cout << "start session" << std::endl;

    for (std::string data; stream >> LengthPrefixed{ZLIB{data}};) {
        std::cout << "data: " << std::quoted(data) << std::endl;
        stream << LengthPrefixed{ZLIB{data}} << std::flush;
    }
}

void client(std::string data) {
    tcp::endpoint ep(ip::address::from_string(host), port);
    tcp::iostream stream(ep);

    stream << LengthPrefixed{ZLIB{data}} << std::flush;
    std::cout << "sent: " << std::quoted(data) << std::endl;

    stream >> LengthPrefixed{ZLIB{data}};
    std::cout << "result: " << std::quoted(data) << std::endl;
}

int main(int argc, const char**) {
    try {
        if (argc > 1)
            client("hello world");
        else
            server();
    } catch (const std::exception& e) {
        std::cerr << e.what() << '\n';
    }
}

【讨论】:

  • 我将T 模板化,因此它可以根据需要存储std::string&amp; 或std::string const&amp;.. 如果我没记错的话,它可以存储任何类型重载operator&lt;&lt; 和operator&gt;&gt;,不是吗?
  • 你也说过但它并不决定另一个人不会跟随。如果我明白了,operator&gt;&gt;(std::istream &amp;is, LengthPrefixed lp) 可以帮助流了解它必须读取多少数据。如果我错了,请纠正我。
  • 我接受你的回答,因为我同意你的看法。您的回答更笼统,对有相同问题的人更有帮助。
  • “如果我明白了,运算符>>(std::istream &is, LengthPrefixed lp) 帮助流了解它必须读取多少数据” .此外,我们将其输入到长度有限的流中。
  • 我会建议boost::iostreams::restrict 如果不是我知道issues when used recursively 我不确定不会干涉,这就是为什么我通过std::stringstream“容忍”额外的副本跨度>
【解决方案2】:

使用boost::serialization解决问题,步骤如下:

  1. 首先,我将 zipping 移至如下函数:
namespace io = boost::iostreams;

namespace my {
std::string compress(const std::string &data)
{
    std::stringstream input, output;

    input << data;

    io::filtering_ostream io_out;
    io_out.push(io::zlib_compressor());
    io_out.push(output);

    io::copy(input, io_out);

    return output.str();
}

std::string decompress(const std::string &data)
{
    std::stringstream input, output;

    input << data;

    io::filtering_istream io_in;
    io_in.push(io::zlib_decompressor());
    io_in.push(input);

    io::copy(io_in, output);

    return output.str();
}
} // namespace my
  1. 然后我为这样的字符串缓冲区创建了一个包装器(按照documentation 的教程):
class Package
{
public:
    Package(const std::string &buffer) : buffer(buffer) {}

private:
    std::string buffer;

    friend class boost::serialization::access;

    template<class Archive>
    void serialize(Archive &ar, const unsigned int)
    {
        ar & buffer;
    }

};
  1. 最后我在读取之后和发送之前附加了序列化。
/**
 * receiver
 */
Package request;

{
    boost::archive::text_iarchive ia(*stream);
    ia >> request;
}

std::string data = my::decompress(request.buffer);

// do something with data

Package response(my::compress(data));

{
    boost::archive::text_oarchive oa(*stream);
    oa << response;
}

/**
 * sender
 */
std::string data = "hello world";
Package package(my::compress(data));

// send request
{
    boost::archive::text_oarchive oa(*m_stream);
    oa << package;
}

// waiting for a response
{
    boost::archive::text_iarchive ia(*m_stream);
    ia >> package;
}

// decompress response buffer
result = my::decompress(package.get_buffer());

【讨论】:

  • 您意识到这通过使用my:::compress/my::decompress 的双缓冲来回避问题,而Boost Serialization 与它无关,对吧?事实上,这与我建议的解决方案基本相同,尽管我的回答更笼统。
  • @sehe,Boost.Serialization 被用作仅序列化Package 对象的工具,并且该对象已经在发送方使用my::compress 打包,在接收方中使用反向操作
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-07-19
  • 2011-08-05
  • 2014-07-29
  • 2012-12-25
  • 1970-01-01
  • 2012-06-21
  • 2013-09-05
相关资源
最近更新 更多