Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 20 additions & 8 deletions lib/splitclient-rb/sse/event_source/client.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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?
Expand All @@ -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)
Expand All @@ -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)
Expand All @@ -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}")
Expand All @@ -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
Expand Down
7 changes: 6 additions & 1 deletion lib/splitclient-rb/sse/workers/segments_worker.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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 StandardError => 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

Expand Down
22 changes: 17 additions & 5 deletions lib/splitclient-rb/sse/workers/splits_worker.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand All @@ -137,16 +137,28 @@ 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 StandardError => 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(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
end
Expand Down
2 changes: 1 addition & 1 deletion lib/splitclient-rb/version.rb
Original file line number Diff line number Diff line change
@@ -1,3 +1,3 @@
module SplitIoClient
VERSION = '8.11.1'
VERSION = '8.11.2'
end
107 changes: 107 additions & 0 deletions spec/sse/event_source/client_spec.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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')
Expand Down
10 changes: 10 additions & 0 deletions spec/sse/workers/segments_worker_spec.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
25 changes: 25 additions & 0 deletions spec/sse/workers/splits_worker_spec.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down