Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions runtime-light/server/rpc/init-functions.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -194,6 +194,7 @@ void init_server(kphp::component::stream&& request_stream, kphp::stl::vector<std
auto& rpc_server_instance_st{RpcServerInstanceState::get()};
rpc_server_instance_st.request_stream = std::move(request_stream);
rpc_server_instance_st.query_id = invoke_rpc.query_id.value;
rpc_server_instance_st.tl_storer.clear();
rpc_server_instance_st.tl_storer.store_bytes(invoke_rpc.query);
rpc_server_instance_st.tl_fetcher = tl::fetcher{rpc_server_instance_st.tl_storer.view()};

Expand Down
60 changes: 40 additions & 20 deletions runtime-light/stdlib/rpc/rpc-api.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,15 @@ namespace kphp::rpc {

namespace detail {

// store bytes for `kphp::rpc::dest_actor_flags_header` in RpcServerInstanceState::tl_storer.
// we do this to avoid allocating new buffer for regularized rpc extra headers and copying whole request.
void reserve_header() noexcept {
auto& rpc_server_instance_st{RpcServerInstanceState::get()};
kphp::rpc::dest_actor_flags_header reserved_header{};
static_assert(sizeof(reserved_header) == RESERVED_HEADER_SIZE);
rpc_server_instance_st.tl_storer.store_bytes({reinterpret_cast<const std::byte*>(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 {};
Expand Down Expand Up @@ -270,37 +279,48 @@ kphp::coro::task<class_instance<C$VK$TL$RpcResponse>> 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<double> 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<double>{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<kphp::stl::vector<std::byte, kphp::memory::script_allocator>> opt_request_vec{};
std::span<const std::byte> 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<std::byte> 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<const std::byte> new_header{reinterpret_cast<const std::byte*>(std::addressof(*opt_new_extra_header)),
sizeof(std::remove_cvref_t<decltype(*opt_new_extra_header)>)};
std::span<const std::byte> request_body{rpc_server_instance_st.tl_storer.view().subspan(cur_extra_header_size)};

std::span<std::byte> 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<std::byte>{reinterpret_cast<std::byte*>(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<std::byte>{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()};
Expand Down
14 changes: 9 additions & 5 deletions runtime-light/stdlib/rpc/rpc-api.h
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -58,6 +59,8 @@ inline kphp::coro::task<std::expected<void, int32_t>> send_response(std::span<co

namespace detail {

static constexpr size_t RESERVED_HEADER_SIZE{sizeof(kphp::rpc::dest_actor_flags_header)};

kphp::rpc::query_info rpc_tl_query_one_impl(std::string_view actor, const mixed& tl_object, std::optional<double> opt_timeout, bool collect_resp_extra_info,
bool ignore_answer) noexcept;

Expand All @@ -70,6 +73,8 @@ kphp::coro::task<class_instance<C$VK$TL$RpcResponse>> typed_rpc_tl_query_result_

} // namespace detail

void clean_buffers() noexcept;

} // namespace kphp::rpc

// === server =====================================================================================
Expand Down Expand Up @@ -197,9 +202,7 @@ inline void f$fetch_raw_vector_double(array<double>& 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;
}

Expand Down Expand Up @@ -268,8 +271,9 @@ inline kphp::coro::task<> f$rpc_server_store_response(class_instance<C$VK$TL$Rpc
// as we are in a coroutine, we must own the data to prevent it from being overwritten by another coroutine,
// so create a TLBuffer owned by this coroutine
auto& rpc_server_instance_st{RpcServerInstanceState::get()};
tl::K2RpcResponse rpc_response{.value =
tl::k2RpcResponseHeader{.flags = {}, .extra_flags = {}, .extra = {}, .result = rpc_server_instance_st.tl_storer.view()}};
tl::K2RpcResponse rpc_response{
.value = tl::k2RpcResponseHeader{
.flags = {}, .extra_flags = {}, .extra = {}, .result = rpc_server_instance_st.tl_storer.view().subspan(kphp::rpc::detail::RESERVED_HEADER_SIZE)}};
tl::storer tls{rpc_response.footprint()};
rpc_response.store(tls);

Expand Down
4 changes: 4 additions & 0 deletions runtime-light/tl/tl-core.h
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,10 @@ class storer {
return m_buffer;
}

std::span<std::byte> view() noexcept {
return m_buffer;
}

void clear() noexcept {
m_buffer.clear();
}
Expand Down
Loading