Skip to content

chore: solid_queue v1.4.0 → v1.5.0 — update Mongoid models for upstream changes #7

Description

@github-actions

solid_queue v1.5.0 — Upstream Model Changes

Comparing rails/solid_queue models between v1.4.0v1.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

File Upstream Change
claimed_execution.rb 🔄 Modified
job/executable.rb 🔄 Modified
job/schedulable.rb 🔄 Modified
queue.rb 🔄 Modified
queue_selector.rb 🔄 Modified
record.rb 🔄 Modified
recurring_task.rb 🔄 Modified
semaphore.rb 🔄 Modified
record/distinct_values.rb 🆕 Added in upstream

Review Checklist

  • lib/solid_queue_mongoid/models/claimed_execution.rb — review upstream changes
  • lib/solid_queue_mongoid/models/job/executable.rb — review upstream changes
  • lib/solid_queue_mongoid/models/job/schedulable.rb — review upstream changes
  • lib/solid_queue_mongoid/models/queue.rb — review upstream changes
  • lib/solid_queue_mongoid/models/queue_selector.rb — review upstream changes
  • lib/solid_queue_mongoid/models/record.rb — review upstream changes
  • lib/solid_queue_mongoid/models/recurring_task.rb — review upstream changes
  • lib/solid_queue_mongoid/models/semaphore.rb — review upstream changes
  • lib/solid_queue_mongoid/models/record/distinct_values.rb — new upstream file, consider adding

Detailed Diffs


claimed_execution.rb — 🔄 Modified upstream

📥 What changed in upstream (v1.4.0 → v1.5.0)
--- solid_queue@v1.4.0/claimed_execution.rb
+++ solid_queue@v1.5.0/claimed_execution.rb
@@ -43,7 +43,6 @@
         SolidQueue.instrument(:fail_many_claimed) do |payload|
           executions.each do |execution|
             execution.failed_with(error)
-            execution.unblock_next_job
           end
 
           payload[:process_ids] = executions.map(&:process_id).uniq
@@ -71,13 +70,11 @@
       failed_with(result.error)
       raise result.error
     end
-  ensure
-    unblock_next_job
   end
 
   def release
     SolidQueue.instrument(:release_claimed, job_id: job.id, process_id: process_id) do
-      transaction do
+      unless_already_finalized do
         job.dispatch_bypassing_concurrency_limits
         destroy!
       end
@@ -89,14 +86,7 @@
   end
 
   def failed_with(error)
-    transaction do
-      job.failed_with(error)
-      destroy!
-    end
-  end
-
-  def unblock_next_job
-    job.unblock_next_blocked_job
+    finalize { job.failed_with(error) }
   end
 
   private
@@ -108,9 +98,28 @@
     end
 
     def finished
-      transaction do
-        job.finished!
+      finalize { job.finished! }
+    end
+
+    def finalize
+      finalized = unless_already_finalized do
+        yield
         destroy!
+        true
+      end
+
+      # Unblock the next job outside the finalize transaction so a failure while
+      # releasing the concurrency lock or dispatching the next job can't roll back
+      # a job that already finished or failed. Only the actor that owned and
+      # finalized the claim gets here, so the lock is released exactly once.
+      job.unblock_next_blocked_job if finalized
+    end
+
+    def unless_already_finalized
+      transaction do
+        return false unless self.class.unscoped.lock.find_by(id: id)
+
+        yield
       end
     end
 end
📄 Our current Mongoid model — `lib/solid_queue_mongoid/models/claimed_execution.rb`
# frozen_string_literal: true

