Class: Kaal::SchedulerFileLoader::JobApplier
- Inherits:
-
Object
- Object
- Kaal::SchedulerFileLoader::JobApplier
- 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
- #active_job_dispatch?(job_class, queue) ⇒ Boolean
- #apply(job) ⇒ nil, { key: Kaal::rbs_any, existing_definition: Kaal::rbs_any, existing_registry_entry: Kaal::rbs_any }
- #build_callback(job, job_class) ⇒ Kaal::rbs_any
- #callback_for(key:, job_class_name:, queue:, args_template:, kwargs_template:) ⇒ Kaal::rbs_any
- #conflict?(key:, existing_definition:) ⇒ Boolean
- #dispatch_job(job_class, queue, args, kwargs, key) ⇒ Kaal::rbs_any
-
#initialize(configuration:, definition_registry:, registry:, logger:, helper_bundle:) ⇒ JobApplier
constructor
A new instance of JobApplier.
- #persisted_metadata(job, job_class) ⇒ Kaal::rbs_any
- #resolved_job_class(job_class_name:, key:, queue: nil) ⇒ Kaal::rbs_any
- #rollback_job(key:, existing_definition:, existing_registry_entry:) ⇒ Kaal::rbs_any
- #rollback_jobs(applied_job_contexts) ⇒ Kaal::rbs_any
- #validate_keyword_keys(raw_kwargs, key) ⇒ nil
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.
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
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 }
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
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
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
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
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
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
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
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.}") end |
#rollback_jobs(applied_job_contexts) ⇒ 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
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 |