From e212c5a087d9ecd97cb9eb65951d95ac917786e1 Mon Sep 17 00:00:00 2001 From: deleteLater Date: Tue, 1 Sep 2026 11:38:20 +0800 Subject: [PATCH 1/9] Wait for WebSocket handshake before reconnect reset --- lib/featbit/web_socket_data_synchronizer.rb | 141 +++++++++++++++----- spec/web_socket_data_synchronizer_spec.rb | 78 +++++++++++ 2 files changed, 183 insertions(+), 36 deletions(-) diff --git a/lib/featbit/web_socket_data_synchronizer.rb b/lib/featbit/web_socket_data_synchronizer.rb index 7bd024d..f83666d 100644 --- a/lib/featbit/web_socket_data_synchronizer.rb +++ b/lib/featbit/web_socket_data_synchronizer.rb @@ -7,6 +7,71 @@ require "websocket-client-simple" module FeatBit + class WebSocketConnectionAttempt + attr_reader :socket + + def initialize(timeout:, clock:) + @deadline = clock.call + timeout + @clock = clock + @mutex = Mutex.new + @condition = ConditionVariable.new + @result = nil + end + + def connect(connector, url, headers) + configured = false + @socket = connector.call(url, headers) do |connected_socket| + yield connected_socket + configured = true + end + yield @socket unless configured + @socket + end + + def signal(result) + @mutex.synchronize do + return if @result + + @result = result + @condition.broadcast + end + end + + def wait(stopped:) + @mutex.synchronize do + until @result || stopped.call + remaining = @deadline - @clock.call + return :timeout unless remaining.positive? + + @condition.wait(@mutex, [remaining, 0.05].min) + end + stopped.call ? :stopped : @result + end + end + + def monitor(stopped:, ping:, interval:) + next_ping = @clock.call + interval + until stopped.call || socket_closed? + if @clock.call >= next_ping + ping.call(@socket) + next_ping = @clock.call + interval + end + sleep(0.05) + end + end + + private + + def socket_closed? + return @socket.closed? if @socket.respond_to?(:closed?) + return !@socket.open? if @socket.respond_to?(:open?) + + false + rescue StandardError + true + end + end + class WebSocketDataSynchronizer ALPHABETS = { "0" => "Q", "1" => "B", "2" => "W", "3" => "S", "4" => "P", "5" => "H", "6" => "D", "7" => "X", "8" => "Z", "9" => "U" }.freeze @@ -87,32 +152,32 @@ def process_message(message) def run delay = @options.reconnect_delay until @closed - socket = nil + socket = attempt = nil begin - configured = false - socket = @connector.call(websocket_url, headers) do |connected_socket| - configure_socket(connected_socket) - configured = true + attempt = WebSocketConnectionAttempt.new(timeout: @options.connect_timeout, clock: method(:monotonic_time)) + socket = attempt.connect(@connector, websocket_url, headers) do |connected_socket| + configure_socket(connected_socket, attempt) end @socket_mutex.synchronize { @socket = socket } break if @closed - configure_socket(socket) unless configured - delay = @options.reconnect_delay - next_ping = monotonic_time + PING_INTERVAL - until @closed || socket_closed?(socket) - if monotonic_time >= next_ping - socket.send(JSON.generate(messageType: "ping", data: nil)) - next_ping = monotonic_time + PING_INTERVAL - end - interruptible_sleep(0.05) + connection_result = attempt.wait(stopped: -> { @closed }) + break if connection_result == :stopped + + unless connection_result == :opened + delay = retry_unopened_socket(socket, connection_result, delay) + next end + + delay = @options.reconnect_delay + attempt.monitor(stopped: -> { @closed }, ping: method(:send_ping), interval: PING_INTERVAL) break if @closed @status_provider.update(Status::INTERRUPTED, message: "WebSocket disconnected") interruptible_sleep(delay) delay = [delay * 2, 30.0].min rescue StandardError => e + socket ||= attempt&.socket fail_status(e, interrupted: true) safe_close_socket(socket) interruptible_sleep(delay) unless @closed @@ -124,28 +189,40 @@ def run end end + def retry_unopened_socket(socket, result, delay) + fail_status(Timeout::Error.new("WebSocket handshake timed out"), interrupted: true) if result == :timeout + safe_close_socket(socket) + interruptible_sleep(delay) unless @closed + [delay * 2, 30.0].min + end + + def send_ping(socket) = socket.send(JSON.generate(messageType: "ping", data: nil)) + def connect(url, request_headers, &configure) Timeout.timeout(@options.connect_timeout) do WebSocket::Client::Simple.connect(url, headers: request_headers, &configure) end end - def configure_socket(socket) - open_handler = method(:handle_socket_open) + def configure_socket(socket, attempt = nil) + open_handler = -> { handle_socket_open(socket, attempt) } message_handler = method(:handle_socket_message) - error_handler = method(:handle_socket_error) - close_handler = method(:handle_socket_close) + error_handler = ->(event) { handle_socket_error(socket, event, attempt) } + close_handler = ->(event) { handle_socket_close(event, attempt) } - socket.on(:open) { open_handler.call(socket) } + socket.on(:open) { open_handler.call } socket.on(:message) { |event| message_handler.call(socket, event) } - socket.on(:error) { |event| error_handler.call(socket, event) } + socket.on(:error) { |event| error_handler.call(event) } socket.on(:close) { |event| close_handler.call(event) } end - def handle_socket_open(socket) + def handle_socket_open(socket, attempt = nil) socket.send(JSON.generate(messageType: "data-sync", data: { timestamp: @data_store.version })) + attempt&.signal(:opened) rescue StandardError => e fail_status(e) + attempt&.signal(:failed) + safe_close_socket(socket) end def handle_socket_message(socket, event) @@ -154,26 +231,20 @@ def handle_socket_message(socket, event) safe_close_socket(socket) unless @closed end - def handle_socket_error(socket, event) + def handle_socket_error(socket, event, attempt = nil) return if @closed - fail_status(event.respond_to?(:message) ? event.message : event, interrupted: true) + error = event.respond_to?(:message) ? event.message : event + fail_status(error, interrupted: true) + attempt&.signal(:failed) safe_close_socket(socket) end - def handle_socket_close(_event) + def handle_socket_close(_event, attempt = nil) + attempt&.signal(:closed) @status_provider.update(Status::INTERRUPTED, message: "WebSocket closed") unless @closed end - def socket_closed?(socket) - return socket.closed? if socket.respond_to?(:closed?) - return !socket.open? if socket.respond_to?(:open?) - - false - rescue StandardError - true - end - def safe_close_socket(socket) return true unless socket @@ -193,9 +264,7 @@ def interruptible_sleep(duration) end end - def monotonic_time - Process.clock_gettime(Process::CLOCK_MONOTONIC) - end + def monotonic_time = Process.clock_gettime(Process::CLOCK_MONOTONIC) def process_patch(data) items = Array(fetch(data, "featureFlags", [])).map { |flag| [:flags, flag] } diff --git a/spec/web_socket_data_synchronizer_spec.rb b/spec/web_socket_data_synchronizer_spec.rb index 8a24408..cfbf247 100644 --- a/spec/web_socket_data_synchronizer_spec.rb +++ b/spec/web_socket_data_synchronizer_spec.rb @@ -250,6 +250,84 @@ def close = @closed = true expect(Process.clock_gettime(Process::CLOCK_MONOTONIC) - started).to be < 1 end + it "times out when the WebSocket handshake never opens" do + closed = Queue.new + socket = Class.new do + def initialize(closed) + @closed_event = closed + @handlers = {} + @closed = false + end + + def on(event, &block) = @handlers[event] = block + def closed? = @closed + + def close + return true if @closed + + @closed = true + @closed_event << true + true + end + end.new(closed) + timeout_options = FeatBit::Options.new(env_secret: "secret", connect_timeout: 0.05, reconnect_delay: 10) + synchronizer = described_class.new( + options: timeout_options, + data_store: store, + status_provider: status, + connector: ->(*) { socket } + ) + + synchronizer.start + + expect(Timeout.timeout(2) { closed.pop }).to be(true) + expect(status.status).to eq(FeatBit::Status::INTERRUPTED) + expect(status.message).to eq("WebSocket handshake timed out") + expect(synchronizer.close).to be(true) + end + + it "backs off consecutive failures that happen before the WebSocket opens" do + delays = Queue.new + sockets = [] + connector = lambda do |_url, _headers, &configure| + socket = Class.new do + attr_reader :handlers + + def initialize + @handlers = {} + @closed = false + end + + def on(event, &block) = @handlers[event] = block + def closed? = @closed + def close = @closed = true + end.new + sockets << socket + configure.call(socket) + Thread.new { socket.handlers.fetch(:error).call(StandardError.new("handshake failed")) } + socket + end + synchronizer = described_class.new( + options: options, + data_store: store, + status_provider: status, + connector: connector + ) + delay_count = 0 + allow(synchronizer).to receive(:interruptible_sleep) do |duration| + delay_count += 1 + delays << duration + synchronizer.instance_variable_set(:@closed, true) if delay_count == 2 + end + + synchronizer.start + + expect(Timeout.timeout(2) { delays.pop }).to eq(1.0) + expect(Timeout.timeout(2) { delays.pop }).to eq(2.0) + expect(Timeout.timeout(2) { sleep(0.01) while synchronizer.instance_variable_get(:@thread)&.alive? }).to be_nil + expect(sockets.length).to eq(2) + end + it "closes a failed socket so the reconnect loop can recover" do fake_socket = Class.new do attr_reader :handlers From dd1583390afa3749e2310d3e6043627d043e5d1b Mon Sep 17 00:00:00 2001 From: deleteLater Date: Tue, 1 Sep 2026 11:39:11 +0800 Subject: [PATCH 2/9] Stop reconnecting after WebSocket rejection --- lib/featbit/web_socket_data_synchronizer.rb | 77 ++++++++++++++++-- spec/web_socket_data_synchronizer_spec.rb | 88 +++++++++++++++++++++ 2 files changed, 157 insertions(+), 8 deletions(-) diff --git a/lib/featbit/web_socket_data_synchronizer.rb b/lib/featbit/web_socket_data_synchronizer.rb index f83666d..1592f46 100644 --- a/lib/featbit/web_socket_data_synchronizer.rb +++ b/lib/featbit/web_socket_data_synchronizer.rb @@ -7,6 +7,55 @@ require "websocket-client-simple" module FeatBit + class WebSocketClosePolicy + SERVER_REJECTED_CODE = 4003 + + def initialize(status_provider) + @status_provider = status_provider + @mutex = Mutex.new + @rejected = false + end + + def close_frame?(event) + event.respond_to?(:type) && event.type.to_sym == :close + rescue StandardError + false + end + + def reject?(event) + return false unless close_code(event) == SERVER_REJECTED_CODE + + @mutex.synchronize { @rejected = true } + reason = close_reason(event) + message = "WebSocket connection rejected by server (4003)" + message = "#{message}: #{reason}" unless reason.empty? + @status_provider.update(Status::FAILED, message: message) + true + end + + def rejected? + @mutex.synchronize { @rejected } + end + + private + + def close_code(event) + return event.code.to_i if event.respond_to?(:code) && event.code + return event[:code].to_i if event.is_a?(Hash) && event[:code] + + event["code"].to_i if event.is_a?(Hash) && event["code"] + end + + def close_reason(event) + value = if event.respond_to?(:reason) && event.reason + event.reason + elsif event.respond_to?(:data) && event.data + event.data + end + value.to_s + end + end + class WebSocketConnectionAttempt attr_reader :socket @@ -81,6 +130,7 @@ def initialize(options:, data_store:, status_provider:, on_flags_changed: nil, c @options = options @data_store = data_store @status_provider = status_provider + @close_policy = WebSocketClosePolicy.new(status_provider) @on_flags_changed = on_flags_changed @connector = connector || method(:connect) @closed = false @@ -151,7 +201,7 @@ def process_message(message) def run delay = @options.reconnect_delay - until @closed + until @closed || @close_policy.rejected? socket = attempt = nil begin attempt = WebSocketConnectionAttempt.new(timeout: @options.connect_timeout, clock: method(:monotonic_time)) @@ -159,10 +209,10 @@ def run configure_socket(connected_socket, attempt) end @socket_mutex.synchronize { @socket = socket } - break if @closed + break if @closed || @close_policy.rejected? connection_result = attempt.wait(stopped: -> { @closed }) - break if connection_result == :stopped + break if %i[stopped rejected].include?(connection_result) unless connection_result == :opened delay = retry_unopened_socket(socket, connection_result, delay) @@ -171,12 +221,14 @@ def run delay = @options.reconnect_delay attempt.monitor(stopped: -> { @closed }, ping: method(:send_ping), interval: PING_INTERVAL) - break if @closed + break if @closed || @close_policy.rejected? @status_provider.update(Status::INTERRUPTED, message: "WebSocket disconnected") interruptible_sleep(delay) delay = [delay * 2, 30.0].min rescue StandardError => e + break if @close_policy.rejected? + socket ||= attempt&.socket fail_status(e, interrupted: true) safe_close_socket(socket) @@ -226,13 +278,19 @@ def handle_socket_open(socket, attempt = nil) end def handle_socket_message(socket, event) + if @close_policy.close_frame?(event) + handle_socket_close(event) + safe_close_socket(socket) + return + end + return if process_message(event.respond_to?(:data) ? event.data : event.to_s) safe_close_socket(socket) unless @closed end def handle_socket_error(socket, event, attempt = nil) - return if @closed + return if @closed || @close_policy.rejected? error = event.respond_to?(:message) ? event.message : event fail_status(error, interrupted: true) @@ -240,9 +298,12 @@ def handle_socket_error(socket, event, attempt = nil) safe_close_socket(socket) end - def handle_socket_close(_event, attempt = nil) - attempt&.signal(:closed) - @status_provider.update(Status::INTERRUPTED, message: "WebSocket closed") unless @closed + def handle_socket_close(event, attempt = nil) + rejected = @close_policy.reject?(event) + attempt&.signal(rejected ? :rejected : :closed) + return if @closed || rejected || @close_policy.rejected? + + @status_provider.update(Status::INTERRUPTED, message: "WebSocket closed") end def safe_close_socket(socket) diff --git a/spec/web_socket_data_synchronizer_spec.rb b/spec/web_socket_data_synchronizer_spec.rb index cfbf247..1a316cd 100644 --- a/spec/web_socket_data_synchronizer_spec.rb +++ b/spec/web_socket_data_synchronizer_spec.rb @@ -361,6 +361,94 @@ def close = @closed = true expect(synchronizer.close).to be(true) end + it "stops reconnecting when the server rejects the connection with close code 4003" do + raw_frame = WebSocket::Frame::Outgoing::Server.new( + version: 13, type: :close, code: 4003, data: "invalid environment secret" + ).to_s + close_frame = WebSocket::Frame::Incoming::Client.new(version: 13, data: raw_frame).next + attempts = 0 + socket = Class.new do + attr_reader :handlers + + def initialize + @handlers = {} + @closed = false + end + + def on(event, &block) = @handlers[event] = block + def send(*) = nil + def closed? = @closed + def close = @closed = true + end.new + connector = lambda do |_url, _headers, &configure| + attempts += 1 + configure.call(socket) + socket + end + rejection_options = FeatBit::Options.new(env_secret: "secret", reconnect_delay: 0.01) + synchronizer = described_class.new( + options: rejection_options, + data_store: store, + status_provider: status, + connector: connector + ) + + synchronizer.start + Timeout.timeout(2) { sleep(0.01) until socket.handlers.key?(:message) } + socket.handlers.fetch(:open).call + socket.handlers.fetch(:message).call(close_frame) + Timeout.timeout(2) { sleep(0.01) while synchronizer.instance_variable_get(:@thread)&.alive? } + + expect(attempts).to eq(1) + expect(socket).to be_closed + expect(status.status).to eq(FeatBit::Status::FAILED) + expect(status.message).to eq("WebSocket connection rejected by server (4003): invalid environment secret") + expect(synchronizer.close).to be(true) + end + + it "reconnects after a non-rejected close frame" do + raw_frame = WebSocket::Frame::Outgoing::Server.new( + version: 13, type: :close, code: 1000, data: "service restart" + ).to_s + close_frame = WebSocket::Frame::Incoming::Client.new(version: 13, data: raw_frame).next + connected = Queue.new + connector = lambda do |_url, _headers, &configure| + socket = Class.new do + attr_reader :handlers + + def initialize + @handlers = {} + @closed = false + end + + def on(event, &block) = @handlers[event] = block + def send(*) = nil + def closed? = @closed + def close = @closed = true + end.new + configure.call(socket) + connected << socket + socket + end + reconnect_options = FeatBit::Options.new(env_secret: "secret", reconnect_delay: 0.01) + synchronizer = described_class.new( + options: reconnect_options, + data_store: store, + status_provider: status, + connector: connector + ) + + synchronizer.start + first_socket = Timeout.timeout(2) { connected.pop } + first_socket.handlers.fetch(:open).call + first_socket.handlers.fetch(:message).call(close_frame) + second_socket = Timeout.timeout(2) { connected.pop } + + expect(first_socket).to be_closed + expect(second_socket).not_to be_closed + expect(synchronizer.close).to be(true) + end + it "encodes a current timestamp into the connection token without falling back to the secret" do secret = "abcdefghijklmnopqrstuvwxyz0123456789ABCDEFG=" synchronizer = described_class.new(options: options, data_store: store, status_provider: status) From cd71512742e3f522563034aaad5f787a14cbabdb Mon Sep 17 00:00:00 2001 From: deleteLater Date: Wed, 2 Sep 2026 10:38:26 +0800 Subject: [PATCH 3/9] Return synchronization result details for unchanged data --- lib/featbit/web_socket_data_synchronizer.rb | 44 +++++++++++++++------ spec/web_socket_data_synchronizer_spec.rb | 36 +++++++++++++++-- 2 files changed, 64 insertions(+), 16 deletions(-) diff --git a/lib/featbit/web_socket_data_synchronizer.rb b/lib/featbit/web_socket_data_synchronizer.rb index 1592f46..4c49f09 100644 --- a/lib/featbit/web_socket_data_synchronizer.rb +++ b/lib/featbit/web_socket_data_synchronizer.rb @@ -7,6 +7,27 @@ require "websocket-client-simple" module FeatBit + class SynchronizationResult + attr_reader :valid, :changed + + def initialize(valid:, changed:) + @valid = valid + @changed = changed + freeze + end + + def valid? = valid + def changed? = changed + + INVALID = new(valid: false, changed: false) + UNCHANGED = new(valid: true, changed: false) + CHANGED = new(valid: true, changed: true) + + def self.valid(changed:) + changed ? CHANGED : UNCHANGED + end + end + class WebSocketClosePolicy SERVER_REJECTED_CODE = 4003 @@ -170,11 +191,11 @@ def close def process_message(message) envelope = message.is_a?(String) ? JSON.parse(message) : message - return false unless fetch(envelope, "messageType") == "data-sync" + return SynchronizationResult::INVALID unless fetch(envelope, "messageType") == "data-sync" data = fetch(envelope, "data", {}) event_type = fetch(data, "eventType") - return false unless valid_data?(data, event_type) + return SynchronizationResult::INVALID unless valid_data?(data, event_type) old_keys = @data_store.all_flags.keys if event_type == "full" @@ -185,16 +206,14 @@ def process_message(message) end changed_keys.each { |key| safely_notify(key) } - return false unless changed - @status_provider.update(Status::READY) - true + SynchronizationResult.valid(changed: changed) rescue JSON::ParserError => e @status_provider.update(Status::FAILED, message: "invalid data: #{e.message}") - false + SynchronizationResult::INVALID rescue StandardError => e fail_status(e) - false + SynchronizationResult::INVALID end private @@ -284,7 +303,8 @@ def handle_socket_message(socket, event) return end - return if process_message(event.respond_to?(:data) ? event.data : event.to_s) + result = process_message(event.respond_to?(:data) ? event.data : event.to_s) + return if result.valid? safe_close_socket(socket) unless @closed end @@ -330,16 +350,18 @@ def monotonic_time = Process.clock_gettime(Process::CLOCK_MONOTONIC) def process_patch(data) items = Array(fetch(data, "featureFlags", [])).map { |flag| [:flags, flag] } items.concat(Array(fetch(data, "segments", [])).map { |segment| [:segments, segment] }) + changed = false changed_keys = [] items.sort_by { |_kind, item| item_version(item) }.each do |kind, item| applied = @data_store.upsert(kind, item, version: item_version(item)) + changed = true if applied if applied && kind == :flags changed_keys << fetch(item, "key").to_s elsif applied changed_keys.concat(flags_referencing_segment(fetch(item, "id").to_s)) end end - [true, changed_keys.uniq] + [changed, changed_keys.uniq] end def valid_data?(data, event_type) @@ -378,9 +400,7 @@ def item_version(item) @data_store.version + 1 end - def websocket_url - "#{@options.streaming_uri}?token=#{build_token(@options.env_secret)}&type=server" - end + def websocket_url = "#{@options.streaming_uri}?token=#{build_token(@options.env_secret)}&type=server" def headers { diff --git a/spec/web_socket_data_synchronizer_spec.rb b/spec/web_socket_data_synchronizer_spec.rb index 1a316cd..402104b 100644 --- a/spec/web_socket_data_synchronizer_spec.rb +++ b/spec/web_socket_data_synchronizer_spec.rb @@ -13,7 +13,9 @@ synchronizer = described_class.new(options: options, data_store: store, status_provider: status, on_flags_changed: lambda { |key| changes << key }) - expect(synchronizer.process_message(test_bootstrap(test_flag))).to be(true) + result = synchronizer.process_message(test_bootstrap(test_flag)) + expect(result).to be_valid + expect(result).to be_changed expect(store.flag("welcome")).not_to be_nil expect(status.status).to eq(FeatBit::Status::READY) expect(changes).to include("welcome") @@ -21,7 +23,9 @@ it "rejects malformed data without raising" do synchronizer = described_class.new(options: options, data_store: store, status_provider: status) - expect { synchronizer.process_message("not-json") }.not_to raise_error + result = synchronizer.process_message("not-json") + expect(result).not_to be_valid + expect(result).not_to be_changed expect(status.status).to eq(FeatBit::Status::FAILED) end @@ -42,7 +46,9 @@ "messageType" => "data-sync", "data" => { "eventType" => "patch", "featureFlags" => [patched], "segments" => [] } } - expect(synchronizer.process_message(message)).to be(true) + result = synchronizer.process_message(message) + expect(result).to be_valid + expect(result).to be_changed expect(store.flag("welcome")["name"]).to eq("updated") expect(changes).to eq(["welcome"]) end @@ -65,7 +71,9 @@ "data" => { "eventType" => "patch", "featureFlags" => [stale_flag, fresh_flag], "segments" => [] } } - expect(synchronizer.process_message(message)).to be(true) + result = synchronizer.process_message(message) + expect(result).to be_valid + expect(result).to be_changed expect(store.flag("fresh")).not_to be_nil expect(changes).to eq(["fresh"]) end @@ -81,10 +89,30 @@ socket = instance_double("Socket", close: true) synchronizer = described_class.new(options: options, data_store: store, status_provider: status) + result = synchronizer.process_message(message) synchronizer.send(:handle_socket_message, socket, JSON.generate(message)) + expect(result).to be_valid + expect(result).not_to be_changed expect(socket).not_to have_received(:close) expect(store.version).to eq(10) + expect(status.status).to eq(FeatBit::Status::READY) + end + + it "accepts an unchanged full synchronization and reports ready" do + message = test_bootstrap(test_flag) + expect(store.init(message)).to be(true) + status.update(FeatBit::Status::INTERRUPTED, message: "disconnected") + socket = instance_double("Socket", close: true) + synchronizer = described_class.new(options: options, data_store: store, status_provider: status) + + result = synchronizer.process_message(message) + synchronizer.send(:handle_socket_message, socket, JSON.generate(message)) + + expect(result).to be_valid + expect(result).not_to be_changed + expect(socket).not_to have_received(:close) + expect(status.status).to eq(FeatBit::Status::READY) end it "rejects an entire malformed patch before applying valid siblings" do From 9cec5baafe0ab309b7ab69ab92d1103599ae6cb9 Mon Sep 17 00:00:00 2001 From: deleteLater Date: Wed, 2 Sep 2026 11:06:17 +0800 Subject: [PATCH 4/9] Make WebSocket shutdown completion reliable --- README.md | 1 + lib/featbit/client.rb | 8 +- lib/featbit/web_socket_data_synchronizer.rb | 186 ++++++------ lib/featbit/web_socket_lifecycle.rb | 110 +++++++ spec/client_spec.rb | 58 +++- spec/web_socket_data_synchronizer_spec.rb | 8 +- spec/web_socket_shutdown_spec.rb | 302 ++++++++++++++++++++ 7 files changed, 570 insertions(+), 103 deletions(-) create mode 100644 lib/featbit/web_socket_lifecycle.rb create mode 100644 spec/web_socket_shutdown_spec.rb diff --git a/README.md b/README.md index f25efa7..2294b84 100644 --- a/README.md +++ b/README.md @@ -230,6 +230,7 @@ Call `track` after evaluating the related experiment flag. `numeric_value` defau - Event enqueue is non-blocking and bounded; overload increments `dropped_events`. - Status and flag callbacks execute outside internal locks. - `close` is thread-safe and idempotent, flushes accepted events, stops background workers, and prevents new events. +- A successful `close` confirms WebSocket connection work and callbacks have finished. It waits up to five seconds for the synchronizer; blocked application callbacks or cleanup failures return `false`, which may be retried. A concurrent or callback-initiated `close` can also return `false` while shutdown is incomplete; call again from outside callbacks to wait for completion. The `closed` status means shutdown was requested, not necessarily that cleanup succeeded. - Public client methods contain ordinary internal failures and return fallbacks or `false`. ## Supported Ruby versions diff --git a/lib/featbit/client.rb b/lib/featbit/client.rb index d5a73fb..170fd11 100644 --- a/lib/featbit/client.rb +++ b/lib/featbit/client.rb @@ -122,10 +122,12 @@ def remove_flag_change_listener(id) end def close + owns_close = false @lifecycle_mutex.synchronize do - return @close_result unless @close_result.nil? - return true if @closed + return true if @close_result == true + return false if @closing + @closing = owns_close = true @closed = true end @@ -135,6 +137,8 @@ def close rescue StandardError => e safe_log(:warn, "FeatBit client close failed: #{e.message}") false + ensure + @lifecycle_mutex.synchronize { @closing = false } if owns_close end alias stop close diff --git a/lib/featbit/web_socket_data_synchronizer.rb b/lib/featbit/web_socket_data_synchronizer.rb index 4c49f09..b3c7a33 100644 --- a/lib/featbit/web_socket_data_synchronizer.rb +++ b/lib/featbit/web_socket_data_synchronizer.rb @@ -5,6 +5,7 @@ require "time" require "timeout" require "websocket-client-simple" +require_relative "web_socket_lifecycle" module FeatBit class SynchronizationResult @@ -80,26 +81,51 @@ def close_reason(event) class WebSocketConnectionAttempt attr_reader :socket - def initialize(timeout:, clock:) + def initialize(timeout:, clock:, lifecycle:, closer:) @deadline = clock.call + timeout @clock = clock @mutex = Mutex.new @condition = ConditionVariable.new @result = nil + @lifecycle = lifecycle + @closer = closer + @finished = false end def connect(connector, url, headers) - configured = false - @socket = connector.call(url, headers) do |connected_socket| - yield connected_socket - configured = true + @connector_thread = Thread.new do + configured = false + @socket = connector.call(url, headers) do |connected_socket| + @socket = connected_socket + yield connected_socket + configured = true + end + yield @socket unless configured + rescue StandardError => e + @connect_error = e + end + until @connector_thread.join(0.01) + break if @lifecycle.stopped? || @clock.call >= @deadline end - yield @socket unless configured + return @socket if @lifecycle.stopped? + + raise Timeout::Error, "WebSocket connection timed out" if @connector_thread.alive? + raise @connect_error if @connect_error + @socket end + def cleanup + # Cancellation is confined to the owned connector, never an application + # callback. Joining before close prevents a late connector creating a socket. + @connector_thread&.kill if @connector_thread&.alive? + @connector_thread&.join + @closer.call(@socket) + end + def signal(result) @mutex.synchronize do + @finished = true unless result == :opened return if @result @result = result @@ -121,7 +147,7 @@ def wait(stopped:) def monitor(stopped:, ping:, interval:) next_ping = @clock.call + interval - until stopped.call || socket_closed? + until stopped.call || @mutex.synchronize { @finished } || socket_closed? if @clock.call >= next_ping ping.call(@socket) next_ping = @clock.call + interval @@ -154,42 +180,26 @@ def initialize(options:, data_store:, status_provider:, on_flags_changed: nil, c @close_policy = WebSocketClosePolicy.new(status_provider) @on_flags_changed = on_flags_changed @connector = connector || method(:connect) - @closed = false - @socket_mutex = Mutex.new - @lifecycle_mutex = Mutex.new - @thread = nil - @close_result = nil + @lifecycle = WebSocketLifecycle.new end def start - @lifecycle_mutex.synchronize do - return false if @closed || @thread&.alive? - - @thread = Thread.new { run } - @thread.name = "featbit-websocket-sync" if @thread.respond_to?(:name=) - true - end + @lifecycle.start { run } rescue StandardError => e fail_status(e) false end def close - @lifecycle_mutex.synchronize do - return @close_result unless @close_result.nil? - - @closed = true - socket = @socket_mutex.synchronize { @socket } - socket_closed = safe_close_socket(socket) - @thread&.join(5) unless Thread.current.equal?(@thread) - @close_result = socket_closed && !@thread&.alive? - end + @lifecycle.close rescue StandardError => e @options.logger&.warn("FeatBit synchronizer close failed: #{e.message}") false end def process_message(message) + return SynchronizationResult::INVALID if @lifecycle.stopped? + envelope = message.is_a?(String) ? JSON.parse(message) : message return SynchronizationResult::INVALID unless fetch(envelope, "messageType") == "data-sync" @@ -205,8 +215,8 @@ def process_message(message) changed, changed_keys = process_patch(data) end - changed_keys.each { |key| safely_notify(key) } - @status_provider.update(Status::READY) + changed_keys.each { |key| safely_notify(key) unless @lifecycle.stopped? } + @status_provider.update(Status::READY) unless @lifecycle.stopped? SynchronizationResult.valid(changed: changed) rescue JSON::ParserError => e @status_provider.update(Status::FAILED, message: "invalid data: #{e.message}") @@ -220,71 +230,56 @@ def process_message(message) def run delay = @options.reconnect_delay - until @closed || @close_policy.rejected? - socket = attempt = nil - begin - attempt = WebSocketConnectionAttempt.new(timeout: @options.connect_timeout, clock: method(:monotonic_time)) - socket = attempt.connect(@connector, websocket_url, headers) do |connected_socket| - configure_socket(connected_socket, attempt) - end - @socket_mutex.synchronize { @socket = socket } - break if @closed || @close_policy.rejected? - - connection_result = attempt.wait(stopped: -> { @closed }) - break if %i[stopped rejected].include?(connection_result) - - unless connection_result == :opened - delay = retry_unopened_socket(socket, connection_result, delay) - next - end + until @lifecycle.stopped? || @close_policy.rejected? + opened = run_attempt + break if @lifecycle.stopped? || @close_policy.rejected? - delay = @options.reconnect_delay - attempt.monitor(stopped: -> { @closed }, ping: method(:send_ping), interval: PING_INTERVAL) - break if @closed || @close_policy.rejected? - - @status_provider.update(Status::INTERRUPTED, message: "WebSocket disconnected") - interruptible_sleep(delay) - delay = [delay * 2, 30.0].min - rescue StandardError => e - break if @close_policy.rejected? - - socket ||= attempt&.socket - fail_status(e, interrupted: true) - safe_close_socket(socket) - interruptible_sleep(delay) unless @closed - delay = [delay * 2, 30.0].min - ensure - safe_close_socket(socket) if @closed - @socket_mutex.synchronize { @socket = nil } - end + delay = @options.reconnect_delay if opened + interruptible_sleep(delay) + delay = [delay * 2, 30.0].min end end - def retry_unopened_socket(socket, result, delay) - fail_status(Timeout::Error.new("WebSocket handshake timed out"), interrupted: true) if result == :timeout - safe_close_socket(socket) - interruptible_sleep(delay) unless @closed - [delay * 2, 30.0].min + def run_attempt + attempt = WebSocketConnectionAttempt.new( + timeout: @options.connect_timeout, clock: method(:monotonic_time), lifecycle: @lifecycle, closer: method(:safe_close_socket) + ) + @lifecycle.activate(attempt) + attempt.connect(@connector, websocket_url, headers) { |socket| configure_socket(socket, attempt) } + result = attempt.wait(stopped: @lifecycle.method(:stopped?)) + raise Timeout::Error, "WebSocket handshake timed out" if result == :timeout + return false unless result == :opened + + attempt.monitor(stopped: @lifecycle.method(:stopped?), ping: method(:send_ping), interval: PING_INTERVAL) + handle_socket_close(nil) unless @lifecycle.stopped? || @close_policy.rejected? + true + rescue StandardError => e + fail_status(e, interrupted: true) unless @close_policy.rejected? + false + ensure + @lifecycle.finish(attempt) if attempt end def send_ping(socket) = socket.send(JSON.generate(messageType: "ping", data: nil)) def connect(url, request_headers, &configure) - Timeout.timeout(@options.connect_timeout) do - WebSocket::Client::Simple.connect(url, headers: request_headers, &configure) - end + socket = ClosableWebSocketClient.new + configure.call(socket) + socket.connect(url, headers: request_headers) + socket end def configure_socket(socket, attempt = nil) - open_handler = -> { handle_socket_open(socket, attempt) } - message_handler = method(:handle_socket_message) - error_handler = ->(event) { handle_socket_error(socket, event, attempt) } - close_handler = ->(event) { handle_socket_close(event, attempt) } - - socket.on(:open) { open_handler.call } - socket.on(:message) { |event| message_handler.call(socket, event) } - socket.on(:error) { |event| error_handler.call(event) } - socket.on(:close) { |event| close_handler.call(event) } + lifecycle = @lifecycle + handlers = { + open: ->(_) { handle_socket_open(socket, attempt) }, + message: ->(event) { handle_socket_message(socket, event, attempt) }, + error: ->(event) { handle_socket_error(socket, event, attempt) }, + close: ->(event) { handle_socket_close(event, attempt) } + } + handlers.each do |name, handler| + socket.on(name) { |event| lifecycle.dispatch(attempt) { handler.call(event) } } + end end def handle_socket_open(socket, attempt = nil) @@ -293,35 +288,38 @@ def handle_socket_open(socket, attempt = nil) rescue StandardError => e fail_status(e) attempt&.signal(:failed) - safe_close_socket(socket) + safe_close_socket(socket) unless attempt end - def handle_socket_message(socket, event) + def handle_socket_message(socket, event, attempt = nil) if @close_policy.close_frame?(event) - handle_socket_close(event) - safe_close_socket(socket) + handle_socket_close(event, attempt) + safe_close_socket(socket) unless attempt return end result = process_message(event.respond_to?(:data) ? event.data : event.to_s) return if result.valid? - safe_close_socket(socket) unless @closed + attempt&.signal(:failed) + safe_close_socket(socket) unless attempt || @lifecycle.stopped? end def handle_socket_error(socket, event, attempt = nil) - return if @closed || @close_policy.rejected? + return if @lifecycle.stopped? || @close_policy.rejected? error = event.respond_to?(:message) ? event.message : event fail_status(error, interrupted: true) attempt&.signal(:failed) - safe_close_socket(socket) + safe_close_socket(socket) unless attempt end def handle_socket_close(event, attempt = nil) + return if @lifecycle.stopped? + rejected = @close_policy.reject?(event) attempt&.signal(rejected ? :rejected : :closed) - return if @closed || rejected || @close_policy.rejected? + return if rejected || @close_policy.rejected? @status_provider.update(Status::INTERRUPTED, message: "WebSocket closed") end @@ -329,7 +327,7 @@ def handle_socket_close(event, attempt = nil) def safe_close_socket(socket) return true unless socket - socket.close != false + Timeout.timeout(2) { (socket.is_a?(ClosableWebSocketClient) ? socket.close(drain: true) : socket.close) != false } rescue StandardError => e @options.logger&.warn("FeatBit WebSocket close failed: #{e.message}") false @@ -337,7 +335,7 @@ def safe_close_socket(socket) def interruptible_sleep(duration) deadline = monotonic_time + duration.to_f - until @closed + until @lifecycle.stopped? remaining = deadline - monotonic_time break unless remaining.positive? @@ -431,6 +429,8 @@ def safely_notify(flag_key) end def fail_status(error, interrupted: false) + return if @lifecycle.stopped? + message = error.respond_to?(:message) ? error.message : error.to_s @options.logger&.warn("FeatBit WebSocket synchronization failed: #{message}") @status_provider.update(interrupted ? Status::INTERRUPTED : Status::FAILED, message: message) diff --git a/lib/featbit/web_socket_lifecycle.rb b/lib/featbit/web_socket_lifecycle.rb new file mode 100644 index 0000000..4381ef4 --- /dev/null +++ b/lib/featbit/web_socket_lifecycle.rb @@ -0,0 +1,110 @@ +# frozen_string_literal: true + +require "timeout" +require "websocket-client-simple" + +module FeatBit + # Owns shutdown and drains callbacks without holding locks around user code. + class WebSocketLifecycle + CLOSE_WAIT = 5.0 + + def initialize + @mutex = Mutex.new + @condition = ConditionVariable.new + @callbacks = Hash.new(0) + @stopped = false + @clean = true + @attempt = nil + end + + def start(&block) + @mutex.synchronize do + return false if @stopped || @thread&.alive? + + @thread = Thread.new(&block) + @thread.name = "featbit-websocket-sync" if @thread.respond_to?(:name=) + true + end + end + + def stopped? + @mutex.synchronize { @stopped } + end + + def close + thread, in_callback = @mutex.synchronize do + @stopped = true + [@thread, @callbacks.key?(Thread.current)] + end + return false if in_callback || Thread.current.equal?(thread) + + thread&.join(CLOSE_WAIT) + @mutex.synchronize { !thread&.alive? && @clean } + end + + def activate(attempt) + @mutex.synchronize { @attempt = attempt } + end + + def dispatch(attempt) + entered = @mutex.synchronize do + unless @stopped || !@attempt.equal?(attempt) + @callbacks[Thread.current] += 1 + true + end + end + yield if entered + ensure + if entered + @mutex.synchronize do + @callbacks[Thread.current] -= 1 + @callbacks.delete(Thread.current) if @callbacks[Thread.current].zero? + @condition.broadcast + end + end + end + + def finish(attempt) + @mutex.synchronize do + @attempt = nil + @condition.wait(@mutex) until @callbacks.empty? + end + clean = attempt.cleanup + @mutex.synchronize do + @clean &&= clean + @stopped = true unless clean + end + end + end + + # websocket-client-simple has no abort API. Keep its internal cleanup here, + # after the connector has exited and SDK callbacks have drained. + class ClosableWebSocketClient < WebSocket::Client::Simple::Client + CLOSE_TIMEOUT = 1.0 + + def close(drain: false) + # The library also calls close from its own I/O error callbacks. Only the + # SDK owner may drain/kill the reader, after application callbacks return. + unless drain + @closed = true + return true + end + + begin + Timeout.timeout(CLOSE_TIMEOUT) { super() } + rescue Timeout::Error + # A stalled close-frame write must not prevent closing the TCP socket. + ensure + @closed = true + begin + @socket&.close + @socket = nil + ensure + @thread&.kill + @thread&.join unless Thread.current.equal?(@thread) + end + end + true + end + end +end diff --git a/spec/client_spec.rb b/spec/client_spec.rb index 928bd7e..8a60190 100644 --- a/spec/client_spec.rb +++ b/spec/client_spec.rb @@ -165,11 +165,15 @@ def broken_store.segment(_key) = raise("boom") expect(client.close).to be(true) end - it "allows component shutdown to re-enter close without locking" do + it "reports a reentrant close as incomplete without deadlocking" do client = nil + reentrant_results = [] synchronizer = instance_double("Synchronizer", start: true, close: true) processor = Object.new - processor.define_singleton_method(:close) { client.close } + processor.define_singleton_method(:close) do + reentrant_results << client.close + true + end options = FeatBit::Options.new( env_secret: "secret", start_wait: 0.001, @@ -179,9 +183,10 @@ def broken_store.segment(_key) = raise("boom") client = described_class.new(options) expect(client.close).to be(true) + expect(reentrant_results).to eq([false]) end - it "attempts to close every component once and never raises" do + it "retries failed shutdown without skipping components or raising" do synchronizer = instance_double("Synchronizer", start: true) processor = instance_double("EventProcessor") allow(synchronizer).to receive(:close).and_raise("synchronizer failed") @@ -196,8 +201,51 @@ def broken_store.segment(_key) = raise("boom") expect(client.close).to be(false) expect(client.close).to be(false) - expect(synchronizer).to have_received(:close).once - expect(processor).to have_received(:close).once + expect(synchronizer).to have_received(:close).twice + expect(processor).to have_received(:close).twice expect(client.status_provider.status).to eq(FeatBit::Status::CLOSED) end + + it "can complete a previously incomplete close and caches only success" do + synchronizer = instance_double("Synchronizer", start: true) + allow(synchronizer).to receive(:close).and_return(false, true) + client = described_class.new(FeatBit::Options.new( + env_secret: "secret", disable_events: true, start_wait: 0.001, + synchronizer_factory: ->(*) { synchronizer } + )) + + expect(client.close).to be(false) + expect(client.close).to be(true) + expect(client.close).to be(true) + expect(synchronizer).to have_received(:close).twice + end + + it "does not report concurrent or status-listener shutdown as completed" do + entered = Queue.new + release = Queue.new + synchronizer = instance_double("Synchronizer", start: true) + allow(synchronizer).to receive(:close) do + entered << true + release.pop + true + end + options = FeatBit::Options.new( + env_secret: "secret", disable_events: true, start_wait: 0.001, synchronizer_factory: ->(*) { synchronizer } + ) + client = described_class.new(options) + nested = [] + client.status_provider.add_listener { |state, _| nested << client.close if state == FeatBit::Status::CLOSED } + closer = Thread.new { client.close } + Timeout.timeout(2) { entered.pop } + + expect(client.close).to be(false) + release << true + expect(closer.join(2).value).to be(true) + expect(nested).to eq([false]) + expect(client.close).to be(true) + expect(synchronizer).to have_received(:close).once + ensure + release << true if release + closer&.join(2) + end end diff --git a/spec/web_socket_data_synchronizer_spec.rb b/spec/web_socket_data_synchronizer_spec.rb index 402104b..5deac97 100644 --- a/spec/web_socket_data_synchronizer_spec.rb +++ b/spec/web_socket_data_synchronizer_spec.rb @@ -345,14 +345,14 @@ def close = @closed = true allow(synchronizer).to receive(:interruptible_sleep) do |duration| delay_count += 1 delays << duration - synchronizer.instance_variable_set(:@closed, true) if delay_count == 2 + synchronizer.close if delay_count == 2 end synchronizer.start expect(Timeout.timeout(2) { delays.pop }).to eq(1.0) expect(Timeout.timeout(2) { delays.pop }).to eq(2.0) - expect(Timeout.timeout(2) { sleep(0.01) while synchronizer.instance_variable_get(:@thread)&.alive? }).to be_nil + expect(synchronizer.close).to be(true) expect(sockets.length).to eq(2) end @@ -385,6 +385,7 @@ def close = @closed = true fake_socket.handlers.fetch(:error).call(StandardError.new("connection failed")) + Timeout.timeout(2) { sleep(0.01) until fake_socket.closed? } expect(fake_socket).to be_closed expect(synchronizer.close).to be(true) end @@ -425,7 +426,8 @@ def close = @closed = true Timeout.timeout(2) { sleep(0.01) until socket.handlers.key?(:message) } socket.handlers.fetch(:open).call socket.handlers.fetch(:message).call(close_frame) - Timeout.timeout(2) { sleep(0.01) while synchronizer.instance_variable_get(:@thread)&.alive? } + worker = synchronizer.instance_variable_get(:@lifecycle).instance_variable_get(:@thread) + expect(worker.join(2)).to eq(worker) expect(attempts).to eq(1) expect(socket).to be_closed diff --git a/spec/web_socket_shutdown_spec.rb b/spec/web_socket_shutdown_spec.rb new file mode 100644 index 0000000..1af52a0 --- /dev/null +++ b/spec/web_socket_shutdown_spec.rb @@ -0,0 +1,302 @@ +# frozen_string_literal: true + +require "spec_helper" +require "socket" +require "timeout" + +class ShutdownTestSocket + attr_reader :handlers, :sent, :close_calls + + def initialize + @handlers = {} + @sent = [] + @close_calls = 0 + end + + def on(name, &block) = @handlers[name] = block + def send(message) = @sent << message + def closed? = @close_calls.positive? + + def close + @close_calls += 1 + @handlers[:close]&.call(nil) + true + end + + def emit(name, event = nil) = @handlers.fetch(name).call(event) +end + +RSpec.describe "WebSocket shutdown" do + let(:store) { FeatBit::InMemoryDataStore.new } + let(:status) { FeatBit::StatusProvider.new } + let(:socket) { ShutdownTestSocket.new } + let(:connected) { Queue.new } + let(:options) { FeatBit::Options.new(env_secret: "secret", connect_timeout: 30, reconnect_delay: 0.01) } + + def build_sync(connector:, **args) + @synchronizer = FeatBit::WebSocketDataSynchronizer.new( + options: options, data_store: store, status_provider: status, connector: connector, **args + ) + end + + def wait_for(queue) = Timeout.timeout(2) { queue.pop } + + def configured_connector + lambda do |*, &configure| + configure.call(socket) + connected << Thread.current + socket + end + end + + after { @synchronizer&.close } + + it "can close before start and cannot subsequently restart" do + sync = build_sync(connector: configured_connector) + expect(sync.close).to be(true) + expect(sync.close).to be(true) + expect(sync.start).to be(false) + expect(connected).to be_empty + end + + it "does not claim successful cleanup or reconnect after a socket close failure" do + allow(socket).to receive(:close).and_return(false) + sync = build_sync(connector: configured_connector) + sync.start + wait_for(connected) + + expect(sync.close).to be(false) + expect(sync.close).to be(false) + expect(sync.start).to be(false) + expect(connected).to be_empty + expect(socket).to have_received(:close).once + end + + it "cancels a connector that has not returned or published a socket" do + sync = build_sync(connector: lambda { |*| + connected << Thread.current + Queue.new.pop + }) + sync.start + connector_thread = wait_for(connected) + + expect(Timeout.timeout(2) { sync.close }).to be(true) + expect(connector_thread).not_to be_alive + expect(sync.close).to be(true) + expect(sync.start).to be(false) + end + + it "closes a partially connected socket only after its connector has exited" do + connector_thread = nil + allow(socket).to receive(:close).and_wrap_original do |original| + expect(connector_thread).not_to be_alive + original.call + end + sync = build_sync(connector: lambda { |*, &configure| + configure.call(socket) + connected << Thread.current + Queue.new.pop + }) + sync.start + connector_thread = wait_for(connected) + + expect(Timeout.timeout(2) { sync.close }).to be(true) + expect(socket.close_calls).to eq(1) + end + + it "interrupts handshake waiting and ignores every late callback" do + sync = build_sync(connector: configured_connector) + sync.start + wait_for(connected) + expect(Timeout.timeout(2) { sync.close }).to be(true) + status.update(FeatBit::Status::CLOSED) + + socket.emit(:open) + socket.emit(:message, JSON.generate(test_bootstrap(test_flag))) + socket.emit(:error, StandardError.new("late error")) + socket.emit(:close, Struct.new(:code, :reason).new(4003, "late rejection")) + + expect(socket.sent).to be_empty + expect(store).not_to be_initialized + expect(status.status).to eq(FeatBit::Status::CLOSED) + expect(socket.close_calls).to eq(1) + end + + it "waits for an in-flight callback without killing it or reporting early success" do + entered = Queue.new + release = Queue.new + stub_const("FeatBit::WebSocketLifecycle::CLOSE_WAIT", 0.05) + sync = build_sync(connector: configured_connector, on_flags_changed: lambda { |_| + entered << true + release.pop + }) + sync.start + wait_for(connected) + socket.emit(:open) + callback = Thread.new { socket.emit(:message, JSON.generate(test_bootstrap(test_flag))) } + wait_for(entered) + + expect(sync.close).to be(false) + expect(callback).to be_alive + expect(socket.close_calls).to eq(0) + release << true + expect(callback.join(2)).to eq(callback) + expect(Timeout.timeout(2) { sleep(0.01) until sync.close }).to be_nil + expect(status.status).not_to eq(FeatBit::Status::READY) + expect(socket.close_calls).to eq(1) + ensure + release << true if release + callback&.join(2) + end + + it "allows a callback to request close without joining itself or continuing notifications" do + calls = [] + close_results = [] + sync = build_sync(connector: configured_connector, on_flags_changed: lambda { |key| + calls << key + close_results << @synchronizer.close + status.update(FeatBit::Status::CLOSED) + }) + sync.start + wait_for(connected) + socket.emit(:open) + socket.emit(:message, JSON.generate(test_bootstrap(test_flag, test_flag(key: "second")))) + + expect(close_results).to eq([false]) + expect(calls).to eq(["welcome"]) + expect(Timeout.timeout(2) { sync.close }).to be(true) + expect(status.status).to eq(FeatBit::Status::CLOSED) + end + + it "ignores callbacks from an old connection after reconnecting" do + sockets = Queue.new + sync = build_sync(connector: lambda { |*, &configure| + candidate = ShutdownTestSocket.new + configure.call(candidate) + sockets << candidate + candidate + }) + sync.start + first = wait_for(sockets) + first.emit(:open) + first.emit(:error, StandardError.new("disconnect")) + second = wait_for(sockets) + second.emit(:open) + second.emit(:message, JSON.generate(test_bootstrap(test_flag))) + + first.emit(:error, StandardError.new("old error")) + first.emit(:close, Struct.new(:code).new(4003)) + expect(status.status).to eq(FeatBit::Status::READY) + expect(second).not_to be_closed + end + + %w[ws wss].each do |scheme| + it "cancels a real stalled #{scheme == 'wss' ? 'TLS' : 'HTTP Upgrade'} connection" do + server = TCPServer.new("127.0.0.1", 0) + accepted = Queue.new + eof = Queue.new + peer = nil + server_thread = Thread.new do + peer = server.accept + peer.readpartial(4096) # Observe the actual ClientHello or Upgrade request. + accepted << true + peer.read # Never send TLS/Upgrade response; wait for client EOF. + eof << true + end + real_options = FeatBit::Options.new( + env_secret: "secret", streaming_url: "#{scheme}://127.0.0.1:#{server.addr[1]}", connect_timeout: 30 + ) + @synchronizer = FeatBit::WebSocketDataSynchronizer.new(options: real_options, data_store: store, status_provider: status) + @synchronizer.start + wait_for(accepted) + attempt = @synchronizer.instance_variable_get(:@lifecycle).instance_variable_get(:@attempt) + actual_socket = attempt.socket + connector_thread = attempt.instance_variable_get(:@connector_thread) + + expect(Timeout.timeout(2) { @synchronizer.close }).to be(true) + expect(wait_for(eof)).to be(true) + expect(connector_thread).not_to be_alive + expect(actual_socket.thread&.alive?).not_to be(true) + expect(server_thread.join(2)).to eq(server_thread) + ensure + @synchronizer&.close + peer&.close + server&.close + server_thread&.kill + server_thread&.join + end + end + + it "returns from close inside a real receive callback before draining the reader" do + server = TCPServer.new("127.0.0.1", 0) + completed = Queue.new + peer = nil + server_thread = Thread.new do + peer = server.accept + handshake = WebSocket::Handshake::Server.new + handshake << peer.readpartial(4096) until handshake.finished? + peer.write(handshake.to_s) + peer.write(WebSocket::Frame::Outgoing::Server.new( + version: 13, type: :text, data: JSON.generate(test_bootstrap(test_flag)) + ).to_s) + peer.read + end + real_options = FeatBit::Options.new(env_secret: "secret", streaming_url: "ws://127.0.0.1:#{server.addr[1]}") + @synchronizer = FeatBit::WebSocketDataSynchronizer.new( + options: real_options, data_store: store, status_provider: status, + on_flags_changed: ->(_) { completed << [@synchronizer.close, Thread.current] } + ) + @synchronizer.start + + result, reader = wait_for(completed) + expect(result).to be(false) + expect(Timeout.timeout(2) { @synchronizer.close }).to be(true) + expect(reader).not_to be_alive + expect(server_thread.join(2)).to eq(server_thread) + ensure + @synchronizer&.close + peer&.close + server&.close + server_thread&.kill + server_thread&.join + end + + it "forces transport cleanup when sending the close frame stalls" do + stub_const("FeatBit::ClosableWebSocketClient::CLOSE_TIMEOUT", 0.05) + client = FeatBit::ClosableWebSocketClient.new + transport = instance_double(TCPSocket, close: nil) + reader = Thread.new { sleep } + client.instance_variable_set(:@socket, transport) + client.instance_variable_set(:@thread, reader) + allow(client).to receive(:send) { Queue.new.pop } + + expect(Timeout.timeout(2) { client.close(drain: true) }).to be(true) + expect(transport).to have_received(:close) + expect(reader).not_to be_alive + ensure + reader&.kill + reader&.join + end + + it "defers library-initiated reader shutdown to the owner" do + client = FeatBit::ClosableWebSocketClient.new + transport = instance_double(TCPSocket, close: nil) + client.instance_variable_set(:@socket, transport) + returned = Queue.new + reader = Thread.new do + client.instance_variable_set(:@thread, Thread.current) + returned << client.close + sleep + end + + expect(wait_for(returned)).to be(true) + expect(reader).to be_alive + expect(transport).not_to have_received(:close) + expect(client.close(drain: true)).to be(true) + expect(transport).to have_received(:close) + expect(reader).not_to be_alive + ensure + reader&.kill + reader&.join + end +end From 0db3d64b17f564ddf719913cfb036a42a54319d7 Mon Sep 17 00:00:00 2001 From: deleteLater Date: Wed, 2 Sep 2026 11:34:42 +0800 Subject: [PATCH 5/9] Refactor WebSocket data synchronization into data_sync --- lib/featbit.rb | 2 +- .../data_sync/closable_web_socket_client.rb | 37 ++++ .../data_sync/synchronization_result.rb | 24 +++ .../data_sync/web_socket_close_policy.rb | 54 +++++ .../web_socket_connection_attempt.rb | 95 +++++++++ .../web_socket_data_synchronizer.rb | 165 +-------------- .../{ => data_sync}/web_socket_lifecycle.rb | 34 ---- .../web_socket_connection_spec.rb} | 190 +++--------------- .../web_socket_data_synchronizer_spec.rb | 144 +++++++++++++ .../web_socket_shutdown_spec.rb | 0 10 files changed, 387 insertions(+), 358 deletions(-) create mode 100644 lib/featbit/data_sync/closable_web_socket_client.rb create mode 100644 lib/featbit/data_sync/synchronization_result.rb create mode 100644 lib/featbit/data_sync/web_socket_close_policy.rb create mode 100644 lib/featbit/data_sync/web_socket_connection_attempt.rb rename lib/featbit/{ => data_sync}/web_socket_data_synchronizer.rb (70%) rename lib/featbit/{ => data_sync}/web_socket_lifecycle.rb (65%) rename spec/{web_socket_data_synchronizer_spec.rb => data_sync/web_socket_connection_spec.rb} (69%) create mode 100644 spec/data_sync/web_socket_data_synchronizer_spec.rb rename spec/{ => data_sync}/web_socket_shutdown_spec.rb (100%) diff --git a/lib/featbit.rb b/lib/featbit.rb index c389d8e..713562a 100644 --- a/lib/featbit.rb +++ b/lib/featbit.rb @@ -8,7 +8,7 @@ require_relative "featbit/data_store" require_relative "featbit/evaluator" require_relative "featbit/event_processor" -require_relative "featbit/web_socket_data_synchronizer" +require_relative "featbit/data_sync/web_socket_data_synchronizer" require_relative "featbit/client" module FeatBit diff --git a/lib/featbit/data_sync/closable_web_socket_client.rb b/lib/featbit/data_sync/closable_web_socket_client.rb new file mode 100644 index 0000000..ec8314f --- /dev/null +++ b/lib/featbit/data_sync/closable_web_socket_client.rb @@ -0,0 +1,37 @@ +# frozen_string_literal: true + +require "timeout" +require "websocket-client-simple" + +module FeatBit + # websocket-client-simple has no abort API. Keep its internal cleanup here, + # after the connector has exited and SDK callbacks have drained. + class ClosableWebSocketClient < WebSocket::Client::Simple::Client + CLOSE_TIMEOUT = 1.0 + + def close(drain: false) + # The library also calls close from its own I/O error callbacks. Only the + # SDK owner may drain/kill the reader, after application callbacks return. + unless drain + @closed = true + return true + end + + begin + Timeout.timeout(CLOSE_TIMEOUT) { super() } + rescue Timeout::Error + # A stalled close-frame write must not prevent closing the TCP socket. + ensure + @closed = true + begin + @socket&.close + @socket = nil + ensure + @thread&.kill + @thread&.join unless Thread.current.equal?(@thread) + end + end + true + end + end +end diff --git a/lib/featbit/data_sync/synchronization_result.rb b/lib/featbit/data_sync/synchronization_result.rb new file mode 100644 index 0000000..65b03ea --- /dev/null +++ b/lib/featbit/data_sync/synchronization_result.rb @@ -0,0 +1,24 @@ +# frozen_string_literal: true + +module FeatBit + class SynchronizationResult + attr_reader :valid, :changed + + def initialize(valid:, changed:) + @valid = valid + @changed = changed + freeze + end + + def valid? = valid + def changed? = changed + + INVALID = new(valid: false, changed: false) + UNCHANGED = new(valid: true, changed: false) + CHANGED = new(valid: true, changed: true) + + def self.valid(changed:) + changed ? CHANGED : UNCHANGED + end + end +end diff --git a/lib/featbit/data_sync/web_socket_close_policy.rb b/lib/featbit/data_sync/web_socket_close_policy.rb new file mode 100644 index 0000000..bd9ab92 --- /dev/null +++ b/lib/featbit/data_sync/web_socket_close_policy.rb @@ -0,0 +1,54 @@ +# frozen_string_literal: true + +require_relative "../status" + +module FeatBit + class WebSocketClosePolicy + SERVER_REJECTED_CODE = 4003 + + def initialize(status_provider) + @status_provider = status_provider + @mutex = Mutex.new + @rejected = false + end + + def close_frame?(event) + event.respond_to?(:type) && event.type.to_sym == :close + rescue StandardError + false + end + + def reject?(event) + return false unless close_code(event) == SERVER_REJECTED_CODE + + @mutex.synchronize { @rejected = true } + reason = close_reason(event) + message = "WebSocket connection rejected by server (4003)" + message = "#{message}: #{reason}" unless reason.empty? + @status_provider.update(Status::FAILED, message: message) + true + end + + def rejected? + @mutex.synchronize { @rejected } + end + + private + + def close_code(event) + return event.code.to_i if event.respond_to?(:code) && event.code + return event[:code].to_i if event.is_a?(Hash) && event[:code] + + event["code"].to_i if event.is_a?(Hash) && event["code"] + end + + def close_reason(event) + value = if event.respond_to?(:reason) && event.reason + event.reason + elsif event.respond_to?(:data) && event.data + event.data + end + value.to_s + end + end +end diff --git a/lib/featbit/data_sync/web_socket_connection_attempt.rb b/lib/featbit/data_sync/web_socket_connection_attempt.rb new file mode 100644 index 0000000..9b605f4 --- /dev/null +++ b/lib/featbit/data_sync/web_socket_connection_attempt.rb @@ -0,0 +1,95 @@ +# frozen_string_literal: true + +require "timeout" + +module FeatBit + class WebSocketConnectionAttempt + attr_reader :socket + + def initialize(timeout:, clock:, lifecycle:, closer:) + @deadline = clock.call + timeout + @clock = clock + @mutex = Mutex.new + @condition = ConditionVariable.new + @result = nil + @lifecycle = lifecycle + @closer = closer + @finished = false + end + + def connect(connector, url, headers) + @connector_thread = Thread.new do + configured = false + @socket = connector.call(url, headers) do |connected_socket| + @socket = connected_socket + yield connected_socket + configured = true + end + yield @socket unless configured + rescue StandardError => e + @connect_error = e + end + until @connector_thread.join(0.01) + break if @lifecycle.stopped? || @clock.call >= @deadline + end + return @socket if @lifecycle.stopped? + + raise Timeout::Error, "WebSocket connection timed out" if @connector_thread.alive? + raise @connect_error if @connect_error + + @socket + end + + def cleanup + # Cancellation is confined to the owned connector, never an application + # callback. Joining before close prevents a late connector creating a socket. + @connector_thread&.kill if @connector_thread&.alive? + @connector_thread&.join + @closer.call(@socket) + end + + def signal(result) + @mutex.synchronize do + @finished = true unless result == :opened + return if @result + + @result = result + @condition.broadcast + end + end + + def wait(stopped:) + @mutex.synchronize do + until @result || stopped.call + remaining = @deadline - @clock.call + return :timeout unless remaining.positive? + + @condition.wait(@mutex, [remaining, 0.05].min) + end + stopped.call ? :stopped : @result + end + end + + def monitor(stopped:, ping:, interval:) + next_ping = @clock.call + interval + until stopped.call || @mutex.synchronize { @finished } || socket_closed? + if @clock.call >= next_ping + ping.call(@socket) + next_ping = @clock.call + interval + end + sleep(0.05) + end + end + + private + + def socket_closed? + return @socket.closed? if @socket.respond_to?(:closed?) + return !@socket.open? if @socket.respond_to?(:open?) + + false + rescue StandardError + true + end + end +end diff --git a/lib/featbit/web_socket_data_synchronizer.rb b/lib/featbit/data_sync/web_socket_data_synchronizer.rb similarity index 70% rename from lib/featbit/web_socket_data_synchronizer.rb rename to lib/featbit/data_sync/web_socket_data_synchronizer.rb index b3c7a33..17b1ff4 100644 --- a/lib/featbit/web_socket_data_synchronizer.rb +++ b/lib/featbit/data_sync/web_socket_data_synchronizer.rb @@ -4,170 +4,13 @@ require "securerandom" require "time" require "timeout" -require "websocket-client-simple" +require_relative "synchronization_result" +require_relative "web_socket_close_policy" +require_relative "web_socket_connection_attempt" require_relative "web_socket_lifecycle" +require_relative "closable_web_socket_client" module FeatBit - class SynchronizationResult - attr_reader :valid, :changed - - def initialize(valid:, changed:) - @valid = valid - @changed = changed - freeze - end - - def valid? = valid - def changed? = changed - - INVALID = new(valid: false, changed: false) - UNCHANGED = new(valid: true, changed: false) - CHANGED = new(valid: true, changed: true) - - def self.valid(changed:) - changed ? CHANGED : UNCHANGED - end - end - - class WebSocketClosePolicy - SERVER_REJECTED_CODE = 4003 - - def initialize(status_provider) - @status_provider = status_provider - @mutex = Mutex.new - @rejected = false - end - - def close_frame?(event) - event.respond_to?(:type) && event.type.to_sym == :close - rescue StandardError - false - end - - def reject?(event) - return false unless close_code(event) == SERVER_REJECTED_CODE - - @mutex.synchronize { @rejected = true } - reason = close_reason(event) - message = "WebSocket connection rejected by server (4003)" - message = "#{message}: #{reason}" unless reason.empty? - @status_provider.update(Status::FAILED, message: message) - true - end - - def rejected? - @mutex.synchronize { @rejected } - end - - private - - def close_code(event) - return event.code.to_i if event.respond_to?(:code) && event.code - return event[:code].to_i if event.is_a?(Hash) && event[:code] - - event["code"].to_i if event.is_a?(Hash) && event["code"] - end - - def close_reason(event) - value = if event.respond_to?(:reason) && event.reason - event.reason - elsif event.respond_to?(:data) && event.data - event.data - end - value.to_s - end - end - - class WebSocketConnectionAttempt - attr_reader :socket - - def initialize(timeout:, clock:, lifecycle:, closer:) - @deadline = clock.call + timeout - @clock = clock - @mutex = Mutex.new - @condition = ConditionVariable.new - @result = nil - @lifecycle = lifecycle - @closer = closer - @finished = false - end - - def connect(connector, url, headers) - @connector_thread = Thread.new do - configured = false - @socket = connector.call(url, headers) do |connected_socket| - @socket = connected_socket - yield connected_socket - configured = true - end - yield @socket unless configured - rescue StandardError => e - @connect_error = e - end - until @connector_thread.join(0.01) - break if @lifecycle.stopped? || @clock.call >= @deadline - end - return @socket if @lifecycle.stopped? - - raise Timeout::Error, "WebSocket connection timed out" if @connector_thread.alive? - raise @connect_error if @connect_error - - @socket - end - - def cleanup - # Cancellation is confined to the owned connector, never an application - # callback. Joining before close prevents a late connector creating a socket. - @connector_thread&.kill if @connector_thread&.alive? - @connector_thread&.join - @closer.call(@socket) - end - - def signal(result) - @mutex.synchronize do - @finished = true unless result == :opened - return if @result - - @result = result - @condition.broadcast - end - end - - def wait(stopped:) - @mutex.synchronize do - until @result || stopped.call - remaining = @deadline - @clock.call - return :timeout unless remaining.positive? - - @condition.wait(@mutex, [remaining, 0.05].min) - end - stopped.call ? :stopped : @result - end - end - - def monitor(stopped:, ping:, interval:) - next_ping = @clock.call + interval - until stopped.call || @mutex.synchronize { @finished } || socket_closed? - if @clock.call >= next_ping - ping.call(@socket) - next_ping = @clock.call + interval - end - sleep(0.05) - end - end - - private - - def socket_closed? - return @socket.closed? if @socket.respond_to?(:closed?) - return !@socket.open? if @socket.respond_to?(:open?) - - false - rescue StandardError - true - end - end - class WebSocketDataSynchronizer ALPHABETS = { "0" => "Q", "1" => "B", "2" => "W", "3" => "S", "4" => "P", "5" => "H", "6" => "D", "7" => "X", "8" => "Z", "9" => "U" }.freeze diff --git a/lib/featbit/web_socket_lifecycle.rb b/lib/featbit/data_sync/web_socket_lifecycle.rb similarity index 65% rename from lib/featbit/web_socket_lifecycle.rb rename to lib/featbit/data_sync/web_socket_lifecycle.rb index 4381ef4..17def09 100644 --- a/lib/featbit/web_socket_lifecycle.rb +++ b/lib/featbit/data_sync/web_socket_lifecycle.rb @@ -1,8 +1,5 @@ # frozen_string_literal: true -require "timeout" -require "websocket-client-simple" - module FeatBit # Owns shutdown and drains callbacks without holding locks around user code. class WebSocketLifecycle @@ -76,35 +73,4 @@ def finish(attempt) end end end - - # websocket-client-simple has no abort API. Keep its internal cleanup here, - # after the connector has exited and SDK callbacks have drained. - class ClosableWebSocketClient < WebSocket::Client::Simple::Client - CLOSE_TIMEOUT = 1.0 - - def close(drain: false) - # The library also calls close from its own I/O error callbacks. Only the - # SDK owner may drain/kill the reader, after application callbacks return. - unless drain - @closed = true - return true - end - - begin - Timeout.timeout(CLOSE_TIMEOUT) { super() } - rescue Timeout::Error - # A stalled close-frame write must not prevent closing the TCP socket. - ensure - @closed = true - begin - @socket&.close - @socket = nil - ensure - @thread&.kill - @thread&.join unless Thread.current.equal?(@thread) - end - end - true - end - end end diff --git a/spec/web_socket_data_synchronizer_spec.rb b/spec/data_sync/web_socket_connection_spec.rb similarity index 69% rename from spec/web_socket_data_synchronizer_spec.rb rename to spec/data_sync/web_socket_connection_spec.rb index 5deac97..39d7804 100644 --- a/spec/web_socket_data_synchronizer_spec.rb +++ b/spec/data_sync/web_socket_connection_spec.rb @@ -8,168 +8,6 @@ let(:store) { FeatBit::InMemoryDataStore.new } let(:status) { FeatBit::StatusProvider.new(logger: options.logger) } - it "processes full synchronization messages and reports ready" do - changes = [] - synchronizer = described_class.new(options: options, data_store: store, status_provider: status, on_flags_changed: lambda { |key| - changes << key - }) - result = synchronizer.process_message(test_bootstrap(test_flag)) - expect(result).to be_valid - expect(result).to be_changed - expect(store.flag("welcome")).not_to be_nil - expect(status.status).to eq(FeatBit::Status::READY) - expect(changes).to include("welcome") - end - - it "rejects malformed data without raising" do - synchronizer = described_class.new(options: options, data_store: store, status_provider: status) - result = synchronizer.process_message("not-json") - expect(result).not_to be_valid - expect(result).not_to be_changed - expect(status.status).to eq(FeatBit::Status::FAILED) - end - - it "applies patches in timestamp order and only reports affected flags" do - other = test_flag(key: "other") - expect(store.init(test_bootstrap(test_flag, other), version: 1)).to be(true) - changes = [] - synchronizer = described_class.new( - options: options, - data_store: store, - status_provider: status, - on_flags_changed: ->(key) { changes << key } - ) - patched = test_flag - patched["name"] = "updated" - patched["updatedAt"] = "2026-01-02T00:00:00Z" - message = { - "messageType" => "data-sync", - "data" => { "eventType" => "patch", "featureFlags" => [patched], "segments" => [] } - } - result = synchronizer.process_message(message) - expect(result).to be_valid - expect(result).to be_changed - expect(store.flag("welcome")["name"]).to eq("updated") - expect(changes).to eq(["welcome"]) - end - - it "ignores stale patch items without hiding fresh changes" do - expect(store.init(test_bootstrap(test_flag), version: 10)).to be(true) - changes = [] - synchronizer = described_class.new( - options: options, - data_store: store, - status_provider: status, - on_flags_changed: ->(key) { changes << key } - ) - stale_flag = test_flag - stale_flag["timestamp"] = 9 - fresh_flag = test_flag(key: "fresh") - fresh_flag["timestamp"] = 11 - message = { - "messageType" => "data-sync", - "data" => { "eventType" => "patch", "featureFlags" => [stale_flag, fresh_flag], "segments" => [] } - } - - result = synchronizer.process_message(message) - expect(result).to be_valid - expect(result).to be_changed - expect(store.flag("fresh")).not_to be_nil - expect(changes).to eq(["fresh"]) - end - - it "accepts an all-stale patch without closing the socket" do - expect(store.init(test_bootstrap(test_flag), version: 10)).to be(true) - stale_flag = test_flag - stale_flag["timestamp"] = 9 - message = { - "messageType" => "data-sync", - "data" => { "eventType" => "patch", "featureFlags" => [stale_flag], "segments" => [] } - } - socket = instance_double("Socket", close: true) - synchronizer = described_class.new(options: options, data_store: store, status_provider: status) - - result = synchronizer.process_message(message) - synchronizer.send(:handle_socket_message, socket, JSON.generate(message)) - - expect(result).to be_valid - expect(result).not_to be_changed - expect(socket).not_to have_received(:close) - expect(store.version).to eq(10) - expect(status.status).to eq(FeatBit::Status::READY) - end - - it "accepts an unchanged full synchronization and reports ready" do - message = test_bootstrap(test_flag) - expect(store.init(message)).to be(true) - status.update(FeatBit::Status::INTERRUPTED, message: "disconnected") - socket = instance_double("Socket", close: true) - synchronizer = described_class.new(options: options, data_store: store, status_provider: status) - - result = synchronizer.process_message(message) - synchronizer.send(:handle_socket_message, socket, JSON.generate(message)) - - expect(result).to be_valid - expect(result).not_to be_changed - expect(socket).not_to have_received(:close) - expect(status.status).to eq(FeatBit::Status::READY) - end - - it "rejects an entire malformed patch before applying valid siblings" do - fresh_flag = test_flag(key: "fresh") - fresh_flag["timestamp"] = 11 - malformed_flag = test_flag - malformed_flag.delete("key") - message = { - "messageType" => "data-sync", - "data" => { "eventType" => "patch", "featureFlags" => [malformed_flag, fresh_flag], "segments" => [] } - } - socket = instance_double("Socket", close: true) - synchronizer = described_class.new(options: options, data_store: store, status_provider: status) - - synchronizer.send(:handle_socket_message, socket, JSON.generate(message)) - - expect(socket).to have_received(:close) - expect(store.flag("fresh")).to be_nil - end - - it "closes a socket when setup fails before reconnecting" do - closed = Queue.new - fake_socket = Class.new do - def initialize(closed) - @closed = closed - end - - def on(*) = raise("handler setup failed") - - def close - @closed << true - true - end - end.new(closed) - slow_options = FeatBit::Options.new(env_secret: "secret", reconnect_delay: 10) - synchronizer = described_class.new( - options: slow_options, - data_store: store, - status_provider: status, - connector: ->(*) { fake_socket } - ) - - synchronizer.start - - expect(Timeout.timeout(2) { closed.pop }).to be(true) - expect(synchronizer.close).to be(true) - end - - it "closes the socket when a synchronization message is rejected" do - socket = instance_double("Socket", close: true) - synchronizer = described_class.new(options: options, data_store: store, status_provider: status) - - synchronizer.send(:handle_socket_message, socket, "not-json") - - expect(socket).to have_received(:close) - end - it "connects with FeatBit authentication and requests the local version" do fake_socket = Class.new do attr_reader :handlers, :sent @@ -257,6 +95,34 @@ def close = @closed = true expect(synchronizer.close).to be(true) end + it "closes a socket when setup fails before reconnecting" do + closed = Queue.new + fake_socket = Class.new do + def initialize(closed) + @closed = closed + end + + def on(*) = raise("handler setup failed") + + def close + @closed << true + true + end + end.new(closed) + slow_options = FeatBit::Options.new(env_secret: "secret", reconnect_delay: 10) + synchronizer = described_class.new( + options: slow_options, + data_store: store, + status_provider: status, + connector: ->(*) { fake_socket } + ) + + synchronizer.start + + expect(Timeout.timeout(2) { closed.pop }).to be(true) + expect(synchronizer.close).to be(true) + end + it "interrupts reconnect backoff promptly when closed" do attempted = Queue.new connector = lambda do |_url, _headers| diff --git a/spec/data_sync/web_socket_data_synchronizer_spec.rb b/spec/data_sync/web_socket_data_synchronizer_spec.rb new file mode 100644 index 0000000..f3b7514 --- /dev/null +++ b/spec/data_sync/web_socket_data_synchronizer_spec.rb @@ -0,0 +1,144 @@ +# frozen_string_literal: true + +require "spec_helper" +require "timeout" + +RSpec.describe FeatBit::WebSocketDataSynchronizer do + let(:options) { FeatBit::Options.new(env_secret: "secret") } + let(:store) { FeatBit::InMemoryDataStore.new } + let(:status) { FeatBit::StatusProvider.new(logger: options.logger) } + + it "processes full synchronization messages and reports ready" do + changes = [] + synchronizer = described_class.new(options: options, data_store: store, status_provider: status, on_flags_changed: lambda { |key| + changes << key + }) + result = synchronizer.process_message(test_bootstrap(test_flag)) + expect(result).to be_valid + expect(result).to be_changed + expect(store.flag("welcome")).not_to be_nil + expect(status.status).to eq(FeatBit::Status::READY) + expect(changes).to include("welcome") + end + + it "rejects malformed data without raising" do + synchronizer = described_class.new(options: options, data_store: store, status_provider: status) + result = synchronizer.process_message("not-json") + expect(result).not_to be_valid + expect(result).not_to be_changed + expect(status.status).to eq(FeatBit::Status::FAILED) + end + + it "applies patches in timestamp order and only reports affected flags" do + other = test_flag(key: "other") + expect(store.init(test_bootstrap(test_flag, other), version: 1)).to be(true) + changes = [] + synchronizer = described_class.new( + options: options, + data_store: store, + status_provider: status, + on_flags_changed: ->(key) { changes << key } + ) + patched = test_flag + patched["name"] = "updated" + patched["updatedAt"] = "2026-01-02T00:00:00Z" + message = { + "messageType" => "data-sync", + "data" => { "eventType" => "patch", "featureFlags" => [patched], "segments" => [] } + } + result = synchronizer.process_message(message) + expect(result).to be_valid + expect(result).to be_changed + expect(store.flag("welcome")["name"]).to eq("updated") + expect(changes).to eq(["welcome"]) + end + + it "ignores stale patch items without hiding fresh changes" do + expect(store.init(test_bootstrap(test_flag), version: 10)).to be(true) + changes = [] + synchronizer = described_class.new( + options: options, + data_store: store, + status_provider: status, + on_flags_changed: ->(key) { changes << key } + ) + stale_flag = test_flag + stale_flag["timestamp"] = 9 + fresh_flag = test_flag(key: "fresh") + fresh_flag["timestamp"] = 11 + message = { + "messageType" => "data-sync", + "data" => { "eventType" => "patch", "featureFlags" => [stale_flag, fresh_flag], "segments" => [] } + } + + result = synchronizer.process_message(message) + expect(result).to be_valid + expect(result).to be_changed + expect(store.flag("fresh")).not_to be_nil + expect(changes).to eq(["fresh"]) + end + + it "accepts an all-stale patch without closing the socket" do + expect(store.init(test_bootstrap(test_flag), version: 10)).to be(true) + stale_flag = test_flag + stale_flag["timestamp"] = 9 + message = { + "messageType" => "data-sync", + "data" => { "eventType" => "patch", "featureFlags" => [stale_flag], "segments" => [] } + } + socket = instance_double("Socket", close: true) + synchronizer = described_class.new(options: options, data_store: store, status_provider: status) + + result = synchronizer.process_message(message) + synchronizer.send(:handle_socket_message, socket, JSON.generate(message)) + + expect(result).to be_valid + expect(result).not_to be_changed + expect(socket).not_to have_received(:close) + expect(store.version).to eq(10) + expect(status.status).to eq(FeatBit::Status::READY) + end + + it "accepts an unchanged full synchronization and reports ready" do + message = test_bootstrap(test_flag) + expect(store.init(message)).to be(true) + status.update(FeatBit::Status::INTERRUPTED, message: "disconnected") + socket = instance_double("Socket", close: true) + synchronizer = described_class.new(options: options, data_store: store, status_provider: status) + + result = synchronizer.process_message(message) + synchronizer.send(:handle_socket_message, socket, JSON.generate(message)) + + expect(result).to be_valid + expect(result).not_to be_changed + expect(socket).not_to have_received(:close) + expect(status.status).to eq(FeatBit::Status::READY) + end + + it "rejects an entire malformed patch before applying valid siblings" do + fresh_flag = test_flag(key: "fresh") + fresh_flag["timestamp"] = 11 + malformed_flag = test_flag + malformed_flag.delete("key") + message = { + "messageType" => "data-sync", + "data" => { "eventType" => "patch", "featureFlags" => [malformed_flag, fresh_flag], "segments" => [] } + } + socket = instance_double("Socket", close: true) + synchronizer = described_class.new(options: options, data_store: store, status_provider: status) + + synchronizer.send(:handle_socket_message, socket, JSON.generate(message)) + + expect(socket).to have_received(:close) + expect(store.flag("fresh")).to be_nil + end + + it "closes the socket when a synchronization message is rejected" do + socket = instance_double("Socket", close: true) + synchronizer = described_class.new(options: options, data_store: store, status_provider: status) + + synchronizer.send(:handle_socket_message, socket, "not-json") + + expect(socket).to have_received(:close) + end +end diff --git a/spec/web_socket_shutdown_spec.rb b/spec/data_sync/web_socket_shutdown_spec.rb similarity index 100% rename from spec/web_socket_shutdown_spec.rb rename to spec/data_sync/web_socket_shutdown_spec.rb From 9da4fc00816f76dedfd6bb3ba8ad54fecac81d39 Mon Sep 17 00:00:00 2001 From: deleteLater Date: Wed, 9 Sep 2026 15:43:08 +0800 Subject: [PATCH 6/9] Return UNCHANGED for non data-sync messages --- .../data_sync/web_socket_data_synchronizer.rb | 2 +- .../web_socket_data_synchronizer_spec.rb | 30 +++++++++++++++++++ 2 files changed, 31 insertions(+), 1 deletion(-) diff --git a/lib/featbit/data_sync/web_socket_data_synchronizer.rb b/lib/featbit/data_sync/web_socket_data_synchronizer.rb index 17b1ff4..b669421 100644 --- a/lib/featbit/data_sync/web_socket_data_synchronizer.rb +++ b/lib/featbit/data_sync/web_socket_data_synchronizer.rb @@ -44,7 +44,7 @@ def process_message(message) return SynchronizationResult::INVALID if @lifecycle.stopped? envelope = message.is_a?(String) ? JSON.parse(message) : message - return SynchronizationResult::INVALID unless fetch(envelope, "messageType") == "data-sync" + return SynchronizationResult::UNCHANGED if fetch(envelope, "messageType") != "data-sync" data = fetch(envelope, "data", {}) event_type = fetch(data, "eventType") diff --git a/spec/data_sync/web_socket_data_synchronizer_spec.rb b/spec/data_sync/web_socket_data_synchronizer_spec.rb index f3b7514..a597829 100644 --- a/spec/data_sync/web_socket_data_synchronizer_spec.rb +++ b/spec/data_sync/web_socket_data_synchronizer_spec.rb @@ -8,6 +8,36 @@ let(:store) { FeatBit::InMemoryDataStore.new } let(:status) { FeatBit::StatusProvider.new(logger: options.logger) } + [FeatBit::Status::STARTING, FeatBit::Status::READY, FeatBit::Status::INTERRUPTED].each do |initial_status| + it "ignores unknown message types while #{initial_status}" do + store.init(test_bootstrap(test_flag)) unless initial_status == FeatBit::Status::STARTING + status.update(initial_status, message: "existing status") + original_state = [store.initialized?, store.version, store.all_flags] + socket = instance_double("Socket", close: true) + attempt = instance_double(FeatBit::WebSocketConnectionAttempt, signal: nil) + changes = [] + synchronizer = described_class.new( + options: options, data_store: store, status_provider: status, + on_flags_changed: ->(key) { changes << key } + ) + message = JSON.generate(messageType: "future-message", data: { private: "payload" }) + + synchronizer.send(:handle_socket_message, socket, message, attempt) + + expect(socket).not_to have_received(:close) + expect(attempt).not_to have_received(:signal) + expect([store.initialized?, store.version, store.all_flags]).to eq(original_state) + expect([status.status, status.message]).to eq([initial_status, "existing status"]) + expect(changes).to be_empty + + synchronizer.send(:handle_socket_message, socket, JSON.generate(test_bootstrap(test_flag)), attempt) + + expect(status.status).to eq(FeatBit::Status::READY) + expect(store.flag("welcome")).not_to be_nil + expect(attempt).not_to have_received(:signal) + end + end + it "processes full synchronization messages and reports ready" do changes = [] synchronizer = described_class.new(options: options, data_store: store, status_provider: status, on_flags_changed: lambda { |key| From bd7388cdc7e27df0e56fc77fceda681d7862c640 Mon Sep 17 00:00:00 2001 From: deleteLater Date: Wed, 9 Sep 2026 17:34:51 +0800 Subject: [PATCH 7/9] Remove unnecessary note --- README.md | 1 - 1 file changed, 1 deletion(-) diff --git a/README.md b/README.md index 2294b84..f25efa7 100644 --- a/README.md +++ b/README.md @@ -230,7 +230,6 @@ Call `track` after evaluating the related experiment flag. `numeric_value` defau - Event enqueue is non-blocking and bounded; overload increments `dropped_events`. - Status and flag callbacks execute outside internal locks. - `close` is thread-safe and idempotent, flushes accepted events, stops background workers, and prevents new events. -- A successful `close` confirms WebSocket connection work and callbacks have finished. It waits up to five seconds for the synchronizer; blocked application callbacks or cleanup failures return `false`, which may be retried. A concurrent or callback-initiated `close` can also return `false` while shutdown is incomplete; call again from outside callbacks to wait for completion. The `closed` status means shutdown was requested, not necessarily that cleanup succeeded. - Public client methods contain ordinary internal failures and return fallbacks or `false`. ## Supported Ruby versions From c90ff57adfd2a8fe103c09147713b3b23bddcb17 Mon Sep 17 00:00:00 2001 From: deleteLater Date: Wed, 9 Sep 2026 18:38:15 +0800 Subject: [PATCH 8/9] Ignore invalid WebSocket messages without disrupting synchronization --- .../data_sync/web_socket_data_synchronizer.rb | 28 +++++----- .../web_socket_data_synchronizer_spec.rb | 53 +++++++++++++++++-- 2 files changed, 64 insertions(+), 17 deletions(-) diff --git a/lib/featbit/data_sync/web_socket_data_synchronizer.rb b/lib/featbit/data_sync/web_socket_data_synchronizer.rb index b669421..757a51b 100644 --- a/lib/featbit/data_sync/web_socket_data_synchronizer.rb +++ b/lib/featbit/data_sync/web_socket_data_synchronizer.rb @@ -44,11 +44,14 @@ def process_message(message) return SynchronizationResult::INVALID if @lifecycle.stopped? envelope = message.is_a?(String) ? JSON.parse(message) : message - return SynchronizationResult::UNCHANGED if fetch(envelope, "messageType") != "data-sync" + message_type = fetch(envelope, "messageType") + return ignore_message("missing message type") unless message_type.is_a?(String) && !message_type.empty? + + return SynchronizationResult::UNCHANGED if message_type != "data-sync" data = fetch(envelope, "data", {}) event_type = fetch(data, "eventType") - return SynchronizationResult::INVALID unless valid_data?(data, event_type) + return ignore_message("invalid data") unless valid_data?(data, event_type) old_keys = @data_store.all_flags.keys if event_type == "full" @@ -61,16 +64,19 @@ def process_message(message) changed_keys.each { |key| safely_notify(key) unless @lifecycle.stopped? } @status_provider.update(Status::READY) unless @lifecycle.stopped? SynchronizationResult.valid(changed: changed) - rescue JSON::ParserError => e - @status_provider.update(Status::FAILED, message: "invalid data: #{e.message}") - SynchronizationResult::INVALID - rescue StandardError => e - fail_status(e) - SynchronizationResult::INVALID + rescue JSON::ParserError + ignore_message("invalid JSON") + rescue StandardError + ignore_message("processing error") end private + def ignore_message(reason) + @options.logger&.error("FeatBit ignored invalid sync message: #{reason}") + SynchronizationResult::INVALID + end + def run delay = @options.reconnect_delay until @lifecycle.stopped? || @close_policy.rejected? @@ -141,11 +147,7 @@ def handle_socket_message(socket, event, attempt = nil) return end - result = process_message(event.respond_to?(:data) ? event.data : event.to_s) - return if result.valid? - - attempt&.signal(:failed) - safe_close_socket(socket) unless attempt || @lifecycle.stopped? + process_message(event.respond_to?(:data) ? event.data : event.to_s) end def handle_socket_error(socket, event, attempt = nil) diff --git a/spec/data_sync/web_socket_data_synchronizer_spec.rb b/spec/data_sync/web_socket_data_synchronizer_spec.rb index a597829..6fdb8c6 100644 --- a/spec/data_sync/web_socket_data_synchronizer_spec.rb +++ b/spec/data_sync/web_socket_data_synchronizer_spec.rb @@ -56,7 +56,7 @@ result = synchronizer.process_message("not-json") expect(result).not_to be_valid expect(result).not_to be_changed - expect(status.status).to eq(FeatBit::Status::FAILED) + expect(status.status).to eq(FeatBit::Status::STARTING) end it "applies patches in timestamp order and only reports affected flags" do @@ -159,16 +159,61 @@ synchronizer.send(:handle_socket_message, socket, JSON.generate(message)) - expect(socket).to have_received(:close) + expect(socket).not_to have_received(:close) expect(store.flag("fresh")).to be_nil end - it "closes the socket when a synchronization message is rejected" do + it "keeps the socket open when a synchronization message is rejected" do socket = instance_double("Socket", close: true) synchronizer = described_class.new(options: options, data_store: store, status_provider: status) synchronizer.send(:handle_socket_message, socket, "not-json") - expect(socket).to have_received(:close) + expect(socket).not_to have_received(:close) + end + + [false, true].each do |with_attempt| + it "preserves state and processes later messages after invalid input, with attempt: #{with_attempt}" do + logger = instance_double(Logger, error: nil) + configured_options = FeatBit::Options.new(env_secret: "secret", logger: logger) + socket = instance_double("Socket", close: true) + attempt = with_attempt ? instance_double(FeatBit::WebSocketConnectionAttempt, signal: nil) : nil + changes = [] + synchronizer = described_class.new( + options: configured_options, data_store: store, status_provider: status, + on_flags_changed: ->(key) { changes << key } + ) + invalid_messages = [ + "not-json-private-payload", "null", "[]", "{}", + JSON.generate(messageType: "data-sync", data: nil), + JSON.generate(messageType: "data-sync", data: { eventType: "unknown", featureFlags: [], segments: [] }), + JSON.generate(messageType: "data-sync", data: { eventType: "full", featureFlags: {}, segments: [] }) + ] + + [FeatBit::Status::STARTING, FeatBit::Status::READY, FeatBit::Status::INTERRUPTED].each do |state| + store.init(test_bootstrap(test_flag)) unless state == FeatBit::Status::STARTING + status.update(state, message: "existing status") + original = [store.initialized?, store.version, store.all_flags] + invalid_messages.each { |message| synchronizer.send(:handle_socket_message, socket, message, attempt) } + + expect([store.initialized?, store.version, store.all_flags]).to eq(original) + expect([status.status, status.message]).to eq([state, "existing status"]) + expect(changes).to be_empty + end + expect(logger).to have_received(:error).exactly(invalid_messages.length * 3).times + expect(logger).not_to have_received(:error).with(include("private-payload")) + + patched = test_flag + patched["updatedAt"] = "2026-01-02T00:00:00Z" + patched["name"] = "recovered" + message = { messageType: "data-sync", data: { eventType: "patch", featureFlags: [patched], segments: [] } } + synchronizer.send(:handle_socket_message, socket, JSON.generate(message), attempt) + + expect(store.flag("welcome")["name"]).to eq("recovered") + expect(status.status).to eq(FeatBit::Status::READY) + expect(changes).to eq(["welcome"]) + expect(socket).not_to have_received(:close) + expect(attempt).not_to have_received(:signal) if attempt + end end end From 2a5833e6872bb6f17be968fd10c5d686f07da453 Mon Sep 17 00:00:00 2001 From: deleteLater Date: Wed, 9 Sep 2026 20:19:13 +0800 Subject: [PATCH 9/9] Harden shutdown cleanup and retain failed close results --- README.md | 3 +- lib/featbit/client.rb | 7 +-- .../data_sync/closable_web_socket_client.rb | 4 +- spec/client_spec.rb | 63 ------------------- spec/data_sync/web_socket_shutdown_spec.rb | 32 ++++++++++ 5 files changed, 36 insertions(+), 73 deletions(-) diff --git a/README.md b/README.md index f25efa7..5832d66 100644 --- a/README.md +++ b/README.md @@ -223,13 +223,12 @@ client.flush Call `track` after evaluating the related experiment flag. `numeric_value` defaults to `1.0`. -### Thread safety and shutdown +### Thread safety - Share one client across application threads; do not create a client per request. - Flag evaluation performs no network I/O. - Event enqueue is non-blocking and bounded; overload increments `dropped_events`. - Status and flag callbacks execute outside internal locks. -- `close` is thread-safe and idempotent, flushes accepted events, stops background workers, and prevents new events. - Public client methods contain ordinary internal failures and return fallbacks or `false`. ## Supported Ruby versions diff --git a/lib/featbit/client.rb b/lib/featbit/client.rb index 170fd11..34aaa4d 100644 --- a/lib/featbit/client.rb +++ b/lib/featbit/client.rb @@ -122,12 +122,9 @@ def remove_flag_change_listener(id) end def close - owns_close = false @lifecycle_mutex.synchronize do - return true if @close_result == true - return false if @closing + return @close_result == true if @closed - @closing = owns_close = true @closed = true end @@ -137,8 +134,6 @@ def close rescue StandardError => e safe_log(:warn, "FeatBit client close failed: #{e.message}") false - ensure - @lifecycle_mutex.synchronize { @closing = false } if owns_close end alias stop close diff --git a/lib/featbit/data_sync/closable_web_socket_client.rb b/lib/featbit/data_sync/closable_web_socket_client.rb index ec8314f..9b68f29 100644 --- a/lib/featbit/data_sync/closable_web_socket_client.rb +++ b/lib/featbit/data_sync/closable_web_socket_client.rb @@ -19,8 +19,8 @@ def close(drain: false) begin Timeout.timeout(CLOSE_TIMEOUT) { super() } - rescue Timeout::Error - # A stalled close-frame write must not prevent closing the TCP socket. + rescue IOError, SystemCallError, OpenSSL::SSL::SSLError, Timeout::Error + # A failed close handshake is harmless if the ensure cleanup succeeds. ensure @closed = true begin diff --git a/spec/client_spec.rb b/spec/client_spec.rb index 8a60190..b2c1217 100644 --- a/spec/client_spec.rb +++ b/spec/client_spec.rb @@ -185,67 +185,4 @@ def broken_store.segment(_key) = raise("boom") expect(client.close).to be(true) expect(reentrant_results).to eq([false]) end - - it "retries failed shutdown without skipping components or raising" do - synchronizer = instance_double("Synchronizer", start: true) - processor = instance_double("EventProcessor") - allow(synchronizer).to receive(:close).and_raise("synchronizer failed") - allow(processor).to receive(:close).and_return(true) - options = FeatBit::Options.new( - env_secret: "secret", - start_wait: 0.001, - synchronizer_factory: ->(*) { synchronizer }, - event_processor_factory: ->(*) { processor } - ) - client = described_class.new(options) - - expect(client.close).to be(false) - expect(client.close).to be(false) - expect(synchronizer).to have_received(:close).twice - expect(processor).to have_received(:close).twice - expect(client.status_provider.status).to eq(FeatBit::Status::CLOSED) - end - - it "can complete a previously incomplete close and caches only success" do - synchronizer = instance_double("Synchronizer", start: true) - allow(synchronizer).to receive(:close).and_return(false, true) - client = described_class.new(FeatBit::Options.new( - env_secret: "secret", disable_events: true, start_wait: 0.001, - synchronizer_factory: ->(*) { synchronizer } - )) - - expect(client.close).to be(false) - expect(client.close).to be(true) - expect(client.close).to be(true) - expect(synchronizer).to have_received(:close).twice - end - - it "does not report concurrent or status-listener shutdown as completed" do - entered = Queue.new - release = Queue.new - synchronizer = instance_double("Synchronizer", start: true) - allow(synchronizer).to receive(:close) do - entered << true - release.pop - true - end - options = FeatBit::Options.new( - env_secret: "secret", disable_events: true, start_wait: 0.001, synchronizer_factory: ->(*) { synchronizer } - ) - client = described_class.new(options) - nested = [] - client.status_provider.add_listener { |state, _| nested << client.close if state == FeatBit::Status::CLOSED } - closer = Thread.new { client.close } - Timeout.timeout(2) { entered.pop } - - expect(client.close).to be(false) - release << true - expect(closer.join(2).value).to be(true) - expect(nested).to eq([false]) - expect(client.close).to be(true) - expect(synchronizer).to have_received(:close).once - ensure - release << true if release - closer&.join(2) - end end diff --git a/spec/data_sync/web_socket_shutdown_spec.rb b/spec/data_sync/web_socket_shutdown_spec.rb index 1af52a0..de85682 100644 --- a/spec/data_sync/web_socket_shutdown_spec.rb +++ b/spec/data_sync/web_socket_shutdown_spec.rb @@ -278,6 +278,38 @@ def configured_connector reader&.join end + it "completes cleanup when writing the close frame fails" do + client = FeatBit::ClosableWebSocketClient.new + transport = instance_double(TCPSocket, close: nil) + reader = Thread.new { sleep } + client.instance_variable_set(:@socket, transport) + client.instance_variable_set(:@thread, reader) + allow(client).to receive(:send).and_raise(Errno::EPIPE) + + expect(client.close(drain: true)).to be(true) + expect(transport).to have_received(:close) + expect(reader).not_to be_alive + ensure + reader&.kill + reader&.join + end + + it "reports transport cleanup failure even when the close frame also fails" do + client = FeatBit::ClosableWebSocketClient.new + transport = instance_double(TCPSocket) + allow(transport).to receive(:close).and_raise(IOError, "cleanup failed") + reader = Thread.new { sleep } + client.instance_variable_set(:@socket, transport) + client.instance_variable_set(:@thread, reader) + allow(client).to receive(:send).and_raise(Errno::EPIPE) + + expect { client.close(drain: true) }.to raise_error(IOError, "cleanup failed") + expect(reader).not_to be_alive + ensure + reader&.kill + reader&.join + end + it "defers library-initiated reader shutdown to the owner" do client = FeatBit::ClosableWebSocketClient.new transport = instance_double(TCPSocket, close: nil)