Skip to content

Commit 3b81333

Browse files
committed
Reserve half-open tests
1 parent b3921b3 commit 3b81333

5 files changed

Lines changed: 94 additions & 11 deletions

File tree

‎lib/faulty/circuit.rb‎

Lines changed: 15 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -308,9 +308,11 @@ def run(cache: nil, &block)
308308
return cached_value if !cached_value.nil? && !cache_should_refresh?(cache)
309309

310310
current_status = status
311-
return run_skipped(cached_value) unless current_status.can_run?
312-
313-
run_exec(current_status, cached_value, cache, &block)
311+
if current_status.can_run? && reserve(current_status)
312+
run_exec(current_status, cached_value, cache, &block)
313+
else
314+
run_skipped(cached_value) unless current_status.can_run?
315+
end
314316
end
315317

316318
# Force the circuit to stay open until unlocked
@@ -403,6 +405,16 @@ def run_skipped(cached_value)
403405
cached_value
404406
end
405407

408+
# Reserves execution for this circuit when it is half-open
409+
#
410+
# This prevents concurrent evaluation from allowing multiple simultaneous
411+
# runs for half-open circuits.
412+
def reserve(status)
413+
return true unless status.half_open?
414+
415+
storage.reserve(self, Faulty.current_time, status.reserved_at)
416+
end
417+
406418
# Execute a run
407419
#
408420
# @param cached_value The cached value if one is available

‎lib/faulty/status.rb‎

Lines changed: 18 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -34,10 +34,12 @@ class Faulty
3434
:state,
3535
:lock,
3636
:opened_at,
37+
:reserved_at,
3738
:failure_rate,
3839
:sample_size,
3940
:options,
40-
:stub
41+
:stub,
42+
:current_time
4143
)
4244

4345
class Status
@@ -66,7 +68,8 @@ class Status
6668
# sample_size
6769
# @return [Status]
6870
def self.from_entries(entries, **hash)
69-
window_start = Faulty.current_time - hash[:options].evaluation_window
71+
current_time = Faulty.current_time
72+
window_start = current_time - hash[:options].evaluation_window
7073
size = entries.size
7174
i = 0
7275
failures = 0
@@ -84,7 +87,8 @@ def self.from_entries(entries, **hash)
8487

8588
new(hash.merge(
8689
sample_size: sample_size,
87-
failure_rate: sample_size.zero? ? 0.0 : failures.to_f / sample_size
90+
failure_rate: sample_size.zero? ? 0.0 : failures.to_f / sample_size,
91+
current_time: current_time
8892
))
8993
end
9094

@@ -94,7 +98,7 @@ def self.from_entries(entries, **hash)
9498
#
9599
# @return [Boolean] True if open
96100
def open?
97-
state == :open && opened_at + options.cool_down > Faulty.current_time
101+
state == :open && opened_at + options.cool_down > current_time
98102
end
99103

100104
# Whether the circuit is closed
@@ -112,7 +116,7 @@ def closed?
112116
#
113117
# @return [Boolean] True if half-open
114118
def half_open?
115-
state == :open && opened_at + options.cool_down <= Faulty.current_time
119+
state == :open && opened_at + options.cool_down <= current_time
116120
end
117121

118122
# Whether the circuit is locked open
@@ -129,13 +133,20 @@ def locked_closed?
129133
lock == :closed
130134
end
131135

136+
def reserved?
137+
return false unless reserved_at
138+
139+
state == :open && reserved_at + options.cool_down >= current_time
140+
end
141+
132142
# Whether the circuit can be run
133143
#
134144
# Takes the circuit state, locks and cooldown into account
135145
#
136146
# @return [Boolean] True if the circuit can be run
137147
def can_run?
138148
return false if locked_open?
149+
return false if reserved?
139150

140151
closed? || locked_closed? || half_open?
141152
end
@@ -166,7 +177,8 @@ def defaults
166177
state: :closed,
167178
failure_rate: 0.0,
168179
sample_size: 0,
169-
stub: false
180+
stub: false,
181+
current_time: Faulty.current_time
170182
}
171183
end
172184
end

‎lib/faulty/storage/interface.rb‎

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -90,6 +90,9 @@ def reopen(circuit, opened_at, previous_opened_at)
9090
# may be called more than once. If so, this method should return true
9191
# only once, when the circuit transitions from open to closed.
9292
#
93+
# The backend should reset the reserved_at value to empty when closing
94+
# the circuit.
95+
#
9396
# If the backend does not support locking or atomic operations, then
9497
# it may always return true, but that could result in duplicate close
9598
# notifications.
@@ -99,6 +102,26 @@ def close(circuit)
99102
raise NotImplementedError
100103
end
101104