module SolidQueue
  class ClaimedExecution < Execution
    assumes_attributes_from_job # inherits queue_name and priority from job

    field :process_id, type: BSON::ObjectId

    belongs_to :process, class_name: "SolidQueue::Process", optional: true

    # Executions whose process_id references a process that no longer exists.
    scope :orphaned, lambda {
      existing_process_ids = SolidQueue::Process.all.pluck(:id)
      existing_process_ids.empty? ? all : where(:process_id.nin => existing_process_ids)
    }

    index({ process_id: 1 })

    Result = Struct.new(:success, :error) do
      def success?
        success
      end
    end

    class << self
      # Atomically creates ClaimedExecution records for the given job_ids and
      # yields the claimed set to the block (which deletes the ReadyExecutions).
      def claiming(job_ids, process_id, &block)
        job_data = Array(job_ids).map { |job_id| { job_id: job_id, process_id: process_id } }

        SolidQueue.instrument(:claim, process_id: process_id, job_ids: job_ids) do |payload|
          claimed = job_data.filter_map do |attrs|
            create!(attrs)
          rescue Mongoid::Errors::Validations, Mongo::Error::OperationFailure
            nil
          end

          block.call(claimed)

          payload[:size] = claimed.size
          payload[:claimed_job_ids] = claimed.map(&:job_id)
        end
      end

      def release_all
        SolidQueue.instrument(:release_many_claimed) do |payload|
          executions = all.to_a
          executions.each do |execution|
            execution.release
          rescue Mongoid::Errors::Validations, Mongo::Error::OperationFailure
            # If ReadyExecution already exists, that's fine
          end
          payload[:size] = executions.size
        end
      end

      def fail_all_with(error)
        executions = includes(:job).to_a
        return if executions.empty?

        SolidQueue.instrument(:fail_many_claimed) do |payload|
          executions.each do |execution|
            execution.failed_with(error)
            execution.unblock_next_job
          end
          payload[:process_ids] = executions.map(&:process_id).uniq
          payload[:job_ids] = executions.map(&:job_id).uniq
          payload[:size] = executions.size
        end
      end

      def discard_all_in_batches(*)
        raise Execution::UndiscardableError, "Can't discard jobs in progress"
      end

      def discard_all_from_jobs(*)
        raise Execution::UndiscardableError, "Can't discard jobs in progress"
      end
    end

    # Called by Pool thread — executes the job and marks it finished or failed.
    def perform
      result = execute

      if result.success?
        finished
      else
        failed_with(result.error)
        raise result.error
      end
    ensure
      unblock_next_job
    end

    # Release this execution back to ready (called by process deregister / prune).
    def release
      SolidQueue.instrument(:release_claimed, job_id: job.id, process_id: process_id) do
        job.dispatch_bypassing_concurrency_limits
        destroy!
      end
    end

    def discard
      raise Execution::UndiscardableError, "Can't discard a job in progress"
    end

    def failed_with(error)
      Mongoid.transaction do
        job.failed_with(error)
        destroy!
      end
    end

    def unblock_next_job
      job.unblock_next_blocked_job
    end

    private

    def execute
      ActiveJob::Base.execute(job.arguments.merge("provider_job_id" => job.id.to_s))
      Result.new(true, nil)
    rescue Exception => e # rubocop:disable Lint/RescueException
      Result.new(false, e)
    end

    def finished
      Mongoid.transaction do
        job.finished!
        destroy!
      end
    end
  end
end

job/executable.rb — 🔄 Modified upstream

📥 What changed in upstream (v1.4.0 → v1.5.0)
--- solid_queue@v1.4.0/job/executable.rb
+++ solid_queue@v1.5.0/job/executable.rb
@@ -41,15 +41,16 @@
           end
 
           def successfully_dispatched(jobs)
-            dispatched_and_ready(jobs) + dispatched_and_blocked(jobs)
+            jobs_by_id = jobs.index_by(&:id)
+            dispatched_and_ready(jobs_by_id) + dispatched_and_blocked(jobs_by_id)
           end
 
-          def dispatched_and_ready(jobs)
-            where(id: ReadyExecution.where(job_id: jobs.map(&:id)).pluck(:job_id))
+          def dispatched_and_ready(jobs_by_id)
+            ReadyExecution.where(job_id: jobs_by_id.keys).pluck(:job_id).map { |id| jobs_by_id[id] }
           end
 
-          def dispatched_and_blocked(jobs)
-            where(id: BlockedExecution.where(job_id: jobs.map(&:id)).pluck(:job_id))
+          def dispatched_and_blocked(jobs_by_id)
+            BlockedExecution.where(job_id: jobs_by_id.keys).pluck(:job_id).map { |id| jobs_by_id[id] }
           end
       end
📄 Our current Mongoid model — `lib/solid_queue_mongoid/models/job/executable.rb`
# frozen_string_literal: true

