diff --git a/README.md b/README.md index 90a1b1097..a1b67c513 100644 --- a/README.md +++ b/README.md @@ -279,6 +279,7 @@ Here's an overview of the different options: It is recommended to set this value less than or equal to the queue database's connection pool size minus 2, as each worker thread uses one connection, and two additional connections are reserved for polling and heartbeat. - `processes`: this is the number of worker processes that will be forked by the supervisor with the settings given. By default, this is `1`, just a single process. This setting is useful if you want to dedicate more than one CPU core to a queue or queues with the same configuration. Only workers have this setting. **Note**: this option will be ignored if [running in `async` mode](#fork-vs-async-mode). - `concurrency_maintenance`: whether the dispatcher will perform the concurrency maintenance work. This is `true` by default, and it's useful if you don't use any [concurrency controls](#concurrency-controls) and want to disable it or if you run multiple dispatchers and want some of them to just dispatch jobs without doing anything else. +- `batch_maintenance`: whether the dispatcher will sweep stalled [batches](#batch-jobs) as part of its maintenance work, on the same timer as concurrency maintenance (see [batch maintenance](#batch-maintenance)). This is `true` by default; disable it if you don't use batches, or if you run multiple dispatchers and want only some of them doing maintenance work. ### Optional scheduler configuration @@ -653,6 +654,9 @@ and optionally trigger callbacks based on their status. It supports the followin - If a job is part of a batch, it can enqueue more jobs for that batch using `batch#enqueue` - Attaching arbitrary metadata to a batch +Callback jobs are regular jobs: they don't receive any extra arguments, and they can access +the batch they belong to through the `batch` accessor: + ```rb class SleepyJob < ApplicationJob def perform(seconds_to_sleep) @@ -662,20 +666,20 @@ class SleepyJob < ApplicationJob end class BatchFinishJob < ApplicationJob - def perform(batch) # batch is always the default first argument - Rails.logger.info "Good job finishing all jobs" + def perform + Rails.logger.info "Finished all #{batch.total_jobs} jobs" end end class BatchSuccessJob < ApplicationJob - def perform(batch) # batch is always the default first argument - Rails.logger.info "Good job finishing all jobs, and all of them worked!" + def perform + Rails.logger.info "All #{batch.completed_jobs} jobs worked!" end end class BatchFailureJob < ApplicationJob - def perform(batch) # batch is always the default first argument - Rails.logger.info "At least one job failed, sorry!" + def perform + Rails.logger.info "#{batch.failed_jobs} jobs failed, sorry!" end end @@ -689,11 +693,24 @@ SolidQueue::Batch.enqueue( end ``` +A job joins the batch that's active when its enqueue is requested. This also works when +Rails defers the actual enqueue until after the surrounding transaction commits. + +- A job created outside a batch and enqueued inside one joins that batch. +- Creating a job inside a batch without enqueueing it doesn't keep the batch open. +- If a job already carries a batch ID but is enqueued inside another active batch, the + active batch takes precedence. + +Callbacks can be given as a job class or as a configured job instance, e.g. +`on_finish: BatchFinishJob.new.set(queue: :batches)`. Note that the job is serialized when +the batch is created, so options resolved at that point (like `wait_until:` timestamps) are +relative to batch creation, not to when the callback is eventually enqueued. + ### Batch options In the case of an empty batch, a `SolidQueue::Batch::EmptyJob` is enqueued. -By default, this jobs run on the `default` queue. You can specify an alternative queue for it in an initializer: +By default, this job runs on the `default` queue. You can specify an alternative queue for it in an initializer: ```rb Rails.application.config.after_initialize do # or to_prepare @@ -701,6 +718,96 @@ Rails.application.config.after_initialize do # or to_prepare end ``` +The empty job and batch callback jobs always enqueue through Solid Queue, even when the +job classes involved (or the application default) use a different Active Job adapter. + +### Batch progress and counters + +Batches track `total_jobs`, `completed_jobs`, `failed_jobs` and `pending_jobs`, plus a +`progress_percentage` helper. A couple of accounting details to be aware of: + +- Every *attempt* counts: when a job is retried via `retry_on`, each retry is enqueued as a + new job in the batch, so a job that fails twice and then succeeds contributes 3 to + `total_jobs` (2 completed retries + 1 success). +- Jobs discarded via `discard_on`, concurrency's `on_conflict: :discard`, or manual + discarding count as completed, not failed. +- Manually retrying a failed job (via `SolidQueue::FailedExecution#retry`) doesn't re-add it + to its batch: if the batch already finished as failed, a successful manual retry won't + change the batch's status. + +### Batch maintenance + +Batch completion is normally detected as jobs finish, without ever locking the batch row +outside a single once-per-batch moment. A few edge cases can't trigger that detection: jobs +removed via bulk discards (which delete jobs without callbacks), a process that crashed +after enqueueing jobs but before starting its batch, or a completion whose callback +enqueueing failed and rolled back. + +The dispatcher sweeps these up automatically via `SolidQueue::Batch.sweep_stalled`, as part +of its regular maintenance (every `concurrency_maintenance_interval` seconds, sharing a +single maintenance timer and database connection). If you disable `batch_maintenance` (or +don't run a dispatcher), you can run the sweep yourself, for example as a +[recurring task](#recurring-tasks): + +```yml +batch_maintenance: + command: "SolidQueue::Batch.sweep_stalled" + schedule: every 5 minutes +``` + +### Clearing batches + +Finished, non-failed batches are cleared after `config.solid_queue.clear_finished_jobs_after`, +but only when you invoke it: like jobs, batches are cleared with +`SolidQueue::Batch.clear_finished_in_batches`, which you'd typically run periodically +alongside `SolidQueue::Job.clear_finished_in_batches`. Failed batches are kept, like failed +jobs, so you can inspect them. + +### Upgrading existing installations + +If you installed Solid Queue before batches existed, add the new tables with a migration in +`db/queue_migrate`: + +```ruby +class AddSolidQueueBatches < ActiveRecord::Migration[7.1] + def change + create_table :solid_queue_batches do |t| + t.string :active_job_batch_id + t.string :description + t.text :on_finish + t.text :on_success + t.text :on_failure + t.text :metadata + t.integer :total_jobs, default: 0, null: false + t.integer :completed_jobs, default: 0, null: false + t.integer :failed_jobs, default: 0, null: false + t.datetime :enqueued_at + t.datetime :finished_at + t.datetime :failed_at + t.timestamps + + t.index :active_job_batch_id, unique: true + t.index :finished_at + end + + create_table :solid_queue_batch_executions do |t| + t.bigint :job_id, null: false + t.bigint :batch_id, null: false + t.datetime :created_at, null: false + + t.index :job_id, unique: true + t.index :batch_id + end + + add_column :solid_queue_jobs, :batch_id, :bigint + add_index :solid_queue_jobs, :batch_id + + add_foreign_key :solid_queue_batch_executions, :solid_queue_batches, column: :batch_id, on_delete: :cascade + add_foreign_key :solid_queue_batch_executions, :solid_queue_jobs, column: :job_id, on_delete: :cascade + end +end +``` + ## Puma plugin We provide a Puma plugin if you want to run the Solid Queue's supervisor together with Puma and have Puma monitor and manage it. You just need to add diff --git a/app/jobs/solid_queue/batch/empty_job.rb b/app/jobs/solid_queue/batch/empty_job.rb index d29e1ad00..e3fac1b90 100644 --- a/app/jobs/solid_queue/batch/empty_job.rb +++ b/app/jobs/solid_queue/batch/empty_job.rb @@ -3,6 +3,9 @@ module SolidQueue class Batch class EmptyJob < (defined?(ApplicationJob) ? ApplicationJob : ActiveJob::Base) + # Always use Solid Queue, even when ApplicationJob uses another adapter. + self.queue_adapter = :solid_queue + def perform # This job does nothing - it just exists to trigger batch completion # The batch completion will be handled by the normal job_finished! flow diff --git a/app/models/solid_queue/batch.rb b/app/models/solid_queue/batch.rb index 3b433d1ec..633c59548 100644 --- a/app/models/solid_queue/batch.rb +++ b/app/models/solid_queue/batch.rb @@ -2,7 +2,11 @@ module SolidQueue class Batch < Record - class AlreadyFinished < StandardError; end + class AlreadyFinished < StandardError + def initialize(message = "You cannot enqueue a batch that is already finished") + super + end + end include Trackable, Clearable @@ -18,7 +22,9 @@ class AlreadyFinished < StandardError; end end end - after_initialize :set_active_job_batch_id + # Provider-agnostic batch identifier, analogous to jobs.active_job_id. + before_create :set_active_job_batch_id + after_commit :start_batch, on: :create, unless: -> { ActiveRecord.respond_to?(:after_all_transactions_commit) } class << self @@ -50,7 +56,9 @@ def wrap_in_batch_context(batch_id) end def enqueue(&block) - raise AlreadyFinished, "You cannot enqueue a batch that is already finished" if finished? + # Fast-fail for the common case. create_all_from_jobs atomically guards + # concurrent additions when it creates their tracking rows. + raise AlreadyFinished if finished? transaction do save! if new_record? @@ -73,37 +81,81 @@ def metadata def check_completion return if finished? || !enqueued? - return if batch_executions.any? - rows = Batch - .where(id: id) - .unfinished - .empty_executions - .update_all(finished_at: Time.current) - - return if rows.zero? - - with_lock do - failed = jobs.joins(:failed_execution).count - finished_attributes = {} - if failed > 0 - finished_attributes[:failed_at] = Time.current - finished_attributes[:failed_jobs] = failed + return if batch_executions.exists? + + transaction do + finished_rows = Batch.where(id: id).unfinished.enqueued.empty_executions.update_all(finished_at: Time.current) + finalize_completion if finished_rows.positive? + end + end + + COMPLETION_GRACE = 3.seconds + + def self.sweep_stalled(stalled_for: 5.minutes, batch_size: 500) + SolidQueue.instrument(:sweep_stalled_batches, stalled_for: stalled_for, size: 0, started: 0, repaired: 0) do |payload| + # BatchExecution rows represent outstanding work. A row for a resolved + # job violates that invariant, so remove it immediately; destroy's + # after_commit callback retries the batch completion check. + [ BatchExecution.for_finished_jobs, BatchExecution.for_failed_jobs ].each do |leaked| + leaked.find_each(batch_size: batch_size) do |batch_execution| + payload[:repaired] += 1 + batch_execution.destroy + end + end + + # A started batch with no tracking rows can finish, but allow time for a + # transaction-deferred EmptyJob enqueue to become visible. + unfinished.empty_executions.where(enqueued_at: ...COMPLETION_GRACE.ago).find_each(batch_size: batch_size) do |batch| + payload[:size] += 1 + batch.check_completion end - finished_attributes[:completed_jobs] = total_jobs - failed - update!(finished_attributes) - enqueue_callback_jobs + unfinished.where(enqueued_at: nil).where(created_at: ...stalled_for.ago).find_each(batch_size: batch_size) do |batch| + payload[:started] += 1 + batch.start_batch + end end end + def start_batch + # Single-winner start so concurrent sweepers can't enqueue duplicate empty jobs + transaction do + if Batch.where(id: id, enqueued_at: nil).update_all(enqueued_at: Time.current).positive? + enqueue_empty_job if reload.total_jobs == 0 + end + end + + check_completion + end + private def set_active_job_batch_id self.active_job_batch_id ||= SecureRandom.uuid end - def as_active_job(active_job_klass) - active_job_klass.is_a?(ActiveJob::Base) ? active_job_klass : active_job_klass.new + def finalize_completion + reload + + # PostgreSQL can let a blocked CAS win from a stale NOT EXISTS snapshot. + # Re-check in a new statement while this transaction holds the row lock. + raise ActiveRecord::Rollback if batch_executions.exists? + + SolidQueue.instrument(:finish_batch, batch_id: id) do |payload| + failed = jobs.failed.count + finished_attributes = { completed_jobs: total_jobs - failed } + if failed > 0 + finished_attributes[:failed_at] = Time.current + finished_attributes[:failed_jobs] = failed + end + + update_columns(finished_attributes) + enqueue_callback_jobs + + payload[:total_jobs] = total_jobs + payload[:completed_jobs] = self[:completed_jobs] + payload[:failed_jobs] = failed + end end def serialize_callback(value) @@ -118,7 +170,11 @@ def serialize_callback(value) def enqueue_callback_job(callback_name) active_job = ActiveJob::Base.deserialize(send(callback_name)) active_job.callback_batch_id = id - active_job.enqueue + # Bypass the job class's adapter so callbacks stay in Solid Queue and + # their enqueue stays in this transaction, while honoring enqueue callbacks. + active_job.run_callbacks(:enqueue) do + Job.enqueue(active_job, scheduled_at: active_job.scheduled_at || Time.current) + end end def enqueue_callback_jobs @@ -136,10 +192,5 @@ def enqueue_empty_job EmptyJob.perform_later end end - - def start_batch - enqueue_empty_job if reload.total_jobs == 0 - update!(enqueued_at: Time.current) - end end end diff --git a/app/models/solid_queue/batch/trackable.rb b/app/models/solid_queue/batch/trackable.rb index 4ee7e5369..5df157ba7 100644 --- a/app/models/solid_queue/batch/trackable.rb +++ b/app/models/solid_queue/batch/trackable.rb @@ -10,24 +10,18 @@ module Trackable scope :succeeded, -> { finished.where(failed_at: nil) } scope :unfinished, -> { where(finished_at: nil) } scope :failed, -> { where.not(failed_at: nil) } - scope :empty_executions, -> { - where(<<~SQL) - NOT EXISTS ( - SELECT 1 FROM solid_queue_batch_executions - WHERE solid_queue_batch_executions.batch_id = solid_queue_batches.id - LIMIT 1 - ) - SQL - } + scope :enqueued, -> { where.not(enqueued_at: nil) } + # Join-free so update_all keeps this condition in the completion update's own WHERE + scope :empty_executions, -> { where.not(id: BatchExecution.select(:batch_id)) } end def status if finished? - failed? ? "failed" : "completed" + failed? ? :failed : :completed elsif enqueued? - "enqueued" + :enqueued else - "pending" + :pending end end @@ -47,12 +41,13 @@ def enqueued? enqueued_at.present? end + # Failed jobs no longer have tracking rows, so exclude them from the completed count. def completed_jobs - finished? ? self[:completed_jobs] : total_jobs - batch_executions.count + finished? ? self[:completed_jobs] : total_jobs - pending_jobs - failed_jobs end def failed_jobs - finished? ? self[:failed_jobs] : jobs.joins(:failed_execution).count + finished? ? self[:failed_jobs] : jobs.failed.count end def pending_jobs @@ -61,7 +56,7 @@ def pending_jobs def progress_percentage return 0 if total_jobs == 0 - ((completed_jobs + failed_jobs) * 100.0 / total_jobs).round(2) + ((total_jobs - pending_jobs) * 100.0 / total_jobs).round(2) end end end diff --git a/app/models/solid_queue/batch_execution.rb b/app/models/solid_queue/batch_execution.rb index 1957f2e59..95e82adbe 100644 --- a/app/models/solid_queue/batch_execution.rb +++ b/app/models/solid_queue/batch_execution.rb @@ -5,13 +5,10 @@ class BatchExecution < Record belongs_to :job, optional: true belongs_to :batch - after_commit :check_completion, on: :destroy + scope :for_finished_jobs, -> { joins(:job).merge(SolidQueue::Job.finished) } + scope :for_failed_jobs, -> { joins(job: :failed_execution) } - private - def check_completion - batch = Batch.find_by(id: batch_id) - batch.check_completion if batch.present? - end + after_commit :check_completion, on: :destroy class << self def create_all_from_jobs(jobs) @@ -19,14 +16,24 @@ def create_all_from_jobs(jobs) return if batch_jobs.empty? batch_jobs.group_by(&:batch_id).each do |batch_id, jobs| + # Increment first: inserting tracking rows takes a shared FK lock on + # the batch row, then incrementing can deadlock concurrent MySQL adders. + total = jobs.size + updated = SolidQueue::Batch.where(id: batch_id).unfinished.update_all([ "total_jobs = total_jobs + ?", total ]) + raise Batch::AlreadyFinished if updated.zero? + BatchExecution.insert_all!(jobs.map { |job| { batch_id:, job_id: job.respond_to?(:provider_job_id) ? job.provider_job_id : job.id } }) - - total = jobs.size - SolidQueue::Batch.where(id: batch_id).update_all([ "total_jobs = total_jobs + ?", total ]) end end end + + private + def check_completion + # Skip the serialized callback and metadata columns on this hot path + batch = Batch.select(:id, :finished_at, :enqueued_at).find_by(id: batch_id) + batch.check_completion if batch.present? + end end end diff --git a/app/models/solid_queue/execution/batchable.rb b/app/models/solid_queue/execution/batchable.rb deleted file mode 100644 index fe9aa6ad3..000000000 --- a/app/models/solid_queue/execution/batchable.rb +++ /dev/null @@ -1,23 +0,0 @@ -# frozen_string_literal: true - -module SolidQueue - class Execution - module Batchable - extend ActiveSupport::Concern - - included do - after_create :update_batch_progress, if: -> { job.batch_id? } - end - - private - def update_batch_progress - if is_a?(FailedExecution) - # FailedExecutions are only created when the job is done retrying - job.batch_execution&.destroy! - end - rescue => e - Rails.logger.error "[SolidQueue] Failed to notify batch #{job.batch_id} about job #{job.id} failure: #{e.message}" - end - end - end -end diff --git a/app/models/solid_queue/failed_execution/batchable.rb b/app/models/solid_queue/failed_execution/batchable.rb new file mode 100644 index 000000000..64e32ef4c --- /dev/null +++ b/app/models/solid_queue/failed_execution/batchable.rb @@ -0,0 +1,22 @@ +# frozen_string_literal: true + +module SolidQueue + class FailedExecution + # A FailedExecution is created only after retries are exhausted, when the + # job stops counting as pending in its batch. + module Batchable + extend ActiveSupport::Concern + + included do + after_create :destroy_job_batch_execution, if: -> { job.batch_id? } + end + + private + def destroy_job_batch_execution + job.batch_execution&.destroy! + rescue ActiveRecord::ActiveRecordError => e + SolidQueue.instrument(:batch_progress_error, batch_id: job.batch_id, job_id: job.id, error: e) + end + end + end +end diff --git a/app/models/solid_queue/job.rb b/app/models/solid_queue/job.rb index 6cb59e12f..f595d2769 100644 --- a/app/models/solid_queue/job.rb +++ b/app/models/solid_queue/job.rb @@ -10,7 +10,13 @@ class EnqueueError < StandardError; end class << self def enqueue_all(active_jobs) - active_jobs.each { |job| job.scheduled_at ||= Time.current } + # Bulk enqueues bypass ActiveJob#enqueue, so batch membership is captured here + current_batch_id = Batch.current_batch_id + + active_jobs.each do |job| + job.scheduled_at ||= Time.current + job.batch_id = current_batch_id || job.batch_id + end active_jobs_by_job_id = active_jobs.index_by(&:job_id) transaction do diff --git a/app/models/solid_queue/job/batchable.rb b/app/models/solid_queue/job/batchable.rb index 5ab1bae44..1afdfcf20 100644 --- a/app/models/solid_queue/job/batchable.rb +++ b/app/models/solid_queue/job/batchable.rb @@ -29,8 +29,8 @@ def update_batch_progress return unless batch_id.present? batch_execution&.destroy! - rescue => e - Rails.logger.error "[SolidQueue] Failed to update batch #{batch_id} progress for job #{id}: #{e.message}" + rescue ActiveRecord::ActiveRecordError => e + SolidQueue.instrument(:batch_progress_error, batch_id: batch_id, job_id: id, error: e) end end end diff --git a/app/models/solid_queue/job/executable.rb b/app/models/solid_queue/job/executable.rb index bd3625820..1e89ca42a 100644 --- a/app/models/solid_queue/job/executable.rb +++ b/app/models/solid_queue/job/executable.rb @@ -18,10 +18,11 @@ module Executable class_methods do def prepare_all_for_execution(jobs) + # Track before dispatch so conflict-discarded jobs count like single enqueues. + batch_all(jobs) + due, not_yet_due = jobs.partition(&:due?) - (dispatch_all(due) + schedule_all(not_yet_due)).tap do |jobs| - batch_all(jobs.select { |job| job.batch_id.present? }) - end + dispatch_all(due) + schedule_all(not_yet_due) end def dispatch_all(jobs) diff --git a/lib/active_job/batch_id.rb b/lib/active_job/batch_id.rb index d9cb803cb..5ab621513 100644 --- a/lib/active_job/batch_id.rb +++ b/lib/active_job/batch_id.rb @@ -11,13 +11,26 @@ module BatchId attr_accessor :callback_batch_id end + # Seed membership at construction for enqueue paths deferred until after the + # batch context ends, notably bulk enqueue. Enqueueing inside another batch + # rebinds membership; enqueueing without a batch preserves it. + # SolidQueue::Job.enqueue_all repeats this capture because bulk enqueue + # bypasses #enqueue. def initialize(*arguments, **kwargs) super self.batch_id = SolidQueue::Batch.current_batch_id if solid_queue_job? end + def enqueue(options = {}) + self.batch_id = SolidQueue::Batch.current_batch_id || batch_id if solid_queue_job? + super + end + def serialize - super.merge("batch_id" => batch_id, "callback_batch_id" => callback_batch_id) + super.tap do |data| + data["batch_id"] = batch_id if batch_id + data["callback_batch_id"] = callback_batch_id if callback_batch_id + end end def deserialize(job_data) @@ -27,7 +40,12 @@ def deserialize(job_data) end def batch - @batch ||= SolidQueue::Batch.find_by(id: callback_batch_id || batch_id) + batch_id_to_load = callback_batch_id || batch_id + return if batch_id_to_load.nil? + return @batch if defined?(@batch) && @loaded_batch_id == batch_id_to_load + + @loaded_batch_id = batch_id_to_load + @batch = SolidQueue::Batch.find_by(id: batch_id_to_load) end private diff --git a/lib/generators/solid_queue/install/templates/db/queue_schema.rb b/lib/generators/solid_queue/install/templates/db/queue_schema.rb index 1047e9eed..f9a71dabb 100644 --- a/lib/generators/solid_queue/install/templates/db/queue_schema.rb +++ b/lib/generators/solid_queue/install/templates/db/queue_schema.rb @@ -149,6 +149,8 @@ t.index [ "batch_id" ], name: "index_solid_queue_batch_executions_on_batch_id" end + add_foreign_key "solid_queue_batch_executions", "solid_queue_batches", column: "batch_id", on_delete: :cascade + add_foreign_key "solid_queue_batch_executions", "solid_queue_jobs", column: "job_id", on_delete: :cascade add_foreign_key "solid_queue_blocked_executions", "solid_queue_jobs", column: "job_id", on_delete: :cascade add_foreign_key "solid_queue_claimed_executions", "solid_queue_jobs", column: "job_id", on_delete: :cascade add_foreign_key "solid_queue_failed_executions", "solid_queue_jobs", column: "job_id", on_delete: :cascade diff --git a/lib/solid_queue/configuration.rb b/lib/solid_queue/configuration.rb index f88ce3ed7..d63fbd59d 100644 --- a/lib/solid_queue/configuration.rb +++ b/lib/solid_queue/configuration.rb @@ -27,7 +27,8 @@ def instantiate batch_size: 500, polling_interval: 1, concurrency_maintenance: true, - concurrency_maintenance_interval: 600 + concurrency_maintenance_interval: 600, + batch_maintenance: true } SCHEDULER_DEFAULTS = { diff --git a/lib/solid_queue/dispatcher.rb b/lib/solid_queue/dispatcher.rb index 461ce8036..112399ef5 100644 --- a/lib/solid_queue/dispatcher.rb +++ b/lib/solid_queue/dispatcher.rb @@ -7,8 +7,8 @@ class Dispatcher < Processes::Poller attr_reader :batch_size after_boot :run_start_hooks - after_boot :start_concurrency_maintenance - before_shutdown :stop_concurrency_maintenance + after_boot :start_maintenance + before_shutdown :stop_maintenance before_shutdown :run_stop_hooks after_shutdown :run_exit_hooks @@ -17,17 +17,21 @@ def initialize(**options) @batch_size = options[:batch_size] - @concurrency_maintenance = ConcurrencyMaintenance.new(options[:concurrency_maintenance_interval], options[:batch_size]) if options[:concurrency_maintenance] + # Run both maintenance routines on one timer instead of another thread. + if options[:concurrency_maintenance] || options[:batch_maintenance] + @maintenance = Maintenance.new(options[:concurrency_maintenance_interval], options[:batch_size], + concurrency: options[:concurrency_maintenance], batches: options[:batch_maintenance]) + end super(**options) end def metadata - super.merge(batch_size: batch_size, concurrency_maintenance_interval: concurrency_maintenance&.interval) + super.merge(batch_size: batch_size).merge(maintenance&.metadata || {}) end private - attr_reader :concurrency_maintenance + attr_reader :maintenance def poll batch = dispatch_next_batch @@ -41,12 +45,12 @@ def dispatch_next_batch end end - def start_concurrency_maintenance - concurrency_maintenance&.start + def start_maintenance + maintenance&.start end - def stop_concurrency_maintenance - concurrency_maintenance&.stop + def stop_maintenance + maintenance&.stop end def all_work_completed? diff --git a/lib/solid_queue/dispatcher/concurrency_maintenance.rb b/lib/solid_queue/dispatcher/concurrency_maintenance.rb index 81cf770cc..0174a63c0 100644 --- a/lib/solid_queue/dispatcher/concurrency_maintenance.rb +++ b/lib/solid_queue/dispatcher/concurrency_maintenance.rb @@ -1,44 +1,11 @@ # frozen_string_literal: true module SolidQueue - class Dispatcher::ConcurrencyMaintenance - include AppExecutor - - attr_reader :interval, :batch_size - + # Kept for compatibility: concurrency maintenance runs on the shared + # Dispatcher::Maintenance timer, together with batch maintenance. + class Dispatcher::ConcurrencyMaintenance < Dispatcher::Maintenance def initialize(interval, batch_size) - @interval = interval - @batch_size = batch_size - end - - def start - @concurrency_maintenance_task = Concurrent::TimerTask.new(run_now: true, execution_interval: interval) do - expire_semaphores - unblock_blocked_executions - end - - @concurrency_maintenance_task.add_observer do |_, _, error| - handle_thread_error(error) if error - end - - @concurrency_maintenance_task.execute - end - - def stop - @concurrency_maintenance_task&.shutdown + super(interval, batch_size, concurrency: true, batches: false) end - - private - def expire_semaphores - wrap_in_app_executor do - Semaphore.expired.in_batches(of: batch_size, &:delete_all) - end - end - - def unblock_blocked_executions - wrap_in_app_executor do - BlockedExecution.unblock(batch_size) - end - end end end diff --git a/lib/solid_queue/dispatcher/maintenance.rb b/lib/solid_queue/dispatcher/maintenance.rb new file mode 100644 index 000000000..e0183ab3d --- /dev/null +++ b/lib/solid_queue/dispatcher/maintenance.rb @@ -0,0 +1,68 @@ +# frozen_string_literal: true + +module SolidQueue + class Dispatcher::Maintenance + include AppExecutor + + attr_reader :interval, :batch_size + + def initialize(interval, batch_size, concurrency:, batches:) + @interval = interval + @batch_size = batch_size + @concurrency = concurrency + @batches = batches + end + + def concurrency? + @concurrency + end + + def batches? + @batches + end + + def metadata + { concurrency_maintenance_interval: (interval if concurrency?), batch_maintenance: batches? } + end + + def start + @maintenance_task = Concurrent::TimerTask.new(run_now: true, execution_interval: interval) do + if concurrency? + expire_semaphores + unblock_blocked_executions + end + + sweep_stalled_batches if batches? + end + + @maintenance_task.add_observer do |_, _, error| + handle_thread_error(error) if error + end + + @maintenance_task.execute + end + + def stop + @maintenance_task&.shutdown + end + + private + def expire_semaphores + wrap_in_app_executor do + Semaphore.expired.in_batches(of: batch_size, &:delete_all) + end + end + + def unblock_blocked_executions + wrap_in_app_executor do + BlockedExecution.unblock(batch_size) + end + end + + def sweep_stalled_batches + wrap_in_app_executor do + Batch.sweep_stalled(batch_size: batch_size) + end + end + end +end diff --git a/lib/solid_queue/engine.rb b/lib/solid_queue/engine.rb index 1a7448b82..e030aa507 100644 --- a/lib/solid_queue/engine.rb +++ b/lib/solid_queue/engine.rb @@ -41,7 +41,10 @@ class Engine < ::Rails::Engine initializer "solid_queue.active_job.extensions" do ActiveSupport.on_load :active_job do include ActiveJob::ConcurrencyControls - include ActiveJob::BatchId + + ActiveSupport.on_load :active_record do + ActiveJob::Base.include ActiveJob::BatchId + end end end diff --git a/lib/solid_queue/log_subscriber.rb b/lib/solid_queue/log_subscriber.rb index 96fb19bf5..edea19672 100644 --- a/lib/solid_queue/log_subscriber.rb +++ b/lib/solid_queue/log_subscriber.rb @@ -39,6 +39,18 @@ def discard(event) debug formatted_event(event, action: "Discard job", **event.payload.slice(:job_id, :status)) end + def finish_batch(event) + info formatted_event(event, action: "Finish batch", **event.payload.slice(:batch_id, :total_jobs, :completed_jobs, :failed_jobs)) + end + + def sweep_stalled_batches(event) + debug formatted_event(event, action: "Sweep stalled batches", **event.payload.slice(:size, :started, :repaired)) + end + + def batch_progress_error(event) + error formatted_event(event, action: "Error updating batch progress", **event.payload.slice(:batch_id, :job_id), error: formatted_error(event.payload[:error])) + end + def release_many_blocked(event) debug formatted_event(event, action: "Unblock jobs", **event.payload.slice(:limit, :size)) end diff --git a/test/dummy/db/queue_schema.rb b/test/dummy/db/queue_schema.rb index 299d40494..4feed9f46 100644 --- a/test/dummy/db/queue_schema.rb +++ b/test/dummy/db/queue_schema.rb @@ -161,6 +161,8 @@ t.index ["batch_id"], name: "index_solid_queue_batch_executions_on_batch_id" end + add_foreign_key "solid_queue_batch_executions", "solid_queue_batches", column: "batch_id", on_delete: :cascade + add_foreign_key "solid_queue_batch_executions", "solid_queue_jobs", column: "job_id", on_delete: :cascade add_foreign_key "solid_queue_blocked_executions", "solid_queue_jobs", column: "job_id", on_delete: :cascade add_foreign_key "solid_queue_claimed_executions", "solid_queue_jobs", column: "job_id", on_delete: :cascade add_foreign_key "solid_queue_failed_executions", "solid_queue_jobs", column: "job_id", on_delete: :cascade diff --git a/test/integration/batch_lifecycle_test.rb b/test/integration/batch_lifecycle_test.rb index cc2836891..0899b382f 100644 --- a/test/integration/batch_lifecycle_test.rb +++ b/test/integration/batch_lifecycle_test.rb @@ -9,7 +9,8 @@ class BatchLifecycleTest < ActiveSupport::TestCase @_on_thread_error = SolidQueue.on_thread_error SolidQueue.on_thread_error = silent_on_thread_error_for([ FailingJobError ], @_on_thread_error) @worker = SolidQueue::Worker.new(queues: "background", threads: 3) - @dispatcher = SolidQueue::Dispatcher.new(batch_size: 10, polling_interval: 0.2) + # Fast maintenance so leaked tracking rows get repaired within the test windows + @dispatcher = SolidQueue::Dispatcher.new(batch_size: 10, polling_interval: 0.2, concurrency_maintenance_interval: 1) SolidQueue::Batch::EmptyJob.queue_as "background" end @@ -97,8 +98,8 @@ def perform @dispatcher.start @worker.start - wait_for_batches_to_finish_for(2.seconds) - wait_for_jobs_to_finish_for(1.second) + wait_for_batches_to_finish_for(5.seconds) + wait_for_jobs_to_finish_for(5.seconds) expected_values = [ "1: 1 jobs succeeded!", "1.1: 1 jobs succeeded!", "2: 1 jobs succeeded!", "3: 1 jobs succeeded!" ] assert_equal expected_values.sort, JobBuffer.values.sort @@ -119,7 +120,7 @@ def perform @dispatcher.start @worker.start - wait_for_batches_to_finish_for(2.seconds) + wait_for_batches_to_finish_for(5.seconds) assert_equal [ "added from inside 1", "added from inside 2", "added from inside 3", "hey", "ho" ], JobBuffer.values.sort assert_equal 3, SolidQueue::Batch.finished.count @@ -128,6 +129,32 @@ def perform assert_finished_in_order(job!(job1), batch1.reload) end + test "prebuilt jobs capture their batch before enqueue is deferred" do + skip if Rails::VERSION::MAJOR == 7 && Rails::VERSION::MINOR == 1 + + ApplicationJob.enqueue_after_transaction_commit = true + + job = AddToBufferJob.new("prebuilt") + assert_nil job.batch_id + + batch = nil + JobResult.transaction do + # Materialize the transaction so Active Job defers the adapter push. + JobResult.create!(queue_name: "default", status: "") + + batch = SolidQueue::Batch.enqueue do + job.enqueue + end + + assert_nil SolidQueue::Job.find_by(active_job_id: job.job_id) + end + + persisted_job = job!(job) + assert_equal batch.id, persisted_job.batch_id + assert_equal 1, batch.reload.total_jobs + assert_equal 1, batch.batch_executions.count + end + test "when self.enqueue_after_transaction_commit = true" do skip if Rails::VERSION::MAJOR == 7 && Rails::VERSION::MINOR == 1 @@ -293,8 +320,8 @@ def perform @dispatcher.start @worker.start - wait_for_batches_to_finish_for(2.seconds) - wait_for_jobs_to_finish_for(1.second) + wait_for_batches_to_finish_for(5.seconds) + wait_for_jobs_to_finish_for(5.seconds) assert_equal [ "Hi finish #{batch.id}!", "Hi success #{batch.id}!", "hey" ].sort, JobBuffer.values.sort assert_equal 1, batch.reload.completed_jobs diff --git a/test/models/solid_queue/batch_test.rb b/test/models/solid_queue/batch_test.rb index 4046a7aa5..b1064acf7 100644 --- a/test/models/solid_queue/batch_test.rb +++ b/test/models/solid_queue/batch_test.rb @@ -10,7 +10,7 @@ class SolidQueue::BatchTest < ActiveSupport::TestCase class BatchWithArgumentsJob < ApplicationJob def perform(arg1, arg2) - Rails.logger.info "Hi #{batch.batch_id}, #{arg1}, #{arg2}!" + Rails.logger.info "Hi #{batch.id}, #{arg1}, #{arg2}!" end end @@ -131,4 +131,367 @@ def perform(arg) assert_equal 2, batch.jobs.count end + + class OtherAdapterCallbackJob < ApplicationJob + self.queue_adapter = :test + + def perform; end + end + + test "empty job stays on solid_queue regardless of the app's default adapter" do + original = ApplicationJob.queue_adapter + ApplicationJob.queue_adapter = :test + + assert_equal "solid_queue", SolidQueue::Batch::EmptyJob.queue_adapter_name + ensure + ApplicationJob.queue_adapter = original + end + + class HookedCallbackJob < ApplicationJob + cattr_accessor :enqueue_hook_ran, default: false + + before_enqueue { self.class.enqueue_hook_ran = true } + + def perform; end + end + + class AbortingCallbackJob < ApplicationJob + before_enqueue { throw :abort } + + def perform; end + end + + test "callback jobs run their Active Job enqueue callbacks" do + HookedCallbackJob.enqueue_hook_ran = false + batch = SolidQueue::Batch.enqueue(on_finish: HookedCallbackJob) { NiceJob.perform_later("world") } + + batch.jobs.sole.finished! + + assert batch.reload.finished? + assert HookedCallbackJob.enqueue_hook_ran + assert_equal 1, SolidQueue::Job.where(class_name: HookedCallbackJob.name).count + end + + test "callback jobs honor an aborting before_enqueue without breaking completion" do + batch = SolidQueue::Batch.enqueue(on_finish: AbortingCallbackJob) { NiceJob.perform_later("world") } + + batch.jobs.sole.finished! + + assert batch.reload.finished? + assert_equal 0, SolidQueue::Job.where(class_name: AbortingCallbackJob.name).count + end + + test "callback jobs enqueue through solid_queue regardless of their class adapter" do + batch = SolidQueue::Batch.enqueue(on_finish: OtherAdapterCallbackJob) do + NiceJob.perform_later("world") + end + + batch.jobs.sole.finished! + + assert batch.reload.finished? + assert_equal 1, SolidQueue::Job.where(class_name: OtherAdapterCallbackJob.name).count + end + + test "assigns a reserved active_job_batch_id on create" do + batch = SolidQueue::Batch.enqueue { NiceJob.perform_later("world") } + + assert batch.active_job_batch_id.present? + end + + test "cannot enqueue when the batch was finished concurrently" do + batch = SolidQueue::Batch.enqueue(on_finish: BatchCompletionJob) do + NiceJob.perform_later("world") + end + + stale = SolidQueue::Batch.find(batch.id) + SolidQueue::Batch.where(id: batch.id).update_all(finished_at: Time.current) + + assert_raises(SolidQueue::Batch::AlreadyFinished) do + stale.enqueue { NiceJob.perform_later("another") } + end + end + + test "jobs enqueued inside the block join the batch even when instantiated outside" do + job = NiceJob.new("outside") + + batch = SolidQueue::Batch.enqueue do + job.enqueue + end + + assert_equal batch.id, SolidQueue::Job.find_by!(active_job_id: job.job_id).batch_id + end + + test "jobs instantiated inside the block keep its batch when enqueued outside any context" do + job = nil + batch = SolidQueue::Batch.enqueue { job = NiceJob.new("inside") } + + job.enqueue + + assert_equal batch.id, SolidQueue::Job.find_by!(active_job_id: job.job_id).batch_id + end + + test "in-flight counters do not double count failed jobs" do + batch = SolidQueue::Batch.enqueue do + 3.times { |i| NiceJob.perform_later(i) } + end + + jobs = batch.jobs.order(:id).to_a + jobs.first.failed_with(RuntimeError.new("boom")) + + batch.reload + assert_equal 3, batch.total_jobs + assert_equal 2, batch.pending_jobs + assert_equal 1, batch.failed_jobs + assert_equal 0, batch.completed_jobs + assert_equal 33.33, batch.progress_percentage + + jobs.second.finished! + + batch.reload + assert_equal 1, batch.pending_jobs + assert_equal 1, batch.completed_jobs + assert_equal 66.67, batch.progress_percentage + end + + test "start_batch completes batches whose jobs finished before the batch was started" do + batch = SolidQueue::Batch.enqueue(on_finish: BatchCompletionJob) do + NiceJob.perform_later("world") + end + + # Simulate the race where all jobs finish before enqueued_at is committed + batch.update_columns(enqueued_at: nil) + batch.jobs.sole.finished! + + assert_not batch.reload.finished? + + batch.send(:start_batch) + + assert batch.reload.finished? + end + + # Competing completion checks must have one CAS winner, one final counter + # update, and one set of callback jobs. + test "concurrent completion checks finish the batch exactly once" do + batch = SolidQueue::Batch.enqueue(on_finish: BatchCompletionJob) do + NiceJob.perform_later("world") + end + + # Leave the batch as a completion candidate that no check has picked up yet: + # tracking rows removed without firing their destroy callbacks, as happens + # with bulk discards. + SolidQueue::BatchExecution.where(batch_id: batch.id).delete_all + + concurrency = 8 + barrier = Concurrent::CyclicBarrier.new(concurrency) + threads = concurrency.times.map do + Thread.new do + SolidQueue::Record.connection_pool.with_connection do + barrier.wait + 3.times { SolidQueue::Batch.find(batch.id).check_completion } + end + end + end + threads.each(&:join) + + batch.reload + assert batch.finished? + assert_equal batch.total_jobs, batch.completed_jobs + assert_equal 1, SolidQueue::Job.where(class_name: "BatchCompletionJob").count + + batch.check_completion + assert_equal 1, SolidQueue::Job.where(class_name: "BatchCompletionJob").count + end + + # The increment must precede tracking inserts. On MySQL, FK validation takes + # a shared batch-row lock; upgrading it afterward can deadlock concurrent adders. + test "concurrent adders to the same batch keep exact accounting" do + batch = SolidQueue::Batch.enqueue { NiceJob.perform_later("seed") } + + concurrency, adds_per_thread = 8, 5 + barrier = Concurrent::CyclicBarrier.new(concurrency) + errors = Queue.new + threads = concurrency.times.map do + Thread.new do + SolidQueue::Record.connection_pool.with_connection do + barrier.wait + adds_per_thread.times do + SolidQueue::Batch.find(batch.id).enqueue { NiceJob.perform_later("added") } + rescue => e + errors << e + end + end + end + end + threads.each(&:join) + + raised = [] + raised << errors.pop until errors.empty? + assert_empty raised + + expected = 1 + concurrency * adds_per_thread + assert_equal expected, SolidQueue::Job.where(batch_id: batch.id).count + assert_equal expected, batch.reload.total_jobs + end + + # Guards the execution re-check after winning the finishing update: on + # PostgreSQL READ COMMITTED, a completion check that blocked on a concurrent + # adder's row lock re-evaluates its NOT EXISTS against the original snapshot, + # so it can win despite the adder's freshly committed executions. Without the + # re-check, this finishes a batch that still has work. + test "a completion check that races a concurrent adder does not finish the batch" do + batch = SolidQueue::Batch.enqueue(on_finish: BatchCompletionJob) { NiceJob.perform_later("world") } + job = SolidQueue::Job.where(batch_id: batch.id).sole + + # Make the batch a completion candidate, then re-add the job from a + # transaction that holds the batch row lock while completion runs. + SolidQueue::BatchExecution.where(batch_id: batch.id).delete_all + + adder_started = Queue.new + adder = Thread.new do + SolidQueue::Record.connection_pool.with_connection do + SolidQueue::Record.transaction do + # Same statements, same order as BatchExecution.create_all_from_jobs + SolidQueue::Batch.where(id: batch.id).unfinished.update_all("total_jobs = total_jobs + 1") + SolidQueue::BatchExecution.insert_all!([ { batch_id: batch.id, job_id: job.id } ]) + adder_started << true + sleep 0.5 # hold the row lock so the completion check blocks on it + end + end + end + + adder_started.pop + SolidQueue::Batch.find(batch.id).check_completion + adder.join + + assert_not batch.reload.finished? + assert_equal 1, SolidQueue::BatchExecution.where(batch_id: batch.id).count + end + + test "start_batch is single-winner: stale instances cannot restart a started batch" do + batch = SolidQueue::Batch.create!(on_finish: BatchCompletionJob) + batch.update_columns(enqueued_at: nil, total_jobs: 0) + SolidQueue::Job.where(batch_id: batch.id).destroy_all + + stale_a = SolidQueue::Batch.find(batch.id) + stale_b = SolidQueue::Batch.find(batch.id) + + stale_a.start_batch + started_at = batch.reload.enqueued_at + assert_equal 1, batch.total_jobs + + travel 1.second do + stale_b.start_batch + end + + assert_equal 1, batch.reload.total_jobs + assert_equal started_at, batch.enqueued_at + end + + test "batch capture runs before deferred enqueues" do + ancestors = ApplicationJob.ancestors + assert_includes ancestors, ActiveJob::BatchId + + if defined?(ActiveJob::EnqueueAfterTransactionCommit) + assert_operator ancestors.index(ActiveJob::BatchId), :<, ancestors.index(ActiveJob::EnqueueAfterTransactionCommit) + end + end + + test "reused job instances join the currently active batch" do + job = NiceJob.new("reused") + batch_a = SolidQueue::Batch.enqueue { job.enqueue } + batch_b = SolidQueue::Batch.enqueue { job.enqueue } + + assert_equal [ batch_a.id, batch_b.id ], + SolidQueue::Job.where(active_job_id: job.job_id).order(:id).pluck(:batch_id) + end + + test "batch accessor reflects a batch assigned after a nil read" do + job = NiceJob.new("late") + assert_nil job.batch + + batch = SolidQueue::Batch.enqueue { job.enqueue } + + assert_equal batch.id, job.batch.id + end + + # Removing a tracking row can fail mid-flight and be swallowed (e.g. SQLite + # busy inside the finishing transaction), leaving a resolved job with a live + # row and a batch that can never finish. The sweep repairs exactly that state. + test "sweep_stalled repairs tracking rows leaked by swallowed removal errors" do + batch = SolidQueue::Batch.enqueue(on_finish: BatchCompletionJob) do + 2.times { |i| NiceJob.perform_later(i) } + end + + jobs = batch.jobs.order(:id).to_a + # Simulate both leak flavors by resolving the jobs without callbacks + jobs.first.update_columns(finished_at: Time.current) + SolidQueue::FailedExecution.insert_all!([ { job_id: jobs.second.id, error: { exception_class: "RuntimeError" }.to_json } ]) + + assert_equal 2, SolidQueue::BatchExecution.where(batch_id: batch.id).count + assert_not batch.reload.finished? + + SolidQueue::Batch.sweep_stalled + + batch.reload + assert batch.finished? + assert batch.failed? + assert_equal 1, batch.failed_jobs + assert_equal 1, batch.completed_jobs + assert_equal 1, SolidQueue::Job.where(class_name: "BatchCompletionJob").count + end + + test "sweep_stalled finishes batches whose jobs were bulk discarded" do + batch = SolidQueue::Batch.enqueue do + 3.times { |i| NiceJob.perform_later(i) } + end + + # Bulk discards delete jobs without callbacks; the foreign key cascade removes + # the batch executions, and the sweep picks up the completion. + SolidQueue::ReadyExecution.discard_all_in_batches + + assert_equal 0, SolidQueue::BatchExecution.count + assert_not batch.reload.finished? + + # Age the batch past the completion grace so the sweep will consider it + batch.update_columns(enqueued_at: 5.seconds.ago) + SolidQueue::Batch.sweep_stalled + + batch.reload + assert batch.finished? + assert_equal 0, batch.pending_jobs + assert_equal 3, batch.completed_jobs + end + + test "sweep_stalled starts batches whose creating process died before starting them" do + batch = SolidQueue::Batch.enqueue { NiceJob.perform_later("world") } + + # Simulate a process that crashed after committing jobs but before start_batch + batch.update_columns(enqueued_at: nil, created_at: 10.minutes.ago) + batch.jobs.sole.finished! + + assert_not batch.reload.finished? + + SolidQueue::Batch.sweep_stalled + + assert batch.reload.finished? + end + + test "conflict-discarded jobs count the same for single and bulk enqueues" do + result1 = JobResult.create!(queue_name: "default", status: "") + batch1 = SolidQueue::Batch.enqueue do + DiscardableUpdateResultJob.perform_later(result1, name: "A") + DiscardableUpdateResultJob.perform_later(result1, name: "B") + end + + result2 = JobResult.create!(queue_name: "default", status: "") + batch2 = SolidQueue::Batch.enqueue do + ActiveJob.perform_all_later([ + DiscardableUpdateResultJob.new(result2, name: "A"), + DiscardableUpdateResultJob.new(result2, name: "B") + ]) + end + + assert_equal batch1.reload.total_jobs, batch2.reload.total_jobs + assert_equal batch1.pending_jobs, batch2.pending_jobs + end end diff --git a/test/unit/dispatcher_test.rb b/test/unit/dispatcher_test.rb index 7df0591f5..316fc622b 100644 --- a/test/unit/dispatcher_test.rb +++ b/test/unit/dispatcher_test.rb @@ -31,11 +31,34 @@ class DispatcherTest < ActiveSupport::TestCase process = SolidQueue::Process.first assert_equal "Dispatcher", process.kind - assert_metadata process, polling_interval: 0.1, batch_size: 10 + assert_metadata process, polling_interval: 0.1, batch_size: 10, batch_maintenance: true + assert_nil process.metadata["concurrency_maintenance_interval"] ensure no_concurrency_maintenance_dispatcher.stop end + test "batch maintenance is optional" do + no_batch_maintenance_dispatcher = SolidQueue::Dispatcher.new(polling_interval: 0.1, batch_size: 10, batch_maintenance: false) + no_batch_maintenance_dispatcher.start + + wait_for_registered_processes(1, timeout: 1.second) + + process = SolidQueue::Process.first + assert_equal "Dispatcher", process.kind + assert_metadata process, concurrency_maintenance_interval: 600, batch_maintenance: false + ensure + no_batch_maintenance_dispatcher.stop + end + + test "ConcurrencyMaintenance remains constructible with its original signature" do + maintenance = SolidQueue::Dispatcher::ConcurrencyMaintenance.new(600, 100) + + assert_equal 600, maintenance.interval + assert_equal 100, maintenance.batch_size + assert maintenance.concurrency? + assert_not maintenance.batches? + end + test "polling queries are logged" do log = StringIO.new with_active_record_logger(ActiveSupport::Logger.new(log)) do