Class: Agent::Lock::Store::RedisStore

Inherits:
Object
  • Object
show all
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

Class Method Summary collapse

Instance Method Summary collapse

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.

Returns:



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

Parameters:

Returns:

  • (Boolean)


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.

Parameters:



136
# File 'lib/agent/lock/store/redis_store.rb', line 136

def delete(record) = client.del(key_for(record.id))

#describe ⇒ String

Returns:

  • (String)


64
# File 'lib/agent/lock/store/redis_store.rb', line 64

def describe = "redis #{url} (#{namespace})"

#find(scope) ⇒ Record?

Parameters:

Returns:



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

Parameters:

  • id (String)

Returns:

  • (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.

Returns:

  • (String)


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.

Yields:

  • the critical section

Returns:

  • (Object) —

    whatever the block returns

Raises:

  • (Error) —

    when the mutex stayed held for longer than the timeout, or the block outlived its lease and ran unprotected for a while



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.

Parameters:



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)