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
2 changes: 2 additions & 0 deletions contract-tests/service.rb
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,8 @@
'flag-change-listeners',
'flag-value-change-listeners',
'fdv1-fallback',
'retry-conformance-fdv1-streaming',
'retry-conformance-fdv1-polling',
],
}.to_json
end
Expand Down
2 changes: 1 addition & 1 deletion launchdarkly-server-sdk.gemspec
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@ Gem::Specification.new do |spec|

spec.add_runtime_dependency "benchmark", "~> 0.1", ">= 0.1.1"
spec.add_runtime_dependency "concurrent-ruby", "~> 1.1"
spec.add_runtime_dependency "ld-eventsource", "2.6.0"
spec.add_runtime_dependency "ld-eventsource", "3.0.0"
spec.add_runtime_dependency "observer", "~> 0.1.2"
spec.add_runtime_dependency "openssl", ">= 3.1.2", "< 5.0"
spec.add_runtime_dependency "semantic", "~> 1.6"
Expand Down
4 changes: 4 additions & 0 deletions lib/ldclient-rb/impl/data_source.rb
Original file line number Diff line number Diff line change
Expand Up @@ -138,6 +138,10 @@ def update_status(new_state, new_error)
@mutex.synchronize do
old_status = @current_status

# A poll or stream connection that was still in flight when the data source stopped must not
# report after OFF.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This comment is true, but also perhaps too specific. It doesn't have to be in flight or stopped, just that nothing else can be reported after Off.

return if old_status.state == LaunchDarkly::Interfaces::DataSource::Status::OFF

if new_state == LaunchDarkly::Interfaces::DataSource::Status::INTERRUPTED && old_status.state == LaunchDarkly::Interfaces::DataSource::Status::INITIALIZING
# See {LaunchDarkly::Interfaces::DataSource::UpdateSink#update_status} for more information
new_state = LaunchDarkly::Interfaces::DataSource::Status::INITIALIZING
Expand Down
28 changes: 11 additions & 17 deletions lib/ldclient-rb/impl/data_source/polling.rb
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
require "ldclient-rb/impl/data_source"
require "ldclient-rb/impl/repeating_task"
require "ldclient-rb/impl/retry_state"
require "ldclient-rb/impl/util"

require "concurrent/atomics"
Expand All @@ -16,7 +17,8 @@ def initialize(config, requestor)
@initialized = Concurrent::AtomicBoolean.new(false)
@started = Concurrent::AtomicBoolean.new(false)
@ready = Concurrent::Event.new
@task = Impl::RepeatingTask.new(@config.poll_interval, 0, -> { self.poll }, @config.logger, 'LD/PollingDataSource')
@retry_state = Impl::RetryState.for_polling(@config.poll_interval, @config.logger)
@task = Impl::RepeatingTask.new(-> { @retry_state.next_delay }, 0, -> { self.poll }, @config.logger, 'LD/PollingDataSource')
end

def initialized?
Expand Down Expand Up @@ -45,9 +47,11 @@ def poll
@ready.set
end
end
@retry_state.record_success
@config.data_source_update_sink&.update_status(LaunchDarkly::Interfaces::DataSource::Status::VALID, nil)
rescue JSON::ParserError => e
@config.logger.error { "[LDClient] JSON parsing failed for polling response." }
@retry_state.record_failure(:normal)
@config.logger.error { "[LDClient] JSON parsing failed for polling response - #{Util.retry_message(@retry_state.next_delay)}" }
error_info = LaunchDarkly::Interfaces::DataSource::ErrorInfo.new(
LaunchDarkly::Interfaces::DataSource::ErrorInfo::INVALID_DATA,
0,
Expand All @@ -56,25 +60,15 @@ def poll
)
@config.data_source_update_sink&.update_status(LaunchDarkly::Interfaces::DataSource::Status::INTERRUPTED, error_info)
rescue Impl::DataSource::UnexpectedResponseError => e
@retry_state.record_failure(Impl::RetryState.classify_http_status(e.status))
error_info = LaunchDarkly::Interfaces::DataSource::ErrorInfo.new(
LaunchDarkly::Interfaces::DataSource::ErrorInfo::ERROR_RESPONSE, e.status, nil, Time.now)
message = Util.http_error_message(e.status, "polling request", "will retry")
message = Util.http_error_retry_message(e.status, "polling request", @retry_state.next_delay)
@config.logger.error { "[LDClient] #{message}" }

