From 012431d302f5bea2decce968a332e807d8eac9ef Mon Sep 17 00:00:00 2001 From: jsonbailey Date: Mon, 5 Oct 2026 17:33:23 -0500 Subject: [PATCH] fix: Reuse an in-flight big segment status poll instead of querying twice The status getter polled the store whenever no status was cached yet, so a request arriving while the startup poll was still in flight started a second metadata query of its own. Status requests now wait on the same lock the poll task holds and re-check the cached status after acquiring it, so a poll that is already running satisfies them. The poll schedule is unchanged: the first poll still fires immediately. Observers are notified outside the lock, because a listener that calls back into the manager would otherwise deadlock on a non-reentrant mutex. --- lib/ldclient-rb/impl/big_segments.rb | 35 +++++++++++++++---- spec/impl/big_segments_spec.rb | 51 ++++++++++++++++++++++++++++ 2 files changed, 79 insertions(+), 7 deletions(-) diff --git a/lib/ldclient-rb/impl/big_segments.rb b/lib/ldclient-rb/impl/big_segments.rb index cdccbeaa..add2c11d 100644 --- a/lib/ldclient-rb/impl/big_segments.rb +++ b/lib/ldclient-rb/impl/big_segments.rb @@ -21,6 +21,7 @@ def initialize(big_segments_config, logger) @status_provider = BigSegmentStoreStatusProviderImpl.new(-> { get_status }) @logger = logger @last_status = nil + @poll_lock = Mutex.new unless @store.nil? @cache = ExpiringCache.new(big_segments_config.context_cache_size, big_segments_config.context_cache_time) @@ -49,18 +50,41 @@ def get_context_membership(context_key) return BigSegmentMembershipResult.new(nil, BigSegmentsStatus::STORE_ERROR) end end - poll_store_and_update_status unless @last_status - unless @last_status.available + status = get_status + unless status.available return BigSegmentMembershipResult.new(membership, BigSegmentsStatus::STORE_ERROR) end - BigSegmentMembershipResult.new(membership, @last_status.stale ? BigSegmentsStatus::STALE : BigSegmentsStatus::HEALTHY) + BigSegmentMembershipResult.new(membership, status.stale ? BigSegmentsStatus::STALE : BigSegmentsStatus::HEALTHY) end def get_status - @last_status || poll_store_and_update_status + status = @last_status + return status if status + + new_status = @poll_lock.synchronize do + # Another caller may have finished a poll while we waited for the lock. + status = @last_status + return status if status + + query_store_status + end + @status_provider.update_status(new_status) + + new_status end def poll_store_and_update_status + new_status = @poll_lock.synchronize { query_store_status } + @status_provider.update_status(new_status) + + new_status + end + + # + # Queries the store and caches the result. Callers hold @poll_lock, so this must not notify + # observers - a listener that calls back into the manager would deadlock on the mutex. + # + private def query_store_status new_status = Interfaces::BigSegmentStoreStatus.new(false, false) # default to "unavailable" if we don't get a new status below unless @store.nil? begin @@ -71,9 +95,6 @@ def poll_store_and_update_status end end @last_status = new_status - @status_provider.update_status(new_status) - - new_status end def stale?(timestamp) diff --git a/spec/impl/big_segments_spec.rb b/spec/impl/big_segments_spec.rb index e384199f..9bf9b090 100644 --- a/spec/impl/big_segments_spec.rb +++ b/spec/impl/big_segments_spec.rb @@ -228,6 +228,57 @@ def next_status_matching(statuses, timeout = 5) next_status_matching(statuses) { |status| !status.stale } end end + + # + # A request that arrives while a poll is already running should wait for that poll instead of + # starting its own. The store signals when it has entered a query and then stays there long + # enough for the request to overlap it, so the race is forced rather than left to timing. + # + context "request during an in-flight poll" do + let(:long_poll_interval) { 30 } + let(:query_time) { 0.3 } + + def slow_metadata_store(queries, query_started) + store = double + allow(store).to receive(:get_metadata) do + queries.increment + query_started.set + sleep(query_time) + always_up_to_date + end + allow(store).to receive(:stop) + store + end + + it "the status getter reuses the result instead of querying again" do + queries = Concurrent::AtomicFixnum.new + query_started = Concurrent::Event.new + store = slow_metadata_store(queries, query_started) + + with_manager(BigSegmentsConfig.new(store: store, status_poll_interval: long_poll_interval)) do |m| + expect(query_started.wait(5)).to be(true) # the startup poll is now inside the store + + expect(m.status_provider.status.available).to be(true) + expect(queries.value).to eq(1) + end + end + + it "a membership query reuses the result instead of querying again" do + queries = Concurrent::AtomicFixnum.new + query_started = Concurrent::Event.new + store = slow_metadata_store(queries, query_started) + expected_membership = { 'key1' => true } + allow(store).to receive(:get_membership).with(context_hash).and_return(expected_membership) + + with_manager(BigSegmentsConfig.new(store: store, status_poll_interval: long_poll_interval)) do |m| + expect(query_started.wait(5)).to be(true) # the startup poll is now inside the store + + expected_result = BigSegmentMembershipResult.new(expected_membership, BigSegmentsStatus::HEALTHY) + expect(m.get_context_membership(context_key)).to eq(expected_result) + expect(queries.value).to eq(1) + end + end + end end end end