This commit is contained in:
Craig Tiller 2024-06-11 09:36:08 -07:00
parent b15f503cc9
commit 248a4ebe78
5 changed files with 70 additions and 45 deletions

View File

@ -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",

View File

@ -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();
}));
}

View File

@ -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();

View File

@ -14,6 +14,8 @@
#include "src/core/lib/transport/call_spine.h"
#include "absl/functional/any_invocable.h"
#include <grpc/support/port_platform.h>
#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<void(ServerMetadata&)>
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<ServerMetadataHandle> 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<ServerMetadataHandle> 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(

View File

@ -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<void(ServerMetadata&)>
on_server_trailing_metadata_from_initiator = [](ServerMetadata&) {});
} // namespace grpc_core