Skip to content
Merged
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
37 changes: 35 additions & 2 deletions lib/ldclient-rb/impl/data_source/status_provider.rb
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

require "concurrent"
require "forwardable"
require "ldclient-rb/impl/util"
require "ldclient-rb/interfaces"

module LaunchDarkly
Expand All @@ -19,18 +20,24 @@ class StatusProviderV2
extend Forwardable
def_delegators :@status_broadcaster, :add_listener, :remove_listener

MONOTONIC_CLOCK = -> { Impl::Util.monotonic_seconds }
private_constant :MONOTONIC_CLOCK

#
# Creates a new status provider.
#
# @param status_broadcaster [LaunchDarkly::Impl::Broadcaster] Broadcaster for status changes
# @param clock [#call] returns seconds on a monotonic clock; injectable for tests
#
def initialize(status_broadcaster)
def initialize(status_broadcaster, clock: MONOTONIC_CLOCK)
@status_broadcaster = status_broadcaster
@clock = clock
@status = LaunchDarkly::Interfaces::DataSource::Status.new(
LaunchDarkly::Interfaces::DataSource::Status::INITIALIZING,
Time.now,
nil
)
@state_since_monotonic = @clock.call
@lock = Concurrent::ReadWriteLock.new
end

Expand All @@ -41,8 +48,29 @@ def status
end
end

#
# The current status together with the seconds the data source has spent
# in its current state, read atomically under one lock acquisition so the
# duration can never be paired with a stale state.
#
# The duration is measured on the monotonic clock, so a wall-clock step
# cannot distort it. `status.state_since` remains the wall-clock time
# reported to applications; this duration is the value to use when it
# decides behavior.
#
# @private
# @return [Array(LaunchDarkly::Interfaces::DataSource::Status, Float)]
#
def status_and_seconds_in_state
@lock.with_read_lock do
[@status, @clock.call - @state_since_monotonic]
end
end

# (see LaunchDarkly::Interfaces::DataSource::UpdateSink#update_status)
def update_status(new_state, new_error)
return if new_state.nil?

status_to_broadcast = nil

@lock.with_write_lock do
Expand All @@ -62,7 +90,12 @@ def update_status(new_state, new_error)
# No change if state is the same and no error
return if new_state == old_status.state && new_error.nil?

new_since = new_state == old_status.state ? @status.state_since : Time.now
if new_state == old_status.state
new_since = @status.state_since
else
new_since = Time.now
@state_since_monotonic = @clock.call
end
new_error = @status.last_error if new_error.nil?

