Class: Kaal::Coordinator

Inherits:
Object
  • Object
show all
Defined in:
lib/kaal/core/coordinator.rb,
sig/kaal/core/coordinator.rbs

Overview

Coordinator manages the main scheduler loop that calculates due fire times and dispatches cron work safely across multiple nodes using backend lease coordination.

The coordinator:

  1. Runs a background thread on tick_interval
  2. Calculates due cron fire times and acquires distributed leases for them
  3. Dispatches claimed work and supports graceful shutdown and test re-entrancy

Constant Summary collapse

DELAYED_JOB_BATCH_SIZE =

Returns:

  • (100)
100
DELAYED_JOB_MAX_BATCHES_PER_TICK =

Returns:

  • (10)
10
DELAYED_JOB_DELETE_CONFIRMATION_JITTER_MAX =

Returns:

  • (::Float)
0.05

Instance Method Summary collapse

Constructor Details

#initialize(configuration:, registry:) ⇒ Coordinator

Initialize a new Coordinator instance.

Parameters:

  • configuration (Configuration) —

    the scheduler configuration

  • registry (Registry) —

    the registered crons registry

  • configuration: (Kaal::rbs_any)
  • registry: (Kaal::rbs_any)


30
31
32
33
34
35
36
37
38
# File 'lib/kaal/core/coordinator.rb', line 30

def initialize(configuration:, registry:)
  @configuration = configuration
  @registry = registry
  @thread = nil
  @running = false
  @stop_requested = false
  @mutex = Mutex.new
  @tick_cv = ConditionVariable.new
end

Instance Method Details

#acquire_lock(lock_key) ⇒ Kaal::rbs_any

