【问题标题】:How to implement an IPC protocol using Boost ASIO?如何使用 Boost ASIO 实现 IPC 协议?
【发布时间】:2021-10-14 12:24:09
【问题描述】:

我正在尝试为将使用 Boost ASIO 构建的项目实现一个简单的 IPC 协议。这个想法是通过 IP/TCP 进行通信,服务器带有后端,客户端将使用从服务器接收到的数据来构建前端。整个会话会这样进行:

  1. 连接已建立
  2. 客户端发送一个 2 字节的数据包,其中包含一些信息,服务器将使用这些信息来构建其响应(存储为结构 propertiesPacket
  3. 服务器处理接收到的数据并将输出存储在一个名为processedData 的可变大小结构中
  4. 服务器发送一个 2 字节的无符号整数,指示客户端它将接收到的结构的大小(假设结构的大小为 n 字节)
  5. 服务器将结构数据作为n字节包发送
  6. 连接已结束

我尝试自己实现这一点,遵循 Boost ASIO 文档中的精彩教程,以及库中包含的示例和我在 Github 上找到的一些存储库,但由于这是我第一次使用网络和 IPC,我无法让它工作,我的客户端返回一个异常,说连接已被对等方重置。

我现在拥有的是这样的:

// File client.cpp
int main(int argc, char *argv[])
{
    try {
        propertiesPacket properties;
        // ...
        // We set the data inside the properties struct
        // ...

        boost::asio::io_context io;
        boost::asio::ip::tcp::socket socket(io);
        boost::asio::ip::tcp::resolver resolver(io);

        boost::asio::connect(socket, resolver.resolve(argv[1], argv[2]));
        boost::asio::write(socket, boost::asio::buffer(&properties, sizeof(propertiesPacket)));

        unsigned short responseSize {};
        boost::asio::read(socket, boost::asio::buffer(&responseSize, sizeof(short)));

        processedData* response = reinterpret_cast<processedData*>(malloc(responseSize));
        boost::asio::read(socket, boost::asio::buffer(response, responseSize));

        // ...
        // The client handles the data
        // ...

        return 0;
    } catch (std::exception &e) {
        std::cerr << e.what() << std::endl;
    }
}
// File server.cpp
class ServerConnection
    : public std::enable_shared_from_this<ServerConnection>
{
    public:
        using TCPSocket = boost::asio::ip::tcp::socket;

        ServerConnection::ServerConnection(TCPSocket socket)
          : socket_(std::move(socket)),
            properties_(nullptr),
            filePacket_(nullptr),
            filePacketSize_(0)
        {
        }


        void start() { doRead(); }

    private:
        void doRead()
        {
            auto self(shared_from_this());
            socket_.async_read_some(boost::asio::buffer(properties_, sizeof(propertiesPacket)),
                                    [this, self](boost::system::error_code ec, std::size_t /*length*/)
                                    {
                                        if (!ec) {
                                            processData();
                                            doWrite(&filePacketSize_, sizeof(short));

                                            doWrite(filePacket_, sizeof(*filePacket_));
                                        }
                                    });

        }

        void doWrite(void* data, size_t length)
        {
            auto self(shared_from_this());
            boost::asio::async_write(socket_, boost::asio::buffer(data, length),
                                     [this, self](boost::system::error_code ec, std::size_t /*length*/)
                                     {
                                         if (!ec) { doRead(); }
                                     });
        }

        void processData()
        { /* Data is processed */ }

        TCPSocket socket_;

        propertiesPacket* properties_;
        processedData* filePacket_;
        short filePacketSize_;
};

class Server
{
    public:
        using IOContext = boost::asio::io_context;
        using TCPSocket = boost::asio::ip::tcp::socket;
        using TCPAcceptor = boost::asio::ip::tcp::acceptor;

        Server::Server(IOContext& io, short port)
            : socket_(io),
              acceptor_(io, boost::asio::ip::tcp::endpoint(boost::asio::ip::tcp::v4(), port))
        {
            doAccept();
        }


    private:
        void doAccept()
        {
            acceptor_.async_accept(socket_,
                [this](boost::system::error_code ec)
                {
                    if (!ec) {
                        std::make_shared<ServerConnection>(std::move(socket_))->start();
                    }

                    doAccept();
                });
        }

        TCPSocket socket_;
        TCPAcceptor acceptor_;
};

我做错了什么?我的猜测是,在doRead 函数内部,多次调用doWrite 函数,然后该函数还调用doRead 部分是导致问题的原因,但我不知道异步写入数据的正确方法是什么多次是。但我也确信这不是我的代码中唯一表现不佳的部分。

【问题讨论】:

  • ServerConnection 中,成员properties_filePacketSize_ 永远不会被赋予除nullptr 之外的其他值,因此您的程序具有UB。要么包含完整的独立代码,要么在您的代码中修复它:)。
  • 如果您尝试上传(多堆)文件,请参阅例如这里:stackoverflow.com/a/66607723/85371(向下滚动到"TCP Socket Version"
  • 感谢 cmets!它们的值在我遗漏的processData 函数中分配。另外,当我读取客户端发送的数据时,properties_ 不应该被赋值吗?
  • 无论如何,我会检查您评论的链接,谢谢! :)

