Class: Kaal::SchedulerFileLoader

Inherits:
Object
  • Object
show all
Includes:
SchedulerHashTransform, SchedulerPlaceholderSupport, Kaal::Support::HashTools
Defined in:
lib/kaal/scheduler_file/loader.rb,
lib/kaal/scheduler_file/job_applier.rb,
lib/kaal/scheduler_file/helper_bundle.rb,
lib/kaal/scheduler_file/job_normalizer.rb,
lib/kaal/scheduler_file/payload_loader.rb,
sig/kaal/scheduler_file/loader.rbs,
sig/kaal/scheduler_file/job_applier.rbs,
sig/kaal/scheduler_file/helper_bundle.rbs,
sig/kaal/scheduler_file/job_normalizer.rbs,
sig/kaal/scheduler_file/payload_loader.rbs

Overview

Loads scheduler definitions from config/kaal-scheduler.yml and registers them.

Defined Under Namespace

Classes: HelperBundle, JobApplier, JobNormalizer, PayloadLoader

Constant Summary collapse

PLACEHOLDER_PATTERN =

Returns:

  • (::Regexp)
/\{\{\s*([a-zA-Z0-9_.]+)\s*\}\}/
ALLOWED_PLACEHOLDERS =

Returns:

  • (::Hash[::String, Kaal::rbs_any])
{
  'fire_time.iso8601' => ->(ctx) { ctx.fetch(:fire_time).iso8601 },
  'fire_time.unix' => ->(ctx) { ctx.fetch(:fire_time).to_i },
  'idempotency_key' => ->(ctx) { ctx.fetch(:idempotency_key) },
  'key' => ->(ctx) { ctx.fetch(:key) }
}.freeze

Instance Method Summary collapse

Methods included from Kaal::Support::HashTools

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

Methods included from SchedulerPlaceholderSupport

#placeholder_token_anchors, #replace_placeholders, #resolve_placeholders, #validate_placeholder_key, #validate_placeholder_syntax, #validate_placeholders

Methods included from SchedulerHashTransform

#stringify_keys, #symbolize_keys_deep

Constructor Details

#initialize(configuration:, definition_registry:, registry:, logger:, runtime_context: RuntimeContext.default) ⇒ SchedulerFileLoader

Returns a new instance of SchedulerFileLoader.

Parameters:

  • configuration: (Kaal::rbs_any)
  • definition_registry: (Kaal::rbs_any)
  • registry: (Kaal::rbs_any)
  • logger: (Kaal::rbs_any)
  • runtime_context: (Kaal::rbs_any) (defaults to: RuntimeContext.default)


31
32
33
34
35
36
37
38
39
40
41
42
43
44
# File 'lib/kaal/scheduler_file/loader.rb', line 31

def initialize(
  configuration:,
  definition_registry:,
  registry:,
  logger:,
  runtime_context: RuntimeContext.default
)
  @configuration = configuration
  @definition_registry = definition_registry
  @registry = registry
  @logger = logger
  @runtime_context = runtime_context
  @placeholder_resolvers = ALLOWED_PLACEHOLDERS
end

Instance Method Details

#apply_job(job) ⇒ Kaal::rbs_any

