diff --git a/runtime-light/server/rpc/init-functions.cpp b/runtime-light/server/rpc/init-functions.cpp index 58ae0c767e..9416bebba6 100644 --- a/runtime-light/server/rpc/init-functions.cpp +++ b/runtime-light/server/rpc/init-functions.cpp @@ -194,6 +194,7 @@ void init_server(kphp::component::stream&& request_stream, kphp::stl::vector(std::addressof(reserved_header)), sizeof(reserved_header)}); +} + mixed mixed_array_get_value(const mixed& arr, const string& str_key, int64_t num_key) noexcept { if (!arr.is_array()) [[unlikely]] { return {}; @@ -270,37 +279,48 @@ kphp::coro::task> typed_rpc_tl_query_result_ } // namespace detail +void clean_buffers() noexcept { + auto& rpc_server_instance_st{RpcServerInstanceState::get()}; + rpc_server_instance_st.tl_storer.clear(); + kphp::rpc::detail::reserve_header(); + // TODO we need this just because we have one buffer for f$store_* functions and f$fetch_* functions. + // if we make another buffer for `rpc_server_instance_st.tl_fetcher`, then we don't need this. + rpc_server_instance_st.tl_fetcher = tl::fetcher{rpc_server_instance_st.tl_storer.view().subspan(kphp::rpc::detail::RESERVED_HEADER_SIZE)}; +} + kphp::rpc::query_info send_request(std::string_view actor, std::optional opt_timeout, bool ignore_answer, bool collect_responses_extra_info) noexcept { auto& rpc_client_instance_st{RpcClientInstanceState::get()}; auto& rpc_server_instance_st{RpcServerInstanceState::get()}; const auto timestamp{std::chrono::duration{std::chrono::system_clock::now().time_since_epoch()}.count()}; - // if we have to allocate memory for request buffer, it will be in this vector, and it will be freed at the end of the function - std::optional> opt_request_vec{}; - std::span request_buffer{rpc_server_instance_st.tl_storer.view()}; + // We have reserved place for one `kphp::rpc::dest_actor_flags_header` in `rpc_server_instance_st.tl_storer` + // before storing request and calling `send_request(...)`, + // so the real serialized request starts after `RESERVED_HEADER_SIZE` bytes in tl storer. + // We do this to have enough place for regularized header after `kphp::rpc::regularize_extra_headers(...)` call. + // This optimization helps us avoid allocating and copying the whole request. + std::span request_buffer{rpc_server_instance_st.tl_storer.view().subspan(detail::RESERVED_HEADER_SIZE)}; if (const auto& [opt_new_extra_header, cur_extra_header_size]{kphp::rpc::regularize_extra_headers(request_buffer, ignore_answer)}; opt_new_extra_header) { std::span new_header{reinterpret_cast(std::addressof(*opt_new_extra_header)), sizeof(std::remove_cvref_t)}; - std::span request_body{rpc_server_instance_st.tl_storer.view().subspan(cur_extra_header_size)}; - - std::span new_request_buffer{}; - size_t request_and_headers_size{new_header.size() + request_body.size()}; - if (request_and_headers_size <= StringLibContext::STATIC_BUFFER_LENGTH) { - // we have enough space in static buffer to store request with regularized headers - auto& string_lib_ctx{StringLibContext::get()}; - new_request_buffer = std::span{reinterpret_cast(string_lib_ctx.static_buf.get()), request_and_headers_size}; - } else { - // we have to allocate buffer for request with regularized headers - opt_request_vec.emplace(request_and_headers_size); - new_request_buffer = std::span{opt_request_vec->data(), request_and_headers_size}; - } - - std::ranges::copy(new_header, new_request_buffer.subspan(0, new_header.size()).begin()); - std::ranges::copy(request_body, new_request_buffer.subspan(new_header.size()).begin()); - request_buffer = new_request_buffer; + // Let's name `request_body` as `request_buffer.subspan(cur_extra_header_size)}` + + // If `regularize_extra_headers` gave us new header, then we must serialize `new_header` before `request_body`. + // + // here serialized request (business logic) starts + // \/ + // tl_storer was: |reserved dest-actor-flags-header| [optional old header] |request-body| + // + // We want to serialize |our new header| right before |request-body| : + // + // tl_storer will be: ... may be some bytes leaved here ... |our new header| |request-body| + + // we do always have enough bytes for `new_header` before `request_body`, because we have reserved it in `f$rpc_clean(...)` call. + size_t new_header_offset{detail::RESERVED_HEADER_SIZE + cur_extra_header_size - new_header.size()}; + request_buffer = rpc_server_instance_st.tl_storer.view().subspan(new_header_offset); + std::ranges::copy(new_header, request_buffer.data()); } const size_t request_size{request_buffer.size_bytes()}; diff --git a/runtime-light/stdlib/rpc/rpc-api.h b/runtime-light/stdlib/rpc/rpc-api.h index d02358bd44..fdef374024 100644 --- a/runtime-light/stdlib/rpc/rpc-api.h +++ b/runtime-light/stdlib/rpc/rpc-api.h @@ -24,6 +24,7 @@ #include "runtime-light/stdlib/rpc/rpc-client-state.h" #include "runtime-light/stdlib/rpc/rpc-constants.h" #include "runtime-light/stdlib/rpc/rpc-exceptions.h" +#include "runtime-light/stdlib/rpc/rpc-extra-headers.h" #include "runtime-light/stdlib/rpc/rpc-extra-info.h" #include "runtime-light/stdlib/rpc/rpc-tl-error.h" #include "runtime-light/stdlib/rpc/rpc-tl-function.h" @@ -58,6 +59,8 @@ inline kphp::coro::task> send_response(std::span opt_timeout, bool collect_resp_extra_info, bool ignore_answer) noexcept; @@ -70,6 +73,8 @@ kphp::coro::task> typed_rpc_tl_query_result_ } // namespace detail +void clean_buffers() noexcept; + } // namespace kphp::rpc // === server ===================================================================================== @@ -197,9 +202,7 @@ inline void f$fetch_raw_vector_double(array& vector, int64_t num_elems) } inline bool f$rpc_clean() noexcept { - auto& rpc_server_instance_st{RpcServerInstanceState::get()}; - rpc_server_instance_st.tl_storer.clear(); - rpc_server_instance_st.tl_fetcher = tl::fetcher{rpc_server_instance_st.tl_storer.view()}; + kphp::rpc::clean_buffers(); return true; } @@ -268,8 +271,9 @@ inline kphp::coro::task<> f$rpc_server_store_response(class_instance view() noexcept { + return m_buffer; + } + void clear() noexcept { m_buffer.clear(); }