Module: Resque::JobChain::Persistence

Defined in:
lib/resque/job_chain/persistence.rb

Constant Summary collapse

START_SCRIPT =
File.read(File.expand_path('../../../lua/start_chain.lua', __dir__))
ADVANCE_SCRIPT =
File.read(File.expand_path('../../../lua/advance_chain.lua', __dir__))
FAIL_SCRIPT =
File.read(File.expand_path('../../../lua/fail_step.lua', __dir__))

Class Method Summary collapse

Class Method Details

.active_chainsObject



82
83
84
# File 'lib/resque/job_chain/persistence.rb', line 82

def active_chains
  raw_redis.smembers(active_key)
end

.active_keyObject



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, error_message, strategy, max_attempts)
  raw_redis.eval(
    FAIL_SCRIPT,
    keys: [state_key(chain_id), active_key],
    argv: [
      step_index.to_s,
      error_message.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

.namespaceObject



102
103
104
# File 'lib/resque/job_chain/persistence.rb', line 102

def namespace
  Resque.redis.namespace
end

.raw_redisObject



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