Class: Agent::Lock::Store::RedisStore
- Inherits:
-
Object
- Object
- Agent::Lock::Store::RedisStore
- Defined in:
- lib/agent/lock/store/redis_store.rb
Overview
The same locks in a local Redis, for the machines that already run one.
Redis buys two things the filesystem cannot: SET NX is atomic across
machines rather than only across processes on one, and a TTL expires an
abandoned lock without anybody having to reason about whether its holder
is still alive. It costs the thing that makes the file store pleasant,
which is that you can cat a lock, so the value stored is the same
markdown document either way.
Picked with AGENT_LOCK_BACKEND=redis, or by default when one answers on REDIS_URL. See Store's moduledoc for how that default is decided.
Constant Summary collapse
- NAMESPACE =
"agent-lock"- MUTEX_LEASE_MS =
How long the store mutex outlives a holder that died holding it. A claim's critical section is a SCAN and one SET, so ten seconds is somebody gone, not somebody busy.
10_000- MUTEX_TIMEOUT =
Seconds to wait for the mutex. Longer than the lease, so that a crashed holder is waited out rather than reported.
15- LEASE_LOST =
"the lock store's mutex lease (#{MUTEX_LEASE_MS}ms) ran out before the claim finished, " \ "so another process may have claimed an overlapping scope: check the listing".freeze
- RELEASE =
<<~LUA if redis.call("GET", KEYS[1]) == ARGV[1] then return redis.call("DEL", KEYS[1]) end return 0 LUA
Instance Attribute Summary collapse
-
#client ⇒ Object
readonly
Returns the value of attribute client.
-
#tree ⇒ Object
readonly
Returns the value of attribute tree.
Class Method Summary collapse
-
.create_client ⇒ Object
or the error that prevented it.
- .url ⇒ Object
Instance Method Summary collapse
-
#all ⇒ Array<Record>
SCAN rather than KEYS: this runs on whatever Redis the machine already has, which may be somebody's shared development instance, and KEYS blocks the server for the length of the scan.
-
#create(record) ⇒ Boolean
Named to match the Store interface FileSystemStore shares: an action with a boolean outcome, not a pure predicate.
- #delete(record) ⇒ void
- #describe ⇒ String
- #find(scope) ⇒ Record?
-
#initialize(tree, client: nil) ⇒ RedisStore
constructor
A new instance of RedisStore.
- #key_for(id) ⇒ String
-
#mutex_key ⇒ String
Outside the namespace on purpose.
-
#synchronize { ... } ⇒ Object
Runs the block with every other process in this store shut out, so a scan for conflicts and the write it justifies cannot be interleaved.
- #update(record) ⇒ void
Constructor Details
#initialize(tree, client: nil) ⇒ RedisStore
Returns a new instance of RedisStore.
48 49 50 51 52 53 54 55 56 57 58 59 60 61 |
# File 'lib/agent/lock/store/redis_store.rb', line 48 def initialize(tree, client: nil) @tree = tree if client.nil? client, error = self.class.create_client raise(error) if client.nil? && error @client = client if client else @client = client end raise "no functional Redis client could be created for #{self.class.url}" unless @client end |
Instance Attribute Details
#client ⇒ Object (readonly)
Returns the value of attribute client.
46 47 48 |
# File 'lib/agent/lock/store/redis_store.rb', line 46 def client @client end |
#tree ⇒ Object (readonly)
Returns the value of attribute tree.
46 47 48 |
# File 'lib/agent/lock/store/redis_store.rb', line 46 def tree @tree end |
Class Method Details
.create_client ⇒ Object
or the error that prevented it
213 214 215 216 217 218 219 220 |
# File 'lib/agent/lock/store/redis_store.rb', line 213 def create_client @client ||= ::Redis.new(url: url).tap do |client| _version = client.info["redis_version"] end [@client, nil] rescue Redis::CannotConnectError, Redis::BaseError, Errno::ECONNREFUSED, SocketError => e [nil, e] end |
.url ⇒ Object
209 |
# File 'lib/agent/lock/store/redis_store.rb', line 209 def url = ENV.fetch("REDIS_URL", "redis://127.0.0.1:6379/0") |
Instance Method Details
#all ⇒ Array<Record>
SCAN rather than KEYS: this runs on whatever Redis the machine already has, which may be somebody's shared development instance, and KEYS blocks the server for the length of the scan.
71 72 73 74 |
# File 'lib/agent/lock/store/redis_store.rb', line 71 def all keys = client.scan_each(match: "#{namespace}:*").to_a.uniq.sort keys.filter_map { |key| parse(client.get(key), key) } end |
#create(record) ⇒ Boolean
Named to match the Store interface FileSystemStore shares: an action with a boolean outcome, not a pure predicate.
rubocop:disable-next Naming/PredicateMethod
124 125 126 127 128 |
# File 'lib/agent/lock/store/redis_store.rb', line 124 def create(record) args = { nx: true } args[:ex] = ttl_seconds if ttl_seconds.positive? !!client.set(key_for(record.id), record.to_markdown, **args) end |
#delete(record) ⇒ void
This method returns an undefined value.
136 |
# File 'lib/agent/lock/store/redis_store.rb', line 136 def delete(record) = client.del(key_for(record.id)) |
#describe ⇒ String
64 |
# File 'lib/agent/lock/store/redis_store.rb', line 64 def describe = "redis #{url} (#{namespace})" |
#find(scope) ⇒ Record?
116 |
# File 'lib/agent/lock/store/redis_store.rb', line 116 def find(scope) = parse(client.get(key_for(Record.id_for(tree, scope))), nil) |
#key_for(id) ⇒ String
140 |
# File 'lib/agent/lock/store/redis_store.rb', line 140 def key_for(id) = "#{namespace}:#{id}" |
#mutex_key ⇒ String
Outside the namespace on purpose. all scans <namespace>:*, and a
mutex inside it would be read back as a lock on every scan.
112 |
# File 'lib/agent/lock/store/redis_store.rb', line 112 def mutex_key = "#{NAMESPACE}-mutex:#{digest}" |
#synchronize { ... } ⇒ Object
Runs the block with every other process in this store shut out, so a
scan for conflicts and the write it justifies cannot be interleaved.
SET NX alone settles a race for one scope, but lib/** and
lib/a1.rb are two keys, and two agents that both scanned an empty
namespace before either wrote both won.
The mutex is a key set with NX to a token only this call knows, and leased rather than held: a holder that dies mid-claim cannot release it, so the lease does, within MUTEX_LEASE_MS. The wait for it is bounded, and longer than the lease, so a crashed holder costs the next agent a pause and never an error.
Not re-entrant. A nested call waits on its own mutex until the lease frees it, and the outer call then finds it gone and raises.
95 96 97 98 99 100 101 102 103 104 105 106 |
# File 'lib/agent/lock/store/redis_store.rb', line 95 def synchronize token = SecureRandom.hex(16) wait_for(token) begin outcome = yield ensure released = release(token) end raise Error, LEASE_LOST unless released outcome end |
#update(record) ⇒ void
This method returns an undefined value.
132 |
# File 'lib/agent/lock/store/redis_store.rb', line 132 def update(record) = client.set(key_for(record.id), record.to_markdown, keepttl: true) |