From 6f57364a367a3c720292676f9875d73aa8efd07b Mon Sep 17 00:00:00 2001 From: Bilal Al-Shahwany Date: Mon, 17 Aug 2026 15:19:14 -0700 Subject: [PATCH 1/2] fixed sse thread leak, workers thread crash and avoid retry if shutdown in progress --- lib/splitclient-rb/sse/event_source/client.rb | 28 +++-- .../sse/workers/segments_worker.rb | 7 +- .../sse/workers/splits_worker.rb | 24 +++- lib/splitclient-rb/version.rb | 2 +- spec/sse/event_source/client_spec.rb | 107 ++++++++++++++++++ spec/sse/workers/segments_worker_spec.rb | 10 ++ spec/sse/workers/splits_worker_spec.rb | 25 ++++ 7 files changed, 188 insertions(+), 15 deletions(-) diff --git a/lib/splitclient-rb/sse/event_source/client.rb b/lib/splitclient-rb/sse/event_source/client.rb index cad1b4b6..74ef024d 100644 --- a/lib/splitclient-rb/sse/event_source/client.rb +++ b/lib/splitclient-rb/sse/event_source/client.rb @@ -33,12 +33,15 @@ def initialize(config, @read_timeout = read_timeout @connected = Concurrent::AtomicBoolean.new(false) @first_event = Concurrent::AtomicBoolean.new(true) + @shutdown = Concurrent::AtomicBoolean.new(false) @socket = nil + @connect_stream_thread = nil end def close(status = nil) return if @socket.nil? + @shutdown.make_true @config.logger.debug("Closing SSEClient socket") if @config.debug_enabled push_status(status) @connected.make_false @@ -75,7 +78,9 @@ def connected? private def connect_thread(latch) + @connect_stream_thread.join if @connect_stream_thread @config.threads[:connect_stream] = Thread.new do + @connect_stream_thread = Thread.current @config.logger.info('Starting connect_stream thread ...') new_status = connect_stream(latch) push_status(new_status) unless new_status.nil? @@ -84,7 +89,7 @@ def connect_thread(latch) end def connect_stream(latch) - return Constants::PUSH_RETRYABLE_ERROR unless socket_write(latch) + return return_retry_if_not_shutdown unless socket_write(latch) while connected? || @first_event.value begin if IO.select([@socket], nil, nil, @read_timeout) @@ -103,27 +108,27 @@ def connect_stream(latch) retry rescue Errno::ETIMEDOUT => e @config.logger.error("SSE read operation timed out!: #{e.inspect}") - return Constants::PUSH_RETRYABLE_ERROR + return return_retry_if_not_shutdown rescue EOFError => e @config.logger.error("SSE read operation EOF, server closed the connection, will reconnect: #{e.inspect}") - return Constants::PUSH_RETRYABLE_ERROR + return return_retry_if_not_shutdown rescue Errno::EBADF, IOError => e @config.logger.error("SSE read operation EBADF or IOError: #{e.inspect}") - return Constants::PUSH_RETRYABLE_ERROR + return return_retry_if_not_shutdown rescue StandardError => e @config.logger.error("SSE read operation StandardError: #{e.inspect}") return nil if ENV['SPLITCLIENT_ENV'] == 'test' @config.logger.error("Error reading partial data: #{e.inspect}") - return Constants::PUSH_RETRYABLE_ERROR + return return_retry_if_not_shutdown end else @config.logger.error("SSE read operation timed out, no data available.") - return Constants::PUSH_RETRYABLE_ERROR + return return_retry_if_not_shutdown end rescue Exception => e @config.logger.debug("SSE socket is not connected: #{e.inspect}") if @config.debug_enabled - return Constants::PUSH_RETRYABLE_ERROR + return return_retry_if_not_shutdown end process_data(partial_data) @@ -133,10 +138,17 @@ def connect_stream(latch) nil end + def return_retry_if_not_shutdown + return nil if @shutdown.value + + Constants::PUSH_RETRYABLE_ERROR + end + def socket_write(latch) @first_event.make_true @socket = socket_connect @socket.puts(build_request(@uri)) + @shutdown.make_false true rescue StandardError => e @config.logger.error("Error during connecting to #{@uri.host}. Error: #{e.inspect}") @@ -155,7 +167,7 @@ def read_first_event(data, latch) if response_code != OK_CODE @config.logger.error("SSE first event failed, code: #{response_code}") latch.count_down - return Constants::PUSH_RETRYABLE_ERROR + return return_retry_if_not_shutdown end @connected.make_true diff --git a/lib/splitclient-rb/sse/workers/segments_worker.rb b/lib/splitclient-rb/sse/workers/segments_worker.rb index 828e7c2d..cc1c7eae 100644 --- a/lib/splitclient-rb/sse/workers/segments_worker.rb +++ b/lib/splitclient-rb/sse/workers/segments_worker.rb @@ -46,7 +46,12 @@ def perform cn = item[:change_number] @config.logger.debug("SegmentsWorker change_number dequeue #{segment_name}, #{cn}") if @config.debug_enabled - @synchronizer.fetch_segment(segment_name, cn) + begin + @synchronizer.fetch_segment(segment_name, cn) + rescue Exception => e + @config.logger.error('Error fetching segments ') + @config.logger.debug("Segment Worker failed to fetch segment: #{e.inspect}") if @config.debug_enabled + end end end diff --git a/lib/splitclient-rb/sse/workers/splits_worker.rb b/lib/splitclient-rb/sse/workers/splits_worker.rb index 7fffc1f0..77ee9172 100644 --- a/lib/splitclient-rb/sse/workers/splits_worker.rb +++ b/lib/splitclient-rb/sse/workers/splits_worker.rb @@ -58,10 +58,10 @@ def perform case notification.data['type'] when SSE::EventSource::EventTypes::SPLIT_UPDATE success = update_feature_flag(notification) - @synchronizer.fetch_splits(notification.data['changeNumber'], 0) unless success + fetch_splits(notification.data['changeNumber'], 0) unless success when SSE::EventSource::EventTypes::RB_SEGMENT_UPDATE success = update_rule_based_segment(notification) - @synchronizer.fetch_splits(0, notification.data['changeNumber']) unless success + fetch_splits(0, notification.data['changeNumber']) unless success when SSE::EventSource::EventTypes::SPLIT_KILL kill_feature_flag(notification) end @@ -125,7 +125,7 @@ def kill_feature_flag(notification) @feature_flags_repository.kill(notification.data['changeNumber'], notification.data['splitName'], notification.data['defaultTreatment']) - @synchronizer.fetch_splits(notification.data['changeNumber'], 0) + fetch_splits(notification.data['changeNumber'], 0) end def return_object_from_json(notification) @@ -137,16 +137,30 @@ def fetch_segments_if_not_exists(segment_names, object_repository) return if segment_names.nil? object_repository.set_segment_names(segment_names) - @segment_fetcher.fetch_segments_if_not_exists(segment_names) + begin + @segment_fetcher.fetch_segments_if_not_exists(segment_names) + rescue Exception => e + @config.logger.error('Error fetching segments ') + @config.logger.debug("Split Worker failed to fetch segment: #{e.inspect}") if @config.debug_enabled + end end def fetch_rule_based_segments_if_not_exists(segment_names, change_number) return false if segment_names.nil? || segment_names.empty? || @rule_based_segment_repository.contains?(segment_names.to_a) - @synchronizer.fetch_splits(0, change_number) + fetch_splits(0, change_number) true end + + def fetch_splits(cn, rbs_cn) + begin + @synchronizer.fetch_splits(cn, rbs_cn) + rescue Exception => e + @config.logger.error('Error fetching feature flags ') + @config.logger.debug("Split Worker failed to fetch feature flags: #{e.inspect}") if @config.debug_enabled + end + end end end end diff --git a/lib/splitclient-rb/version.rb b/lib/splitclient-rb/version.rb index efc0fb8e..ec80d1f1 100644 --- a/lib/splitclient-rb/version.rb +++ b/lib/splitclient-rb/version.rb @@ -1,3 +1,3 @@ module SplitIoClient - VERSION = '8.11.1' + VERSION = '8.11.2' end diff --git a/spec/sse/event_source/client_spec.rb b/spec/sse/event_source/client_spec.rb index 39827dc4..c6476d66 100644 --- a/spec/sse/event_source/client_spec.rb +++ b/spec/sse/event_source/client_spec.rb @@ -49,6 +49,113 @@ let(:event_occupancy) { "d4\r\nevent: message\ndata: {\"id\":\"123\",\"timestamp\":1586803930362,\"encoding\":\"json\",\"channel\":\"[?occupancy=metrics.publishers]control_pri\",\"data\":\"{\\\"metrics\\\":{\\\"publishers\\\":2}}\",\"name\":\"[meta]occupancy\"}\n\n\r\n" } let(:event_error) { "d4\r\nevent: error\ndata: {\"message\":\"Token expired\",\"code\":40142,\"statusCode\":401,\"href\":\"https://help.ably.io/error/40142\"}" } + context 'check connect_stream thread leak via busy process_data' do + let(:log) { StringIO.new } + let(:events_queue) { Queue.new } + let(:config) { SplitIoClient::SplitConfig.new(logger: Logger.new(log), debug_enabled: false) } + let(:telemetry_runtime_producer) { SplitIoClient::Telemetry::RuntimeProducer.new(config) } + let(:api_token) { 'api-token-test' } + let(:event_parser) { SplitIoClient::SSE::EventSource::EventParser.new(config) } + let(:push_status_queue) { Queue.new } + let(:notification_manager_keeper) { SplitIoClient::SSE::NotificationManagerKeeper.new(config, telemetry_runtime_producer, push_status_queue) } + + let(:keepalive) { "c\r\n:keepalive\n\n\r\n" } + + it 'Avoid push retryable when normal shutdown is called' do + mock_server do |server| + server.setup_response('/') do |_, res| + res.content_type = 'text/event-stream' + res.status = 200 + res.chunked = true + rd, wr = IO.pipe + wr.write(keepalive) + res.body = rd + Thread.new do + # keep dribbling data so a live reader always has something to consume + 20.times { sleep 0.5; (wr.write(keepalive) rescue nil) } + wr.close rescue nil + end + end + + sse_client = SplitIoClient::SSE::EventSource::Client.new( + config, api_token, telemetry_runtime_producer, event_parser, + notification_manager_keeper, double(process: true), push_status_queue + ) + + expect(sse_client.start(server.base_uri)).to eq(true) + thread_a = config.threads[:connect_stream] + expect(thread_a.alive?).to eq(true) + + # close + restart + sse_client.close + expect(sse_client.start(server.base_uri)).to eq(true) + thread_b = config.threads[:connect_stream] + expect(thread_b).not_to eq(thread_a) + + push_queue = [] + begin + loop { push_queue << push_status_queue.pop(true) } + rescue ThreadError + # Queue is now empty + end + + expect(thread_a.alive?).to eq(false), 'thread A leaked: it is still running on thread B\'s socket' + expect(push_queue).not_to include(SplitIoClient::Constants::PUSH_RETRYABLE_ERROR) + + sse_client.close + end + end + + it 'old thread survives and reads the new socket' do + mock_server do |server| + server.setup_response('/') do |_, res| + res.content_type = 'text/event-stream' + res.status = 200 + res.chunked = true + rd, wr = IO.pipe + wr.write(keepalive) + res.body = rd + Thread.new do + # keep dribbling data so a live reader always has something to consume + 20.times { sleep 0.5; (wr.write(keepalive) rescue nil) } + wr.close rescue nil + end + end + + sse_client = SplitIoClient::SSE::EventSource::Client.new( + config, api_token, telemetry_runtime_producer, event_parser, + notification_manager_keeper, double(process: true), push_status_queue + ) + + # Make process_data slow, simulating the real SDK doing an HTTP splitChanges + # fetch inside the connect_stream thread. + in_process_data = Queue.new + sse_client.define_singleton_method(:process_data) do |_partial| + in_process_data.push(Thread.current) + sleep 2 + end + + expect(sse_client.start(server.base_uri)).to eq(true) + thread_a = config.threads[:connect_stream] + + # wait until thread A is parked inside process_data + in_process_data.pop + expect(thread_a.alive?).to eq(true) + + # storm pattern: close + immediate restart while A is busy + sse_client.close + expect(sse_client.start(server.base_uri)).to eq(true) + thread_b = config.threads[:connect_stream] + expect(thread_b).not_to eq(thread_a) + sleep 4 # A has long since returned from its 2s process_data + + expect(thread_a.alive?).to eq(false), 'thread A leaked: it is still running on thread B\'s socket' + + sse_client.close + end + end + end + context 'tests' do it 'receive split update event' do stub_request(:get, 'https://sdk.split.io/api/splitChanges?s=1.3&since=-1&rbSince=-1') diff --git a/spec/sse/workers/segments_worker_spec.rb b/spec/sse/workers/segments_worker_spec.rb index 918e4329..2d7d8c60 100644 --- a/spec/sse/workers/segments_worker_spec.rb +++ b/spec/sse/workers/segments_worker_spec.rb @@ -100,6 +100,16 @@ expect(a_request(:get, 'https://sdk.split.io/api/segmentChanges/segment1?since=1470947453877')).to have_been_made.times(1) end + it 'recover from possible fetch exception' do + allow(synchronizer).to receive(:fetch_segment).and_raise(StandardError) + worker = subject.new(synchronizer, config, segments_repository) + worker.start + worker.add_to_queue(1_506_703_262_918, 'segment1') + + sleep 1 + expect(config.threads[:segment_update_worker].status).not_to eq(nil) + end + private def mock_split_changes(splits_json) diff --git a/spec/sse/workers/splits_worker_spec.rb b/spec/sse/workers/splits_worker_spec.rb index b3318365..24848306 100644 --- a/spec/sse/workers/splits_worker_spec.rb +++ b/spec/sse/workers/splits_worker_spec.rb @@ -100,6 +100,17 @@ expect(a_request(:get, 'https://sdk.split.io/api/splitChanges?s=1.3&since=1506703262916&rbSince=-1')).to have_been_made.times(0) end + + it 'recover from possible fetch exception' do + allow(synchronizer).to receive(:fetch_splits).and_raise(StandardError) + + worker = subject.new(synchronizer, config, splits_repository, telemetry_runtime_producer, segment_fetcher, rule_based_segments_repository) + worker.start + worker.add_to_queue(SplitIoClient::SSE::EventSource::StreamData.new("SPLIT_UPDATE", 123, JSON.parse('{"type":"SPLIT_UPDATE","changeNumber":1506703262918}'), 'test')) + sleep 1 + + expect(config.threads[:split_update_worker].status).not_to eq(nil) + end end context 'kill split notification' do @@ -373,6 +384,20 @@ expect(a_request(:get, 'https://sdk.split.io/api/segmentChanges/segment1?since=-1')).to have_been_made.once expect(segments_repository.used_segment_names[1]).to eq('segment1') end + + it 'recover from possible segment fetch exception.' do + stub_request(:get, 'https://sdk.split.io/api/splitChanges?s=1.3&since=1234&rbSince=-1').to_return(status: 200, body: '{"ff":{"d": [],"s": 1234,"t": 1234}, "rbs":{"d":[],"s":-1,"t":-1}}') + stub_request(:get, 'https://sdk.split.io/api/segmentChanges/maur-2?since=-1').to_return(status: 200, body: '{"name":"maur-2","added":["admin"],"removed":[],"since":-1,"till":-1}}') + allow(synchronizer).to receive(:fetch_segment).and_raise(StandardError) + worker = subject.new(synchronizer, config, splits_repository, telemetry_runtime_producer, segment_fetcher, rule_based_segments_repository) + worker.start + + splits_repository.set_change_number(1234) + worker.add_to_queue(event_split_update_segments) + sleep 1 + + expect(config.threads[:split_update_worker].status).not_to eq(nil) + end end private From 87bad31b11bb04287ec51b0aa5cbb574332e2c01 Mon Sep 17 00:00:00 2001 From: Bilal Al-Shahwany Date: Mon, 17 Aug 2026 21:02:30 -0700 Subject: [PATCH 2/2] polishing --- lib/splitclient-rb/sse/workers/segments_worker.rb | 2 +- lib/splitclient-rb/sse/workers/splits_worker.rb | 14 ++++++-------- 2 files changed, 7 insertions(+), 9 deletions(-) diff --git a/lib/splitclient-rb/sse/workers/segments_worker.rb b/lib/splitclient-rb/sse/workers/segments_worker.rb index cc1c7eae..6141e22a 100644 --- a/lib/splitclient-rb/sse/workers/segments_worker.rb +++ b/lib/splitclient-rb/sse/workers/segments_worker.rb @@ -48,7 +48,7 @@ def perform begin @synchronizer.fetch_segment(segment_name, cn) - rescue Exception => e + rescue StandardError => e @config.logger.error('Error fetching segments ') @config.logger.debug("Segment Worker failed to fetch segment: #{e.inspect}") if @config.debug_enabled end diff --git a/lib/splitclient-rb/sse/workers/splits_worker.rb b/lib/splitclient-rb/sse/workers/splits_worker.rb index 77ee9172..816b52b4 100644 --- a/lib/splitclient-rb/sse/workers/splits_worker.rb +++ b/lib/splitclient-rb/sse/workers/splits_worker.rb @@ -139,7 +139,7 @@ def fetch_segments_if_not_exists(segment_names, object_repository) object_repository.set_segment_names(segment_names) begin @segment_fetcher.fetch_segments_if_not_exists(segment_names) - rescue Exception => e + rescue StandardError => e @config.logger.error('Error fetching segments ') @config.logger.debug("Split Worker failed to fetch segment: #{e.inspect}") if @config.debug_enabled end @@ -153,13 +153,11 @@ def fetch_rule_based_segments_if_not_exists(segment_names, change_number) true end - def fetch_splits(cn, rbs_cn) - begin - @synchronizer.fetch_splits(cn, rbs_cn) - rescue Exception => e - @config.logger.error('Error fetching feature flags ') - @config.logger.debug("Split Worker failed to fetch feature flags: #{e.inspect}") if @config.debug_enabled - end + def fetch_splits(change_number, rbs_change_number) + @synchronizer.fetch_splits(change_number, rbs_change_number) + rescue StandardError => e + @config.logger.error('Error fetching feature flags ') + @config.logger.debug("Split Worker failed to fetch feature flags: #{e.inspect}") if @config.debug_enabled end end end