Skip to content

Commit 3d299f7

Browse files
authored
BugFix: Fix some problems at stream invoking (#60)
- fix stream crash at ServerContext because it has been moved - fix stream transinfo not pass into ServerContext - delete deprecated interface and useless code in trpc/client/make_client_context.h/.cc
1 parent 3cc8ba5 commit 3d299f7

7 files changed

Lines changed: 14 additions & 49 deletions

File tree

trpc/client/make_client_context.cc

Lines changed: 0 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -21,10 +21,6 @@
2121
#include "trpc/util/hash_util.h"
2222
#include "trpc/util/log/logging.h"
2323
#include "trpc/util/time.h"
24-
// #ifdef TRPC_BUILD_INCLUDE_RPCZ
25-
// #include "trpc/rpcz/filter/rpcz_filter_index.h"
26-
// #include "trpc/rpcz/span.h"
27-
// #endif
2824

2925
namespace trpc {
3026

@@ -75,39 +71,6 @@ void RegisterMakeClientContextCallback(MakeClientContextCallback&& callback) {
7571
callbacks.emplace_back(std::move(callback));
7672
}
7773

78-
ClientContextPtr MakeClientContext(const ServerContextPtr& ctx) {
79-
ClientContextPtr client_ctx = MakeRefCounted<ClientContext>();
80-
81-
const auto& trans_info = ctx->GetPbReqTransInfo();
82-
if (trans_info.size() > 0) {
83-
client_ctx->SetReqTransInfo(trans_info.begin(), trans_info.end());
84-
}
85-
86-
client_ctx->SetMessageType(ctx->GetMessageType());
87-
client_ctx->SetCallerName(ctx->GetCalleeName());
88-
client_ctx->SetCallerFuncName(ctx->GetFuncName());
89-
90-
RunMakeClientContextCallbacks(ctx, client_ctx);
91-
92-
// Calculate the remaining timeout and set it to the client context.
93-
int64_t nowms = static_cast<int64_t>(trpc::time::GetMilliSeconds());
94-
int64_t cost_time = nowms - ctx->GetRecvTimestamp();
95-
int64_t left_time = static_cast<int64_t>(ctx->GetTimeout()) - cost_time;
96-
if (left_time < 0) {
97-
left_time = 0;
98-
}
99-
100-
client_ctx->SetTimeout(left_time);
101-
102-
if (ctx->IsUseFullLinkTimeout()) {
103-
client_ctx->SetFullLinkTimeout(left_time);
104-
} else {
105-
client_ctx->SetTimeout(left_time);
106-
}
107-
108-
return client_ctx;
109-
}
110-
11174
ClientContextPtr MakeClientContext(const ServiceProxyPtr& proxy) {
11275
TRPC_ASSERT(proxy && "proxy must be created before create a ClientContext obj");
11376
return MakeRefCounted<ClientContext>(proxy->GetClientCodec());

trpc/client/make_client_context.h

Lines changed: 0 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -68,12 +68,4 @@ void BackFillServerTransInfo(const ClientContextPtr& client_context, const Serve
6868
using MakeClientContextCallback = std::function<void(const ServerContextPtr&, ClientContextPtr&)>;
6969
void RegisterMakeClientContextCallback(MakeClientContextCallback&& callback);
7070

71-
/// @brief Create client context based on server context
72-
/// @param ctx server context
73-
/// @return client context
74-
/// @deprecated Use MakeClientContext(const ServerContextPtr& ctx, const ServiceProxyPtr& proxy) instead.
75-
/// @private
76-
[[deprecated("use the interface with params as 'ctx' and 'proxy' ")]]
77-
ClientContextPtr MakeClientContext(const ServerContextPtr& ctx);
78-
7971
} // namespace trpc

trpc/client/service_proxy.cc

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -692,6 +692,10 @@ stream::StreamReaderWriterProviderPtr ServiceProxy::SelectStreamProvider(const C
692692
TRPC_ASSERT(codec_->Name() == "trpc" || codec_->Name() == "http");
693693
TRPC_ASSERT(thread_model_ != nullptr);
694694

695+
if (context->GetResponse() == nullptr) {
696+
context->SetResponse(codec_->CreateResponsePtr());
697+
}
698+
695699
FillClientContext(context);
696700

697701
// Get one address of peer server

trpc/stream/grpc/grpc_server_stream_handler.cc

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -269,7 +269,7 @@ StreamReaderWriterProviderPtr FiberGrpcServerStreamHandler::CreateStream(StreamO
269269
TRPC_LOG_ERROR("stream " << stream_id << " already exists.");
270270
return nullptr;
271271
}
272-
ServerContextPtr& context = std::any_cast<ServerContextPtr&>(options.context.context);
272+
ServerContextPtr context = std::any_cast<ServerContextPtr>(options.context.context);
273273
auto stream = MakeRefCounted<GrpcServerStream>(std::move(options), session_.get(), &session_mutex_, &session_cv_);
274274
// Set the stream for the context so that it can be obtained in stream_rpc_method_handler.
275275
context->SetStreamReaderWriterProvider(stream);

trpc/stream/http/async/server/stream_handler.cc

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -29,7 +29,7 @@ StreamReaderWriterProviderPtr HttpServerAsyncStreamHandler::CreateStream(StreamO
2929
return nullptr;
3030
}
3131

32-
ServerContextPtr& context = std::any_cast<ServerContextPtr&>(options.context.context);
32+
ServerContextPtr context = std::any_cast<ServerContextPtr>(options.context.context);
3333
options.connection_id = GetMutableStreamOptions()->connection_id;
3434
stream_ = MakeRefCounted<HttpServerAsyncStream>(std::move(options));
3535
context->SetStreamReaderWriterProvider(stream_);

trpc/stream/trpc/trpc_server_stream.cc

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -118,6 +118,8 @@ RetCode TrpcServerStream::HandleInit(StreamRecvMessage&& msg) {
118118
context->SetCallerName(request_metadata.caller());
119119
context->SetFuncName(request_metadata.func());
120120
context->SetReqEncodeType(init_frame.stream_init_metadata.content_type());
121+
context->GetMutablePbReqTransInfo()->insert(request_metadata.trans_info().begin(),
122+
request_metadata.trans_info().end());
121123

122124
std::string var_path = "trpc/stream_rpc/server" + request_metadata.func();
123125
auto stream_var = StreamVarHelper::GetInstance()->GetOrCreateStreamVar(var_path);
@@ -272,6 +274,11 @@ RetCode TrpcServerStream::SendClose(StreamSendMessage&& msg) {
272274
return RetCode::kError;
273275
}
274276

277+
const auto& context = std::any_cast<const ServerContextPtr&>(GetMutableStreamOptions()->context.context);
278+
279+
auto& close_meta = std::any_cast<TrpcStreamCloseMeta&>(msg.metadata);
280+
close_meta.mutable_trans_info()->insert(context->GetPbRspTransInfo().begin(), context->GetPbRspTransInfo().end());
281+
275282
RetCode ret = TrpcStream::SendClose(std::move(msg));
276283
if (TRPC_UNLIKELY(ret != RetCode::kSuccess)) {
277284
return ret;
@@ -281,7 +288,6 @@ RetCode TrpcServerStream::SendClose(StreamSendMessage&& msg) {
281288

282289
Stop();
283290

284-
const auto& context = std::any_cast<const ServerContextPtr&>(GetMutableStreamOptions()->context.context);
285291
RunMessageFilter(FilterPoint::SERVER_PRE_SEND_MSG, context);
286292

287293
TRPC_FMT_DEBUG("server stream, send close end, stream id: {}", GetId());

trpc/stream/trpc/trpc_server_stream_handler.cc

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -35,7 +35,7 @@ StreamReaderWriterProviderPtr TrpcServerStreamHandler::CreateStream(StreamOption
3535
// Since IsNewStream has previously determined that the stream does not exist, it is likely that it still does not
3636
// exist, so create the stream first.
3737
options.connection_id = options_.connection_id;
38-
ServerContextPtr& context = std::any_cast<ServerContextPtr&>(options.context.context);
38+
ServerContextPtr context = std::any_cast<ServerContextPtr>(options.context.context);
3939
auto stream = MakeRefCounted<TrpcServerStream>(std::move(options));
4040
context->SetStreamReaderWriterProvider(stream);
4141

0 commit comments

Comments
 (0)