From e3b88e7d86aff7707819ddd010850123e730a91a Mon Sep 17 00:00:00 2001 From: Craig Tiller Date: Tue, 11 Jun 2024 07:13:40 -0700 Subject: [PATCH] x --- .../thready_event_engine.cc | 38 ++++++++++++++++--- test/cpp/end2end/BUILD | 1 + tools/bazel.rc | 1 - 3 files changed, 33 insertions(+), 7 deletions(-) diff --git a/src/core/lib/event_engine/thready_event_engine/thready_event_engine.cc b/src/core/lib/event_engine/thready_event_engine/thready_event_engine.cc index f448108a390..3be7a0e0b32 100644 --- a/src/core/lib/event_engine/thready_event_engine/thready_event_engine.cc +++ b/src/core/lib/event_engine/thready_event_engine/thready_event_engine.cc @@ -22,6 +22,7 @@ #include #include "src/core/lib/gprpp/crash.h" +#include "src/core/lib/gprpp/sync.h" #include "src/core/lib/gprpp/thd.h" namespace grpc_event_engine { @@ -39,20 +40,45 @@ ThreadyEventEngine::CreateListener( absl::AnyInvocable on_shutdown, const EndpointConfig& config, std::unique_ptr memory_allocator_factory) { + struct AcceptState { + grpc_core::Mutex mu_; + grpc_core::CondVar cv_; + int pending_accepts_ ABSL_GUARDED_BY(mu_) = 0; + }; + auto accept_state = std::make_shared(); return impl_->CreateListener( - [this, on_accept = std::make_shared( - std::move(on_accept))](std::unique_ptr endpoint, - MemoryAllocator memory_allocator) { + [this, accept_state, + on_accept = std::make_shared( + std::move(on_accept))](std::unique_ptr endpoint, + MemoryAllocator memory_allocator) { + { + grpc_core::MutexLock lock(&accept_state->mu_); + ++accept_state->pending_accepts_; + } Asynchronously( - [on_accept, endpoint = std::move(endpoint), + [on_accept, accept_state, endpoint = std::move(endpoint), memory_allocator = std::move(memory_allocator)]() mutable { (*on_accept)(std::move(endpoint), std::move(memory_allocator)); + { + grpc_core::MutexLock lock(&accept_state->mu_); + --accept_state->pending_accepts_; + if (accept_state->pending_accepts_ == 0) { + accept_state->cv_.Signal(); + } + } }); }, - [this, + [this, accept_state, on_shutdown = std::move(on_shutdown)](absl::Status status) mutable { - Asynchronously([on_shutdown = std::move(on_shutdown), + Asynchronously([accept_state, on_shutdown = std::move(on_shutdown), status = std::move(status)]() mutable { + while (true) { + grpc_core::MutexLock lock(&accept_state->mu_); + if (accept_state->pending_accepts_ == 0) { + break; + } + accept_state->cv_.Wait(&accept_state->mu_); + } on_shutdown(std::move(status)); }); }, diff --git a/test/cpp/end2end/BUILD b/test/cpp/end2end/BUILD index 6ada7da27c1..4a424354c39 100644 --- a/test/cpp/end2end/BUILD +++ b/test/cpp/end2end/BUILD @@ -112,6 +112,7 @@ grpc_cc_test( tags = [ "cpp_end2end_test", "no_test_ios", + "thready_tsan", ], deps = [ "//:gpr", diff --git a/tools/bazel.rc b/tools/bazel.rc index 36ec153f52a..3d98df37012 100644 --- a/tools/bazel.rc +++ b/tools/bazel.rc @@ -133,7 +133,6 @@ build:thready_tsan --copt=-DGPR_NO_DIRECT_SYSCALLS build:thready_tsan --copt=-DGRPC_TSAN build:thready_tsan --copt=-DGRPC_MAXIMIZE_THREADYNESS build:thready_tsan --linkopt=-fsanitize=thread -build:thready_tsan --test_tag_filters=thready_tsan build:thready_tsan --action_env=TSAN_OPTIONS=suppressions=test/core/test_util/tsan_suppressions.txt:halt_on_error=1:second_deadlock_stack=1 build:tsan_macos --strip=never