Parameters:

  • lock_key (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


343
344
345
346
347
348
349
350
351
352
353
354
# File 'lib/kaal/core/coordinator.rb', line 343

def acquire_lock(lock_key)
  backend = @configuration.backend
  logger = @configuration.logger

  # No adapter = no locking (dev/test)
  return true unless backend

  backend.acquire(lock_key, @configuration.lease_ttl)
rescue StandardError => e
  logger&.error("Lock acquisition failed for #{lock_key}: #{e.message}")
  false
end

#already_dispatched?(key, fire_time) ⇒ Boolean

Check if a cron job was already dispatched.

Parameters:

  • key (String) —

    the cron job key

  • fire_time (Time) —

    the fire time to check

Returns:

  • (Boolean) —

    true if already dispatched, false otherwise



333
334
335
336
337
338
339
340
341
# File 'lib/kaal/core/coordinator.rb', line 333

def already_dispatched?(key, fire_time)
  adapter = @configuration.backend
  return false if adapter.nil? || !adapter.respond_to?(:dispatch_registry)

  adapter.dispatch_registry.dispatched?(key, fire_time)
rescue StandardError => e
  @configuration.logger&.warn("Error checking dispatch status for #{key}: #{e.message}")
  false
end

#apply_delayed_job_claim_jitter_if_needed(delayed_store) ⇒ nil, Kaal::rbs_any

Parameters:

  • delayed_store (Kaal::rbs_any)

Returns:

  • (nil, Kaal::rbs_any)


425
426
427
428
429
430
# File 'lib/kaal/core/coordinator.rb', line 425

def apply_delayed_job_claim_jitter_if_needed(delayed_store)
  return unless delayed_store.claim_strategy == :delete_confirmation

  jitter = rand * DELAYED_JOB_DELETE_CONFIRMATION_JITTER_MAX
  sleep(jitter) if jitter.positive?
end

#calculate_and_dispatch_due_times(entry) ⇒ nil, Kaal::rbs_any

Parameters:

  • entry (Kaal::rbs_any)

Returns:

  • (nil, Kaal::rbs_any)


176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
# File 'lib/kaal/core/coordinator.rb', line 176

def calculate_and_dispatch_due_times(entry)
  now = Time.now.utc
  window_start = now - @configuration.window_lookback
  window_end = now + @configuration.window_lookahead

  # Parse cron expression using fugit
  cron = parse_cron(entry.cron)
  return unless cron

  # Find all occurrences within the window
  occurrences = find_occurrences(cron, window_start, window_end)
  @configuration.logger&.debug("Coordinator: Found #{occurrences.length} occurrences for #{entry.key} in window [#{window_start}, #{window_end}]")

  # For each occurrence that's due (in the past or now), try to dispatch
  occurrences.each do |fire_time|
    dispatch_if_due(entry, fire_time, now)
  end
end

#cleanup_old_dispatch_records(recovery_window) ⇒ void

This method returns an undefined value.

Clean up old dispatch records to prevent database bloat.

Called after recovery completes. Deletes dispatch records older than the recovery window, since they are no longer needed for future recovery.

Parameters:

  • recovery_window (Integer) —

    seconds - records older than this are deleted



313
314
315
316
317
318
319
320
321
322
323
324
325
# File 'lib/kaal/core/coordinator.rb', line 313

def cleanup_old_dispatch_records(recovery_window)
  logger = @configuration.logger
  adapter = @configuration.backend
  return if adapter.nil? || !adapter.respond_to?(:dispatch_registry)

  registry = adapter.dispatch_registry
  return unless registry.respond_to?(:cleanup)

  deleted_count = registry.cleanup(recovery_window: recovery_window)
  logger&.debug("Cleaned up #{deleted_count} old dispatch records") if deleted_count.positive?
rescue StandardError => e
  logger&.warn("Error cleaning up old dispatch records: #{e.message}")
end

#delayed_store_for_tick ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


416
417
418
419
# File 'lib/kaal/core/coordinator.rb', line 416

def delayed_store_for_tick
  backend = @configuration.backend
  backend.respond_to?(:delayed_store) ? backend.delayed_store : nil
end

#dispatch_delayed_job(job, delayed_store) ⇒ Kaal::rbs_any

Parameters:

  • job (Kaal::rbs_any)
  • delayed_store (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
# File 'lib/kaal/core/coordinator.rb', line 393

def dispatch_delayed_job(job, delayed_store)
  if delayed_store.requires_dispatch_lock?
    lock_key = generate_delayed_lock_key(job.fetch(:job_id))
    return unless acquire_lock(lock_key)
  end

  job_class = Kaal::JobDispatcher.resolve_job_class(
    job_class_name: job.fetch(:job_class),
    key: job.fetch(:job_id),
    queue: job[:queue]
  )
  Kaal::JobDispatcher.dispatch(job_class:, queue: job[:queue], args: job.fetch(:args))
  @configuration.logger&.debug("Dispatched delayed job #{job.fetch(:job_id)} for #{job.fetch(:run_at)}")
  true
rescue StandardError => e
  Kaal::DelayedJob::DispatchFailureLogger.log_claimed_dispatch_failure(
    logger: @configuration.logger,
    job:,
    error: e
  )
  nil
end

#dispatch_due_delayed_jobs ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
# File 'lib/kaal/core/coordinator.rb', line 372

def dispatch_due_delayed_jobs
  delayed_store = delayed_store_for_tick
  return unless delayed_store

  DELAYED_JOB_MAX_BATCHES_PER_TICK.times do
    break if stop_delayed_dispatch?

    apply_delayed_job_claim_jitter_if_needed(delayed_store)
    due_jobs = delayed_store.pop_due(now: Time.now.utc, limit: DELAYED_JOB_BATCH_SIZE)
    break if due_jobs.empty?

    due_jobs.each do |job|
      break if stop_delayed_dispatch?

      dispatch_delayed_job(job, delayed_store)
    end
  end
rescue StandardError => e
  @configuration.logger&.error("Delayed job dispatch failed: #{e.message}")
end

#dispatch_if_due(entry, fire_time, now) ⇒ Kaal::rbs_any

Parameters:

  • entry (Kaal::rbs_any)
  • fire_time (Kaal::rbs_any)
  • now (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
# File 'lib/kaal/core/coordinator.rb', line 209

def dispatch_if_due(entry, fire_time, now)
  # Only dispatch if fire_time is in the past or now
  return if fire_time > now

  logger = @configuration.logger
  cron_key = entry.key
  return if Kaal.configuration.enable_log_dispatch_registry && already_dispatched?(cron_key, fire_time)

  # Generate a unique backend coordination key for this fire time
  lock_key = generate_lock_key(cron_key, fire_time)

  # Try to acquire the coordination lease
  if acquire_lock(lock_key)
    dispatch_work(entry, fire_time)
  elsif logger
    logger.debug("Failed to acquire lock for #{lock_key}")
  end
rescue StandardError => e
  cron_key ||= 'unknown'
  logger&.error("Error dispatching work for #{cron_key}: #{e.message}")
end

#dispatch_work(entry, fire_time) ⇒ Kaal::rbs_any

Parameters:

  • entry (Kaal::rbs_any)
  • fire_time (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


360
361
362
363
364
365
366
367
368
369
370
# File 'lib/kaal/core/coordinator.rb', line 360

def dispatch_work(entry, fire_time)
  # Call the enqueue callback with fire_time and idempotency_key
  cron_key = entry.key
  logger = @configuration.logger

  idempotency_key = generate_idempotency_key(cron_key, fire_time)
  entry.enqueue.call(fire_time:, idempotency_key:)
  logger&.debug("Dispatched work for #{cron_key} at #{fire_time}")
rescue StandardError => e
  logger&.error("Work dispatch failed for #{cron_key}: #{e.message}")
end

#each_enabled_entry {|arg0| ... } ⇒ Kaal::rbs_any

Yields:

Yield Parameters:

  • arg0 (Kaal::rbs_any)

Yield Returns:

  • (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


356
357
358
# File 'lib/kaal/core/coordinator.rb', line 356

def each_enabled_entry(&)
  enabled_entry_enumerator.each(&)
end

#enabled_entry_enumerator ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


454
455
456
# File 'lib/kaal/core/coordinator.rb', line 454

def enabled_entry_enumerator
  @enabled_entry_enumerator ||= EnabledEntryEnumerator.new(configuration: @configuration, registry: @registry)
end

#execute_tick ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


164
165
166
167
168
169
170
171
172
173
174
# File 'lib/kaal/core/coordinator.rb', line 164

def execute_tick
  each_enabled_entry do |entry|
    calculate_and_dispatch_due_times(entry)
  end
  dispatch_due_delayed_jobs
rescue ConfigurationError => e
  log_configuration_error('Kaal coordinator tick failed', e)
  raise
rescue StandardError => e
  log_runtime_error('Kaal coordinator tick failed', e)
end

#find_occurrences(cron, start_time, end_time) ⇒ Kaal::rbs_any

Parameters:

  • cron (Kaal::rbs_any)
  • start_time (Kaal::rbs_any)
  • end_time (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


205
206
207
# File 'lib/kaal/core/coordinator.rb', line 205

def find_occurrences(cron, start_time, end_time)
  occurrence_finder.call(cron:, start_time:, end_time:)
end

#generate_delayed_lock_key(job_id) ⇒ ::String

Parameters:

  • job_id (Kaal::rbs_any)

Returns:

  • (::String)


436
# File 'lib/kaal/core/coordinator.rb', line 436

def generate_delayed_lock_key(job_id) = "#{@configuration.namespace || 'kaal'}:delayed_dispatch:#{job_id}"

#generate_idempotency_key(cron_key, fire_time) ⇒ Kaal::rbs_any

Parameters:

  • cron_key (Kaal::rbs_any)
  • fire_time (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


432
# File 'lib/kaal/core/coordinator.rb', line 432

def generate_idempotency_key(cron_key, fire_time) = Kaal::IdempotencyKeyGenerator.call(cron_key, fire_time, configuration: @configuration)

#generate_lock_key(cron_key, fire_time) ⇒ ::String

Parameters:

  • cron_key (Kaal::rbs_any)
  • fire_time (Kaal::rbs_any)

Returns:

  • (::String)


434
# File 'lib/kaal/core/coordinator.rb', line 434

def generate_lock_key(cron_key, fire_time) = "#{@configuration.namespace || 'kaal'}:dispatch:#{cron_key}:#{fire_time.to_i}"

#log_configuration_error(prefix, error, logger: @configuration.logger) ⇒ Kaal::rbs_any

Parameters:

  • prefix (Kaal::rbs_any)
  • error (Kaal::rbs_any)
  • logger: (Kaal::rbs_any) (defaults to: @configuration.logger)

Returns:

  • (Kaal::rbs_any)


458
459
460
# File 'lib/kaal/core/coordinator.rb', line 458

def log_configuration_error(prefix, error, logger: @configuration.logger)
  logger&.error("#{prefix} due to configuration error: #{error.message}")
end

#log_runtime_error(prefix, error, logger: @configuration.logger) ⇒ Kaal::rbs_any

Parameters:

  • prefix (Kaal::rbs_any)
  • error (Kaal::rbs_any)
  • logger: (Kaal::rbs_any) (defaults to: @configuration.logger)

Returns:

  • (Kaal::rbs_any)


462
463
464
# File 'lib/kaal/core/coordinator.rb', line 462

def log_runtime_error(prefix, error, logger: @configuration.logger)
  logger&.error("#{prefix}: #{error.message}")
end

#occurrence_finder ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


450
451
452
# File 'lib/kaal/core/coordinator.rb', line 450

def occurrence_finder
  @occurrence_finder ||= OccurrenceFinder.new(configuration: @configuration)
end

#parse_cron(cron_expression) ⇒ Kaal::rbs_any

Parameters:

  • cron_expression (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


195
196
197
198
199
200
201
202
203
# File 'lib/kaal/core/coordinator.rb', line 195

def parse_cron(cron_expression)
  result = Fugit.parse_cron("#{cron_expression} #{scheduler_time_zone_resolver.time_zone_identifier}")
  raise ArgumentError, "Invalid cron expression: #{cron_expression}" unless result

  result
rescue ArgumentError => e
  @configuration.logger&.warn("Failed to parse cron expression '#{cron_expression}': #{e.message}")
  nil
end

#recover_entry(entry, start_time, end_time) ⇒ Integer

Recover missed runs for a single cron entry.

Parameters:

  • entry (Kaal::Registry::Entry) —

    the cron job entry

  • start_time (Time) —

    the start of the recovery window

  • end_time (Time) —

    the end of the recovery window

Returns:

  • (Integer) —

    number of occurrences attempted to dispatch



278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
# File 'lib/kaal/core/coordinator.rb', line 278

def recover_entry(entry, start_time, end_time)
  logger = @configuration.logger
  entry_key = entry.key
  cron = parse_cron(entry.cron)
  return 0 unless cron

  occurrences = find_occurrences(cron, start_time, end_time)

  # Filter out already-dispatched runs if dispatch logging is enabled
  occurrences.reject! { |fire_time| already_dispatched?(entry_key, fire_time) } if @configuration.enable_log_dispatch_registry
  occurrences_size = occurrences.size
  logger&.info("Recovering #{occurrences_size} missed runs for #{entry_key}")

  # Attempt to dispatch each missed occurrence
  occurrences.each do |fire_time|
    dispatch_if_due(entry, fire_time, Time.now.utc)
  end

  occurrences_size
rescue ConfigurationError => e
  log_configuration_error("Error recovering entry #{entry_key}", e, logger:)
  raise
rescue StandardError => e
  log_runtime_error("Error recovering entry #{entry_key}", e, logger:)
  0
end

#recover_missed_runs ⇒ void

This method returns an undefined value.

Recover missed cron runs after downtime.

Looks back over the recovery window to find cron jobs that should have executed but were missed due to downtime. Uses dispatch records (if enabled) to skip already-dispatched jobs, and relies on lease coordination for duplicate prevention.



239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
# File 'lib/kaal/core/coordinator.rb', line 239

def recover_missed_runs
  return unless @configuration.enable_dispatch_recovery

  # Add random jitter to reduce lock contention when multiple nodes restart simultaneously
  jitter = rand(0..@configuration.recovery_startup_jitter)
  sleep(jitter) if jitter.positive?

  current_time = Time.now.utc
  recovery_window = @configuration.recovery_window
  recovery_start = current_time - recovery_window
  recovery_end = current_time

  logger = @configuration.logger
  logger&.info("Starting missed-run recovery for window: #{recovery_start} to #{recovery_end}")

  total_recovered = 0
  each_enabled_entry do |entry|
    recovered = recover_entry(entry, recovery_start, recovery_end)
    total_recovered += recovered
  end

  logger&.info("Missed-run recovery completed: attempted #{total_recovered} dispatches")

  # Clean up old dispatch records after recovery completes
  cleanup_old_dispatch_records(recovery_window)
rescue ConfigurationError => e
  log_configuration_error('Kaal missed-run recovery failed', e, logger:)
  raise
rescue StandardError => e
  log_runtime_error('Error during missed-run recovery', e, logger:)
end

#request_stop ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


140
141
142
143
144
145
146
147
# File 'lib/kaal/core/coordinator.rb', line 140

def request_stop
  @mutex.synchronize do
    return unless @running

    @stop_requested = true
    @tick_cv.signal
  end
end

#reset! ⇒ void

This method returns an undefined value.

Reset coordinator state for re-entrancy in tests.

Stops any running thread before clearing state to avoid orphaning it. Raises an error if the thread cannot be stopped within the timeout.

Raises:

  • (RuntimeError) —

    if stop! times out



123
124
125
126
127
128
129
130
131
132
133
134
135
136
# File 'lib/kaal/core/coordinator.rb', line 123

def reset!
  # Stop any running thread first to prevent orphaned threads
  if running?
    stopped = stop!
    raise 'Failed to stop coordinator thread within timeout' unless stopped
  end

  # Now safe to reset all state
  @mutex.synchronize do
    @running = false
    @stop_requested = false
    @thread = nil
  end
end

#restart! ⇒ Thread

Restart the coordinator (stop then start).

Returns:

  • (Thread) —

    the started thread



97
98
99
100
# File 'lib/kaal/core/coordinator.rb', line 97

def restart!
  stop!
  start!
end

#run_loop ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


149
150
151
152
153
154
155
156
157
158
# File 'lib/kaal/core/coordinator.rb', line 149

def run_loop
  loop do
    break if stop_requested?

    execute_tick
    sleep_until_next_tick
  end
ensure
  @mutex.synchronize { @running = false }
end

#running? ⇒ Boolean

Check if the coordinator is currently running.

Returns:

  • (Boolean) —

    true if running, false otherwise



88
89
90
# File 'lib/kaal/core/coordinator.rb', line 88

def running?
  @mutex.synchronize { @running }
end

#scheduler_time_zone_resolver ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


446
447
448
# File 'lib/kaal/core/coordinator.rb', line 446

def scheduler_time_zone_resolver
  @scheduler_time_zone_resolver ||= SchedulerTimeZoneResolver.new(configuration: @configuration)
end

#sleep_until_next_tick ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


438
439
440
441
442
443
444
# File 'lib/kaal/core/coordinator.rb', line 438

def sleep_until_next_tick
  @mutex.synchronize do
    @tick_cv.wait(@mutex, @configuration.tick_interval)
  end
rescue StandardError => e
  @configuration.logger&.error("Sleep interrupted: #{e.message}")
end

#start! ⇒ Thread

Start the coordinator background thread.

Returns:

  • (Thread) —

    the started thread, or nil if already running



45
46
47
48
49
50
51
52
53
54
55
56
57
58
# File 'lib/kaal/core/coordinator.rb', line 45

def start!
  @mutex.synchronize do
    return nil if @running

    # Run recovery before starting the main loop
    recover_missed_runs

    @running = true
    @stop_requested = false
    @thread = Thread.new { run_loop }
    @thread.abort_on_exception = true
    return @thread
  end
end

#stop!(timeout: 30) ⇒ Boolean

Stop the coordinator gracefully.

Signals the coordinator to stop after the current tick completes, then waits for the thread to finish.

Parameters:

  • timeout (Integer) (defaults to: 30) —

    seconds to wait for graceful shutdown (default: 30)

  • timeout: (::Integer) (defaults to: 30)

Returns:

  • (Boolean) —

    true if stopped, false if timeout



69
70
71
72
73
74
75
76
77
78
79
80
81
82
# File 'lib/kaal/core/coordinator.rb', line 69

def stop!(timeout: 30) # rubocop:disable Naming/PredicateMethod
  request_stop

  # Wait for thread to finish outside the lock
  result = @thread&.join(timeout)

  # If we had a thread and join timed out, thread is still alive
  return false if @thread && result.nil?

  @thread = nil
  @mutex.synchronize { @running = false }

  true
end

#stop_delayed_dispatch? ⇒ Boolean

Returns:

  • (Boolean)


421
422
423
# File 'lib/kaal/core/coordinator.rb', line 421

def stop_delayed_dispatch?
  stop_requested?
end

#stop_requested? ⇒ Boolean

Returns:

  • (Boolean)


160
161
162
# File 'lib/kaal/core/coordinator.rb', line 160

def stop_requested?
  @mutex.synchronize { @stop_requested }
end

#tick! ⇒ void

This method returns an undefined value.

Execute a single tick manually.

This is useful for testing and Rake tasks that want to trigger the scheduler without running the background loop.



110
111
112
# File 'lib/kaal/core/coordinator.rb', line 110

def tick!
  execute_tick
end