标签: c++ networking boost ipc boost-asio


【解决方案1】:

除了我mentioned in the comments显示的代码有问题外,确实有你怀疑的问题:

我的猜测是,在 doRead 函数内部,多次调用 doWrite 函数,然后该函数也调用 doRead 部分是导致问题的原因

“doRead”在同一个函数中的事实不一定是问题(这只是全双工套接字 IO)。但是“多次调用”是。见docs

此操作是通过对流的async_write_some 函数的零次或多次调用来实现的,称为组合操作。程序必须确保流不执行其他写入操作(例如async_write、流的async_write_some 函数,或执行写入的任何其他组合操作),直到此操作完成。

通常的方法是将整个消息放在一个缓冲区中,但如果复制“昂贵”,您可以使用称为scatter/gather buffers 的 BufferSequence。

具体来说,你会替换

doWrite(&filePacketSize_, sizeof(short));
doWrite(filePacket_, sizeof(*filePacket_));

类似的东西

std::vector<boost::asio::const_buffer> msg{
    boost::asio::buffer(&filePacketSize_, sizeof(short)),
    boost::asio::buffer(filePacket_, sizeof(*filePacket_)),
};

doWrite(msg);

请注意,这假定 filePacketSizefilePacket 已分配正确的值!

您当然可以修改do_write 以接受缓冲序列:

template <typename Buffers> void doWrite(Buffers msg)
{
    auto self(shared_from_this());
    boost::asio::async_write(
        socket_, msg,
        [this, self](boost::system::error_code ec, std::size_t /*length*/) {
            if (!ec) {
                doRead();
            }
        });
}

但在你的情况下,我会通过内联正文来简化(现在你不会多次调用它)。

旁注

不要使用newdelete。切勿在 C++ 中使用 malloc。永远不要使用reinterpret_cast&lt;&gt;(除了标准允许的极少数例外情况!)。而不是

    processedData* response = reinterpret_cast<processedData*>(malloc(responseSize));

随便用

    processedData response;

(可选添加{} 用于聚合的值初始化)。如果您需要可变长度消息,请考虑在消息中放置一个向量或数组。当然,数组是固定长度的,但它保留了 POD 特性,因此使用起来可能更容易。如果你使用向量,你会想要一个分散/聚集读入一个缓冲区序列,就像我在上面展示的写端一样。

与其在不一致的shortunsigned short 类型之间重新解释,不如只用标准大小拼写类型:std::uint16_t。 请记住,您没有考虑字节顺序,因此您的协议将无法跨编译器/架构移植。

临时修复

这是我在查看您共享的代码后最终得到的清单。

Live On Coliru

#include <boost/asio.hpp>
#include <iostream>

namespace ba = boost::asio;
using boost::asio::ip::tcp;
using boost::system::error_code;
using TCPSocket = tcp::socket;

struct processedData { };
struct propertiesPacket { };

// File server.cpp
class ServerConnection : public std::enable_shared_from_this<ServerConnection> {
  public:
    ServerConnection(TCPSocket socket) : socket_(std::move(socket))
    { }

    void start() {
        std::clog << __PRETTY_FUNCTION__ << std::endl;
        doRead();
    }

  private:
    void doRead()
    {
        std::clog << __PRETTY_FUNCTION__ << std::endl;
        auto self(shared_from_this());
        socket_.async_read_some(
            ba::buffer(&properties_, sizeof(properties_)),
            [this, self](error_code ec, std::size_t length) {
                std::clog << "received: " << length << std::endl;
                if (!ec) {
                    processData();

                    std::vector<ba::const_buffer> msg{
                        ba::buffer(&filePacketSize_, sizeof(uint16_t)),
                        ba::buffer(&filePacket_, filePacketSize_),
                    };

                    ba::async_write(socket_, msg,
                        [this, self = shared_from_this()](
                            error_code ec, std::size_t length) {
                            std::clog << " written: " << length
                                      << std::endl;
                            if (!ec) {
                                doRead();
                            }
                        });
                }
            });
    }

