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

如何在含同步/异步接口的库中正确使用asio::strand?

基于ASIO的UDP设备发现Acceptor问题解答

我正在通过编写一个基于UDP连接特定设备的库来学习ASIO。设备开机后会进入广播模式,向网络发送身份信息,因此主机需要监听UDP端口,匹配设备格式的数据包,并跟踪已检测到的设备(因为设备连接后仍会持续广播)。我曾用Python异步生成器完成原型开发,在ASIO中,类似的实现类比为asio::ip::tcp::acceptor。由于每个设备数据率较高,我希望自定义acceptor运行在多线程io_context中,因此需要为acceptor及每个设备控制器使用strand来序列化操作。

以下是我当前的示例代码:

template <typename Executor>
class Acceptor {
    std::unordered_set<asio::ip::udp::endpoint> accepted_connections;
    asio::basic_datagram_socket<asio::ip::udp, Executor> receive_sock;

    std::shared_ptr<Device<Executor>> accept(asio::error_code& ec) {
        std::array<std::byte, buffer_size> buffer;
        typename asio::basic_datagram_socket<asio::ip::udp, Executor>::endpoint_type remote_endpoint;
        receive_sock.receive_from(asio::buffer(buffer), remote_endpoint, {}, ec);
        if (!ec) {
            return {};
        }
        if (accepted_connections.contains(remote_endpoint)) {
            return {};
        }
        if (auto ret = std::start_lifetime_as<Header>(buffer.data())->validate(); !ret) {
            ec = ret.error();
            return {};
        }
        accepted_connections.emplace(remote_endpoint.remote_endpoint());
        return MakeDevice(remote_endpoint, buffer);
    }
};

问题1:如何修改阻塞式accept函数,使其在strand中运行,确保对accepted_connections集合的访问被正确序列化?

我发现strand的defer、dispatch、execute、post等操作都会将函数推入队列延后执行,这会导致无法即时检查error_code(因为函数延后执行时,accept返回时ec变量尚未赋值)。查看asio::ip::tcp::acceptor的accept实现源码也未得到有效信息。

解决方案

要让阻塞式accept的所有逻辑(socket接收、集合访问)都在strand上下文中执行,同时能同步获取结果,核心是用内部strand+同步Promise的方式实现:

  1. 在Acceptor类内部持有asio::strand<Executor>,将socket绑定到strand执行器,确保socket操作也受strand序列化。
  2. 用std::promise同步获取strand中执行的操作结果,解决即时获取error_code和返回值的问题。

修改后的代码:

template <typename Executor>
class Acceptor {
    asio::strand<Executor> strand_;
    std::unordered_set<asio::ip::udp::endpoint> accepted_connections;
    asio::basic_datagram_socket<asio::ip::udp, decltype(strand_)> receive_sock;

public:
    Acceptor(Executor ex)
        : strand_(std::move(ex)),
          receive_sock(strand_) {}

    std::shared_ptr<Device<decltype(strand_)>> accept(asio::error_code& ec) {
        std::promise<std::pair<asio::error_code, std::shared_ptr<Device<decltype(strand_)>>>> prom;
        auto fut = prom.get_future();

        // 将accept逻辑投递到strand执行
        asio::dispatch(strand_, [this, &prom]() {
            asio::error_code local_ec;
            std::array<std::byte, buffer_size> buffer;
            asio::ip::udp::endpoint remote_endpoint;
            receive_sock.receive_from(asio::buffer(buffer), remote_endpoint, {}, local_ec);

            std::shared_ptr<Device<decltype(strand_)>> dev;
            if (!local_ec) {
                if (!accepted_connections.contains(remote_endpoint)) {
                    if (auto ret = std::start_lifetime_as<Header>(buffer.data())->validate(); ret) {
                        accepted_connections.emplace(remote_endpoint);
                        dev = MakeDevice(remote_endpoint, buffer);
                    } else {
                        local_ec = ret.error();
                    }
                }
            }

            prom.set_value({local_ec, std::move(dev)});
        });

        // 同步等待结果返回
        auto result = fut.get();
        ec = result.first;
        return result.second;
    }
};

