# frozen_string_literal: true
require "fugit"
module SolidQueue
class RecurringTask < Record
field :key, type: String
field :schedule, type: String
field :command, type: String
field :class_name, type: String
field :arguments, type: Array, default: []
field :queue_name, type: String
field :priority, type: Integer, default: 0
field :description, type: String
field :static, type: Boolean, default: false
index({ key: 1 }, { unique: true })
scope :static, -> { where(static: true) }
scope :dynamic, -> { where(static: false) }
validates :key, presence: true
validate :ensure_schedule_supported
validate :ensure_command_or_class_present
validate :ensure_existing_job_class
has_many :recurring_executions, foreign_key: :task_key, primary_key: :key,
class_name: "SolidQueue::RecurringExecution"
mattr_accessor :default_job_class
self.default_job_class = "SolidQueue::RecurringJob".safe_constantize
class << self
def wrap(args)
args.is_a?(self) ? args : from_configuration(args.first, **args.second)
end
def from_configuration(key, **options)
new(
key: key,
class_name: options[:class],
command: options[:command],
arguments: Array(options[:args]),
schedule: options[:schedule],
queue_name: options[:queue].presence,
priority: options[:priority].presence,
description: options[:description],
static: options.fetch(:static, true)
)
end
def create_dynamic_task(key, **options)
from_configuration(key, **options.merge(static: false)).save!
end
def delete_dynamic_task(key)
RecurringTask.dynamic.find_by!(key: key).destroy
end
# Upsert all static tasks; used by Scheduler::RecurringSchedule#persist_tasks.
def create_or_update_all(tasks)
tasks.each do |task|
existing = where(key: task.key).first
if existing
existing.update!(task.attributes_for_upsert)
else
create!(task.attributes_for_upsert.merge(key: task.key))
end
end
end
end
def delay_from_now
[(next_time - Time.current).to_f, 0.1].max
end
def next_time
parsed_schedule.next_time.utc
end
def previous_time
parsed_schedule.previous_time.utc
end
def last_enqueued_time
recurring_executions.maximum(:run_at)
end
def enqueue(at:)
SolidQueue.instrument(:enqueue_recurring_task, task: key, at: at) do |payload|
active_job = if using_solid_queue_adapter?
enqueue_and_record(run_at: at)
else
payload[:other_adapter] = true
perform_later.tap do |job|
payload[:enqueue_error] = job.enqueue_error&.message unless job.successfully_enqueued?
end
end
active_job.tap do |enqueued_job|
payload[:active_job_id] = enqueued_job.job_id if enqueued_job
end
rescue RecurringExecution::AlreadyRecorded
payload[:skipped] = true
false
rescue Job::EnqueueError => e
payload[:enqueue_error] = e.message
false
end
end
def to_s
"#{class_name}.perform_later(#{arguments.map(&:inspect).join(",")}) [ #{parsed_schedule.original} ]"
end
def attributes_for_upsert
attrs = attributes.except("_id", "id", "created_at", "updated_at")
attrs.delete("key")
attrs
end
private
def ensure_schedule_supported
unless parsed_schedule.instance_of?(Fugit::Cron)
errors.add :schedule, :unsupported, message: "is not a supported recurring schedule"
end
rescue ArgumentError => e
message = if e.message.include?("multiple crons")
"generates multiple cron schedules. Please use separate recurring tasks for each schedule, " \
"or use explicit cron syntax (e.g., '40 0,15 * * *' for multiple times with the same minutes)"
else
e.message
end
errors.add :schedule, :unsupported, message: message
end
def ensure_command_or_class_present
return if command.present? || class_name.present?
errors.add :base, :command_and_class_blank, message: "either command or class must be present"
end
def ensure_existing_job_class
return unless class_name.present? && job_class.nil?
errors.add :class_name, :undefined, message: "doesn't correspond to an existing class"
end
def using_solid_queue_adapter?
job_class.respond_to?(:queue_adapter_name) &&
job_class.queue_adapter_name.inquiry.solid_queue?
end
def enqueue_and_record(run_at:)
RecurringExecution.record(key, run_at) do
job_class.new(*arguments_with_kwargs).set(enqueue_options).tap do |active_job|
active_job.run_callbacks(:enqueue) do
Job.enqueue(active_job)
end
end
end
end
def perform_later
job_class.new(*arguments_with_kwargs).tap do |active_job|
active_job.enqueue(enqueue_options)
end
end
def arguments_with_kwargs
if class_name.nil?
command
elsif arguments.last.is_a?(Hash)
arguments[0...-1] + [Hash.ruby2_keywords_hash(arguments.last)]
else
arguments
end
end
def parsed_schedule
@parsed_schedule ||= Fugit.parse(schedule, multi: :fail)
end
def job_class
@job_class ||= class_name.present? ? class_name.safe_constantize : self.class.default_job_class
end
def enqueue_options
{ queue: queue_name, priority: priority }.compact
end
end
end
solid_queue
v1.5.0— Upstream Model ChangesComparing
rails/solid_queuemodels betweenv1.4.0→v1.5.0.Each section shows what changed upstream alongside our current Mongoid model for context.
Review each diff and decide if the corresponding Mongoid model needs updating.
Summary
claimed_execution.rbjob/executable.rbjob/schedulable.rbqueue.rbqueue_selector.rbrecord.rbrecurring_task.rbsemaphore.rbrecord/distinct_values.rbReview Checklist
lib/solid_queue_mongoid/models/claimed_execution.rb— review upstream changeslib/solid_queue_mongoid/models/job/executable.rb— review upstream changeslib/solid_queue_mongoid/models/job/schedulable.rb— review upstream changeslib/solid_queue_mongoid/models/queue.rb— review upstream changeslib/solid_queue_mongoid/models/queue_selector.rb— review upstream changeslib/solid_queue_mongoid/models/record.rb— review upstream changeslib/solid_queue_mongoid/models/recurring_task.rb— review upstream changeslib/solid_queue_mongoid/models/semaphore.rb— review upstream changeslib/solid_queue_mongoid/models/record/distinct_values.rb— new upstream file, consider addingDetailed Diffs
claimed_execution.rb— 🔄 Modified upstream📥 What changed in upstream (v1.4.0 → v1.5.0)
📄 Our current Mongoid model — `lib/solid_queue_mongoid/models/claimed_execution.rb`
job/executable.rb— 🔄 Modified upstream📥 What changed in upstream (v1.4.0 → v1.5.0)
📄 Our current Mongoid model — `lib/solid_queue_mongoid/models/job/executable.rb`
job/schedulable.rb— 🔄 Modified upstream📥 What changed in upstream (v1.4.0 → v1.5.0)
📄 Our current Mongoid model — `lib/solid_queue_mongoid/models/job/schedulable.rb`
queue.rb— 🔄 Modified upstream📥 What changed in upstream (v1.4.0 → v1.5.0)
📄 Our current Mongoid model — `lib/solid_queue_mongoid/models/queue.rb`
queue_selector.rb— 🔄 Modified upstream📥 What changed in upstream (v1.4.0 → v1.5.0)
📄 Our current Mongoid model — `lib/solid_queue_mongoid/models/queue_selector.rb`
record.rb— 🔄 Modified upstream📥 What changed in upstream (v1.4.0 → v1.5.0)
📄 Our current Mongoid model — `lib/solid_queue_mongoid/models/record.rb`
record/distinct_values.rb— 🆕 New upstream file📥 What changed in upstream (v1.4.0 → v1.5.0)
# frozen_string_literal: true module SolidQueue class Record module DistinctValues extend ActiveSupport::Concern # PostgreSQL has no native loose index scan, so a plain DISTINCT on a leading # index column degrades to a full index scan on large tables. We emulate one # with a recursive CTE that walks the index jumping between distinct values. class_methods do def distinct_values_of(column) if loose_index_scan_emulation_needed? loose_distinct_via_recursive_cte(column) else distinct.pluck(column) end end private def loose_index_scan_emulation_needed? connection.adapter_name == "PostgreSQL" end # Emulates a loose index scan, honoring the current scope (e.g. LIKE prefixes) # by building the anchor and the recursive step as scoped relations, whose # #to_sql inlines any bind parameters so they can be embedded in the raw CTE. def loose_distinct_via_recursive_cte(column) col = connection.quote_column_name(column) connection.select_values(<<~SQL.squish) WITH RECURSIVE t AS ( (#{next_distinct_value(col, "#{col} IS NOT NULL")}) UNION ALL SELECT (#{next_distinct_value(col, "#{col} > t.#{col}")}) FROM t WHERE t.#{col} IS NOT NULL ) SELECT #{col} FROM t WHERE #{col} IS NOT NULL SQL end # Smallest value of `col` within the current scope that matches `condition`. def next_distinct_value(col, condition) all.where(Arel.sql(condition)).reorder(Arel.sql(col)).limit(1).select(Arel.sql(col)).to_sql end end end end end📄 Our current Mongoid model — `lib/solid_queue_mongoid/models/record/distinct_values.rb`
No local equivalent exists for this file.
recurring_task.rb— 🔄 Modified upstream📥 What changed in upstream (v1.4.0 → v1.5.0)
📄 Our current Mongoid model — `lib/solid_queue_mongoid/models/recurring_task.rb`
semaphore.rb— 🔄 Modified upstream📥 What changed in upstream (v1.4.0 → v1.5.0)
📄 Our current Mongoid model — `lib/solid_queue_mongoid/models/semaphore.rb`