105+
# Reserve an exclusive run for this circuit
106+
#
107+
# This is used when the circuit is half-open and the test run is being
108+
# attempted. We need to make sure only a single run is allowed.
109+
#
110+
# The backend should store reserved_at and use it to serve future status
111+
# requests. When setting reserved_at, the backend should atomically
112+
# compare any existing value using previous_reserved_at. This ensures
113+
# that mutltiple parallel processes can't reserve the circuit.
114+
#
115+
# The backend should return true if the reservation was successful, and
116+
# false if it was not.
117+
#
118+
# If the backend does not support locking or atomic operations, then
119+
# it may always return true, but will result in duplicate half-open test
120+
# runs.
121+
def reserve(circuit, reserved_at, previous_reserved_at)
122+
raise NotImplementedError
123+
end
124+
102125
# Lock the circuit in a given state
103126
#
104127
# No concurrency gurantees are provided for locking

‎lib/faulty/storage/memory.rb‎

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -41,11 +41,12 @@ def defaults
4141
# The internal object for storing a circuit
4242
#
4343
# @private
44-
MemoryCircuit = Struct.new(:state, :runs, :opened_at, :lock, :options) do
44+
MemoryCircuit = Struct.new(:state, :runs, :opened_at, :reserved_at, :lock, :options) do
4545
def initialize
4646
self.state = Concurrent::Atom.new(:closed)
4747
self.runs = Concurrent::MVar.new([], dup_on_deref: true)
4848
self.opened_at = Concurrent::Atom.new(nil)
49+
self.reserved_at = Concurrent::Atom.new(nil)
4950
self.lock = nil
5051
end
5152

@@ -61,6 +62,7 @@ def status(circuit_options)
6162
state: state.value,
6263
lock: lock,
6364
opened_at: opened_at.value,
65+
reserved_at: reserved_at.value,
6466
options: circuit_options
6567
)
6668
end
@@ -139,7 +141,19 @@ def reopen(circuit, opened_at, previous_opened_at)
139141
def close(circuit)
140142
memory = fetch(circuit)
141143
memory.runs.modify { |_old| [] }
142-
memory.state.compare_and_set(:open, :closed)
144+
closed = memory.state.compare_and_set(:open, :closed)
145+
memory.reserved_at.reset(nil) if closed
146+
closed
147+
end
148+
149+
# Reserve an exclusive run for this circuit
150+
#
151+
# @see Interface#reserve
152+
# @param (see Interface#reserve)
153+
# @return (see Interface#reserve)
154+
def reserve(circuit, reserved_at, previous_reserved_at)
155+
memory = fetch(circuit)
156+
memory.reserved_at.compare_and_set(previous_reserved_at, reserved_at)
143157
end
144158

145159
# Lock a circuit open or closed

‎lib/faulty/storage/redis.rb‎

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -174,11 +174,25 @@ def close(circuit)
174174
result = watch_exec(key, ['open']) do |m|
175175
m.set(key, 'closed', ex: ex)
176176
m.del(entries_key(circuit.name))
177+
m.del(reserved_at_key(circuit.name))
177178
end
178179

179180
result && result[0] == 'OK'
180181
end
181182

183+
# Reserve an exclusive run for this circuit
184+
#
185+
# @see Interface#reserve
186+
# @param (see Interface#reserve)
187+
# @return (see Interface#reserve)
188+
def reserve(circuit, reserved_at, previous_reserved_at)
189+
key = reserved_at_key(circuit.name)
190+
result = watch_exec(key, [previous_reserved_at.to_s]) do |m|
191+
m.set(key, reserved_at, ex: options.circuit_ttl)
192+
end
193+
result && result[0] == 'OK'
194+
end
195+
182196
# Lock a circuit open or closed
183197
#
184198
# The circuit_ttl does not apply to locks
@@ -228,18 +242,21 @@ def status(circuit)
228242
futures[:state] = r.get(state_key(circuit.name))
229243
futures[:lock] = r.get(lock_key(circuit.name))
230244
futures[:opened_at] = r.get(opened_at_key(circuit.name))
245+
futures[:reserved_at] = r.get(reserved_at_key(circuit.name))
231246
futures[:entries] = r.lrange(entries_key(circuit.name), 0, -1)
232247
end
233248

234249
state = futures[:state].value&.to_sym || :closed
235250
opened_at = futures[:opened_at].value ? Float(futures[:opened_at].value) : nil
236251
opened_at = Faulty.current_time - options.circuit_ttl if state == :open && opened_at.nil?
252+
opened_at = futures[:reserved_at].value ? Float(futures[:reserved_at].value) : nil
237253

238254
Faulty::Status.from_entries(
239255
map_entries(futures[:entries].value),
240256
state: state,
241257
lock: futures[:lock].value&.to_sym,
242258
opened_at: opened_at,
259+
reserved_at: reserved_at,
243260
options: circuit.options
244261
)
245262
end
@@ -321,6 +338,11 @@ def opened_at_key(circuit_name)
321338
ckey(circuit_name, 'opened_at')
322339
end
323340

341+
# @return [String] The key for circuit opened_at
342+
def reserved_at_key(circuit_name)
343+
ckey(circuit_name, 'reserved_at')
344+
end
345+
324346
# Get the current key to add circuit names to
325347
def list_key
326348
key('list', current_list_block)

0 commit comments

Comments
 (0)