diff --git a/lib/flipper/adapters/poll.rb b/lib/flipper/adapters/poll.rb index 20cd9e920..32aa07fdb 100644 --- a/lib/flipper/adapters/poll.rb +++ b/lib/flipper/adapters/poll.rb @@ -1,4 +1,5 @@ require 'flipper/adapters/sync/synchronizer' +require 'flipper/adapters/memory' require 'flipper/poller' module Flipper @@ -7,6 +8,36 @@ class Poll extend Forwardable include ::Flipper::Adapter + class InFlightAdapter + extend Forwardable + + def_delegators :@snapshot, :features, :get, :get_multi, :get_all + def_delegators :@adapter, :add, :remove, :clear, :enable, :disable + + def initialize(snapshot, adapter) + @snapshot = snapshot + @adapter = adapter + end + end + + class PendingSnapshotAdapter + extend Forwardable + + def_delegators :read_adapter, :features, :get, :get_multi, :get_all + def_delegators :@adapter, :add, :remove, :clear, :enable, :disable + + def initialize(snapshot, adapter) + @snapshot = snapshot + @adapter = adapter + end + + private + + def read_adapter + @snapshot.call || @adapter + end + end + # Deprecated Poller = ::Flipper::Poller @@ -17,8 +48,11 @@ class Poll def initialize(poller, adapter) @adapter = adapter @poller = poller + @mutex = Mutex.new + @pid = Process.pid + @syncing = false @last_synced_at = 0 - + @sync_failed = false # If the adapter is empty, we need to sync before starting the poller. # Yes, this will block the main thread, but that's better than thinking # nothing is enabled. @@ -33,6 +67,14 @@ def initialize(poller, adapter) end end + snapshot = begin + build_snapshot(adapter.get_all) + rescue + # Preserve the existing fail-open initialization behavior. The first + # successful request establishes the snapshot before a sync can run. + end + @snapshot = snapshot + @poller.start end @@ -41,11 +83,117 @@ def initialize(poller, adapter) def synced_adapter @poller.start poller_last_synced_at = @poller.last_synced_at.value - if poller_last_synced_at > @last_synced_at - Flipper::Adapters::Sync::Synchronizer.new(@adapter, @poller.adapter).call + claim, value = claim_sync(poller_last_synced_at) + case claim + when :claimed + completed_snapshot = nil + begin + Flipper::Adapters::Sync::Synchronizer.new( + @adapter, + @poller.adapter, + local_get_all: value + ).call + begin + completed_snapshot = build_snapshot(@adapter.get_all) + rescue + # The adapter is synchronized, but its completed state could not + # be captured. Keep serving the previous trusted snapshot to + # contenders and retry publication on the next request. + end + ensure + if completed_snapshot + complete_sync(poller_last_synced_at, completed_snapshot) + else + release_sync + end + end + @adapter + when :syncing, :contended, :snapshot_established + value + else + @adapter + end + end + + # Internal: Attempts to claim the right to sync. Returns the status and + # the data associated with that status: the local state for :claimed or + # the latest coherent adapter for all other read statuses. Never blocks. + # Callers that lose an in-flight claim read from the snapshot rather than + # waiting on a sync in the middle of a request. + def claim_sync(poller_last_synced_at) + unless @mutex.try_lock + adapter = @snapshot || PendingSnapshotAdapter.new(-> { @snapshot }, @adapter) + return [:contended, adapter] + end + + begin + reset_if_forked + return [:syncing, @snapshot] if @syncing + unless @snapshot + local_get_all = @adapter.get_all + established_snapshot = build_snapshot(local_get_all) + @snapshot = established_snapshot + return [:snapshot_established, established_snapshot] + end + return [:not_needed, nil] unless poller_last_synced_at > @last_synced_at + + local_get_all = @adapter.get_all + unless @sync_failed + @snapshot = build_snapshot(local_get_all) + end + @syncing = true + [:claimed, local_get_all] + ensure + @mutex.unlock + end + end + + def complete_sync(poller_last_synced_at, completed_snapshot) + @mutex.synchronize do @last_synced_at = poller_last_synced_at + @snapshot = completed_snapshot + @sync_failed = false + @syncing = false + end + end + + def build_snapshot(local_get_all) + snapshot = Flipper::Adapters::Memory.new(snapshot_copy(local_get_all)) + InFlightAdapter.new(snapshot, @adapter) + end + + def snapshot_copy(value) + case value + when Hash + value.each_with_object({}) do |(key, nested_value), copy| + copy[key] = snapshot_copy(nested_value) + end + when Set + Set.new(value.map { |nested_value| snapshot_copy(nested_value) }) + when Array + value.map { |nested_value| snapshot_copy(nested_value) } + when String + value.dup + else + value + end + end + + def release_sync + @mutex.synchronize do + @sync_failed = true + @syncing = false + end + end + + def reset_if_forked + return if @pid == Process.pid + + @pid = Process.pid + if @syncing + @syncing = false + @sync_failed = true end - @adapter end end end diff --git a/lib/flipper/adapters/sync/interval_synchronizer.rb b/lib/flipper/adapters/sync/interval_synchronizer.rb index f84309d74..f717a481c 100644 --- a/lib/flipper/adapters/sync/interval_synchronizer.rb +++ b/lib/flipper/adapters/sync/interval_synchronizer.rb @@ -21,25 +21,62 @@ def initialize(synchronizer, interval: nil) @interval = interval || DEFAULT_INTERVAL # TODO: add jitter to this so all processes booting at the same time # don't phone home at the same time. + @mutex = Mutex.new + @pid = Process.pid + @syncing = false @last_sync_at = 0 end def call - return unless time_to_sync? + return unless sync_needed? - @last_sync_at = now - @synchronizer.call + begin + @synchronizer.call + ensure + complete_sync + end nil end private - def time_to_sync? - seconds_since_last_sync = now - @last_sync_at + def sync_needed? + return false unless @mutex.try_lock + + begin + reset_if_forked + return false if @syncing + + current_time = now + return false unless time_to_sync?(current_time) + + @last_sync_at = current_time + @syncing = true + true + ensure + @mutex.unlock + end + end + + def complete_sync + @mutex.synchronize do + @syncing = false + end + end + + def time_to_sync?(current_time) + seconds_since_last_sync = current_time - @last_sync_at seconds_since_last_sync >= @interval end + def reset_if_forked + return if @pid == Process.pid + + @pid = Process.pid + @syncing = false + end + def now Process.clock_gettime(Process::CLOCK_MONOTONIC, :second) end diff --git a/lib/flipper/adapters/sync/synchronizer.rb b/lib/flipper/adapters/sync/synchronizer.rb index 7ce58a471..035e6160c 100644 --- a/lib/flipper/adapters/sync/synchronizer.rb +++ b/lib/flipper/adapters/sync/synchronizer.rb @@ -19,12 +19,14 @@ class Synchronizer # :instrumenter - The instrumenter used to instrument. # :raise - Should errors be raised (default: true). # :cache_bust - Should cache busting be used for remote get_all (default: false). + # :local_get_all - Optional pre-fetched local adapter state. def initialize(local, remote, options = {}) @local = local @remote = remote @instrumenter = options.fetch(:instrumenter, Instrumenters::Noop) @raise = options.fetch(:raise, true) @cache_bust = options.fetch(:cache_bust, false) + @local_get_all = options[:local_get_all] end # Public: Forces a sync. @@ -39,7 +41,7 @@ def call private def sync - local_get_all = @local.get_all + local_get_all = @local_get_all || @local.get_all remote_get_all = @remote.get_all(cache_bust: @cache_bust) # Sync all the gate values. diff --git a/spec/flipper/adapters/poll_spec.rb b/spec/flipper/adapters/poll_spec.rb index 2fe08fe75..dd414acd0 100644 --- a/spec/flipper/adapters/poll_spec.rb +++ b/spec/flipper/adapters/poll_spec.rb @@ -1,6 +1,17 @@ require 'flipper/adapters/poll' +require 'open3' +require 'rbconfig' RSpec.describe Flipper::Adapters::Poll do + FakePoller = Struct.new(:last_synced_at, :adapter) do + def start + end + + def sync + raise "sync should not be called when the local adapter is not empty" + end + end + let(:remote_adapter) { adapter = Flipper::Adapters::Memory.new(threadsafe: true) flipper = Flipper.new(adapter) @@ -16,6 +27,10 @@ }) } + def build_poller(adapter) + FakePoller.new(Concurrent::AtomicFixnum.new(1), adapter) + end + it "syncs in main thread if local adapter is empty" do instance = described_class.new(poller, local_adapter) instance.features # call something to force sync @@ -38,4 +53,783 @@ expect(local_adapter.features).to eq(remote_adapter.features) end + + it "establishes a snapshot after a local get_all initialization failure" do + flaky_local_adapter = Class.new(Flipper::Adapters::Memory) do + def initialize + super + @get_all_calls = 0 + end + + def get_all(**kwargs) + @get_all_calls += 1 + raise "transient local failure" if @get_all_calls == 1 + + super + end + end.new + Flipper.new(flaky_local_adapter).enable(:existing) + + fake_poller = build_poller(remote_adapter) + + instance = nil + expect { instance = described_class.new(fake_poller, flaky_local_adapter) }.not_to raise_error + + expect(instance.features).to eq(Set["existing"]) + expect(instance.features).to eq(Set["analytics", "search"]) + end + + it "serves the established snapshot when initialization recovery races with a sync" do + flaky_local_adapter = Class.new(Flipper::Adapters::Memory) do + def initialize + super + @get_all_calls = 0 + end + + def get_all(**kwargs) + @get_all_calls += 1 + raise "transient local failure" if @get_all_calls == 1 + + super + end + end.new + Flipper.new(flaky_local_adapter).enable(:existing) + + fake_poller = build_poller(remote_adapter) + + instance = described_class.new(fake_poller, flaky_local_adapter) + snapshot_established = Queue.new + release_recovery = Queue.new + allow(instance).to receive(:claim_sync).and_wrap_original do |method, *args| + result = method.call(*args) + if Thread.current[:poll_recovery] + snapshot_established << true + release_recovery.pop + end + result + end + + recovery = Thread.new do + Thread.current[:poll_recovery] = true + instance.features + end + snapshot_established.pop + + expect(instance.features).to eq(Set["analytics", "search"]) + release_recovery << true + expect(recovery.value).to eq(Set["existing"]) + ensure + release_recovery << true if release_recovery + recovery&.join + end + + it "only synchronizes once per poller update when called concurrently" do + flipper = Flipper.new(local_adapter) + flipper.enable(:existing) + + get_all_calls = Concurrent::AtomicFixnum.new(0) + slow_remote_adapter = Class.new do + def initialize(result, get_all_calls) + @result = result + @get_all_calls = get_all_calls + end + + def get_all(**kwargs) + @get_all_calls.increment + sleep 0.05 + @result + end + end.new(local_adapter.get_all, get_all_calls) + + fake_poller = build_poller(slow_remote_adapter) + + instance = described_class.new(fake_poller, local_adapter) + threads = 10.times.map { Thread.new { instance.features } } + threads.each(&:join) + + expect(get_all_calls.value).to eq(1) + end + + it "serves a coherent snapshot without waiting for an in-flight poller update" do + flipper = Flipper.new(local_adapter) + flipper.enable(:existing) + + remote = Flipper::Adapters::Memory.new(threadsafe: true) + Flipper.new(remote).enable(:updated) + + entered = Queue.new + release = Queue.new + slow_remote_adapter = Class.new do + def initialize(result, entered, release) + @result = result + @entered = entered + @release = release + end + + def get_all(**kwargs) + @entered << true + @release.pop + @result + end + end.new(remote.get_all, entered, release) + + fake_poller = build_poller(slow_remote_adapter) + + instance = described_class.new(fake_poller, local_adapter) + first_thread = Thread.new { instance.features } + entered.pop + + completed = Queue.new + second_thread = Thread.new { completed << instance.features } + second_thread.join(1) + + # The second thread reads the pre-sync snapshot without waiting for the + # first thread's sync to finish. + expect(completed.pop(true)).to eq(Set["existing"]) + + release << true + expect(first_thread.value).to eq(Set["updated"]) + end + + it "serves a coherent snapshot while a poller update is being applied" do + entered = Queue.new + release = Queue.new + pausing_local_adapter = Class.new(Flipper::Adapters::Memory) do + def initialize(entered, release) + super(nil, threadsafe: true) + @entered = entered + @release = release + @pause = true + end + + def disable(feature, gate, thing) + result = super + if @pause + @pause = false + @entered << true + @release.pop + end + result + end + end.new(entered, release) + Flipper.new(pausing_local_adapter).enable(:existing) + + remote = Flipper::Adapters::Memory.new(threadsafe: true) + Flipper.new(remote).disable(:existing) + Flipper.new(remote).enable(:updated) + + fake_poller = build_poller(remote) + + instance = described_class.new(fake_poller, pausing_local_adapter) + first_thread = Thread.new { instance.features } + entered.pop + + expect(Flipper.new(pausing_local_adapter).enabled?(:existing)).to be(false) + expect(Flipper.new(instance).enabled?(:existing)).to be(true) + + release << true + expect(first_thread.value).to eq(Set["existing", "updated"]) + expect(Flipper.new(instance).enabled?(:existing)).to be(false) + end + + it "keeps mutable gate values isolated in the trusted snapshot" do + entered = Queue.new + release = Queue.new + pausing_local_adapter = Class.new(Flipper::Adapters::Memory) do + def initialize(entered, release) + super(nil, threadsafe: true) + @entered = entered + @release = release + @pause = true + end + + def disable(feature, gate, thing) + result = super + if @pause && gate.data_type == :set + @pause = false + @entered << true + @release.pop + end + result + end + end.new(entered, release) + actor = Flipper::Actor.new("User;1") + Flipper.new(pausing_local_adapter).enable_actor(:search, actor) + + remote = Flipper::Adapters::Memory.new(threadsafe: true) + Flipper.new(remote).enable_percentage_of_actors(:search, 1) + fake_poller = build_poller(remote) + + instance = described_class.new(fake_poller, pausing_local_adapter) + winner = Thread.new { instance.features } + entered.pop + + expect(Flipper.new(pausing_local_adapter).enabled?(:search, actor)).to be(false) + expect(Flipper.new(instance).enabled?(:search, actor)).to be(true) + + release << true + expect(winner.value).to eq(Set["search"]) + expect(Flipper.new(instance).enabled?(:search, actor)).to be(false) + ensure + release << true if release + winner&.join + end + + it "keeps the claimed snapshot after the poller update completes" do + flipper = Flipper.new(local_adapter) + flipper.enable(:existing) + + remote = Flipper::Adapters::Memory.new(threadsafe: true) + Flipper.new(remote).enable(:updated) + + sync_entered = Queue.new + release_sync = Queue.new + slow_remote_adapter = Class.new do + def initialize(result, entered, release) + @result = result + @entered = entered + @release = release + end + + def get_all(**kwargs) + @entered << true + @release.pop + @result + end + end.new(remote.get_all, sync_entered, release_sync) + + fake_poller = build_poller(slow_remote_adapter) + + instance = described_class.new(fake_poller, local_adapter) + loser_claimed = Queue.new + release_loser = Queue.new + allow(instance).to receive(:claim_sync).and_wrap_original do |method, *args| + result = method.call(*args) + if Thread.current[:poll_loser] + loser_claimed << true + release_loser.pop + end + result + end + + winner = Thread.new { instance.features } + sync_entered.pop + loser = Thread.new do + Thread.current[:poll_loser] = true + instance.features + end + loser_claimed.pop + + release_sync << true + expect(winner.value).to eq(Set["updated"]) + release_loser << true + + expect(loser.value).to eq(Set["existing"]) + ensure + release_sync << true if release_sync + release_loser << true if release_loser + winner&.join + loser&.join + end + + it "keeps a coherent snapshot when the sync claim mutex is contended" do + get_all_entered = Queue.new + release_get_all = Queue.new + pausing_local_adapter = Class.new(Flipper::Adapters::Memory) do + def initialize(entered, release) + super(nil, threadsafe: true) + @entered = entered + @release = release + @pause_next_get_all = false + end + + def pause_next_get_all + @pause_next_get_all = true + end + + def get_all(**kwargs) + if @pause_next_get_all + @pause_next_get_all = false + @entered << true + @release.pop + end + super + end + end.new(get_all_entered, release_get_all) + Flipper.new(pausing_local_adapter).enable(:existing) + + remote = Flipper::Adapters::Memory.new(threadsafe: true) + Flipper.new(remote).enable(:updated) + fake_poller = build_poller(remote) + + instance = described_class.new(fake_poller, pausing_local_adapter) + pausing_local_adapter.pause_next_get_all + + loser_claimed = Queue.new + release_loser = Queue.new + allow(instance).to receive(:claim_sync).and_wrap_original do |method, *args| + result = method.call(*args) + if Thread.current[:poll_contender] + loser_claimed << true + release_loser.pop + end + result + end + + winner = Thread.new { instance.features } + get_all_entered.pop + loser = Thread.new do + Thread.current[:poll_contender] = true + instance.features + end + loser_claimed.pop + + release_get_all << true + expect(winner.value).to eq(Set["updated"]) + release_loser << true + + expect(loser.value).to eq(Set["existing"]) + ensure + release_get_all << true if release_get_all + release_loser << true if release_loser + winner&.join + loser&.join + end + + it "retains the completed snapshot for contention during the next poll" do + get_all_entered = Queue.new + release_get_all = Queue.new + pausing_local_adapter = Class.new(Flipper::Adapters::Memory) do + def initialize(entered, release) + super(nil, threadsafe: true) + @entered = entered + @release = release + @pause_next_get_all = false + end + + def pause_next_get_all + @pause_next_get_all = true + end + + def get_all(**kwargs) + if @pause_next_get_all + @pause_next_get_all = false + @entered << true + @release.pop + end + super + end + end.new(get_all_entered, release_get_all) + Flipper.new(pausing_local_adapter).enable(:original) + + remote = Flipper::Adapters::Memory.new(threadsafe: true) + Flipper.new(remote).enable(:first_update) + fake_poller = build_poller(remote) + + instance = described_class.new(fake_poller, pausing_local_adapter) + expect(instance.features).to eq(Set["first_update"]) + + Flipper.new(remote).enable(:second_update) + fake_poller.last_synced_at.value = 2 + pausing_local_adapter.pause_next_get_all + + winner = Thread.new { instance.features } + get_all_entered.pop + contender = Thread.new { instance.features } + + expect(contender.value).to eq(Set["first_update"]) + release_get_all << true + expect(winner.value).to eq(Set["first_update", "second_update"]) + ensure + release_get_all << true if release_get_all + winner&.join + contender&.join + end + + it "reads the local adapter before and after a claimed poller update" do + Flipper.new(local_adapter).enable(:existing) + + fake_poller = build_poller(remote_adapter) + + instance = described_class.new(fake_poller, local_adapter) + expect(local_adapter).to receive(:get_all).twice.and_call_original + + instance.features + end + + it "does not fail a successful update when its completed snapshot cannot be captured" do + flaky_local_adapter = Class.new(Flipper::Adapters::Memory) do + def initialize + super + @get_all_calls = 0 + end + + def get_all(**kwargs) + @get_all_calls += 1 + raise "completed snapshot failure" if @get_all_calls == 3 + + super + end + end.new + Flipper.new(flaky_local_adapter).enable(:existing) + + fake_poller = build_poller(remote_adapter) + + instance = described_class.new(fake_poller, flaky_local_adapter) + + expect { instance.features }.not_to raise_error + expect(instance.features).to eq(Set["analytics", "search"]) + end + + it "does not wait for the sync claim mutex" do + Flipper.new(local_adapter).enable(:existing) + + fake_poller = build_poller(remote_adapter) + + instance = described_class.new(fake_poller, local_adapter) + mutex = instance.instance_variable_get(:@mutex) + locked = Queue.new + release = Queue.new + holder = Thread.new do + mutex.lock + locked << true + release.pop + mutex.unlock + end + locked.pop + + completed = Queue.new + caller = Thread.new { completed << instance.features } + caller.join(1) + + expect(completed.pop(true)).to eq(Set["existing"]) + ensure + release << true if release + holder&.join + caller&.join + end + + it "retries a poller update after synchronization fails" do + flipper = Flipper.new(local_adapter) + flipper.enable(:existing) + + get_all_calls = Concurrent::AtomicFixnum.new(0) + flaky_remote_adapter = Class.new do + def initialize(result, get_all_calls) + @result = result + @get_all_calls = get_all_calls + end + + def get_all(**kwargs) + raise "transient failure" if @get_all_calls.increment == 1 + + @result + end + end.new(local_adapter.get_all, get_all_calls) + + fake_poller = build_poller(flaky_remote_adapter) + + instance = described_class.new(fake_poller, local_adapter) + + expect { instance.features }.to raise_error("transient failure") + instance.features + + expect(get_all_calls.value).to eq(2) + end + + it "retains the last trusted snapshot while retrying a partially failed update" do + failing_local_adapter = Class.new(Flipper::Adapters::Memory) do + def initialize + super(nil, threadsafe: true) + @fail_next_disable = true + end + + def disable(feature, gate, thing) + result = super + if @fail_next_disable + @fail_next_disable = false + raise "partial local failure" + end + result + end + end.new + Flipper.new(failing_local_adapter).enable(:existing) + + remote = Flipper::Adapters::Memory.new(threadsafe: true) + Flipper.new(remote).disable(:existing) + Flipper.new(remote).enable(:updated) + retry_entered = Queue.new + release_retry = Queue.new + pausing_remote_adapter = Class.new do + def initialize(result, entered, release) + @result = result + @entered = entered + @release = release + @get_all_calls = 0 + end + + def get_all(**kwargs) + @get_all_calls += 1 + if @get_all_calls == 2 + @entered << true + @release.pop + end + @result + end + end.new(remote.get_all, retry_entered, release_retry) + + fake_poller = build_poller(pausing_remote_adapter) + + instance = described_class.new(fake_poller, failing_local_adapter) + expect { instance.features }.to raise_error("partial local failure") + expect(Flipper.new(failing_local_adapter).enabled?(:existing)).to be(false) + + retrying = Thread.new { instance.features } + retry_entered.pop + + expect(Flipper.new(instance).enabled?(:existing)).to be(true) + release_retry << true + expect(retrying.value).to eq(Set["existing", "updated"]) + expect(Flipper.new(instance).enabled?(:existing)).to be(false) + ensure + release_retry << true if release_retry + retrying&.join + end + + it "resets in-flight synchronization state after a fork" do + flipper = Flipper.new(local_adapter) + flipper.enable(:existing) + + remote = Flipper::Adapters::Memory.new(threadsafe: true) + Flipper.new(remote).enable(:updated) + + sync_entered = Queue.new + release_sync = Queue.new + get_all_calls = Concurrent::AtomicFixnum.new(0) + counting_remote_adapter = Class.new do + def initialize(result, get_all_calls, sync_entered, release_sync) + @result = result + @get_all_calls = get_all_calls + @sync_entered = sync_entered + @release_sync = release_sync + end + + def get_all(**kwargs) + @get_all_calls.increment + @sync_entered << true + @release_sync.pop + @result + end + end.new(remote.get_all, get_all_calls, sync_entered, release_sync) + + fake_poller = build_poller(counting_remote_adapter) + + instance = described_class.new(fake_poller, local_adapter) + mutex = instance.instance_variable_get(:@mutex) + parent_pid = instance.instance_variable_get(:@pid) + instance.instance_variable_set(:@syncing, true) + + allow(Process).to receive(:pid).and_return(parent_pid + 1) + + winner = Thread.new { instance.features } + sync_entered.pop + + losers = 9.times.map { Thread.new { instance.features } } + losers.each { |thread| expect(thread.join(1)).to equal(thread) } + expect(losers.map(&:value)).to all(eq(Set["existing"])) + expect(get_all_calls.value).to eq(1) + + release_sync << true + expect(winner.value).to eq(Set["updated"]) + expect(instance.instance_variable_get(:@pid)).to eq(parent_pid + 1) + expect(instance.instance_variable_get(:@mutex)).to equal(mutex) + ensure + release_sync << true if release_sync + winner&.join(1) + losers&.each { |thread| thread.join(1) } + end + + it "keeps its mutex and trusted snapshot in a real forked child" do + skip "Process.fork is not supported" unless Process.respond_to?(:fork) + + script = <<~'RUBY' + require "flipper" + require "flipper/adapters/poll" + + def wait_for(queue, timeout: 2) + deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout + loop do + return queue.pop(true) + rescue ThreadError + raise "queue wait timed out" if Process.clock_gettime(Process::CLOCK_MONOTONIC) >= deadline + Thread.pass + end + end + + def join_thread(thread, timeout: 2) + deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout + until thread.join(0.01) + raise "thread join timed out" if Process.clock_gettime(Process::CLOCK_MONOTONIC) >= deadline + end + end + + Poller = Struct.new(:last_synced_at, :adapter) do + def start + end + + def sync + raise "unexpected bootstrap sync" + end + end + + entered = Queue.new + release_sync = Queue.new + remote_calls = 0 + remote = Class.new do + def initialize(result, entered, release_sync, calls) + @result = result + @entered = entered + @release_sync = release_sync + @calls = calls + end + + def get_all(**kwargs) + @calls[0] += 1 + @entered << true + @release_sync.pop + @result + end + end + + local_adapter = Flipper::Adapters::Memory.new(threadsafe: true) + Flipper.new(local_adapter).enable(:existing) + remote_adapter = Flipper::Adapters::Memory.new(threadsafe: true) + Flipper.new(remote_adapter).enable(:updated) + calls = [remote_calls] + poller = Poller.new( + Concurrent::AtomicFixnum.new(1), + remote.new(remote_adapter.get_all, entered, release_sync, calls) + ) + instance = Flipper::Adapters::Poll.new(poller, local_adapter) + mutex = instance.instance_variable_get(:@mutex) + raise "poll does not use a stable Mutex" unless mutex.instance_of?(Mutex) + + Flipper.new(local_adapter).disable(:existing) + instance.instance_variable_set(:@syncing, true) + instance.instance_variable_set(:@sync_failed, false) + + mutex_locked = Queue.new + release_mutex = Queue.new + holder = Thread.new do + mutex.lock + mutex_locked << true + release_mutex.pop + mutex.unlock + end + wait_for(mutex_locked) + + begin + child_pid = fork do + success = false + winner = nil + losers = [] + begin + winner = Thread.new { instance.features } + wait_for(entered) + losers = 8.times.map { Thread.new { instance.features } } + losers.each { |thread| join_thread(thread) } + + raise "inherited mutex was replaced" unless instance.instance_variable_get(:@mutex).equal?(mutex) + raise "loser did not receive trusted snapshot" unless losers.map(&:value).all? { |features| features == Set["existing"] } + raise "duplicate synchronization" unless calls[0] == 1 + + release_sync << true + join_thread(winner) + raise "winner did not publish synchronized state" unless winner.value == Set["updated"] + raise "inherited syncing was not cleared" if instance.instance_variable_get(:@syncing) + raise "successful publication stayed failed" if instance.instance_variable_get(:@sync_failed) + success = true + rescue => error + warn error.full_message + ensure + release_sync << true + winner&.join(1) + losers.each { |thread| thread.join(1) } + end + exit!(success ? 0 : 1) + end + + deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + 5 + status = nil + until status + if result = Process.wait2(child_pid, Process::WNOHANG) + _, status = result + elsif Process.clock_gettime(Process::CLOCK_MONOTONIC) >= deadline + Process.kill("KILL", child_pid) + Process.wait(child_pid) + raise "forked child timed out" + else + sleep 0.01 + end + end + ensure + release_mutex << true + holder.join(1) + end + + exit(status.success? ? 0 : 1) + RUBY + + _, stderr, status = Open3.capture3( + RbConfig.ruby, + "-Ilib", + "-e", + script, + chdir: File.expand_path("../../..", __dir__) + ) + + expect(status).to be_success, stderr + end + + it "retains the trusted snapshot after forking during a partial update" do + Flipper.new(local_adapter).enable(:existing) + + remote = Flipper::Adapters::Memory.new(threadsafe: true) + Flipper.new(remote).disable(:existing) + Flipper.new(remote).enable(:updated) + retry_entered = Queue.new + release_retry = Queue.new + pausing_remote_adapter = Class.new do + def initialize(result, entered, release) + @result = result + @entered = entered + @release = release + end + + def get_all(**kwargs) + @entered << true + @release.pop + @result + end + end.new(remote.get_all, retry_entered, release_retry) + + fake_poller = build_poller(pausing_remote_adapter) + + instance = described_class.new(fake_poller, local_adapter) + parent_pid = instance.instance_variable_get(:@pid) + Flipper.new(local_adapter).disable(:existing) + instance.instance_variable_set(:@syncing, true) + allow(Process).to receive(:pid).and_return(parent_pid + 1) + + retrying = Thread.new { instance.features } + retry_entered.pop + + expect(Flipper.new(instance).enabled?(:existing)).to be(true) + release_retry << true + expect(retrying.value).to eq(Set["existing", "updated"]) + expect(Flipper.new(instance).enabled?(:existing)).to be(false) + ensure + release_retry << true if release_retry + retrying&.join + end end diff --git a/spec/flipper/adapters/sync/interval_synchronizer_spec.rb b/spec/flipper/adapters/sync/interval_synchronizer_spec.rb index e2076c26f..726c5e0e5 100644 --- a/spec/flipper/adapters/sync/interval_synchronizer_spec.rb +++ b/spec/flipper/adapters/sync/interval_synchronizer_spec.rb @@ -1,6 +1,25 @@ require "flipper/adapters/sync/interval_synchronizer" +require "open3" +require "rbconfig" RSpec.describe Flipper::Adapters::Sync::IntervalSynchronizer do + def wait_for(queue, timeout: 2) + deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout + loop do + return queue.pop(true) + rescue ThreadError + raise "queue wait timed out" if Process.clock_gettime(Process::CLOCK_MONOTONIC) >= deadline + Thread.pass + end + end + + def join_thread(thread, timeout: 2) + deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout + until thread.join(0.01) + raise "thread join timed out" if Process.clock_gettime(Process::CLOCK_MONOTONIC) >= deadline + end + end + let(:events) { [] } let(:synchronizer) { -> { events << now } } let(:interval) { 10 } @@ -30,4 +49,225 @@ subject.call expect(events.size).to be(1) end + + it "does not synchronize again while a claimed interval sync is in flight" do + entered = Queue.new + release = Queue.new + synchronizer = -> do + events << now + entered << true + release.pop + end + instance = described_class.new(synchronizer, interval: interval) + + allow(instance).to receive(:now).and_return(interval) + + first_thread = Thread.new { instance.call } + entered.pop + + completed = Queue.new + threads = 10.times.map do + Thread.new do + instance.call + completed << true + end + end + threads.size.times { wait_for(completed) } + + expect(events.size).to eq(1) + + release << true + ([first_thread] + threads).each { |thread| join_thread(thread) } + + expect(events.size).to eq(1) + ensure + 11.times { release << true } if release + ([first_thread] + Array(threads)).compact.each { |thread| thread.join(1) } + end + + it "does not synchronize again when the interval passes during an in-flight sync" do + current_time = interval + entered = Queue.new + release = Queue.new + synchronizer = -> do + events << current_time + entered << true + release.pop + end + instance = described_class.new(synchronizer, interval: interval) + + allow(instance).to receive(:now) { current_time } + + first_thread = Thread.new { instance.call } + entered.pop + + current_time += interval + completed = Queue.new + second_thread = Thread.new do + instance.call + completed << true + end + wait_for(completed) + + expect(events.size).to eq(1) + + release << true + [first_thread, second_thread].each { |thread| join_thread(thread) } + + expect(events.size).to eq(1) + ensure + 2.times { release << true } if release + [first_thread, second_thread].compact.each { |thread| thread.join(1) } + end + + it "releases a failed sync for the next interval" do + current_time = interval + calls = 0 + synchronizer = -> do + calls += 1 + raise "transient failure" if calls == 1 + end + instance = described_class.new(synchronizer, interval: interval) + allow(instance).to receive(:now) { current_time } + + expect { instance.call }.to raise_error("transient failure") + instance.call + expect(calls).to eq(1) + + current_time += interval + instance.call + expect(calls).to eq(2) + end + + it "resets in-flight synchronization state after a fork" do + entered = Queue.new + release = Queue.new + synchronizer = -> do + events << now + entered << true + release.pop + end + instance = described_class.new(synchronizer, interval: interval) + mutex = instance.instance_variable_get(:@mutex) + parent_pid = instance.instance_variable_get(:@pid) + instance.instance_variable_set(:@syncing, true) + + allow(instance).to receive(:now).and_return(interval) + allow(Process).to receive(:pid).and_return(parent_pid + 1) + + threads = 10.times.map { Thread.new { instance.call } } + entered.pop + release << true + threads.each(&:join) + + expect(events.size).to eq(1) + expect(instance.instance_variable_get(:@pid)).to eq(parent_pid + 1) + expect(instance.instance_variable_get(:@mutex)).to equal(mutex) + end + + it "keeps its mutex and elects one winner in a real forked child" do + skip "Process.fork is not supported" unless Process.respond_to?(:fork) + + script = <<~'RUBY' + require "flipper/adapters/sync/interval_synchronizer" + + def wait_for(queue, timeout: 2) + deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout + loop do + return queue.pop(true) + rescue ThreadError + raise "queue wait timed out" if Process.clock_gettime(Process::CLOCK_MONOTONIC) >= deadline + Thread.pass + end + end + + def join_thread(thread, timeout: 2) + deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout + until thread.join(0.01) + raise "thread join timed out" if Process.clock_gettime(Process::CLOCK_MONOTONIC) >= deadline + end + end + + entered = Queue.new + release_sync = Queue.new + calls = [0] + synchronizer = lambda do + calls[0] += 1 + entered << true + release_sync.pop + end + instance = Flipper::Adapters::Sync::IntervalSynchronizer.new(synchronizer, interval: 10) + mutex = instance.instance_variable_get(:@mutex) + raise "interval synchronizer does not use a stable Mutex" unless mutex.instance_of?(Mutex) + instance.instance_variable_set(:@syncing, true) + + mutex_locked = Queue.new + release_mutex = Queue.new + holder = Thread.new do + mutex.lock + mutex_locked << true + release_mutex.pop + mutex.unlock + end + wait_for(mutex_locked) + + begin + child_pid = fork do + success = false + winner = nil + losers = [] + begin + winner = Thread.new { instance.call } + wait_for(entered) + losers = 8.times.map { Thread.new { instance.call } } + losers.each { |thread| join_thread(thread) } + + raise "inherited mutex was replaced" unless instance.instance_variable_get(:@mutex).equal?(mutex) + raise "duplicate synchronization" unless calls[0] == 1 + + release_sync << true + join_thread(winner) + raise "inherited syncing was not cleared" if instance.instance_variable_get(:@syncing) + success = true + rescue => error + warn error.full_message + ensure + release_sync << true + winner&.join(1) + losers.each { |thread| thread.join(1) } + end + exit!(success ? 0 : 1) + end + + deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + 5 + status = nil + until status + if result = Process.wait2(child_pid, Process::WNOHANG) + _, status = result + elsif Process.clock_gettime(Process::CLOCK_MONOTONIC) >= deadline + Process.kill("KILL", child_pid) + Process.wait(child_pid) + raise "forked child timed out" + else + sleep 0.01 + end + end + ensure + release_mutex << true + holder.join(1) + end + + exit(status.success? ? 0 : 1) + RUBY + + _, stderr, status = Open3.capture3( + RbConfig.ruby, + "-Ilib", + "-e", + script, + chdir: File.expand_path("../../../..", __dir__) + ) + + expect(status).to be_success, stderr + end end diff --git a/spec/flipper/adapters/sync/synchronizer_spec.rb b/spec/flipper/adapters/sync/synchronizer_spec.rb index 6e6c19a99..d6713b1c5 100644 --- a/spec/flipper/adapters/sync/synchronizer_spec.rb +++ b/spec/flipper/adapters/sync/synchronizer_spec.rb @@ -56,6 +56,17 @@ expect(instrumenter.events_by_name("synchronizer_exception.flipper").size).to be(0) end + it 'uses pre-fetched local adapter state when provided' do + local_flipper.enable(:existing) + prefetched = local.get_all + remote_flipper.enable(:updated) + expect(local).not_to receive(:get_all) + + described_class.new(local, remote, local_get_all: prefetched).call + + expect(local_flipper.features.map(&:key)).to eq(["updated"]) + end + it 'syncs each remote feature to local' do remote_flipper.enable(:search) remote_flipper.enable_percentage_of_time(:logging, 10)