Class: Kaal::SchedulerFileLoader::JobApplier

Inherits:
Object
  • Object
show all
Includes:
Kaal::Support::HashTools
Defined in:
lib/kaal/scheduler_file/job_applier.rb,
sig/kaal/scheduler_file/job_applier.rbs

Overview

Applies normalized scheduler jobs and rolls them back on failure.

Instance Method Summary collapse

Methods included from Kaal::Support::HashTools

constantize, deep_dup, deep_merge, duplicable?, stringify_keys, symbolize_keys

Constructor Details

#initialize(configuration:, definition_registry:, registry:, logger:, helper_bundle:) ⇒ JobApplier

Returns a new instance of JobApplier.

Parameters:

  • configuration: (Kaal::rbs_any)
  • definition_registry: (Kaal::rbs_any)
  • registry: (Kaal::rbs_any)
  • logger: (Kaal::rbs_any)
  • helper_bundle: (Kaal::rbs_any)


15
16
17
18
19
20
21
# File 'lib/kaal/scheduler_file/job_applier.rb', line 15

def initialize(configuration:, definition_registry:, registry:, logger:, helper_bundle:)
  @configuration = configuration
  @definition_registry = definition_registry
  @registry = registry
  @logger = logger
  @helper_bundle = helper_bundle
end

Instance Method Details

#active_job_dispatch?(job_class, queue) ⇒ Boolean

Parameters:

  • job_class (Kaal::rbs_any)
  • queue (Kaal::rbs_any)

Returns:

  • (Boolean)


212
213
214
# File 'lib/kaal/scheduler_file/job_applier.rb', line 212

def active_job_dispatch?(job_class, queue)
  Kaal::JobDispatcher.active_job_dispatch?(job_class, queue)
end

#apply(job) ⇒ nil, { key: Kaal::rbs_any, existing_definition: Kaal::rbs_any, existing_registry_entry: Kaal::rbs_any }

Parameters:

  • job (Kaal::rbs_any)

Returns:

  • (nil, { key: Kaal::rbs_any, existing_definition: Kaal::rbs_any, existing_registry_entry: Kaal::rbs_any })


23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
# File 'lib/kaal/scheduler_file/job_applier.rb', line 23

def apply(job)
  key = job.fetch(:key)
  cron = job.fetch(:cron)
  job_class_name = job.fetch(:job_class_name)
  queue = job.fetch(:queue)
  existing_definition = @definition_registry.find_definition(key)
  existing_registry_entry = @registry.find(key)
  return nil if conflict?(key:, existing_definition:)

  job_class = resolved_job_class(job_class_name:, key:, queue:)
  callback = callback_for(
    key: key,
    job_class_name: job_class_name,
    queue: queue,
    args_template: job.fetch(:args),
    kwargs_template: job.fetch(:kwargs)
  )
   = (job, job_class)

  @definition_registry.upsert_definition(
    key: key,
    cron: cron,
    enabled: job.fetch(:enabled),
    source: 'file',
    metadata: 
  )

  begin
    @registry.upsert(key: key, cron: cron, enqueue: callback)
  rescue StandardError
    rollback_job(key:, existing_definition:, existing_registry_entry:)
    raise
  end

  { key: key, existing_definition: existing_definition, existing_registry_entry: existing_registry_entry }
end

#build_callback(job, job_class) ⇒ Kaal::rbs_any

