Skip to content

Commit e27a7e9

Browse files
author
huangjun
committed
Fix selective_channel late SubDone retry after EndRPC (#3358)
Ignore late selective_channel SubDone callbacks once the main RPC enters EndRPC, and keep a defensive balancer null check in Sender::IssueRPC. Add a regression test covering timeout plus delayed sub-call completion.
1 parent 35682ff commit e27a7e9

4 files changed

Lines changed: 80 additions & 2 deletions

File tree

‎src/brpc/controller.cpp‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -625,6 +625,16 @@ void Controller::OnVersionedRPCReturned(const CompletionInfo& info,
625625
return;
626626
}
627627

628+
if (is_ending_rpc()) {
629+
// SelectiveChannel may still deliver late SubDone callbacks after the
630+
// main RPC has entered EndRPC(). Ignore those callbacks instead of
631+
// letting them re-enter retry/backup on partially torn-down state.
632+
_error_code = saved_error;
633+
response_attachment().clear();
634+
CHECK_EQ(0, bthread_id_unlock(info.id));
635+
return;
636+
}
637+
628638
if ((!_error_code && _retry_policy == NULL) ||
629639
_current_call.nretry >= _max_retry) {
630640
goto END_OF_RPC;
@@ -881,6 +891,8 @@ void Controller::Call::OnComplete(
881891
}
882892

883893
void Controller::EndRPC(const CompletionInfo& info) {
894+
add_flag(FLAGS_ENDING_RPC);
895+
884896
if (_timeout_id != 0) {
885897
bthread_timer_del(_timeout_id);
886898
_timeout_id = 0;

‎src/brpc/controller.h‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -152,6 +152,7 @@ friend void policy::ProcessThriftRequest(InputMessageBase*);
152152
static const uint32_t FLAGS_PB_SINGLE_REPEATED_TO_ARRAY = (1 << 20);
153153
static const uint32_t FLAGS_MANAGE_HTTP_BODY_ON_ERROR = (1 << 21);
154154
static const uint32_t FLAGS_WRITE_TO_SOCKET_IN_BACKGROUND = (1 << 22);
155+
static const uint32_t FLAGS_ENDING_RPC = (1 << 23);
155156

156157
public:
157158
struct Inheritable {
@@ -796,6 +797,8 @@ friend void policy::ProcessThriftRequest(InputMessageBase*);
796797
return has_flag(FLAGS_ENABLED_CIRCUIT_BREAKER);
797798
}
798799

800+
bool is_ending_rpc() const { return has_flag(FLAGS_ENDING_RPC); }
801+
799802
std::string& protocol_param() { return _thrift_method_name; }
800803
const std::string& protocol_param() const { return _thrift_method_name; }
801804

‎src/brpc/selective_channel.cpp‎

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -306,14 +306,20 @@ Sender::Sender(Controller* cntl,
306306

307307
int Sender::IssueRPC(int64_t start_realtime_us) {
308308
_main_cntl->_current_call.need_feedback = false;
309+
ChannelBalancer* balancer =
310+
static_cast<ChannelBalancer*>(_main_cntl->_lb.get());
311+
if (balancer == NULL) {
312+
_main_cntl->SetFailed(ECANCELED,
313+
"SelectiveChannel balancer is unavailable");
314+
return -1;
315+
}
309316
LoadBalancer::SelectIn sel_in = { start_realtime_us,
310317
true,
311318
_main_cntl->has_request_code(),
312319
_main_cntl->_request_code,
313320
_main_cntl->_accessed };
314321
ChannelBalancer::SelectOut sel_out;
315-
const int rc = static_cast<ChannelBalancer*>(_main_cntl->_lb.get())
316-
->SelectChannel(sel_in, &sel_out);
322+
const int rc = balancer->SelectChannel(sel_in, &sel_out);
317323
if (rc != 0) {
318324
_main_cntl->SetFailed(rc, "Fail to select channel, %s", berror(rc));
319325
return -1;

‎test/brpc_channel_unittest.cpp‎

Lines changed: 57 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -214,6 +214,21 @@ class MyEchoService : public ::test::EchoService {
214214
std::function<MockFuncType> mockfunc_;
215215
};
216216

217+
class DelayedCloseEchoService : public ::test::EchoService {
218+
public:
219+
void Echo(google::protobuf::RpcController* cntl_base,
220+
const ::test::EchoRequest* req,
221+
::test::EchoResponse*,
222+
google::protobuf::Closure* done) override {
223+
brpc::ClosureGuard done_guard(done);
224+
if (req->sleep_us() > 0) {
225+
bthread_usleep(req->sleep_us());
226+
}
227+
static_cast<brpc::Controller*>(cntl_base)->CloseConnection(
228+
"Close connection after delay");
229+
}
230+
};
231+
217232
pthread_once_t register_mock_protocol = PTHREAD_ONCE_INIT;
218233

219234
class ChannelTest : public ::testing::Test{
@@ -2938,6 +2953,48 @@ TEST_F(ChannelTest, backup_request_policy) {
29382953
}
29392954
}
29402955

2956+
TEST_F(ChannelTest, selective_channel_ignores_late_subdone_after_timeout) {
2957+
DelayedCloseEchoService service;
2958+
brpc::Server server;
2959+
ASSERT_EQ(0, server.AddService(&service, brpc::SERVER_DOESNT_OWN_SERVICE));
2960+
ASSERT_EQ(0, server.Start("127.0.0.1:0", NULL));
2961+
2962+
brpc::SelectiveChannel channel;
2963+
ASSERT_EQ(0, channel.Init("rr", NULL));
2964+
2965+
brpc::ChannelOptions options;
2966+
options.timeout_ms = 100;
2967+
for (int i = 0; i < 2; ++i) {
2968+
brpc::Channel* sub_channel = new brpc::Channel;
2969+
ASSERT_EQ(0, sub_channel->Init(server.listen_address(), &options));
2970+
ASSERT_EQ(0, channel.AddChannel(sub_channel, NULL));
2971+
}
2972+
2973+
brpc::Controller cntl;
2974+
cntl.set_max_retry(3);
2975+
cntl.set_backup_request_ms(1);
2976+
cntl.set_timeout_ms(10);
2977+
2978+
test::EchoRequest req;
2979+
test::EchoResponse res;
2980+
req.set_message(__FUNCTION__);
2981+
req.set_sleep_us(50000);
2982+
CallMethod(&channel, &cntl, &req, &res, true);
2983+
2984+
ASSERT_EQ(brpc::ERPCTIMEDOUT, cntl.ErrorCode()) << cntl.ErrorText();
2985+
ASSERT_TRUE(cntl.has_backup_request());
2986+
ASSERT_GE(cntl.retried_count(), 1);
2987+
2988+
// Let the delayed sub-calls close their connections after the main RPC
2989+
// has already timed out. This used to re-enter retry/backup from
2990+
// SubDone::Run() on a partially torn-down controller.
2991+
bthread_usleep(120000);
2992+
2993+
EXPECT_EQ(brpc::ERPCTIMEDOUT, cntl.ErrorCode()) << cntl.ErrorText();
2994+
server.Stop(0);
2995+
server.Join();
2996+
}
2997+
29412998
TEST_F(ChannelTest, multiple_threads_single_channel) {
29422999
srand(time(NULL));
29433000
ASSERT_EQ(0, StartAccept(_ep));

0 commit comments

Comments
 (0)