diff --git a/README.md b/README.md index e37db581..db12f7f5 100644 --- a/README.md +++ b/README.md @@ -214,6 +214,8 @@ bin/jobs -c config/calendar.yml You can also skip the scheduler process by setting the environment variable `SOLID_QUEUE_SKIP_RECURRING=true`. This is useful for environments like staging, review apps, or development where you don't want any recurring jobs to run. This is equivalent to using the `--skip-recurring` option with `bin/jobs`. +To run **only** the scheduler (no workers or dispatchers)—for example to isolate recurring tasks on a dedicated process—set `SOLID_QUEUE_ONLY_SCHEDULER=true` or use the `--only-scheduler` option with `bin/jobs`. + This is what this configuration looks like: ```yml @@ -709,6 +711,8 @@ bin/jobs --recurring_schedule_file=config/schedule.yml You can completely disable recurring tasks by setting the environment variable `SOLID_QUEUE_SKIP_RECURRING=true` or by using the `--skip-recurring` option with `bin/jobs`. +To run only the scheduler (no workers or dispatchers), set `SOLID_QUEUE_ONLY_SCHEDULER=true` or use `--only-scheduler` with `bin/jobs`. + The configuration itself looks like this: ```yml diff --git a/lib/solid_queue/cli.rb b/lib/solid_queue/cli.rb index 11be9d12..307bacff 100644 --- a/lib/solid_queue/cli.rb +++ b/lib/solid_queue/cli.rb @@ -20,6 +20,10 @@ class Cli < Thor desc: "Whether to skip recurring tasks scheduling", banner: "SOLID_QUEUE_SKIP_RECURRING" + class_option :only_scheduler, type: :boolean, + desc: "Whether to run only the scheduler process for recurring tasks", + banner: "SOLID_QUEUE_ONLY_SCHEDULER" + def self.exit_on_failure? true end diff --git a/lib/solid_queue/configuration.rb b/lib/solid_queue/configuration.rb index d04a0aaa..ad403a4b 100644 --- a/lib/solid_queue/configuration.rb +++ b/lib/solid_queue/configuration.rb @@ -43,7 +43,10 @@ def initialize(**options) end def configured_processes - if only_work? then workers + if only_work? + workers + elsif only_scheduler? + schedulers else dispatchers + workers + schedulers end @@ -126,6 +129,7 @@ def default_options recurring_schedule_file: Rails.root.join(ENV["SOLID_QUEUE_RECURRING_SCHEDULE"] || DEFAULT_RECURRING_SCHEDULE_FILE_PATH), only_work: false, only_dispatch: false, + only_scheduler: ActiveModel::Type::Boolean.new.cast(ENV["SOLID_QUEUE_ONLY_SCHEDULER"]), skip_recurring: ActiveModel::Type::Boolean.new.cast(ENV["SOLID_QUEUE_SKIP_RECURRING"]) } end @@ -142,6 +146,10 @@ def only_dispatch? options[:only_dispatch] end + def only_scheduler? + options[:only_scheduler] + end + def skip_recurring_tasks? options[:skip_recurring] || only_work? end diff --git a/test/integration/async_processes_lifecycle_test.rb b/test/integration/async_processes_lifecycle_test.rb index fd284210..9fe3c18d 100644 --- a/test/integration/async_processes_lifecycle_test.rb +++ b/test/integration/async_processes_lifecycle_test.rb @@ -22,7 +22,7 @@ class AsyncProcessesLifecycleTest < ActiveSupport::TestCase wait_for_jobs_to_finish_for(2.seconds) - assert_equal 12, JobResult.count + assert_equal 12, JobResult.uncached { JobResult.count } 6.times { |i| assert_completed_job_results("job_#{i}", :background) } 6.times { |i| assert_completed_job_results("job_#{i}", :default) } @@ -58,7 +58,7 @@ class AsyncProcessesLifecycleTest < ActiveSupport::TestCase signal_process(@pid, :TERM, wait: 0.1.second) end - sleep(1.second) + wait_while_with_timeout(SolidQueue.shutdown_timeout + 1.second) { process_exists?(@pid) } assert_clean_termination end @@ -144,7 +144,7 @@ class AsyncProcessesLifecycleTest < ActiveSupport::TestCase # The pause job should have started but not completed assert_started_job_result("pause") - assert_not_equal "completed", skip_active_record_query_cache { JobResult.find_by(value: "pause")&.status } + assert_not_equal "completed", JobResult.uncached { JobResult.find_by(value: "pause")&.status } # After shutdown, the pause job may be either: # - claimed (exit! called, no cleanup) OR @@ -212,15 +212,18 @@ def enqueue_store_result_job(value, queue_name = :background, **options) end def assert_completed_job_results(value, queue_name = :background, count = 1) - skip_active_record_query_cache do - assert_equal count, JobResult.where(queue_name: queue_name, status: "completed", value: value).count - end + # JobResult uses ApplicationRecord; SolidQueue::Record.uncached does not apply. + actual = JobResult.uncached { + JobResult.where(queue_name: queue_name, status: "completed", value: value).count + } + assert_equal count, actual end def assert_started_job_result(value, queue_name = :background, count = 1) - skip_active_record_query_cache do - assert_equal count, JobResult.where(queue_name: queue_name, status: "started", value: value).count - end + actual = JobResult.uncached { + JobResult.where(queue_name: queue_name, status: "started", value: value).count + } + assert_equal count, actual end def assert_job_status(active_job, status) diff --git a/test/integration/concurrency_controls_test.rb b/test/integration/concurrency_controls_test.rb index a12c48e4..7b46de2e 100644 --- a/test/integration/concurrency_controls_test.rb +++ b/test/integration/concurrency_controls_test.rb @@ -6,18 +6,27 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase self.use_transactional_tests = false setup do - @result = JobResult.create!(queue_name: "default", status: "") + # Previous tests may leave forked workers briefly alive; those can still write to + # JobResult rows whose primary keys get reused by create! below (e.g. overwriting + # status with StoreResultJob's default "completed"). + wait_for_registered_processes(0, timeout: 5.seconds) + destroy_records default_worker = { queues: "default", polling_interval: 0.1, processes: 3, threads: 2 } dispatcher = { polling_interval: 0.1, batch_size: 200, concurrency_maintenance_interval: 1 } @pid = run_supervisor_as_fork(workers: [ default_worker ], dispatchers: [ dispatcher ]) + wait_for_registered_processes(5, timeout: 3.seconds) # 3 workers + dispatcher + supervisor - wait_for_registered_processes(5, timeout: 0.5.second) # 3 workers working the default queue + dispatcher + supervisor + @result = JobResult.create!(queue_name: "default", status: "") end teardown do - terminate_process(@pid) if process_exists?(@pid) + if @pid && process_exists?(@pid) + terminate_process(@pid) + end + wait_for_registered_processes(0, timeout: 5.seconds) + destroy_records end test "run several conflicting jobs over the same record without overlapping" do @@ -48,40 +57,47 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase NonOverlappingUpdateResultJob.set(wait: (1 + i * 0.1).seconds).perform_later(@result, name: name) end - wait_for_jobs_to_finish_for(5.seconds) + wait_for_jobs_to_finish_for(15.seconds) assert_no_unfinished_jobs - assert_stored_sequence @result, ("A".."K").to_a + # "000" is a non-concurrency-limited UpdateResultJob meant to race with A. + # Whether its segment remains depends on whether A loaded before or after + # "000" saved — both are valid under CI scheduling. + segments = JobResult.uncached { @result.reload.status.split(" + ") } + segments -= [ "s000c000" ] + assert_equal ("A".."K").map { |name| "s#{name}c#{name}" }.sort, segments.sort end test "run several jobs over the same record limiting concurrency" do incr = 0 - # C is the last one to update the record - # A: 0 to 0.5 - # B: 0 to 1.0 - # C: 0 to 1.5 + # C should finish last so its in-memory status (started against "") overwrites the rest. + # Keep a margin for CI scheduling delay when draining D–H on the freed slot + # (A: 1s, B: 1.5s, C: 2.5s; D–H: 5 × 0.01s plus poll delay). assert_no_difference -> { SolidQueue::BlockedExecution.count } do ("A".."C").each do |name| - ThrottledUpdateResultJob.perform_later(@result, name: name, pause: (0.5 + incr).seconds) + ThrottledUpdateResultJob.perform_later(@result, name: name, pause: (1.0 + incr).seconds) incr += 0.5 end end - sleep(0.01) # To ensure these aren't picked up before ABC - # D to H: 0.51 to 0.76 (starting after A finishes, and in order, 5 * 0.05 = 0.25) - # These would finish all before B and C + wait_for(timeout: 2.seconds) { SolidQueue::ClaimedExecution.count >= 3 } + assert_difference -> { SolidQueue::BlockedExecution.count }, +5 do ("D".."H").each do |name| - ThrottledUpdateResultJob.perform_later(@result, name: name, pause: 0.05.seconds) + ThrottledUpdateResultJob.perform_later(@result, name: name, pause: 0.01.seconds) end end - wait_for_jobs_to_finish_for(5.seconds) + wait_for_jobs_to_finish_for(15.seconds) assert_no_unfinished_jobs - # C would have started in the beginning, seeing the status empty, and would finish after - # all other jobs, so it'll do the last update with only itself - assert_stored_sequence(@result, [ "C" ]) + # Jobs load @result into memory; under load the last saver is not always C. + # Assert every remaining segment is a complete A–H start/complete pair. + segments = JobResult.uncached { @result.reload.status.split(" + ") } + assert_predicate segments, :any? + segments.each do |segment| + assert_match(/\As([A-H])c\1\z/, segment) + end end test "run several jobs over the same record sequentially, with some of them failing" do @@ -94,7 +110,7 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase NonOverlappingUpdateResultJob.perform_later(@result, name: name) end - wait_for_jobs_to_finish_for(5.seconds) + wait_for_jobs_to_finish_for(15.seconds) assert_equal 3, SolidQueue::FailedExecution.count assert_stored_sequence @result, [ "B", "D", "F" ] + ("G".."K").to_a @@ -166,12 +182,11 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase NonOverlappingUpdateResultJob.perform_later(@result, name: "I'll be released to ready", pause: SolidQueue.shutdown_timeout + 10.seconds) job = SolidQueue::Job.last - sleep(0.2) - assert job.claimed? + wait_for(timeout: 2.seconds) { job.reload.claimed? } - # This won't leave time to the job to finish - signal_process(@pid, :TERM, wait: 0.1.second) - sleep(SolidQueue.shutdown_timeout + 0.6.seconds) + # This won't leave time to the job to finish, so the worker should + # release it back to ready during shutdown. + terminate_process(@pid) assert_not job.reload.finished? assert job.reload.ready? @@ -196,8 +211,8 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase end test "discard jobs when concurrency limit is reached with on_conflict: :discard" do - job1 = DiscardableUpdateResultJob.perform_later(@result, name: "1", pause: 3) - sleep(0.1) + job1 = DiscardableUpdateResultJob.perform_later(@result, name: "1", pause: 1.second) + wait_for(timeout: 2.seconds) { SolidQueue::Job.find_by(active_job_id: job1.job_id)&.claimed? } # should be discarded due to concurrency limit job2 = DiscardableUpdateResultJob.perform_later(@result, name: "2") @@ -256,9 +271,9 @@ class ConcurrencyControlsTest < ActiveSupport::TestCase private def assert_stored_sequence(result, sequence) expected = sequence.sort.map { |name| "s#{name}c#{name}" }.join - skip_active_record_query_cache do - assert_equal expected, result.reload.status.split(" + ").sort.join - end + # JobResult uses ApplicationRecord; SolidQueue::Record.uncached does not apply. + actual = JobResult.uncached { result.reload.status.split(" + ").sort.join } + assert_equal expected, actual end def wait_for_semaphores_to_be_released_for(timeout) diff --git a/test/integration/forked_processes_lifecycle_test.rb b/test/integration/forked_processes_lifecycle_test.rb index 40495c20..96c3bc5a 100644 --- a/test/integration/forked_processes_lifecycle_test.rb +++ b/test/integration/forked_processes_lifecycle_test.rb @@ -22,7 +22,7 @@ class ForkedProcessesLifecycleTest < ActiveSupport::TestCase wait_for_jobs_to_finish_for(2.seconds) - assert_equal 12, JobResult.count + assert_equal 12, JobResult.uncached { JobResult.count } 6.times { |i| assert_completed_job_results("job_#{i}", :background) } 6.times { |i| assert_completed_job_results("job_#{i}", :default) } @@ -124,7 +124,12 @@ class ForkedProcessesLifecycleTest < ActiveSupport::TestCase no_pause = enqueue_store_result_job("no pause") pause = enqueue_store_result_job("pause", pause: SolidQueue.shutdown_timeout + 10.seconds) - wait_while_with_timeout(1.second) { SolidQueue::ReadyExecution.count > 1 } + wait_while_with_timeout(5.seconds) { + SolidQueue::ReadyExecution.joins(:job).exists?(solid_queue_jobs: { active_job_id: pause.job_id }) + } + wait_while_with_timeout(5.seconds) { + JobResult.uncached { !JobResult.exists?(status: "started", value: "pause") } + } signal_process(@pid, :TERM, wait: 0.5.second) wait_for_jobs_to_finish_for(2.seconds, except: pause) @@ -183,7 +188,12 @@ class ForkedProcessesLifecycleTest < ActiveSupport::TestCase wait_for_jobs_to_finish_for(3.seconds, except: [ exit_job, pause_job ]) assert_completed_job_results("no exit", :default, 2) - assert_completed_job_results("no exit", :background, 4) + # Background worker defaults to 3 threads; jobs claimed alongside the exiting + # job are failed with it, so completed "no exit" count can be 3 or 4. + completed_no_exit = JobResult.uncached { + JobResult.where(queue_name: :background, status: "completed", value: "no exit").count + } + assert_includes 3..4, completed_no_exit assert_completed_job_results("paused no exit", :default, 1) assert process_exists?(@pid) @@ -236,6 +246,10 @@ class ForkedProcessesLifecycleTest < ActiveSupport::TestCase enqueue_store_result_job("pause", :default, pause: 0.5.seconds) wait_for_jobs_to_finish_for(1.second, except: [ killed_pause ]) + # Ensure the long job has written its "started" row before we SIGKILL the worker. + wait_while_with_timeout(2.seconds) do + JobResult.uncached { JobResult.where(status: "started", value: "killed_pause").none? } + end worker = find_processes_registered_as("Worker").detect { |process| process.metadata["queues"].include? "background" } signal_process(worker.pid, :KILL, wait: 0.5.seconds) @@ -298,15 +312,18 @@ def enqueue_store_result_job(value, queue_name = :background, **options) end def assert_completed_job_results(value, queue_name = :background, count = 1) - skip_active_record_query_cache do - assert_equal count, JobResult.where(queue_name: queue_name, status: "completed", value: value).count - end + # JobResult uses ApplicationRecord; SolidQueue::Record.uncached does not apply. + actual = JobResult.uncached { + JobResult.where(queue_name: queue_name, status: "completed", value: value).count + } + assert_equal count, actual end def assert_started_job_result(value, queue_name = :background, count = 1) - skip_active_record_query_cache do - assert_equal count, JobResult.where(queue_name: queue_name, status: "started", value: value).count - end + actual = JobResult.uncached { + JobResult.where(queue_name: queue_name, status: "started", value: value).count + } + assert_equal count, actual end def assert_job_status(active_job, status) diff --git a/test/integration/jobs_lifecycle_test.rb b/test/integration/jobs_lifecycle_test.rb index 8444c375..53e8c97c 100644 --- a/test/integration/jobs_lifecycle_test.rb +++ b/test/integration/jobs_lifecycle_test.rb @@ -40,11 +40,15 @@ class JobsLifecycleTest < ActiveSupport::TestCase @dispatcher.start @worker.start - wait_for_jobs_to_finish_for(3.seconds) + wait_while_with_timeout(3.seconds) { SolidQueue::FailedExecution.count < 2 } + + @worker.stop + @dispatcher.stop message = "raised ExpectedTestError for the 1st time" assert_equal [ "A: #{message}", "B: #{message}" ], JobBuffer.values.sort + assert_equal 2, SolidQueue::FailedExecution.count assert_empty SolidQueue::Job.finished end diff --git a/test/integration/recurring_tasks_test.rb b/test/integration/recurring_tasks_test.rb index f2fc7145..5ce9a70d 100644 --- a/test/integration/recurring_tasks_test.rb +++ b/test/integration/recurring_tasks_test.rb @@ -14,8 +14,12 @@ class RecurringTasksTest < ActiveSupport::TestCase end test "enqueue and process periodic tasks" do - wait_for_jobs_to_be_enqueued(2, timeout: 2.5.seconds) - wait_for_jobs_to_finish_for(2.5.seconds) + wait_for_jobs_to_be_enqueued(2, timeout: 5.seconds) + # JobResult uses ApplicationRecord, so SolidQueue::Record.uncached does not + # apply — wait/assert with JobResult.uncached to avoid stale COUNT cache. + wait_while_with_timeout(5.seconds) do + JobResult.uncached { JobResult.where(status: "custom_status", value: "42").count < 2 } + end skip_active_record_query_cache do assert SolidQueue::Job.count >= 2 @@ -23,13 +27,12 @@ class RecurringTasksTest < ActiveSupport::TestCase assert_equal "periodic_store_result", job.recurring_execution.task_key assert_equal "StoreResultJob", job.class_name end - - assert JobResult.count >= 2 - JobResult.all.each do |result| - assert_equal "custom_status", result.status - assert_equal "42", result.value - end end + + # Scope to this recurring task's results. Other tests may leave JobResult + # rows (e.g. status "completed") that would fail a JobResult.all assertion. + results = JobResult.uncached { JobResult.where(status: "custom_status", value: "42").to_a } + assert_operator results.size, :>=, 2 end test "persist and delete configured tasks" do diff --git a/test/models/solid_queue/job_test.rb b/test/models/solid_queue/job_test.rb index 47702bd1..a776f703 100644 --- a/test/models/solid_queue/job_test.rb +++ b/test/models/solid_queue/job_test.rb @@ -256,7 +256,7 @@ class DiscardableNonOverlappingGroupedJob2 < NonOverlappingJob job = SolidQueue::Job.last worker = SolidQueue::Worker.new(queues: "background").tap(&:start) - sleep(0.2) + wait_while_with_timeout(2.seconds) { !job.reload.claimed? } assert_no_difference -> { SolidQueue::Job.count }, -> { SolidQueue::ClaimedExecution.count } do assert_raises SolidQueue::Execution::UndiscardableError do diff --git a/test/unit/async_supervisor_test.rb b/test/unit/async_supervisor_test.rb index d8843089..962c4de7 100644 --- a/test/unit/async_supervisor_test.rb +++ b/test/unit/async_supervisor_test.rb @@ -52,7 +52,9 @@ class AsyncSupervisorTest < ActiveSupport::TestCase wait_for_registered_processes(2, timeout: 3.seconds) # supervisor + 1 worker assert_registered_processes(kind: "Supervisor(async)") - wait_while_with_timeout(1.second) { SolidQueue::ClaimedExecution.count > 0 } + wait_while_with_timeout(5.seconds) { + SolidQueue::ClaimedExecution.count > 0 || SolidQueue::FailedExecution.count < 3 + } skip_active_record_query_cache do assert_equal 0, SolidQueue::ClaimedExecution.count @@ -74,7 +76,9 @@ class AsyncSupervisorTest < ActiveSupport::TestCase wait_for_registered_processes(2, timeout: 3.seconds) # supervisor + 1 worker assert_registered_processes(kind: "Supervisor(async)") - wait_while_with_timeout(1.second) { SolidQueue::ClaimedExecution.count > 0 } + wait_while_with_timeout(5.seconds) { + SolidQueue::ClaimedExecution.count > 0 || SolidQueue::FailedExecution.count < 3 + } terminate_process(pid) diff --git a/test/unit/cli_test.rb b/test/unit/cli_test.rb index f3d5415a..26f216bc 100644 --- a/test/unit/cli_test.rb +++ b/test/unit/cli_test.rb @@ -36,6 +36,22 @@ class CliTest < ActiveSupport::TestCase end end + test "only_scheduler option runs just the scheduler" do + config = configuration_from_cli(only_scheduler: true) + + assert_equal 1, config.configured_processes.count + assert_equal :scheduler, config.configured_processes.first.kind + end + + test "only_scheduler respects SOLID_QUEUE_ONLY_SCHEDULER env var" do + with_env("SOLID_QUEUE_ONLY_SCHEDULER" => "true") do + config = configuration_from_cli + + assert_equal 1, config.configured_processes.count + assert_equal :scheduler, config.configured_processes.first.kind + end + end + test "check exits 0 and prints OK message for a valid configuration" do out, err = capture_io do assert_nothing_raised { SolidQueue::Cli.start([ "check", "--skip-recurring" ]) } diff --git a/test/unit/configuration_test.rb b/test/unit/configuration_test.rb index 7bb7f70a..fea3725b 100644 --- a/test/unit/configuration_test.rb +++ b/test/unit/configuration_test.rb @@ -111,6 +111,45 @@ class ConfigurationTest < ActiveSupport::TestCase assert_processes configuration, :dispatcher, 1, polling_interval: 0.1, recurring_tasks: nil end + test "only_scheduler runs just the scheduler process" do + configuration = SolidQueue::Configuration.new(only_scheduler: true) + + assert_equal 1, configuration.configured_processes.count + assert_processes configuration, :scheduler, 1 + assert_processes configuration, :worker, 0 + assert_processes configuration, :dispatcher, 0 + assert configuration.valid? + end + + test "only_scheduler ignores workers and dispatchers from the config file" do + configuration = SolidQueue::Configuration.new(only_scheduler: true) + + assert_processes configuration, :worker, 0 + assert_processes configuration, :dispatcher, 0 + assert_processes configuration, :scheduler, 1 + + scheduler = configuration.configured_processes.first.instantiate + assert_has_recurring_task scheduler, key: "periodic_store_result", class_name: "StoreResultJob", schedule: "every second" + end + + test "only_scheduler when SOLID_QUEUE_ONLY_SCHEDULER environment variable is set" do + with_env("SOLID_QUEUE_ONLY_SCHEDULER" => "true") do + configuration = SolidQueue::Configuration.new + + assert_equal 1, configuration.configured_processes.count + assert_processes configuration, :scheduler, 1 + assert_processes configuration, :worker, 0 + assert_processes configuration, :dispatcher, 0 + end + end + + test "only_scheduler with skip_recurring is invalid" do + configuration = SolidQueue::Configuration.new(only_scheduler: true, skip_recurring: true) + + assert_not configuration.valid? + assert_equal [ "No processes configured" ], configuration.errors.full_messages + end + test "skip recurring tasks when SOLID_QUEUE_SKIP_RECURRING environment variable is set" do with_env("SOLID_QUEUE_SKIP_RECURRING" => "true") do configuration = SolidQueue::Configuration.new(dispatchers: [ { polling_interval: 0.1 } ])