chore: import upstream snapshot with attribution
This commit is contained in:
@@ -0,0 +1,10 @@
|
||||
Note:
|
||||
Only thrift framed transport supported now, in another words, only working on thrift nonblocking mode.
|
||||
|
||||
summary:
|
||||
echo_client/echo_server:
|
||||
brpc + thrift protocol version
|
||||
native_client/native_server:
|
||||
native thrift cpp version
|
||||
|
||||
|
||||
Executable
+88
@@ -0,0 +1,88 @@
|
||||
// Licensed to the Apache Software Foundation (ASF) under one
|
||||
// or more contributor license agreements. See the NOTICE file
|
||||
// distributed with this work for additional information
|
||||
// regarding copyright ownership. The ASF licenses this file
|
||||
// to you under the Apache License, Version 2.0 (the
|
||||
// "License"); you may not use this file except in compliance
|
||||
// with the License. You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing,
|
||||
// software distributed under the License is distributed on an
|
||||
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
|
||||
// KIND, either express or implied. See the License for the
|
||||
// specific language governing permissions and limitations
|
||||
// under the License.
|
||||
|
||||
// A client sending thrift requests to server every 1 second.
|
||||
|
||||
#include <gflags/gflags.h>
|
||||
|
||||
#include "gen-cpp/echo_types.h"
|
||||
|
||||
#include <butil/logging.h>
|
||||
#include <butil/time.h>
|
||||
#include <brpc/channel.h>
|
||||
#include <brpc/thrift_message.h>
|
||||
#include <bvar/bvar.h>
|
||||
|
||||
bvar::LatencyRecorder g_latency_recorder("client");
|
||||
|
||||
DEFINE_string(server, "0.0.0.0:8019", "IP Address of server");
|
||||
DEFINE_string(load_balancer, "", "The algorithm for load balancing");
|
||||
DEFINE_int32(timeout_ms, 100, "RPC timeout in milliseconds");
|
||||
DEFINE_int32(max_retry, 3, "Max retries(not including the first RPC)");
|
||||
|
||||
int main(int argc, char* argv[]) {
|
||||
// Parse gflags. We recommend you to use gflags as well.
|
||||
google::ParseCommandLineFlags(&argc, &argv, true);
|
||||
|
||||
// A Channel represents a communication line to a Server. Notice that
|
||||
// Channel is thread-safe and can be shared by all threads in your program.
|
||||
brpc::Channel channel;
|
||||
|
||||
// Initialize the channel, NULL means using default options.
|
||||
brpc::ChannelOptions options;
|
||||
options.protocol = brpc::PROTOCOL_THRIFT;
|
||||
options.timeout_ms = FLAGS_timeout_ms/*milliseconds*/;
|
||||
options.max_retry = FLAGS_max_retry;
|
||||
if (channel.Init(FLAGS_server.c_str(), FLAGS_load_balancer.c_str(), &options) != 0) {
|
||||
LOG(ERROR) << "Fail to initialize channel";
|
||||
return -1;
|
||||
}
|
||||
|
||||
brpc::ThriftStub stub(&channel);
|
||||
|
||||
// Send a request and wait for the response every 1 second.
|
||||
while (!brpc::IsAskedToQuit()) {
|
||||
brpc::Controller cntl;
|
||||
example::EchoRequest req;
|
||||
example::EchoResponse res;
|
||||
|
||||
req.__set_data("hello");
|
||||
req.__set_need_by_proxy(10);
|
||||
|
||||
stub.CallMethod("Echo", &cntl, &req, &res, NULL);
|
||||
|
||||
if (cntl.Failed()) {
|
||||
LOG(ERROR) << "Fail to send thrift request, " << cntl.ErrorText();
|
||||
sleep(1); // Remove this sleep in production code.
|
||||
} else {
|
||||
g_latency_recorder << cntl.latency_us();
|
||||
LOG(INFO) << "Thrift Response: " << res;
|
||||
}
|
||||
|
||||
LOG_EVERY_SECOND(INFO)
|
||||
<< "Sending thrift requests at qps=" << g_latency_recorder.qps(1)
|
||||
<< " latency=" << g_latency_recorder.latency(1);
|
||||
|
||||
sleep(1);
|
||||
|
||||
}
|
||||
|
||||
LOG(INFO) << "EchoClient is going to quit";
|
||||
return 0;
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,147 @@
|
||||
// Licensed to the Apache Software Foundation (ASF) under one
|
||||
// or more contributor license agreements. See the NOTICE file
|
||||
// distributed with this work for additional information
|
||||
// regarding copyright ownership. The ASF licenses this file
|
||||
// to you under the Apache License, Version 2.0 (the
|
||||
// "License"); you may not use this file except in compliance
|
||||
// with the License. You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing,
|
||||
// software distributed under the License is distributed on an
|
||||
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
|
||||
// KIND, either express or implied. See the License for the
|
||||
// specific language governing permissions and limitations
|
||||
// under the License.
|
||||
|
||||
// A client sending requests to server by multiple threads.
|
||||
|
||||
#include "gen-cpp/echo_types.h"
|
||||
|
||||
#include <gflags/gflags.h>
|
||||
#include <bthread/bthread.h>
|
||||
#include <butil/logging.h>
|
||||
#include <brpc/server.h>
|
||||
#include <brpc/channel.h>
|
||||
#include <brpc/thrift_message.h>
|
||||
#include <bvar/bvar.h>
|
||||
|
||||
DEFINE_int32(thread_num, 50, "Number of threads to send requests");
|
||||
DEFINE_bool(use_bthread, false, "Use bthread to send requests");
|
||||
DEFINE_int32(request_size, 16, "Bytes of each request");
|
||||
DEFINE_string(connection_type, "", "Connection type. Available values: single, pooled, short");
|
||||
DEFINE_string(server, "0.0.0.0:8019", "IP Address of server");
|
||||
DEFINE_string(load_balancer, "", "The algorithm for load balancing");
|
||||
DEFINE_int32(timeout_ms, 100, "RPC timeout in milliseconds");
|
||||
DEFINE_int32(max_retry, 3, "Max retries(not including the first RPC)");
|
||||
DEFINE_bool(dont_fail, false, "Print fatal when some call failed");
|
||||
DEFINE_int32(dummy_port, -1, "Launch dummy server at this port");
|
||||
|
||||
std::string g_request;
|
||||
|
||||
bvar::LatencyRecorder g_latency_recorder("client");
|
||||
bvar::Adder<int> g_error_count("client_error_count");
|
||||
|
||||
static void* sender(void* arg) {
|
||||
// Normally, you should not call a Channel directly, but instead construct
|
||||
// a stub Service wrapping it. stub can be shared by all threads as well.
|
||||
brpc::ThriftStub stub(static_cast<brpc::Channel*>(arg));
|
||||
|
||||
while (!brpc::IsAskedToQuit()) {
|
||||
// We will receive response synchronously, safe to put variables
|
||||
// on stack.
|
||||
example::EchoRequest req;
|
||||
example::EchoResponse res;
|
||||
brpc::Controller cntl;
|
||||
|
||||
req.__set_data(g_request);
|
||||
req.__set_need_by_proxy(10);
|
||||
|
||||
// Because `done'(last parameter) is NULL, this function waits until
|
||||
// the response comes back or error occurs(including timedout).
|
||||
stub.CallMethod("Echo", &cntl, &req, &res, NULL);
|
||||
if (!cntl.Failed()) {
|
||||
g_latency_recorder << cntl.latency_us();
|
||||
} else {
|
||||
g_error_count << 1;
|
||||
CHECK(brpc::IsAskedToQuit() || !FLAGS_dont_fail)
|
||||
<< "error=" << cntl.ErrorText() << " latency=" << cntl.latency_us();
|
||||
// We can't connect to the server, sleep a while. Notice that this
|
||||
// is a specific sleeping to prevent this thread from spinning too
|
||||
// fast. You should continue the business logic in a production
|
||||
// server rather than sleeping.
|
||||
bthread_usleep(50000);
|
||||
}
|
||||
}
|
||||
return NULL;
|
||||
}
|
||||
|
||||
int main(int argc, char* argv[]) {
|
||||
// Parse gflags. We recommend you to use gflags as well.
|
||||
GFLAGS_NAMESPACE::ParseCommandLineFlags(&argc, &argv, true);
|
||||
|
||||
// A Channel represents a communication line to a Server. Notice that
|
||||
// Channel is thread-safe and can be shared by all threads in your program.
|
||||
brpc::Channel channel;
|
||||
|
||||
// Initialize the channel, NULL means using default options.
|
||||
brpc::ChannelOptions options;
|
||||
options.protocol = brpc::PROTOCOL_THRIFT;
|
||||
options.connection_type = FLAGS_connection_type;
|
||||
options.connect_timeout_ms = std::min(FLAGS_timeout_ms / 2, 100);
|
||||
options.timeout_ms = FLAGS_timeout_ms;
|
||||
options.max_retry = FLAGS_max_retry;
|
||||
if (channel.Init(FLAGS_server.c_str(), FLAGS_load_balancer.c_str(), &options) != 0) {
|
||||
LOG(ERROR) << "Fail to initialize channel";
|
||||
return -1;
|
||||
}
|
||||
|
||||
if (FLAGS_request_size <= 0) {
|
||||
LOG(ERROR) << "Bad request_size=" << FLAGS_request_size;
|
||||
return -1;
|
||||
}
|
||||
g_request.resize(FLAGS_request_size, 'r');
|
||||
|
||||
if (FLAGS_dummy_port >= 0) {
|
||||
brpc::StartDummyServerAt(FLAGS_dummy_port);
|
||||
}
|
||||
|
||||
std::vector<bthread_t> bids;
|
||||
std::vector<pthread_t> pids;
|
||||
if (!FLAGS_use_bthread) {
|
||||
pids.resize(FLAGS_thread_num);
|
||||
for (int i = 0; i < FLAGS_thread_num; ++i) {
|
||||
if (pthread_create(&pids[i], NULL, sender, &channel) != 0) {
|
||||
LOG(ERROR) << "Fail to create pthread";
|
||||
return -1;
|
||||
}
|
||||
}
|
||||
} else {
|
||||
bids.resize(FLAGS_thread_num);
|
||||
for (int i = 0; i < FLAGS_thread_num; ++i) {
|
||||
if (bthread_start_background(
|
||||
&bids[i], NULL, sender, &channel) != 0) {
|
||||
LOG(ERROR) << "Fail to create bthread";
|
||||
return -1;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
while (!brpc::IsAskedToQuit()) {
|
||||
sleep(1);
|
||||
LOG(INFO) << "Sending EchoRequest at qps=" << g_latency_recorder.qps(1)
|
||||
<< " latency=" << g_latency_recorder.latency(1);
|
||||
}
|
||||
|
||||
LOG(INFO) << "EchoClient is going to quit";
|
||||
for (int i = 0; i < FLAGS_thread_num; ++i) {
|
||||
if (!FLAGS_use_bthread) {
|
||||
pthread_join(pids[i], NULL);
|
||||
} else {
|
||||
bthread_join(bids[i], NULL);
|
||||
}
|
||||
}
|
||||
|
||||
return 0;
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
/*
|
||||
* Licensed to the Apache Software Foundation (ASF) under one
|
||||
* or more contributor license agreements. See the NOTICE file
|
||||
* distributed with this work for additional information
|
||||
* regarding copyright ownership. The ASF licenses this file
|
||||
* to you under the Apache License, Version 2.0 (the
|
||||
* "License"); you may not use this file except in compliance
|
||||
* with the License. You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing,
|
||||
* software distributed under the License is distributed on an
|
||||
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
|
||||
* KIND, either express or implied. See the License for the
|
||||
* specific language governing permissions and limitations
|
||||
* under the License.
|
||||
*/
|
||||
|
||||
namespace cpp example
|
||||
|
||||
struct EchoRequest {
|
||||
1: optional string data;
|
||||
2: optional i32 need_by_proxy;
|
||||
}
|
||||
|
||||
struct ProxyRequest {
|
||||
2: optional i32 need_by_proxy;
|
||||
}
|
||||
|
||||
struct EchoResponse {
|
||||
1: required string data;
|
||||
}
|
||||
|
||||
service EchoService {
|
||||
EchoResponse Echo(1:EchoRequest request);
|
||||
}
|
||||
|
||||
@@ -0,0 +1,78 @@
|
||||
// Licensed to the Apache Software Foundation (ASF) under one
|
||||
// or more contributor license agreements. See the NOTICE file
|
||||
// distributed with this work for additional information
|
||||
// regarding copyright ownership. The ASF licenses this file
|
||||
// to you under the Apache License, Version 2.0 (the
|
||||
// "License"); you may not use this file except in compliance
|
||||
// with the License. You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing,
|
||||
// software distributed under the License is distributed on an
|
||||
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
|
||||
// KIND, either express or implied. See the License for the
|
||||
// specific language governing permissions and limitations
|
||||
// under the License.
|
||||
|
||||
// A thrift client sending requests to server every 1 second.
|
||||
|
||||
#include <gflags/gflags.h>
|
||||
#include "gen-cpp/EchoService.h"
|
||||
#include "gen-cpp/echo_types.h"
|
||||
#include <thrift/transport/TSocket.h>
|
||||
#include <thrift/transport/TBufferTransports.h>
|
||||
#include <thrift/protocol/TBinaryProtocol.h>
|
||||
|
||||
#include <butil/logging.h>
|
||||
|
||||
// _THRIFT_STDCXX_H_ is defined by thrift/stdcxx.h which was added since thrift 0.11.0
|
||||
// but deprecated after 0.13.0
|
||||
#ifndef THRIFT_STDCXX
|
||||
#if defined(_THRIFT_STDCXX_H_)
|
||||
# define THRIFT_STDCXX apache::thrift::stdcxx
|
||||
#elif defined(_THRIFT_VERSION_LOWER_THAN_0_11_0_)
|
||||
# define THRIFT_STDCXX boost
|
||||
# include <boost/make_shared.hpp>
|
||||
#else
|
||||
# define THRIFT_STDCXX std
|
||||
#endif
|
||||
#endif
|
||||
|
||||
DEFINE_string(server, "0.0.0.0", "IP Address of server");
|
||||
DEFINE_int32(port, 8019, "Port of server");
|
||||
|
||||
int main(int argc, char **argv) {
|
||||
|
||||
// Parse gflags. We recommend you to use gflags as well.
|
||||
google::ParseCommandLineFlags(&argc, &argv, true);
|
||||
|
||||
THRIFT_STDCXX::shared_ptr<apache::thrift::transport::TSocket> socket(
|
||||
new apache::thrift::transport::TSocket(FLAGS_server, FLAGS_port));
|
||||
THRIFT_STDCXX::shared_ptr<apache::thrift::transport::TTransport> transport(
|
||||
new apache::thrift::transport::TFramedTransport(socket));
|
||||
THRIFT_STDCXX::shared_ptr<apache::thrift::protocol::TProtocol> protocol(
|
||||
new apache::thrift::protocol::TBinaryProtocol(transport));
|
||||
|
||||
example::EchoServiceClient client(protocol);
|
||||
transport->open();
|
||||
|
||||
example::EchoRequest req;
|
||||
req.__set_data("hello");
|
||||
req.__set_need_by_proxy(10);
|
||||
|
||||
example::EchoResponse res;
|
||||
|
||||
while (1) {
|
||||
try {
|
||||
client.Echo(res, req);
|
||||
LOG(INFO) << "Req=" << req << " Res=" << res;
|
||||
} catch (std::exception& e) {
|
||||
LOG(ERROR) << "Fail to rpc, " << e.what();
|
||||
}
|
||||
sleep(1);
|
||||
}
|
||||
transport->close();
|
||||
|
||||
return 0;
|
||||
}
|
||||
+109
@@ -0,0 +1,109 @@
|
||||
// Licensed to the Apache Software Foundation (ASF) under one
|
||||
// or more contributor license agreements. See the NOTICE file
|
||||
// distributed with this work for additional information
|
||||
// regarding copyright ownership. The ASF licenses this file
|
||||
// to you under the Apache License, Version 2.0 (the
|
||||
// "License"); you may not use this file except in compliance
|
||||
// with the License. You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing,
|
||||
// software distributed under the License is distributed on an
|
||||
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
|
||||
// KIND, either express or implied. See the License for the
|
||||
// specific language governing permissions and limitations
|
||||
// under the License.
|
||||
|
||||
// A thrift server to receive EchoRequest and send back EchoResponse.
|
||||
|
||||
#include <gflags/gflags.h>
|
||||
|
||||
#include <butil/logging.h>
|
||||
|
||||
#include "gen-cpp/EchoService.h"
|
||||
|
||||
#include <thrift/protocol/TBinaryProtocol.h>
|
||||
#include <thrift/server/TSimpleServer.h>
|
||||
#include <thrift/transport/TServerSocket.h>
|
||||
#include <thrift/transport/TTransportUtils.h>
|
||||
#include <thrift/server/TNonblockingServer.h>
|
||||
|
||||
// _THRIFT_STDCXX_H_ is defined by thrift/stdcxx.h which was added since thrift 0.11.0
|
||||
// but deprecated after 0.13.0, PosixThreadFactory was also deprecated in 0.13.0
|
||||
#include <thrift/TProcessor.h> // to include stdcxx.h if present
|
||||
#ifndef THRIFT_STDCXX
|
||||
#if defined(_THRIFT_STDCXX_H_)
|
||||
# define THRIFT_STDCXX apache::thrift::stdcxx
|
||||
#include <thrift/transport/TNonblockingServerSocket.h>
|
||||
#include <thrift/concurrency/PosixThreadFactory.h>
|
||||
#elif defined(_THRIFT_VERSION_LOWER_THAN_0_11_0_)
|
||||
# define THRIFT_STDCXX boost
|
||||
#include <boost/make_shared.hpp>
|
||||
#include <thrift/concurrency/PosixThreadFactory.h>
|
||||
#else
|
||||
# define THRIFT_STDCXX std
|
||||
#include <thrift/concurrency/ThreadFactory.h>
|
||||
#include <thrift/transport/TNonblockingServerSocket.h>
|
||||
#endif
|
||||
#endif
|
||||
|
||||
DEFINE_int32(port, 8019, "Port of server");
|
||||
|
||||
class EchoServiceHandler : virtual public example::EchoServiceIf {
|
||||
public:
|
||||
EchoServiceHandler() {}
|
||||
|
||||
void Echo(example::EchoResponse& res, const example::EchoRequest& req) {
|
||||
// Process request, just attach a simple string.
|
||||
res.data = req.data + " world";
|
||||
return;
|
||||
}
|
||||
|
||||
};
|
||||
|
||||
int main(int argc, char *argv[]) {
|
||||
// Parse gflags. We recommend you to use gflags as well.
|
||||
google::ParseCommandLineFlags(&argc, &argv, true);
|
||||
|
||||
THRIFT_STDCXX::shared_ptr<EchoServiceHandler> handler(new EchoServiceHandler());
|
||||
#if THRIFT_STDCXX != std
|
||||
// For thrift version less than 0.13.0
|
||||
THRIFT_STDCXX::shared_ptr<apache::thrift::concurrency::PosixThreadFactory> thread_factory(
|
||||
new apache::thrift::concurrency::PosixThreadFactory(
|
||||
apache::thrift::concurrency::PosixThreadFactory::ROUND_ROBIN,
|
||||
apache::thrift::concurrency::PosixThreadFactory::NORMAL, 1, false));
|
||||
#else
|
||||
// For thrift version greater equal than 0.13.0
|
||||
THRIFT_STDCXX::shared_ptr<apache::thrift::concurrency::ThreadFactory> thread_factory(
|
||||
new apache::thrift::concurrency::ThreadFactory(false));
|
||||
#endif
|
||||
|
||||
THRIFT_STDCXX::shared_ptr<apache::thrift::server::TProcessor> processor(
|
||||
new example::EchoServiceProcessor(handler));
|
||||
THRIFT_STDCXX::shared_ptr<apache::thrift::protocol::TProtocolFactory> protocol_factory(
|
||||
new apache::thrift::protocol::TBinaryProtocolFactory());
|
||||
THRIFT_STDCXX::shared_ptr<apache::thrift::transport::TTransportFactory> transport_factory(
|
||||
new apache::thrift::transport::TBufferedTransportFactory());
|
||||
THRIFT_STDCXX::shared_ptr<apache::thrift::concurrency::ThreadManager> thread_mgr(
|
||||
apache::thrift::concurrency::ThreadManager::newSimpleThreadManager(2));
|
||||
|
||||
thread_mgr->threadFactory(thread_factory);
|
||||
|
||||
thread_mgr->start();
|
||||
|
||||
#if defined(_THRIFT_STDCXX_H_) || !defined (_THRIFT_VERSION_LOWER_THAN_0_11_0_)
|
||||
THRIFT_STDCXX::shared_ptr<apache::thrift::transport::TNonblockingServerSocket> server_transport =
|
||||
THRIFT_STDCXX::make_shared<apache::thrift::transport::TNonblockingServerSocket>(FLAGS_port);
|
||||
|
||||
apache::thrift::server::TNonblockingServer server(processor,
|
||||
transport_factory, transport_factory, protocol_factory,
|
||||
protocol_factory, server_transport);
|
||||
#else
|
||||
apache::thrift::server::TNonblockingServer server(processor,
|
||||
transport_factory, transport_factory, protocol_factory,
|
||||
protocol_factory, FLAGS_port);
|
||||
#endif
|
||||
server.serve();
|
||||
return 0;
|
||||
}
|
||||
Executable
+81
@@ -0,0 +1,81 @@
|
||||
// Licensed to the Apache Software Foundation (ASF) under one
|
||||
// or more contributor license agreements. See the NOTICE file
|
||||
// distributed with this work for additional information
|
||||
// regarding copyright ownership. The ASF licenses this file
|
||||
// to you under the Apache License, Version 2.0 (the
|
||||
// "License"); you may not use this file except in compliance
|
||||
// with the License. You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing,
|
||||
// software distributed under the License is distributed on an
|
||||
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
|
||||
// KIND, either express or implied. See the License for the
|
||||
// specific language governing permissions and limitations
|
||||
// under the License.
|
||||
|
||||
// A server to receive EchoRequest and send back EchoResponse.
|
||||
|
||||
#include <gflags/gflags.h>
|
||||
#include <butil/logging.h>
|
||||
#include <brpc/server.h>
|
||||
#include <brpc/thrift_service.h>
|
||||
#include "gen-cpp/echo_types.h"
|
||||
|
||||
DEFINE_int32(port, 8019, "TCP Port of this server");
|
||||
DEFINE_int32(idle_timeout_s, -1, "Connection will be closed if there is no "
|
||||
"read/write operations during the last `idle_timeout_s'");
|
||||
DEFINE_int32(max_concurrency, 0, "Limit of request processing in parallel");
|
||||
|
||||
// Adapt your own thrift-based protocol to use brpc
|
||||
class EchoServiceImpl : public brpc::ThriftService {
|
||||
public:
|
||||
void ProcessThriftFramedRequest(brpc::Controller* cntl,
|
||||
brpc::ThriftFramedMessage* req,
|
||||
brpc::ThriftFramedMessage* res,
|
||||
google::protobuf::Closure* done) override {
|
||||
// Dispatch calls to different methods
|
||||
if (cntl->thrift_method_name() == "Echo") {
|
||||
return Echo(cntl, req->Cast<example::EchoRequest>(),
|
||||
res->Cast<example::EchoResponse>(), done);
|
||||
} else {
|
||||
cntl->SetFailed(brpc::ENOMETHOD, "Fail to find method=%s",
|
||||
cntl->thrift_method_name().c_str());
|
||||
done->Run();
|
||||
}
|
||||
}
|
||||
|
||||
void Echo(brpc::Controller* cntl,
|
||||
const example::EchoRequest* req,
|
||||
example::EchoResponse* res,
|
||||
google::protobuf::Closure* done) {
|
||||
// This object helps you to call done->Run() in RAII style. If you need
|
||||
// to process the request asynchronously, pass done_guard.release().
|
||||
brpc::ClosureGuard done_guard(done);
|
||||
|
||||
res->data = req->data + " (Echo)";
|
||||
}
|
||||
};
|
||||
|
||||
int main(int argc, char* argv[]) {
|
||||
// Parse gflags. We recommend you to use gflags as well.
|
||||
google::ParseCommandLineFlags(&argc, &argv, true);
|
||||
|
||||
brpc::Server server;
|
||||
brpc::ServerOptions options;
|
||||
|
||||
options.thrift_service = new EchoServiceImpl;
|
||||
options.idle_timeout_sec = FLAGS_idle_timeout_s;
|
||||
options.max_concurrency = FLAGS_max_concurrency;
|
||||
|
||||
// Start the server.
|
||||
if (server.Start(FLAGS_port, &options) != 0) {
|
||||
LOG(ERROR) << "Fail to start EchoServer";
|
||||
return -1;
|
||||
}
|
||||
|
||||
// Wait until Ctrl-C is pressed, then Stop() and Join() the server.
|
||||
server.RunUntilAskedToQuit();
|
||||
return 0;
|
||||
}
|
||||
Executable
+104
@@ -0,0 +1,104 @@
|
||||
// Licensed to the Apache Software Foundation (ASF) under one
|
||||
// or more contributor license agreements. See the NOTICE file
|
||||
// distributed with this work for additional information
|
||||
// regarding copyright ownership. The ASF licenses this file
|
||||
// to you under the Apache License, Version 2.0 (the
|
||||
// "License"); you may not use this file except in compliance
|
||||
// with the License. You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing,
|
||||
// software distributed under the License is distributed on an
|
||||
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
|
||||
// KIND, either express or implied. See the License for the
|
||||
// specific language governing permissions and limitations
|
||||
// under the License.
|
||||
|
||||
// A server to receive EchoRequest and send back EchoResponse.
|
||||
|
||||
#include <gflags/gflags.h>
|
||||
#include <butil/logging.h>
|
||||
#include <brpc/server.h>
|
||||
#include <brpc/thrift_message.h>
|
||||
#include <brpc/channel.h>
|
||||
#include <brpc/thrift_service.h>
|
||||
#include "gen-cpp/echo_types.h"
|
||||
|
||||
DEFINE_int32(port, 8019, "TCP Port of this server");
|
||||
DEFINE_int32(idle_timeout_s, -1, "Connection will be closed if there is no "
|
||||
"read/write operations during the last `idle_timeout_s'");
|
||||
DEFINE_int32(max_concurrency, 0, "Limit of request processing in parallel");
|
||||
|
||||
// Adapt your own thrift-based protocol to use brpc
|
||||
class EchoServiceImpl : public brpc::ThriftService {
|
||||
public:
|
||||
EchoServiceImpl() {
|
||||
// Initialize the channel, NULL means using default options.
|
||||
brpc::ChannelOptions options;
|
||||
options.protocol = brpc::PROTOCOL_THRIFT;
|
||||
if (_channel.Init("0.0.0.0", FLAGS_port , &options) != 0) {
|
||||
LOG(ERROR) << "Fail to initialize channel";
|
||||
}
|
||||
}
|
||||
|
||||
void ProcessThriftFramedRequest(brpc::Controller* cntl,
|
||||
brpc::ThriftFramedMessage* req,
|
||||
brpc::ThriftFramedMessage* res,
|
||||
google::protobuf::Closure* done) override {
|
||||
// Dispatch calls to different methods
|
||||
if (cntl->thrift_method_name() == "Echo") {
|
||||
// Proxy request/response to RealEcho, note that as a proxy we
|
||||
// don't need to Cast the messages to native types.
|
||||
brpc::Controller cntl;
|
||||
brpc::ThriftStub stub(&_channel);
|
||||
// TODO: Following Cast<> drops data field from ProxyRequest which
|
||||
// does not recognize the field, should be debugged further.
|
||||
// LOG(INFO) << "req=" << *req->Cast<example::ProxyRequest>();
|
||||
stub.CallMethod("RealEcho", &cntl, req, res, NULL);
|
||||
done->Run();
|
||||
} else if (cntl->thrift_method_name() == "RealEcho") {
|
||||
return RealEcho(cntl, req->Cast<example::EchoRequest>(),
|
||||
res->Cast<example::EchoResponse>(), done);
|
||||
} else {
|
||||
cntl->SetFailed(brpc::ENOMETHOD, "Fail to find method=%s",
|
||||
cntl->thrift_method_name().c_str());
|
||||
done->Run();
|
||||
}
|
||||
}
|
||||
|
||||
void RealEcho(brpc::Controller* cntl,
|
||||
const example::EchoRequest* req,
|
||||
example::EchoResponse* res,
|
||||
google::protobuf::Closure* done) {
|
||||
// This object helps you to call done->Run() in RAII style. If you need
|
||||
// to process the request asynchronously, pass done_guard.release().
|
||||
brpc::ClosureGuard done_guard(done);
|
||||
|
||||
res->data = req->data + " (RealEcho)";
|
||||
}
|
||||
private:
|
||||
brpc::Channel _channel;
|
||||
};
|
||||
|
||||
int main(int argc, char* argv[]) {
|
||||
// Parse gflags. We recommend you to use gflags as well.
|
||||
google::ParseCommandLineFlags(&argc, &argv, true);
|
||||
|
||||
brpc::Server server;
|
||||
brpc::ServerOptions options;
|
||||
|
||||
options.thrift_service = new EchoServiceImpl;
|
||||
options.idle_timeout_sec = FLAGS_idle_timeout_s;
|
||||
options.max_concurrency = FLAGS_max_concurrency;
|
||||
|
||||
// Start the server.
|
||||
if (server.Start(FLAGS_port, &options) != 0) {
|
||||
LOG(ERROR) << "Fail to start EchoServer";
|
||||
return -1;
|
||||
}
|
||||
|
||||
// Wait until Ctrl-C is pressed, then Stop() and Join() the server.
|
||||
server.RunUntilAskedToQuit();
|
||||
return 0;
|
||||
}
|
||||
Reference in New Issue
Block a user