在Asio中,取消后如何重新开始异步读取与异步写入?
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风格:
#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导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。