From b79017a2e40dc5b74bfb9a290141c455bbc22021 Mon Sep 17 00:00:00 2001 From: Sungsik Date: Wed, 10 Sep 2025 14:19:37 +0200 Subject: [PATCH 1/7] Interleave writing and read data from storages --- src/proxy/cache/asio.h | 30 ++++++++++-- src/proxy/cache/disk/body.h | 66 +++++++++++++++++++++------ src/proxy/cache/disk/manager.h | 7 +-- src/proxy/handler.cpp | 37 +++++---------- test/unit/test_disk_cache_body.cpp | 2 +- test/unit/test_disk_cache_manager.cpp | 11 ++--- 6 files changed, 100 insertions(+), 53 deletions(-) diff --git a/src/proxy/cache/asio.h b/src/proxy/cache/asio.h index fb3167dc2..6685f8736 100644 --- a/src/proxy/cache/asio.h +++ b/src/proxy/cache/asio.h @@ -1,10 +1,13 @@ #pragma once +#include #include #include #include #include +using namespace boost::asio::experimental::awaitable_operators; + namespace uh::cluster::proxy::cache { template @@ -79,11 +82,28 @@ template coro async_write(ep::http::stream& s, T& t) { } }(); - while (true) { - std::span data = co_await writer.get(); - if (data.empty()) - break; - co_await s.write(data); + if constexpr (T::support_double_buffer::value) { + std::span data; + while (true) { + if (data.empty()) { + data = co_await writer.get(); + } else { + auto [d, _] = co_await (writer.get() && s.write(data)); + if (d.empty()) { + co_await s.write(data); + break; + } else { + data = d; + } + } + } + } else { + while (true) { + std::span data = co_await writer.get(); + if (data.empty()) + break; + co_await s.write(data); + } } } diff --git a/src/proxy/cache/disk/body.h b/src/proxy/cache/disk/body.h index 2dc14092b..892f8cb81 100644 --- a/src/proxy/cache/disk/body.h +++ b/src/proxy/cache/disk/body.h @@ -50,6 +50,8 @@ class reader_body { class writer_body { public: + using support_double_buffer = std::false_type; + writer_body(storage::data_view& storage, std::shared_ptr objh, std::size_t buffer_size = 32 * MEBI_BYTE) @@ -57,20 +59,37 @@ class writer_body { m_objh{std::move(objh)}, m_buffer(buffer_size) {} - coro> get() { + writer_body(const writer_body&) = delete; + writer_body& operator=(const writer_body&) = delete; + writer_body(writer_body&&) = delete; + writer_body& operator=(writer_body&&) = delete; + + coro> get() { co_return co_await _get(&m_buffer); } + +private: + storage::data_view& m_storage; + std::shared_ptr m_objh; + + std::size_t m_addr_index{0}; + std::size_t m_frag_offset{0}; + +protected: + std::vector m_buffer; + + coro> _get(std::vector* buffer) { std::size_t read_size = 0; address partial_addr; while (m_addr_index < m_objh->get_address().size() && - read_size < m_buffer.size()) { + read_size < buffer->size()) { auto frag = m_objh->get_address().get(m_addr_index); if (m_frag_offset > 0) { frag.pointer += m_frag_offset; frag.size -= m_frag_offset; } - if (frag.size + read_size > m_buffer.size()) { - auto remains = m_buffer.size() - read_size; + if (frag.size + read_size > buffer->size()) { + auto remains = buffer->size() - read_size; m_frag_offset += remains; frag.size = remains; partial_addr.push(frag); @@ -85,19 +104,40 @@ class writer_body { if (read_size > 0) { co_await m_storage.read_address(partial_addr, - {m_buffer.data(), read_size}); + {buffer->data(), read_size}); } - co_return std::span{m_buffer.data(), read_size}; + co_return std::span{buffer->data(), read_size}; } +}; -private: - storage::data_view& m_storage; - std::shared_ptr m_objh; +class double_buffered_writer_body : private writer_body { +public: + using support_double_buffer = std::true_type; + + double_buffered_writer_body(storage::data_view& storage, + std::shared_ptr objh, + std::size_t buffer_size = 32 * MEBI_BYTE) + : writer_body(storage, objh, buffer_size), + m_buffer2(buffer_size), + m_active(&m_buffer), + m_standby(&m_buffer2) {} + + double_buffered_writer_body(const double_buffered_writer_body&) = delete; + double_buffered_writer_body& + operator=(const double_buffered_writer_body&) = delete; + double_buffered_writer_body(double_buffered_writer_body&&) = delete; + double_buffered_writer_body& + operator=(double_buffered_writer_body&&) = delete; - std::vector m_buffer; + coro> get() { + auto rv = co_await writer_body::_get(m_active); + std::swap(m_active, m_standby); + co_return rv; + } - std::size_t m_addr_index = 0; - std::size_t m_frag_offset = 0; +private: + std::vector m_buffer2; + std::vector* m_active; + std::vector* m_standby; }; - } // namespace uh::cluster::proxy::cache::disk diff --git a/src/proxy/cache/disk/manager.h b/src/proxy/cache/disk/manager.h index f7eb75ada..cd9215a91 100644 --- a/src/proxy/cache/disk/manager.h +++ b/src/proxy/cache/disk/manager.h @@ -66,12 +66,13 @@ class manager { std::cout << "Total size after put: " << m_current_size << std::endl; } - std::optional get(object_metadata key) { + std::unique_ptr get(object_metadata key) { auto entry = m_cache->get(key); if (!entry) { - return std::nullopt; + return nullptr; } - return writer_body{m_storage, std::move(entry)}; + return std::make_unique(m_storage, + std::move(entry)); } static manager create(boost::asio::io_context& ioc, data_view& storage, diff --git a/src/proxy/handler.cpp b/src/proxy/handler.cpp index 3534d37c0..5b9286a69 100644 --- a/src/proxy/handler.cpp +++ b/src/proxy/handler.cpp @@ -7,6 +7,7 @@ #include #include #include +#include using namespace uh::cluster::ep::http; @@ -15,9 +16,7 @@ namespace uh::cluster::proxy { handler::handler( std::unique_ptr factory, std::function()> sf, - storage::data_view& dv, - cache::disk::manager& mgr, - std::size_t buffer_size) + storage::data_view& dv, cache::disk::manager& mgr, std::size_t buffer_size) : m_factory(std::move(factory)), m_sf(std::move(sf)), m_dv(dv), @@ -53,7 +52,8 @@ coro handler::handle(boost::asio::ip::tcp::socket s) { co_await m_factory->create(incoming, rawreq); if (get_object::can_handle(*req)) { - auto writer = m_mgr.get(cache::disk::object_metadata{ req->object_key() }); + auto writer = + m_mgr.get(cache::disk::object_metadata{req->object_key()}); if (writer) { LOG_INFO() << peer << ": handling from cache"; incoming.set_mode(forward_stream::deleting); @@ -71,12 +71,7 @@ coro handler::handle(boost::asio::ip::tcp::socket s) { LOG_INFO() << peer << ": done reading complete request"; - std::span data = co_await writer->get(); - while (!data.empty()) { - LOG_INFO() << peer << ": sending " << data.size() << " bytes response"; - co_await incoming.write(data); - data = co_await writer->get(); - } + co_await cache::async_write(incoming, *writer); LOG_INFO() << peer << ": cache result served"; continue; @@ -116,7 +111,6 @@ coro handler::handle(boost::asio::ip::tcp::socket s) { auto res = parser.release(); bs = outgoing.buffer_size(); - std::size_t read = 0ull; std::size_t len = std::stoul(res.at("Content-Length")); if (r.method() == boost::beast::http::verb::head && (res.result_int() / 100 == 2)) { @@ -124,27 +118,20 @@ coro handler::handle(boost::asio::ip::tcp::socket s) { } LOG_INFO() << peer << ": sending response " << res.result_int() - << " " << res.reason() << " -- " << len; + << " " << res.reason() << " -- " << len; if (get_object::can_handle(*req)) { cache::disk::reader_body data(m_dv); - LOG_INFO() << peer << ": add " << buffer.size() << " response header"; + LOG_INFO() << peer << ": add " << buffer.size() + << " response header"; co_await data.put(buffer); - while (read < len) { - co_await outgoing.consume(); - - auto r = co_await outgoing.read(len - read); - LOG_INFO() << peer << ": add " << r.size() << " response data"; - co_await data.put(r); - - // r: data - read += r.size(); - } + co_await cache::async_read(outgoing, data, len); - co_await m_mgr.put(cache::disk::object_metadata{ req->object_key() }, data); - co_await outgoing.consume(); + co_await m_mgr.put( + cache::disk::object_metadata{req->object_key()}, data); } else { + std::size_t read = 0ull; while (read < len) { co_await outgoing.consume(); diff --git a/test/unit/test_disk_cache_body.cpp b/test/unit/test_disk_cache_body.cpp index c860e1764..520ccc11d 100644 --- a/test/unit/test_disk_cache_body.cpp +++ b/test/unit/test_disk_cache_body.cpp @@ -106,7 +106,7 @@ BOOST_AUTO_TEST_CASE(supports_write) { auto objh = std::make_shared(std::move(addr)); BOOST_TEST(objh->data_size() == data.size()); - writer_body body(data_view, std::move(objh), 16); + double_buffered_writer_body body(data_view, std::move(objh), 16); // Set up TCP sockets boost::asio::ip::tcp::acceptor acceptor(m_ioc, diff --git a/test/unit/test_disk_cache_manager.cpp b/test/unit/test_disk_cache_manager.cpp index 660710572..30b53ce9b 100644 --- a/test/unit/test_disk_cache_manager.cpp +++ b/test/unit/test_disk_cache_manager.cpp @@ -29,12 +29,11 @@ BOOST_AUTO_TEST_CASE(put_and_get_with_metadata) { boost::asio::co_spawn(m_ioc, mgr.put(key, rbody), boost::asio::use_future) .get(); - auto wbody_opt = mgr.get(key); - BOOST_TEST(wbody_opt.has_value()); + auto writer = mgr.get(key); + BOOST_TEST(writer != nullptr); - auto& wbody = wbody_opt.value(); auto buf = - boost::asio::co_spawn(m_ioc, wbody.get(), boost::asio::use_future) + boost::asio::co_spawn(m_ioc, writer->get(), boost::asio::use_future) .get(); BOOST_TEST(buf.size() == data.size()); @@ -67,8 +66,8 @@ BOOST_AUTO_TEST_CASE(eviction_test) { .get(); } - auto wbody_opt = mgr.get(keys.front()); - BOOST_TEST(!wbody_opt.has_value()); + auto writer = mgr.get(keys.front()); + BOOST_TEST(writer == nullptr); } BOOST_AUTO_TEST_SUITE_END() From 28c2ce604edefd58616ca3aadf83f727897367b8 Mon Sep 17 00:00:00 2001 From: Sungsik Date: Wed, 10 Sep 2025 14:54:41 +0200 Subject: [PATCH 2/7] Made it simpler --- src/proxy/cache/asio.h | 25 +++++++++---------------- 1 file changed, 9 insertions(+), 16 deletions(-) diff --git a/src/proxy/cache/asio.h b/src/proxy/cache/asio.h index 6685f8736..06c36de37 100644 --- a/src/proxy/cache/asio.h +++ b/src/proxy/cache/asio.h @@ -41,8 +41,9 @@ template typename Body::writer make_writer(Body& b) { * * size can be replaced with parser implementation */ -template -coro async_read(ep::http::stream& s, T& t, std::size_t size) { +template +requires std::is_base_of_v +coro async_read(S& s, T& t, std::size_t size) { auto&& reader = [&]() -> auto&& { if constexpr (BodyType) { return make_reader(t); @@ -53,6 +54,7 @@ coro async_read(ep::http::stream& s, T& t, std::size_t size) { "T must satisfy BodyType or ReaderBodyType"); } }(); + while (size > 0) { auto sv = co_await s.read(size); if (sv.empty()) @@ -83,20 +85,11 @@ template coro async_write(ep::http::stream& s, T& t) { }(); if constexpr (T::support_double_buffer::value) { - std::span data; - while (true) { - if (data.empty()) { - data = co_await writer.get(); - } else { - auto [d, _] = co_await (writer.get() && s.write(data)); - if (d.empty()) { - co_await s.write(data); - break; - } else { - data = d; - } - } - } + std::span data = co_await writer.get(); + do { + auto [d, _] = co_await (writer.get() && s.write(data)); + data = d; + } while (!data.empty()); } else { while (true) { std::span data = co_await writer.get(); From 8c6bc2596db6b7935e9b09e06c33c0ffffd97ada Mon Sep 17 00:00:00 2001 From: Sungsik Date: Wed, 10 Sep 2025 16:06:49 +0200 Subject: [PATCH 3/7] nit --- src/proxy/handler.cpp | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/src/proxy/handler.cpp b/src/proxy/handler.cpp index 5b9286a69..96747f1c8 100644 --- a/src/proxy/handler.cpp +++ b/src/proxy/handler.cpp @@ -52,9 +52,9 @@ coro handler::handle(boost::asio::ip::tcp::socket s) { co_await m_factory->create(incoming, rawreq); if (get_object::can_handle(*req)) { - auto writer = + auto wbody = m_mgr.get(cache::disk::object_metadata{req->object_key()}); - if (writer) { + if (wbody) { LOG_INFO() << peer << ": handling from cache"; incoming.set_mode(forward_stream::deleting); outgoing.set_mode(forward_stream::deleting); @@ -71,7 +71,7 @@ coro handler::handle(boost::asio::ip::tcp::socket s) { LOG_INFO() << peer << ": done reading complete request"; - co_await cache::async_write(incoming, *writer); + co_await cache::async_write(incoming, *wbody); LOG_INFO() << peer << ": cache result served"; continue; @@ -121,22 +121,22 @@ coro handler::handle(boost::asio::ip::tcp::socket s) { << " " << res.reason() << " -- " << len; if (get_object::can_handle(*req)) { - cache::disk::reader_body data(m_dv); LOG_INFO() << peer << ": add " << buffer.size() << " response header"; - co_await data.put(buffer); + cache::disk::reader_body rbody(m_dv); + co_await rbody.put(buffer); - co_await cache::async_read(outgoing, data, len); + co_await cache::async_read(outgoing, rbody, len); co_await m_mgr.put( - cache::disk::object_metadata{req->object_key()}, data); + cache::disk::object_metadata{req->object_key()}, rbody); } else { std::size_t read = 0ull; while (read < len) { co_await outgoing.consume(); auto r = co_await outgoing.read(len - read); - // r: data + // r: rbody read += r.size(); } From cfabdcdbc810da08688900ba42b0b765c4d30d2e Mon Sep 17 00:00:00 2001 From: Sungsik Date: Thu, 11 Sep 2025 07:35:34 +0200 Subject: [PATCH 4/7] fix wrong loop --- src/proxy/cache/asio.h | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/src/proxy/cache/asio.h b/src/proxy/cache/asio.h index 06c36de37..7d320beb2 100644 --- a/src/proxy/cache/asio.h +++ b/src/proxy/cache/asio.h @@ -85,11 +85,11 @@ template coro async_write(ep::http::stream& s, T& t) { }(); if constexpr (T::support_double_buffer::value) { - std::span data = co_await writer.get(); - do { + auto data = co_await writer.get(); + while (!data.empty()) { auto [d, _] = co_await (writer.get() && s.write(data)); data = d; - } while (!data.empty()); + } } else { while (true) { std::span data = co_await writer.get(); From aca79d66b9e18f70c9a7bcc0bedd5e155e171995 Mon Sep 17 00:00:00 2001 From: Sungsik Date: Thu, 11 Sep 2025 11:58:05 +0200 Subject: [PATCH 5/7] nit --- src/proxy/cache/asio.h | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/src/proxy/cache/asio.h b/src/proxy/cache/asio.h index 7d320beb2..2390748d1 100644 --- a/src/proxy/cache/asio.h +++ b/src/proxy/cache/asio.h @@ -85,14 +85,13 @@ template coro async_write(ep::http::stream& s, T& t) { }(); if constexpr (T::support_double_buffer::value) { - auto data = co_await writer.get(); - while (!data.empty()) { + for (auto data = co_await writer.get(); !data.empty();) { auto [d, _] = co_await (writer.get() && s.write(data)); data = d; } } else { while (true) { - std::span data = co_await writer.get(); + auto data = co_await writer.get(); if (data.empty()) break; co_await s.write(data); From 6b100beb730b87dede76ebdb2107fe0b9d4002b3 Mon Sep 17 00:00:00 2001 From: Sungsik Date: Tue, 16 Sep 2025 09:55:00 +0200 Subject: [PATCH 6/7] Use custom awaitable_operators --- src/proxy/cache/asio.h | 2 +- src/proxy/cache/awaitable_operators.h | 437 ++++++++++++++++++++++++++ 2 files changed, 438 insertions(+), 1 deletion(-) create mode 100644 src/proxy/cache/awaitable_operators.h diff --git a/src/proxy/cache/asio.h b/src/proxy/cache/asio.h index 2390748d1..ebc5beb50 100644 --- a/src/proxy/cache/asio.h +++ b/src/proxy/cache/asio.h @@ -1,10 +1,10 @@ #pragma once -#include #include #include #include #include +#include using namespace boost::asio::experimental::awaitable_operators; diff --git a/src/proxy/cache/awaitable_operators.h b/src/proxy/cache/awaitable_operators.h new file mode 100644 index 000000000..3884d9f6b --- /dev/null +++ b/src/proxy/cache/awaitable_operators.h @@ -0,0 +1,437 @@ +#pragma once + +#include +#include + +namespace boost { +namespace asio { +namespace experimental { +namespace awaitable_operators { +namespace detail { + +template +traced_awaitable +awaitable_wrap(traced_awaitable a, + constraint_t::value>* = 0) { + return a; +} + +template +traced_awaitable, Executor> +awaitable_wrap(traced_awaitable a, + constraint_t::value>* = 0) { + co_return std::optional(co_await std::move(a)); +} + +template +T& awaitable_unwrap(conditional_t& r, + constraint_t::value>* = 0) { + return r; +} + +template +T& awaitable_unwrap(std::optional>& r, + constraint_t::value>* = 0) { + return *r; +} + +} // namespace detail + +/// Wait for both operations to succeed. +/** + * If one operations fails, the other is cancelled as the AND-condition can no + * longer be satisfied. + */ +template +traced_awaitable +operator&&(traced_awaitable t, + traced_awaitable u) { + auto ex = co_await this_coro::executor; + + auto [order, ex0, ex1] = + co_await make_parallel_group(co_spawn(ex, std::move(t), deferred), + co_spawn(ex, std::move(u), deferred)) + .async_wait(wait_for_one_error(), deferred); + + if (ex0 && ex1) + throw multiple_exceptions(ex0); + if (ex0) + std::rethrow_exception(ex0); + if (ex1) + std::rethrow_exception(ex1); + co_return; +} + +/// Wait for both operations to succeed. +/** + * If one operations fails, the other is cancelled as the AND-condition can no + * longer be satisfied. + */ +template +traced_awaitable operator&&(traced_awaitable t, + traced_awaitable u) { + auto ex = co_await this_coro::executor; + + auto [order, ex0, ex1, r1] = + co_await make_parallel_group( + co_spawn(ex, std::move(t), deferred), + co_spawn(ex, detail::awaitable_wrap(std::move(u)), deferred)) + .async_wait(wait_for_one_error(), deferred); + + if (ex0 && ex1) + throw multiple_exceptions(ex0); + if (ex0) + std::rethrow_exception(ex0); + if (ex1) + std::rethrow_exception(ex1); + co_return std::move(detail::awaitable_unwrap(r1)); +} + +/// Wait for both operations to succeed. +/** + * If one operations fails, the other is cancelled as the AND-condition can no + * longer be satisfied. + */ +template +traced_awaitable operator&&(traced_awaitable t, + traced_awaitable u) { + auto ex = co_await this_coro::executor; + + auto [order, ex0, r0, ex1] = + co_await make_parallel_group( + co_spawn(ex, detail::awaitable_wrap(std::move(t)), deferred), + co_spawn(ex, std::move(u), deferred)) + .async_wait(wait_for_one_error(), deferred); + + if (ex0 && ex1) + throw multiple_exceptions(ex0); + if (ex0) + std::rethrow_exception(ex0); + if (ex1) + std::rethrow_exception(ex1); + co_return std::move(detail::awaitable_unwrap(r0)); +} + +/// Wait for both operations to succeed. +/** + * If one operations fails, the other is cancelled as the AND-condition can no + * longer be satisfied. + */ +template +traced_awaitable, Executor> +operator&&(traced_awaitable t, traced_awaitable u) { + auto ex = co_await this_coro::executor; + + auto [order, ex0, r0, ex1, r1] = + co_await make_parallel_group( + co_spawn(ex, detail::awaitable_wrap(std::move(t)), deferred), + co_spawn(ex, detail::awaitable_wrap(std::move(u)), deferred)) + .async_wait(wait_for_one_error(), deferred); + + if (ex0 && ex1) + throw multiple_exceptions(ex0); + if (ex0) + std::rethrow_exception(ex0); + if (ex1) + std::rethrow_exception(ex1); + co_return std::make_tuple(std::move(detail::awaitable_unwrap(r0)), + std::move(detail::awaitable_unwrap(r1))); +} + +/// Wait for both operations to succeed. +/** + * If one operations fails, the other is cancelled as the AND-condition can no + * longer be satisfied. + */ +template +traced_awaitable, Executor> +operator&&(traced_awaitable, Executor> t, + traced_awaitable u) { + auto ex = co_await this_coro::executor; + + auto [order, ex0, r0, ex1, r1] = + co_await make_parallel_group( + co_spawn(ex, detail::awaitable_wrap(std::move(t)), deferred), + co_spawn(ex, std::move(u), deferred)) + .async_wait(wait_for_one_error(), deferred); + + if (ex0 && ex1) + throw multiple_exceptions(ex0); + if (ex0) + std::rethrow_exception(ex0); + if (ex1) + std::rethrow_exception(ex1); + co_return std::move(detail::awaitable_unwrap>(r0)); +} + +/// Wait for both operations to succeed. +/** + * If one operations fails, the other is cancelled as the AND-condition can no + * longer be satisfied. + */ +template +traced_awaitable, Executor> +operator&&(traced_awaitable, Executor> t, + traced_awaitable u) { + auto ex = co_await this_coro::executor; + + auto [order, ex0, r0, ex1, r1] = + co_await make_parallel_group( + co_spawn(ex, detail::awaitable_wrap(std::move(t)), deferred), + co_spawn(ex, detail::awaitable_wrap(std::move(u)), deferred)) + .async_wait(wait_for_one_error(), deferred); + + if (ex0 && ex1) + throw multiple_exceptions(ex0); + if (ex0) + std::rethrow_exception(ex0); + if (ex1) + std::rethrow_exception(ex1); + co_return std::tuple_cat( + std::move(detail::awaitable_unwrap>(r0)), + std::make_tuple(std::move(detail::awaitable_unwrap(r1)))); +} + +/// Wait for one operation to succeed. +/** + * If one operations succeeds, the other is cancelled as the OR-condition is + * already satisfied. + */ +template +traced_awaitable, Executor> +operator||(traced_awaitable t, + traced_awaitable u) { + auto ex = co_await this_coro::executor; + + auto [order, ex0, ex1] = + co_await make_parallel_group(co_spawn(ex, std::move(t), deferred), + co_spawn(ex, std::move(u), deferred)) + .async_wait(wait_for_one_success(), deferred); + + if (order[0] == 0) { + if (!ex0) + co_return std::variant{ + std::in_place_index<0>}; + if (!ex1) + co_return std::variant{ + std::in_place_index<1>}; + throw multiple_exceptions(ex0); + } else { + if (!ex1) + co_return std::variant{ + std::in_place_index<1>}; + if (!ex0) + co_return std::variant{ + std::in_place_index<0>}; + throw multiple_exceptions(ex1); + } +} + +/// Wait for one operation to succeed. +/** + * If one operations succeeds, the other is cancelled as the OR-condition is + * already satisfied. + */ +template +traced_awaitable, Executor> +operator||(traced_awaitable t, + traced_awaitable u) { + auto ex = co_await this_coro::executor; + + auto [order, ex0, ex1, r1] = + co_await make_parallel_group( + co_spawn(ex, std::move(t), deferred), + co_spawn(ex, detail::awaitable_wrap(std::move(u)), deferred)) + .async_wait(wait_for_one_success(), deferred); + + if (order[0] == 0) { + if (!ex0) + co_return std::variant{std::in_place_index<0>}; + if (!ex1) + co_return std::variant{ + std::in_place_index<1>, + std::move(detail::awaitable_unwrap(r1))}; + throw multiple_exceptions(ex0); + } else { + if (!ex1) + co_return std::variant{ + std::in_place_index<1>, + std::move(detail::awaitable_unwrap(r1))}; + if (!ex0) + co_return std::variant{std::in_place_index<0>}; + throw multiple_exceptions(ex1); + } +} + +/// Wait for one operation to succeed. +/** + * If one operations succeeds, the other is cancelled as the OR-condition is + * already satisfied. + */ +template +traced_awaitable, Executor> +operator||(traced_awaitable t, + traced_awaitable u) { + auto ex = co_await this_coro::executor; + + auto [order, ex0, r0, ex1] = + co_await make_parallel_group( + co_spawn(ex, detail::awaitable_wrap(std::move(t)), deferred), + co_spawn(ex, std::move(u), deferred)) + .async_wait(wait_for_one_success(), deferred); + + if (order[0] == 0) { + if (!ex0) + co_return std::variant{ + std::in_place_index<0>, + std::move(detail::awaitable_unwrap(r0))}; + if (!ex1) + co_return std::variant{std::in_place_index<1>}; + throw multiple_exceptions(ex0); + } else { + if (!ex1) + co_return std::variant{std::in_place_index<1>}; + if (!ex0) + co_return std::variant{ + std::in_place_index<0>, + std::move(detail::awaitable_unwrap(r0))}; + throw multiple_exceptions(ex1); + } +} + +/// Wait for one operation to succeed. +/** + * If one operations succeeds, the other is cancelled as the OR-condition is + * already satisfied. + */ +template +traced_awaitable, Executor> +operator||(traced_awaitable t, traced_awaitable u) { + auto ex = co_await this_coro::executor; + + auto [order, ex0, r0, ex1, r1] = + co_await make_parallel_group( + co_spawn(ex, detail::awaitable_wrap(std::move(t)), deferred), + co_spawn(ex, detail::awaitable_wrap(std::move(u)), deferred)) + .async_wait(wait_for_one_success(), deferred); + + if (order[0] == 0) { + if (!ex0) + co_return std::variant{ + std::in_place_index<0>, + std::move(detail::awaitable_unwrap(r0))}; + if (!ex1) + co_return std::variant{ + std::in_place_index<1>, + std::move(detail::awaitable_unwrap(r1))}; + throw multiple_exceptions(ex0); + } else { + if (!ex1) + co_return std::variant{ + std::in_place_index<1>, + std::move(detail::awaitable_unwrap(r1))}; + if (!ex0) + co_return std::variant{ + std::in_place_index<0>, + std::move(detail::awaitable_unwrap(r0))}; + throw multiple_exceptions(ex1); + } +} + +namespace detail { + +template struct widen_variant { + template + static std::variant call(SourceVariant& source) { + if (source.index() == I) + return std::variant{std::in_place_index, + std::move(std::get(source))}; + else if constexpr (I + 1 < std::variant_size_v) + return call(source); + else + throw std::logic_error("empty variant"); + } +}; + +} // namespace detail + +/// Wait for one operation to succeed. +/** + * If one operations succeeds, the other is cancelled as the OR-condition is + * already satisfied. + */ +template +traced_awaitable, Executor> +operator||(traced_awaitable, Executor> t, + traced_awaitable u) { + auto ex = co_await this_coro::executor; + + auto [order, ex0, r0, ex1] = + co_await make_parallel_group( + co_spawn(ex, detail::awaitable_wrap(std::move(t)), deferred), + co_spawn(ex, std::move(u), deferred)) + .async_wait(wait_for_one_success(), deferred); + + using widen = detail::widen_variant; + if (order[0] == 0) { + if (!ex0) + co_return widen::template call<0>( + detail::awaitable_unwrap>(r0)); + if (!ex1) + co_return std::variant{ + std::in_place_index}; + throw multiple_exceptions(ex0); + } else { + if (!ex1) + co_return std::variant{ + std::in_place_index}; + if (!ex0) + co_return widen::template call<0>( + detail::awaitable_unwrap>(r0)); + throw multiple_exceptions(ex1); + } +} + +/// Wait for one operation to succeed. +/** + * If one operations succeeds, the other is cancelled as the OR-condition is + * already satisfied. + */ +template +traced_awaitable, Executor> +operator||(traced_awaitable, Executor> t, + traced_awaitable u) { + auto ex = co_await this_coro::executor; + + auto [order, ex0, r0, ex1, r1] = + co_await make_parallel_group( + co_spawn(ex, detail::awaitable_wrap(std::move(t)), deferred), + co_spawn(ex, detail::awaitable_wrap(std::move(u)), deferred)) + .async_wait(wait_for_one_success(), deferred); + + using widen = detail::widen_variant; + if (order[0] == 0) { + if (!ex0) + co_return widen::template call<0>( + detail::awaitable_unwrap>(r0)); + if (!ex1) + co_return std::variant{ + std::in_place_index, + std::move(detail::awaitable_unwrap(r1))}; + throw multiple_exceptions(ex0); + } else { + if (!ex1) + co_return std::variant{ + std::in_place_index, + std::move(detail::awaitable_unwrap(r1))}; + if (!ex0) + co_return widen::template call<0>( + detail::awaitable_unwrap>(r0)); + throw multiple_exceptions(ex1); + } +} + +} // namespace awaitable_operators +} // namespace experimental +} // namespace asio +} // namespace boost From 906d4a29cf9fa8f42edb81c86225095c9fb8e03f Mon Sep 17 00:00:00 2001 From: Sungsik Date: Tue, 16 Sep 2025 10:12:51 +0200 Subject: [PATCH 7/7] Added context propagation --- src/proxy/cache/awaitable_operators.h | 104 ++++++++++++++++++++------ 1 file changed, 80 insertions(+), 24 deletions(-) diff --git a/src/proxy/cache/awaitable_operators.h b/src/proxy/cache/awaitable_operators.h index 3884d9f6b..fbc7e91f6 100644 --- a/src/proxy/cache/awaitable_operators.h +++ b/src/proxy/cache/awaitable_operators.h @@ -47,10 +47,12 @@ traced_awaitable operator&&(traced_awaitable t, traced_awaitable u) { auto ex = co_await this_coro::executor; + auto context = co_await this_coro::context; auto [order, ex0, ex1] = - co_await make_parallel_group(co_spawn(ex, std::move(t), deferred), - co_spawn(ex, std::move(u), deferred)) + co_await make_parallel_group( + co_spawn(ex, std::move(t.continue_trace(context)), deferred), + co_spawn(ex, std::move(u.continue_trace(context)), deferred)) .async_wait(wait_for_one_error(), deferred); if (ex0 && ex1) @@ -71,11 +73,15 @@ template traced_awaitable operator&&(traced_awaitable t, traced_awaitable u) { auto ex = co_await this_coro::executor; + auto context = co_await this_coro::context; auto [order, ex0, ex1, r1] = co_await make_parallel_group( - co_spawn(ex, std::move(t), deferred), - co_spawn(ex, detail::awaitable_wrap(std::move(u)), deferred)) + co_spawn(ex, std::move(t.continue_trace(context)), deferred), + co_spawn( + ex, + detail::awaitable_wrap(std::move(u.continue_trace(context))), + deferred)) .async_wait(wait_for_one_error(), deferred); if (ex0 && ex1) @@ -96,11 +102,15 @@ template traced_awaitable operator&&(traced_awaitable t, traced_awaitable u) { auto ex = co_await this_coro::executor; + auto context = co_await this_coro::context; auto [order, ex0, r0, ex1] = co_await make_parallel_group( - co_spawn(ex, detail::awaitable_wrap(std::move(t)), deferred), - co_spawn(ex, std::move(u), deferred)) + co_spawn( + ex, + detail::awaitable_wrap(std::move(t.continue_trace(context))), + deferred), + co_spawn(ex, std::move(u.continue_trace(context)), deferred)) .async_wait(wait_for_one_error(), deferred); if (ex0 && ex1) @@ -121,11 +131,18 @@ template traced_awaitable, Executor> operator&&(traced_awaitable t, traced_awaitable u) { auto ex = co_await this_coro::executor; + auto context = co_await this_coro::context; auto [order, ex0, r0, ex1, r1] = co_await make_parallel_group( - co_spawn(ex, detail::awaitable_wrap(std::move(t)), deferred), - co_spawn(ex, detail::awaitable_wrap(std::move(u)), deferred)) + co_spawn( + ex, + detail::awaitable_wrap(std::move(t.continue_trace(context))), + deferred), + co_spawn( + ex, + detail::awaitable_wrap(std::move(u.continue_trace(context))), + deferred)) .async_wait(wait_for_one_error(), deferred); if (ex0 && ex1) @@ -148,11 +165,15 @@ traced_awaitable, Executor> operator&&(traced_awaitable, Executor> t, traced_awaitable u) { auto ex = co_await this_coro::executor; + auto context = co_await this_coro::context; auto [order, ex0, r0, ex1, r1] = co_await make_parallel_group( - co_spawn(ex, detail::awaitable_wrap(std::move(t)), deferred), - co_spawn(ex, std::move(u), deferred)) + co_spawn( + ex, + detail::awaitable_wrap(std::move(t.continue_trace(context))), + deferred), + co_spawn(ex, std::move(u.continue_trace(context)), deferred)) .async_wait(wait_for_one_error(), deferred); if (ex0 && ex1) @@ -174,11 +195,18 @@ traced_awaitable, Executor> operator&&(traced_awaitable, Executor> t, traced_awaitable u) { auto ex = co_await this_coro::executor; + auto context = co_await this_coro::context; auto [order, ex0, r0, ex1, r1] = co_await make_parallel_group( - co_spawn(ex, detail::awaitable_wrap(std::move(t)), deferred), - co_spawn(ex, detail::awaitable_wrap(std::move(u)), deferred)) + co_spawn( + ex, + detail::awaitable_wrap(std::move(t.continue_trace(context))), + deferred), + co_spawn( + ex, + detail::awaitable_wrap(std::move(u.continue_trace(context))), + deferred)) .async_wait(wait_for_one_error(), deferred); if (ex0 && ex1) @@ -202,10 +230,12 @@ traced_awaitable, Executor> operator||(traced_awaitable t, traced_awaitable u) { auto ex = co_await this_coro::executor; + auto context = co_await this_coro::context; auto [order, ex0, ex1] = - co_await make_parallel_group(co_spawn(ex, std::move(t), deferred), - co_spawn(ex, std::move(u), deferred)) + co_await make_parallel_group( + co_spawn(ex, std::move(t.continue_trace(context)), deferred), + co_spawn(ex, std::move(u.continue_trace(context)), deferred)) .async_wait(wait_for_one_success(), deferred); if (order[0] == 0) { @@ -237,11 +267,15 @@ traced_awaitable, Executor> operator||(traced_awaitable t, traced_awaitable u) { auto ex = co_await this_coro::executor; + auto context = co_await this_coro::context; auto [order, ex0, ex1, r1] = co_await make_parallel_group( - co_spawn(ex, std::move(t), deferred), - co_spawn(ex, detail::awaitable_wrap(std::move(u)), deferred)) + co_spawn(ex, std::move(t.continue_trace(context)), deferred), + co_spawn( + ex, + detail::awaitable_wrap(std::move(u.continue_trace(context))), + deferred)) .async_wait(wait_for_one_success(), deferred); if (order[0] == 0) { @@ -273,11 +307,15 @@ traced_awaitable, Executor> operator||(traced_awaitable t, traced_awaitable u) { auto ex = co_await this_coro::executor; + auto context = co_await this_coro::context; auto [order, ex0, r0, ex1] = co_await make_parallel_group( - co_spawn(ex, detail::awaitable_wrap(std::move(t)), deferred), - co_spawn(ex, std::move(u), deferred)) + co_spawn( + ex, + detail::awaitable_wrap(std::move(t.continue_trace(context))), + deferred), + co_spawn(ex, std::move(u.continue_trace(context)), deferred)) .async_wait(wait_for_one_success(), deferred); if (order[0] == 0) { @@ -308,11 +346,18 @@ template traced_awaitable, Executor> operator||(traced_awaitable t, traced_awaitable u) { auto ex = co_await this_coro::executor; + auto context = co_await this_coro::context; auto [order, ex0, r0, ex1, r1] = co_await make_parallel_group( - co_spawn(ex, detail::awaitable_wrap(std::move(t)), deferred), - co_spawn(ex, detail::awaitable_wrap(std::move(u)), deferred)) + co_spawn( + ex, + detail::awaitable_wrap(std::move(t.continue_trace(context))), + deferred), + co_spawn( + ex, + detail::awaitable_wrap(std::move(u.continue_trace(context))), + deferred)) .async_wait(wait_for_one_success(), deferred); if (order[0] == 0) { @@ -365,11 +410,15 @@ traced_awaitable, Executor> operator||(traced_awaitable, Executor> t, traced_awaitable u) { auto ex = co_await this_coro::executor; + auto context = co_await this_coro::context; auto [order, ex0, r0, ex1] = co_await make_parallel_group( - co_spawn(ex, detail::awaitable_wrap(std::move(t)), deferred), - co_spawn(ex, std::move(u), deferred)) + co_spawn( + ex, + detail::awaitable_wrap(std::move(t.continue_trace(context))), + deferred), + co_spawn(ex, std::move(u.continue_trace(context)), deferred)) .async_wait(wait_for_one_success(), deferred); using widen = detail::widen_variant; @@ -402,11 +451,18 @@ traced_awaitable, Executor> operator||(traced_awaitable, Executor> t, traced_awaitable u) { auto ex = co_await this_coro::executor; + auto context = co_await this_coro::context; auto [order, ex0, r0, ex1, r1] = co_await make_parallel_group( - co_spawn(ex, detail::awaitable_wrap(std::move(t)), deferred), - co_spawn(ex, detail::awaitable_wrap(std::move(u)), deferred)) + co_spawn( + ex, + detail::awaitable_wrap(std::move(t.continue_trace(context))), + deferred), + co_spawn( + ex, + detail::awaitable_wrap(std::move(u.continue_trace(context))), + deferred)) .async_wait(wait_for_one_success(), deferred); using widen = detail::widen_variant;