module SolidQueue
  class Job
    module Executable
      extend ActiveSupport::Concern

      included do
        include ConcurrencyControls, Schedulable, Retryable

        has_one :ready_execution,   class_name: "SolidQueue::ReadyExecution",   dependent: :destroy
        has_one :claimed_execution, class_name: "SolidQueue::ClaimedExecution", dependent: :destroy

        after_create :prepare_for_execution

        scope :finished, -> { where(:finished_at.ne => nil) }
        scope :failed,   -> { where(:id.in => SolidQueue::FailedExecution.all.pluck(:job_id)) }
        scope :pending,  -> { where(finished_at: nil) }
      end

      class_methods do # rubocop:disable Metrics/BlockLength
        # Dispatch a collection of jobs, partitioned by schedule and concurrency.
        def prepare_all_for_execution(jobs)
          due, not_yet_due = jobs.partition(&:due?)
          dispatch_all(due) + schedule_all(not_yet_due)
        end

        def dispatch_all(jobs)
          with_concurrency_limits, without_concurrency_limits = jobs.partition(&:concurrency_limited?)

          dispatch_all_at_once(without_concurrency_limits)
          dispatch_all_one_by_one(with_concurrency_limits)

          successfully_dispatched(jobs)
        end

        private

        def dispatch_all_at_once(jobs)
          ReadyExecution.create_all_from_jobs(jobs)
        end

        def dispatch_all_one_by_one(jobs)
          jobs.each(&:dispatch)
        end

        def successfully_dispatched(jobs)
          jobs.map(&:id)
          dispatched_and_ready(jobs) + dispatched_and_blocked(jobs)
        end

        def dispatched_and_ready(jobs)
          job_ids = jobs.map(&:id)
          where(:id.in => ReadyExecution.where(:job_id.in => job_ids).pluck(:job_id))
        end

        def dispatched_and_blocked(jobs)
          job_ids = jobs.map(&:id)
          where(:id.in => BlockedExecution.where(:job_id.in => job_ids).pluck(:job_id))
        end
      end

      # status helpers matching SolidQueue runtime expectations
      %w[ready claimed failed].each do |status|
        define_method("#{status}?") { public_send("#{status}_execution").present? }
      end

      def prepare_for_execution
        if due?
          dispatch
        else
          schedule
        end
      end

      def dispatch
        if due?
          if acquire_concurrency_lock
            ready
          else
            handle_concurrency_conflict
          end
        else
          schedule
        end
      end

      # Called by ClaimedExecution#release — bypasses the semaphore check.
      def dispatch_bypassing_concurrency_limits
        ready
      end

      def finished!
        if SolidQueue.preserve_finished_jobs?
          update(finished_at: Time.current)
          # Clean up the claimed execution if still present (e.g. called directly
          # outside of ClaimedExecution#finished which does its own destroy!).
          claimed_execution&.destroy
        else
          destroy!
        end
      end

      alias finish finished!

      def finished?
        finished_at.present?
      end

      def status
        if finished?
          :finished
        elsif (exec = execution)
          exec.type
        end
      end

      def discard
        execution&.discard
      end

      def ready
        existing = ReadyExecution.where(job_id: id).first
        return existing if existing

        re = ReadyExecution.new(job_id: id)
        re.queue_name = queue_name
        re.priority   = priority
        re.save!
        re
      rescue Mongoid::Errors::Validations, Mongo::Error::OperationFailure
        ReadyExecution.where(job_id: id).first
      end

      def execution
        %w[ready claimed failed].reduce(nil) do |acc, status|
          acc || public_send("#{status}_execution")
        end
      end
    end
  end
end

job/schedulable.rb — 🔄 Modified upstream

📥 What changed in upstream (v1.4.0 → v1.5.0)
--- solid_queue@v1.4.0/job/schedulable.rb
+++ solid_queue@v1.5.0/job/schedulable.rb
@@ -23,7 +23,8 @@
           end
 
           def successfully_scheduled(jobs)
-            where(id: ScheduledExecution.where(job_id: jobs.map(&:id)).pluck(:job_id))
+            jobs_by_id = jobs.index_by(&:id)
+            ScheduledExecution.where(job_id: jobs_by_id.keys).pluck(:job_id).map { |id| jobs_by_id[id] }
           end
       end
