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

Boost.Cobalt:从C API回调使用通道时出现空指针恢复问题

协程管道中调用C API导致Asio内部空指针恢复问题

我首次尝试使用协程,经验有限,想在通过channel通信的协程管道里调用C API。管道结构如下:

--------    1    ---------    2    ------
| source | ----> | process | ----> | sink |
 --------         ---------         ------

每个方块是协程,箭头是通道。C API用于process协程,签名大致为bool start_work(consumer_callback),属于同步API,会为每个生成的数据调用一次consumer_callback。

最初想在回调里直接写入通道2,但会改变回调签名,不可行。于是改为向回调传递协程句柄,由回调恢复协程,再由恢复后的协程写入通道2。简化代码如下:

#include <coroutine>
#include <optional>
#include <string>
#include <boost/cobalt/channel.hpp>
#include <boost/cobalt/main.hpp>
#include <boost/cobalt/promise.hpp>
#include <boost/cobalt/join.hpp>

namespace cobalt = boost::cobalt;

// Data to communicate between the callback and the channel writer.
struct Data {
   std::optional<int> result;
   bool done = false;
   std::coroutine_handle<> coro_handle;
};

using Callback = void (*)(int, void*, bool);

void consumer_callback(int i, void* data, bool done) {
   Data& data_ = *reinterpret_cast<Data*>(data);
   data_.done = done;
   if (!done) {
      data_.result = i;
   }
   data_.coro_handle.resume();
}

// C API that produces results and calls the callback to consume each result.
// Results are integers.
void start_work(void* data, Callback cb) {
    bool done = false;
    for (int i = 0; i < 10; ++i) {
       cb(i, data, done); // !done case
    }
    done = true;
    cb(0, data, done); // done case
}

struct Awaiter : std::suspend_always {
    Data& data;
    bool first;

    bool await_ready() {
        return data.result.has_value();
    }

    void await_suspend(std::coroutine_handle<> h) {
        data.coro_handle = h;
        if (first) start_work(&data, consumer_callback);
    }

    int await_resume() {
        assert(data.result.has_value());
        auto opt = std::exchange(data.result, std::nullopt);
        return opt.value();
    }
};

Awaiter nextResult(Data& data, bool first) {
    return {{}, data, first};
}

cobalt::promise<void> source(cobalt::channel<std::string>& out) {
    co_await out.write("Hello world!");
    out.close();
}

cobalt::promise<void> process(cobalt::channel<std::string>& in, cobalt::channel<int>& out) {
    Data data;
    while (in.is_open() && out.is_open()) {
        auto _ = co_await in.read(); // ignore result for now
        auto first = true;
        while (!data.done || data.result.has_value()) {
            auto i = co_await nextResult(data, first);
            co_await out.write(i);
            first = false;
        }
    }
    in.close();
    out.close();
}

cobalt::promise<void> sink(cobalt::channel<int>& in) {
    while (in.is_open()) {
        auto i = co_await in.read(); // ignore result for now
    }
    in.close();
}

cobalt::main co_main(int argc, char* argv[]) {
    cobalt::channel<std::string> a;
    cobalt::channel<int> b;
    co_await cobalt::join(
        source(a),
        process(a, b),
        sink(b)
    );
    co_return 0;
}

sink能正确接收所有数据,但process协程结束时,Asio内部出现恢复空指针的情况,请问哪里出错了?

环境:Ubuntu 20.04、Boost 1.85、g13 -std=gnu2a


问题根源

核心问题出在回调对协程句柄的重复调用以及Data结构体状态管理不当:

  • 同步C APIstart_work会在await_suspend中阻塞并一次性触发所有回调,每次回调都调用resume(),但协程仅需被恢复一次,后续重复调用会访问已失效的协程句柄。
  • Data结构体的coro_handle未在协程执行完当前批次处理后重置,若外层循环再次执行,会复用无效的句柄。
  • 直接在协程调度线程中执行阻塞的start_work,会阻塞Boost.Cobalt调度器,引发线程安全问题。

修复方案

1. 调整数据传递逻辑,用队列存储批量结果

修改Data结构体,用队列缓存所有生成的结果,避免单次结果覆盖:

struct Data {
   std::queue<int> results;
   bool done = false;
   std::coroutine_handle<> coro_handle;
   std::mutex mutex; // 线程安全保护
};

2. 修改回调函数,仅在所有结果生成完毕后恢复协程

void consumer_callback(int i, void* data, bool done) {
   Data& data_ = *reinterpret_cast<Data*>(data);
   std::lock_guard<std::mutex> lock(data_.mutex);
   data_.done = done;
   if (!done) {
      data_.results.push(i);
   }
   // 仅在全部数据生成完成后恢复协程
   if (done) {
      data_.coro_handle.resume();
   }
}

3. 重构Awaiter,避免阻塞调度器线程

将同步C API调用放到单独线程执行,防止阻塞协程调度器:

struct Awaiter : std::suspend_always {
    Data& data;
    bool first;

    bool await_ready() {
        std::lock_guard<std::mutex> lock(data.mutex);
        return !data.results.empty() || data.done;
    }

    void await_suspend(std::coroutine_handle<> h) {
        data.coro_handle = h;
        if (first) {
            // 启动独立线程执行同步C API
            std::jthread([&data = data]() {
                start_work(&data, consumer_callback);
            }).detach();
        }
    }

    std::optional<int> await_resume() {
        std::lock_guard<std::mutex> lock(data.mutex);
        if (!data.results.empty()) {
            auto val = data.results.front();
            data.results.pop();
            return val;
        }
        return std::nullopt;
    }
};

4. 适配process协程逻辑

cobalt::promise<void> process(cobalt::channel<std::string>& in, cobalt::channel<int>& out) {
    while (in.is_open() && out.is_open()) {
        auto _ = co_await in.read();
        Data data;
        auto first = true;
        while (!data.done || !data.results.empty()) {
            auto opt_i = co_await nextResult(data, first);
            if (opt_i) {
                co_await out.write(*opt_i);
            }
            first = false;
        }
    }
    in.close();
    out.close();
}

关键注意事项

  • 永远不要在协程调度线程中执行阻塞操作,必须将同步C API放到独立线程运行。
  • 用互斥锁保护跨线程访问的Data结构体,避免数据竞争。
  • 确保协程句柄仅在协程挂起状态时被调用,防止重复恢复或访问已销毁的协程帧。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 00:50:57