Parameters:

  • job (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


91
92
93
# File 'lib/kaal/scheduler_file/loader.rb', line 91

def apply_job(job)
  job_applier.apply(job)
end

#build_callback(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)


107
108
109
110
111
112
113
114
115
# File 'lib/kaal/scheduler_file/loader.rb', line 107

def build_callback(key:, job_class_name:, queue:, args_template:, kwargs_template:)
  job_applier.callback_for(
    key: key,
    job_class_name: job_class_name,
    queue: queue,
    args_template: args_template,
    kwargs_template: kwargs_template
  )
end

#extract_job_options(payload, key:) ⇒ Kaal::rbs_any

Parameters:

  • payload (Kaal::rbs_any)
  • key: (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


87
88
89
# File 'lib/kaal/scheduler_file/loader.rb', line 87

def extract_job_options(payload, key:)
  job_normalizer.send(:extract_job_options, payload, key:)
end

#extract_jobs(payload) ⇒ Kaal::rbs_any

Parameters:

  • payload (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


75
76
77
# File 'lib/kaal/scheduler_file/loader.rb', line 75

def extract_jobs(payload)
  payload_loader.extract_jobs(payload)
end

#handle_missing_file(path) ⇒ Kaal::rbs_any

Parameters:

  • path (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


71
72
73
# File 'lib/kaal/scheduler_file/loader.rb', line 71

def handle_missing_file(path)
  payload_loader.handle_missing_file(path)
end

#helper_bundle ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


148
149
150
# File 'lib/kaal/scheduler_file/loader.rb', line 148

def helper_bundle
  @helper_bundle ||= HelperBundle.new(loader: self)
end

#job_applier ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


138
139
140
141
142
143
144
145
146
# File 'lib/kaal/scheduler_file/loader.rb', line 138

def job_applier
  @job_applier ||= JobApplier.new(
    configuration: @configuration,
    definition_registry: @definition_registry,
    registry: @registry,
    logger: @logger,
    helper_bundle: helper_bundle
  )
end

#job_normalizer ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


130
131
132
133
134
135
136
# File 'lib/kaal/scheduler_file/loader.rb', line 130

def job_normalizer
  @job_normalizer ||= JobNormalizer.new(
    hash_transform: helper_bundle,
    placeholder_support: helper_bundle,
    cron_validator: ->(cron) { Kaal.valid?(cron) }
  )
end

#load ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
# File 'lib/kaal/scheduler_file/loader.rb', line 46

def load
  applied_job_contexts = []
  path, payload = payload_loader.load
  return handle_missing_file(path) unless payload

  jobs = extract_jobs(payload)
  validate_unique_keys(jobs)
  normalized_jobs = jobs.map { |job_payload| normalize_job(job_payload) }
  applied_jobs = []
  normalized_jobs.each do |job|
    applied_job_context = apply_job(job)
    next unless applied_job_context

    applied_jobs << job
    applied_job_contexts << applied_job_context
  end

  applied_jobs
rescue StandardError
  rollback_applied_jobs(applied_job_contexts)
  raise
end

#normalize_job(job_payload) ⇒ Kaal::rbs_any

Parameters:

  • job_payload (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


83
84
85
# File 'lib/kaal/scheduler_file/loader.rb', line 83

def normalize_job(job_payload)
  job_normalizer.call(job_payload)
end

#payload_loader ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


121
122
123
124
125
126
127
128
# File 'lib/kaal/scheduler_file/loader.rb', line 121

def payload_loader
  @payload_loader ||= PayloadLoader.new(
    configuration: @configuration,
    runtime_context: @runtime_context,
    logger: @logger,
    hash_transform: helper_bundle
  )
end

#resolve_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)


117
118
119
# File 'lib/kaal/scheduler_file/loader.rb', line 117

def resolve_job_class(job_class_name:, key:, queue: nil)
  job_applier.resolved_job_class(job_class_name:, key:, queue:)
end

#rollback_applied_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)


99
100
101
# File 'lib/kaal/scheduler_file/loader.rb', line 99

def rollback_applied_job(key:, existing_definition:, existing_registry_entry:)
  job_applier.rollback_job(key:, existing_definition:, existing_registry_entry:)
end

#rollback_applied_jobs(applied_job_contexts = []) ⇒ Kaal::rbs_any

Parameters:

  • applied_job_contexts (Kaal::rbs_any) (defaults to: [])

Returns:

  • (Kaal::rbs_any)


95
96
97
# File 'lib/kaal/scheduler_file/loader.rb', line 95

def rollback_applied_jobs(applied_job_contexts = [])
  job_applier.rollback_jobs(applied_job_contexts)
end

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

Parameters:

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

Returns:

  • (Boolean)


103
104
105
# File 'lib/kaal/scheduler_file/loader.rb', line 103

def skip_due_to_conflict?(key:, existing_definition:)
  job_applier.conflict?(key:, existing_definition:)
end

#validate_unique_keys(jobs) ⇒ Kaal::rbs_any

Parameters:

  • jobs (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


79
80
81
# File 'lib/kaal/scheduler_file/loader.rb', line 79

def validate_unique_keys(jobs)
  payload_loader.validate_unique_keys(jobs)
end