📄 Our current Mongoid model — `lib/solid_queue_mongoid/models/job/schedulable.rb`
# frozen_string_literal: true

module SolidQueue
  class Job
    module Schedulable
      extend ActiveSupport::Concern

      included do
        field :scheduled_at, type: Time

        index({ scheduled_at: 1 }, { sparse: true })

        has_one :scheduled_execution, class_name: "SolidQueue::ScheduledExecution", dependent: :destroy

        scope :scheduled, -> { where(finished_at: nil) }
      end

      class_methods do
        def schedule_all(jobs)
          schedule_all_at_once(jobs)
          successfully_scheduled(jobs)
        end

        private

        def schedule_all_at_once(jobs)
          ScheduledExecution.create_all_from_jobs(jobs)
        end

        def successfully_scheduled(jobs)
          job_ids = jobs.map(&:id)
          where(:id.in => ScheduledExecution.where(:job_id.in => job_ids).pluck(:job_id))
        end
      end

      # A job is due if it has no scheduled_at or it's in the past/present.
      def due?
        scheduled_at.nil? || scheduled_at <= Time.current
      end

      # True when a ScheduledExecution document exists.
      def scheduled?
        scheduled_execution.present?
      end

      def schedule
        ScheduledExecution.create_or_find_by!(job_id: id)
      end

      def execution
        super || scheduled_execution
      end
    end
  end
end

queue.rb — 🔄 Modified upstream

📥 What changed in upstream (v1.4.0 → v1.5.0)
--- solid_queue@v1.4.0/queue.rb
+++ solid_queue@v1.5.0/queue.rb
@@ -6,9 +6,7 @@
 
     class << self
       def all
-        Job.select(:queue_name).distinct.collect do |job|
-          new(job.queue_name)
-        end
+        Job.distinct_values_of(:queue_name).map { |name| new(name) }
       end
 
       def find_by_name(name)
📄 Our current Mongoid model — `lib/solid_queue_mongoid/models/queue.rb`
# frozen_string_literal: true

module SolidQueue
  # Plain Ruby class — mirrors upstream SolidQueue::Queue (1.3.x).
  # State is derived from actual Job/Pause/ReadyExecution documents;
  # no separate Queue collection is maintained.
  class Queue
    attr_accessor :name

    class << self
      def all
        Job.distinct(:queue_name).map { |queue_name| new(queue_name) }
      end

      def find_by_name(name)
        new(name)
      end
    end

    def initialize(name)
      @name = name
    end

    def paused?
      Pause.where(queue_name: name).exists?
    end

    def pause
      Pause.pause_queue(name)
    end

    def resume
      Pause.resume_queue(name)
    end

    def clear
      ReadyExecution.queued_as(name).discard_all_in_batches
    end

    def size
      @size ||= ReadyExecution.queued_as(name).count
    end

    def latency
      @latency ||= begin
        now = Time.current
        oldest_enqueued_at = ReadyExecution.queued_as(name).min(:created_at) || now
        (now - oldest_enqueued_at).to_i
      end
    end

    def human_latency
      ActiveSupport::Duration.build(latency).inspect
    end

    def ==(other)
      name == (other.respond_to?(:name) ? other.name : other)
    end
    alias eql? ==

    def hash
      name.hash
    end
  end
end

queue_selector.rb — 🔄 Modified upstream

📥 What changed in upstream (v1.4.0 → v1.5.0)
--- solid_queue@v1.4.0/queue_selector.rb
+++ solid_queue@v1.5.0/queue_selector.rb
@@ -43,7 +43,7 @@
       end
 
       def all_queues
-        relation.distinct(:queue_name).pluck(:queue_name)
+        relation.distinct_values_of(:queue_name)
       end
 
       def exact_names
@@ -53,7 +53,7 @@
       def prefixed_names
         if prefixes.empty? then []
         else
-          relation.where(([ "queue_name LIKE ?" ] * prefixes.count).join(" OR "), *prefixes).distinct(:queue_name).pluck(:queue_name)
+          relation.where(([ "queue_name LIKE ?" ] * prefixes.count).join(" OR "), *prefixes).distinct_values_of(:queue_name)
         end
       end
