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/entrypoint/commands/s3/delete_objects.cpp b/src/entrypoint/commands/s3/delete_objects.cpp index 0dfb0f0d8..7e070f2e3 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); @@ -73,30 +70,50 @@ coro delete_objects::handle(request& req) { throw command_exception(status::bad_request, "MalformedXML", "XML is invalid."); - auto bucket_id = req.bucket(); - std::vector success; - std::vector failure; + std::vector targets; for (const auto& obj : object_nodes) { - auto key = obj.get().get_optional("Key"); - if (!key) { + 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) { + metric::increase(1); + + co_await m_dir.bucket_exists(req.bucket()); + + auto targets = co_await get_delete_object_keys(req); + + auto bucket_id = req.bucket(); + std::vector success; + std::vector failure; + 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 589acb75b..9ca92b38a 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,17 @@ class delete_objects : public command { std::string action_id() const override; + 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: directory& m_dir; static constexpr std::size_t MAXIMUM_DELETE_KEYS = 1000; 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..b92acced3 100644 --- a/src/proxy/cache/disk/manager.h +++ b/src/proxy/cache/disk/manager.h @@ -1,19 +1,19 @@ #pragma once #include -#include +#include #include #include #include +#include +#include #include #include #include -#include - namespace uh::cluster::proxy::cache::disk { class manager { @@ -27,8 +27,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 = @@ -58,15 +57,24 @@ 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::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; + } + + 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)); + } } 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..c60d8f0ad 100644 --- a/src/proxy/handler.h +++ b/src/proxy/handler.h @@ -13,7 +13,10 @@ #include #include +#include +#include #include +#include #include #include @@ -106,10 +109,25 @@ 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 targets = + co_await delete_objects::get_delete_object_keys(*req); + for (const auto& t : targets) { + m_mgr.remove(cache::disk::object_metadata{ + t.key, t.version.value_or("")}); + } + } if (get_object::can_handle(*req)) { - auto d_source = - m_mgr.get(cache::disk::object_metadata{req->object_key()}); - if (d_source) { + 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)}; LOG_INFO() << peer << ": handling from cache"; incoming.set_mode(decltype(incoming)::deleting); @@ -137,7 +155,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 +215,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()