    void processData() {
        std::clog << __PRETTY_FUNCTION__ << std::endl;
        /* Data is processed */
    }
    TCPSocket socket_;

    propertiesPacket properties_{};
    processedData    filePacket_{};
    uint16_t         filePacketSize_ = sizeof(filePacket_);
};

class Server
{
  public:
    using IOContext   = ba::io_context;
    using TCPAcceptor = tcp::acceptor;

    Server(IOContext& io, uint16_t port)
        : socket_(io)
        , acceptor_(io, {tcp::v4(), port})
    {
        doAccept();
    }

  private:
    void doAccept()
    {
        std::clog << __PRETTY_FUNCTION__ << std::endl;
        acceptor_.async_accept(socket_, [this](error_code ec) {
            if (!ec) {
                std::clog << "Accepted " << socket_.remote_endpoint()
                          << std::endl;
                std::make_shared<ServerConnection>(std::move(socket_))->start();
                doAccept();
            } else {
                std::clog << "Accept " << ec.message() << std::endl;
            }
        });
    }

    TCPSocket   socket_;
    TCPAcceptor acceptor_;
};

// File client.cpp
int main(int argc, char *argv[])
{
    ba::io_context io;
    Server         s{io, 6869};

    std::thread server_thread{[&io] {
        io.run();
    }};

    // always check argc!
    std::vector<std::string> args(argv, argv + argc);

    if (args.size() == 1)
        args = {"demo", "127.0.0.1", "6869"};

    // avoid race with server accept thread
    post(io, [&io, args] {
        try {
            propertiesPacket properties;
            // ...
            // We set the data inside the properties struct
            // ...

            tcp::socket   socket(io);
            tcp::resolver resolver(io);

            connect(socket, resolver.resolve(args.at(1), args.at(2)));
            write(socket, ba::buffer(&properties, sizeof(properties)));

            uint16_t responseSize{};
            ba::read(socket, ba::buffer(&responseSize, sizeof(uint16_t)));

            std::clog << "Client responseSize: " << responseSize << std::endl;
            processedData response{};
            assert(responseSize <= sizeof(response));
            ba::read(socket, ba::buffer(&response, responseSize));

            // ...
            // The client handles the data
            // ...
            
            // for online demo:
            io.stop();
        } catch (std::exception const& e) {
            std::clog << e.what() << std::endl;
        }
    });

    io.run_one();
    server_thread.join();
}

打印类似的东西

void Server::doAccept()
Server::doAccept()::<lambda(boost::system::error_code)> Success
void ServerConnection::start()
void ServerConnection::doRead()
void Server::doAccept()
received: 1
void ServerConnection::processData()
 written: 3
void ServerConnection::doRead()
Client responseSize: 1

【讨论】:

  • 我一直在尝试实现它,并按照您的建议,我删除了我正在使用的所有指针,现在它们只是普通变量,我还将数据存储在我的结构中现在是一个向量。似乎发送消息的部分工作得很好,但是当我尝试从客户端接收我的结构时,我得到一个Segment violation ('core' generated)。我想这是因为我没有初始化我的向量,但我仍然得到违规:(。有什么想法吗?
  • std::vector 不是一个平凡的类型,所以你不能根据语言规则“不初始化”它。你读过这篇笔记"If you use vector, you'd want a scatter/gather read into a buffer sequence like I showed above for the write side."吗?
  • 这听起来有点像您忘记了 C++ 不是 C,并且假设每种类型都是普通/标准布局(又名“POD”类型,或按位可复制数据)是不安全的在 C++ 中。如果您展示一个独立的代码示例(可能作为一个单独的问题),我们可以帮助您解决问题。
  • 是的,这就是我所做的,虽然我只是将整个结构作为缓冲区发送,而不是每个单独的成员。我应该这样做吗?
  • "是的,我做了 ABC,虽然我只是做了 XYZ" :) 不平凡的类型不能按位复制,更不用说从跨进程空间甚至架构接收消息的问题(word大小,字节顺序)。是的,您应该以安全的方式序列化数据。给这只猫剥皮的方法有很多,所以也许你可以发布你拥有的代码,以便我们提供帮助。
猜你喜欢
  • 2016-08-03
  • 1970-01-01
  • 1970-01-01
  • 2016-01-07
  • 2016-08-09
  • 1970-01-01
  • 2011-06-19
  • 1970-01-01
相关资源
最近更新 更多