📄 Our current Mongoid model — `lib/solid_queue_mongoid/models/queue_selector.rb`
# frozen_string_literal: true

module SolidQueue
  # Mirrors SolidQueue::QueueSelector but uses Mongo queries for wildcard resolution.
  class QueueSelector
    attr_reader :raw_queues, :relation

    def initialize(queue_list, relation)
      @raw_queues = Array(queue_list).map { |q| q.to_s.strip }.presence || ["*"]
      @relation = relation
    end

    # Returns an array of Mongoid criteria scoped to individual queue names,
    # or a single all/none criteria when appropriate.
    def scoped_relations
      if all?
        [relation.all]
      elsif none?
        []
      else
        queue_names.map { |queue_name| relation.queued_as(queue_name) }
      end
    end

    private

    def all?
      include_all_queues? && paused_queues.empty?
    end

    def none?
      queue_names.empty?
    end

    def queue_names
      @queue_names ||= eligible_queues - paused_queues
    end

    def eligible_queues
      if include_all_queues?
        all_queues
      else
        in_raw_order(exact_names + prefixed_names)
      end
    end

    def include_all_queues?
      raw_queues.include?("*")
    end

    # Pull all distinct queue names currently present in this relation.
    def all_queues
      relation.distinct(:queue_name)
    end

    def exact_names
      raw_queues.select { |q| exact_name?(q) }
    end

    def prefixed_names
      return [] if prefixes.empty?

      prefixes.flat_map do |prefix|
        # Use anchored regex for mongo prefix match
        relation.where(queue_name: /\A#{Regexp.escape(prefix)}/).distinct(:queue_name)
      end.uniq
    end

    def prefixes
      @prefixes ||= raw_queues.select { |q| prefixed_name?(q) }.map { |q| q.chomp("*") }
    end

    def exact_name?(queue)
      !queue.include?("*")
    end

    def prefixed_name?(queue)
      queue.end_with?("*")
    end

    def paused_queues
      @paused_queues ||= Pause.all.pluck(:queue_name)
    end

    def in_raw_order(queues)
      return queues if queues.size <= 1 || prefixes.empty?

      queues = queues.dup
      raw_queues.flat_map { |raw| delete_in_order(raw, queues) }.compact
    end

    def delete_in_order(raw_queue, queues)
      if exact_name?(raw_queue)
        queues.delete(raw_queue)
      elsif prefixed_name?(raw_queue)
        prefix = raw_queue.chomp("*")
        queues.select { |q| q.start_with?(prefix) }.tap { |matches| queues -= matches }
      end
    end
  end
end

record.rb — 🔄 Modified upstream

📥 What changed in upstream (v1.4.0 → v1.5.0)
--- solid_queue@v1.4.0/record.rb
+++ solid_queue@v1.5.0/record.rb
@@ -3,6 +3,9 @@
 module SolidQueue
   class Record < ActiveRecord::Base
     self.abstract_class = true
+    self.strict_loading_by_default = false
+
+    include DistinctValues
 
     connects_to(**SolidQueue.connects_to) if SolidQueue.connects_to
📄 Our current Mongoid model — `lib/solid_queue_mongoid/models/record.rb`
# frozen_string_literal: true

module SolidQueue
  class Record
    include Mongoid::Document
    include Mongoid::Timestamps

    # Subclasses declare their MySQL index name → MongoDB field spec mappings here.
    INDEX_HINTS = {}.freeze

    # Override Mongoid's index_specifications to use per-class instance variables
    # instead of the shared cattr_accessor class variable.
    # This prevents index cross-contamination between models.
    class << self
      def index_specifications
        @index_specifications ||= []
      end

      def index_specifications=(val)
        @_sq_index_specs = val
      end

      # Override Mongoid's index() to use our per-class storage.
      def index(spec, options = nil)
        specification = Mongoid::Indexable::Specification.new(self, spec, options)
        return if index_specifications.include?(specification)

        index_specifications.push(specification)
      end
    end

    # Dynamic collection naming with prefix
    def self.inherited(subclass)
      super

      collection_name = subclass.name.demodulize.tableize
      prefixed_name = "#{SolidQueue.collection_prefix}#{collection_name}"

      subclass.store_in collection: prefixed_name, client: -> { SolidQueue.client.to_s }

      # Each subclass gets its own empty index_specifications array.
      # Indexes must be explicitly defined in each subclass (not inherited from parent)
      # because Mongoid::Indexable::Specification stores a reference to the klass,
      # and copying parent specs would create indexes on the parent's collection.
      subclass.instance_variable_set(:@_sq_index_specs, [])
    end

    class << self
      # MongoDB has no row-level locking; this is a no-op stub so SolidQueue
      # code that chains .non_blocking_lock still works.
      def non_blocking_lock
        all
      end

      # Translate solid_queue MySQL index names to MongoDB hint specs and apply
      # via Mongoid's .hint(). Unknown names are ignored (no hint applied).
      # Subclasses may override INDEX_HINTS to register their own mappings.
      def use_index(*indexes)
        specs = indexes.filter_map do |name|
          name.is_a?(Hash) ? name : self::INDEX_HINTS[name.to_sym]
        end
        specs.any? ? hint(specs.first) : all
      end

      # MongoDB supports unique indexes which serve the same purpose.
      def supports_insert_conflict_target?
        true
      end

      # Mongoid 9 supports multi-document transactions via replica sets.
      # We wrap in a MongoDB session transaction when available; fall back to a
      # plain yield for non-replica-set environments (e.g. tests with a standalone).
      def transaction(_requires_new: false, &block)
        Mongoid::QueryCache.clear_cache
        Mongoid.default_client.with_session do |session|
          session.start_transaction
          result = yield
          session.commit_transaction
          result
        end
      rescue Mongo::Error::InvalidSession, Mongo::Error::OperationFailure => e
        # Not in a replica set or session not supported — execute without transaction
        raise if e.message.to_s.include?("Transaction numbers are only allowed")

        yield
      rescue StandardError
        yield
      end

      # Mongoid equivalent of ActiveRecord's create_or_find_by!.
      # Tries to create; on duplicate-key error or uniqueness validation error
      # finds the existing record. Falls back to finding by job_id alone when
      # the full attrs lookup would miss the existing record (e.g. queue_name
      # is not in attrs but was set via a before_create callback).
      def create_or_find_by!(attrs, &block)
        record = new(attrs)
        block.call(record) if block_given?
        record.save!
        record
      rescue Mongoid::Errors::Validations => e
        # If the only errors are uniqueness-related, fall back to find the existing record
        raise unless uniqueness_only_error?(e.document)

        find_by_unique_key(attrs) || where(attrs).first || record
      rescue Mongo::Error::OperationFailure => e
        raise unless duplicate_key_error?(e)

        find_by_unique_key(attrs) || where(attrs).first || raise(e)
      end

      # find_by that raises Mongoid::Errors::DocumentNotFound when missing.
      def find_by!(attrs)
        find_by(attrs) || raise(Mongoid::Errors::DocumentNotFound.new(self, attrs))
      end

      private

      def duplicate_key_error?(err)
        msg = err.respond_to?(:message) ? err.message.to_s : err.to_s
        msg.include?("E11000") || msg.include?("duplicate key")
      end

      # Try to find an existing record using just the unique key field(s).
      # Used as fallback when find_by(full_attrs) misses because some fields
      # (e.g. queue_name) are only set by callbacks, not passed in attrs.
      # Uses where().first to avoid DocumentNotFound exceptions.
      def find_by_unique_key(attrs)
        return where(job_id: attrs[:job_id]).first if attrs[:job_id]
        return where(key: attrs[:key]).first if attrs[:key]

        nil
      end

      def uniqueness_only_error?(document)
        return false unless document.respond_to?(:errors)

        document.errors.all? do |error|
          error.type == :taken || error.message.to_s.include?("already been taken") ||
            (error.attribute.to_s != "base" &&
              document.class.validators
                      .select { |v| v.is_a?(Mongoid::Validatable::UniquenessValidator) }
                      .any? { |v| v.attributes.include?(error.attribute.to_sym) })
        end
      end
    end
  end
end

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)
--- solid_queue@v1.4.0/recurring_task.rb
+++ solid_queue@v1.5.0/recurring_task.rb
@@ -56,16 +56,17 @@
       end
     end
 