if Util.http_error_recoverable?(e.status)
@config.data_source_update_sink&.update_status(
LaunchDarkly::Interfaces::DataSource::Status::INTERRUPTED,
error_info
)
else
# Publish the OFF status before releasing anyone waiting on the
# ready event, so a client that returns from start can rely on the
# data source status already reflecting the failure.
stop_with_error_info error_info
@ready.set # if client was waiting on us, make it stop waiting - has no effect if already set
end
@config.data_source_update_sink&.update_status(LaunchDarkly::Interfaces::DataSource::Status::INTERRUPTED, error_info)
rescue StandardError => e
Impl::Util.log_exception(@config.logger, "Exception while polling", e)
@retry_state.record_failure(:normal)
Impl::Util.log_exception(@config.logger, "Exception while polling - #{Util.retry_message(@retry_state.next_delay)}", e)
@config.data_source_update_sink&.update_status(
LaunchDarkly::Interfaces::DataSource::Status::INTERRUPTED,
LaunchDarkly::Interfaces::DataSource::ErrorInfo.new(LaunchDarkly::Interfaces::DataSource::ErrorInfo::UNKNOWN, 0, e.to_s, Time.now)
Expand Down
95 changes: 62 additions & 33 deletions lib/ldclient-rb/impl/data_source/stream.rb
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
require "ldclient-rb/impl/data_source"
require "ldclient-rb/impl/model/serialization"
require "ldclient-rb/impl/retry_state"
require "ldclient-rb/impl/util"
require "ldclient-rb/in_memory_store"

Expand Down Expand Up @@ -31,6 +32,8 @@ def initialize(sdk_key, config, diagnostic_accumulator = nil)
@started = Concurrent::AtomicBoolean.new(false)
@stopped = Concurrent::AtomicBoolean.new(false)
@ready = Concurrent::Event.new
@stop_event = Concurrent::Event.new
@retry_state = Impl::RetryState.for_streaming(@config.initial_reconnect_delay, @config.logger)
@connection_attempt_start_time = 0
end

Expand All @@ -49,46 +52,17 @@ def start
read_timeout: READ_TIMEOUT_SECONDS,
logger: @config.logger,
socket_factory: @config.socket_factory,
reconnect_time: @config.initial_reconnect_delay,
# The SDK waits in the failure handlers instead. This must be an Integer: the SSE client
# multiplies it by a power of two, and a Float 0.0 becomes NaN once that power overflows.
reconnect_time: 0,

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The mixture of floats and ints throughout the logic is a tad concerning if we anticipate edge cases like this comment describes.

}
log_connection_started

uri = Impl::Util.add_payload_filter_key(@config.stream_uri + "/all", @config)
@es = SSE::Client.new(uri, **opts) do |conn|
conn.on_connect { |response_headers| DataSource.record_environment_id(@data_source_update_sink, response_headers) }
conn.on_event { |event| process_message(event) }
conn.on_error { |err|
log_connection_result(false)
case err
when SSE::Errors::HTTPStatusError
status = err.status
error_info = LaunchDarkly::Interfaces::DataSource::ErrorInfo.new(
LaunchDarkly::Interfaces::DataSource::ErrorInfo::ERROR_RESPONSE, status, nil, Time.now)
message = Util.http_error_message(status, "streaming connection", "will retry")
@config.logger.error { "[LDClient] #{message}" }

if Util.http_error_recoverable?(status)
@data_source_update_sink&.update_status(
LaunchDarkly::Interfaces::DataSource::Status::INTERRUPTED,
error_info
)
else
@ready.set # if client was waiting on us, make it stop waiting - has no effect if already set
stop_with_error_info error_info
end
when SSE::Errors::HTTPContentTypeError, SSE::Errors::HTTPProxyError, SSE::Errors::ReadTimeoutError
@data_source_update_sink&.update_status(
LaunchDarkly::Interfaces::DataSource::Status::INTERRUPTED,
LaunchDarkly::Interfaces::DataSource::ErrorInfo.new(LaunchDarkly::Interfaces::DataSource::ErrorInfo::NETWORK_ERROR, 0, err.to_s, Time.now)
)

else
@data_source_update_sink&.update_status(
LaunchDarkly::Interfaces::DataSource::Status::INTERRUPTED,
LaunchDarkly::Interfaces::DataSource::ErrorInfo.new(LaunchDarkly::Interfaces::DataSource::ErrorInfo::UNKNOWN, 0, err.to_s, Time.now)
)
end
}
conn.on_error { |err| handle_error(err) }
end

@ready
Expand All @@ -106,11 +80,65 @@ def stop
def stop_with_error_info(error_info = nil)
if @stopped.make_true
@es.close
@stop_event.set
@data_source_update_sink&.update_status(LaunchDarkly::Interfaces::DataSource::Status::OFF, error_info)
@config.logger.info { "[LDClient] Stream connection stopped" }
end
end

def handle_error(err)
if err.is_a?(SSE::Errors::HTTPStatusError)
status = err.status
error_info = LaunchDarkly::Interfaces::DataSource::ErrorInfo.new(
LaunchDarkly::Interfaces::DataSource::ErrorInfo::ERROR_RESPONSE, status, nil, Time.now)
handle_failure(Impl::RetryState.classify_http_status(status), error_info) do |delay|
@config.logger.error { "[LDClient] #{Util.http_error_retry_message(status, 'streaming connection', delay)}" }
end
return
end

