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/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 diff --git a/app/models/solid_queue/job.rb b/app/models/solid_queue/job.rb index f9386db16..9b115acd2 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.debug "[SolidQueue] Enqueued job #{kwargs}" + end end private 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 diff --git a/lib/solid_queue.rb b/lib/solid_queue.rb index 014ced5ed..261141165 100644 --- a/lib/solid_queue.rb +++ b/lib/solid_queue.rb @@ -3,5 +3,12 @@ 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/supervisor" + module SolidQueue + mattr_accessor :logger, default: ActiveSupport::Logger.new($stdout) 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 84% rename from app/models/solid_queue/dispatcher.rb rename to lib/solid_queue/dispatcher.rb index 19698d00e..8b16ace83 100644 --- a/app/models/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/app/models/concerns/solid_queue/runnable.rb b/lib/solid_queue/runnable.rb similarity index 61% rename from app/models/concerns/solid_queue/runnable.rb rename to lib/solid_queue/runnable.rb index 7f6c92d58..3c6680717 100644 --- a/app/models/concerns/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/app/models/solid_queue/scheduler.rb b/lib/solid_queue/scheduler.rb similarity index 74% rename from app/models/solid_queue/scheduler.rb rename to lib/solid_queue/scheduler.rb index 0c4eedd6d..0c68c54f3 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 @@ -10,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 @@ -20,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/supervisor.rb b/lib/solid_queue/supervisor.rb new file mode 100644 index 000000000..afca725eb --- /dev/null +++ b/lib/solid_queue/supervisor.rb @@ -0,0 +1,36 @@ +# frozen_string_literal: true + +class SolidQueue::Supervisor + include SolidQueue::Runnable + + attr_accessor :dispatchers, :scheduler + + def self.start + configuration = SolidQueue::Configuration.new + dispatchers = configuration.queues.values.map { |queue_options| SolidQueue::Dispatcher.new(**queue_options) } + 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/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." 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