-    def delay_from_now
-      [ (next_time - Time.current).to_f, 0.1 ].max
+
+    def next_time_after(time)
+      parsed_schedule_with_time_zone.next_time(time).utc
     end
 
     def next_time
-      parsed_schedule.next_time.utc
+      parsed_schedule_with_time_zone.next_time.utc
     end
 
     def previous_time
-      parsed_schedule.previous_time.utc
+      parsed_schedule_with_time_zone.previous_time.utc
     end
 
     def last_enqueued_time
@@ -85,6 +86,7 @@
 
           perform_later.tap do |job|
             unless job.successfully_enqueued?
+              report_enqueue_error(job.enqueue_error, at: at)
               payload[:enqueue_error] = job.enqueue_error&.message
             end
           end
@@ -97,6 +99,7 @@
         payload[:skipped] = true
         false
       rescue Job::EnqueueError => error
+        report_enqueue_error(error, at: at)
         payload[:enqueue_error] = error.message
         false
       end
@@ -168,11 +171,24 @@
         end
       end
 
+      def parsed_schedule_with_time_zone
+        @parsed_schedule_with_time_zone ||= apply_default_time_zone_to(parsed_schedule)
+      end
 
       def parsed_schedule
         @parsed_schedule ||= Fugit.parse(schedule, multi: :fail)
       end
 