if err.is_a?(SSE::Errors::StreamClosedByServerError)
error_info = LaunchDarkly::Interfaces::DataSource::ErrorInfo.new(
LaunchDarkly::Interfaces::DataSource::ErrorInfo::NETWORK_ERROR, 0, err.to_s, Time.now)
handle_failure(:normal, error_info) do |delay|
@config.logger.warn { "[LDClient] The server closed the streaming connection - #{Util.retry_message(delay)}" }
end
return
end

error_kind = case err
when SSE::Errors::HTTPContentTypeError, SSE::Errors::HTTPProxyError, SSE::Errors::ReadTimeoutError
LaunchDarkly::Interfaces::DataSource::ErrorInfo::NETWORK_ERROR
else
LaunchDarkly::Interfaces::DataSource::ErrorInfo::UNKNOWN
end
error_info = LaunchDarkly::Interfaces::DataSource::ErrorInfo.new(error_kind, 0, err.to_s, Time.now)
handle_failure(:normal, error_info) do |delay|
@config.logger.warn { "[LDClient] Error on streaming connection: #{err} - #{Util.retry_message(delay)}" }
end
end

#
# Records a failure, reports it, and waits before the next attempt. The SSE client calls this on its
# worker thread before it reconnects, so the wait is the reconnect delay, and {#stop} ends it at once.
#
# @param kind [Symbol] `:normal` or `:unexpected`
# @param error_info [LaunchDarkly::Interfaces::DataSource::ErrorInfo]
# @yieldparam delay [Numeric] seconds until the next attempt, for the log message
#
def handle_failure(kind, error_info)
return if @stopped.value

log_connection_result(false)
@retry_state.record_failure(kind)
delay = @retry_state.next_delay
yield delay
@data_source_update_sink&.update_status(LaunchDarkly::Interfaces::DataSource::Status::INTERRUPTED, error_info)

@stop_event.wait([delay, Impl::RetryState::MAX_WAIT].min)
log_connection_started
end

#
# The original implementation of this class relied on the feature store
# directly, which we are trying to move away from. Customers who might have
Expand Down Expand Up @@ -162,6 +190,7 @@ def process_message(message)
@config.logger.warn { "[LDClient] Unknown message received: #{method}" }
end

@retry_state.record_success
@data_source_update_sink&.update_status(LaunchDarkly::Interfaces::DataSource::Status::VALID, nil)
rescue JSON::ParserError => e
@config.logger.error { "[LDClient] JSON parsing failed for method #{method}. Ignoring event." }
Expand Down
3 changes: 2 additions & 1 deletion lib/ldclient-rb/impl/data_system/streaming.rb
Original file line number Diff line number Diff line change
Expand Up @@ -332,7 +332,8 @@ def stop

@logger.warn { "[LDClient] #{error_info.message}" }

when SSE::Errors::HTTPContentTypeError, SSE::Errors::HTTPProxyError, SSE::Errors::ReadTimeoutError
when SSE::Errors::HTTPContentTypeError, SSE::Errors::HTTPProxyError, SSE::Errors::ReadTimeoutError,
SSE::Errors::StreamClosedByServerError
@logger.warn { "[LDClient] Network error on stream connection: #{error}, will retry" }

update = LaunchDarkly::Interfaces::DataSystem::Update.new(
Expand Down
14 changes: 8 additions & 6 deletions lib/ldclient-rb/impl/repeating_task.rb
Original file line number Diff line number Diff line change
@@ -1,12 +1,15 @@
require "ldclient-rb/impl/retry_state"
require "ldclient-rb/impl/util"

require "concurrent/atomics"

module LaunchDarkly
module Impl
#
# Runs a task again and again on a worker thread, with a fixed interval
# between runs.
# Runs a task again and again on a worker thread.
#
# The interval is read after each run, and the wait starts when the run
# returns.
#
# The worker waits on an event instead of calling `sleep`, so `stop` can
# wake it at once even if `stop` runs before the worker starts waiting.
Expand All @@ -17,7 +20,7 @@ class RepeatingTask
attr_reader :name

#
# @param interval [Numeric] seconds between the start of one run and the start of the next
# @param interval [Numeric, #call] seconds between runs, or an object that returns them
# @param start_delay [Numeric, nil] seconds to wait before the first run
# @param task [Proc] the code to run
# @param logger [Logger]
Expand All @@ -39,14 +42,13 @@ def start
@stop_event.wait(@start_delay) unless @start_delay.nil? || @start_delay == 0

until @stopped.value do
started_at = Time.now
begin
@task.call
rescue => e
Impl::Util.log_exception(@logger, "Uncaught exception from repeating task", e)
end
delta = @interval - (Time.now - started_at)
@stop_event.wait(delta) if delta > 0
delay = @interval.respond_to?(:call) ? @interval.call : @interval
@stop_event.wait([delay, RetryState::MAX_WAIT].min) if delay > 0
end
end

Expand Down
Loading
Loading