MCPcopy Create free account
hub / github.com/AlexInLog/ReactivePlusPlus / Service

Class Service

src/examples/rppgrpc/communication/server.cpp:10–51  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

8#include "protocol.pb.h"
9
10class Service : public TestService::CallbackService
11{
12public:
13 Service()
14 {
15 client_side_requests.get_observable().subscribe([](const Request& s) { std::cout << "[ClientSideRequest]: " << s.ShortDebugString() << std::endl; });
16 }
17
18 grpc::ServerBidiReactor<::Request, ::Response>* Bidirectional(::grpc::CallbackServerContext* /*context*/) override
19 {
20 rpp::subjects::publish_subject<Response> response{};
21 rpp::subjects::publish_subject<Request> request{};
22 request.get_observable()
23 | rpp::ops::subscribe([](const Request& s) { std::cout << "[BidireactionalRequest]: " << s.ShortDebugString() << std::endl; });
24 request.get_observable()
25 | rpp::ops::map([](const Request& request) {
26 Response response{};
27 response.set_value(std::string{"BidiResponse "} + request.value());
28 return response;
29 })
30 | rpp::ops::subscribe(response.get_observer());
31 return rppgrpc::make_server_reactor(response.get_observable(), request.get_observer());
32 }
33
34 ::grpc::ServerReadReactor<::Request>* ClientSide(::grpc::CallbackServerContext* /*context*/, ::Response* /*response*/) override
35 {
36 return rppgrpc::make_server_reactor(client_side_requests.get_observer());
37 }
38
39 ::grpc::ServerWriteReactor<::Response>* ServerSide(::grpc::CallbackServerContext* /*context*/, const ::Request* /*request*/) override
40 {
41 return rppgrpc::make_server_reactor(client_side_requests.get_observable()
42 | rpp::ops::map([](const Request& v) {
43 Response response{};
44 response.set_value(std::string{"ServerSideResponse "} + v.value());
45 return response;
46 }));
47 }
48
49private:
50 rpp::subjects::publish_subject<Request> client_side_requests{};
51};
52
53int main()
54{

Callers

nothing calls this directly

Calls

no outgoing calls

Tested by

no test coverage detected