关键说明

  • 内部strand确保所有对accepted_connections的读写和socket操作都被序列化,多线程环境下不会出现并发冲突。
  • std::promise和std::future解决了strand队列操作无法即时返回结果的问题,调用线程会同步等待strand内的操作完成。
  • socket绑定到strand执行器,保证即使是阻塞的receive_from也在strand上下文中执行,避免跨线程操作socket的潜在问题。

问题2:实现async_accept时,如何将strand融入其中以保证accepted_connections集合的操作有序?

我需要使用async_initiate来兼容各种完成令牌,但不确定是在协程中await strand,还是需要修改async_initiate让协程在strand中运行?

解决方案

实现async_accept时,要让整个异步流程(接收、验证、集合操作、回调)都在strand上下文中执行,结合async_initiate和协程的实现方式如下:

  1. 编写异步协程async_accept_impl,利用绑定到strand的socket执行异步接收,后续逻辑自然在strand中执行。
  2. 在async_initiate内部,通过strand调度协程的启动和完成回调,确保所有操作都被序列化。

修改后的代码:

template <typename Executor>
class Acceptor {
    asio::strand<Executor> strand_;
    std::unordered_set<asio::ip::udp::endpoint> accepted_connections;
    asio::basic_datagram_socket<asio::ip::udp, decltype(strand_)> receive_sock;

    // 异步接收处理协程
    asio::awaitable<std::shared_ptr<Device<decltype(strand_)>>> async_accept_impl() {
        std::array<std::byte, buffer_size> buffer;
        asio::ip::udp::endpoint remote_endpoint;
        asio::error_code ec;

        // 异步接收,socket绑定到strand,操作自动在strand上下文执行
        co_await receive_sock.async_receive_from(
            asio::buffer(buffer), remote_endpoint,
            asio::redirect_error(asio::use_awaitable, ec)
        );

        if (ec) {
            co_return nullptr;
        }

        // 集合访问在strand上下文内,无需额外同步
        if (accepted_connections.contains(remote_endpoint)) {
            co_return nullptr;
        }

        if (auto ret = std::start_lifetime_as<Header>(buffer.data())->validate(); !ret) {
            co_return nullptr;
        }

        accepted_connections.emplace(remote_endpoint);
        co_return MakeDevice(remote_endpoint, buffer);
    }

public:
    Acceptor(Executor ex) : strand_(std::move(ex)), receive_sock(strand_) {}

    // 兼容任意完成令牌的async_accept
    template <typename CompletionToken>
    auto async_accept(CompletionToken&& token) {
        return asio::async_initiate<CompletionToken, void(asio::error_code, std::shared_ptr<Device<decltype(strand_)>>)>(
            [this](auto&& handler) {
                // 确保协程在strand上下文启动
                asio::dispatch(strand_, [this, handler = std::forward<decltype(handler)>(handler)]() mutable {
                    asio::co_spawn(strand_, async_accept_impl(),
                        [handler = std::move(handler)](std::exception_ptr e, std::shared_ptr<Device<decltype(strand_)>> dev) mutable {
                            asio::error_code ec;
                            if (e) {
                                try {
                                    std::rethrow_exception(e);
                                } catch (const asio::system_error& se) {
                                    ec = se.code();
                                } catch (...) {
                                    ec = asio::error::operation_aborted;
                                }
                            }
                            // 回调也通过strand调度,保证序列化
                            asio::dispatch(strand_, [handler = std::move(handler), ec, dev]() mutable {
                                handler(ec, dev);
                            });
                        }
                    );
                });
            },
            token
        );
    }
};

关键说明

  • 协程async_accept_impl直接使用绑定到strand的socket,异步接收操作自动在strand上下文中执行,后续的集合访问和验证逻辑无需额外同步。
  • 在async_initiate内部,通过asio::dispatch(strand_, ...)启动协程,确保整个异步流程都处于strand的序列化上下文。
  • 完成回调同样通过strand调度,避免回调在非strand线程执行时引发并发问题。
  • 如果不使用协程,也可以采用链式异步操作:先调用socket.async_receive_from,在其完成回调中通过strand调度集合检查、设备创建逻辑,最后调用完成令牌。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 17:10:54