Module: Resque::JobChain::Persistence
- Defined in:
- lib/resque/job_chain/persistence.rb
Constant Summary collapse
- START_SCRIPT =
File.read(File.('../../../lua/start_chain.lua', __dir__))
- ADVANCE_SCRIPT =
File.read(File.('../../../lua/advance_chain.lua', __dir__))
- FAIL_SCRIPT =
File.read(File.('../../../lua/fail_step.lua', __dir__))
Class Method Summary collapse
- .active_chains ⇒ Object
- .active_key ⇒ Object
- .advance(chain_id, step_index) ⇒ Object
- .context_key(chain_id) ⇒ Object
- .fail_step(chain_id, step_index, error_message, strategy, max_attempts) ⇒ Object
- .load_context(chain_id) ⇒ Object
- .load_state(chain_id) ⇒ Object
- .load_step(chain_id, step_index) ⇒ Object
- .merge_context(chain_id, updates) ⇒ Object
- .namespace ⇒ Object
- .raw_redis ⇒ Object
- .start_chain(chain_id:, steps:, context:) ⇒ Object
- .state_key(chain_id) ⇒ Object
- .steps_key(chain_id) ⇒ Object
- .stringify_keys(hash) ⇒ Object
Class Method Details
.active_chains ⇒ Object
82 83 84 |
# File 'lib/resque/job_chain/persistence.rb', line 82 def active_chains raw_redis.smembers(active_key) end |
.active_key ⇒ Object
98 99 100 |
# File 'lib/resque/job_chain/persistence.rb', line 98 def active_key "#{namespace}:job_chain:active" end |
.advance(chain_id, step_index) ⇒ Object
28 29 30 31 32 33 34 35 36 37 38 39 |
# File 'lib/resque/job_chain/persistence.rb', line 28 def advance(chain_id, step_index) raw_redis.eval( ADVANCE_SCRIPT, keys: [state_key(chain_id), active_key], argv: [ step_index.to_s, Time.now.utc.iso8601, chain_id, JobChain.configuration.completed_chain_ttl.to_s ] ) end |
.context_key(chain_id) ⇒ Object
94 95 96 |
# File 'lib/resque/job_chain/persistence.rb', line 94 def context_key(chain_id) "#{namespace}:job_chain:#{chain_id}:context" end |
.fail_step(chain_id, step_index, error_message, strategy, max_attempts) ⇒ Object
41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 |
# File 'lib/resque/job_chain/persistence.rb', line 41 def fail_step(chain_id, step_index, , strategy, max_attempts) raw_redis.eval( FAIL_SCRIPT, keys: [state_key(chain_id), active_key], argv: [ step_index.to_s, .to_s, strategy.to_s, max_attempts.to_s, Time.now.utc.iso8601, chain_id, JobChain.configuration.completed_chain_ttl.to_s ] ) end |
.load_context(chain_id) ⇒ Object
68 69 70 71 72 73 |
# File 'lib/resque/job_chain/persistence.rb', line 68 def load_context(chain_id) json = raw_redis.get(context_key(chain_id)) return {} unless json JSON.parse(json) end |
.load_state(chain_id) ⇒ Object
57 58 59 |
# File 'lib/resque/job_chain/persistence.rb', line 57 def load_state(chain_id) raw_redis.hgetall(state_key(chain_id)) end |
.load_step(chain_id, step_index) ⇒ Object
61 62 63 64 65 66 |
# File 'lib/resque/job_chain/persistence.rb', line 61 def load_step(chain_id, step_index) json = raw_redis.lindex(steps_key(chain_id), step_index) return nil unless json Step.from_h(JSON.parse(json)) end |
.merge_context(chain_id, updates) ⇒ Object
75 76 77 78 79 80 |
# File 'lib/resque/job_chain/persistence.rb', line 75 def merge_context(chain_id, updates) current = load_context(chain_id) merged = current.merge(stringify_keys(updates)) raw_redis.set(context_key(chain_id), JSON.generate(merged)) merged end |
.namespace ⇒ Object
102 103 104 |
# File 'lib/resque/job_chain/persistence.rb', line 102 def namespace Resque.redis.namespace end |
.raw_redis ⇒ Object
106 107 108 |
# File 'lib/resque/job_chain/persistence.rb', line 106 def raw_redis Resque.redis.redis end |
.start_chain(chain_id:, steps:, context:) ⇒ Object
12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 |
# File 'lib/resque/job_chain/persistence.rb', line 12 def start_chain(chain_id:, steps:, context:) steps_json = steps.map { |s| JSON.generate(s.to_h) } result = raw_redis.eval( START_SCRIPT, keys: [state_key(chain_id), steps_key(chain_id), context_key(chain_id), active_key], argv: [ steps.length.to_s, JSON.generate(stringify_keys(context)), chain_id, JSON.generate(steps_json), Time.now.utc.iso8601 ] ) result == 1 end |
.state_key(chain_id) ⇒ Object
86 87 88 |
# File 'lib/resque/job_chain/persistence.rb', line 86 def state_key(chain_id) "#{namespace}:job_chain:#{chain_id}:state" end |
.steps_key(chain_id) ⇒ Object
90 91 92 |
# File 'lib/resque/job_chain/persistence.rb', line 90 def steps_key(chain_id) "#{namespace}:job_chain:#{chain_id}:steps" end |
.stringify_keys(hash) ⇒ Object
110 111 112 |
# File 'lib/resque/job_chain/persistence.rb', line 110 def stringify_keys(hash) hash.each_with_object({}) { |(k, v), h| h[k.to_s] = v } end |