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
14 changes: 7 additions & 7 deletions examples/cpp/rpc/ServerApp.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -136,7 +136,7 @@ void ServerApp::create_server(
}

calculator_example::detail::Calculator_representation_limits_Out ServerApp::ServerImpl::representation_limits(
const CalculatorServer_ClientContext& info)
const eprosima::fastdds::dds::rpc::RpcRequest& info)
{
static_cast<void>(info);
calculator_example::detail::Calculator_representation_limits_Out limits;
Expand All @@ -146,7 +146,7 @@ calculator_example::detail::Calculator_representation_limits_Out ServerApp::Serv
}

int32_t ServerApp::ServerImpl::addition(
const CalculatorServer_ClientContext& info,
const eprosima::fastdds::dds::rpc::RpcRequest& info,
/*in*/ int32_t value1,
/*in*/ int32_t value2)
{
Expand All @@ -166,7 +166,7 @@ int32_t ServerApp::ServerImpl::addition(
}

int32_t ServerApp::ServerImpl::subtraction(
const CalculatorServer_ClientContext& info,
const eprosima::fastdds::dds::rpc::RpcRequest& info,
/*in*/ int32_t value1,
/*in*/ int32_t value2)
{
Expand All @@ -186,7 +186,7 @@ int32_t ServerApp::ServerImpl::subtraction(
}

void ServerApp::ServerImpl::fibonacci_seq(
const CalculatorServer_ClientContext& info,
const eprosima::fastdds::dds::rpc::RpcRequest& info,
/*in*/ uint32_t n_results,
/*result*/ eprosima::fastdds::dds::rpc::RpcServerWriter<int32_t>& result_writer)
{
Expand All @@ -213,7 +213,7 @@ void ServerApp::ServerImpl::fibonacci_seq(
}

int32_t ServerApp::ServerImpl::sum_all(
const CalculatorServer_ClientContext& info,
const eprosima::fastdds::dds::rpc::RpcRequest& info,
/*in*/ eprosima::fastdds::dds::rpc::RpcServerReader<int32_t>& value)
{
static_cast<void>(info);
Expand Down Expand Up @@ -246,7 +246,7 @@ int32_t ServerApp::ServerImpl::sum_all(
}

void ServerApp::ServerImpl::accumulator(
const CalculatorServer_ClientContext& info,
const eprosima::fastdds::dds::rpc::RpcRequest& info,
/*in*/ eprosima::fastdds::dds::rpc::RpcServerReader<int32_t>& value,
/*result*/ eprosima::fastdds::dds::rpc::RpcServerWriter<int32_t>& result_writer)
{
Expand Down Expand Up @@ -279,7 +279,7 @@ void ServerApp::ServerImpl::accumulator(
}

void ServerApp::ServerImpl::filter(
const CalculatorServer_ClientContext& info,
const eprosima::fastdds::dds::rpc::RpcRequest& info,
/*in*/ eprosima::fastdds::dds::rpc::RpcServerReader<int32_t>& value,
/*in*/ calculator_example::FilterKind filter_kind,
/*result*/ eprosima::fastdds::dds::rpc::RpcServerWriter<int32_t>& result_writer)
Expand Down
16 changes: 8 additions & 8 deletions examples/cpp/rpc/ServerApp.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -80,42 +80,42 @@ class ServerApp : public Application
explicit ServerImpl() = default;

calculator_example::detail::Calculator_representation_limits_Out representation_limits(
const calculator_example::CalculatorServer_ClientContext& info) override;
const eprosima::fastdds::dds::rpc::RpcRequest& info) override;

int32_t addition(
const calculator_example::CalculatorServer_ClientContext& info,
const eprosima::fastdds::dds::rpc::RpcRequest& info,
/*in*/ int32_t value1,
/*in*/ int32_t value2) override;

int32_t subtraction(
const calculator_example::CalculatorServer_ClientContext& info,
const eprosima::fastdds::dds::rpc::RpcRequest& info,
/*in*/ int32_t value1,
/*in*/ int32_t value2) override;

void fibonacci_seq(
const calculator_example::CalculatorServer_ClientContext& info,
const eprosima::fastdds::dds::rpc::RpcRequest& info,
/*in*/ uint32_t n_results,
/*result*/ eprosima::fastdds::dds::rpc::RpcServerWriter<int32_t>& result_writer) override;

int32_t sum_all(
const calculator_example::CalculatorServer_ClientContext& info,
const eprosima::fastdds::dds::rpc::RpcRequest& info,
/*in*/ eprosima::fastdds::dds::rpc::RpcServerReader<int32_t>& value) override;

void accumulator(
const calculator_example::CalculatorServer_ClientContext& info,
const eprosima::fastdds::dds::rpc::RpcRequest& info,
/*in*/ eprosima::fastdds::dds::rpc::RpcServerReader<int32_t>& value,
/*result*/ eprosima::fastdds::dds::rpc::RpcServerWriter<int32_t>& result_writer) override;

void filter(
const calculator_example::CalculatorServer_ClientContext& info,
const eprosima::fastdds::dds::rpc::RpcRequest& info,
/*in*/ eprosima::fastdds::dds::rpc::RpcServerReader<int32_t>& value,
/*in*/ calculator_example::FilterKind filter_kind,
/*result*/ eprosima::fastdds::dds::rpc::RpcServerWriter<int32_t>& result_writer) override;

};

std::shared_ptr<ServerImpl> server_impl_;
std::shared_ptr<calculator_example::CalculatorServer> server_;
std::shared_ptr<eprosima::fastdds::dds::rpc::RpcServer> server_;
dds::DomainParticipant* participant_;
size_t thread_pool_size_;
std::atomic<bool> stop_;
Expand Down
2 changes: 2 additions & 0 deletions examples/cpp/rpc/types/calculatorClient.cxx
Original file line number Diff line number Diff line change
Expand Up @@ -205,6 +205,7 @@ class CalculatorClient : public Calculator
std::shared_ptr<T> result,
std::promise<TResult>& promise)
{
result->info.related_sample_identity.writer_guid(requester_->get_requester_reader()->guid());
if (fdds::RETCODE_OK == requester_->send_request((void*)&request, result->info))
{
std::lock_guard<std::mutex> _ (mtx_);
Expand All @@ -224,6 +225,7 @@ class CalculatorClient : public Calculator
const RequestType& request,
std::shared_ptr<T> result)
{
result->info.related_sample_identity.writer_guid(requester_->get_requester_reader()->guid());
if (fdds::RETCODE_OK == requester_->send_request((void*)&request, result->info))
{
std::lock_guard<std::mutex> _ (mtx_);
Expand Down
10 changes: 10 additions & 0 deletions examples/cpp/rpc/types/calculatorPubSubTypes.cxx
Original file line number Diff line number Diff line change
Expand Up @@ -175,6 +175,11 @@ namespace calculator_example {
delete pData;
}

eProsima_user_DllExport void register_type_object_representation() override
{
register_Calculator_Request_type_identifier(type_identifiers_);
}

};

class Calculator_ReplyPubSubType : public eprosima::fastdds::dds::TopicDataType
Expand Down Expand Up @@ -318,6 +323,11 @@ namespace calculator_example {
delete pData;
}

eProsima_user_DllExport void register_type_object_representation() override
{
register_Calculator_Reply_type_identifier(type_identifiers_);
}

};

eprosima::fastdds::dds::rpc::ServiceTypeSupport create_Calculator_service_type_support()
Expand Down
122 changes: 104 additions & 18 deletions examples/cpp/rpc/types/calculatorServer.cxx
Original file line number Diff line number Diff line change
Expand Up @@ -63,22 +63,38 @@ namespace frpc = eprosima::fastdds::dds::rpc;
namespace frtps = eprosima::fastdds::rtps;

class CalculatorServerLogic
: public CalculatorServer
: public frpc::RpcServer
, public std::enable_shared_from_this<CalculatorServerLogic>
{
using RequestType = Calculator_Request;
using ReplyType = Calculator_Reply;

public:

CalculatorServerLogic(
eprosima::fastdds::dds::DomainParticipant& part,
fdds::DomainParticipant& part,
const char* service_name,
const eprosima::fastdds::dds::ReplierQos& qos,
const fdds::ReplierQos& qos,
size_t thread_pool_size,
std::shared_ptr<CalculatorServer_IServerImplementation> implementation)
: CalculatorServer()
: CalculatorServerLogic(
part,
service_name,
qos,
std::make_shared<ThreadPool>(*this, thread_pool_size),
std::move(implementation))
{
}

CalculatorServerLogic(
fdds::DomainParticipant& part,
const char* service_name,
const fdds::ReplierQos& qos,
std::shared_ptr<frpc::RpcServerSchedulingStrategy> scheduler,
std::shared_ptr<CalculatorServer_IServerImplementation> implementation)
: frpc::RpcServer()
, participant_(part)
, thread_pool_(*this, thread_pool_size)
, request_scheduler_(scheduler)
, implementation_(std::move(implementation))
{
// Register the service type support
Expand Down Expand Up @@ -106,8 +122,6 @@ class CalculatorServerLogic

~CalculatorServerLogic() override
{
stop();

if (nullptr != replier_)
{
participant_.delete_service_replier(service_->get_service_name(), replier_);
Expand Down Expand Up @@ -174,7 +188,21 @@ class CalculatorServerLogic
}

// Wait for all threads to finish
thread_pool_.stop();
request_scheduler_->server_stopped(shared_from_this());
}

void execute_request(
const std::shared_ptr<frpc::RpcRequest>& request) override
{
auto ctx = std::dynamic_pointer_cast<RequestContext>(request);
if (ctx)
{
execute_request(ctx);
}
else
{
throw std::runtime_error("Invalid request context type");
}
}

private:
Expand Down Expand Up @@ -724,7 +752,7 @@ class CalculatorServerLogic

//} operation filter

struct RequestContext : CalculatorServer_ClientContext
struct RequestContext : frpc::RpcRequest
{
RequestType request;
frpc::RequestInfo info;
Expand Down Expand Up @@ -999,6 +1027,7 @@ class CalculatorServerLogic
};

struct ThreadPool
: public frpc::RpcServerSchedulingStrategy
{
ThreadPool(
CalculatorServerLogic& server,
Expand All @@ -1015,7 +1044,7 @@ class CalculatorServerLogic
{
while (!finished_)
{
std::shared_ptr<RequestContext> req;
std::shared_ptr<frpc::RpcRequest> req;
{
std::unique_lock<std::mutex> lock(mtx_);
cv_.wait(lock, [this]()
Expand All @@ -1041,9 +1070,12 @@ class CalculatorServerLogic
}
}

void new_request(
const std::shared_ptr<RequestContext>& req)
void schedule_request(
const std::shared_ptr<frpc::RpcRequest>& req,
const std::shared_ptr<frpc::RpcServer>& server) override
{
static_cast<void>(server);

std::lock_guard<std::mutex> lock(mtx_);
if (!finished_)
{
Expand All @@ -1052,8 +1084,11 @@ class CalculatorServerLogic
}
}

void stop()
void server_stopped(
const std::shared_ptr<frpc::RpcServer>& server) override
{
static_cast<void>(server);

// Notify all threads in the pool to stop
{
std::lock_guard<std::mutex> lock(mtx_);
Expand All @@ -1075,7 +1110,7 @@ class CalculatorServerLogic
CalculatorServerLogic& server_;
std::mutex mtx_;
std::condition_variable cv_;
std::queue<std::shared_ptr<RequestContext>> requests_;
std::queue<std::shared_ptr<frpc::RpcRequest>> requests_;
bool finished_{ false };
std::vector<std::thread> threads_;
};
Expand Down Expand Up @@ -1107,7 +1142,7 @@ class CalculatorServerLogic
processing_requests_[id] = ctx;
}

thread_pool_.new_request(ctx);
request_scheduler_->schedule_request(ctx, shared_from_this());
}

void execute_request(
Expand Down Expand Up @@ -1312,22 +1347,73 @@ class CalculatorServerLogic
fdds::GuardCondition finish_condition_;
std::mutex mtx_;
std::map<frtps::SampleIdentity, std::shared_ptr<RequestContext>> processing_requests_;
ThreadPool thread_pool_;
std::shared_ptr<frpc::RpcServerSchedulingStrategy> request_scheduler_;
std::shared_ptr<CalculatorServer_IServerImplementation> implementation_;

};

struct CalculatorServerProxy
: public frpc::RpcServer
{
CalculatorServerProxy(
std::shared_ptr<frpc::RpcServer> impl)
: impl_(std::move(impl))
{
}

~CalculatorServerProxy() override
{
if (impl_)
{
impl_->stop();
}
}

void run() override
{
impl_->run();
}

void stop() override
{
impl_->stop();
}

void execute_request(
const std::shared_ptr<frpc::RpcRequest>& request) override
{
impl_->execute_request(request);
}

private:

std::shared_ptr<frpc::RpcServer> impl_;
};

} // namespace detail

std::shared_ptr<CalculatorServer> create_CalculatorServer(
std::shared_ptr<eprosima::fastdds::dds::rpc::RpcServer> create_CalculatorServer(
eprosima::fastdds::dds::DomainParticipant& part,
const char* service_name,
const eprosima::fastdds::dds::ReplierQos& qos,
size_t thread_pool_size,
std::shared_ptr<CalculatorServer_IServerImplementation> implementation)
{
return std::make_shared<detail::CalculatorServerLogic>(
auto ptr = std::make_shared<detail::CalculatorServerLogic>(
part, service_name, qos, thread_pool_size, implementation);
return std::make_shared<detail::CalculatorServerProxy>(ptr);
}

std::shared_ptr<eprosima::fastdds::dds::rpc::RpcServer> create_CalculatorServer(
eprosima::fastdds::dds::DomainParticipant& part,
const char* service_name,
const eprosima::fastdds::dds::ReplierQos& qos,
std::shared_ptr<eprosima::fastdds::dds::rpc::RpcServerSchedulingStrategy> scheduler,
std::shared_ptr<CalculatorServer_IServerImplementation> implementation)
{
auto ptr = std::make_shared<detail::CalculatorServerLogic>(
part, service_name, qos, scheduler, implementation);
return std::make_shared<detail::CalculatorServerProxy>(ptr);
}

//} interface Calculator
Expand Down
Loading
Loading