From 248a4ebe78ef0919f292a7137ae0f955e444c496 Mon Sep 17 00:00:00 2001 From: Craig Tiller Date: Tue, 11 Jun 2024 09:36:08 -0700 Subject: [PATCH] x --- src/core/BUILD | 8 +- .../ext/transport/inproc/inproc_transport.cc | 6 +- src/core/lib/surface/client_call.cc | 1 + src/core/lib/transport/call_spine.cc | 95 +++++++++++-------- src/core/lib/transport/call_spine.h | 5 +- 5 files changed, 70 insertions(+), 45 deletions(-) diff --git a/src/core/BUILD b/src/core/BUILD index e34ac606f0f..2098aad343a 100644 --- a/src/core/BUILD +++ b/src/core/BUILD @@ -7010,8 +7010,8 @@ grpc_cc_library( "ext/transport/inproc/legacy_inproc_transport.h", ], external_deps = [ + "absl/log", "absl/log:check", - "absl/log:log", "absl/status", "absl/status:statusor", "absl/strings", @@ -7028,6 +7028,7 @@ grpc_cc_library( "error", "experiments", "iomgr_fwd", + "metadata", "metadata_batch", "resource_quota", "slice", @@ -7554,7 +7555,10 @@ grpc_cc_library( hdrs = [ "lib/transport/call_spine.h", ], - external_deps = ["absl/log:check"], + external_deps = [ + "absl/functional:any_invocable", + "absl/log:check", + ], deps = [ "1999", "call_arena_allocator", diff --git a/src/core/ext/transport/inproc/inproc_transport.cc b/src/core/ext/transport/inproc/inproc_transport.cc index 78dfb2761c7..5a3d62243b7 100644 --- a/src/core/ext/transport/inproc/inproc_transport.cc +++ b/src/core/ext/transport/inproc/inproc_transport.cc @@ -32,6 +32,7 @@ #include "src/core/lib/promise/try_seq.h" #include "src/core/lib/resource_quota/resource_quota.h" #include "src/core/lib/surface/channel_create.h" +#include "src/core/lib/transport/metadata.h" #include "src/core/lib/transport/transport.h" #include "src/core/server/server.h" @@ -149,7 +150,10 @@ class InprocClientTransport final : public ClientTransport { return server_call_initiator.status(); } ForwardCall(child_call_handler, - std::move(*server_call_initiator)); + std::move(*server_call_initiator), + [](ServerMetadata& md) { + md.Set(GrpcStatusFromWire(), true); + }); return absl::OkStatus(); })); } diff --git a/src/core/lib/surface/client_call.cc b/src/core/lib/surface/client_call.cc index 6ca22f87b4d..ee06b0e0b99 100644 --- a/src/core/lib/surface/client_call.cc +++ b/src/core/lib/surface/client_call.cc @@ -338,6 +338,7 @@ void ClientCall::CommitBatch(const grpc_op* ops, size_t nops, void* notify_tag, [this, out_status, out_status_details, out_error_string, out_trailing_metadata]( ServerMetadataHandle server_trailing_metadata) { + ResetDeadline(); GRPC_TRACE_LOG(call, INFO) << DebugTag() << "RecvStatusOnClient " << server_trailing_metadata->DebugString(); diff --git a/src/core/lib/transport/call_spine.cc b/src/core/lib/transport/call_spine.cc index 6d72e435c84..e75f7d1b621 100644 --- a/src/core/lib/transport/call_spine.cc +++ b/src/core/lib/transport/call_spine.cc @@ -14,6 +14,8 @@ #include "src/core/lib/transport/call_spine.h" +#include "absl/functional/any_invocable.h" + #include #include "src/core/lib/promise/for_each.h" @@ -21,7 +23,9 @@ namespace grpc_core { -void ForwardCall(CallHandler call_handler, CallInitiator call_initiator) { +void ForwardCall(CallHandler call_handler, CallInitiator call_initiator, + absl::AnyInvocable + on_server_trailing_metadata_from_initiator) { // Read messages from handler into initiator. call_handler.SpawnGuarded("read_messages", [call_handler, call_initiator]() mutable { @@ -68,47 +72,56 @@ void ForwardCall(CallHandler call_handler, CallInitiator call_initiator) { return Empty(); }); }); - call_initiator.SpawnInfallible("read_the_things", [call_initiator, - call_handler]() mutable { - return Seq( - call_initiator.CancelIfFails(TrySeq( - call_initiator.PullServerInitialMetadata(), + call_initiator.SpawnInfallible( + "read_the_things", + [call_initiator, call_handler, + on_server_trailing_metadata_from_initiator = + std::move(on_server_trailing_metadata_from_initiator)]() mutable { + return Seq( + call_initiator.CancelIfFails(TrySeq( + call_initiator.PullServerInitialMetadata(), + [call_handler, call_initiator]( + absl::optional md) mutable { + const bool has_md = md.has_value(); + return If( + has_md, + [&call_handler, &call_initiator, + md = std::move(md)]() mutable { + call_handler.SpawnGuarded( + "recv_initial_metadata", + [md = std::move(*md), call_handler]() mutable { + return call_handler.PushServerInitialMetadata( + std::move(md)); + }); + return ForEach( + OutgoingMessages(call_initiator), + [call_handler](MessageHandle msg) mutable { + return call_handler.SpawnWaitable( + "recv_message", [msg = std::move(msg), + call_handler]() mutable { + return call_handler.CancelIfFails( + call_handler.PushMessage( + std::move(msg))); + }); + }); + }, + []() -> StatusFlag { return Success{}; }); + })), + call_initiator.PullServerTrailingMetadata(), [call_handler, - call_initiator](absl::optional md) mutable { - const bool has_md = md.has_value(); - return If( - has_md, - [&call_handler, &call_initiator, - md = std::move(md)]() mutable { - call_handler.SpawnGuarded( - "recv_initial_metadata", - [md = std::move(*md), call_handler]() mutable { - return call_handler.PushServerInitialMetadata( - std::move(md)); - }); - return ForEach( - OutgoingMessages(call_initiator), - [call_handler](MessageHandle msg) mutable { - return call_handler.SpawnWaitable( - "recv_message", - [msg = std::move(msg), call_handler]() mutable { - return call_handler.CancelIfFails( - call_handler.PushMessage(std::move(msg))); - }); - }); - }, - []() -> StatusFlag { return Success{}; }); - })), - call_initiator.PullServerTrailingMetadata(), - [call_handler](ServerMetadataHandle md) mutable { - call_handler.SpawnInfallible( - "recv_trailing", [call_handler, md = std::move(md)]() mutable { - call_handler.PushServerTrailingMetadata(std::move(md)); - return Empty{}; - }); - return Empty{}; - }); - }); + on_server_trailing_metadata_from_initiator = + std::move(on_server_trailing_metadata_from_initiator)]( + ServerMetadataHandle md) mutable { + on_server_trailing_metadata_from_initiator(*md); + call_handler.SpawnInfallible( + "recv_trailing", + [call_handler, md = std::move(md)]() mutable { + call_handler.PushServerTrailingMetadata(std::move(md)); + return Empty{}; + }); + return Empty{}; + }); + }); } CallInitiatorAndHandler MakeCallPair( diff --git a/src/core/lib/transport/call_spine.h b/src/core/lib/transport/call_spine.h index 31b792e6f99..b3c4f173e86 100644 --- a/src/core/lib/transport/call_spine.h +++ b/src/core/lib/transport/call_spine.h @@ -441,7 +441,10 @@ auto OutgoingMessages(CallHalf h) { // Forward a call from `call_handler` to `call_initiator` (with initial metadata // `client_initial_metadata`) -void ForwardCall(CallHandler call_handler, CallInitiator call_initiator); +void ForwardCall( + CallHandler call_handler, CallInitiator call_initiator, + absl::AnyInvocable + on_server_trailing_metadata_from_initiator = [](ServerMetadata&) {}); } // namespace grpc_core