Skip to content
Draft
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
4 changes: 4 additions & 0 deletions runtime-common/core/allocator/platform-allocator.h
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
#pragma once

#include <cstddef>
#include <type_traits>

#include "runtime-common/core/allocator/runtime-allocator.h"

Expand All @@ -13,6 +14,9 @@ namespace kphp::memory {
template<typename T>
struct platform_allocator {
using value_type = T;
using propagate_on_container_copy_assignment = std::true_type;
using propagate_on_container_move_assignment = std::true_type;
using is_always_equal = std::true_type;

platform_allocator() noexcept = default;

Expand Down
3 changes: 3 additions & 0 deletions runtime-common/core/allocator/script-allocator.h
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,9 @@ namespace memory {
template<typename T>
struct script_allocator {
using value_type = T;
using propagate_on_container_copy_assignment = std::true_type;
using propagate_on_container_move_assignment = std::true_type;
using is_always_equal = std::true_type;

script_allocator() noexcept = default;

Expand Down
52 changes: 25 additions & 27 deletions runtime-light/coroutine/await-set.h
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@
#pragma once

#include <cstddef>
#include <memory>
#include <optional>
#include <type_traits>

Expand All @@ -14,62 +13,61 @@
#include "runtime-light/coroutine/coroutine-state.h"
#include "runtime-light/coroutine/detail/await-set.h"
#include "runtime-light/coroutine/type-traits.h"
#include "runtime-light/stdlib/diagnostics/logs.h"

namespace kphp::coro {

template<typename return_type>
class await_set {
std::unique_ptr<detail::await_set::await_broker<return_type>> m_await_broker;
detail::await_set::await_broker<return_type> m_await_broker;
kphp::coro::async_stack_root& m_coroutine_stack_root;

public:
await_set() noexcept
: m_await_broker(std::make_unique<detail::await_set::await_broker<return_type>>()),
m_coroutine_stack_root(CoroutineInstanceState::get().coroutine_stack_root) {}

await_set(await_set&& other) noexcept
: m_await_broker(std::move(other.m_await_broker)),
m_coroutine_stack_root(other.m_coroutine_stack_root) {}

await_set& operator=(await_set&& other) noexcept {
if (this != std::addressof(other)) {
m_await_broker = std::move(other.m_await_broker);
m_coroutine_stack_root = other.m_coroutine_stack_root;
}
return *this;
}
: m_coroutine_stack_root(CoroutineInstanceState::get().coroutine_stack_root) {}

await_set(const await_set&) = delete;
await_set(await_set&& other) = delete;

await_set& operator=(const await_set&) = delete;
await_set& operator=(await_set&& other) = delete;

~await_set() = default;

template<typename... Args>
void* operator new(size_t n, [[maybe_unused]] Args&&... args) noexcept {
return kphp::memory::script::alloc(n);
}

template<typename... Args>
auto operator new(size_t n, std::align_val_t al, [[maybe_unused]] Args&&... args) noexcept -> void* {
return kphp::memory::script::alloc_aligned(n, al);
}

void operator delete(void* ptr, [[maybe_unused]] size_t n) noexcept {
kphp::memory::script::free(ptr);
}

template<typename awaitable_type>
requires kphp::coro::concepts::awaitable<awaitable_type> && std::is_same_v<typename awaitable_traits<awaitable_type>::awaiter_return_type, return_type>
void push(awaitable_type awaitable) noexcept {
kphp::log::assertion(m_await_broker != nullptr);
m_await_broker->start_task(detail::await_set::make_await_set_task(std::move(awaitable)), m_coroutine_stack_root, STACK_RETURN_ADDRESS);
m_await_broker.start_task(detail::await_set::make_await_set_task(std::move(awaitable)), m_coroutine_stack_root, STACK_RETURN_ADDRESS);
}

auto next() noexcept {
kphp::log::assertion(m_await_broker != nullptr);
return detail::await_set::await_set_awaitable<return_type>{*m_await_broker};
return detail::await_set::await_set_awaitable<return_type>{m_await_broker};
}

auto try_next() noexcept {
using result_type = std::optional<decltype(std::declval<typename detail::await_set::await_set_task<return_type>::promise_type>().result())>;
if (m_await_broker == nullptr) [[unlikely]] {
return result_type{std::nullopt};
}
return result_type{m_await_broker->try_get_result()};
return result_type{m_await_broker.try_get_result()};
}

bool empty() const noexcept {
return size() == 0;
}

size_t size() const noexcept {
kphp::log::assertion(m_await_broker != nullptr);
return m_await_broker->size();
return m_await_broker.size();
}
};

Expand Down
16 changes: 1 addition & 15 deletions runtime-light/coroutine/detail/await-set.h
Original file line number Diff line number Diff line change
Expand Up @@ -51,20 +51,6 @@ class await_broker {
await_broker& operator=(const await_broker&) = delete;
await_broker& operator=(await_broker&& other) = delete;

template<typename... Args>
void* operator new(size_t n, [[maybe_unused]] Args&&... args) noexcept {
return kphp::memory::script::alloc(n);
}

template<typename... Args>
auto operator new(size_t n, std::align_val_t al, [[maybe_unused]] Args&&... args) noexcept -> void* {
return kphp::memory::script::alloc_aligned(n, al);
}

void operator delete(void* ptr, [[maybe_unused]] size_t n) noexcept {
kphp::memory::script::free(ptr);
}

void start_task(await_set_task<return_type>&& task, kphp::coro::async_stack_root& coroutine_stack_root, void* return_address) noexcept {
auto& promise{task.m_promise};
m_tasks_storage.push_front(promise.m_coroutine_node);
Expand Down Expand Up @@ -154,7 +140,7 @@ class await_broker {
}
}

size_t size() noexcept {
size_t size() const noexcept {
return m_tasks_count;
}

Expand Down
51 changes: 25 additions & 26 deletions runtime-light/coroutine/event.h
Original file line number Diff line number Diff line change
Expand Up @@ -6,22 +6,21 @@

#include <concepts>
#include <coroutine>
#include <memory>
#include <utility>
#include <variant>

#include "common/containers/intrusive-list.h"
#include "common/mixin/not_copyable.h"
#include "common/wrappers/overloaded.h"
#include "runtime-common/core/allocator/script-allocator-managed.h"
#include "runtime-common/core/allocator/script-malloc-interface.h"
#include "runtime-light/coroutine/async-stack.h"
#include "runtime-light/coroutine/coroutine-state.h"
#include "runtime-light/stdlib/diagnostics/logs.h"

namespace kphp::coro {

class event {
struct event_controller : kphp::memory::script_allocator_managed, vk::not_copyable {
struct event_controller : vk::not_copyable {
// 1) std::monostate => not set and no coroutines are waiting
// 2) non empty list => linked list of coroutines waiting for the event to trigger
// 3) empty list => the event is triggered and all coroutines are resumed
Expand Down Expand Up @@ -55,28 +54,32 @@ class event {
auto await_resume() noexcept -> void;
};

std::unique_ptr<event_controller> m_controller;
event_controller m_controller;

public:
event() noexcept
: m_controller(std::make_unique<event_controller>()) {
kphp::log::assertion(m_controller != nullptr);
}
event() = default;

event(event&& other) noexcept
: m_controller(std::move(other.m_controller)) {}
event(const event&) = delete;
event(event&& other) = delete;

event& operator=(event&& other) noexcept {
if (this != std::addressof(other)) {
m_controller = std::move(other.m_controller);
}
return *this;
}
event& operator=(const event&) = delete;
event& operator=(event&& other) = delete;

~event() = default;

event(const event&) = delete;
event& operator=(const event&) = delete;
template<typename... Args>
void* operator new(size_t n, [[maybe_unused]] Args&&... args) noexcept {
return kphp::memory::script::alloc(n);
}

template<typename... Args>
auto operator new(size_t n, std::align_val_t al, [[maybe_unused]] Args&&... args) noexcept -> void* {
return kphp::memory::script::alloc_aligned(n, al);
}

void operator delete(void* ptr, [[maybe_unused]] size_t n) noexcept {
kphp::memory::script::free(ptr);
}

auto set() noexcept -> void;
auto unset() noexcept -> void;
Expand Down Expand Up @@ -148,23 +151,19 @@ inline auto event::event_controller::is_set() const noexcept -> bool {
}

inline auto event::set() noexcept -> void {
kphp::log::assertion(m_controller != nullptr);
m_controller->set();
m_controller.set();
}

inline auto event::unset() noexcept -> void {
kphp::log::assertion(m_controller != nullptr);
m_controller->unset();
m_controller.unset();
}

inline auto event::is_set() const noexcept -> bool {
kphp::log::assertion(m_controller != nullptr);
return m_controller->is_set();
return m_controller.is_set();
}

inline auto event::operator co_await() noexcept {
kphp::log::assertion(m_controller != nullptr);
return event::awaiter{*this->m_controller};
return event::awaiter{m_controller};
}

} // namespace kphp::coro
8 changes: 5 additions & 3 deletions runtime-light/state/instance-state.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -235,9 +235,11 @@ kphp::coro::task<> InstanceState::run_instance_epilogue() noexcept {
*/
{
auto& rpc_client_instance_st{RpcClientInstanceState::get()};
auto ignore_answer_request_await_set{std::exchange(rpc_client_instance_st.ignore_answer_request_awaiter_tasks, kphp::coro::await_set<void>{})};
while (!ignore_answer_request_await_set.empty()) {
co_await ignore_answer_request_await_set.next();
auto ignore_answer_request_await_set{
std::exchange(rpc_client_instance_st.ignore_answer_request_awaiter_tasks, std::make_unique<kphp::coro::await_set<void>>())};
kphp::log::assertion(ignore_answer_request_await_set != nullptr);
while (!ignore_answer_request_await_set->empty()) {
co_await ignore_answer_request_await_set->next();
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,7 @@ class client final {

// Wait until transport is available
if (is_occupied) [[unlikely]] {
transport_readiness_notifier.emplace(qid, kphp::coro::event{});
transport_readiness_notifier.try_emplace(qid);
queue.push(qid);
co_await transport_readiness_notifier[qid];
}
Expand Down Expand Up @@ -105,7 +105,7 @@ class client final {
}

auto write(shared_transport_type t, query_id_type qid, std::span<const std::byte> payload) noexcept -> kphp::coro::task<std::expected<void, int32_t>> {
req_finish_notifier.emplace(qid, kphp::coro::event{});
req_finish_notifier.try_emplace(qid);

// The protocol design assumes that interrupting the transfer in the middle of a frame leads to critical error.
// Therefore, we need to write the request in a separate coroutine.
Expand Down Expand Up @@ -230,7 +230,7 @@ class client final {
ctx.get()->query2resp_buffer_provider.emplace(qid, std::move(buffer_provider));

// Register notifier
ctx.get()->resp_finish_notifier.emplace(qid, kphp::coro::event{});
ctx.get()->resp_finish_notifier.try_emplace(qid);
}
} reader;

Expand Down
2 changes: 1 addition & 1 deletion runtime-light/stdlib/fork/wait-queue-state.h
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ class WaitQueueInstanceState {

[[nodiscard]] int64_t create_queue() noexcept {
const int64_t wait_queue_id{m_next_wait_queue_id++};
m_queues.emplace(wait_queue_id, kphp::coro::await_set<int64_t>{});
m_queues.try_emplace(wait_queue_id);
return wait_queue_id;
}

Expand Down
3 changes: 2 additions & 1 deletion runtime-light/stdlib/rpc/rpc-api.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -367,7 +367,8 @@ kphp::coro::task<kphp::rpc::query_info> send_request(std::string_view actor, std
auto ignore_answer_awaiter_task{ignore_answer_awaiter_coroutine(std::move(stream), timeout)};
kphp::log::assertion(kphp::coro::io_scheduler::get().start(ignore_answer_awaiter_task));

rpc_client_instance_st.ignore_answer_request_awaiter_tasks.push(std::move(ignore_answer_awaiter_task));
kphp::log::assertion(rpc_client_instance_st.ignore_answer_request_awaiter_tasks != nullptr);
rpc_client_instance_st.ignore_answer_request_awaiter_tasks->push(std::move(ignore_answer_awaiter_task));
co_return kphp::rpc::query_info{.id = kphp::rpc::IGNORED_ANSWER_QUERY_ID, .request_size = request_size, .timestamp = timestamp};
}
// start awaiter task
Expand Down
3 changes: 2 additions & 1 deletion runtime-light/stdlib/rpc/rpc-client-state.h
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
#pragma once

#include <cstdint>
#include <memory>
#include <optional>
#include <utility>

Expand All @@ -27,7 +28,7 @@ struct RpcClientInstanceState final : private vk::not_copyable {
kphp::stl::unordered_map<int64_t, class_instance<RpcTlQuery>, kphp::memory::script_allocator> response_fetcher_instances;
kphp::stl::unordered_map<int64_t, std::pair<kphp::rpc::response_extra_info_status, kphp::rpc::response_extra_info>, kphp::memory::script_allocator>
rpc_responses_extra_info;
kphp::coro::await_set<void> ignore_answer_request_awaiter_tasks;
std::unique_ptr<kphp::coro::await_set<void>> ignore_answer_request_awaiter_tasks{std::make_unique<kphp::coro::await_set<void>>()};

RpcClientInstanceState() noexcept = default;

Expand Down
2 changes: 1 addition & 1 deletion runtime-light/stdlib/rpc/rpc-queue-state.h
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ class RpcQueueInstanceState final : private vk::not_copyable {

[[nodiscard]] int64_t create_queue() noexcept {
const int64_t wait_queue_id{m_rpc_wait_queue_id++};
m_queues.emplace(wait_queue_id, kphp::coro::await_set<int64_t>{});
m_queues.try_emplace(wait_queue_id);
return wait_queue_id;
}

Expand Down
8 changes: 5 additions & 3 deletions runtime-light/streams/connection.h
Original file line number Diff line number Diff line change
Expand Up @@ -24,12 +24,14 @@ class connection {
std::optional<kphp::coro::event> m_unwatch_event;

shared_state() noexcept = default;
shared_state(shared_state&&) noexcept = default;
shared_state& operator=(shared_state&&) noexcept = default;
~shared_state() = default;

shared_state(shared_state&&) noexcept = delete;
shared_state(const shared_state&) = delete;

shared_state& operator=(shared_state&&) noexcept = delete;
shared_state operator=(const shared_state&) = delete;

~shared_state() = default;
};

class_instance<shared_state> m_shared_state;
Expand Down
Loading