Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 10 additions & 1 deletion app/models/solid_queue/claimed_execution.rb
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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

Expand All @@ -21,12 +26,16 @@ def finished
job.finished
destroy!
end

SolidQueue.logger.info("[SolidQueue] Performed job #{job.id} - #{job.active_job_id}")
end

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

SolidQueue.logger.info("[SolidQueue] Failed job #{job.id} - #{job.active_job_id}")
end
end
3 changes: 1 addition & 2 deletions app/models/solid_queue/execution.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
4 changes: 3 additions & 1 deletion app/models/solid_queue/job.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions app/models/solid_queue/scheduled_execution.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
7 changes: 7 additions & 0 deletions lib/solid_queue.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -33,7 +38,7 @@ def run
workers_pool.post { job.perform }
end
else
sleep(polling_interval)
interruptable_sleep(polling_interval)
end
end
end
Expand Down
10 changes: 10 additions & 0 deletions lib/solid_queue/engine.rb
Original file line number Diff line number Diff line change
@@ -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
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,10 @@

module SolidQueue::Runnable
def start
trap_signals
@stopping = false
@thread = Thread.new { run }

log "Started #{self}"
end

def stop
Expand All @@ -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
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
# frozen_string_literal: true

class SolidQueue::Scheduler
include SolidQueue::Runnable

Expand All @@ -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
Expand All @@ -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
Expand Down
36 changes: 36 additions & 0 deletions lib/solid_queue/supervisor.rb
Original file line number Diff line number Diff line change
@@ -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
6 changes: 6 additions & 0 deletions lib/solid_queue/tasks.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
namespace :solid_queue do
desc "start solid_queue supervisor"
task start: :environment do
SolidQueue::Supervisor.start
end
end
4 changes: 2 additions & 2 deletions solid_queue.gemspec
Original file line number Diff line number Diff line change
Expand Up @@ -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."
Expand Down
Empty file added test/models/solid_queue/.keep
Empty file.
Original file line number Diff line number Diff line change
@@ -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
Expand Down