Parameters:

  • job (Kaal::rbs_any)
  • job_class (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
# File 'lib/kaal/scheduler_file/job_applier.rb', line 146

def build_callback(job, job_class)
  key = job.fetch(:key)
  queue = job.fetch(:queue)
  args_template = job.fetch(:args)
  kwargs_template = job.fetch(:kwargs)

  lambda do |fire_time:, idempotency_key:|
    context = {
      fire_time: fire_time,
      idempotency_key: idempotency_key,
      key: key
    }
    resolved_args = @helper_bundle.resolve_placeholders(deep_dup(args_template), context)
    raw_kwargs = @helper_bundle.resolve_placeholders(deep_dup(kwargs_template), context) || {}
    raise SchedulerConfigError, "kwargs for scheduler job '#{key}' must be a mapping, got #{raw_kwargs.class}" unless raw_kwargs.is_a?(Hash)

    validate_keyword_keys(raw_kwargs, key)

    resolved_kwargs = raw_kwargs.transform_keys(&:to_sym)
    dispatch_job(job_class, queue, resolved_args, resolved_kwargs, key)
  end
end

#callback_for(key:, job_class_name:, queue:, args_template:, kwargs_template:) ⇒ Kaal::rbs_any

Parameters:

  • key: (Kaal::rbs_any)
  • job_class_name: (Kaal::rbs_any)
  • queue: (Kaal::rbs_any)
  • args_template: (Kaal::rbs_any)
  • kwargs_template: (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


66
67
68
69
70
71
72
73
74
75
76
77
# File 'lib/kaal/scheduler_file/job_applier.rb', line 66

def callback_for(key:, job_class_name:, queue:, args_template:, kwargs_template:)
  job_class = resolved_job_class(job_class_name:, key:, queue:)
  build_callback(
    {
      key: key,
      queue: queue,
      args: args_template,
      kwargs: kwargs_template
    },
    job_class
  )
end

#conflict?(key:, existing_definition:) ⇒ Boolean

Parameters:

  • key: (Kaal::rbs_any)
  • existing_definition: (Kaal::rbs_any)

Returns:

  • (Boolean)


88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
# File 'lib/kaal/scheduler_file/job_applier.rb', line 88

def conflict?(key:, existing_definition:)
  existing_source = existing_definition&.[](:source)
  return false unless existing_source && existing_source.to_s != 'file'

  policy = @configuration.scheduler_conflict_policy
  case policy
  when :error
    raise SchedulerConfigError, "Scheduler key conflict for '#{key}' with existing source '#{existing_source}'"
  when :code_wins
    @logger&.warn("Skipping scheduler file job '#{key}' because scheduler_conflict_policy is :code_wins")
    true
  when :file_wins
    false
  else
    raise SchedulerConfigError, "Unsupported scheduler_conflict_policy '#{policy}'"
  end
end

#dispatch_job(job_class, queue, args, kwargs, key) ⇒ Kaal::rbs_any

Parameters:

  • job_class (Kaal::rbs_any)
  • queue (Kaal::rbs_any)
  • args (Kaal::rbs_any)
  • kwargs (Kaal::rbs_any)
  • key (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
# File 'lib/kaal/scheduler_file/job_applier.rb', line 188

def dispatch_job(job_class, queue, args, kwargs, key)
  if kwargs.empty?
    Kaal::JobDispatcher.dispatch(job_class:, queue:, args:, key:)
  else
    job_class_name = job_class.name

    if queue && !job_class.respond_to?(:set)
      raise SchedulerConfigError,
            "job_class '#{job_class_name}' must respond to .set to use queue #{queue.inspect} for scheduler job '#{key}'"
    end

    if queue
      job_class.set(queue: queue).perform_later(*args, **kwargs)
    elsif job_class.respond_to?(:perform_later)
      job_class.perform_later(*args, **kwargs)
    elsif job_class.respond_to?(:perform)
      job_class.perform(*args, **kwargs)
    else
      raise SchedulerConfigError,
            "job_class '#{job_class_name}' must respond to .perform, .perform_later, or .set(...).perform_later for scheduler job '#{key}'"
    end
  end
end

#persisted_metadata(job, job_class) ⇒ Kaal::rbs_any

Parameters:

  • job (Kaal::rbs_any)
  • job_class (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
# File 'lib/kaal/scheduler_file/job_applier.rb', line 130

def (job, job_class)
  , job_class_name, queue, args, kwargs =
    job.values_at(:metadata, :job_class_name, :queue, :args, :kwargs)
   = @helper_bundle.stringify_keys(deep_dup( || {}))
  Kaal::Support::HashTools.deep_merge(
    ,
    'execution' => {
      'target' => active_job_dispatch?(job_class, queue) ? 'active_job' : 'ruby',
      'job_class' => job_class_name,
      'queue' => queue,
      'args' => args,
      'kwargs' => kwargs
    }
  )
end

#resolved_job_class(job_class_name:, key:, queue: nil) ⇒ Kaal::rbs_any

Parameters:

  • job_class_name: (Kaal::rbs_any)
  • key: (Kaal::rbs_any)
  • queue: (Kaal::rbs_any, nil) (defaults to: nil)

Returns:

  • (Kaal::rbs_any)


79
80
81
82
83
84
85
86
# File 'lib/kaal/scheduler_file/job_applier.rb', line 79

def resolved_job_class(job_class_name:, key:, queue: nil)
  Kaal::JobDispatcher.resolve_job_class(
    job_class_name:,
    key:,
    queue:,
    apply_delayed_job_allow_list: false
  )
end

#rollback_job(key:, existing_definition:, existing_registry_entry:) ⇒ Kaal::rbs_any

Parameters:

  • key: (Kaal::rbs_any)
  • existing_definition: (Kaal::rbs_any)
  • existing_registry_entry: (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
# File 'lib/kaal/scheduler_file/job_applier.rb', line 106

def rollback_job(key:, existing_definition:, existing_registry_entry:)
  if existing_definition
    @definition_registry.upsert_definition(
      **Definition::AttributeHelpers.definition_attributes(existing_definition), enabled: existing_definition[:enabled]
    )
  else
    @definition_registry.remove_definition(key)
  end

  @registry.remove(key) if @registry.registered?(key)

  return unless existing_registry_entry

  @registry.upsert(
    key: existing_registry_entry.key,
    cron: existing_registry_entry.cron,
    enqueue: existing_registry_entry.enqueue
  )
rescue StandardError => e
  @logger&.error("Failed to rollback scheduler file application for #{key}: #{e.message}")
end

#rollback_jobs(applied_job_contexts) ⇒ Kaal::rbs_any

Parameters:

  • applied_job_contexts (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


60
61
62
63
64
# File 'lib/kaal/scheduler_file/job_applier.rb', line 60

def rollback_jobs(applied_job_contexts)
  applied_job_contexts.reverse_each do |applied_job_context|
    rollback_job(**applied_job_context)
  end
end

#validate_keyword_keys(raw_kwargs, key) ⇒ nil

Parameters:

  • raw_kwargs (Kaal::rbs_any)
  • key (Kaal::rbs_any)

Returns:

  • (nil)


169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
# File 'lib/kaal/scheduler_file/job_applier.rb', line 169

def validate_keyword_keys(raw_kwargs, key)
  keys = raw_kwargs.keys
  index = 0
  while index < keys.length
    kwargs_key = keys[index]
    if kwargs_key.is_a?(String) || kwargs_key.is_a?(Symbol)
      index += 1
      next
    end

    raise SchedulerConfigError,
          "Invalid keyword argument key #{kwargs_key.inspect} (#{kwargs_key.class}) for scheduler job '#{key}'"
  end

  nil
end