From 7c69595218d589523b3d3f716d837506206b0226 Mon Sep 17 00:00:00 2001 From: jsonbailey Date: Fri, 25 Sep 2026 15:24:51 -0500 Subject: [PATCH 1/3] feat: Retry indefinitely after a data source failure instead of stopping permanently The FDv1 streaming and polling data sources follow the RETRY spec. No HTTP or transport failure stops them. An unexpected status such as 401 moves them to a longer backoff, and the SDK keeps retrying. OFF is reported only after shutdown. --- contract-tests/service.rb | 2 + lib/ldclient-rb/impl/data_source.rb | 4 + lib/ldclient-rb/impl/data_source/polling.rb | 28 +- lib/ldclient-rb/impl/data_source/stream.rb | 95 ++++--- lib/ldclient-rb/impl/data_system/streaming.rb | 3 +- lib/ldclient-rb/impl/repeating_task.rb | 14 +- lib/ldclient-rb/impl/retry_state.rb | 264 ++++++++++++++++++ lib/ldclient-rb/impl/util.rb | 19 ++ lib/ldclient-rb/interfaces/data_source.rb | 11 +- lib/ldclient-rb/ldclient.rb | 5 +- spec/impl/data_source/polling_spec.rb | 100 ++++++- spec/impl/data_source/stream_spec.rb | 219 +++++++++++++-- spec/impl/data_source_spec.rb | 27 ++ .../streaming_synchronizer_spec.rb | 14 + spec/impl/repeating_task_spec.rb | 57 ++++ spec/impl/retry_state_spec.rb | 229 +++++++++++++++ spec/ldclient_end_to_end_spec.rb | 8 +- 17 files changed, 1001 insertions(+), 98 deletions(-) create mode 100644 lib/ldclient-rb/impl/retry_state.rb create mode 100644 spec/impl/retry_state_spec.rb diff --git a/contract-tests/service.rb b/contract-tests/service.rb index 779ed37d..1f9d45e0 100644 --- a/contract-tests/service.rb +++ b/contract-tests/service.rb @@ -55,6 +55,8 @@ 'flag-change-listeners', 'flag-value-change-listeners', 'fdv1-fallback', + 'retry-conformance-fdv1-streaming', + 'retry-conformance-fdv1-polling', ], }.to_json end diff --git a/lib/ldclient-rb/impl/data_source.rb b/lib/ldclient-rb/impl/data_source.rb index e64f860c..82139a1a 100644 --- a/lib/ldclient-rb/impl/data_source.rb +++ b/lib/ldclient-rb/impl/data_source.rb @@ -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. + 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 diff --git a/lib/ldclient-rb/impl/data_source/polling.rb b/lib/ldclient-rb/impl/data_source/polling.rb index 6f01138b..ae45ac91 100644 --- a/lib/ldclient-rb/impl/data_source/polling.rb +++ b/lib/ldclient-rb/impl/data_source/polling.rb @@ -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" @@ -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? @@ -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, @@ -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) diff --git a/lib/ldclient-rb/impl/data_source/stream.rb b/lib/ldclient-rb/impl/data_source/stream.rb index c7b4cb2d..5ed9dea9 100644 --- a/lib/ldclient-rb/impl/data_source/stream.rb +++ b/lib/ldclient-rb/impl/data_source/stream.rb @@ -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" @@ -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 @@ -49,7 +52,9 @@ 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, } log_connection_started @@ -57,38 +62,7 @@ def start @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 @@ -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::StreamClosedError) + 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 @@ -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." } diff --git a/lib/ldclient-rb/impl/data_system/streaming.rb b/lib/ldclient-rb/impl/data_system/streaming.rb index 78c74d63..c8642245 100644 --- a/lib/ldclient-rb/impl/data_system/streaming.rb +++ b/lib/ldclient-rb/impl/data_system/streaming.rb @@ -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::StreamClosedError @logger.warn { "[LDClient] Network error on stream connection: #{error}, will retry" } update = LaunchDarkly::Interfaces::DataSystem::Update.new( diff --git a/lib/ldclient-rb/impl/repeating_task.rb b/lib/ldclient-rb/impl/repeating_task.rb index f11e4dc2..bd5a6df9 100644 --- a/lib/ldclient-rb/impl/repeating_task.rb +++ b/lib/ldclient-rb/impl/repeating_task.rb @@ -1,3 +1,4 @@ +require "ldclient-rb/impl/retry_state" require "ldclient-rb/impl/util" require "concurrent/atomics" @@ -5,8 +6,10 @@ 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. @@ -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] @@ -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 diff --git a/lib/ldclient-rb/impl/retry_state.rb b/lib/ldclient-rb/impl/retry_state.rb new file mode 100644 index 00000000..53e673ee --- /dev/null +++ b/lib/ldclient-rb/impl/retry_state.rb @@ -0,0 +1,264 @@ +module LaunchDarkly + module Impl + # + # Computes how long to wait before a failed operation is tried again. + # + # Each failure is `:normal` or `:unexpected`. After a normal failure, the wait starts at the normal + # initial delay and doubles with each failure, up to the normal ceiling. An unexpected failure moves + # the state to the extended regime: the wait starts at the extended initial delay and doubles up to + # the extended ceiling. The extended bounds stay in place until the reset policy is satisfied, so a + # normal failure that follows cannot lower them. When the policy is satisfied, the state returns to + # the normal regime and the delay sequence starts over. + # + # A random jitter of up to half of each delay is subtracted, so that many callers do not all try + # again at the same moment. The wait never falls below the operating cadence. + # + # After each outcome, call {#record_failure} or {#record_success}, then read {#next_delay} for the + # wait before the next operation. An instance is not thread-safe. + # + # @private + # + class RetryState + # The longest normal delay for streaming, in seconds. + NORMAL_STREAMING_CEILING_DELAY = 30 + + # The delay bounds of the extended regime, in seconds. + EXTENDED_INITIAL_DELAY = 5 * 60 + EXTENDED_CEILING_DELAY = 60 * 60 + + # How long a stream must operate without a failure before its retry state resets, in seconds. + STREAMING_RESET_INTERVAL = 60 + + # How many polls in a row must succeed before polling's retry state resets. + POLLING_RESET_SUCCESSES = 2 + + # The longest wait a caller should pass to `Concurrent::Event#wait`, in seconds. A much longer + # wait raises `RangeError`. This is about 31 years, so the bound has no effect in practice. + MAX_WAIT = 1_000_000_000 + + # The 4xx statuses that are still normal failures. Every other 4xx is unexpected. + NORMAL_4XX_STATUSES = [400, 408, 429].freeze + private_constant :NORMAL_4XX_STATUSES + + MONOTONIC_CLOCK = -> { Process.clock_gettime(Process::CLOCK_MONOTONIC) } + private_constant :MONOTONIC_CLOCK + + # + # Classifies an HTTP status. + # + # `400`, `408`, `429` and every `5xx` are normal. Every other `4xx`, including `401` and `403`, + # is unexpected. + # + # @param status [Integer] + # @return [Symbol] `:normal` or `:unexpected` + # + def self.classify_http_status(status) + if status >= 400 && status < 500 && !NORMAL_4XX_STATUSES.include?(status) + :unexpected + else + :normal + end + end + + # + # Builds the retry state for a streaming data source. + # + # A healthy stream never waits, so the operating cadence is zero. A delay that is not a + # positive, finite number is replaced by the default. A delay longer than a ceiling raises that + # ceiling. + # + # @param initial_reconnect_delay [Numeric] seconds + # @param logger [Logger] + # @param clock [#call] returns monotonic seconds + # @param random [#rand] returns a Float in [0, 1) + # @return [RetryState] + # + def self.for_streaming(initial_reconnect_delay, logger, clock: MONOTONIC_CLOCK, random: Random.new) + delay = usable_delay(initial_reconnect_delay, Config.default_initial_reconnect_delay, "initial_reconnect_delay", logger) + new( + normal_initial_delay: delay, + normal_ceiling_delay: NORMAL_STREAMING_CEILING_DELAY, + extended_initial_delay: [EXTENDED_INITIAL_DELAY, delay].max, + extended_ceiling_delay: EXTENDED_CEILING_DELAY, + reset_policy: AfterHealthyFor.new(STREAMING_RESET_INTERVAL, clock), + operating_cadence: 0, + random: random + ) + end + + # + # Builds the retry state for a polling data source. + # + # The poll interval is both the operating cadence and the normal ceiling, so a normal failure + # waits the interval. An interval that is not a positive, finite number is replaced by the + # default. No wait is shorter than the interval, so the interval wins over the extended ceiling. + # + # @param poll_interval [Numeric] seconds + # @param logger [Logger] + # @param random [#rand] returns a Float in [0, 1) + # @return [RetryState] + # + def self.for_polling(poll_interval, logger, random: Random.new) + interval = usable_delay(poll_interval, Config.default_poll_interval, "poll_interval", logger) + new( + normal_initial_delay: interval, + normal_ceiling_delay: interval, + extended_initial_delay: [EXTENDED_INITIAL_DELAY, interval].max, + extended_ceiling_delay: EXTENDED_CEILING_DELAY, + reset_policy: AfterConsecutiveSuccesses.new(POLLING_RESET_SUCCESSES), + operating_cadence: interval, + random: random + ) + end + + private_class_method def self.usable_delay(value, default, name, logger) + return value if value.is_a?(Numeric) && value.real? && value > 0 && value.finite? + + logger.warn { "[LDClient] #{name} must be a positive, finite number of seconds; using the default of #{default}s" } + default + end + + # + # @param normal_initial_delay [Numeric] the first normal delay, in seconds + # @param normal_ceiling_delay [Numeric] the longest normal delay, in seconds + # @param extended_initial_delay [Numeric] the first extended delay, in seconds + # @param extended_ceiling_delay [Numeric] the longest extended delay, in seconds + # @param reset_policy [AfterHealthyFor, AfterConsecutiveSuccesses] decides when the state resets + # @param operating_cadence [Numeric] the wait between healthy operations, in seconds; no wait is + # shorter than this + # @param random [#rand] returns a Float in [0, 1) + # + def initialize(normal_initial_delay:, normal_ceiling_delay:, extended_initial_delay:, extended_ceiling_delay:, + reset_policy:, operating_cadence: 0, random: Random.new) + @normal_initial_delay = normal_initial_delay + @normal_ceiling_delay = normal_ceiling_delay + @extended_initial_delay = extended_initial_delay + @extended_ceiling_delay = extended_ceiling_delay + @reset_policy = reset_policy + @operating_cadence = operating_cadence + @random = random + + @attempts = 0 + @extended = false + @min_delay = @normal_initial_delay + @max_delay = [@normal_ceiling_delay, @normal_initial_delay].max + @next_delay = @operating_cadence + end + + # + # @return [Numeric] the wait before the next operation, in seconds, as the last recorded + # outcome decided it + # + attr_reader :next_delay + + # + # Records a failed attempt and decides the wait before the next one. + # + # @param kind [Symbol] `:normal` or `:unexpected` + # @return [void] + # + def record_failure(kind) + # A caller can record nothing while it is healthy, so a time-based reset can only be noticed here. + reset_if_due + @reset_policy.note_failure + + if kind == :unexpected && !@extended + # Only the move to the extended regime starts the sequence over. A later unexpected failure + # keeps counting up. + @extended = true + @min_delay = @extended_initial_delay + @max_delay = [@extended_ceiling_delay, @min_delay].max + @attempts = 1 + else + @attempts += 1 + end + + # Integer exponentiation does not overflow, and a Float product that becomes Infinity is + # still cut down to the ceiling. + delay = [@min_delay * (2**(@attempts - 1)), @max_delay].min + jitter = @random.rand * delay / 2 + @next_delay = [delay - jitter, @operating_cadence].max + end + + # + # Records a successful operation, and resets the retry state if the reset policy is satisfied. + # + # The next wait returns to the operating cadence even when the state has not reset, because a + # backoff delay applies to a retry and not to every operation. + # + # @return [void] + # + def record_success + @reset_policy.note_healthy + reset_if_due + @next_delay = @operating_cadence + end + + private def reset_if_due + return unless @reset_policy.satisfied? + + @attempts = 0 + @extended = false + @min_delay = @normal_initial_delay + @max_delay = [@normal_ceiling_delay, @normal_initial_delay].max + end + + # + # Resets once the component has operated without a failure for a number of seconds. + # + # @private + # + class AfterHealthyFor + # + # @param seconds [Numeric] + # @param clock [#call] returns monotonic seconds + # + def initialize(seconds, clock) + @healthy_seconds = seconds + @clock = clock + @healthy_since = nil + end + + # A later call during the same healthy stretch does not move its start. + def note_healthy + @healthy_since = @clock.call if @healthy_since.nil? + end + + def note_failure + @healthy_since = nil + end + + def satisfied? + !@healthy_since.nil? && @clock.call - @healthy_since >= @healthy_seconds + end + end + + # + # Resets once a number of operations in a row have succeeded. + # + # @private + # + class AfterConsecutiveSuccesses + # + # @param count [Integer] + # + def initialize(count) + @count = count + @successes = 0 + end + + def note_healthy + @successes += 1 + end + + def note_failure + @successes = 0 + end + + def satisfied? + @successes >= @count + end + end + end + end +end diff --git a/lib/ldclient-rb/impl/util.rb b/lib/ldclient-rb/impl/util.rb index 6748295a..5c84109a 100644 --- a/lib/ldclient-rb/impl/util.rb +++ b/lib/ldclient-rb/impl/util.rb @@ -164,6 +164,25 @@ def self.http_error_message(status, context, recoverable_message) message = http_error_recoverable?(status) ? recoverable_message : "giving up permanently" "HTTP error #{status}#{desc} for #{context} - #{message}" end + + # + # @param status [Integer] + # @param context [String] what failed, such as "polling request" + # @param delay [Numeric] seconds until the next attempt + # @return [String] + # + def self.http_error_retry_message(status, context, delay) + desc = (status == 401 || status == 403) ? " (invalid SDK key)" : "" + "HTTP error #{status}#{desc} for #{context} - #{retry_message(delay)}" + end + + # + # @param delay [Numeric] seconds until the next attempt + # @return [String] + # + def self.retry_message(delay) + format("will retry in %.1fs", delay) + end end end end diff --git a/lib/ldclient-rb/interfaces/data_source.rb b/lib/ldclient-rb/interfaces/data_source.rb index 27ee5769..0f34384c 100644 --- a/lib/ldclient-rb/interfaces/data_source.rb +++ b/lib/ldclient-rb/interfaces/data_source.rb @@ -183,16 +183,19 @@ class Status # # In streaming mode, this means that the stream connection failed, or had to be dropped due to some # other error, and will be retried after a backoff delay. In polling mode, it means that the last poll - # request failed, and a new poll request will be made after the configured polling interval. + # request failed, and a new poll request will be made after the polling interval, or after a longer + # backoff delay for unexpected error types. # INTERRUPTED = :interrupted # # Indicates that the data source has been permanently shut down. # - # This could be because it encountered an unrecoverable error (for instance, the LaunchDarkly service - # rejected the SDK key; an invalid SDK key will never become valid), or because the SDK client was - # explicitly shut down. + # This could be because the SDK client was explicitly shut down. In the FDv2 data system, it can also + # mean the data source stopped after an error it does not retry. + # + # Unless the SDK is configured with {LaunchDarkly::Config#data_system_config}, no other state follows + # this one. # OFF = :off diff --git a/lib/ldclient-rb/ldclient.rb b/lib/ldclient-rb/ldclient.rb index bbe65e48..63df8c9b 100644 --- a/lib/ldclient-rb/ldclient.rb +++ b/lib/ldclient-rb/ldclient.rb @@ -303,8 +303,9 @@ def secure_mode_hash(context) # If this returns false, it means that the client did not succeed in connecting to # LaunchDarkly within the time limit that you specified in the constructor. It could # still succeed in connecting at a later time (on another thread), or it could have - # given up permanently (for instance, if your SDK key is invalid). In the meantime, - # any call to {#variation} or {#variation_detail} will behave as follows: + # received an error that needs to be fixed (for instance, if your SDK key is + # invalid). In the meantime, any call to {#variation} or {#variation_detail} will + # behave as follows: # # 1. It will check whether the feature store already contains data (that is, you # are using a database-backed store and it was populated by a previous run of this diff --git a/spec/impl/data_source/polling_spec.rb b/spec/impl/data_source/polling_spec.rb index 3dcfdc30..8723d318 100644 --- a/spec/impl/data_source/polling_spec.rb +++ b/spec/impl/data_source/polling_spec.rb @@ -136,22 +136,34 @@ def with_processor(store, initialize_to_valid = false) end describe 'HTTP errors' do - def verify_unrecoverable_http_error(status) - allow(requestor).to receive(:request_all_data).and_raise(Impl::DataSource::UnexpectedResponseError.new(status)) + # Delays short enough that the task polls again at once. + let(:fast_retry_state) { + Impl::RetryState.new(normal_initial_delay: 0.001, normal_ceiling_delay: 0.001, extended_initial_delay: 0.002, + extended_ceiling_delay: 0.002, reset_policy: Impl::RetryState::AfterConsecutiveSuccesses.new(2), + operating_cadence: 0.001) + } + + def verify_unexpected_http_error_keeps_retrying(status) + allow(Impl::RetryState).to receive(:for_polling).and_return(fast_retry_state) + attempts = Concurrent::CountDownLatch.new(3) + allow(requestor).to receive(:request_all_data) do + attempts.count_down + raise Impl::DataSource::UnexpectedResponseError.new(status) + end listener = ListenerSpy.new status_broadcaster.add_listener(listener) - with_processor(InMemoryFeatureStore.new) do |processor| + with_processor(InMemoryFeatureStore.new, true) do |processor| ready = processor.start - finished = ready.wait(1) - expect(finished).to be true + expect(attempts.wait(1)).to be true + expect(ready.set?).to be false expect(processor.initialized?).to be false - expect(listener.statuses.count).to eq(1) - - s = listener.statuses[0] - expect(s.state).to eq(Interfaces::DataSource::Status::OFF) - expect(s.last_error.status_code).to eq(status) + # The first status is the VALID that with_processor sets. + states = listener.statuses.map(&:state) + expect(states[1..2]).to eq([Interfaces::DataSource::Status::INTERRUPTED] * 2) + expect(states).not_to include(Interfaces::DataSource::Status::OFF) + expect(listener.statuses[1].last_error.status_code).to eq(status) end end @@ -174,12 +186,16 @@ def verify_recoverable_http_error(status) end end - it 'stops immediately for error 401' do - verify_unrecoverable_http_error(401) + it 'keeps retrying after error 401' do + verify_unexpected_http_error_keeps_retrying(401) + end + + it 'keeps retrying after error 403' do + verify_unexpected_http_error_keeps_retrying(403) end - it 'stops immediately for error 403' do - verify_unrecoverable_http_error(403) + it 'keeps retrying after error 404' do + verify_unexpected_http_error_keeps_retrying(404) end it 'does not stop immediately for error 408' do @@ -195,6 +211,62 @@ def verify_recoverable_http_error(status) end end + describe 'retry delay' do + let(:logger) { double("logger").as_null_object } + let(:all_data) { { Impl::DataStore::FEATURES => {}, Impl::DataStore::SEGMENTS => {} } } + + def make_processor + config = Config.new(feature_store: InMemoryFeatureStore.new, logger: logger) + config.data_source_update_sink = Impl::DataSource::UpdateSink.new(config.feature_store, status_broadcaster, flag_change_broadcaster) + subject.new(config, requestor) + end + + def next_delay(processor) + processor.instance_variable_get(:@retry_state).next_delay + end + + it 'waits in the extended regime after an unexpected error, then the poll interval after a success' do + responses = [:unauthorized, :ok] + allow(requestor).to receive(:request_all_data) do + raise Impl::DataSource::UnexpectedResponseError.new(401) if responses.shift == :unauthorized + all_data + end + processor = make_processor + + processor.poll + expect(next_delay(processor)).to be_between(150, 300) + processor.poll + expect(next_delay(processor)).to eq(Config.default_poll_interval) + end + + it 'waits the poll interval after a recoverable error' do + allow(requestor).to receive(:request_all_data).and_raise(Impl::DataSource::UnexpectedResponseError.new(503)) + processor = make_processor + + processor.poll + expect(next_delay(processor)).to eq(Config.default_poll_interval) + end + + it 'waits the poll interval after a network error' do + allow(requestor).to receive(:request_all_data).and_raise(StandardError.new("test error")) + processor = make_processor + + processor.poll + expect(next_delay(processor)).to eq(Config.default_poll_interval) + end + + it 'logs the real delay for an HTTP error' do + allow(Impl::RetryState).to receive(:for_polling).and_return( + Impl::RetryState.for_polling(30, logger, random: double("random", rand: 0.0))) + allow(requestor).to receive(:request_all_data).and_raise(Impl::DataSource::UnexpectedResponseError.new(401)) + expect(logger).to receive(:error) do |&block| + expect(block.call).to eq("[LDClient] HTTP error 401 (invalid SDK key) for polling request - will retry in 300.0s") + end + + make_processor.poll + end + end + describe 'stop' do it 'stops promptly rather than continuing to wait for poll interval' do listener = ListenerSpy.new diff --git a/spec/impl/data_source/stream_spec.rb b/spec/impl/data_source/stream_spec.rb index 12080697..572aba5a 100644 --- a/spec/impl/data_source/stream_spec.rb +++ b/spec/impl/data_source/stream_spec.rb @@ -68,39 +68,220 @@ module LaunchDarkly end end - describe 'environment ID' do - def with_connect_handler(processor) - connection = double("connection", on_event: nil, on_error: nil) - handler = nil - allow(connection).to receive(:on_connect) { |&block| handler = block } - allow(SSE::Client).to receive(:new) do |_uri, **_opts, &block| - block.call(connection) - double("SSE::Client", close: nil) - end + # Replaces the SSE client with a double, and yields the handlers the processor registers on it. + def with_handlers(processor) + handlers = {} + connection = double("connection") + %i[on_connect on_event on_error].each do |name| + allow(connection).to receive(name) { |&block| handlers[name] = block } + end + sse_client = double("SSE::Client", close: nil) + allow(SSE::Client).to receive(:new) do |_uri, **opts, &block| + handlers[:opts] = opts + block.call(connection) + sse_client + end - processor.start - begin - yield handler - ensure - processor.stop - end + processor.start + begin + yield handlers, sse_client + ensure + processor.stop end + end + describe 'environment ID' do it 'is recorded from the connection response headers' do - with_connect_handler(processor) do |handler| - handler.call({ "X-LD-EnvID" => "env-abc" }) + with_handlers(processor) do |handlers| + handlers[:on_connect].call({ "X-LD-EnvID" => "env-abc" }) expect(config.data_source_update_sink.environment_id).to eq("env-abc") end end it 'is not recorded when the header is absent' do - with_connect_handler(processor) do |handler| - handler.call({}) + with_handlers(processor) do |handlers| + handlers[:on_connect].call({}) expect(config.data_source_update_sink.environment_id).to be_nil end end end + describe 'retry' do + let(:put_message) { SSE::StreamEvent.new(:put, '{"data":{"flags":{},"segments":{}}}') } + let(:listener) { ListenerSpy.new } + let(:logger) { double("logger").as_null_object } + let(:config) { + config = Config.new(logger: logger) + config.data_source_update_sink = Impl::DataSource::UpdateSink.new(config.feature_store, status_broadcaster, flag_change_broadcaster) + config.data_source_update_sink.update_status(Interfaces::DataSource::Status::VALID, nil) + status_broadcaster.add_listener(listener) + config + } + # Delays short enough that a failure handler returns at once. + let(:fast_retry_state) { + Impl::RetryState.new(normal_initial_delay: 0.001, normal_ceiling_delay: 0.001, extended_initial_delay: 0.002, + extended_ceiling_delay: 0.002, reset_policy: Impl::RetryState::AfterConsecutiveSuccesses.new(2)) + } + + def use_fast_retry_state + allow(Impl::RetryState).to receive(:for_streaming).and_return(fast_retry_state) + end + + def http_error(status) + SSE::Errors::HTTPStatusError.new(status, "") + end + + def states + listener.statuses.map(&:state) + end + + it 'passes an Integer zero reconnect time to the SSE client' do + with_handlers(processor) do |handlers| + expect(handlers[:opts][:reconnect_time]).to eql(0) + end + end + + [401, 403, 404].each do |status| + it "keeps retrying after error #{status}" do + use_fast_retry_state + with_handlers(processor) do |handlers, sse_client| + 3.times { handlers[:on_error].call(http_error(status)) } + + expect(sse_client).not_to have_received(:close) + expect(states).to eq([Interfaces::DataSource::Status::INTERRUPTED] * 3) + expect(listener.statuses.last.last_error.status_code).to eq(status) + expect(processor.instance_variable_get(:@ready).set?).to be false + end + end + end + + it 'waits in the extended regime after an unexpected error' do + with_handlers(processor) do |handlers| + waiter = Thread.new { handlers[:on_error].call(http_error(401)) } + begin + expect(waiter.join(0.2)).to be_nil + expect(processor.instance_variable_get(:@retry_state).next_delay).to be_between(150, 300) + expect(states).to eq([Interfaces::DataSource::Status::INTERRUPTED]) + ensure + processor.stop + waiter.join(1) + end + end + end + + it 'ends a long wait at once when stopped' do + with_handlers(processor) do |handlers| + waiter = Thread.new { handlers[:on_error].call(http_error(401)) } + expect(waiter.join(0.2)).to be_nil + + started_at = Time.now + processor.stop + expect(waiter.join(1)).not_to be_nil + expect(Time.now - started_at).to be < 1 + expect(states.last).to eq(Interfaces::DataSource::Status::OFF) + end + end + + it 'waits the normal delay after a recoverable error' do + with_handlers(processor) do |handlers| + waiter = Thread.new { handlers[:on_error].call(http_error(503)) } + begin + expect(waiter.join(0.2)).to be_nil + expect(processor.instance_variable_get(:@retry_state).next_delay).to be_between(0.5, 1) + ensure + processor.stop + waiter.join(1) + end + end + end + + it 'backs off and warns when the server closes the stream' do + use_fast_retry_state + expect(logger).to receive(:warn) do |&block| + expect(block.call).to match(/server closed the streaming connection - will retry in 0\.0s/) + end + expect(fast_retry_state).to receive(:record_failure).with(:normal).and_call_original + + with_handlers(processor) do |handlers| + handlers[:on_error].call(SSE::Errors::StreamClosedError.new) + expect(states).to eq([Interfaces::DataSource::Status::INTERRUPTED]) + expect(listener.statuses.last.last_error.kind).to eq(Interfaces::DataSource::ErrorInfo::NETWORK_ERROR) + end + end + + it 'logs the real delay for an HTTP error' do + allow(Impl::RetryState).to receive(:for_streaming).and_return( + Impl::RetryState.new(normal_initial_delay: 0.001, normal_ceiling_delay: 0.001, extended_initial_delay: 0.3, + extended_ceiling_delay: 0.3, reset_policy: Impl::RetryState::AfterConsecutiveSuccesses.new(2), + random: double("random", rand: 0.0))) + expect(logger).to receive(:error) do |&block| + expect(block.call).to eq("[LDClient] HTTP error 401 (invalid SDK key) for streaming connection - will retry in 0.3s") + end + + with_handlers(processor) do |handlers| + handlers[:on_error].call(http_error(401)) + end + end + + it 'resets to the normal regime 60 seconds after the first healthy event that follows a failure' do + now = 0.0 + retry_state = Impl::RetryState.for_streaming(1, logger, clock: -> { now }, random: double("random", rand: 0.0)) + allow(Impl::RetryState).to receive(:for_streaming).and_return(retry_state) + retry_state.record_failure(:unexpected) + + with_handlers(processor) do |handlers| + handlers[:on_connect].call({}) + handlers[:on_event].call(put_message) + now += 30 + handlers[:on_event].call(put_message) + now += 30 + + waiter = Thread.new { handlers[:on_error].call(http_error(503)) } + begin + expect(waiter.join(0.2)).to be_nil + expect(retry_state.next_delay).to eq(1) + ensure + processor.stop + waiter.join(1) + end + end + end + + it 'does not record a success for an event that fails' do + use_fast_retry_state + expect(fast_retry_state).not_to receive(:record_success) + + with_handlers(processor) do |handlers| + handlers[:on_connect].call({}) + expect { handlers[:on_event].call(SSE::StreamEvent.new(:put, '{Hi there}')) }.to raise_error(JSON::ParserError) + end + end + + it 'does not report a status after OFF from a handler that is already running' do + use_fast_retry_state + with_handlers(processor) do |handlers| + # The handler passed its stopped check just before stop reported OFF. + config.data_source_update_sink.update_status(Interfaces::DataSource::Status::OFF, nil) + handlers[:on_error].call(http_error(503)) + handlers[:on_error].call(SSE::Errors::StreamClosedError.new) + + expect(states).to eq([Interfaces::DataSource::Status::OFF]) + end + end + + it 'does not report a status or wait after it is stopped' do + with_handlers(processor) do |handlers| + processor.stop + started_at = Time.now + handlers[:on_error].call(http_error(401)) + handlers[:on_error].call(SSE::Errors::StreamClosedError.new) + + expect(Time.now - started_at).to be < 1 + expect(states).to eq([Interfaces::DataSource::Status::OFF]) + end + end + end + describe '#log_connection_result' do it "logs successful connection when diagnostic_accumulator is provided" do diagnostic_accumulator = double("DiagnosticAccumulator") diff --git a/spec/impl/data_source_spec.rb b/spec/impl/data_source_spec.rb index cef347b5..05ff23ec 100644 --- a/spec/impl/data_source_spec.rb +++ b/spec/impl/data_source_spec.rb @@ -10,6 +10,33 @@ module Impl let(:flag_change_broadcaster) { LaunchDarkly::Impl::Broadcaster.new(executor, $null_log) } let(:sink) { subject.new(store, status_broadcaster, flag_change_broadcaster) } + it "ignores every status after OFF" do + listener = ListenerSpy.new + status_broadcaster.add_listener(listener) + error = LaunchDarkly::Interfaces::DataSource::ErrorInfo.new( + LaunchDarkly::Interfaces::DataSource::ErrorInfo::NETWORK_ERROR, 0, "late", Time.now) + + sink.update_status(LaunchDarkly::Interfaces::DataSource::Status::OFF, nil) + sink.update_status(LaunchDarkly::Interfaces::DataSource::Status::VALID, nil) + sink.update_status(LaunchDarkly::Interfaces::DataSource::Status::INTERRUPTED, error) + sink.update_status(LaunchDarkly::Interfaces::DataSource::Status::OFF, error) + + expect(listener.statuses.map(&:state)).to eq([LaunchDarkly::Interfaces::DataSource::Status::OFF]) + expect(sink.current_status.state).to eq(LaunchDarkly::Interfaces::DataSource::Status::OFF) + expect(sink.current_status.last_error).to be_nil + end + + it "ignores a store error after OFF" do + listener = ListenerSpy.new + status_broadcaster.add_listener(listener) + allow(store).to receive(:init).and_raise(StandardError.new("store failure")) + + sink.update_status(LaunchDarkly::Interfaces::DataSource::Status::OFF, nil) + expect { sink.init({}) }.to raise_error(StandardError) + + expect(listener.statuses.map(&:state)).to eq([LaunchDarkly::Interfaces::DataSource::Status::OFF]) + end + it "defaults to initializing" do expect(sink.current_status.state).to eq(LaunchDarkly::Interfaces::DataSource::Status::INITIALIZING) expect(sink.current_status.last_error).to be_nil diff --git a/spec/impl/data_system/streaming_synchronizer_spec.rb b/spec/impl/data_system/streaming_synchronizer_spec.rb index b1ff865a..c2dfcb15 100644 --- a/spec/impl/data_system/streaming_synchronizer_spec.rb +++ b/spec/impl/data_system/streaming_synchronizer_spec.rb @@ -375,6 +375,20 @@ def initialize(type, data = nil) end end + describe '#handle_error' do + let(:synchronizer) { LaunchDarkly::DataSystem::StreamingDataSourceBuilder.new.build(sdk_key, config) } + + it "reports a server-initiated stream close as a network error" do + expect(logger).to receive(:warn) + + update = synchronizer.send(:handle_error, SSE::Errors::StreamClosedError.new, "env-abc", false) + + expect(update.state).to eq(LaunchDarkly::Interfaces::DataSource::Status::INTERRUPTED) + expect(update.error.kind).to eq(LaunchDarkly::Interfaces::DataSource::ErrorInfo::NETWORK_ERROR) + expect(update.environment_id).to eq("env-abc") + end + end + describe 'diagnostic event recording' do let(:synchronizer) { LaunchDarkly::DataSystem::StreamingDataSourceBuilder.new.build(sdk_key, config) } diff --git a/spec/impl/repeating_task_spec.rb b/spec/impl/repeating_task_spec.rb index a754103b..d46fb9d7 100644 --- a/spec/impl/repeating_task_spec.rb +++ b/spec/impl/repeating_task_spec.rb @@ -59,6 +59,63 @@ def null_logger expect(no_more_items).to be true end + it "reads a callable interval after each run" do + runs = Queue.new + run_count = 0 + intervals = [] + interval = -> { + intervals << run_count + 0.01 + } + task = RepeatingTask.new(interval, 0, -> { run_count += 1; runs << run_count }, null_logger, "test") + begin + task.start + 3.times { runs.pop } + ensure + task.stop + end + expect(intervals.take(2)).to eq([1, 2]) + end + + [["a numeric", 0.1], ["a callable", -> { 0.1 }]].each do |desc, interval| + it "starts #{desc} interval when the run returns" do + ends = Queue.new + starts = Queue.new + task = RepeatingTask.new(interval, 0, + -> { + starts << Time.now + sleep(0.1) + ends << Time.now + }, + null_logger, "test") + begin + task.start + starts.pop + first_end = ends.pop + second_start = starts.pop + expect(second_start - first_end).to be >= 0.09 + ensure + task.stop + end + end + end + + it "stops promptly when stopped during a long callable interval" do + ran = Concurrent::Event.new + task = RepeatingTask.new(-> { 1e20 }, 0, -> { ran.set }, null_logger, "test") + begin + task.start + expect(ran.wait(1)).to be true + sleep(0.05) + expect(task.instance_variable_get(:@worker).alive?).to be true + started_at = Time.now + task.stop + expect(Time.now - started_at).to be < 1 + ensure + task.stop + end + end + it "stops promptly when stopped during a long start delay" do ran = Concurrent::Event.new task = RepeatingTask.new(10, 10, -> { ran.set }, null_logger, "test") diff --git a/spec/impl/retry_state_spec.rb b/spec/impl/retry_state_spec.rb new file mode 100644 index 00000000..ad2f6ba0 --- /dev/null +++ b/spec/impl/retry_state_spec.rb @@ -0,0 +1,229 @@ +require "ldclient-rb/impl/retry_state" + +require "spec_helper" + +module LaunchDarkly + module Impl + describe RetryState do + # A random source with no jitter, so a delay is the computed value exactly. + let(:no_jitter) { double("random", rand: 0.0) } + let(:logger) { double.as_null_object } + + before { @now = 0.0 } + let(:clock) { -> { @now } } + + def streaming(delay = 1, random: no_jitter) + RetryState.for_streaming(delay, logger, clock: clock, random: random) + end + + def polling(interval = 30, random: no_jitter) + RetryState.for_polling(interval, logger, random: random) + end + + def delays_after(state, *kinds) + kinds.map do |kind| + state.record_failure(kind) + state.next_delay + end + end + + describe ".classify_http_status" do + [400, 408, 429, 500, 502, 503, 504, 599].each do |status| + it "classifies #{status} as normal" do + expect(RetryState.classify_http_status(status)).to eq(:normal) + end + end + + [401, 403, 404, 405, 409, 410, 499].each do |status| + it "classifies #{status} as unexpected" do + expect(RetryState.classify_http_status(status)).to eq(:unexpected) + end + end + end + + describe "streaming" do + it "waits nothing before any failure" do + expect(streaming.next_delay).to eq(0) + end + + it "doubles the normal delay up to 30 seconds" do + expect(delays_after(streaming, *[:normal] * 7)).to eq([1, 2, 4, 8, 16, 30, 30]) + end + + it "moves to the extended regime after an unexpected failure" do + expect(delays_after(streaming, :unexpected, :unexpected, :unexpected, :unexpected, :unexpected, :unexpected)) + .to eq([300, 600, 1200, 2400, 3600, 3600]) + end + + it "starts the extended sequence over on the move, and keeps counting after it" do + state = streaming + expect(delays_after(state, :normal, :normal, :normal)).to eq([1, 2, 4]) + expect(delays_after(state, :unexpected, :unexpected)).to eq([300, 600]) + end + + it "keeps the extended bounds for a normal failure that follows" do + state = streaming + state.record_failure(:unexpected) + expect(delays_after(state, :normal, :normal, :normal)).to eq([600, 1200, 2400]) + end + + it "uses a configured delay above the ceilings as the bounds" do + expect(delays_after(streaming(600), :normal, :normal)).to eq([600, 600]) + expect(delays_after(streaming(7200), :unexpected, :unexpected)).to eq([7200, 7200]) + end + + it "does not reset before 60 seconds of health" do + state = streaming + state.record_failure(:unexpected) + state.record_success + @now += 59 + state.record_failure(:normal) + expect(state.next_delay).to eq(600) + end + + it "resets after 60 seconds of health, noticed at the next failure" do + state = streaming + state.record_failure(:unexpected) + state.record_success + @now += 60 + state.record_failure(:normal) + expect(state.next_delay).to eq(1) + end + + it "measures health from the first success, not the last" do + state = streaming + state.record_failure(:normal) + state.record_failure(:normal) + state.record_success + @now += 30 + state.record_success + @now += 30 + state.record_failure(:normal) + expect(state.next_delay).to eq(1) + end + + it "starts a new healthy stretch after a failure" do + state = streaming + state.record_failure(:normal) + state.record_success + @now += 50 + state.record_failure(:normal) + state.record_success + @now += 50 + state.record_failure(:normal) + expect(state.next_delay).to eq(4) + end + + it "waits nothing after a success" do + state = streaming + state.record_failure(:unexpected) + state.record_success + expect(state.next_delay).to eq(0) + end + + it "survives a very long outage" do + state = streaming(1.0) + 2000.times { state.record_failure(:normal) } + expect(state.next_delay).to eq(30) + 2000.times { state.record_failure(:unexpected) } + expect(state.next_delay).to eq(3600) + end + end + + describe "polling" do + it "waits the poll interval before any outcome" do + expect(polling.next_delay).to eq(30) + end + + it "waits the poll interval after a normal failure" do + expect(delays_after(polling, :normal, :normal, :normal)).to eq([30, 30, 30]) + end + + it "backs off in the extended regime after an unexpected failure" do + expect(delays_after(polling, :unexpected, :unexpected, :unexpected, :unexpected, :unexpected, :unexpected)) + .to eq([300, 600, 1200, 2400, 3600, 3600]) + end + + it "restores the poll interval after one success" do + state = polling + state.record_failure(:unexpected) + state.record_success + expect(state.next_delay).to eq(30) + end + + it "keeps the extended bounds after one success" do + state = polling + state.record_failure(:unexpected) + state.record_success + state.record_failure(:normal) + expect(state.next_delay).to eq(600) + end + + it "resets after two successes in a row" do + state = polling + state.record_failure(:unexpected) + state.record_success + state.record_success + state.record_failure(:normal) + expect(state.next_delay).to eq(30) + end + + it "does not count successes across a failure" do + state = polling + state.record_failure(:unexpected) + state.record_success + state.record_failure(:normal) + state.record_success + state.record_failure(:normal) + expect(state.next_delay).to eq(1200) + end + + it "never waits less than the poll interval, even with the most jitter" do + state = polling(3000, random: double("random", rand: 0.999)) + expect(delays_after(state, :normal, :unexpected, :unexpected)).to all(eq(3000)) + end + + it "uses a poll interval above the extended ceiling as every bound" do + expect(delays_after(polling(7200), :unexpected, :unexpected, :normal)).to eq([7200, 7200, 7200]) + end + end + + describe "jitter" do + it "subtracts up to half of the delay" do + state = streaming(random: double("random", rand: 0.5)) + expect(delays_after(state, :normal, :normal, :unexpected)).to eq([0.75, 1.5, 225.0]) + end + + it "stays within bounds for a real random source" do + state = streaming(random: Random.new(1234)) + 10.times { state.record_failure(:normal) } + 100.times do + state.record_failure(:normal) + expect(state.next_delay).to be_between(15, 30).inclusive + end + end + end + + describe "invalid input" do + [0, -1, Float::NAN, Float::INFINITY, -Float::INFINITY, nil, "5"].each do |value| + it "uses the default reconnect delay for #{value.inspect} and warns" do + expect(logger).to receive(:warn).once + expect(delays_after(streaming(value), :normal)).to eq([Config.default_initial_reconnect_delay]) + end + + it "uses the default poll interval for #{value.inspect} and warns" do + expect(logger).to receive(:warn).once + state = polling(value) + expect(state.next_delay).to eq(Config.default_poll_interval) + end + end + + it "does not warn for a valid value" do + expect(logger).not_to receive(:warn) + streaming(0.5) + polling(60) + end + end + end + end +end diff --git a/spec/ldclient_end_to_end_spec.rb b/spec/ldclient_end_to_end_spec.rb index 449ba491..b38b3e9b 100644 --- a/spec/ldclient_end_to_end_spec.rb +++ b/spec/ldclient_end_to_end_spec.rb @@ -26,12 +26,16 @@ module LaunchDarkly end end - it "fails in polling mode with 401 error" do + it "waits the full start time and stays uninitialized in polling mode with 401 error" do with_server do |poll_server| poll_server.setup_status_response("/sdk/latest-all", 401) - with_client(test_config(stream: false, data_source: nil, base_uri: poll_server.base_uri.to_s)) do |client| + config = test_config(stream: false, data_source: nil, base_uri: poll_server.base_uri.to_s) + started_at = Time.now + ensure_close(LDClient.new(sdk_key, config, 1)) do |client| + expect(Time.now - started_at).to be >= 1 expect(client.initialized?).to be false + expect(client.data_source_status_provider.status.state).not_to eq(Interfaces::DataSource::Status::OFF) expect(client.variation(ALWAYS_TRUE_FLAG[:key], basic_context, false)).to be false end end From 96d448aa34d278025cb61171d309e083ad80b38b Mon Sep 17 00:00:00 2001 From: jsonbailey Date: Fri, 25 Sep 2026 15:43:08 -0500 Subject: [PATCH 2/3] fix: Use the renamed StreamClosedByServerError --- lib/ldclient-rb/impl/data_source/stream.rb | 2 +- lib/ldclient-rb/impl/data_system/streaming.rb | 2 +- spec/impl/data_source/stream_spec.rb | 6 +++--- spec/impl/data_system/streaming_synchronizer_spec.rb | 2 +- 4 files changed, 6 insertions(+), 6 deletions(-) diff --git a/lib/ldclient-rb/impl/data_source/stream.rb b/lib/ldclient-rb/impl/data_source/stream.rb index 5ed9dea9..ab386fab 100644 --- a/lib/ldclient-rb/impl/data_source/stream.rb +++ b/lib/ldclient-rb/impl/data_source/stream.rb @@ -97,7 +97,7 @@ def handle_error(err) return end - if err.is_a?(SSE::Errors::StreamClosedError) + 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| diff --git a/lib/ldclient-rb/impl/data_system/streaming.rb b/lib/ldclient-rb/impl/data_system/streaming.rb index c8642245..395ad984 100644 --- a/lib/ldclient-rb/impl/data_system/streaming.rb +++ b/lib/ldclient-rb/impl/data_system/streaming.rb @@ -333,7 +333,7 @@ def stop @logger.warn { "[LDClient] #{error_info.message}" } when SSE::Errors::HTTPContentTypeError, SSE::Errors::HTTPProxyError, SSE::Errors::ReadTimeoutError, - SSE::Errors::StreamClosedError + SSE::Errors::StreamClosedByServerError @logger.warn { "[LDClient] Network error on stream connection: #{error}, will retry" } update = LaunchDarkly::Interfaces::DataSystem::Update.new( diff --git a/spec/impl/data_source/stream_spec.rb b/spec/impl/data_source/stream_spec.rb index 572aba5a..98e149b5 100644 --- a/spec/impl/data_source/stream_spec.rb +++ b/spec/impl/data_source/stream_spec.rb @@ -203,7 +203,7 @@ def states expect(fast_retry_state).to receive(:record_failure).with(:normal).and_call_original with_handlers(processor) do |handlers| - handlers[:on_error].call(SSE::Errors::StreamClosedError.new) + handlers[:on_error].call(SSE::Errors::StreamClosedByServerError.new) expect(states).to eq([Interfaces::DataSource::Status::INTERRUPTED]) expect(listener.statuses.last.last_error.kind).to eq(Interfaces::DataSource::ErrorInfo::NETWORK_ERROR) end @@ -263,7 +263,7 @@ def states # The handler passed its stopped check just before stop reported OFF. config.data_source_update_sink.update_status(Interfaces::DataSource::Status::OFF, nil) handlers[:on_error].call(http_error(503)) - handlers[:on_error].call(SSE::Errors::StreamClosedError.new) + handlers[:on_error].call(SSE::Errors::StreamClosedByServerError.new) expect(states).to eq([Interfaces::DataSource::Status::OFF]) end @@ -274,7 +274,7 @@ def states processor.stop started_at = Time.now handlers[:on_error].call(http_error(401)) - handlers[:on_error].call(SSE::Errors::StreamClosedError.new) + handlers[:on_error].call(SSE::Errors::StreamClosedByServerError.new) expect(Time.now - started_at).to be < 1 expect(states).to eq([Interfaces::DataSource::Status::OFF]) diff --git a/spec/impl/data_system/streaming_synchronizer_spec.rb b/spec/impl/data_system/streaming_synchronizer_spec.rb index c2dfcb15..331b9bdd 100644 --- a/spec/impl/data_system/streaming_synchronizer_spec.rb +++ b/spec/impl/data_system/streaming_synchronizer_spec.rb @@ -381,7 +381,7 @@ def initialize(type, data = nil) it "reports a server-initiated stream close as a network error" do expect(logger).to receive(:warn) - update = synchronizer.send(:handle_error, SSE::Errors::StreamClosedError.new, "env-abc", false) + update = synchronizer.send(:handle_error, SSE::Errors::StreamClosedByServerError.new, "env-abc", false) expect(update.state).to eq(LaunchDarkly::Interfaces::DataSource::Status::INTERRUPTED) expect(update.error.kind).to eq(LaunchDarkly::Interfaces::DataSource::ErrorInfo::NETWORK_ERROR) From e2c223719e4ea1e896dc0e019cb4a5253bba4048 Mon Sep 17 00:00:00 2001 From: jsonbailey Date: Mon, 28 Sep 2026 09:55:18 -0500 Subject: [PATCH 3/3] fix: Update ld-eventsource to 3.0.0 --- launchdarkly-server-sdk.gemspec | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/launchdarkly-server-sdk.gemspec b/launchdarkly-server-sdk.gemspec index 2a326d0a..fde3bd76 100644 --- a/launchdarkly-server-sdk.gemspec +++ b/launchdarkly-server-sdk.gemspec @@ -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"