在Asio中,取消后如何重新开始异步读取与异步写入?

编程语言 2026-07-08

Asio在 awaitable<>::operator|| 的一个分支结束时会取消其他awaitable对象。取消 async_read / async_write 可能会丢失数据。如何在取消时保存读写进度并在之后恢复? 例如,如何实现一个自定义cancellation_type::partial async_read,使用async_read_some,以及使用async_write_some的 async_write? Asio版本为1.38.0

#include <asio.hpp>
#include <asio/experimental/awaitable_operators.hpp>
#include <assert.h>
#include <iostream>
#include <stdio.h>
using asio::as_tuple_t;
using asio::awaitable;
using asio::buffer;
using asio::co_spawn;
using asio::detached;
using asio::io_context;
using asio::steady_timer;
using asio::ip::tcp;
using namespace asio::experimental::awaitable_operators;
using std::chrono::steady_clock;
using namespace std::literals::chrono_literals;
using asio::use_awaitable_t;
using default_token = as_tuple_t<use_awaitable_t<>>;
using tcp_acceptor = default_token::as_default_on_t<tcp::acceptor>;
using tcp_socket = default_token::as_default_on_t<tcp::socket>;
using asio::ip::tcp;
void Reverse(char s[], size_t len) {
  for (size_t i = 0; i < len / 2; ++i) {
    char tmp = s[i];
    s[i] = s[len - 1 - i];
    s[len - 1 - i] = tmp;
  }
}
void Print(const void *p, size_t len) {
  auto q = (const unsigned char *)p;
  for (size_t i = 0; i < len; ++i) {
    printf("%02x ", q[i]);
  }
  putchar('\n');
}
awaitable<void> Timeout(steady_timer &timer, steady_clock::duration duration) {
  timer.expires_after(duration);
  co_await timer.async_wait();
}
#define SomeAwaitable(timer) Timeout(timer, 5s)
awaitable<void> Task(io_context &ctx) {
  try {
    tcp_acceptor acceptor(
        ctx, {tcp::endpoint(asio::ip::address_v4::loopback(), 12345)});
    char data[5];
    steady_timer timer{ctx};
    while (1) {
      auto [e, fd] = co_await acceptor.async_accept();
      if (e) {
        std::cout << "accept err:" << e << '\n';
        continue;
      }
      while (1) {
        auto rr =
            co_await (async_read(fd, buffer(data)) || SomeAwaitable(timer));
        switch (rr.index()) {
        default:
          break;
        case 1: // partial read data is discarded
          std::cout << "timeout\n";
          continue;
        }
        auto [e1, nread] = std::get<0>(rr);
        if (e1) {
          std::cout << "read err:" << e1 << '\n';
          fd.close();
          break;
        }
        assert(nread == 5);
        Print(data, nread);
        Reverse(data, nread);
      BeforeWrite:
        auto wr = co_await (async_write(fd, buffer(data, nread)) ||
                            SomeAwaitable(timer));
        switch (wr.index()) {
        default:
          break;
        case 1: // I guess partial written data is discarded too
          std::cout << "timeout\n";
          goto BeforeWrite;
        }
        auto [e2, nwrite] = std::get<0>(wr);
        if (e2) {
          std::cout << "write err:" << e2 << '\n';
          fd.close();
          break;
        }
      }
    }
  } catch (const std::exception &e) {
    std::cerr << "error: " << e.what() << '\n';
    co_return;
  }
}
int main() {
  io_context ctx;
  co_spawn(ctx, Task(ctx), detached);
  ctx.run();
  std::cout << "Done\n";
}

解决方案

你一开始就提出了一个大胆的说法(取消可能会丢失数据)。我并不认为是这样。

此外,你可以显著简化代码,并同时实现并发连接。就这一点而言,我感觉这是C 风格编码,因此让我展示C++20风格:

Live On Compiler Explorer

#include <boost/asio.hpp>
#include <iostream>
#include <print>
namespace asio = boost::asio;
using namespace std::chrono_literals;
using asio::ip::tcp;
using boost::system::error_code;

asio::awaitable<void> Session(tcp::socket s) try {
    error_code ec;
    auto       timed = cancel_after(5s, asio::redirect_error(ec));

    for (std::array<char, 5> data; !ec;) {
        auto nread = co_await async_read(s, asio::buffer(data), timed);

        if (ec == asio::error::operation_aborted) {
            std::cout << "timeout" << std::endl;
            ec.clear(); // allow loop to resume only when the error was timeout
        }

        // still handle any partial read

        if (auto payload = std::span(data).subspan(0, nread); !payload.empty()) {
            std::print("{::#02x}\n", payload);
            std::ranges::reverse(payload);

            auto nwrite = co_await async_write(s, payload /*, timed*/);
            assert(nwrite == payload.size());
        }
    }
} catch (std::exception const& e) {
    std::cerr << "Session: " << e.what() << std::endl;
}

asio::awaitable<void> Listener(uint16_t port) try {
    tcp::acceptor a(co_await asio::this_coro::executor, {{}, port});

    for (;;) {
        auto [ec, s] = co_await a.async_accept(asio::as_tuple);

        if (ec) 
            std::cout << "accept err:" << ec.message() << std::endl;
        else
            co_spawn(a.get_executor(), Session(std::move(s)), asio::detached);
    }
} catch (std::exception const& e) {
    std::cerr << "Listener: " << e.what() << std::endl;
}

int main() {
    asio::io_context ctx;
    co_spawn(ctx, Listener(12345), asio::detached);
    ctx.run();
    std::cout << "Done" << std::endl;
}

如果你确实想保留缓冲区,在写入任何数据之前就继续读取,请使用动态缓冲区:

asio::awaitable<void> Session(tcp::socket s) try {
    for (std::string data;;) {
        auto [ec, nread] = co_await async_read( //
            s, asio::dynamic_buffer(data, 5), cancel_after(5s, asio::as_tuple));

        if (ec.failed()) {
            if (ec == asio::error::operation_aborted) {
                std::cout << "timeout" << std::endl;
                continue;
            } else {
                std::cout << "read: " << ec.message() << std::endl;
                break;
            }
        }

        if (!data.empty()) {
            std::print("{::#02x}\n", std::span(data));
            std::ranges::reverse(data);

            auto nwrite = co_await async_write(s, asio::buffer(data));
            data.erase(0, nwrite); // effectively clear
        }
    }
} catch (std::exception const& e) {
    std::cerr << "Session: " << e.what() << std::endl;
}

效果类似:

站内所有文章版权归属LeftHeroAI导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。

相关文章