-
Notifications
You must be signed in to change notification settings - Fork 0
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Quick integration with gRPC using protobuf
- Loading branch information
Showing
18 changed files
with
199 additions
and
340 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,2 @@ | ||
*.pb.h | ||
*.pb.cc |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1 @@ | ||
protoc -I. --cpp_out=. --grpc_out=. --plugin=protoc-gen-grpc=/usr/local/bin/grpc_cpp_plugin yess.proto |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,37 @@ | ||
syntax = "proto3"; | ||
|
||
package yess; | ||
|
||
message Event { | ||
int32 id = 1; | ||
int32 stream_id = 2; | ||
string type = 3; | ||
string payload = 4; | ||
int32 version = 5; | ||
} | ||
|
||
message Stream { | ||
int32 id = 1; | ||
string type = 2; | ||
int32 version = 3; | ||
repeated Event events = 4; | ||
} | ||
|
||
message Response { | ||
int32 status = 1; | ||
string msg = 2; | ||
} | ||
|
||
message CreateStreamReq { | ||
string type = 1; | ||
} | ||
|
||
message PushEventReq { | ||
int32 stream_id = 1; | ||
Event event = 2; | ||
} | ||
|
||
service YessService { | ||
rpc CreateStream (CreateStreamReq) returns (Response); | ||
rpc PushEvent (PushEventReq) returns (Response); | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,90 @@ | ||
#pragma once | ||
|
||
#include <iostream> | ||
#include <memory> | ||
#include <string> | ||
|
||
#include <grpcpp/ext/proto_server_reflection_plugin.h> | ||
#include <grpcpp/grpcpp.h> | ||
#include <grpcpp/health_check_service_interface.h> | ||
|
||
#include "../protos/yess.grpc.pb.h" | ||
#include "action_handler.hpp" | ||
#include "msg/response.hpp" | ||
#include "nlohmann/json.hpp" | ||
#include "yess.pb.h" | ||
|
||
using grpc::Server; | ||
using grpc::ServerBuilder; | ||
using grpc::ServerContext; | ||
using grpc::Status; | ||
using yess::CreateStreamReq; | ||
using yess::PushEventReq; | ||
using yess::Response; | ||
using yess::YessService; | ||
using json = nlohmann::json; | ||
|
||
namespace yess | ||
{ | ||
class YessServiceImpl : public YessService::Service | ||
{ | ||
public: | ||
YessServiceImpl(std::string conn_str) | ||
: handler_(std::make_unique<yess::Action_handler>(conn_str)) | ||
{ | ||
} | ||
|
||
private: | ||
Status CreateStream(ServerContext *context, | ||
const CreateStreamReq *request, | ||
Response *reply) override | ||
{ | ||
json cmd = {{"action", "CreateStream"}, {"type", request->type()}}; | ||
|
||
msg::Response resp = handler_->handle(cmd); | ||
reply->set_msg(resp.getMessage()); | ||
reply->set_status((int)resp.getStatus()); | ||
return Status::OK; | ||
} | ||
|
||
Status PushEvent(ServerContext *context, | ||
const PushEventReq *request, | ||
Response *reply) override | ||
{ | ||
|
||
json cmd = {{"action", "PushEvent"}, | ||
{"streamId", request->stream_id()}, | ||
{"type", request->event().type()}, | ||
{"payload", request->event().payload()}, | ||
{"version", request->event().version()}}; | ||
|
||
msg::Response resp = handler_->handle(cmd); | ||
reply->set_msg(resp.getMessage()); | ||
reply->set_status((int)resp.getStatus()); | ||
return Status::OK; | ||
} | ||
std::unique_ptr<yess::Action_handler> handler_; | ||
}; | ||
|
||
void run_server() | ||
{ | ||
std::string server_address("0.0.0.0:2929"); | ||
YessServiceImpl service("yess.db"); | ||
|
||
grpc::EnableDefaultHealthCheckService(true); | ||
grpc::reflection::InitProtoReflectionServerBuilderPlugin(); | ||
ServerBuilder builder; | ||
// Listen on the given address without any authentication mechanism. | ||
builder.AddListeningPort(server_address, grpc::InsecureServerCredentials()); | ||
// Register "service" as the instance through which we'll communicate with | ||
// clients. In this case it corresponds to an *synchronous* service. | ||
builder.RegisterService(&service); | ||
// Finally assemble the server. | ||
std::unique_ptr<Server> server(builder.BuildAndStart()); | ||
std::cout << "Server listening on " << server_address << std::endl; | ||
|
||
// Wait for the server to shutdown. Note that some other thread must be | ||
// responsible for shutting down the server for this call to ever return. | ||
server->Wait(); | ||
} | ||
} // namespace yess |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,7 +1,9 @@ | ||
#include "server.hpp" | ||
#include "grpc_service.hpp" | ||
|
||
int main(int argc, char **argv) | ||
{ | ||
yess::Server server; | ||
return server.run(argc, argv); | ||
yess::run_server(); | ||
|
||
return 0; | ||
} | ||
|
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.