+      def apply_default_time_zone_to(schedule)
+        if schedule.respond_to?(:zone) && schedule.zone.nil? && default_time_zone.present?
+          Fugit.parse("#{schedule.to_cron_s} #{default_time_zone}", multi: :fail)
+        else
+          schedule
+        end
+      rescue ArgumentError
+        schedule
+      end
+
       def job_class
         @job_class ||= class_name.present? ? class_name.safe_constantize : self.class.default_job_class
       end
@@ -180,5 +196,15 @@
       def enqueue_options
         { queue: queue_name, priority: priority }.compact
       end
+
+      def default_time_zone
+        SolidQueue.time_zone
+      end
+
+      def report_enqueue_error(error, at:)
+        if error
+          Rails.error.report(error, handled: true, source: "application.solid_queue", context: { task: key, at: at })
+        end
+      end
   end
 end
📄 Our current Mongoid model — `lib/solid_queue_mongoid/models/recurring_task.rb`
# 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

semaphore.rb — 🔄 Modified upstream

📥 What changed in upstream (v1.4.0 → v1.5.0)
--- solid_queue@v1.4.0/semaphore.rb
+++ solid_queue@v1.5.0/semaphore.rb
@@ -32,7 +32,13 @@
 
     class Proxy
       def self.signal_all(jobs)
-        Semaphore.where(key: jobs.map(&:concurrency_key)).update_all("value = value + 1")
+        # Guard against incrementing a semaphore's value beyond its limit. Jobs can
+        # have different limits, so group them and cap each group with `value < limit`.
+        jobs.group_by { |job| job.concurrency_limit || 1 }.each do |limit, grouped_jobs|
+          Semaphore.where(key: grouped_jobs.map(&:concurrency_key))
+            .where(value: ...limit)
+            .update_all("value = value + 1")
+        end
       end
 
       def initialize(job)
📄 Our current Mongoid model — `lib/solid_queue_mongoid/models/semaphore.rb`
# frozen_string_literal: true

