diff --git a/src/core/lib/surface/client_call.cc b/src/core/lib/surface/client_call.cc index ee06b0e0b99..e931f71a151 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) { + saw_trailing_metadata_.store(true, std::memory_order_relaxed); ResetDeadline(); GRPC_TRACE_LOG(call, INFO) << DebugTag() << "RecvStatusOnClient " diff --git a/src/core/lib/surface/client_call.h b/src/core/lib/surface/client_call.h index 26f27d26b5b..2d97ab9b731 100644 --- a/src/core/lib/surface/client_call.h +++ b/src/core/lib/surface/client_call.h @@ -82,8 +82,9 @@ class ClientCall final void InternalUnref(const char*) override { WeakUnref(); } void Orphaned() override { - // TODO(ctiller): only when we're not already finished - CancelWithError(absl::CancelledError()); + if (!saw_trailing_metadata_.load(std::memory_order_relaxed)) { + CancelWithError(absl::CancelledError()); + } } void SetCompletionQueue(grpc_completion_queue*) override { @@ -164,6 +165,7 @@ class ClientCall final ServerMetadataHandle received_initial_metadata_; ServerMetadataHandle received_trailing_metadata_; bool is_trailers_only_; + std::atomic saw_trailing_metadata_{false}; }; grpc_call* MakeClientCall( diff --git a/src/core/lib/surface/server_call.cc b/src/core/lib/surface/server_call.cc index 92f7db547d1..3eeddfca866 100644 --- a/src/core/lib/surface/server_call.cc +++ b/src/core/lib/surface/server_call.cc @@ -188,6 +188,8 @@ void ServerCall::CommitBatch(const grpc_op* ops, size_t nops, void* notify_tag, [this, cancelled = op->data.recv_close_on_server.cancelled]() { return Map(call_handler_.WasCancelled(), [cancelled, this](bool result) -> Success { + saw_was_cancelled_.store(true, + std::memory_order_relaxed); ResetDeadline(); *cancelled = result ? 1 : 0; return Success{}; diff --git a/src/core/lib/surface/server_call.h b/src/core/lib/surface/server_call.h index 9d571d111ac..fa463c2f6a6 100644 --- a/src/core/lib/surface/server_call.h +++ b/src/core/lib/surface/server_call.h @@ -20,6 +20,7 @@ #include #include +#include #include #include #include @@ -101,8 +102,9 @@ class ServerCall final : public Call, public DualRefCounted { void InternalUnref(const char*) override { WeakUnref(); } void Orphaned() override { - // TODO(ctiller): only when we're not already finished - CancelWithError(absl::CancelledError()); + if (!saw_was_cancelled_.load(std::memory_order_relaxed)) { + CancelWithError(absl::CancelledError()); + } } void SetCompletionQueue(grpc_completion_queue*) override { @@ -155,6 +157,7 @@ class ServerCall final : public Call, public DualRefCounted { ClientMetadataHandle client_initial_metadata_stored_; grpc_completion_queue* const cq_; ServerInterface* const server_; + std::atomic saw_was_cancelled_{false}; }; grpc_call* MakeServerCall(CallHandler call_handler,