From 525de47e69ae000b483b400441eda6512cdb13e4 Mon Sep 17 00:00:00 2001 From: Rosa Gutierrez Date: Wed, 15 Feb 2023 10:11:01 +0100 Subject: [PATCH 1/4] Move everything that should live under lib/ there, from app/models/ --- lib/solid_queue.rb | 6 +++ .../solid_queue/configuration.rb | 4 -- {app/models => lib}/solid_queue/dispatcher.rb | 0 lib/solid_queue/manager.rb | 37 +++++++++++++++++++ .../concerns => lib}/solid_queue/runnable.rb | 0 {app/models => lib}/solid_queue/scheduler.rb | 2 + test/models/solid_queue/.keep | 0 .../configuration_test.rb | 2 +- 8 files changed, 46 insertions(+), 5 deletions(-) rename {app/models => lib}/solid_queue/configuration.rb (97%) rename {app/models => lib}/solid_queue/dispatcher.rb (100%) create mode 100644 lib/solid_queue/manager.rb rename {app/models/concerns => lib}/solid_queue/runnable.rb (100%) rename {app/models => lib}/solid_queue/scheduler.rb (95%) create mode 100644 test/models/solid_queue/.keep rename test/{models/solid_queue => unit}/configuration_test.rb (91%) diff --git a/lib/solid_queue.rb b/lib/solid_queue.rb index 014ced5ed..b4573d56c 100644 --- a/lib/solid_queue.rb +++ b/lib/solid_queue.rb @@ -3,5 +3,11 @@ require "active_job/queue_adapters/solid_queue_adapter" +require "solid_queue/configuration" +require "solid_queue/runnable" +require "solid_queue/dispatcher" +require "solid_queue/scheduler" +require "solid_queue/manager" + module SolidQueue end diff --git a/app/models/solid_queue/configuration.rb b/lib/solid_queue/configuration.rb similarity index 97% rename from app/models/solid_queue/configuration.rb rename to lib/solid_queue/configuration.rb index 4cd4a7a24..aaf21d559 100644 --- a/app/models/solid_queue/configuration.rb +++ b/lib/solid_queue/configuration.rb @@ -29,10 +29,6 @@ def scheduler_disabled? raw_config.dig(:scheduler, :disabled) end - def each_queue(&block) - queues.each(&block) - end - def scheduler_options (raw_config[:scheduler] || {}).with_defaults(SCHEDULER_DEFAULTS) end diff --git a/app/models/solid_queue/dispatcher.rb b/lib/solid_queue/dispatcher.rb similarity index 100% rename from app/models/solid_queue/dispatcher.rb rename to lib/solid_queue/dispatcher.rb diff --git a/lib/solid_queue/manager.rb b/lib/solid_queue/manager.rb new file mode 100644 index 000000000..3afd24edd --- /dev/null +++ b/lib/solid_queue/manager.rb @@ -0,0 +1,37 @@ +# frozen_string_literal: true + +class SolidQueue::Manager + include SolidQueue::Runnable + + attr_accessor :dispatchers, :scheduler + + def self.start + configuration = SolidQueue::Configuration.new + dispatchers = configuration.queues.map { |queue| SolidQueue::Dispatcher.new(queue) } + scheduler = unless configuration.scheduler_disabled? + SolidQueue::Scheduler.new(configuration.scheduler_options) + end + + new(dispatchers, scheduler).start + end + + def initialize(dispatchers, scheduler = nil) + @dispatchers = dispatchers + @scheduler = scheduler + end + + def start + trap_signals + + dispatchers.each(&:start) + scheduler&.start + + Kernel.loop do + sleep 0.1 + break if stopping? + end + + dispatchers.each(&:stop) + scheduler&.stop + end +end diff --git a/app/models/concerns/solid_queue/runnable.rb b/lib/solid_queue/runnable.rb similarity index 100% rename from app/models/concerns/solid_queue/runnable.rb rename to lib/solid_queue/runnable.rb diff --git a/app/models/solid_queue/scheduler.rb b/lib/solid_queue/scheduler.rb similarity index 95% rename from app/models/solid_queue/scheduler.rb rename to lib/solid_queue/scheduler.rb index 0c4eedd6d..5df49e0e9 100644 --- a/app/models/solid_queue/scheduler.rb +++ b/lib/solid_queue/scheduler.rb @@ -1,3 +1,5 @@ +# frozen_string_literal: true + class SolidQueue::Scheduler include SolidQueue::Runnable diff --git a/test/models/solid_queue/.keep b/test/models/solid_queue/.keep new file mode 100644 index 000000000..e69de29bb diff --git a/test/models/solid_queue/configuration_test.rb b/test/unit/configuration_test.rb similarity index 91% rename from test/models/solid_queue/configuration_test.rb rename to test/unit/configuration_test.rb index c20443fec..e3b7ff57b 100644 --- a/test/models/solid_queue/configuration_test.rb +++ b/test/unit/configuration_test.rb @@ -1,6 +1,6 @@ require "test_helper" -class SolidQueue::ConfigurationTest < ActiveSupport::TestCase +class ConfigurationTest < ActiveSupport::TestCase test "read configuration from default file" do configuration = SolidQueue::Configuration.new assert 2, configuration.queues.count From 5163064d108abe73050e7f253778a90e5e7a0332 Mon Sep 17 00:00:00 2001 From: Rosa Gutierrez Date: Thu, 16 Feb 2023 19:20:56 +0100 Subject: [PATCH 2/4] Add a rake task and a supervisor to start dispatchers and scheduler --- app/models/solid_queue/job.rb | 4 +++- lib/solid_queue.rb | 3 ++- lib/solid_queue/dispatcher.rb | 7 ++++++- lib/solid_queue/engine.rb | 10 ++++++++++ lib/solid_queue/runnable.rb | 14 +++++++++++++- lib/solid_queue/scheduler.rb | 7 ++++++- lib/solid_queue/{manager.rb => supervisor.rb} | 7 +++---- lib/solid_queue/tasks.rb | 6 ++++++ solid_queue.gemspec | 4 ++-- 9 files changed, 51 insertions(+), 11 deletions(-) rename lib/solid_queue/{manager.rb => supervisor.rb} (74%) create mode 100644 lib/solid_queue/tasks.rb diff --git a/app/models/solid_queue/job.rb b/app/models/solid_queue/job.rb index f9386db16..3f52279da 100644 --- a/app/models/solid_queue/job.rb +++ b/app/models/solid_queue/job.rb @@ -14,7 +14,9 @@ def enqueue_active_job(active_job, scheduled_at: Time.current) end def enqueue(**kwargs) - create!(**kwargs.compact.with_defaults(defaults)) + create!(**kwargs.compact.with_defaults(defaults)).tap do + SolidQueue.logger.info "[SolidQueue] Enqueued job #{kwargs}" + end end private diff --git a/lib/solid_queue.rb b/lib/solid_queue.rb index b4573d56c..261141165 100644 --- a/lib/solid_queue.rb +++ b/lib/solid_queue.rb @@ -7,7 +7,8 @@ require "solid_queue/runnable" require "solid_queue/dispatcher" require "solid_queue/scheduler" -require "solid_queue/manager" +require "solid_queue/supervisor" module SolidQueue + mattr_accessor :logger, default: ActiveSupport::Logger.new($stdout) end diff --git a/lib/solid_queue/dispatcher.rb b/lib/solid_queue/dispatcher.rb index 19698d00e..8b16ace83 100644 --- a/lib/solid_queue/dispatcher.rb +++ b/lib/solid_queue/dispatcher.rb @@ -21,6 +21,11 @@ def stop super end + def inspect + "Dispatcher(queue=#{queue}, worker_count=#{worker_count}, polling_interval=#{polling_interval})" + end + alias to_s inspect + private def run loop do @@ -33,7 +38,7 @@ def run workers_pool.post { job.perform } end else - sleep(polling_interval) + interruptable_sleep(polling_interval) end end end diff --git a/lib/solid_queue/engine.rb b/lib/solid_queue/engine.rb index 84e68c32f..a32d8315f 100644 --- a/lib/solid_queue/engine.rb +++ b/lib/solid_queue/engine.rb @@ -1,5 +1,15 @@ module SolidQueue class Engine < ::Rails::Engine isolate_namespace SolidQueue + + rake_tasks do + load "solid_queue/tasks.rb" + end + + initializer "solid_queue.logger" do |app| + ActiveSupport.on_load(:solid_queue) do + self.logger = ::Rails.logger + end + end end end diff --git a/lib/solid_queue/runnable.rb b/lib/solid_queue/runnable.rb index 7f6c92d58..3c6680717 100644 --- a/lib/solid_queue/runnable.rb +++ b/lib/solid_queue/runnable.rb @@ -2,9 +2,10 @@ module SolidQueue::Runnable def start - trap_signals @stopping = false @thread = Thread.new { run } + + log "Started #{self}" end def stop @@ -29,4 +30,15 @@ def stopping? def wait @thread&.join end + + def interruptable_sleep(seconds) + while !stopping? && seconds > 0 + Kernel.sleep 0.1 + seconds -= 0.1 + end + end + + def log(message) + SolidQueue.logger.info("[SolidQueue] #{message}") + end end diff --git a/lib/solid_queue/scheduler.rb b/lib/solid_queue/scheduler.rb index 5df49e0e9..0c68c54f3 100644 --- a/lib/solid_queue/scheduler.rb +++ b/lib/solid_queue/scheduler.rb @@ -12,6 +12,11 @@ def initialize(**options) @polling_interval = options[:polling_interval] end + def inspect + "Scheduler(batch_size=#{batch_size}, polling_interval=#{polling_interval})" + end + alias to_s inspect + private def run loop do @@ -22,7 +27,7 @@ def run if batch.size > 0 SolidQueue::ScheduledExecution.prepare_batch(batch) else - sleep(polling_interval) + interruptable_sleep(polling_interval) end end end diff --git a/lib/solid_queue/manager.rb b/lib/solid_queue/supervisor.rb similarity index 74% rename from lib/solid_queue/manager.rb rename to lib/solid_queue/supervisor.rb index 3afd24edd..afca725eb 100644 --- a/lib/solid_queue/manager.rb +++ b/lib/solid_queue/supervisor.rb @@ -1,15 +1,15 @@ # frozen_string_literal: true -class SolidQueue::Manager +class SolidQueue::Supervisor include SolidQueue::Runnable attr_accessor :dispatchers, :scheduler def self.start configuration = SolidQueue::Configuration.new - dispatchers = configuration.queues.map { |queue| SolidQueue::Dispatcher.new(queue) } + dispatchers = configuration.queues.values.map { |queue_options| SolidQueue::Dispatcher.new(**queue_options) } scheduler = unless configuration.scheduler_disabled? - SolidQueue::Scheduler.new(configuration.scheduler_options) + SolidQueue::Scheduler.new(**configuration.scheduler_options) end new(dispatchers, scheduler).start @@ -22,7 +22,6 @@ def initialize(dispatchers, scheduler = nil) def start trap_signals - dispatchers.each(&:start) scheduler&.start diff --git a/lib/solid_queue/tasks.rb b/lib/solid_queue/tasks.rb new file mode 100644 index 000000000..b1c8c08d8 --- /dev/null +++ b/lib/solid_queue/tasks.rb @@ -0,0 +1,6 @@ +namespace :solid_queue do + desc "start solid_queue supervisor" + task start: :environment do + SolidQueue::Supervisor.start + end +end diff --git a/solid_queue.gemspec b/solid_queue.gemspec index 6dab75c91..e84b25a44 100644 --- a/solid_queue.gemspec +++ b/solid_queue.gemspec @@ -3,8 +3,8 @@ require_relative "lib/solid_queue/version" Gem::Specification.new do |spec| spec.name = "solid_queue" spec.version = SolidQueue::VERSION - spec.authors = ["Rosa Gutierrez"] - spec.email = ["rosa@37signals.com"] + spec.authors = [ "Rosa Gutierrez" ] + spec.email = [ "rosa@37signals.com" ] spec.homepage = "https://github.com/basecamp/solid_queue" spec.summary = "Database-backed Active Job backend." spec.description = "Database-backed Active Job backend." From 004a304e6225a946fc5ac5f7816f823396c15e76 Mon Sep 17 00:00:00 2001 From: Rosa Gutierrez Date: Fri, 17 Feb 2023 11:13:48 +0100 Subject: [PATCH 3/4] Ensure priority is assumed from job when creating executions Since priority takes a default value of 0, we weren't updating it in any case when the execution was being created from a job with a different priority, because it already had a value. --- app/models/solid_queue/execution.rb | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/app/models/solid_queue/execution.rb b/app/models/solid_queue/execution.rb index 268405419..cfa8b5d30 100644 --- a/app/models/solid_queue/execution.rb +++ b/app/models/solid_queue/execution.rb @@ -6,7 +6,6 @@ class SolidQueue::Execution < ActiveRecord::Base private def assume_attributes_from_job self.queue_name ||= job&.queue_name - self.priority ||= job&.priority + self.priority = job&.priority if job&.priority.to_i > priority end - end From 032dd08e05aa7670f7f548e9eaf551f2f3d58122 Mon Sep 17 00:00:00 2001 From: Rosa Gutierrez Date: Fri, 17 Feb 2023 11:30:07 +0100 Subject: [PATCH 4/4] Tweak logging for benchmark and production tests --- app/models/solid_queue/claimed_execution.rb | 11 ++++++++++- app/models/solid_queue/job.rb | 2 +- app/models/solid_queue/scheduled_execution.rb | 2 ++ 3 files changed, 13 insertions(+), 2 deletions(-) diff --git a/app/models/solid_queue/claimed_execution.rb b/app/models/solid_queue/claimed_execution.rb index a5ea629df..78571a2a6 100644 --- a/app/models/solid_queue/claimed_execution.rb +++ b/app/models/solid_queue/claimed_execution.rb @@ -1,7 +1,10 @@ class SolidQueue::ClaimedExecution < SolidQueue::Execution def self.claim_batch(job_ids) - rows = job_ids.map { |id| { job_id: id, created_at: Time.current } } + claimed_at = Time.current + rows = job_ids.map { |id| { job_id: id, created_at: claimed_at } } insert_all(rows) if rows.any? + + SolidQueue.logger.info("[SolidQueue] Claimed #{rows.size} jobs at #{claimed_at}") end def perform @@ -13,6 +16,8 @@ def perform private def execute + SolidQueue.logger.info("[SolidQueue] Performing job #{job.id} - #{job.active_job_id}") + ActiveJob::Base.execute(job.arguments) end @@ -21,6 +26,8 @@ def finished job.finished destroy! end + + SolidQueue.logger.info("[SolidQueue] Performed job #{job.id} - #{job.active_job_id}") end def failed_with(error) @@ -28,5 +35,7 @@ def failed_with(error) job.failed_with(error) destroy! end + + SolidQueue.logger.info("[SolidQueue] Failed job #{job.id} - #{job.active_job_id}") end end diff --git a/app/models/solid_queue/job.rb b/app/models/solid_queue/job.rb index 3f52279da..9b115acd2 100644 --- a/app/models/solid_queue/job.rb +++ b/app/models/solid_queue/job.rb @@ -15,7 +15,7 @@ def enqueue_active_job(active_job, scheduled_at: Time.current) def enqueue(**kwargs) create!(**kwargs.compact.with_defaults(defaults)).tap do - SolidQueue.logger.info "[SolidQueue] Enqueued job #{kwargs}" + SolidQueue.logger.debug "[SolidQueue] Enqueued job #{kwargs}" end end diff --git a/app/models/solid_queue/scheduled_execution.rb b/app/models/solid_queue/scheduled_execution.rb index 5eb631c36..5d1a1ff67 100644 --- a/app/models/solid_queue/scheduled_execution.rb +++ b/app/models/solid_queue/scheduled_execution.rb @@ -19,6 +19,8 @@ def prepare_batch(batch) where(id: batch.map(&:id)).delete_all end end + + SolidQueue.logger.info("[SolidQueue] Prepared scheduled batch with #{rows.size} jobs at #{prepared_at}") end end