Skip to content

Commit cbf6fbb

Browse files
authored
Fix: Adjust FiberTransport to lowdown latency in busrt tcp connect scenario (#66)
- Do not capture proxy in version report function - Lowdown log level when TcpConnection ReadIoData error
1 parent 3d299f7 commit cbf6fbb

6 files changed

Lines changed: 181 additions & 98 deletions

File tree

trpc/runtime/iomodel/reactor/default/tcp_connection.cc

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -493,9 +493,9 @@ int TcpConnection::ReadIoData(NoncontiguousBuffer& buff) {
493493
continue;
494494
}
495495
ret = n;
496-
TRPC_LOG_ERROR("TcpConnection::ReadIoData fd:" << socket_.GetFd() << ", ip:" << GetPeerIp()
497-
<< ", port:" << GetPeerPort() << ", is_client:" << IsClient()
498-
<< ", errno:" << errno << ", read failed and connection close.");
496+
TRPC_LOG_WARN("TcpConnection::ReadIoData fd:" << socket_.GetFd() << ", ip:" << GetPeerIp()
497+
<< ", port:" << GetPeerPort() << ", is_client:" << IsClient()
498+
<< ", errno:" << errno << ", read failed and connection close.");
499499
HandleClose(true);
500500
}
501501
break;

trpc/transport/client/fiber/BUILD

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,7 @@ cc_test(
5555
"//trpc/common/future:future_utility",
5656
"//trpc/coroutine:fiber",
5757
"//trpc/coroutine:future",
58+
"//trpc/stream:stream_handler",
5859
"//trpc/transport/client/fiber/testing:fake_server",
5960
"//trpc/transport/client/fiber/testing:thread_model_op",
6061
"@com_google_googletest//:gtest_main",

trpc/transport/client/fiber/fiber_connector_group_manager.cc

Lines changed: 52 additions & 67 deletions
Original file line numberDiff line numberDiff line change
@@ -27,33 +27,31 @@
2727
namespace trpc {
2828

2929
FiberConnectorGroupManager::FiberConnectorGroupManager(TransInfo&& trans_info) : trans_info_(std::move(trans_info)) {
30-
tcp_impl_.store(std::make_unique<TcpImpl>().release());
3130
udp_impl_.store(std::make_unique<UdpImpl>().release());
3231

33-
fiber_transport_state_.store(ClientTransportState::kInitialized, std::memory_order_release);
32+
fiber_transport_state_.store(ClientTransportState::kInitialized);
3433
}
3534

3635
FiberConnectorGroupManager::~FiberConnectorGroupManager() { Stop(); }
3736

3837
void FiberConnectorGroupManager::Stop() {
39-
std::scoped_lock _(mutex_);
40-
41-
if (fiber_transport_state_.load(std::memory_order_acquire) != ClientTransportState::kInitialized) {
38+
ClientTransportState old_value = ClientTransportState::kInitialized;
39+
ClientTransportState new_value = ClientTransportState::kStopped;
40+
if (!fiber_transport_state_.compare_exchange_strong(old_value, new_value)) {
4241
return;
4342
}
4443

4544
{
46-
Hazptr hazptr;
47-
auto tcp_impl = hazptr.Keep(&tcp_impl_);
48-
if (tcp_impl) {
49-
for (auto&& [name, group] : tcp_impl->tcp_connector_groups) {
50-
(void)name;
51-
group->Stop();
52-
}
53-
}
45+
std::unique_lock lock(tcp_mutex_);
46+
tcp_connector_groups_.swap(tcp_connector_groups_to_destroy_);
47+
}
48+
49+
for (auto& group : tcp_connector_groups_to_destroy_) {
50+
group.second->Stop();
5451
}
5552

5653
{
54+
std::scoped_lock _(udp_mutex_);
5755
Hazptr hazptr;
5856
auto udp_impl = hazptr.Keep(&udp_impl_);
5957
if (udp_impl) {
@@ -66,47 +64,45 @@ void FiberConnectorGroupManager::Stop() {
6664
}
6765
}
6866
}
69-
70-
fiber_transport_state_.store(ClientTransportState::kStopped, std::memory_order_release);
7167
}
7268

7369
void FiberConnectorGroupManager::Destroy() {
74-
std::scoped_lock _(mutex_);
75-
76-
if (fiber_transport_state_.load(std::memory_order_acquire) != ClientTransportState::kStopped) {
70+
ClientTransportState old_value = ClientTransportState::kStopped;
71+
ClientTransportState new_value = ClientTransportState::kDestroyed;
72+
if (!fiber_transport_state_.compare_exchange_strong(old_value, new_value)) {
7773
return;
7874
}
7975

80-
auto tcp_impl = tcp_impl_.exchange(nullptr, std::memory_order_relaxed);
81-
if (tcp_impl) {
82-
for (auto&& [name, group] : tcp_impl->tcp_connector_groups) {
83-
(void)name;
84-
group->Destroy();
85-
delete group;
86-
group = nullptr;
87-
}
76+
std::unordered_map<std::string, FiberConnectorGroup*> tmp_tcp;
77+
{
78+
std::unique_lock lock(tcp_mutex_);
79+
tcp_connector_groups_to_destroy_.swap(tmp_tcp);
80+
}
8881

89-
tcp_impl->tcp_connector_groups.clear();
90-
tcp_impl->Retire();
82+
for (auto& group : tmp_tcp) {
83+
group.second->Destroy();
84+
delete group.second;
85+
group.second = nullptr;
9186
}
9287

93-
auto udp_impl = udp_impl_.exchange(nullptr, std::memory_order_relaxed);
94-
if (udp_impl) {
95-
if (udp_impl->udp_connector_groups[0] != nullptr) {
96-
udp_impl->udp_connector_groups[0]->Destroy();
97-
delete udp_impl->udp_connector_groups[0];
98-
udp_impl->udp_connector_groups[0] = nullptr;
99-
}
88+
{
89+
std::scoped_lock _(udp_mutex_);
90+
auto udp_impl = udp_impl_.exchange(nullptr, std::memory_order_relaxed);
91+
if (udp_impl) {
92+
if (udp_impl->udp_connector_groups[0] != nullptr) {
93+
udp_impl->udp_connector_groups[0]->Destroy();
94+
delete udp_impl->udp_connector_groups[0];
95+
udp_impl->udp_connector_groups[0] = nullptr;
96+
}
10097

101-
if (udp_impl->udp_connector_groups[1] != nullptr) {
102-
udp_impl->udp_connector_groups[1]->Destroy();
103-
delete udp_impl->udp_connector_groups[1];
104-
udp_impl->udp_connector_groups[1] = nullptr;
98+
if (udp_impl->udp_connector_groups[1] != nullptr) {
99+
udp_impl->udp_connector_groups[1]->Destroy();
100+
delete udp_impl->udp_connector_groups[1];
101+
udp_impl->udp_connector_groups[1] = nullptr;
102+
}
103+
udp_impl->Retire();
105104
}
106-
udp_impl->Retire();
107105
}
108-
109-
fiber_transport_state_.store(ClientTransportState::kDestroyed, std::memory_order_release);
110106
}
111107

112108
FiberConnectorGroup* FiberConnectorGroupManager::Get(const NodeAddr& node_addr) {
@@ -123,40 +119,29 @@ FiberConnectorGroup* FiberConnectorGroupManager::GetFromTcpGroup(const NodeAddr&
123119
snprintf(const_cast<char*>(endpoint.c_str()), len, "%s:%d", node_addr.ip.c_str(), node_addr.port);
124120

125121
{
126-
Hazptr hazptr;
127-
auto ptr = hazptr.Keep(&tcp_impl_);
128-
auto it = ptr->tcp_connector_groups.find(endpoint);
129-
if (it != ptr->tcp_connector_groups.end()) {
122+
std::shared_lock lock(tcp_mutex_);
123+
auto it = tcp_connector_groups_.find(endpoint);
124+
if (it != tcp_connector_groups_.end()) {
130125
return it->second;
131126
}
132127
}
133128

134-
FiberConnectorGroup* pool{nullptr};
129+
if (TRPC_UNLIKELY(fiber_transport_state_.load(std::memory_order_relaxed) != ClientTransportState::kInitialized)) {
130+
return nullptr;
131+
}
135132

136133
{
137-
auto new_tcp_impl = std::make_unique<TcpImpl>();
138-
139-
std::scoped_lock _(mutex_);
140-
141-
{
142-
Hazptr hazptr;
143-
auto ptr = hazptr.Keep(&tcp_impl_);
144-
auto it = ptr->tcp_connector_groups.find(endpoint);
145-
if (it != ptr->tcp_connector_groups.end()) {
146-
return it->second;
147-
} else {
148-
new_tcp_impl->tcp_connector_groups = ptr->tcp_connector_groups;
149-
}
134+
std::unique_lock lock(tcp_mutex_);
135+
auto it = tcp_connector_groups_.find(endpoint);
136+
if (it != tcp_connector_groups_.end()) {
137+
return it->second;
150138
}
151139

152-
pool = CreateTcpConnectorGroup(node_addr);
140+
FiberConnectorGroup* pool = CreateTcpConnectorGroup(node_addr);
141+
tcp_connector_groups_.emplace(std::move(endpoint), pool);
153142

154-
new_tcp_impl->tcp_connector_groups.emplace(std::move(endpoint), pool);
155-
156-
tcp_impl_.exchange(new_tcp_impl.release(), std::memory_order_acq_rel)->Retire();
143+
return pool;
157144
}
158-
159-
return pool;
160145
}
161146

162147
FiberConnectorGroup* FiberConnectorGroupManager::GetFromUdpGroup(const NodeAddr& node_addr) {
@@ -176,7 +161,7 @@ FiberConnectorGroup* FiberConnectorGroupManager::GetFromUdpGroup(const NodeAddr&
176161
{
177162
auto new_udp_impl = std::make_unique<UdpImpl>();
178163

179-
std::scoped_lock _(mutex_);
164+
std::scoped_lock _(udp_mutex_);
180165
{
181166
Hazptr hazptr;
182167
auto old_udp_impl = hazptr.Keep(&udp_impl_);

trpc/transport/client/fiber/fiber_connector_group_manager.h

Lines changed: 12 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616
#include <atomic>
1717
#include <memory>
1818
#include <mutex>
19+
#include <shared_mutex>
1920
#include <unordered_map>
2021

2122
#include "trpc/transport/client/client_transport_message.h"
@@ -50,20 +51,25 @@ class FiberConnectorGroupManager {
5051
int GetUdpConnectorGroupIndex(bool is_ipv6) const { return is_ipv6 ? 1 : 0; }
5152

5253
private:
53-
struct TcpImpl : HazptrObject<TcpImpl> {
54-
std::unordered_map<std::string, FiberConnectorGroup*> tcp_connector_groups;
55-
};
56-
5754
struct UdpImpl : HazptrObject<UdpImpl> {
5855
// 0: ipv4, 1: ipv6
5956
FiberConnectorGroup* udp_connector_groups[2];
6057
};
6158

62-
std::atomic<TcpImpl*> tcp_impl_{nullptr};
59+
std::unordered_map<std::string, FiberConnectorGroup*> tcp_connector_groups_;
60+
61+
// Initialized by Stop function to store connector groups which will by destroyed by Destroy function.
62+
std::unordered_map<std::string, FiberConnectorGroup*> tcp_connector_groups_to_destroy_;
63+
64+
// With gcc 8.3.1, no fiber yield is allowed in critical section protected by shared mutex,
65+
// as fiber may rescheduled into another thread, making shared mutex unlock an undefined behavior,
66+
// which may leads to deadlock, as specified by cpp cppreference:
67+
// https://en.cppreference.com/w/cpp/thread/shared_mutex/unlock
68+
mutable std::shared_mutex tcp_mutex_;
6369

6470
std::atomic<UdpImpl*> udp_impl_{nullptr};
6571

66-
std::mutex mutex_;
72+
std::mutex udp_mutex_;
6773

6874
std::atomic<ClientTransportState> fiber_transport_state_{ClientTransportState::kUnknown};
6975

trpc/transport/client/fiber/fiber_transport.cc

Lines changed: 38 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -44,12 +44,17 @@ void FiberTransport::Destroy() {
4444

4545
int FiberTransport::SendRecv(CTransportReqMsg* req_msg, CTransportRspMsg* rsp_msg) {
4646
if (IsRunningInFiberWorker()) {
47-
FiberConnectorGroup* connector_group = connector_group_manager_->Get(req_msg->context->GetNodeAddr());
48-
TRPC_ASSERT(connector_group && "connector_group can not be nullptr");
49-
5047
if (!req_msg->context->IsBackupRequest()) {
51-
return connector_group->SendRecv(req_msg, rsp_msg);
48+
FiberConnectorGroup* connector_group = connector_group_manager_->Get(req_msg->context->GetNodeAddr());
49+
if (connector_group != nullptr) {
50+
return connector_group->SendRecv(req_msg, rsp_msg);
51+
}
52+
53+
TRPC_FMT_ERROR("Can't get connector group of {}:{}", req_msg->context->GetNodeAddr().ip,
54+
req_msg->context->GetNodeAddr().port);
55+
return TrpcRetCode::TRPC_INVOKE_UNKNOWN_ERR;
5256
}
57+
5358
return SendRecvForBackupRequest(req_msg, rsp_msg);
5459
}
5560

@@ -89,9 +94,6 @@ int FiberTransport::SendRecvForBackupRequest(CTransportReqMsg* req_msg, CTranspo
8994
NoncontiguousBuffer buff_back(req_msg->send_data);
9095

9196
for (int i = 0; i < 2; ++i) {
92-
FiberConnectorGroup* connector_group = connector_group_manager_->Get(backup_info->backup_addrs[i].addr);
93-
TRPC_ASSERT(connector_group && "connector_group can not be nullptr");
94-
9597
auto cb = [&ret_code, i, backup_info, sync_retry](int err_code, std::string&& err_msg) {
9698
if (sync_retry->IsFinished()) {
9799
return;
@@ -110,7 +112,14 @@ int FiberTransport::SendRecvForBackupRequest(CTransportReqMsg* req_msg, CTranspo
110112
}
111113
};
112114

113-
connector_group->SendRecvForBackupRequest(req_msg, rsp_msg, std::move(cb));
115+
auto& addr = backup_info->backup_addrs[i].addr;
116+
FiberConnectorGroup* connector_group = connector_group_manager_->Get(addr);
117+
if (connector_group != nullptr) {
118+
connector_group->SendRecvForBackupRequest(req_msg, rsp_msg, std::move(cb));
119+
} else {
120+
TRPC_FMT_ERROR("Can't get connector group of {}:{}", addr.ip, addr.port);
121+
cb(TrpcRetCode::TRPC_CLIENT_CONNECT_ERR, "connector_group is nullptr");
122+
}
114123

115124
if (i == 0) {
116125
if (ret_code != TrpcRetCode::TRPC_CLIENT_CONNECT_ERR) {
@@ -132,11 +141,16 @@ int FiberTransport::SendRecvForBackupRequest(CTransportReqMsg* req_msg, CTranspo
132141

133142
Future<CTransportRspMsg> FiberTransport::AsyncSendRecv(CTransportReqMsg* req_msg) {
134143
if (IsRunningInFiberWorker()) {
135-
FiberConnectorGroup* connector_group = connector_group_manager_->Get(req_msg->context->GetNodeAddr());
136-
TRPC_ASSERT(connector_group && "connector_group can not be nullptr");
137-
138144
if (!req_msg->context->IsBackupRequest()) {
139-
return connector_group->AsyncSendRecv(req_msg);
145+
FiberConnectorGroup* connector_group = connector_group_manager_->Get(req_msg->context->GetNodeAddr());
146+
if (connector_group != nullptr) {
147+
return connector_group->AsyncSendRecv(req_msg);
148+
}
149+
150+
TRPC_FMT_ERROR("Can't get connector group of {}:{}", req_msg->context->GetNodeAddr().ip,
151+
req_msg->context->GetNodeAddr().port);
152+
return MakeExceptionFuture<CTransportRspMsg>(
153+
CommonException("not found connector group.", TrpcRetCode::TRPC_INVOKE_UNKNOWN_ERR));
140154
}
141155

142156
return MakeExceptionFuture<CTransportRspMsg>(
@@ -174,9 +188,13 @@ int FiberTransport::SendOnly(CTransportReqMsg* req_msg) {
174188

175189
if (IsRunningInFiberWorker()) {
176190
FiberConnectorGroup* connector_group = connector_group_manager_->Get(req_msg->context->GetNodeAddr());
177-
TRPC_ASSERT(connector_group && "connector_group can not be nullptr");
178-
179-
ret = connector_group->SendOnly(req_msg);
191+
if (connector_group != nullptr) {
192+
ret = connector_group->SendOnly(req_msg);
193+
} else {
194+
TRPC_FMT_ERROR("Can't get connector group of {}:{}", req_msg->context->GetNodeAddr().ip,
195+
req_msg->context->GetNodeAddr().port);
196+
ret = TrpcRetCode::TRPC_INVOKE_UNKNOWN_ERR;
197+
}
180198

181199
object_pool::Delete(req_msg);
182200
} else {
@@ -206,9 +224,12 @@ int FiberTransport::SendOnlyFromOutSide(CTransportReqMsg* req_msg) {
206224
stream::StreamReaderWriterProviderPtr FiberTransport::CreateStream(const NodeAddr& addr,
207225
stream::StreamOptions&& stream_options) {
208226
FiberConnectorGroup* connector_group = connector_group_manager_->Get(addr);
209-
TRPC_ASSERT(connector_group && "connector_group can not be nullptr");
227+
if (connector_group) {
228+
return connector_group->CreateStream(std::move(stream_options));
229+
}
210230

211-
return connector_group->CreateStream(std::move(stream_options));
231+
TRPC_FMT_ERROR("Can't get connector group of {}:{}", addr.ip, addr.port);
232+
return nullptr;
212233
}
213234

214235
} // namespace trpc

0 commit comments

Comments
 (0)