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
3 changes: 1 addition & 2 deletions rpc-lib/rpc-lib/internal/channel/channel.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -43,10 +43,9 @@ namespace vshalygin::rpc::internal {
const auto method_idx = static_cast<uint32_t>(method->index());
m_connection->request_async(create_transfer_msg_req(req_id, method_idx, request))
.then([cg = std::move(cg), request_controller, response, req_id] (auto value) mutable {
auto guard = std::move(cg);
return value.lock().with(
[&](request_result rc, cl::buffer &&buffer) mutable {
auto guard = std::move(cg);

try {
if(is_success(rc)) {
if(!is_response_buffer_valid(buffer)) {
Expand Down
104 changes: 60 additions & 44 deletions rpc-lib/rpc-lib/internal/connection/connection.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -162,17 +162,17 @@ namespace vshalygin::rpc::internal {
add_request_to_map(msg_number, std::move(promise));

auto send_callback = [self = shared_from_this(), msg_number](auto value) {
value.lock().with([&](pipe_op_res res) {
if(is_success(res)) {
self->process_request_sent(msg_number);
} else if(res == pipe_op_res::canceled) {
self->complete_request(msg_number, request_result::send_canceled, {});
} else if(res == pipe_op_res::timeout) {
self->complete_request(msg_number, request_result::send_timeout, {});
} else {
self->complete_request(msg_number, request_result::send_failed, {});
}
});
pipe_op_res result = *value.lock();

if(is_success(result)) {
self->process_request_sent(msg_number);
} else if(result == pipe_op_res::canceled) {
self->complete_request(msg_number, request_result::send_canceled, {});
} else if(result == pipe_op_res::timeout) {
self->complete_request(msg_number, request_result::send_timeout, {});
} else {
self->complete_request(msg_number, request_result::send_failed, {});
}
};

m_transport.send_async(std::move(message))
Expand Down Expand Up @@ -202,14 +202,19 @@ namespace vshalygin::rpc::internal {
{
m_transport.recv_async()
.then([self = shared_from_this()](auto value) mutable {
return value.lock().with([&](pipe_op_res r, cl::buffer &&message) {
if(is_success(r)) {
self->dispatch_receive_event(std::move(message));
self->do_receive_async();
} else {
self->complete_receive_routine(r);
}
});
pipe_op_res result = pipe_op_res::failed;
cl::buffer message;
value.lock().with([&](pipe_op_res value_result, cl::buffer &&value_message) {
result = value_result;
message = std::move(value_message);
});

if(is_success(result)) {
self->dispatch_receive_event(std::move(message));
self->do_receive_async();
} else {
self->complete_receive_routine(result);
}
});
}

Expand Down Expand Up @@ -270,21 +275,27 @@ namespace vshalygin::rpc::internal {

void connection::impl::process_request_sent(uint64_t req_msg_number)
{
req_result_promise to_delete;
req_result_promise promise;
std::optional<request_result> result;

auto map = m_request_map.lock();
auto it = map->find(req_msg_number);
if(it != map->end()) {
auto &req_data = it->second;
req_data.is_req_sent = true;
if(req_data.fail_req_result) {
assert(is_fail(*req_data.fail_req_result));
req_data.promise.set_value(cl::ftuple(*req_data.fail_req_result,
cl::buffer{}));
to_delete = std::move(req_data.promise);
map->erase(it);
{
auto map = m_request_map.lock();
auto it = map->find(req_msg_number);
if(it != map->end()) {
auto &req_data = it->second;
req_data.is_req_sent = true;
if(req_data.fail_req_result) {
assert(is_fail(*req_data.fail_req_result));
result = *req_data.fail_req_result;
promise = std::move(req_data.promise);
map->erase(it);
}
}
}

if(result) {
promise.set_value(cl::ftuple(*result, cl::buffer{}));
}
}

void connection::impl::process_watch_event()
Expand Down Expand Up @@ -345,22 +356,27 @@ namespace vshalygin::rpc::internal {
request_result::canceled :
request_result::failed;

std::vector<req_result_promise> to_delete;
std::vector<req_result_promise> promises;

auto map = m_request_map.lock();
for(auto it = map->begin(); it != map->end(); ) {
auto &req_data = it->second;
m_multiple_timer.cancel(req_data.timer_id);

if(req_data.is_req_sent) {
req_data.promise.set_value(cl::ftuple(req_result, cl::buffer{}));
to_delete.push_back(std::move(req_data.promise));
it = map->erase(it);
} else {
req_data.fail_req_result = req_result;
++it;
{
auto map = m_request_map.lock();
for(auto it = map->begin(); it != map->end(); ) {
auto &req_data = it->second;
m_multiple_timer.cancel(req_data.timer_id);

if(req_data.is_req_sent) {
promises.push_back(std::move(req_data.promise));
it = map->erase(it);
} else {
req_data.fail_req_result = req_result;
++it;
}
}
}

for(auto &promise : promises) {
promise.set_value(cl::ftuple(req_result, cl::buffer{}));
}
}

size_t connection::impl::get_pending_requests_count() const
Expand Down
129 changes: 78 additions & 51 deletions rpc-lib/rpc-lib/pipe/memory-pipe/mem-buffer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@

#include <common-lib/thread/thread-pool/thread-pool.h>

#include <vector>

namespace vshalygin::rpc {
std::shared_ptr<mem_buffer> mem_buffer::create(cl::thread_pool *thread_pool)
{
Expand Down Expand Up @@ -65,24 +67,35 @@ namespace vshalygin::rpc {

void mem_buffer::invalidate(bool cancel_read)
{
std::lock_guard guard(m_mtx);
if(m_is_valid) {
std::vector<read_promise> pending_promises;
{
std::lock_guard guard(m_mtx);
if(!m_is_valid) {
return;
}

m_is_valid = false;
m_read_invalidation_result = cancel_read ? pipe_op_res::canceled
: pipe_op_res::failed;
m_buffer = {};

auto read_promises = m_read_promises->lock();
auto &q = read_promises->get<0>();
pending_promises.reserve(q.size());
for(auto it = q.begin(); it != q.end(); ++it) {
q.modify(it, [cancel_read](read_promise_data &el) {
el.promise.set_value(cl::ftuple(
cancel_read ? pipe_op_res::canceled : pipe_op_res::failed,
cl::buffer{}));
q.modify(it, [&pending_promises](read_promise_data &el) {
pending_promises.push_back(std::move(el.promise));
});
}
q.clear();
m_timer.cancel_all();
}

for(auto &promise : pending_promises) {
promise.set_value(cl::ftuple(
cancel_read ? pipe_op_res::canceled : pipe_op_res::failed,
cl::buffer{}));
}
}

bool mem_buffer::is_valid() const
Expand Down Expand Up @@ -125,72 +138,86 @@ namespace vshalygin::rpc {

void mem_buffer::resolve_read_promise()
{
read_promise to_delete;
read_promise promise;
cl::buffer buffer;

std::lock_guard guard(m_mtx);
auto read_promises = m_read_promises->lock();
auto &q = read_promises->get<0>();
if(!q.empty() && !m_buffer.empty()) {
auto buffer = std::move(m_buffer.front());
{
std::lock_guard guard(m_mtx);
auto read_promises = m_read_promises->lock();
auto &q = read_promises->get<0>();
if(q.empty() || m_buffer.empty()) {
return;
}

buffer = std::move(m_buffer.front());
m_buffer.pop();

if(q.begin()->timer_id) {
m_timer.cancel(*q.begin()->timer_id);
}

q.modify(q.begin(), [&buffer, &to_delete](read_promise_data &el) mutable {
el.promise.set_value(cl::ftuple(pipe_op_res::success, std::move(buffer)));
to_delete = std::move(el.promise);

q.modify(q.begin(), [&promise](read_promise_data &el) {
promise = std::move(el.promise);
});
q.pop_front();
}

promise.set_value(cl::ftuple(pipe_op_res::success, std::move(buffer)));
}

void mem_buffer::read_impl(read_promise promise,
const std::optional<std::chrono::milliseconds> &timeout,
bool started_while_valid)
{
std::lock_guard guard(m_mtx);
std::optional<pipe_op_res> result;
cl::buffer result_buffer;

if(!m_is_valid) {
promise.set_value(cl::ftuple(
started_while_valid ? m_read_invalidation_result : pipe_op_res::failed,
cl::buffer{}));
} else if(!m_buffer.empty()) {
auto msg = std::move(m_buffer.front());
m_buffer.pop();
promise.set_value(cl::ftuple(pipe_op_res::success, std::move(msg)));
} else {
const auto id = m_next_read_promise_id++;

auto read_promises = m_read_promises->lock();
std::optional<uint64_t> timer_id;
if(timeout) {
auto timeout_callback = [id, read_promises_wp = std::weak_ptr(m_read_promises)]() {
if(auto read_promises = read_promises_wp.lock()) {
//avoid deadlock in case future callback stores mem_buffer itself
read_promise promise;
{
auto promises = read_promises->lock();
auto &m = promises->get<1>();
auto it = m.find(id);
if(it != m.end()) {
m.modify(it, [&promise](read_promise_data &el) {
promise = std::move(el.promise);
});
m.erase(it);
{
std::lock_guard guard(m_mtx);

if(!m_is_valid) {
result = started_while_valid ? m_read_invalidation_result
: pipe_op_res::failed;
} else if(!m_buffer.empty()) {
result = pipe_op_res::success;
result_buffer = std::move(m_buffer.front());
m_buffer.pop();
} else {
const auto id = m_next_read_promise_id++;

auto read_promises = m_read_promises->lock();
std::optional<uint64_t> timer_id;
if(timeout) {
auto timeout_callback = [id, read_promises_wp = std::weak_ptr(m_read_promises)]() {
if(auto read_promises = read_promises_wp.lock()) {
//avoid deadlock in case future callback stores mem_buffer itself
read_promise promise;
{
auto promises = read_promises->lock();
auto &m = promises->get<1>();
auto it = m.find(id);
if(it != m.end()) {
m.modify(it, [&promise](read_promise_data &el) {
promise = std::move(el.promise);
});
m.erase(it);
}
}
if(promise.is_valid()) {
promise.set_value(cl::ftuple(pipe_op_res::timeout, cl::buffer{}));
}
}
if(promise.is_valid()) {
promise.set_value(cl::ftuple(pipe_op_res::timeout, cl::buffer{}));
}
}
};
};

timer_id = m_timer.start(std::move(timeout_callback), *timeout);
}

timer_id = m_timer.start(std::move(timeout_callback), *timeout);
read_promises->push_back({ id, timer_id, std::move(promise) });
}
}

read_promises->push_back({ id, timer_id, std::move(promise) });
if(result) {
promise.set_value(cl::ftuple(*result, std::move(result_buffer)));
}
}
}
Loading