@status = LaunchDarkly::Interfaces::DataSource::Status.new(
Expand Down
12 changes: 8 additions & 4 deletions lib/ldclient-rb/impl/data_source/stream.rb
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ def initialize(sdk_key, config, diagnostic_accumulator = nil)
@stopped = Concurrent::AtomicBoolean.new(false)
@ready = Concurrent::Event.new
@connection_attempt_start_time = 0
@connection_attempt_started_monotonic = nil
end

def initialized?
Expand Down Expand Up @@ -184,13 +185,16 @@ def key_for_path(kind, path)

def log_connection_started
@connection_attempt_start_time = Impl::Util::current_time_millis
@connection_attempt_started_monotonic = Impl::Util.monotonic_seconds
end

def log_connection_result(is_success)
if !@diagnostic_accumulator.nil? && @connection_attempt_start_time > 0
@diagnostic_accumulator.record_stream_init(@connection_attempt_start_time, !is_success,
Impl::Util::current_time_millis - @connection_attempt_start_time)
@connection_attempt_start_time = 0
if !@diagnostic_accumulator.nil? && !@connection_attempt_started_monotonic.nil?
# The monotonic stamp is both the "attempt in flight" sentinel and the
# duration base; the wall-clock value is only the reported timestamp.
duration_millis = ((Impl::Util.monotonic_seconds - @connection_attempt_started_monotonic) * 1_000).to_i
@diagnostic_accumulator.record_stream_init(@connection_attempt_start_time, !is_success, duration_millis)
@connection_attempt_started_monotonic = nil
end
end
end
Expand Down
22 changes: 13 additions & 9 deletions lib/ldclient-rb/impl/data_system/fdv2.rb
Original file line number Diff line number Diff line change
Expand Up @@ -459,10 +459,12 @@ def consume_synchronizer_results(synchronizer, check_recovery: false)
break if update == "quit"

if update == "check"
# Check condition periodically
current_status = @data_source_status_provider.status
return SyncResult::RECOVER if check_recovery && recovery_condition(current_status)
return SyncResult::FALLBACK if fallback_condition(current_status)
# Check condition periodically. One atomic read pairs the status with
# its monotonic time-in-state, so a verdict can never pair a stale
# state with a fresh duration.
current_status, seconds_in_state = @data_source_status_provider.status_and_seconds_in_state
return SyncResult::RECOVER if check_recovery && recovery_condition(current_status, seconds_in_state)
return SyncResult::FALLBACK if fallback_condition(current_status, seconds_in_state)
end
next
end
Expand Down Expand Up @@ -514,13 +516,14 @@ def record_environment_id(environment_id)
# Determine if we should fallback to the next synchronizer.
#
# @param status [LaunchDarkly::Interfaces::DataSource::Status] Current data source status
# @param seconds_in_state [Float] monotonic seconds spent in status.state
# @return [Boolean] true if fallback condition is met
#
def fallback_condition(status)
def fallback_condition(status, seconds_in_state)
interrupted_at_runtime = status.state == LaunchDarkly::Interfaces::DataSource::Status::INTERRUPTED &&
Time.now - status.state_since > 60 # 1 minute
seconds_in_state > 60 # 1 minute
cannot_initialize = status.state == LaunchDarkly::Interfaces::DataSource::Status::INITIALIZING &&
Time.now - status.state_since > 10 # 10 seconds
seconds_in_state > 10 # 10 seconds

interrupted_at_runtime || cannot_initialize
end
Expand All @@ -529,11 +532,12 @@ def fallback_condition(status)
# Determine if we should recover to the primary synchronizer.
#
# @param status [LaunchDarkly::Interfaces::DataSource::Status] Current data source status
# @param seconds_in_state [Float] monotonic seconds spent in status.state
# @return [Boolean] true if recovery condition is met (healthy for too long)
#
def recovery_condition(status)
def recovery_condition(status, seconds_in_state)
status.state == LaunchDarkly::Interfaces::DataSource::Status::VALID &&
Time.now - status.state_since > 300 # 5 minutes
seconds_in_state > 300 # 5 minutes
end

#
Expand Down
15 changes: 9 additions & 6 deletions lib/ldclient-rb/impl/data_system/streaming.rb
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ def initialize(sdk_key, http_config, initial_reconnect_delay, config)
@stopped = Concurrent::Event.new
@diagnostic_accumulator = nil
@connection_attempt_start_time = 0
@connection_attempt_started_monotonic = nil
end

#
Expand Down Expand Up @@ -118,7 +119,7 @@ def sync(ss)
next unless update

log_connection_result(true)
@connection_attempt_start_time = 0
@connection_attempt_started_monotonic = nil

yield update

Expand Down Expand Up @@ -366,16 +367,18 @@ def stop

private def log_connection_started
@connection_attempt_start_time = Impl::Util.current_time_millis
@connection_attempt_started_monotonic = Impl::Util.monotonic_seconds
end

private def log_connection_result(is_success)
return unless @diagnostic_accumulator
return unless @connection_attempt_start_time > 0
return if @connection_attempt_started_monotonic.nil?

current_time = Impl::Util.current_time_millis
elapsed = current_time - @connection_attempt_start_time
@diagnostic_accumulator.record_stream_init(@connection_attempt_start_time, !is_success, elapsed >= 0 ? elapsed : 0)
@connection_attempt_start_time = 0
# The monotonic stamp is both the "attempt in flight" sentinel and the
# duration base; the wall-clock value is only the reported timestamp.
elapsed = ((Impl::Util.monotonic_seconds - @connection_attempt_started_monotonic) * 1_000).to_i
@diagnostic_accumulator.record_stream_init(@connection_attempt_start_time, !is_success, elapsed)
@connection_attempt_started_monotonic = nil
end
end

Expand Down
14 changes: 11 additions & 3 deletions lib/ldclient-rb/impl/expiring_cache.rb
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
require "ldclient-rb/impl/util"

module LaunchDarkly
module Impl
Expand All @@ -7,10 +8,17 @@ module Impl
# * made thread-safe
# * removed many unused methods
# * reading a key does not reset its expiration time, only writing
# * expiration is measured on a monotonic clock, so wall-clock steps cannot
# retain entries past their TTL or evict them early
class ExpiringCache
def initialize(max_size, ttl)
MONOTONIC_CLOCK = -> { Impl::Util.monotonic_seconds }
private_constant :MONOTONIC_CLOCK

# @param clock [#call] returns seconds on a monotonic clock; injectable for tests
def initialize(max_size, ttl, clock: MONOTONIC_CLOCK)
@max_size = max_size
@ttl = ttl
@clock = clock
@data_lru = {}
@data_ttl = {}
@lock = Mutex.new
Expand All @@ -31,7 +39,7 @@ def []=(key, val)
@data_ttl.delete(key)

@data_lru[key] = val
@data_ttl[key] = Time.now.to_f
@data_ttl[key] = @clock.call

if @data_lru.size > @max_size
key, _ = @data_lru.first # hashes have a FIFO ordering in Ruby
Expand Down Expand Up @@ -63,7 +71,7 @@ def clear
private

def ttl_evict
ttl_horizon = Time.now.to_f - @ttl
ttl_horizon = @clock.call - @ttl
key, time = @data_ttl.first

until time.nil? || time > ttl_horizon
Expand Down
4 changes: 2 additions & 2 deletions lib/ldclient-rb/impl/migrations/migrator.rb
Original file line number Diff line number Diff line change
Expand Up @@ -270,7 +270,7 @@ def initialize(logger, origin, fn, tracker, measure_latency, measure_errors, pay
# @return [LaunchDarkly::Migrations::OperationResult]
#
def run()
start = Time.now
start = Impl::Util.monotonic_seconds

begin
result = @fn.call(@payload)
Expand All @@ -279,7 +279,7 @@ def run()
result = LaunchDarkly::Result.fail("'#{origin}' operation raised an exception", e)
end

@tracker.latency(@origin, (Time.now - start) * 1_000) if @measure_latency
@tracker.latency(@origin, (Impl::Util.monotonic_seconds - start) * 1_000) if @measure_latency
@tracker.error(@origin) if @measure_errors && !result.success?
@tracker.invoked(@origin)

Expand Down
6 changes: 4 additions & 2 deletions lib/ldclient-rb/impl/repeating_task.rb
Original file line number Diff line number Diff line change
Expand Up @@ -39,13 +39,15 @@ def start
@stop_event.wait(@start_delay) unless @start_delay.nil? || @start_delay == 0

until @stopped.value do
started_at = Time.now
started_at = Impl::Util.monotonic_seconds
begin
@task.call
rescue => e
Impl::Util.log_exception(@logger, "Uncaught exception from repeating task", e)
end
delta = @interval - (Time.now - started_at)
# The run must be measured on the monotonic clock: a wall-clock step would
# make this wait negative (immediate re-run) or stretch it by the step size.
delta = @interval - (Impl::Util.monotonic_seconds - started_at)
@stop_event.wait(delta) if delta > 0
end
end
Expand Down
11 changes: 11 additions & 0 deletions lib/ldclient-rb/impl/util.rb
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,17 @@ def self.current_time_millis
(Time.now.to_f * 1000).to_i
end

#
# Seconds on the monotonic clock. Use this for measuring intervals and
# durations; use the wall clock only for values reported outside the
# process.
#
# @return [Float]
#
def self.monotonic_seconds
Process.clock_gettime(Process::CLOCK_MONOTONIC, :float_second)
end

def self.default_http_headers(sdk_key, config)
ret = { "Authorization" => sdk_key, "User-Agent" => "RubyClient/" + LaunchDarkly::VERSION }

Expand Down
Loading
Loading