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.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/client.rb b/lib/featbit/client.rb index d5a73fb..34aaa4d 100644 --- a/lib/featbit/client.rb +++ b/lib/featbit/client.rb @@ -123,8 +123,7 @@ def remove_flag_change_listener(id) def close @lifecycle_mutex.synchronize do - return @close_result unless @close_result.nil? - return true if @closed + return @close_result == true if @closed @closed = true end 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..9b68f29 --- /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 IOError, SystemCallError, OpenSSL::SSL::SSLError, Timeout::Error + # A failed close handshake is harmless if the ensure cleanup succeeds. + 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 55% rename from lib/featbit/web_socket_data_synchronizer.rb rename to lib/featbit/data_sync/web_socket_data_synchronizer.rb index 7bd024d..757a51b 100644 --- a/lib/featbit/web_socket_data_synchronizer.rb +++ b/lib/featbit/data_sync/web_socket_data_synchronizer.rb @@ -4,7 +4,11 @@ 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 WebSocketDataSynchronizer @@ -16,50 +20,38 @@ 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 - @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 false unless 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 false 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" @@ -69,115 +61,118 @@ def process_message(message) changed, changed_keys = process_patch(data) end - changed_keys.each { |key| safely_notify(key) } - return false unless changed - - @status_provider.update(Status::READY) - true - rescue JSON::ParserError => e - @status_provider.update(Status::FAILED, message: "invalid data: #{e.message}") - false - rescue StandardError => e - fail_status(e) - false + 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 + 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 @closed - socket = nil - begin - configured = false - socket = @connector.call(websocket_url, headers) do |connected_socket| - configure_socket(connected_socket) - configured = true - 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) - end - break if @closed - - @status_provider.update(Status::INTERRUPTED, message: "WebSocket disconnected") - interruptible_sleep(delay) - delay = [delay * 2, 30.0].min - rescue StandardError => e - 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 + until @lifecycle.stopped? || @close_policy.rejected? + opened = run_attempt + break if @lifecycle.stopped? || @close_policy.rejected? + + delay = @options.reconnect_delay if opened + interruptible_sleep(delay) + delay = [delay * 2, 30.0].min end end - def connect(url, request_headers, &configure) - Timeout.timeout(@options.connect_timeout) do - WebSocket::Client::Simple.connect(url, headers: request_headers, &configure) - end + 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 configure_socket(socket) - open_handler = method(:handle_socket_open) - message_handler = method(:handle_socket_message) - error_handler = method(:handle_socket_error) - close_handler = method(:handle_socket_close) + def send_ping(socket) = socket.send(JSON.generate(messageType: "ping", data: nil)) - socket.on(:open) { open_handler.call(socket) } - socket.on(:message) { |event| message_handler.call(socket, event) } - socket.on(:error) { |event| error_handler.call(socket, event) } - socket.on(:close) { |event| close_handler.call(event) } + def connect(url, request_headers, &configure) + socket = ClosableWebSocketClient.new + configure.call(socket) + socket.connect(url, headers: request_headers) + socket + end + + def configure_socket(socket, attempt = nil) + 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) + 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) unless attempt end - def handle_socket_message(socket, event) - return if process_message(event.respond_to?(:data) ? event.data : event.to_s) + def handle_socket_message(socket, event, attempt = nil) + if @close_policy.close_frame?(event) + handle_socket_close(event, attempt) + safe_close_socket(socket) unless attempt + return + end - safe_close_socket(socket) unless @closed + process_message(event.respond_to?(:data) ? event.data : event.to_s) end - def handle_socket_error(socket, event) - return if @closed + def handle_socket_error(socket, event, attempt = nil) + return if @lifecycle.stopped? || @close_policy.rejected? - fail_status(event.respond_to?(:message) ? event.message : event, interrupted: true) - safe_close_socket(socket) + error = event.respond_to?(:message) ? event.message : event + fail_status(error, interrupted: true) + attempt&.signal(:failed) + safe_close_socket(socket) unless attempt end - def handle_socket_close(_event) - @status_provider.update(Status::INTERRUPTED, message: "WebSocket closed") unless @closed - end + def handle_socket_close(event, attempt = nil) + return if @lifecycle.stopped? - def socket_closed?(socket) - return socket.closed? if socket.respond_to?(:closed?) - return !socket.open? if socket.respond_to?(:open?) + rejected = @close_policy.reject?(event) + attempt&.signal(rejected ? :rejected : :closed) + return if rejected || @close_policy.rejected? - false - rescue StandardError - true + @status_provider.update(Status::INTERRUPTED, message: "WebSocket closed") end 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 @@ -185,7 +180,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? @@ -193,23 +188,23 @@ 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] } 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) @@ -248,9 +243,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 { @@ -281,6 +274,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/data_sync/web_socket_lifecycle.rb b/lib/featbit/data_sync/web_socket_lifecycle.rb new file mode 100644 index 0000000..17def09 --- /dev/null +++ b/lib/featbit/data_sync/web_socket_lifecycle.rb @@ -0,0 +1,76 @@ +# frozen_string_literal: true + +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 +end diff --git a/spec/client_spec.rb b/spec/client_spec.rb index 928bd7e..b2c1217 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,25 +183,6 @@ def broken_store.segment(_key) = raise("boom") client = described_class.new(options) expect(client.close).to be(true) - end - - it "attempts to close every component once and never raises" 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).once - expect(processor).to have_received(:close).once - expect(client.status_provider.status).to eq(FeatBit::Status::CLOSED) + expect(reentrant_results).to eq([false]) end end diff --git a/spec/web_socket_data_synchronizer_spec.rb b/spec/data_sync/web_socket_connection_spec.rb similarity index 57% rename from spec/web_socket_data_synchronizer_spec.rb rename to spec/data_sync/web_socket_connection_spec.rb index 8a24408..39d7804 100644 --- a/spec/web_socket_data_synchronizer_spec.rb +++ b/spec/data_sync/web_socket_connection_spec.rb @@ -8,140 +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 - }) - expect(synchronizer.process_message(test_bootstrap(test_flag))).to be(true) - 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) - expect { synchronizer.process_message("not-json") }.not_to raise_error - 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" => [] } - } - expect(synchronizer.process_message(message)).to be(true) - 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" => [] } - } - - expect(synchronizer.process_message(message)).to be(true) - 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) - - synchronizer.send(:handle_socket_message, socket, JSON.generate(message)) - - expect(socket).not_to have_received(:close) - expect(store.version).to eq(10) - 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 @@ -229,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| @@ -250,6 +144,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.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(synchronizer.close).to be(true) + 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 @@ -279,10 +251,100 @@ 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 + 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) + 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 + 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) 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..6fdb8c6 --- /dev/null +++ b/spec/data_sync/web_socket_data_synchronizer_spec.rb @@ -0,0 +1,219 @@ +# 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) } + + [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| + 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::STARTING) + 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).not_to have_received(:close) + expect(store.flag("fresh")).to be_nil + end + + 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).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 diff --git a/spec/data_sync/web_socket_shutdown_spec.rb b/spec/data_sync/web_socket_shutdown_spec.rb new file mode 100644 index 0000000..de85682 --- /dev/null +++ b/spec/data_sync/web_socket_shutdown_spec.rb @@ -0,0 +1,334 @@ +# 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 "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) + 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