beast support per-op cancellation

- websocket supports cancellation.
 - Iterating test for ws cancellation.
 - Only terminal cancellation is forwarded by default.
 - basic_stream supports cancellation.
 - supported cancellation is documented.
 - http cancellation additions.
 - Added cancellation_slot tests to http, utils and saved_handler.
 - Added post to write.cpp, to avoid SIGSEV in test.
 - Refresher describes cancellation in more detail.
This commit is contained in:
Klemens
2022-10-27 22:56:19 +08:00
committed by Klemens Morgenstern
parent 0bf3d971a0
commit 3ebff60b1a
35 changed files with 1216 additions and 29 deletions
+1
View File
@@ -48,6 +48,7 @@ add_executable (tests-beast-core
file_posix.cpp
file_stdio.cpp
file_win32.cpp
filtering_cancellation_slot.cpp
flat_buffer.cpp
flat_static_buffer.cpp
flat_stream.cpp
+1
View File
@@ -39,6 +39,7 @@ local SOURCES =
file_posix.cpp
file_stdio.cpp
file_win32.cpp
filtering_cancellation_slot.cpp
flat_buffer.cpp
flat_static_buffer.cpp
flat_stream.cpp
@@ -0,0 +1,48 @@
//
// Copyright (c) 2022 Klemens Morgenstern (klemens.morgenstern@gmx.net)
//
// Distributed under the Boost Software License, Version 1.0. (See accompanying
// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
//
// Test that header file is self-contained.
#include <boost/beast/core/detail/filtering_cancellation_slot.hpp>
#include <boost/beast/_experimental/unit_test/suite.hpp>
namespace boost {
namespace beast {
struct filtering_cancellation_slot_test : beast::unit_test::suite
{
void
run()
{
using ct = net::cancellation_type;
ct fired = ct::none;
auto l = [&fired](ct tp){fired = tp;};
net::cancellation_signal sl;
detail::filtering_cancellation_slot<> slot{ct::terminal, sl.slot()};
slot.type |= ct::total;
slot = sl.slot();
slot.assign(l);
BEAST_EXPECT(fired == ct::none);
sl.emit(ct::total);
BEAST_EXPECT(fired == ct::total);
sl.emit(ct::partial);
BEAST_EXPECT(fired == ct::total);
sl.emit(ct::terminal);
BEAST_EXPECT(fired == ct::terminal);
}
};
BEAST_DEFINE_TESTSUITE(beast,core,filtering_cancellation_slot);
} // beast
} // boost
+83 -5
View File
@@ -9,7 +9,7 @@
// Test that header file is self-contained.
#include <boost/beast/core/saved_handler.hpp>
#include <boost/asio/bind_cancellation_slot.hpp>
#include <boost/beast/_experimental/unit_test/suite.hpp>
#include <stdexcept>
@@ -46,7 +46,7 @@ public:
}
void
operator()()
operator()(system::error_code ec_ = {})
{
failed_ = false;
}
@@ -74,7 +74,7 @@ public:
}
void
operator()()
operator()(system::error_code = {})
{
invoked_ = true;
}
@@ -90,7 +90,7 @@ public:
}
void
operator()()
operator()(system::error_code = {})
{
}
};
@@ -119,7 +119,7 @@ public:
{
saved_handler sh;
try
{
{
sh.emplace(throwing_handler{});
fail();
}
@@ -131,10 +131,88 @@ public:
}
}
void
testSavedHandlerCancellation()
{
{
net::cancellation_signal sig;
saved_handler sh;
BEAST_EXPECT(! sh.has_value());
sh.emplace(
net::bind_cancellation_slot(
sig.slot(), handler{}));
BEAST_EXPECT(sh.has_value());
BEAST_EXPECT(sig.slot().has_handler());
sig.emit(net::cancellation_type::all);
BEAST_EXPECT(! sh.has_value());
BEAST_EXPECT(!sig.slot().has_handler());
sh.emplace(
net::bind_cancellation_slot(
sig.slot(), handler{}));
BEAST_EXPECT(sh.has_value());
BEAST_EXPECT(sig.slot().has_handler());
sig.emit(net::cancellation_type::total);
BEAST_EXPECT(sh.has_value());
BEAST_EXPECT(sig.slot().has_handler());
sig.emit(net::cancellation_type::terminal);
BEAST_EXPECT(! sh.has_value());
BEAST_EXPECT(!sig.slot().has_handler());
sh.emplace(
net::bind_cancellation_slot(
sig.slot(), handler{}),
net::cancellation_type::total);
BEAST_EXPECT(sh.has_value());
BEAST_EXPECT(sig.slot().has_handler());
sig.emit(net::cancellation_type::total);
BEAST_EXPECT(! sh.has_value());
BEAST_EXPECT(!sig.slot().has_handler());
{
saved_handler sh_inner;
sh_inner.emplace(
net::bind_cancellation_slot(
sig.slot(), handler{}));
sh = std::move(sh_inner);
}
BEAST_EXPECT(sh.has_value());
BEAST_EXPECT(sig.slot().has_handler());
sig.emit(net::cancellation_type::all);
BEAST_EXPECT(! sh.has_value());
BEAST_EXPECT(!sig.slot().has_handler());
}
{
saved_handler sh;
net::cancellation_signal sig;
try
{
sh.emplace(
net::bind_cancellation_slot(
sig.slot(),
throwing_handler{}));
fail();
}
catch(std::exception const&)
{
pass();
}
BEAST_EXPECT(!sig.slot().has_handler());
BEAST_EXPECT(! sh.has_value());
}
}
void
run() override
{
testSavedHandler();
testSavedHandlerCancellation();
}
};
+62 -1
View File
@@ -25,6 +25,12 @@
#include <boost/asio/ip/tcp.hpp>
#include <boost/asio/strand.hpp>
#include <boost/asio/write.hpp>
#include <boost/asio/bind_cancellation_slot.hpp>
#include <boost/asio/error.hpp>
#include <boost/asio/io_context.hpp>
#include <boost/asio/connect_pipe.hpp>
#include <boost/asio/readable_pipe.hpp>
#include <boost/asio/writable_pipe.hpp>
#include <atomic>
#if BOOST_ASIO_HAS_CO_AWAIT
@@ -729,7 +735,57 @@ public:
}
}
void
testCancellation(yield_context do_yield)
{
// this is tested on a pipe
// because the test::stream doesn't implement cancellation
{
response<string_body> m;
error_code ec;
net::writable_pipe ts{ioc_};
net::readable_pipe tr{ioc_};
net::connect_pipe(tr, ts);
net::cancellation_signal cl;
net::post(ioc_, [&]{cl.emit(net::cancellation_type::all);});
net::steady_timer timeout(ioc_, std::chrono::seconds(5));
timeout.async_wait(
[&](error_code ec)
{
BEAST_EXPECT(ec == net::error::operation_aborted);
if (!ec) // this means the cancel failed!
ts.close();
});
multi_buffer b;
async_read(tr, b, m, net::bind_cancellation_slot(cl.slot(), do_yield[ec]));
timeout.cancel();
BEAST_EXPECT(ec == net::error::operation_aborted);
}
{
response<string_body> m;
error_code ec;
net::writable_pipe ts{ioc_};
net::readable_pipe tr{ioc_};
net::connect_pipe(tr, ts);
net::cancellation_signal cl;
net::post(ioc_, [&]{cl.emit(net::cancellation_type::all);});
net::steady_timer timeout(ioc_, std::chrono::seconds(5));
timeout.async_wait(
[&](error_code ec)
{
// using BEAST_EXPECT HERE is a race condition, since the test suite might
BEAST_EXPECT(ec == net::error::operation_aborted);
if (!ec) // this means the cancel failed!
ts.close();
});
multi_buffer b;
async_read(tr, b, m, net::bind_cancellation_slot(cl.slot(), do_yield[ec]));
timeout.cancel();
BEAST_EXPECT(ec == net::error::operation_aborted);
}
// the timer handler may be invoked after the test suite is complete if we don't post.
asio::post(ioc_, do_yield);
}
void
run() override
{
@@ -761,6 +817,11 @@ public:
testReadSomeHeader(yield);
});
testReadSomeHeader();
yield_to(
[&](yield_context yield)
{
testCancellation(yield);
});
}
+73
View File
@@ -22,8 +22,13 @@
#include <boost/beast/_experimental/test/stream.hpp>
#include <boost/beast/test/yield_to.hpp>
#include <boost/beast/_experimental/unit_test/suite.hpp>
#include <boost/asio/bind_cancellation_slot.hpp>
#include <boost/asio/error.hpp>
#include <boost/asio/io_context.hpp>
#include <boost/asio/connect_pipe.hpp>
#include <boost/asio/readable_pipe.hpp>
#include <boost/asio/writable_pipe.hpp>
#include <boost/asio/steady_timer.hpp>
#include <boost/asio/strand.hpp>
#include <sstream>
#include <string>
@@ -1050,6 +1055,69 @@ public:
}
#endif
void
testCancellation(yield_context do_yield)
{
// this is tested on a pipe
// because the test::stream doesn't implement cancellation
{
response<string_body> m;
m.version(10);
m.result(status::ok);
m.set(field::server, "test");
m.set(field::content_length, "5");
// make the content big enough so it overflows the buffer
// that'll make the op never complete if we don't cancel
m.body().assign(10000000, '*');
error_code ec;
net::writable_pipe ts{ioc_};
net::readable_pipe tr{ioc_};
net::connect_pipe(tr, ts);
net::cancellation_signal cl;
net::post(ioc_, [&]{cl.emit(net::cancellation_type::all);});
net::steady_timer timeout(ioc_, std::chrono::seconds(5));
timeout.async_wait(
[&](error_code ec)
{
BEAST_EXPECT(ec == net::error::operation_aborted);
if (!ec) // this means the cancel failed!
ts.close();
});
async_write(ts, m, net::bind_cancellation_slot(cl.slot(), do_yield[ec]));
timeout.cancel();
net::post(ioc_, do_yield); // wait for the timeout to finish
BEAST_EXPECT(ec == net::error::operation_aborted);
}
{
response<string_body> m;
m.version(11);
m.result(status::ok);
m.set(field::server, "test");
m.set(field::transfer_encoding, "chunked");
m.body().assign(10000000, '*');
error_code ec;
net::writable_pipe ts{ioc_};
net::readable_pipe tr{ioc_};
net::connect_pipe(tr, ts);
net::cancellation_signal cl;
net::post(ioc_, [&]{cl.emit(net::cancellation_type::all);});
net::steady_timer timeout(ioc_, std::chrono::seconds(5));
timeout.async_wait(
[&](error_code ec)
{
BEAST_EXPECT(ec == net::error::operation_aborted);
if (!ec) // this means the cancel failed!
ts.close();
});
async_write(ts, m, net::bind_cancellation_slot(cl.slot(), do_yield[ec]));
timeout.cancel();
net::post(ioc_, do_yield); // wait for the timeout to finish
BEAST_EXPECT(ec == net::error::operation_aborted);
}
// the timer handler may be invoked after the test suite is complete if we don't post.
asio::post(ioc_, do_yield);
}
void
run() override
@@ -1077,6 +1145,11 @@ public:
#if BOOST_ASIO_HAS_CO_AWAIT
boost::ignore_unused(&write_test::testAwaitableCompiles);
#endif
yield_to(
[&](yield_context yield)
{
testCancellation(yield);
});
}
};
+1
View File
@@ -21,6 +21,7 @@ add_executable (tests-beast-websocket
test.hpp
_detail_prng.cpp
accept.cpp
cancel.cpp
close.cpp
error.cpp
frame.cpp
+1
View File
@@ -12,6 +12,7 @@ local SOURCES =
_detail_impl_base.cpp
_detail_prng.cpp
accept.cpp
cancel.cpp
close.cpp
error.cpp
frame.cpp
+230
View File
@@ -0,0 +1,230 @@
//
// Copyright (c) 2016-2019 Vinnie Falco (vinnie dot falco at gmail dot com)
//
// Distributed under the Boost Software License, Version 1.0. (See accompanying
// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
//
// Official repository: https://github.com/boostorg/beast
//
// Test that header file is self-contained.
#include <boost/beast/websocket/stream.hpp>
#include <boost/beast/core/tcp_stream.hpp>
#include <boost/beast/_experimental/test/stream.hpp>
#include <boost/beast/_experimental/test/tcp.hpp>
#include <boost/beast/_experimental/unit_test/suite.hpp>
#include <boost/asio/strand.hpp>
#include <boost/asio/bind_cancellation_slot.hpp>
#include <boost/asio/deferred.hpp>
#include <boost/asio/redirect_error.hpp>
#include "test.hpp"
namespace boost {
namespace beast {
namespace websocket {
struct async_all_server_op : boost::asio::coroutine
{
stream<asio::ip::tcp::socket> & ws;
async_all_server_op(stream<asio::ip::tcp::socket> & ws) : ws(ws) {}
template<typename Self>
void operator()(Self && self, error_code ec = {}, std::size_t sz = 0)
{
if (ec)
return self.complete(ec);
BOOST_ASIO_CORO_REENTER(*this)
{
self.reset_cancellation_state([](net::cancellation_type ct){return ct;});
BOOST_ASIO_CORO_YIELD ws.async_handshake("test", "/", std::move(self));
BOOST_ASIO_CORO_YIELD ws.async_ping("", std::move(self));
BOOST_ASIO_CORO_YIELD ws.async_write_some(false, net::buffer("FOO", 3), std::move(self));
BOOST_ASIO_CORO_YIELD ws.async_write_some(true, net::buffer("BAR", 3), std::move(self));
BOOST_ASIO_CORO_YIELD ws.async_close("testing", std::move(self));
self.complete({});
}
}
};
template<BOOST_BEAST_ASYNC_TPARAM1 CompletionToken>
BOOST_BEAST_ASYNC_RESULT1(CompletionToken)
async_all_server(
stream<asio::ip::tcp::socket> & ws,
CompletionToken && token)
{
return net::async_compose<CompletionToken, void(error_code)>
(
async_all_server_op{ws},
token, ws
);
}
struct async_all_client_op : boost::asio::coroutine
{
stream<asio::ip::tcp::socket> & ws;
async_all_client_op(stream<asio::ip::tcp::socket> & ws) : ws(ws) {}
struct impl_t
{
impl_t () = default;
std::string res;
net::dynamic_string_buffer<char, std::char_traits<char>, std::allocator<char>> buf =
net::dynamic_buffer(res);
};
std::shared_ptr<impl_t> impl{std::make_shared<impl_t>()};
template<typename Self>
void operator()(Self && self, error_code ec = {}, std::size_t sz = 0)
{
if (ec)
return self.complete(ec);
BOOST_ASIO_CORO_REENTER(*this)
{
// let everything pass
self.reset_cancellation_state([](net::cancellation_type ct){return ct;});
BOOST_ASIO_CORO_YIELD ws.async_accept(std::move(self));
BOOST_ASIO_CORO_YIELD ws.async_pong("", std::move(self));
BOOST_ASIO_CORO_YIELD ws.async_read(impl->buf, std::move(self));
BEAST_EXPECTS(impl->res == "FOOBAR", impl->res);
BOOST_ASIO_CORO_YIELD ws.async_read(impl->buf, std::move(self));
BEAST_EXPECTS(ec == websocket::error::closed
|| ec == net::error::connection_reset
// hard coded winapi error (WSAECONNRESET), same as connection_reset, but asio delivers it with system_category
|| (ec.value() == 10054
&& ec.category() == system::system_category())
|| ec == net::error::not_connected, ec.message());
self.complete({});
}
}
};
template<BOOST_BEAST_ASYNC_TPARAM1 CompletionToken>
BOOST_BEAST_ASYNC_RESULT1(CompletionToken)
async_all_client(
stream<asio::ip::tcp::socket> & ws,
CompletionToken && token)
{
return net::async_compose<CompletionToken, void(error_code)>
(
async_all_client_op{ws},
token, ws
);
}
class cancel_test : public websocket_test_suite
{
public:
std::size_t run_impl(net::cancellation_signal &sl,
net::io_context &ctx,
std::size_t trigger = -1,
net::cancellation_type tp = net::cancellation_type::terminal)
{
std::size_t cnt = 0;
std::size_t res = 0u;
while ((res = ctx.run_one()) != 0)
if (trigger == cnt ++)
{
sl.emit(tp);
}
return cnt;
}
std::size_t testAll(std::size_t & cancel_counter,
bool cancel_server = false,
std::size_t trigger = -1)
{
net::cancellation_signal sig1, sig2;
net::io_context ioc;
using tcp = net::ip::tcp;
stream<tcp::socket> ws1(ioc.get_executor());
stream<tcp::socket> ws2(ioc.get_executor());
test::connect(ws1.next_layer(), ws2.next_layer());
async_all_server(ws1,
net::bind_cancellation_slot(sig1.slot(), [&](system::error_code ec)
{
if (ec)
{
if (ec == net::error::operation_aborted &&
(cancel_server || trigger == static_cast<std::size_t>(-1)))
cancel_counter++;
BEAST_EXPECTS(ec == net::error::operation_aborted
|| ec == net::error::broken_pipe
// winapi WSAECONNRESET, as system_category
|| ec == error_code(10054, boost::system::system_category())
|| ec == net::error::bad_descriptor
|| ec == net::error::eof, ec.message());
get_lowest_layer(ws1).close();
}
}));
async_all_client(
ws2,
net::bind_cancellation_slot(sig2.slot(), [&](system::error_code ec)
{
if (ec)
{
if (ec == net::error::operation_aborted &&
(!cancel_server || trigger == static_cast<std::size_t>(-1)))
cancel_counter++;
BEAST_EXPECTS(ec == net::error::operation_aborted
|| ec == error::closed
|| ec == net::error::broken_pipe
|| ec == net::error::connection_reset
|| ec == net::error::not_connected
// winapi WSAECONNRESET, as system_category
|| ec == error_code(10054, boost::system::system_category())
|| ec == net::error::eof, ec.message());
get_lowest_layer(ws1).close();
}
}));
return run_impl(cancel_server ? sig1 : sig2, ioc, trigger);
}
void brute_force()
{
std::size_t cancel_counter = 0;
const auto init = testAll(cancel_counter);
BEAST_EXPECT(cancel_counter == 0u);
for (std::size_t cnt = 0; cnt < init; cnt ++)
testAll(cancel_counter, true, cnt);
BEAST_EXPECT(cancel_counter > 0u);
cancel_counter = 0u;
for (std::size_t cnt = 0; cnt < init; cnt ++)
testAll(cancel_counter, false, cnt);
BEAST_EXPECT(cancel_counter > 0u);
}
void
run() override
{
brute_force();
}
};
BEAST_DEFINE_TESTSUITE(beast,websocket,cancel);
} // websocket
} // beast
} // boost
+30
View File
@@ -251,6 +251,9 @@ struct handler
using executor_type = boost::asio::io_context::executor_type;
executor_type get_executor() const noexcept;
using cancellation_slot_type = boost::asio::cancellation_slot;
cancellation_slot_type get_cancellation_slot() const noexcept;
void operator()(boost::beast::error_code, std::size_t);
};
//]
@@ -265,6 +268,12 @@ inline auto handler::get_executor() const noexcept ->
static boost::asio::io_context ioc;
return ioc.get_executor();
}
inline auto handler::get_cancellation_slot() const noexcept ->
cancellation_slot_type
{
return cancellation_slot_type();
}
inline void handler::operator()(
boost::beast::error_code, std::size_t)
{
@@ -296,6 +305,18 @@ struct associated_executor<handler, Executor>
Executor const& ex = Executor{}) noexcept;
};
template<class CancellationSlot>
struct associated_cancellation_slot<handler, CancellationSlot>
{
using type = cancellation_slot;
static
type
get(handler const& h,
CancellationSlot const& cs = CancellationSlot{}) noexcept;
};
} // boost
} // asio
//]
@@ -315,6 +336,15 @@ get(handler const&, Executor const&) noexcept -> type
return {};
}
template<class CancellationSlot>
auto
boost::asio::associated_cancellation_slot<handler, CancellationSlot>::
get(handler const&, CancellationSlot const&) noexcept -> type
{
return {};
}
//------------------------------------------------------------------------------
namespace boost {