From 76b81b640276ddcbe0068d2c61a6befbc84565b6 Mon Sep 17 00:00:00 2001 From: Sungsik Nam Date: Wed, 24 Sep 2025 16:52:50 +0200 Subject: [PATCH 01/10] done, with some dev changes (will be removed) --- scripts/start.sh.in | 2 +- src/proxy/CMakeLists.txt | 5 +- src/proxy/forward_stream.cpp | 27 ----- src/proxy/forward_stream.h | 22 +++- src/proxy/handler.cpp | 198 --------------------------------- src/proxy/handler.h | 209 +++++++++++++++++++++++++++++++++-- src/proxy/service.cpp | 58 ++++++---- 7 files changed, 255 insertions(+), 266 deletions(-) delete mode 100644 src/proxy/forward_stream.cpp delete mode 100644 src/proxy/handler.cpp diff --git a/scripts/start.sh.in b/scripts/start.sh.in index 3a903be83..44410f0c4 100755 --- a/scripts/start.sh.in +++ b/scripts/start.sh.in @@ -187,7 +187,7 @@ $UH_CLUSTER entrypoint &> $log_dir/entrypoint.log & pid_entrypoint=$! export OTEL_RESOURCE_ATTRIBUTES="service.name=proxy" -$UH_CLUSTER --downstream-port 8080 --downstream-host localhost proxy &> $log_dir/proxy.log & +$UH_CLUSTER --downstream-port 8081 --downstream-host localhost proxy &> $log_dir/proxy.log & pid_proxy=$! get_running_processes() { diff --git a/src/proxy/CMakeLists.txt b/src/proxy/CMakeLists.txt index 1c929340d..82644719a 100644 --- a/src/proxy/CMakeLists.txt +++ b/src/proxy/CMakeLists.txt @@ -1,2 +1,3 @@ -add_library(proxy service.cpp forward_stream.cpp request_factory.cpp handler.cpp) -target_link_libraries(proxy types utils network entrypoint) +find_package(OpenSSL REQUIRED) +add_library(proxy service.cpp request_factory.cpp) +target_link_libraries(proxy types utils network entrypoint OpenSSL::SSL OpenSSL::Crypto) diff --git a/src/proxy/forward_stream.cpp b/src/proxy/forward_stream.cpp deleted file mode 100644 index fb547b4cc..000000000 --- a/src/proxy/forward_stream.cpp +++ /dev/null @@ -1,27 +0,0 @@ -#include "forward_stream.h" - -#include - -namespace uh::cluster::proxy { - -forward_stream::forward_stream(boost::asio::ip::tcp::socket& s, - boost::asio::ip::tcp::socket& to, - std::size_t buffer_size) - : socket_stream(s, buffer_size), - m_to(to) { -} - -coro forward_stream::consume() { - if (m_mode == forwarding) { - LOG_DEBUG() << peer() << " forwarding " << buffer().size() << " bytes to " << m_to.remote_endpoint(); - co_await boost::asio::async_write(m_to, boost::asio::buffer(buffer())); - } - - co_await socket_stream::consume(); -} - -void forward_stream::set_mode(mode m) { - m_mode = m; -} - -} // namespace uh::cluster::proxy diff --git a/src/proxy/forward_stream.h b/src/proxy/forward_stream.h index 572b91eab..c6e85ab34 100644 --- a/src/proxy/forward_stream.h +++ b/src/proxy/forward_stream.h @@ -1,5 +1,6 @@ #pragma once +#include #include namespace uh::cluster::proxy { @@ -7,24 +8,33 @@ namespace uh::cluster::proxy { /** * Copy read data to additional socket. */ +template class forward_stream : public ep::http::socket_stream { public: /** * Create a stream that reads incoming data from `s` and forwards * it to the configured downstream socket `to`. */ - forward_stream(boost::asio::ip::tcp::socket& s, - boost::asio::ip::tcp::socket& to, - std::size_t buffer_size = 4 * MEBI_BYTE); + forward_stream(boost::asio::ip::tcp::socket& s, OutgoingStream& to, + std::size_t buffer_size = 4 * MEBI_BYTE) + : socket_stream(s, buffer_size), + m_to(to) {} - coro consume() override; + coro consume() override { + if (m_mode == forwarding) { + co_await boost::asio::async_write(m_to, + boost::asio::buffer(buffer())); + } + + co_await socket_stream::consume(); + } enum mode { forwarding, deleting }; - void set_mode(mode m); + void set_mode(mode m) { m_mode = m; } private: - boost::asio::ip::tcp::socket& m_to; + OutgoingStream& m_to; mode m_mode = deleting; }; diff --git a/src/proxy/handler.cpp b/src/proxy/handler.cpp deleted file mode 100644 index f8e1ae5e0..000000000 --- a/src/proxy/handler.cpp +++ /dev/null @@ -1,198 +0,0 @@ -#include "handler.h" - -#include "forward_stream.h" - -#include -#include - -#include -#include -#include - -#include -#include -#include -#include -#include - -using namespace boost::beast; -using namespace boost::beast::http; - -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) - : m_factory(std::move(factory)), - m_sf(std::move(sf)), - m_dv(dv), - m_mgr(mgr), - m_buffer_size(buffer_size) {} - -coro handler::handle(boost::asio::ip::tcp::socket s) { - auto ds = m_sf(); - auto peer = s.remote_endpoint(); - - forward_stream incoming(s, *ds); - auto& outgoing{*ds}; - - constexpr std::size_t buffer_size_to_load = 16_MiB; - - constexpr std::size_t buffer_size_to_relay_and_store = 32_MiB; - constexpr std::size_t buffer_size_to_relay = 4_KiB; - - flat_buffer buffer( - std::max(buffer_size_to_relay, buffer_size_to_relay_and_store)); - - for (;;) { - /* - * Note: lifetime of response must not exceed lifetime of request. - */ - std::string id = generate_unique_id(); - - ep::http::raw_request rawreq; - std::optional resp; - - try { - rawreq = co_await ep::http::raw_request::read(incoming, peer); - - auto& r = rawreq.headers; - LOG_INFO() << peer << ": incoming request: " << r.method_string() - << " " << r.target(); - - incoming.set_mode(forward_stream::forwarding); - std::unique_ptr req = - co_await m_factory->create(incoming, rawreq); - - if (get_object::can_handle(*req)) { - auto d_source = - m_mgr.get(cache::disk::object_metadata{req->object_key()}); - if (d_source) { - LOG_INFO() << peer << ": handling from cache"; - incoming.set_mode(forward_stream::deleting); - - auto& b = req->body(); - auto bs = b.buffer_size(); - - while (!(co_await b.read(bs)).empty()) { - co_await b.consume(); - } - - co_await b.consume(); - - LOG_INFO() << peer << ": done reading complete request"; - - response_parser parser; - response_serializer serializer{ - parser.get()}; - - co_await async_read_header(d_source, parser); - - const char* via_value = PROJECT_NAME " " PROJECT_VERSION; - parser.get().set(field::via, via_value); - - co_await async_write( - async_write_header(s, serializer, - boost::asio::use_awaitable), - s, *d_source); - - LOG_INFO() << peer << ": cache result served"; - continue; - } - } - - LOG_INFO() << peer << ": handling from downstream"; - co_await incoming.consume(); - - if (auto expect = req->header("expect"); - expect && *expect == "100-continue") { - LOG_INFO() << req->peer() << ": forwarding 100 CONTINUE"; - // TODO timeout - response_parser p; - response_serializer sr{p.get()}; - co_await async_read_header(outgoing, buffer, p); - co_await async_write_header(s, sr); - } - - // forwarding request body - auto& b = req->body(); - auto bs = b.buffer_size(); - - while (!(co_await b.read(bs)).empty()) { - co_await b.consume(); - } - - co_await b.consume(); - - // forwarding response - response_parser p; - p.body_limit(std::numeric_limits::max()); - response_serializer sr{p.get()}; - - LOG_INFO() << peer << ": reading header from downstream"; - co_await async_read_header(outgoing, buffer, p); - - if (r.method() == verb::head) { - LOG_INFO() << peer << ": HEAD request, skipping body relay"; - co_await async_write_header(s, sr); - - } else if (get_object::can_handle(*req)) { - auto d_sink = cache::disk::disk_sink{m_dv}; - auto s_sink = socket_sink{s}; - auto body_size = get_content_length(p.get()); - if (!body_size.has_value()) { - throw std::runtime_error("no content length"); - } - LOG_INFO() << peer << ": relaying and storing body of size " - << *body_size; - co_await async_read( - [&]() -> coro { - auto n = co_await async_write_header( - tee(s_sink, d_sink), sr); - d_sink.set_header_size(n); - }, - outgoing, buffer, *body_size, tee(s_sink, d_sink)); - co_await m_mgr.put( - cache::disk::object_metadata{req->object_key()}, d_sink); - - } else { - auto body_size = get_content_length(p.get()); - if (!body_size.has_value()) { - throw std::runtime_error("no content length"); - } - LOG_INFO() << peer << ": relaying body of size " << *body_size; - co_await async_read( - async_write_header(s, sr, boost::asio::use_awaitable), - outgoing, buffer, *body_size, socket_sink(s)); - } - LOG_INFO() << peer << ": done"; - - metric::increase(1); - } catch (const boost::system::system_error& e) { - throw; - } catch (const command_exception& e) { - resp = make_response(e); - } catch (const error_exception& e) { - resp = make_response(command_exception(*e.error())); - } catch (const std::exception& e) { - LOG_ERROR() << s.remote_endpoint() << ": " << e.what(); - resp = make_response(command_exception()); - } - - if (resp) { - co_await write(incoming, std::move(*resp), id); - } - } - - s.shutdown(boost::asio::ip::tcp::socket::shutdown_both); - s.close(); -} - -bool handler::intercept(ep::http::raw_request& r) const { return false; } - -coro handler::handle(ep::http::stream& s, ep::http::raw_request& r) { - co_return; -} - -} // namespace uh::cluster::proxy diff --git a/src/proxy/handler.h b/src/proxy/handler.h index b01476993..0cbd90985 100644 --- a/src/proxy/handler.h +++ b/src/proxy/handler.h @@ -1,29 +1,216 @@ #pragma once -#include -#include +#include +#include #include +#include +#include +#include +#include +#include + +#include +#include +#include + +#include +#include +#include namespace uh::cluster::proxy { namespace http = uh::cluster::ep::http; -class handler : public protocol_handler { +template class handler : public protocol_handler { public: explicit handler(std::unique_ptr factory, - std::function()> sf, - storage::data_view& dv, - cache::disk::manager& mgr, - std::size_t buffer_size); + std::function()> sf, + 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), + m_mgr(mgr), + m_buffer_size(buffer_size) {} + + coro handle(boost::asio::ip::tcp::socket s) override { + using boost::beast::flat_buffer; + using boost::beast::http::empty_body; + using boost::beast::http::field; + using boost::beast::http::fields; + using boost::beast::http::response_parser; + using boost::beast::http::response_serializer; + using boost::beast::http::verb; + auto ds = m_sf(); + auto peer = s.remote_endpoint(); + + auto incoming = forward_stream{s, *ds}; + auto& outgoing{*ds}; + + constexpr std::size_t buffer_size_to_load = 16_MiB; + + constexpr std::size_t buffer_size_to_relay_and_store = 32_MiB; + constexpr std::size_t buffer_size_to_relay = 4_KiB; + + flat_buffer buffer( + std::max(buffer_size_to_relay, buffer_size_to_relay_and_store)); + + for (;;) { + /* + * Note: lifetime of response must not exceed lifetime of request. + */ + std::string id = generate_unique_id(); + + ep::http::raw_request rawreq; + std::optional resp; + + try { + rawreq = co_await ep::http::raw_request::read(incoming, peer); + + auto& r = rawreq.headers; + LOG_INFO() << peer + << ": incoming request: " << r.method_string() << " " + << r.target(); + + incoming.set_mode(decltype(incoming)::forwarding); + std::unique_ptr req = + co_await m_factory->create(incoming, rawreq); + + if (get_object::can_handle(*req)) { + auto d_source = m_mgr.get( + cache::disk::object_metadata{req->object_key()}); + if (d_source) { + LOG_INFO() << peer << ": handling from cache"; + incoming.set_mode(decltype(incoming)::deleting); + + auto& b = req->body(); + auto bs = b.buffer_size(); + + while (!(co_await b.read(bs)).empty()) { + co_await b.consume(); + } + + co_await b.consume(); + + LOG_INFO() << peer << ": done reading complete request"; + + response_parser parser; + response_serializer serializer{ + parser.get()}; + + co_await async_read_header(d_source, parser); + + const char* via_value = + PROJECT_NAME " " PROJECT_VERSION; + parser.get().set(field::via, via_value); + + co_await async_write( + async_write_header(s, serializer, + boost::asio::use_awaitable), + s, *d_source); + + LOG_INFO() << peer << ": cache result served"; + continue; + } + } + + LOG_INFO() << peer << ": handling from downstream"; + co_await incoming.consume(); + + if (auto expect = req->header("expect"); + expect && *expect == "100-continue") { + LOG_INFO() << req->peer() << ": forwarding 100 CONTINUE"; + // TODO timeout + response_parser p; + response_serializer sr{p.get()}; + co_await async_read_header(outgoing, buffer, p); + co_await async_write_header(s, sr); + } + + // forwarding request body + auto& b = req->body(); + auto bs = b.buffer_size(); + + while (!(co_await b.read(bs)).empty()) { + co_await b.consume(); + } + + co_await b.consume(); + + // forwarding response + response_parser p; + p.body_limit(std::numeric_limits::max()); + response_serializer sr{p.get()}; + + LOG_INFO() << peer << ": reading header from downstream"; + co_await async_read_header(outgoing, buffer, p); + + if (r.method() == verb::head) { + LOG_INFO() << peer << ": HEAD request, skipping body relay"; + co_await async_write_header(s, sr); + + } else if (get_object::can_handle(*req)) { + auto d_sink = cache::disk::disk_sink{m_dv}; + auto s_sink = socket_sink{s}; + auto body_size = get_content_length(p.get()); + if (!body_size.has_value()) { + throw std::runtime_error("no content length"); + } + LOG_INFO() << peer << ": relaying and storing body of size " + << *body_size; + co_await async_read( + [&]() -> coro { + auto n = co_await async_write_header( + tee(s_sink, d_sink), sr); + d_sink.set_header_size(n); + }, + outgoing, buffer, *body_size, tee(s_sink, d_sink)); + co_await m_mgr.put( + cache::disk::object_metadata{req->object_key()}, + d_sink); + + } else { + auto body_size = get_content_length(p.get()); + if (!body_size.has_value()) { + throw std::runtime_error("no content length"); + } + LOG_INFO() + << peer << ": relaying body of size " << *body_size; + co_await async_read( + async_write_header(s, sr, boost::asio::use_awaitable), + outgoing, buffer, *body_size, socket_sink(s)); + } + LOG_INFO() << peer << ": done"; + + metric::increase(1); + } catch (const boost::system::system_error& e) { + throw; + } catch (const command_exception& e) { + resp = make_response(e); + } catch (const error_exception& e) { + resp = make_response(command_exception(*e.error())); + } catch (const std::exception& e) { + LOG_ERROR() << s.remote_endpoint() << ": " << e.what(); + resp = make_response(command_exception()); + } + + if (resp) { + co_await write(incoming, std::move(*resp), id); + } + } - coro handle(boost::asio::ip::tcp::socket s) override; + s.shutdown(boost::asio::ip::tcp::socket::shutdown_both); + s.close(); + } - bool intercept(ep::http::raw_request& r) const; - coro handle(ep::http::stream& s, ep::http::raw_request& r); + bool intercept(ep::http::raw_request& r) const { return false; } + coro handle(ep::http::stream& s, ep::http::raw_request& r) { + co_return; + } private: std::unique_ptr m_factory; - std::function()> m_sf; + std::function()> m_sf; storage::data_view& m_dv; cache::disk::manager& m_mgr; std::size_t m_buffer_size; diff --git a/src/proxy/service.cpp b/src/proxy/service.cpp index ea55d0191..3dfbb8ba5 100644 --- a/src/proxy/service.cpp +++ b/src/proxy/service.cpp @@ -2,43 +2,59 @@ #include -#include "request_factory.h" #include "handler.h" +#include "request_factory.h" #include +#include +#include +namespace net = boost::asio; +namespace ssl = net::ssl; namespace uh::cluster::proxy { using tcp = boost::asio::ip::tcp; -std::unique_ptr socket_factory( - boost::asio::io_context& ioc, - const std::string& server, - uint16_t port) { +std::unique_ptr> +socket_factory(boost::asio::io_context& ioc, const std::string& server, + uint16_t port) { + + ssl::context ctx(ssl::context::tls_client); + ctx.set_default_verify_paths(); + ctx.load_verify_file("/home/sungsik/Projects/core/cert.pem"); + + ctx.set_verify_mode(ssl::verify_peer); - auto addr = uh::cluster::resolve(server, port); - if (addr.empty()) { - throw std::runtime_error("lookup failed"); - } + tcp::resolver resolver(ioc); + beast::ssl_stream stream(ioc, ctx); - tcp::socket s(ioc); - boost::asio::connect(s, addr); + auto const results = resolver.resolve(server, std::to_string(port)); - return std::make_unique(std::move(s)); + beast::get_lowest_layer(stream).connect(results); + + stream.set_verify_callback(ssl::host_name_verification(server)); + + stream.handshake(ssl::stream_base::client); + + return std::make_unique(std::move(stream)); } -service::service(boost::asio::io_context& ioc, const service_config& sc, const config& c) +service::service(boost::asio::io_context& ioc, const service_config& sc, + const config& c) : m_ioc(ioc), m_etcd(sc.etcd_config), - m_dv(std::make_unique(ioc, m_etcd, c.gdv)), + m_dv(std::make_unique(ioc, m_etcd, + c.gdv)), m_mgr(cache::disk::manager::create(ioc, *m_dv, 10 * GIBI_BYTE)), - m_server(c.server, std::make_unique( - std::make_unique(), - [this, c]{ return socket_factory(m_ioc, c.downstream_host, c.downstream_port); }, - *m_dv, m_mgr, - c.buffer_size), - m_ioc) { -} + m_server(c.server, + std::make_unique>>( + std::make_unique(), + [this, c] { + return socket_factory(m_ioc, c.downstream_host, + c.downstream_port); + }, + *m_dv, m_mgr, c.buffer_size), + m_ioc) {} } // namespace uh::cluster::proxy From 7f47df4248e07a213292e1f55e56a7798dc10a88 Mon Sep 17 00:00:00 2001 From: Sungsik Nam Date: Thu, 25 Sep 2025 16:18:22 +0200 Subject: [PATCH 02/10] Support both: http and https --- scripts/start.sh.in | 3 +- src/common/utils/common.h | 5 +- src/config/configuration.cpp | 6 + src/proxy/config.h | 7 +- src/proxy/handler.h | 321 ++++++++++++++++++---------------- src/proxy/service.cpp | 53 ++++-- test/unit/test_trace_asio.cpp | 61 +++++++ 7 files changed, 281 insertions(+), 175 deletions(-) diff --git a/scripts/start.sh.in b/scripts/start.sh.in index 44410f0c4..d219265b8 100755 --- a/scripts/start.sh.in +++ b/scripts/start.sh.in @@ -187,7 +187,8 @@ $UH_CLUSTER entrypoint &> $log_dir/entrypoint.log & pid_entrypoint=$! export OTEL_RESOURCE_ATTRIBUTES="service.name=proxy" -$UH_CLUSTER --downstream-port 8081 --downstream-host localhost proxy &> $log_dir/proxy.log & +$UH_CLUSTER --downstream-port 8080 --downstream-host localhost --downstream-insecure proxy &> $log_dir/proxy.log & + pid_proxy=$! get_running_processes() { diff --git a/src/common/utils/common.h b/src/common/utils/common.h index a59a1d2a3..96b509eb3 100644 --- a/src/common/utils/common.h +++ b/src/common/utils/common.h @@ -89,9 +89,12 @@ constexpr const char* ENV_CFG_ETCD_PASSWORD = "UH_ETCD_PASSWORD"; constexpr const char* ENV_CFG_NO_DEDUPE = "UH_NO_DEDUPE"; constexpr const char* ENV_CFG_STORAGE_SERVICE_ID = "UH_STORAGE_INSTANCE_ID"; constexpr const char* ENV_CFG_STORAGE_GROUP_ID = "UH_STORAGE_GROUP_ID"; +constexpr const char* ENV_CFG_DOWNSTREAM_INSECURE = "UH_DOWNSTREAM_INSECURE"; +constexpr const char* ENV_CFG_DOWNSTREAM_CERT_FILE = "UH_DOWNSTREAM_CERT_FILE"; constexpr const char* ENV_CFG_DOWNSTREAM_HOST = "UH_DOWNSTREAM_HOST"; constexpr const char* ENV_CFG_DOWNSTREAM_PORT = "UH_DOWNSTREAM_PORT"; -constexpr const char* ENV_CFG_DOWNSTREAM_CONNECTIONS = "UH_DOWNSTREAM_CONNECTIONS"; +constexpr const char* ENV_CFG_DOWNSTREAM_CONNECTIONS = + "UH_DOWNSTREAM_CONNECTIONS"; constexpr const char* RESERVED_BUCKET_NAME = "ultihash"; diff --git a/src/config/configuration.cpp b/src/config/configuration.cpp index 7e9d67db9..8caaf9d23 100644 --- a/src/config/configuration.cpp +++ b/src/config/configuration.cpp @@ -247,6 +247,12 @@ CLI::App* sub_coordinator(CLI::App& app, coordinator_config& cfg) { CLI::App* sub_proxy(CLI::App& app, proxy::config& cfg) { auto* rv = app.add_subcommand("proxy", "S3 proxy server"); + app.add_flag("--downstream-insecure", cfg.downstream_insecure, + "downstream uses http, instead of https") + ->envname(ENV_CFG_DOWNSTREAM_INSECURE); + app.add_option("--downstream-cert-file", cfg.downstream_cert_file, + "downstream certification file path") + ->envname(ENV_CFG_DOWNSTREAM_CERT_FILE); app.add_option("--downstream-host", cfg.downstream_host, "downstream host") ->envname(ENV_CFG_DOWNSTREAM_HOST); app.add_option("--downstream-port", cfg.downstream_port, "downstream port") diff --git a/src/proxy/config.h b/src/proxy/config.h index abad58192..12883ac0c 100644 --- a/src/proxy/config.h +++ b/src/proxy/config.h @@ -6,11 +6,10 @@ namespace uh::cluster::proxy { struct config { - server_config server = { - .port = 8088, - .bind_address = "0.0.0.0" - }; + server_config server = {.port = 8088, .bind_address = "0.0.0.0"}; + bool downstream_insecure; + std::optional downstream_cert_file; std::string downstream_host; uint16_t downstream_port; std::size_t connections = 16; diff --git a/src/proxy/handler.h b/src/proxy/handler.h index 0cbd90985..39c649cf5 100644 --- a/src/proxy/handler.h +++ b/src/proxy/handler.h @@ -17,14 +17,24 @@ #include #include +#include +#include + namespace uh::cluster::proxy { namespace http = uh::cluster::ep::http; -template class handler : public protocol_handler { +class handler : public protocol_handler { + template + coro _handle(boost::asio::ip::tcp::socket s, StreamType& ds); + public: + using variant_stream = + std::variant>; + explicit handler(std::unique_ptr factory, - std::function()> sf, + std::function()> sf, storage::data_view& dv, cache::disk::manager& mgr, std::size_t buffer_size) : m_factory(std::move(factory)), @@ -33,192 +43,199 @@ template class handler : public protocol_handler { m_mgr(mgr), m_buffer_size(buffer_size) {} - coro handle(boost::asio::ip::tcp::socket s) override { - using boost::beast::flat_buffer; - using boost::beast::http::empty_body; - using boost::beast::http::field; - using boost::beast::http::fields; - using boost::beast::http::response_parser; - using boost::beast::http::response_serializer; - using boost::beast::http::verb; - auto ds = m_sf(); - auto peer = s.remote_endpoint(); + coro handle(boost::asio::ip::tcp::socket s) override; - auto incoming = forward_stream{s, *ds}; - auto& outgoing{*ds}; + bool intercept(ep::http::raw_request& r) const { return false; } + coro handle(ep::http::stream& s, ep::http::raw_request& r) { + co_return; + } - constexpr std::size_t buffer_size_to_load = 16_MiB; +private: + std::unique_ptr m_factory; + std::function()> m_sf; + storage::data_view& m_dv; + cache::disk::manager& m_mgr; + std::size_t m_buffer_size; - constexpr std::size_t buffer_size_to_relay_and_store = 32_MiB; - constexpr std::size_t buffer_size_to_relay = 4_KiB; + coro handle_request(boost::asio::ip::tcp::socket& s, + http::raw_request& rawreq, + const std::string& id, + boost::beast::tcp_stream& ds); +}; - flat_buffer buffer( - std::max(buffer_size_to_relay, buffer_size_to_relay_and_store)); +template +coro handler::_handle(boost::asio::ip::tcp::socket s, StreamType& ds) { + using boost::beast::flat_buffer; + using boost::beast::http::empty_body; + using boost::beast::http::field; + using boost::beast::http::fields; + using boost::beast::http::response_parser; + using boost::beast::http::response_serializer; + using boost::beast::http::verb; + auto peer = s.remote_endpoint(); - for (;;) { - /* - * Note: lifetime of response must not exceed lifetime of request. - */ - std::string id = generate_unique_id(); + auto incoming = forward_stream{s, ds}; + auto& outgoing{ds}; - ep::http::raw_request rawreq; - std::optional resp; + constexpr std::size_t buffer_size_to_load = 16_MiB; - try { - rawreq = co_await ep::http::raw_request::read(incoming, peer); + constexpr std::size_t buffer_size_to_relay_and_store = 32_MiB; + constexpr std::size_t buffer_size_to_relay = 4_KiB; - auto& r = rawreq.headers; - LOG_INFO() << peer - << ": incoming request: " << r.method_string() << " " - << r.target(); + flat_buffer buffer( + std::max(buffer_size_to_relay, buffer_size_to_relay_and_store)); - incoming.set_mode(decltype(incoming)::forwarding); - std::unique_ptr req = - co_await m_factory->create(incoming, rawreq); + for (;;) { + /* + * Note: lifetime of response must not exceed lifetime of request. + */ + std::string id = generate_unique_id(); - if (get_object::can_handle(*req)) { - auto d_source = m_mgr.get( - cache::disk::object_metadata{req->object_key()}); - if (d_source) { - LOG_INFO() << peer << ": handling from cache"; - incoming.set_mode(decltype(incoming)::deleting); + ep::http::raw_request rawreq; + std::optional resp; - auto& b = req->body(); - auto bs = b.buffer_size(); + try { + rawreq = co_await ep::http::raw_request::read(incoming, peer); - while (!(co_await b.read(bs)).empty()) { - co_await b.consume(); - } + auto& r = rawreq.headers; + LOG_INFO() << peer << ": incoming request: " << r.method_string() + << " " << r.target(); - co_await b.consume(); + incoming.set_mode(decltype(incoming)::forwarding); + std::unique_ptr req = + co_await m_factory->create(incoming, rawreq); - LOG_INFO() << peer << ": done reading complete request"; + if (get_object::can_handle(*req)) { + auto d_source = + m_mgr.get(cache::disk::object_metadata{req->object_key()}); + if (d_source) { + LOG_INFO() << peer << ": handling from cache"; + incoming.set_mode(decltype(incoming)::deleting); - response_parser parser; - response_serializer serializer{ - parser.get()}; + auto& b = req->body(); + auto bs = b.buffer_size(); - co_await async_read_header(d_source, parser); + while (!(co_await b.read(bs)).empty()) { + co_await b.consume(); + } - const char* via_value = - PROJECT_NAME " " PROJECT_VERSION; - parser.get().set(field::via, via_value); + co_await b.consume(); - co_await async_write( - async_write_header(s, serializer, - boost::asio::use_awaitable), - s, *d_source); + LOG_INFO() << peer << ": done reading complete request"; - LOG_INFO() << peer << ": cache result served"; - continue; - } - } + response_parser parser; + response_serializer serializer{ + parser.get()}; - LOG_INFO() << peer << ": handling from downstream"; - co_await incoming.consume(); - - if (auto expect = req->header("expect"); - expect && *expect == "100-continue") { - LOG_INFO() << req->peer() << ": forwarding 100 CONTINUE"; - // TODO timeout - response_parser p; - response_serializer sr{p.get()}; - co_await async_read_header(outgoing, buffer, p); - co_await async_write_header(s, sr); - } + co_await async_read_header(d_source, parser); - // forwarding request body - auto& b = req->body(); - auto bs = b.buffer_size(); + const char* via_value = PROJECT_NAME " " PROJECT_VERSION; + parser.get().set(field::via, via_value); - while (!(co_await b.read(bs)).empty()) { - co_await b.consume(); + co_await async_write( + async_write_header(s, serializer, + boost::asio::use_awaitable), + s, *d_source); + + LOG_INFO() << peer << ": cache result served"; + continue; } + } - co_await b.consume(); + LOG_INFO() << peer << ": handling from downstream"; + co_await incoming.consume(); - // forwarding response + if (auto expect = req->header("expect"); + expect && *expect == "100-continue") { + LOG_INFO() << req->peer() << ": forwarding 100 CONTINUE"; + // TODO timeout response_parser p; - p.body_limit(std::numeric_limits::max()); response_serializer sr{p.get()}; - - LOG_INFO() << peer << ": reading header from downstream"; co_await async_read_header(outgoing, buffer, p); + co_await async_write_header(s, sr); + } - if (r.method() == verb::head) { - LOG_INFO() << peer << ": HEAD request, skipping body relay"; - co_await async_write_header(s, sr); + // forwarding request body + auto& b = req->body(); + auto bs = b.buffer_size(); - } else if (get_object::can_handle(*req)) { - auto d_sink = cache::disk::disk_sink{m_dv}; - auto s_sink = socket_sink{s}; - auto body_size = get_content_length(p.get()); - if (!body_size.has_value()) { - throw std::runtime_error("no content length"); - } - LOG_INFO() << peer << ": relaying and storing body of size " - << *body_size; - co_await async_read( - [&]() -> coro { - auto n = co_await async_write_header( - tee(s_sink, d_sink), sr); - d_sink.set_header_size(n); - }, - outgoing, buffer, *body_size, tee(s_sink, d_sink)); - co_await m_mgr.put( - cache::disk::object_metadata{req->object_key()}, - d_sink); - - } else { - auto body_size = get_content_length(p.get()); - if (!body_size.has_value()) { - throw std::runtime_error("no content length"); - } - LOG_INFO() - << peer << ": relaying body of size " << *body_size; - co_await async_read( - async_write_header(s, sr, boost::asio::use_awaitable), - outgoing, buffer, *body_size, socket_sink(s)); - } - LOG_INFO() << peer << ": done"; - - metric::increase(1); - } catch (const boost::system::system_error& e) { - throw; - } catch (const command_exception& e) { - resp = make_response(e); - } catch (const error_exception& e) { - resp = make_response(command_exception(*e.error())); - } catch (const std::exception& e) { - LOG_ERROR() << s.remote_endpoint() << ": " << e.what(); - resp = make_response(command_exception()); + while (!(co_await b.read(bs)).empty()) { + co_await b.consume(); } - if (resp) { - co_await write(incoming, std::move(*resp), id); + co_await b.consume(); + + // forwarding response + response_parser p; + p.body_limit(std::numeric_limits::max()); + response_serializer sr{p.get()}; + + LOG_INFO() << peer << ": reading header from downstream"; + co_await async_read_header(outgoing, buffer, p); + + if (r.method() == verb::head) { + LOG_INFO() << peer << ": HEAD request, skipping body relay"; + co_await async_write_header(s, sr); + + } else if (get_object::can_handle(*req)) { + auto d_sink = cache::disk::disk_sink{m_dv}; + auto s_sink = socket_sink{s}; + auto body_size = get_content_length(p.get()); + if (!body_size.has_value()) { + throw std::runtime_error("no content length"); + } + LOG_INFO() << peer << ": relaying and storing body of size " + << *body_size; + co_await async_read( + [&]() -> coro { + auto n = co_await async_write_header( + tee(s_sink, d_sink), sr); + d_sink.set_header_size(n); + }, + outgoing, buffer, *body_size, tee(s_sink, d_sink)); + co_await m_mgr.put( + cache::disk::object_metadata{req->object_key()}, d_sink); + + } else { + auto body_size = get_content_length(p.get()); + if (!body_size.has_value()) { + throw std::runtime_error("no content length"); + } + LOG_INFO() << peer << ": relaying body of size " << *body_size; + co_await async_read( + async_write_header(s, sr, boost::asio::use_awaitable), + outgoing, buffer, *body_size, socket_sink(s)); } + LOG_INFO() << peer << ": done"; + + metric::increase(1); + } catch (const boost::system::system_error& e) { + throw; + } catch (const command_exception& e) { + resp = make_response(e); + } catch (const error_exception& e) { + resp = make_response(command_exception(*e.error())); + } catch (const std::exception& e) { + LOG_ERROR() << s.remote_endpoint() << ": " << e.what(); + resp = make_response(command_exception()); } - s.shutdown(boost::asio::ip::tcp::socket::shutdown_both); - s.close(); - } - - bool intercept(ep::http::raw_request& r) const { return false; } - coro handle(ep::http::stream& s, ep::http::raw_request& r) { - co_return; + if (resp) { + co_await write(incoming, std::move(*resp), id); + } } -private: - std::unique_ptr m_factory; - std::function()> m_sf; - storage::data_view& m_dv; - cache::disk::manager& m_mgr; - std::size_t m_buffer_size; - - coro handle_request(boost::asio::ip::tcp::socket& s, - http::raw_request& rawreq, - const std::string& id, - boost::beast::tcp_stream& ds); -}; + s.shutdown(boost::asio::ip::tcp::socket::shutdown_both); + s.close(); +} + +coro handler::handle(boost::asio::ip::tcp::socket s) { + auto downstream = m_sf(); + co_await std::visit( + [this, &s](auto& ds) -> coro { + co_await _handle(std::move(s), ds); + }, + *downstream); +} } // namespace uh::cluster::proxy diff --git a/src/proxy/service.cpp b/src/proxy/service.cpp index 3dfbb8ba5..530a80b64 100644 --- a/src/proxy/service.cpp +++ b/src/proxy/service.cpp @@ -7,7 +7,6 @@ #include #include -#include namespace net = boost::asio; namespace ssl = net::ssl; @@ -16,28 +15,47 @@ namespace uh::cluster::proxy { using tcp = boost::asio::ip::tcp; -std::unique_ptr> +std::unique_ptr socket_factory(boost::asio::io_context& ioc, const std::string& server, - uint16_t port) { + uint16_t port, bool insecure, + std::optional cert_file) { + if (insecure) { + LOG_INFO() << "Creating insecure connection to " << server << ":" + << port; + auto addr = uh::cluster::resolve(server, port); + if (addr.empty()) { + throw std::runtime_error("lookup failed"); + } - ssl::context ctx(ssl::context::tls_client); - ctx.set_default_verify_paths(); - ctx.load_verify_file("/home/sungsik/Projects/core/cert.pem"); + tcp::socket s(ioc); + boost::asio::connect(s, addr); - ctx.set_verify_mode(ssl::verify_peer); + return std::make_unique(std::move(s)); - tcp::resolver resolver(ioc); - beast::ssl_stream stream(ioc, ctx); + } else { + LOG_INFO() << "Creating secure connection to " << server << ":" << port; + ssl::context ctx(ssl::context::tls_client); + ctx.set_default_verify_paths(); + if (cert_file.has_value()) { + LOG_INFO() << "Loading cert file " << *cert_file; + ctx.load_verify_file(*cert_file); + } - auto const results = resolver.resolve(server, std::to_string(port)); + ctx.set_verify_mode(ssl::verify_peer); - beast::get_lowest_layer(stream).connect(results); + tcp::resolver resolver(ioc); + beast::ssl_stream stream(ioc, ctx); - stream.set_verify_callback(ssl::host_name_verification(server)); + auto const results = resolver.resolve(server, std::to_string(port)); - stream.handshake(ssl::stream_base::client); + beast::get_lowest_layer(stream).connect(results); - return std::make_unique(std::move(stream)); + stream.set_verify_callback(ssl::host_name_verification(server)); + + stream.handshake(ssl::stream_base::client); + + return std::make_unique(std::move(stream)); + } } service::service(boost::asio::io_context& ioc, const service_config& sc, @@ -48,11 +66,12 @@ service::service(boost::asio::io_context& ioc, const service_config& sc, c.gdv)), m_mgr(cache::disk::manager::create(ioc, *m_dv, 10 * GIBI_BYTE)), m_server(c.server, - std::make_unique>>( + std::make_unique( std::make_unique(), [this, c] { - return socket_factory(m_ioc, c.downstream_host, - c.downstream_port); + return socket_factory( + m_ioc, c.downstream_host, c.downstream_port, + c.downstream_insecure, c.downstream_cert_file); }, *m_dv, m_mgr, c.buffer_size), m_ioc) {} diff --git a/test/unit/test_trace_asio.cpp b/test/unit/test_trace_asio.cpp index 721b16e59..0caf7886e 100644 --- a/test/unit/test_trace_asio.cpp +++ b/test/unit/test_trace_asio.cpp @@ -110,3 +110,64 @@ BOOST_AUTO_TEST_CASE(propagates_context_through_continue) { } BOOST_AUTO_TEST_SUITE_END() + +#include + +using namespace boost::asio; +using boost::asio::experimental::make_parallel_group; +using boost::asio::experimental::wait_for_one_error; + +BOOST_FIXTURE_TEST_SUITE(awaitables, fixture) +awaitable async_timer(io_context& ctx, int ms, std::string& trace) { + steady_timer timer(ctx, std::chrono::milliseconds(ms)); + co_await timer.async_wait(use_awaitable); + trace += "done"; + co_return; +} + +BOOST_AUTO_TEST_CASE(coroutine_parallel_group_deferred) { + + boost::asio::posix::stream_descriptor out(ioc, ::dup(STDOUT_FILENO)); + + using op_type = decltype(out.async_write_some(boost::asio::buffer("", 0))); + + std::cout << boost::core::demangle(typeid(op_type).name()) << std::endl; + // boost::asio::deferred_async_operation::initiate_async_write_some, + // boost::asio::const_buffer const&> + + std::string trace1; + // std::string trace2; + // + auto op1 = co_spawn(ioc, async_timer(ioc, 10, trace1), deferred); + using op_type_2 = decltype(op1); + + std::cout << boost::core::demangle(typeid(op_type_2).name()) << std::endl; + // boost::asio::deferred_async_operation, + // boost::asio::detail::awaitable_as_function > + + // auto op2 = co_spawn(ioc, async_timer(ioc, 20, trace2), deferred); + // + // std::vector ops; + // ops.push_back(std::move(op1)); + // ops.push_back(std::move(op2)); + // + // co_spawn(ioc, + // [&]() -> awaitable { + // auto completion_order = + // co_await make_parallel_group(ops) + // .async_wait(wait_for_one_error(), deferred); + // + // BOOST_CHECK_EQUAL(trace1, "done"); + // BOOST_CHECK_EQUAL(trace2, "done"); + // co_return; + // }(), + // use_future + // ).get(); +} + +BOOST_AUTO_TEST_SUITE_END() From fca551b27dfc1947f54cd39aeb26aa24e788cda4 Mon Sep 17 00:00:00 2001 From: Sungsik Nam Date: Thu, 25 Sep 2025 16:24:39 +0200 Subject: [PATCH 03/10] nit --- scripts/start.sh.in | 1 - test/unit/test_trace_asio.cpp | 61 ----------------------------------- 2 files changed, 62 deletions(-) diff --git a/scripts/start.sh.in b/scripts/start.sh.in index d219265b8..e2d82d099 100755 --- a/scripts/start.sh.in +++ b/scripts/start.sh.in @@ -188,7 +188,6 @@ pid_entrypoint=$! export OTEL_RESOURCE_ATTRIBUTES="service.name=proxy" $UH_CLUSTER --downstream-port 8080 --downstream-host localhost --downstream-insecure proxy &> $log_dir/proxy.log & - pid_proxy=$! get_running_processes() { diff --git a/test/unit/test_trace_asio.cpp b/test/unit/test_trace_asio.cpp index 0caf7886e..721b16e59 100644 --- a/test/unit/test_trace_asio.cpp +++ b/test/unit/test_trace_asio.cpp @@ -110,64 +110,3 @@ BOOST_AUTO_TEST_CASE(propagates_context_through_continue) { } BOOST_AUTO_TEST_SUITE_END() - -#include - -using namespace boost::asio; -using boost::asio::experimental::make_parallel_group; -using boost::asio::experimental::wait_for_one_error; - -BOOST_FIXTURE_TEST_SUITE(awaitables, fixture) -awaitable async_timer(io_context& ctx, int ms, std::string& trace) { - steady_timer timer(ctx, std::chrono::milliseconds(ms)); - co_await timer.async_wait(use_awaitable); - trace += "done"; - co_return; -} - -BOOST_AUTO_TEST_CASE(coroutine_parallel_group_deferred) { - - boost::asio::posix::stream_descriptor out(ioc, ::dup(STDOUT_FILENO)); - - using op_type = decltype(out.async_write_some(boost::asio::buffer("", 0))); - - std::cout << boost::core::demangle(typeid(op_type).name()) << std::endl; - // boost::asio::deferred_async_operation::initiate_async_write_some, - // boost::asio::const_buffer const&> - - std::string trace1; - // std::string trace2; - // - auto op1 = co_spawn(ioc, async_timer(ioc, 10, trace1), deferred); - using op_type_2 = decltype(op1); - - std::cout << boost::core::demangle(typeid(op_type_2).name()) << std::endl; - // boost::asio::deferred_async_operation, - // boost::asio::detail::awaitable_as_function > - - // auto op2 = co_spawn(ioc, async_timer(ioc, 20, trace2), deferred); - // - // std::vector ops; - // ops.push_back(std::move(op1)); - // ops.push_back(std::move(op2)); - // - // co_spawn(ioc, - // [&]() -> awaitable { - // auto completion_order = - // co_await make_parallel_group(ops) - // .async_wait(wait_for_one_error(), deferred); - // - // BOOST_CHECK_EQUAL(trace1, "done"); - // BOOST_CHECK_EQUAL(trace2, "done"); - // co_return; - // }(), - // use_future - // ).get(); -} - -BOOST_AUTO_TEST_SUITE_END() From ed2fc5dcf3fae06917afdca711283280aa7a0caf Mon Sep 17 00:00:00 2001 From: Sungsik Nam Date: Fri, 26 Sep 2025 10:55:57 +0200 Subject: [PATCH 04/10] Support GCC --- src/proxy/handler.h | 30 +++++++++++++++++++----------- 1 file changed, 19 insertions(+), 11 deletions(-) diff --git a/src/proxy/handler.h b/src/proxy/handler.h index 39c649cf5..dd950f3c7 100644 --- a/src/proxy/handler.h +++ b/src/proxy/handler.h @@ -57,6 +57,7 @@ class handler : public protocol_handler { cache::disk::manager& m_mgr; std::size_t m_buffer_size; + friend struct handle_visitor; coro handle_request(boost::asio::ip::tcp::socket& s, http::raw_request& rawreq, const std::string& id, @@ -115,12 +116,13 @@ coro handler::_handle(boost::asio::ip::tcp::socket s, StreamType& ds) { auto& b = req->body(); auto bs = b.buffer_size(); - while (!(co_await b.read(bs)).empty()) { + while (true) { + auto result = co_await b.read(bs); co_await b.consume(); + if (result.empty()) + break; } - co_await b.consume(); - LOG_INFO() << peer << ": done reading complete request"; response_parser parser; @@ -159,12 +161,13 @@ coro handler::_handle(boost::asio::ip::tcp::socket s, StreamType& ds) { auto& b = req->body(); auto bs = b.buffer_size(); - while (!(co_await b.read(bs)).empty()) { + while (true) { + auto result = co_await b.read(bs); co_await b.consume(); + if (result.empty()) + break; } - co_await b.consume(); - // forwarding response response_parser p; p.body_limit(std::numeric_limits::max()); @@ -229,13 +232,18 @@ coro handler::_handle(boost::asio::ip::tcp::socket s, StreamType& ds) { s.close(); } +struct handle_visitor { + handler* h; + boost::asio::ip::tcp::socket s; + + template coro operator()(Downstream& ds) { + co_await h->_handle(std::move(s), ds); + } +}; + coro handler::handle(boost::asio::ip::tcp::socket s) { auto downstream = m_sf(); - co_await std::visit( - [this, &s](auto& ds) -> coro { - co_await _handle(std::move(s), ds); - }, - *downstream); + co_await std::visit(handle_visitor{this, std::move(s)}, *downstream); } } // namespace uh::cluster::proxy From 01a35cb007efecc5f9ee6bb0124fec2a7bfe3b81 Mon Sep 17 00:00:00 2001 From: Sungsik Nam Date: Thu, 25 Sep 2025 17:56:58 +0200 Subject: [PATCH 05/10] wrapper prepared --- test/unit/test_trace_asio.cpp | 71 +++++++++++++++++++++++++++++++++++ 1 file changed, 71 insertions(+) diff --git a/test/unit/test_trace_asio.cpp b/test/unit/test_trace_asio.cpp index 721b16e59..c16d421ed 100644 --- a/test/unit/test_trace_asio.cpp +++ b/test/unit/test_trace_asio.cpp @@ -110,3 +110,74 @@ BOOST_AUTO_TEST_CASE(propagates_context_through_continue) { } BOOST_AUTO_TEST_SUITE_END() + +#include +using namespace boost::asio; +using namespace boost::asio::experimental; + +template +auto wrap_coro_as_deferred(io_context& ioc, Func&& func, + CompletionToken&& token) { + return boost::asio::async_initiate( + boost::asio::co_composed( + [func = std::forward(func)](auto state, + io_context& ioc) -> void { + try { + state.throw_if_cancelled(true); + state.reset_cancellation_state( + boost::asio::enable_terminal_cancellation()); + // Call the provided function and await its result + co_await func(); + } catch (const boost::system::system_error& e) { + co_return {e.code()}; + } + }, + ioc), + token, std::ref(ioc)); +} + +BOOST_FIXTURE_TEST_SUITE(a_traced_coro, fixture) +BOOST_AUTO_TEST_CASE(ranged_parallel_group_streambuf) { + + boost::asio::posix::stream_descriptor out(ioc, ::dup(STDOUT_FILENO)); + boost::asio::posix::stream_descriptor err(ioc, ::dup(STDERR_FILENO)); + + using op_type = decltype(out.async_write_some(boost::asio::buffer("", 0))); + std::cout << boost::core::demangle(typeid(op_type).name()) << "\n"; + auto tmp = wrap_coro_as_deferred( + ioc, + []() -> coro { + // Your coroutine code here + co_return; + }, + deferred); + using op_type_2 = + decltype(out.async_write_some(boost::asio::buffer("", 0))); + std::cout << boost::core::demangle(typeid(op_type_2).name()) << "\n"; + + // std::vector ops; + // + // ops.push_back(out.async_write_some(boost::asio::buffer("first\r\n", 7), + // use_awaitable)); + // + // ops.push_back(err.async_write_some(boost::asio::buffer("second\r\n", 8), + // use_awaitable)); + // + // co_spawn( + // ioc, + // [&]() -> coro { + // auto [completion_order, n] = + // co_await make_parallel_group(ops).async_wait( + // wait_for_one_error(), deferred); + // + // for (std::size_t i = 0; i < completion_order.size(); ++i) { + // std::size_t idx = completion_order[i]; + // std::cout << "operation " << idx << " finished: "; + // std::cout << n[idx] << "\n"; + // } + // }, + // use_future) + // .get(); +} +BOOST_AUTO_TEST_SUITE_END() From 2259536c0389810e8b102af8b745067dd2600883 Mon Sep 17 00:00:00 2001 From: Sungsik Nam Date: Fri, 26 Sep 2025 11:06:16 +0200 Subject: [PATCH 06/10] Refactoring --- src/common/types/common_types.h | 6 +- src/proxy/cache/disk/disk_io.h | 5 - src/proxy/cache/disk/handler_example.cpp | 171 ----------------------- src/proxy/cache/disk/manager.h | 11 +- src/proxy/cache/disk/utils.h | 30 ---- src/proxy/handler.h | 11 +- src/proxy/http.h | 6 +- test/unit/test_disk_cache_manager.cpp | 17 ++- test/unit/test_trace_asio.cpp | 71 ---------- 9 files changed, 30 insertions(+), 298 deletions(-) delete mode 100644 src/proxy/cache/disk/handler_example.cpp delete mode 100644 src/proxy/cache/disk/utils.h diff --git a/src/common/types/common_types.h b/src/common/types/common_types.h index 46432833b..ee77f8f3b 100644 --- a/src/common/types/common_types.h +++ b/src/common/types/common_types.h @@ -40,13 +40,15 @@ template using coro = boost::asio::traced_awaitable; inline coro async_noop() { co_return; }; template struct is_boost_awaitable : std::false_type {}; - template struct is_boost_awaitable> : std::true_type {}; - template constexpr bool is_boost_awaitable_v = is_boost_awaitable::value; +template struct is_coro : std::false_type {}; +template struct is_coro> : std::true_type {}; +template inline constexpr bool is_coro_v = is_coro::value; + template requires is_boost_awaitable_v> inline coro async_wrap(Awaitable&& v) { diff --git a/src/proxy/cache/disk/disk_io.h b/src/proxy/cache/disk/disk_io.h index e03eb9796..168b40d7f 100644 --- a/src/proxy/cache/disk/disk_io.h +++ b/src/proxy/cache/disk/disk_io.h @@ -4,13 +4,8 @@ #pragma once #include -#include -#include #include -#include -#include -#include #include namespace uh::cluster::proxy::cache::disk { diff --git a/src/proxy/cache/disk/handler_example.cpp b/src/proxy/cache/disk/handler_example.cpp deleted file mode 100644 index 235582064..000000000 --- a/src/proxy/cache/disk/handler_example.cpp +++ /dev/null @@ -1,171 +0,0 @@ -#include -#include - -#include - -#include -#include -#include - -namespace beast = boost::beast; -namespace http = beast::http; -namespace asio = boost::asio; - -size_t get_content_length(const http::response& resp) { - auto len_str = resp[http::field::content_length]; // std::string_view 타입 - if (!len_str.empty()) { - return std::stoull(std::string(len_str)); - } -} - -/* - * Read object and create lazy body, send body to incomming - */ -awaitable read_storage(asio::ip::tcp::socket& incomming, - const http::request& req, - const cache_entry& entry, data_view storage) { - http::response res; - res.result(http::status::ok); - res.version(req.version()); - res.set(http::field::content_type, "application/octet-stream"); - res.set(http::field::etag, entry->tag); - res.keep_alive(req.keep_alive()); - res.content_length(entry->size); - - co_await http::async_write_header(incomming, res); - - auto reader = create_reader(storage, entry->addr, 8192); - while (true) { - auto chunk = co_await reader.read(); - if (chunk.empty()) - break; - co_await asio::async_write(incomming, - asio::buffer(chunk.data(), chunk.size())); - } -} - -/* - * Get lazy body from outgoing, save body to storage while fowarding them to - * incomming - */ -awaitable
save_object(asio::ip::tcp::socket& outgoing, - asio::ip::tcp::socket& incomming, - data_view storage) { - std::vector
addresses; - auto writer = create_writer(storage); - while (true) { - co_await http::async_read_some(outgoing, resp_buffer); - if (resp_buffer.empty()) - break; - co_await asio::async_write(incomming, - asio::buffer(chunk.data(), chunk.size())); - auto addr = co_await writer.write(chunk); - addresses.push_back(addr); - } - auto final_addr = concat(addresses); -} - -awaitable -send_conditional_request(asio::ip::tcp::socket& outgoing, - const http::request& req, - const std::string& etag) { - http::request cond_req = req; - cond_req.set(http::field::if_none_match, *etag); - co_await http::async_write(outgoing, cond_req); -} - -coro -reading_response_header(asio::ip::tcp::socket& outgoing, - const http::request& req) { - beast::http::response_parser parser; - parser.body_limit((std::numeric_limits::max)()); - - auto buffer = co_await outgoing.read_until("\r\n\r\n"); - // std::string txt(buffer.data(), buffer.size()); - // boost::replace_all(txt, "\r", "\\r"); - // boost::replace_all(txt, "\n", "\\n"); - - beast::error_code ec; - parser.put(boost::asio::buffer(buffer), ec); - - auto res = parser.release(); - - bs = outgoing.buffer_size(); - std::size_t read = 0ull; - std::size_t len = std::stoul(res.at("Content-Length")); - if (req.method() == boost::beast::http::verb::head && - (res.result_int() / 100 == 2)) { - len = 0; - } - - LOG_INFO() << peer << ": sending response " << res.result_int() << " " - << res.reason() << " -- " << len; - co_return len; -} - -awaitable handler(asio::ip::tcp::socket& incoming, - asio::ip::tcp::socket& outgoing, cache& c, - data_view storage) { - - auto peer = s.remote_endpoint(); - - constexpr size_t storage_size = 1024 * 1024 * 1024; // 1GB - storage_cache sc(c, storage, storage_size); - lock_map locks; - - while (true) { - auto rawreq = co_await raw_request::read(incoming, peer); - auto req = co_await m_factory->create(incoming, rawreq); - - if (req.method() == http::verb::get) { - auto key = object_metadata{req.path(), req.query("versionId")}; - auto pobj = sc.get(key); - if (pobj != nullptr) { - incoming.set_mode(forward_stream::deleting); - outgoing.set_mode(forward_stream::deleting); - write(incoming, get_response(pobj->get_body()), session_id); - co_await incoming.consume(); - co_await outgoing.consume(); - } else { - pobj = co_await locks.acquire(); - if (pobj == nullptr) { - incoming.set_mode(forward_stream::forwarding); - outgoing.set_mode(forward_stream::forwarding); - - // forwarding request - co_await read(req->body()); - - // forwarding response - auto length = reading_response_header(outgoing, req); - auto body = - storage_writer_body{m_dedupe, downstream, length}; - pobj = object_handle::create(body); - co_await outgoing.consume(); - co_await sc.put(key, pobj); - } else { - incoming.set_mode(forward_stream::deleting); - outgoing.set_mode(forward_stream::deleting); - co_await incoming.consume(); - co_await outgoing.consume(); - } - write(incoming, get_response(pobj->get_body()), session_id); - } - - // if (pobj != nullptr) { - // if (pobj->is_fresh()) { - // write(incoming, get_response(pobj->get_body()), - // session_id); - // } else { - // co_await send_conditional_request(outgoing, req, - // entry.tag); - // - // auto resp = co_await outgoing.read_header(); - // if (resp.result() == http::status::not_modified) { - // pobj->revive(); - // write(incoming, get_response(pobj->get_body()), - // session_id); - } else { - // Not a GET request - } - } -} diff --git a/src/proxy/cache/disk/manager.h b/src/proxy/cache/disk/manager.h index cc8f5d688..ca3d0d701 100644 --- a/src/proxy/cache/disk/manager.h +++ b/src/proxy/cache/disk/manager.h @@ -1,12 +1,14 @@ #pragma once #include -#include +#include #include #include #include +#include +#include #include #include @@ -27,8 +29,7 @@ class manager { using stream = ep::http::stream; using body = ep::http::body; - coro put(object_metadata key, disk_sink& w) { - auto objh = w.get_object_handle(); + coro put(object_metadata key, object_handle objh) { auto obj_size = objh.data_size(); auto total_size = @@ -61,12 +62,12 @@ class manager { std::cout << "Total size after put: " << m_current_size << std::endl; } - std::unique_ptr get(object_metadata key) { + std::shared_ptr get(object_metadata key) { auto entry = m_cache->get(key); if (!entry) { return nullptr; } - return std::make_unique(m_storage, std::move(entry)); + return entry; } static manager create(boost::asio::io_context& ioc, data_view& storage, diff --git a/src/proxy/cache/disk/utils.h b/src/proxy/cache/disk/utils.h deleted file mode 100644 index 611ddc8b6..000000000 --- a/src/proxy/cache/disk/utils.h +++ /dev/null @@ -1,30 +0,0 @@ -/* - * Disk utilities: APIs for common disk operations - */ -#pragma once - -#include -#include -#include - -#include - -namespace uh::cluster::proxy::cache::disk::utils { - -inline coro erase(storage::data_view& storage, const address& addr) { - co_await storage.unlink(addr); -} - -inline coro
store(storage::data_view& storage, - std::span sv) { - auto addr = - co_await storage.write(std::string_view{sv.data(), sv.size()}, {0}); - co_return std::move(addr); -} - -inline coro read(storage::data_view& storage, const address& addr, - std::span sv) { - co_await storage.read_address(addr, sv); -} - -} // namespace uh::cluster::proxy::cache::disk::utils diff --git a/src/proxy/handler.h b/src/proxy/handler.h index dd950f3c7..7c80fe3c3 100644 --- a/src/proxy/handler.h +++ b/src/proxy/handler.h @@ -107,9 +107,11 @@ coro handler::_handle(boost::asio::ip::tcp::socket s, StreamType& ds) { co_await m_factory->create(incoming, rawreq); if (get_object::can_handle(*req)) { - auto d_source = + auto objh = m_mgr.get(cache::disk::object_metadata{req->object_key()}); - if (d_source) { + if (objh) { + auto d_source = + cache::disk::disk_source{m_dv, std::move(objh)}; LOG_INFO() << peer << ": handling from cache"; incoming.set_mode(decltype(incoming)::deleting); @@ -137,7 +139,7 @@ coro handler::_handle(boost::asio::ip::tcp::socket s, StreamType& ds) { co_await async_write( async_write_header(s, serializer, boost::asio::use_awaitable), - s, *d_source); + s, d_source); LOG_INFO() << peer << ": cache result served"; continue; @@ -197,7 +199,8 @@ coro handler::_handle(boost::asio::ip::tcp::socket s, StreamType& ds) { }, outgoing, buffer, *body_size, tee(s_sink, d_sink)); co_await m_mgr.put( - cache::disk::object_metadata{req->object_key()}, d_sink); + cache::disk::object_metadata{req->object_key()}, + d_sink.get_object_handle()); } else { auto body_size = get_content_length(p.get()); diff --git a/src/proxy/http.h b/src/proxy/http.h index 23a9c016b..0f0a9f269 100644 --- a/src/proxy/http.h +++ b/src/proxy/http.h @@ -109,9 +109,9 @@ std::optional get_content_length(const Message& msg) { namespace uh::cluster::proxy { template -coro async_read_header(const SourceType& source, Parser& parser) { - auto header_size = std::vector(source->get_header_size()); - auto header = co_await source->get(header_size); +coro async_read_header(SourceType& source, Parser& parser) { + auto header_size = std::vector(source.get_header_size()); + auto header = co_await source.get(header_size); parser.body_limit(std::numeric_limits::max()); boost::system::error_code ec; diff --git a/test/unit/test_disk_cache_manager.cpp b/test/unit/test_disk_cache_manager.cpp index 9933a0a94..23e5c3243 100644 --- a/test/unit/test_disk_cache_manager.cpp +++ b/test/unit/test_disk_cache_manager.cpp @@ -3,6 +3,7 @@ #include #include #include +#include #include #include #include @@ -26,15 +27,17 @@ BOOST_AUTO_TEST_CASE(put_and_get_with_metadata) { key.path = "/foo/bar"; key.version = "v1"; - boost::asio::co_spawn(m_ioc, mgr.put(key, sink), boost::asio::use_future) + boost::asio::co_spawn(m_ioc, mgr.put(key, sink.get_object_handle()), + boost::asio::use_future) .get(); - auto source = mgr.get(key); - BOOST_TEST(source != nullptr); + auto objh = mgr.get(key); + BOOST_TEST(objh != nullptr); + auto source = disk_source{data_view, objh}; auto buf = std::string(128, '\0'); auto sv = - boost::asio::co_spawn(m_ioc, source->get(buf), boost::asio::use_future) + boost::asio::co_spawn(m_ioc, source.get(buf), boost::asio::use_future) .get(); BOOST_TEST(sv.size() == data.size()); @@ -62,13 +65,13 @@ BOOST_AUTO_TEST_CASE(eviction_test) { key.version = "v" + std::to_string(i); keys.push_back(key); - boost::asio::co_spawn(m_ioc, mgr.put(key, sink), + boost::asio::co_spawn(m_ioc, mgr.put(key, sink.get_object_handle()), boost::asio::use_future) .get(); } - auto source = mgr.get(keys.front()); - BOOST_TEST(source == nullptr); + auto objh = mgr.get(keys.front()); + BOOST_TEST(objh == nullptr); } BOOST_AUTO_TEST_SUITE_END() diff --git a/test/unit/test_trace_asio.cpp b/test/unit/test_trace_asio.cpp index c16d421ed..721b16e59 100644 --- a/test/unit/test_trace_asio.cpp +++ b/test/unit/test_trace_asio.cpp @@ -110,74 +110,3 @@ BOOST_AUTO_TEST_CASE(propagates_context_through_continue) { } BOOST_AUTO_TEST_SUITE_END() - -#include -using namespace boost::asio; -using namespace boost::asio::experimental; - -template -auto wrap_coro_as_deferred(io_context& ioc, Func&& func, - CompletionToken&& token) { - return boost::asio::async_initiate( - boost::asio::co_composed( - [func = std::forward(func)](auto state, - io_context& ioc) -> void { - try { - state.throw_if_cancelled(true); - state.reset_cancellation_state( - boost::asio::enable_terminal_cancellation()); - // Call the provided function and await its result - co_await func(); - } catch (const boost::system::system_error& e) { - co_return {e.code()}; - } - }, - ioc), - token, std::ref(ioc)); -} - -BOOST_FIXTURE_TEST_SUITE(a_traced_coro, fixture) -BOOST_AUTO_TEST_CASE(ranged_parallel_group_streambuf) { - - boost::asio::posix::stream_descriptor out(ioc, ::dup(STDOUT_FILENO)); - boost::asio::posix::stream_descriptor err(ioc, ::dup(STDERR_FILENO)); - - using op_type = decltype(out.async_write_some(boost::asio::buffer("", 0))); - std::cout << boost::core::demangle(typeid(op_type).name()) << "\n"; - auto tmp = wrap_coro_as_deferred( - ioc, - []() -> coro { - // Your coroutine code here - co_return; - }, - deferred); - using op_type_2 = - decltype(out.async_write_some(boost::asio::buffer("", 0))); - std::cout << boost::core::demangle(typeid(op_type_2).name()) << "\n"; - - // std::vector ops; - // - // ops.push_back(out.async_write_some(boost::asio::buffer("first\r\n", 7), - // use_awaitable)); - // - // ops.push_back(err.async_write_some(boost::asio::buffer("second\r\n", 8), - // use_awaitable)); - // - // co_spawn( - // ioc, - // [&]() -> coro { - // auto [completion_order, n] = - // co_await make_parallel_group(ops).async_wait( - // wait_for_one_error(), deferred); - // - // for (std::size_t i = 0; i < completion_order.size(); ++i) { - // std::size_t idx = completion_order[i]; - // std::cout << "operation " << idx << " finished: "; - // std::cout << n[idx] << "\n"; - // } - // }, - // use_future) - // .get(); -} -BOOST_AUTO_TEST_SUITE_END() From dc763d5eaca1bb30e9327a1302a6a17da54197db Mon Sep 17 00:00:00 2001 From: Sungsik Nam Date: Fri, 26 Sep 2025 11:02:10 +0200 Subject: [PATCH 07/10] Add remove interface --- src/proxy/cache/disk/manager.h | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/src/proxy/cache/disk/manager.h b/src/proxy/cache/disk/manager.h index ca3d0d701..3e59129f0 100644 --- a/src/proxy/cache/disk/manager.h +++ b/src/proxy/cache/disk/manager.h @@ -70,6 +70,13 @@ class manager { return entry; } + void remove(object_metadata key) { + auto entry = m_cache->remove(key); + if (entry) { + m_deletion_queue.push(std::move(entry)); + } + } + static manager create(boost::asio::io_context& ioc, data_view& storage, std::size_t capacity, std::size_t eviction_margin = 0) { From a8325489567d7c817dc0fc3e554c2d432453428a Mon Sep 17 00:00:00 2001 From: Sungsik Nam Date: Fri, 26 Sep 2025 11:47:51 +0200 Subject: [PATCH 08/10] Handle put/delete object, handle version --- src/entrypoint/commands/s3/delete_objects.cpp | 20 +++++++++----- src/entrypoint/commands/s3/delete_objects.h | 12 ++++++--- src/proxy/handler.h | 27 +++++++++++++++++-- 3 files changed, 47 insertions(+), 12 deletions(-) diff --git a/src/entrypoint/commands/s3/delete_objects.cpp b/src/entrypoint/commands/s3/delete_objects.cpp index 0dfb0f0d8..3a859e4f0 100644 --- a/src/entrypoint/commands/s3/delete_objects.cpp +++ b/src/entrypoint/commands/s3/delete_objects.cpp @@ -51,18 +51,15 @@ response get_response(const std::vector& success, } // namespace -coro delete_objects::handle(request& req) { - metric::increase(1); - - co_await m_dir.bucket_exists(req.bucket()); - +coro>> +delete_objects::get_delete_object_keys(request& req) { LOG_DEBUG() << req.peer() << ": delete_objects::handle(): content-length: " << req.content_length(); std::string buffer = co_await copy_to_buffer(req.body()); - LOG_DEBUG() << req.peer() << ": delete_objects::handle(): request XML: " - << buffer; + LOG_DEBUG() << req.peer() + << ": delete_objects::handle(): request XML: " << buffer; xml_parser xml_parser; bool parsed = xml_parser.parse(buffer); @@ -72,6 +69,15 @@ coro delete_objects::handle(request& req) { object_nodes.size() > MAXIMUM_DELETE_KEYS) throw command_exception(status::bad_request, "MalformedXML", "XML is invalid."); + co_return object_nodes; +} + +coro delete_objects::handle(request& req) { + metric::increase(1); + + co_await m_dir.bucket_exists(req.bucket()); + + auto object_nodes = co_await get_delete_object_keys(req); auto bucket_id = req.bucket(); std::vector success; diff --git a/src/entrypoint/commands/s3/delete_objects.h b/src/entrypoint/commands/s3/delete_objects.h index 589acb75b..8cdef3685 100644 --- a/src/entrypoint/commands/s3/delete_objects.h +++ b/src/entrypoint/commands/s3/delete_objects.h @@ -1,9 +1,12 @@ #pragma once -#include "entrypoint/directory.h" -#include "entrypoint/limits.h" -#include "storage/global/data_view.h" #include +#include +#include + +#include + +#include namespace uh::cluster { @@ -17,6 +20,9 @@ class delete_objects : public command { std::string action_id() const override; + static coro>> + get_delete_object_keys(ep::http::request& req); + private: directory& m_dir; static constexpr std::size_t MAXIMUM_DELETE_KEYS = 1000; diff --git a/src/proxy/handler.h b/src/proxy/handler.h index 7c80fe3c3..791c9aeb0 100644 --- a/src/proxy/handler.h +++ b/src/proxy/handler.h @@ -13,7 +13,10 @@ #include #include +#include +#include #include +#include #include #include @@ -106,9 +109,29 @@ coro handler::_handle(boost::asio::ip::tcp::socket s, StreamType& ds) { std::unique_ptr req = co_await m_factory->create(incoming, rawreq); + if (put_object::can_handle(*req) || + delete_object::can_handle(*req)) { + m_mgr.remove(cache::disk::object_metadata{ + req->object_key(), req->query("versionId").value_or("")}); + } + if (delete_objects::can_handle(*req)) { + auto objs = + co_await delete_objects::get_delete_object_keys(*req); + for (const auto& obj : objs) { + auto key = + obj.get().template get_optional("Key"); + auto ver = obj.get().template get_optional( + "VersionId"); + + if (key.has_value()) { + m_mgr.remove(cache::disk::object_metadata{ + key.value(), ver.value_or("")}); + } + } + } if (get_object::can_handle(*req)) { - auto objh = - m_mgr.get(cache::disk::object_metadata{req->object_key()}); + auto objh = m_mgr.get(cache::disk::object_metadata{ + req->object_key(), req->query("versionId").value_or("")}); if (objh) { auto d_source = cache::disk::disk_source{m_dv, std::move(objh)}; From a91a1c80ae2a870c8cf1e00ac907ab6b4fa9a8e8 Mon Sep 17 00:00:00 2001 From: Sungsik Nam Date: Fri, 26 Sep 2025 13:28:41 +0200 Subject: [PATCH 09/10] nit --- src/proxy/cache/disk/manager.h | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/src/proxy/cache/disk/manager.h b/src/proxy/cache/disk/manager.h index 3e59129f0..ea8989f55 100644 --- a/src/proxy/cache/disk/manager.h +++ b/src/proxy/cache/disk/manager.h @@ -14,8 +14,6 @@ #include #include -#include - namespace uh::cluster::proxy::cache::disk { class manager { @@ -59,7 +57,7 @@ class manager { if (p_prev) { m_deletion_queue.push(std::move(p_prev)); } - std::cout << "Total size after put: " << m_current_size << std::endl; + LOG_INFO() << "Total size after put: " << m_current_size; } std::shared_ptr get(object_metadata key) { From cd0b07bf4994a98987beb4c8a4e1ed3a6ab41679 Mon Sep 17 00:00:00 2001 From: Sungsik Nam Date: Fri, 26 Sep 2025 15:02:39 +0200 Subject: [PATCH 10/10] delete_objects fixed --- src/entrypoint/commands/s3/delete_objects.cpp | 39 ++++++++++++------- src/entrypoint/commands/s3/delete_objects.h | 10 ++++- src/proxy/cache/disk/manager.h | 2 + src/proxy/handler.h | 15 ++----- 4 files changed, 40 insertions(+), 26 deletions(-) diff --git a/src/entrypoint/commands/s3/delete_objects.cpp b/src/entrypoint/commands/s3/delete_objects.cpp index 3a859e4f0..7e070f2e3 100644 --- a/src/entrypoint/commands/s3/delete_objects.cpp +++ b/src/entrypoint/commands/s3/delete_objects.cpp @@ -51,7 +51,7 @@ response get_response(const std::vector& success, } // namespace -coro>> +coro> delete_objects::get_delete_object_keys(request& req) { LOG_DEBUG() << req.peer() << ": delete_objects::handle(): content-length: " << req.content_length(); @@ -69,7 +69,22 @@ delete_objects::get_delete_object_keys(request& req) { object_nodes.size() > MAXIMUM_DELETE_KEYS) throw command_exception(status::bad_request, "MalformedXML", "XML is invalid."); - co_return object_nodes; + + std::vector targets; + for (const auto& obj : object_nodes) { + const auto& pt = obj.get(); + auto key = pt.get_optional("Key"); + if (!key.has_value()) { + throw command_exception(status::bad_request, "MalformedXML", + "XML is invalid."); + } + targets.emplace_back(*key, // + pt.get_optional("VersionId"), + pt.get_optional("ETag"), + pt.get_optional("LastModifiedTime"), + pt.get_optional("Size")); + } + co_return targets; } coro delete_objects::handle(request& req) { @@ -77,32 +92,28 @@ coro delete_objects::handle(request& req) { co_await m_dir.bucket_exists(req.bucket()); - auto object_nodes = co_await get_delete_object_keys(req); + auto targets = co_await get_delete_object_keys(req); auto bucket_id = req.bucket(); std::vector success; std::vector failure; - for (const auto& obj : object_nodes) { - auto key = obj.get().get_optional("Key"); - if (!key) { - throw command_exception(status::bad_request, "MalformedXML", - "XML is invalid."); - } + for (const auto& t : targets) { + auto& key = t.key; try { LOG_DEBUG() << req.peer() << ": delete_objects::handle(): deleting " - << *key; + << key; std::optional ver; - auto boostver = obj.get().get_optional("VersionId"); + auto boostver = t.version; if (boostver) { ver = *boostver; } - co_await m_dir.delete_object(req.bucket(), *key, ver); - success.emplace_back(*key); + co_await m_dir.delete_object(req.bucket(), key, ver); + success.emplace_back(key); } catch (const std::exception& e) { - failure.emplace_back(*key, e.what()); + failure.emplace_back(key, e.what()); } } diff --git a/src/entrypoint/commands/s3/delete_objects.h b/src/entrypoint/commands/s3/delete_objects.h index 8cdef3685..9ca92b38a 100644 --- a/src/entrypoint/commands/s3/delete_objects.h +++ b/src/entrypoint/commands/s3/delete_objects.h @@ -20,7 +20,15 @@ class delete_objects : public command { std::string action_id() const override; - static coro>> + struct delete_target { + std::string key; + boost::optional version; + boost::optional etag; + boost::optional last_modified; + boost::optional size; + }; + + static coro> get_delete_object_keys(ep::http::request& req); private: diff --git a/src/proxy/cache/disk/manager.h b/src/proxy/cache/disk/manager.h index ea8989f55..b92acced3 100644 --- a/src/proxy/cache/disk/manager.h +++ b/src/proxy/cache/disk/manager.h @@ -71,6 +71,8 @@ class manager { void remove(object_metadata key) { auto entry = m_cache->remove(key); if (entry) { + LOG_INFO() << "key: " << key.path << ", version: " << key.version + << " removed from cache"; m_deletion_queue.push(std::move(entry)); } } diff --git a/src/proxy/handler.h b/src/proxy/handler.h index 791c9aeb0..c60d8f0ad 100644 --- a/src/proxy/handler.h +++ b/src/proxy/handler.h @@ -115,18 +115,11 @@ coro handler::_handle(boost::asio::ip::tcp::socket s, StreamType& ds) { req->object_key(), req->query("versionId").value_or("")}); } if (delete_objects::can_handle(*req)) { - auto objs = + auto targets = co_await delete_objects::get_delete_object_keys(*req); - for (const auto& obj : objs) { - auto key = - obj.get().template get_optional("Key"); - auto ver = obj.get().template get_optional( - "VersionId"); - - if (key.has_value()) { - m_mgr.remove(cache::disk::object_metadata{ - key.value(), ver.value_or("")}); - } + for (const auto& t : targets) { + m_mgr.remove(cache::disk::object_metadata{ + t.key, t.version.value_or("")}); } } if (get_object::can_handle(*req)) {