Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions hazelcast/include/hazelcast/client/impl/auto_batcher.h
Original file line number Diff line number Diff line change
Expand Up @@ -98,6 +98,19 @@ class HAZELCAST_API auto_batcher
std::atomic<int32_t> num_returned_;
};

/**
* Tries to satisfy one caller: fast path, otherwise joins (or starts) the
* single in-flight batch fetch and completes the promise from its
* continuation. A caller that loses the race for an ID from the fresh
* batch is retried by posting this function to the executor again, so
* every retry starts from an empty stack. The retry must never be
* expressed as a nested future (then().unwrap() returning new_id()):
* each lost round would add one more unwrap layer to the caller's
* future and Boost completes such a chain recursively, which overflowed
* the executor thread's stack under contention (issue #1491).
*/
void try_get_id(const boost::shared_ptr<boost::promise<int64_t>>& p);

const int32_t batch_size_;
const std::chrono::milliseconds validity_;
util::hz_thread_pool& executor_;
Expand Down
82 changes: 57 additions & 25 deletions hazelcast/src/hazelcast/client/proxy.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,8 @@
#include <unordered_set>
#include <atomic>

#include <boost/asio/post.hpp>

#include "hazelcast/client/impl/ClientLockReferenceIdGenerator.h"
#include "hazelcast/client/proxy/PNCounterImpl.h"
#include "hazelcast/client/spi/ClientContext.h"
Expand All @@ -35,6 +37,7 @@
#include "hazelcast/client/proxy/ReplicatedMapImpl.h"
#include "hazelcast/client/flake_id_generator.h"
#include "hazelcast/client/reliable_topic.h"

#include "hazelcast/util/hz_thread_pool.h"

namespace hazelcast {
Expand Down Expand Up @@ -1163,6 +1166,24 @@ auto_batcher::new_id()
}
}

auto p = boost::make_shared<boost::promise<int64_t>>();
auto f = p->get_future();
try_get_id(p);
return f;
}

void
auto_batcher::try_get_id(const boost::shared_ptr<boost::promise<int64_t>>& p)
{
auto b = block_.load();
if (b) {
int64_t v = b->next();
if (v != INT64_MIN) {
p->set_value(v);
return;
}
}

// Slow path: elect a single fetcher; concurrent callers coalesce onto
// one outstanding batch request (single-flight, no thundering herd).
boost::shared_future<boost::shared_ptr<block>> fetch;
Expand All @@ -1176,7 +1197,8 @@ auto_batcher::new_id()
if (b2 && b2 != b) {
int64_t v = b2->next();
if (v != INT64_MIN) {
return boost::make_ready_future(v);
p->set_value(v);
return;
}
}

Expand Down Expand Up @@ -1209,30 +1231,40 @@ auto_batcher::new_id()
fetch = fetch_in_progress_;
}

return fetch
.then(boost::launch::sync,
[this, gen](boost::shared_future<boost::shared_ptr<block>> f)
-> boost::future<int64_t> {
{
std::lock_guard<std::mutex> g(mutex_);
if (fetch_generation_ == gen &&
fetch_in_progress_.valid()) {
fetch_in_progress_ =
boost::shared_future<boost::shared_ptr<block>>();
}
}
// Rethrows a supplier/server error to every coalesced caller.
boost::shared_ptr<block> nb = f.get();
int64_t v = nb->next();
if (v != INT64_MIN) {
return boost::make_ready_future(v);
}
// More concurrent waiters than the batch could serve: retry
// If a caller that cannot get an ID from the fresh batch
// transparently retries (async analogue of Java's for(;;)).
return new_id();

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

took me a minute to get my head around it, but this is the problematic bit.

})
.unwrap();
// The continuation completes the caller's promise directly. Its own
// future is intentionally discarded: nothing waits on it, so a lost
// race never leaves a future layer behind (see try_get_id docs).
fetch.then(
boost::launch::sync,
[this, gen, p](boost::shared_future<boost::shared_ptr<block>> f) {
{
std::lock_guard<std::mutex> g(mutex_);
if (fetch_generation_ == gen && fetch_in_progress_.valid()) {
fetch_in_progress_ =
boost::shared_future<boost::shared_ptr<block>>();
}
}
boost::shared_ptr<block> nb;
try {
// Rethrows a supplier/server error to every coalesced caller.
nb = f.get();
} catch (...) {
p->set_exception(boost::current_exception());
return;
}
int64_t v = nb->next();
if (v != INT64_MIN) {
p->set_value(v);
return;
}
// More concurrent waiters than the batch could serve: retry from
// an empty stack (async analogue of Java's for(;;)). Posting, rather
// than calling try_get_id inline, also covers the case where the
// next fetch is already complete and a sync continuation would run
// immediately on this thread.
boost::asio::post(executor_.get_executor(),
[this, p]() { try_get_id(p); });
});
}

} // namespace impl
Expand Down
31 changes: 31 additions & 0 deletions hazelcast/test/src/HazelcastTests1.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -2583,6 +2583,37 @@ TEST(AutoBatcherTest, concurrencySmokeTest)
ASSERT_LE(supplier->calls(), NUM_THREADS * IDS_IN_THREAD / 3 + 1000);
}

// Regression test for the stack overflow reported in issue #1491. A waiter
// that loses the race for the freshly fetched batch must retry without
// nesting futures (unwrap-on-unwrap): otherwise the completion cascade
// recurses once per lost round and overflows the executor thread's stack
// (macOS gives secondary threads 512 KiB) after a few hundred lost rounds.
//
// A slow supplier makes every thread coalesce onto the in-flight fetch. When
// it completes, the first three continuations win and their threads wake and
// re-register on the next fetch while the remaining continuations are still
// being run, i.e. ahead of them. The last loser therefore loses every round
// until the other threads have used up their quota.
TEST(AutoBatcherTest, starvedWaiterDoesNotOverflowTheStack)
{
constexpr int NUM_THREADS = 128;
constexpr int IDS_IN_THREAD = 100;
util::hz_thread_pool pool(4);
auto supplier = std::make_shared<sequential_batch_supplier>();
impl::auto_batcher batcher(
3, std::chrono::milliseconds(600000), pool, [supplier](int32_t s) {
std::this_thread::sleep_for(std::chrono::milliseconds(1));
return (*supplier)(s);
});

auto ids = concurrently_generate_ids(
[&]() { return batcher.new_id().get(); }, NUM_THREADS, IDS_IN_THREAD);

for (int64_t i = 0; i < static_cast<int64_t>(ids.size()); ++i) {
ASSERT_TRUE(ids.count(i) > 0) << "Missing ID: " << i;
}
}

TEST(FlakeIdBatchTest, getters)
{
impl::id_batch b(5, 7, 9);
Expand Down
Loading