From 56204de2dd2f0598b0beddee7257a6f4a1cf9c7e Mon Sep 17 00:00:00 2001 From: ihsan Date: Thu, 4 Jun 2026 10:43:20 +0300 Subject: [PATCH 1/2] Do not warn on benign double completion of backup-aware invocations For a backup-aware operation the primary response (notify_response) and the last backup ack (notify_backup) can race to complete the invocation. Both paths complete with the same pending response, so the loser of the race hits boost::promise_already_satisfied. The future is still completed exactly once with the correct value and the duplicate is harmlessly dropped, but complete() logged every occurrence at WARNING, flooding the logs of healthy clients under load. The race became routinely observable after the multi-threaded pipeline work, which made the response and the backup ack run on different threads. Mirror the Java client's AbstractInvocationFuture.warnIfSuspiciousDouble- Completion: only warn when the offered response differs from the canonical pending_response_. A duplicate completion with the same response is logged at finest, while a genuinely conflicting value still warns. This also matches how set_exception() already treats promise_already_satisfied. Add regression tests that drive the private completion path via friend access to assert the benign duplicate is not a warning and a conflicting one still is. Fixes #1474 --- .../client/spi/impl/ClientInvocation.h | 8 + hazelcast/src/hazelcast/client/spi.cpp | 24 +++ hazelcast/test/src/HazelcastTests5.cpp | 173 ++++++++++++++++++ 3 files changed, 205 insertions(+) diff --git a/hazelcast/include/hazelcast/client/spi/impl/ClientInvocation.h b/hazelcast/include/hazelcast/client/spi/impl/ClientInvocation.h index f97ccf002..db7e86774 100644 --- a/hazelcast/include/hazelcast/client/spi/impl/ClientInvocation.h +++ b/hazelcast/include/hazelcast/client/spi/impl/ClientInvocation.h @@ -40,6 +40,12 @@ class logger; namespace client { class address; +namespace test { +// White-box test fixture; granted access to private completion internals to +// deterministically reproduce the response/backup-ack completion race. +class client_invocation_test; +} // namespace test + namespace connection { class Connection; } @@ -140,6 +146,8 @@ class HAZELCAST_API ClientInvocation friend std::ostream& operator<<(std::ostream& os, const ClientInvocation& invocation); + friend class ::hazelcast::client::test::client_invocation_test; + boost::promise& get_promise(); bool detect_and_handle_backup_timeout( diff --git a/hazelcast/src/hazelcast/client/spi.cpp b/hazelcast/src/hazelcast/client/spi.cpp index cec33dafc..53819d82a 100644 --- a/hazelcast/src/hazelcast/client/spi.cpp +++ b/hazelcast/src/hazelcast/client/spi.cpp @@ -2119,6 +2119,30 @@ ClientInvocation::complete(const std::shared_ptr& msg, try { // TODO: move msg content here? this->invocation_promise_.set_value(*msg); + } catch (boost::promise_already_satisfied& e) { + // The promise is already satisfied. This is an expected, benign race + // between notify_response() (primary response) and notify_backup() (the + // last backup ack): for a backup-aware operation both paths can reach + // complete() with the very same pending_response_. The first one wins + // and the duplicate is harmlessly dropped. + // + // This mirrors Java's AbstractInvocationFuture.warnIfSuspiciousDouble- + // Completion, which only warns when the already-set value differs from + // the offered one. Here the canonical value is pending_response_, so a + // completion with the same message is benign (logged at finest) while a + // genuinely different value remains a warning. + auto pending = pending_response_.load(); + auto level = (!pending || *pending != msg) + ? ::hazelcast::logger::level::warning + : ::hazelcast::logger::level::finest; + if (logger_.enabled(level)) { + logger_.log( + level, + boost::str( + boost::format("Failed to set the response for invocation. " + "Dropping the response. %1%, %2% Response: %3%") % + e.what() % *this % *msg)); + } } catch (std::exception& e) { HZ_LOG(logger_, warning, diff --git a/hazelcast/test/src/HazelcastTests5.cpp b/hazelcast/test/src/HazelcastTests5.cpp index 1e4d59c14..cede13a5c 100644 --- a/hazelcast/test/src/HazelcastTests5.cpp +++ b/hazelcast/test/src/HazelcastTests5.cpp @@ -13,15 +13,19 @@ * See the License for the specific language governing permissions and * limitations under the License. */ +#include #include #include #include +#include #include #include #include +#include #include #include +#include #include @@ -46,7 +50,11 @@ #include #include #include +#include +#include #include +#include +#include #include #include #include @@ -4174,6 +4182,171 @@ struct hz_serializer } // namespace client } // namespace hazelcast +namespace hazelcast { +namespace client { +namespace test { + +/** + * Regression tests for the completion of a backup-aware invocation. + * + * For a backup-aware operation both notify_response() (primary response) and + * notify_backup() (last backup ack) can race to call ClientInvocation::complete + * with the very same pending response. The promise can only be satisfied once, + * so the loser of the race hits boost::promise_already_satisfied. This is a + * benign duplicate and must not be reported as a WARNING (it used to flood the + * logs of healthy clients). A double completion carrying a *different* value is + * genuinely suspicious and must still be logged at WARNING. + * + * The double completion only arises from a concurrent interleaving that cannot + * be forced through the public notify()/notify_backup() API, so these tests + * drive the private completion path directly via friend access, after staging + * pending_response_ exactly as notify_response() would. + */ +class client_invocation_test : public ClientTest +{ +protected: + using message_ptr = std::shared_ptr; + using invocation_ptr = std::shared_ptr; + using captured_log = std::pair; + + // A single member is enough and is shared by all cases in this suite; the + // completion logic under test does not depend on cluster state. + static void SetUpTestSuite() + { + member_ = new HazelcastServer{ default_server_factory() }; + } + + static void TearDownTestSuite() + { + delete member_; + member_ = nullptr; + } + + // Stages the canonical pending response, mirroring notify_response(). + static void set_pending_response(const invocation_ptr& invocation, + const message_ptr& msg) + { + invocation->pending_response_.store( + boost::make_shared>(msg)); + } + + // Invokes the private complete(); erase=false so we do not touch the + // invocation registry for this manually created invocation. + static void complete(const invocation_ptr& invocation, + const message_ptr& msg) + { + invocation->complete(msg, false); + } + + // Builds a client whose logger forwards every message (finest and above) + // into the supplied sink. + static hazelcast_client make_capturing_client( + std::mutex& mtx, + std::vector& sink) + { + client_config cfg = get_config(); + cfg.get_logger_config() + .level(logger::level::finest) + .handler([&mtx, &sink](const std::string&, + const std::string&, + logger::level lvl, + const std::string& msg) { + std::lock_guard g{ mtx }; + sink.emplace_back(lvl, msg); + }); + return new_client(std::move(cfg)).get(); + } + + static int count_drop_logs(std::mutex& mtx, + const std::vector& sink, + logger::level lvl) + { + std::lock_guard g{ mtx }; + return static_cast( + std::count_if(sink.begin(), sink.end(), [lvl](const captured_log& e) { + return e.first == lvl && + e.second.find("Failed to set the response") != + std::string::npos; + })); + } + + static HazelcastServer* member_; +}; + +HazelcastServer* client_invocation_test::member_ = nullptr; + +TEST_F(client_invocation_test, + benign_double_completion_with_same_response_is_logged_at_finest) +{ + std::mutex mtx; + std::vector logs; + auto client = make_capturing_client(mtx, logs); + + spi::ClientContext context{ client }; + auto request = protocol::codec::client_ping_encode(); + auto invocation = + spi::impl::ClientInvocation::create(context, request, "ping"); + auto future = invocation->get_promise().get_future(); + + auto response = std::make_shared( + protocol::codec::client_ping_encode()); + + // Reproduce the post-race state: pending_response_ holds the response that + // both the primary-response and the last-backup-ack paths would complete + // with, then let both paths complete with that very same message. + set_pending_response(invocation, response); + complete(invocation, response); // first completion wins + complete(invocation, response); // benign duplicate, must not warn + + // The promise was satisfied exactly once with the response. + ASSERT_EQ(boost::future_status::ready, + future.wait_for(boost::chrono::seconds(0))); + ASSERT_NO_THROW(future.get()); + + ASSERT_EQ(0, count_drop_logs(mtx, logs, logger::level::warning)); + ASSERT_EQ(1, count_drop_logs(mtx, logs, logger::level::finest)); + + client.shutdown().get(); +} + +TEST_F( + client_invocation_test, + suspicious_double_completion_with_different_response_is_logged_at_warning) +{ + std::mutex mtx; + std::vector logs; + auto client = make_capturing_client(mtx, logs); + + spi::ClientContext context{ client }; + auto request = protocol::codec::client_ping_encode(); + auto invocation = + spi::impl::ClientInvocation::create(context, request, "ping"); + auto future = invocation->get_promise().get_future(); + + auto first_response = std::make_shared( + protocol::codec::client_ping_encode()); + auto different_response = std::make_shared( + protocol::codec::client_ping_encode()); + + set_pending_response(invocation, first_response); + complete(invocation, first_response); // first completion wins + complete(invocation, + different_response); // genuinely conflicting, must warn + + ASSERT_EQ(boost::future_status::ready, + future.wait_for(boost::chrono::seconds(0))); + ASSERT_NO_THROW(future.get()); + + ASSERT_EQ(1, count_drop_logs(mtx, logs, logger::level::warning)); + ASSERT_EQ(0, count_drop_logs(mtx, logs, logger::level::finest)); + + client.shutdown().get(); +} + +} // namespace test +} // namespace client +} // namespace hazelcast + namespace std { template<> struct hash From 88e26cb4938adb4b1f4c37eac20a56c67994901b Mon Sep 17 00:00:00 2001 From: ihsan Date: Thu, 4 Jun 2026 12:22:22 +0300 Subject: [PATCH 2/2] Review comment fixes. --- hazelcast/include/hazelcast/client/spi/impl/ClientInvocation.h | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/hazelcast/include/hazelcast/client/spi/impl/ClientInvocation.h b/hazelcast/include/hazelcast/client/spi/impl/ClientInvocation.h index db7e86774..56457e61d 100644 --- a/hazelcast/include/hazelcast/client/spi/impl/ClientInvocation.h +++ b/hazelcast/include/hazelcast/client/spi/impl/ClientInvocation.h @@ -41,7 +41,7 @@ namespace client { class address; namespace test { -// White-box test fixture; granted access to private completion internals to +// Exposed for testing only; granted access to private completion internals to // deterministically reproduce the response/backup-ack completion race. class client_invocation_test; } // namespace test