mirror of
https://github.com/dkam/suo.git
synced 2025-01-29 07:42:43 +00:00
refactor class methods into instance methods
This commit is contained in:
@@ -1,206 +1,188 @@
|
||||
module Suo
|
||||
module Client
|
||||
class Base
|
||||
|
||||
DEFAULT_OPTIONS = {
|
||||
retry_timeout: 0.1,
|
||||
retry_delay: 0.01,
|
||||
acquisition_timeout: 0.1,
|
||||
acquisition_delay: 0.01,
|
||||
stale_lock_expiration: 3600
|
||||
}.freeze
|
||||
|
||||
attr_accessor :client
|
||||
|
||||
include MonitorMixin
|
||||
|
||||
def initialize(options = {})
|
||||
@options = self.class.merge_defaults(options)
|
||||
fail "Client required" unless options[:client]
|
||||
@options = DEFAULT_OPTIONS.merge(options)
|
||||
@retry_count = (@options[:acquisition_timeout] / @options[:acquisition_delay].to_f).ceil
|
||||
@client = @options[:client]
|
||||
super()
|
||||
end
|
||||
|
||||
def lock(key, resources = 1, options = {})
|
||||
options = self.class.merge_defaults(@options.merge(options))
|
||||
token = self.class.lock(key, resources, options)
|
||||
def lock(key, resources = 1)
|
||||
token = acquire_lock(key, resources)
|
||||
|
||||
if token
|
||||
if block_given? && token
|
||||
begin
|
||||
yield if block_given?
|
||||
yield
|
||||
ensure
|
||||
self.class.unlock(key, token, options)
|
||||
unlock(key, token)
|
||||
end
|
||||
|
||||
true
|
||||
else
|
||||
false
|
||||
token
|
||||
end
|
||||
end
|
||||
|
||||
def locked?(key, resources = 1)
|
||||
self.class.locked?(key, resources, @options)
|
||||
locks(key).size >= resources
|
||||
end
|
||||
|
||||
class << self
|
||||
def lock(key, resources = 1, options = {})
|
||||
options = merge_defaults(options)
|
||||
acquisition_token = nil
|
||||
token = SecureRandom.base64(16)
|
||||
def locks(key)
|
||||
val, _ = get(key)
|
||||
locks = deserialize_locks(val)
|
||||
|
||||
retry_with_timeout(key, options) do
|
||||
val, cas = get(key, options)
|
||||
locks
|
||||
end
|
||||
|
||||
if val.nil?
|
||||
set_initial(key, options)
|
||||
next
|
||||
end
|
||||
def refresh(key, acquisition_token)
|
||||
retry_with_timeout(key) do
|
||||
val, cas = get(key)
|
||||
|
||||
locks = deserialize_and_clear_locks(val, options)
|
||||
if val.nil?
|
||||
set_initial(key)
|
||||
next
|
||||
end
|
||||
|
||||
if locks.size < resources
|
||||
add_lock(locks, token)
|
||||
locks = deserialize_and_clear_locks(val)
|
||||
|
||||
newval = serialize_locks(locks)
|
||||
refresh_lock(locks, acquisition_token)
|
||||
|
||||
if set(key, newval, cas, options)
|
||||
acquisition_token = token
|
||||
break
|
||||
end
|
||||
break if set(key, serialize_locks(locks), cas)
|
||||
end
|
||||
end
|
||||
|
||||
def unlock(key, acquisition_token)
|
||||
return unless acquisition_token
|
||||
|
||||
retry_with_timeout(key) do
|
||||
val, cas = get(key)
|
||||
|
||||
break if val.nil?
|
||||
|
||||
locks = deserialize_and_clear_locks(val)
|
||||
|
||||
acquisition_lock = remove_lock(locks, acquisition_token)
|
||||
|
||||
break unless acquisition_lock
|
||||
break if set(key, serialize_locks(locks), cas)
|
||||
end
|
||||
rescue LockClientError => _ # rubocop:disable Lint/HandleExceptions
|
||||
# ignore - assume success due to optimistic locking
|
||||
end
|
||||
|
||||
def clear(key) # rubocop:disable Lint/UnusedMethodArgument
|
||||
fail NotImplementedError
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def acquire_lock(key, resources = 1)
|
||||
acquisition_token = nil
|
||||
token = SecureRandom.base64(16)
|
||||
|
||||
retry_with_timeout(key) do
|
||||
val, cas = get(key)
|
||||
|
||||
if val.nil?
|
||||
set_initial(key)
|
||||
next
|
||||
end
|
||||
|
||||
locks = deserialize_and_clear_locks(val)
|
||||
|
||||
if locks.size < resources
|
||||
add_lock(locks, token)
|
||||
|
||||
newval = serialize_locks(locks)
|
||||
|
||||
if set(key, newval, cas)
|
||||
acquisition_token = token
|
||||
break
|
||||
end
|
||||
end
|
||||
|
||||
acquisition_token
|
||||
end
|
||||
|
||||
def locked?(key, resources = 1, options = {})
|
||||
locks(key, options).size >= resources
|
||||
end
|
||||
acquisition_token
|
||||
end
|
||||
|
||||
def locks(key, options)
|
||||
options = merge_defaults(options)
|
||||
val, _ = get(key, options)
|
||||
locks = deserialize_locks(val)
|
||||
def get(key) # rubocop:disable Lint/UnusedMethodArgument
|
||||
fail NotImplementedError
|
||||
end
|
||||
|
||||
locks
|
||||
end
|
||||
def set(key, newval, oldval) # rubocop:disable Lint/UnusedMethodArgument
|
||||
fail NotImplementedError
|
||||
end
|
||||
|
||||
def refresh(key, acquisition_token, options = {})
|
||||
options = merge_defaults(options)
|
||||
def set_initial(key) # rubocop:disable Lint/UnusedMethodArgument
|
||||
fail NotImplementedError
|
||||
end
|
||||
|
||||
retry_with_timeout(key, options) do
|
||||
val, cas = get(key, options)
|
||||
def synchronize(key) # rubocop:disable Lint/UnusedMethodArgument
|
||||
mon_synchronize { yield }
|
||||
end
|
||||
|
||||
if val.nil?
|
||||
set_initial(key, options)
|
||||
next
|
||||
end
|
||||
def retry_with_timeout(key)
|
||||
start = Time.now.to_f
|
||||
|
||||
locks = deserialize_and_clear_locks(val, options)
|
||||
@retry_count.times do
|
||||
now = Time.now.to_f
|
||||
break if now - start > @options[:acquisition_timeout]
|
||||
|
||||
refresh_lock(locks, acquisition_token)
|
||||
|
||||
break if set(key, serialize_locks(locks), cas, options)
|
||||
synchronize(key) do
|
||||
yield
|
||||
end
|
||||
|
||||
sleep(rand(@options[:acquisition_delay] * 1000).to_f / 1000)
|
||||
end
|
||||
rescue => _
|
||||
raise LockClientError
|
||||
end
|
||||
|
||||
def unlock(key, acquisition_token, options = {})
|
||||
options = merge_defaults(options)
|
||||
def serialize_locks(locks)
|
||||
MessagePack.pack(locks.map { |time, token| [time.to_f, token] })
|
||||
end
|
||||
|
||||
return unless acquisition_token
|
||||
def deserialize_and_clear_locks(val)
|
||||
clear_expired_locks(deserialize_locks(val))
|
||||
end
|
||||
|
||||
retry_with_timeout(key, options) do
|
||||
val, cas = get(key, options)
|
||||
def deserialize_locks(val)
|
||||
unpacked = (val.nil? || val == "") ? [] : MessagePack.unpack(val)
|
||||
|
||||
break if val.nil?
|
||||
|
||||
locks = deserialize_and_clear_locks(val, options)
|
||||
|
||||
acquisition_lock = remove_lock(locks, acquisition_token)
|
||||
|
||||
break unless acquisition_lock
|
||||
break if set(key, serialize_locks(locks), cas, options)
|
||||
end
|
||||
rescue LockClientError => _ # rubocop:disable Lint/HandleExceptions
|
||||
# ignore - assume success due to optimistic locking
|
||||
unpacked.map do |time, token|
|
||||
[Time.at(time), token]
|
||||
end
|
||||
rescue EOFError => _
|
||||
[]
|
||||
end
|
||||
|
||||
def clear(key, options = {}) # rubocop:disable Lint/UnusedMethodArgument
|
||||
fail NotImplementedError
|
||||
end
|
||||
def clear_expired_locks(locks)
|
||||
expired = Time.now - @options[:stale_lock_expiration]
|
||||
locks.reject { |time, _| time < expired }
|
||||
end
|
||||
|
||||
def merge_defaults(options = {})
|
||||
options = self::DEFAULT_OPTIONS.merge(options)
|
||||
def add_lock(locks, token)
|
||||
locks << [Time.now.to_f, token]
|
||||
end
|
||||
|
||||
fail "Client required" unless options[:client]
|
||||
def remove_lock(locks, acquisition_token)
|
||||
lock = locks.find { |_, token| token == acquisition_token }
|
||||
locks.delete(lock)
|
||||
end
|
||||
|
||||
options
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def get(key, options) # rubocop:disable Lint/UnusedMethodArgument
|
||||
fail NotImplementedError
|
||||
end
|
||||
|
||||
def set(key, newval, oldval, options) # rubocop:disable Lint/UnusedMethodArgument
|
||||
fail NotImplementedError
|
||||
end
|
||||
|
||||
def set_initial(key, options) # rubocop:disable Lint/UnusedMethodArgument
|
||||
fail NotImplementedError
|
||||
end
|
||||
|
||||
def synchronize(key, options)
|
||||
yield(key, options)
|
||||
end
|
||||
|
||||
def retry_with_timeout(key, options)
|
||||
count = (options[:retry_timeout] / options[:retry_delay].to_f).ceil
|
||||
|
||||
start = Time.now.to_f
|
||||
|
||||
count.times do
|
||||
now = Time.now.to_f
|
||||
break if now - start > options[:retry_timeout]
|
||||
|
||||
synchronize(key, options) do
|
||||
yield
|
||||
end
|
||||
|
||||
sleep(rand(options[:retry_delay] * 1000).to_f / 1000)
|
||||
end
|
||||
rescue => _
|
||||
raise LockClientError
|
||||
end
|
||||
|
||||
def serialize_locks(locks)
|
||||
MessagePack.pack(locks.map { |time, token| [time.to_f, token] })
|
||||
end
|
||||
|
||||
def deserialize_and_clear_locks(val, options)
|
||||
clear_expired_locks(deserialize_locks(val), options)
|
||||
end
|
||||
|
||||
def deserialize_locks(val)
|
||||
unpacked = (val.nil? || val == "") ? [] : MessagePack.unpack(val)
|
||||
|
||||
unpacked.map do |time, token|
|
||||
[Time.at(time), token]
|
||||
end
|
||||
rescue EOFError => _
|
||||
[]
|
||||
end
|
||||
|
||||
def clear_expired_locks(locks, options)
|
||||
expired = Time.now - options[:stale_lock_expiration]
|
||||
locks.reject { |time, _| time < expired }
|
||||
end
|
||||
|
||||
def add_lock(locks, token)
|
||||
locks << [Time.now.to_f, token]
|
||||
end
|
||||
|
||||
def remove_lock(locks, acquisition_token)
|
||||
lock = locks.find { |_, token| token == acquisition_token }
|
||||
locks.delete(lock)
|
||||
end
|
||||
|
||||
def refresh_lock(locks, acquisition_token)
|
||||
remove_lock(locks, acquisition_token)
|
||||
add_lock(locks, token)
|
||||
end
|
||||
def refresh_lock(locks, acquisition_token)
|
||||
remove_lock(locks, acquisition_token)
|
||||
add_lock(locks, token)
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
@@ -6,25 +6,22 @@ module Suo
|
||||
super
|
||||
end
|
||||
|
||||
class << self
|
||||
def clear(key, options = {})
|
||||
options = merge_defaults(options)
|
||||
options[:client].delete(key)
|
||||
end
|
||||
def clear(key)
|
||||
@client.delete(key)
|
||||
end
|
||||
|
||||
private
|
||||
private
|
||||
|
||||
def get(key, options)
|
||||
options[:client].get_cas(key)
|
||||
end
|
||||
def get(key)
|
||||
@client.get_cas(key)
|
||||
end
|
||||
|
||||
def set(key, newval, cas, options)
|
||||
options[:client].set_cas(key, newval, cas)
|
||||
end
|
||||
def set(key, newval, cas)
|
||||
@client.set_cas(key, newval, cas)
|
||||
end
|
||||
|
||||
def set_initial(key, options)
|
||||
options[:client].set(key, "")
|
||||
end
|
||||
def set_initial(key)
|
||||
@client.set(key, "")
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
@@ -6,37 +6,34 @@ module Suo
|
||||
super
|
||||
end
|
||||
|
||||
class << self
|
||||
def clear(key, options = {})
|
||||
options = merge_defaults(options)
|
||||
options[:client].del(key)
|
||||
def clear(key)
|
||||
@client.del(key)
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def get(key)
|
||||
[@client.get(key), nil]
|
||||
end
|
||||
|
||||
def set(key, newval, _)
|
||||
ret = @client.multi do |multi|
|
||||
multi.set(key, newval)
|
||||
end
|
||||
|
||||
private
|
||||
ret[0] == "OK"
|
||||
end
|
||||
|
||||
def get(key, options)
|
||||
[options[:client].get(key), nil]
|
||||
def synchronize(key)
|
||||
@client.watch(key) do
|
||||
yield
|
||||
end
|
||||
ensure
|
||||
@client.unwatch
|
||||
end
|
||||
|
||||
def set(key, newval, _, options)
|
||||
ret = options[:client].multi do |multi|
|
||||
multi.set(key, newval)
|
||||
end
|
||||
|
||||
ret[0] == "OK"
|
||||
end
|
||||
|
||||
def synchronize(key, options)
|
||||
options[:client].watch(key) do
|
||||
yield
|
||||
end
|
||||
ensure
|
||||
options[:client].unwatch
|
||||
end
|
||||
|
||||
def set_initial(key, options)
|
||||
options[:client].set(key, "")
|
||||
end
|
||||
def set_initial(key)
|
||||
@client.set(key, "")
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
Reference in New Issue
Block a user