module SolidQueue
  # Semaphore for concurrency control.
  #
  # Convention (matches spec expectations):
  #   value = number of USED (acquired) slots   (0 = none in use)
  #   limit = maximum concurrent slots
  #
  #   wait:   acquire a slot — succeeds when value < limit (increments value)
  #   signal: release a slot — increments value (marks another slot returned)
  #
  # NOTE: This intentionally differs from the ActiveRecord original which uses
  # value = remaining available slots.  The specs and BlockedExecution logic here
  # both expect the "used slots" convention.
  class Semaphore < Record
    field :key, type: String
    field :value, type: Integer, default: 0 # number of currently USED slots
    field :limit, type: Integer, default: 1
    field :expires_at, type: Time

    index({ key: 1 }, { unique: true })
    index({ expires_at: 1 })

    validates :key, presence: true, uniqueness: true

    # available: value < limit (at least one slot still free to acquire)
    scope :available, -> { where("$expr" => { "$lt" => ["$value", "$limit"] }) }
    scope :expired, -> { where(:expires_at.lt => Time.current) }

    class << self
      def wait(job)
        Proxy.new(job).wait
      end

      def signal(job)
        Proxy.new(job).signal
      end

      def signal_all(jobs)
        Proxy.signal_all(jobs)
      end

      # Requires a unique index on key.
      # Returns true if created/inserted; false on duplicate.
      def create_unique_by(attributes)
        create!(attributes)
        true
      rescue Mongoid::Errors::Validations, Mongo::Error::OperationFailure => e
        raise unless duplicate_key_error?(e)

        false
      end

      private

      def duplicate_key_error?(err)
        err.message.to_s.include?("E11000") || err.message.to_s.include?("duplicate key")
      end
    end

    # ── Instance methods ──────────────────────────────────────────────────────

    # Atomically acquire one slot (increment value if value < limit).
    # Returns true on success, false when at limit.
    def acquire
      result = self.class.collection.find_one_and_update(
        { _id: id, "value" => { "$lt" => limit } },
        { "$inc" => { "value" => 1 } },
        return_document: :after
      )
      result.present?
    end

    # Release one slot (decrement value). No-op if already at 0.
    def release
      self.class.collection.find_one_and_update(
        { _id: id, "value" => { "$gt" => 0 } },
        { "$inc" => { "value" => -1 } },
        return_document: :after
      )
      true
    end

    # True when there is still room to acquire (value < limit).
    def available?
      reload
      value < limit
    end

    # ── Proxy inner class ─────────────────────────────────────────────────────
    class Proxy
      # Decrement value for every job's semaphore key (signal = release a used slot).
      def self.signal_all(jobs)
        keys = jobs.map(&:concurrency_key)
        return if keys.empty?

        Semaphore.in(key: keys).each do |sem|
          Semaphore.collection.find_one_and_update(
            { _id: sem.id, "value" => { "$gt" => 0 } },
            { "$inc" => { "value" => -1 } }
          )
        end
      end

      def initialize(job)
        @job = job
      end

      # Acquire a slot: succeeds when value < limit.
      # Creates the semaphore document on first use.
      def wait
        semaphore = Semaphore.where(key: key).first

        if semaphore
          # Atomically increment if value < limit
          attempt_acquire(semaphore.id, semaphore.limit)
        else
          attempt_creation
        end
      end

      # Release a slot: decrement value (marks one used slot as freed).
      def signal
        attempt_release_slot
      end

      private

      attr_reader :job

      # Try to create the semaphore with value=1 (one slot in use).
      def attempt_creation
        lim = limit
        if Semaphore.create_unique_by(key: key, value: 1, limit: lim, expires_at: expires_at)
          true
        else
          # Race: someone else created it first — try to acquire from existing
          sem = Semaphore.where(key: key).first
          return false unless sem

          attempt_acquire(sem.id, sem.limit)
        end
      end

      # Atomically increment value if currently < limit.
      def attempt_acquire(semaphore_id, lim)
        result = Semaphore.collection.find_one_and_update(
          { _id: semaphore_id, "value" => { "$lt" => lim } },
          { "$inc" => { "value" => 1 }, "$set" => { "expires_at" => expires_at } },
          return_document: :after
        )
        result.present?
      end

      # Atomically decrement value if currently > 0 (release one used slot).
      def attempt_release_slot
        result = Semaphore.collection.find_one_and_update(
          { "key" => key, "value" => { "$gt" => 0 } },
          { "$inc" => { "value" => -1 }, "$set" => { "expires_at" => expires_at } },
          return_document: :after
        )
        result.present?
      end

      def key
        job.concurrency_key
      end

      def expires_at
        job.respond_to?(:concurrency_duration) ? job.concurrency_duration.from_now : 5.minutes.from_now
      end

      def limit
        (job.respond_to?(:concurrency_limit) ? job.concurrency_limit : nil) || 1
      end
    end
  end
end

Metadata

Metadata

Assignees

No one assigned

    Labels

    upstream-syncsync from solid queue upstream

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions