Skip to content
Open
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
12 changes: 12 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ Solid Queue can be used with SQL databases such as MySQL, PostgreSQL, or SQLite,
- [Validating the configuration](#validating-the-configuration)
- [Lifecycle hooks](#lifecycle-hooks)
- [Errors when enqueuing](#errors-when-enqueuing)
- [Holding jobs while the database is unreachable](#holding-jobs-while-the-database-is-unreachable)
- [Concurrency controls](#concurrency-controls)
- [Performance considerations](#performance-considerations)
- [Failed jobs and retries](#failed-jobs-and-retries)
Expand Down Expand Up @@ -453,6 +454,8 @@ There are several settings that control how Solid Queue works that you can set a
- `preserve_finished_jobs`: whether to keep finished jobs in the `solid_queue_jobs` table—defaults to `true`.
- `clear_finished_jobs_after`: period to keep finished jobs around, in case `preserve_finished_jobs` is true — defaults to 1 day. When installing Solid Queue, [a recurring job](#recurring-tasks) is automatically configured to clear finished jobs every hour on the 12th minute in batches. You can edit the `recurring.yml` configuration to change this as you see fit.
- `default_concurrency_control_period`: the value to be used as the default for the `duration` parameter in [concurrency controls](#concurrency-controls). It defaults to 3 minutes.
- `buffer_enqueues_on_database_error`: whether to hold jobs in memory when the queue database can't be reached, and enqueue them once it's back, instead of raising—defaults to `false`. See [Holding jobs while the database is unreachable](#holding-jobs-while-the-database-is-unreachable).
- `enqueue_buffer_size`: how many jobs each process can hold while the queue database can't be reached, when `buffer_enqueues_on_database_error` is enabled—defaults to 1,000.

### Validating the configuration

Expand Down Expand Up @@ -525,6 +528,15 @@ Solid Queue will raise a `SolidQueue::Job::EnqueueError` for any Active Record e

In the case of recurring tasks, if such error is raised when enqueuing the job corresponding to the task, it'll be handled and logged but it won't bubble up.

### Holding jobs while the database is unreachable

If the queue database goes down, every `perform_later` raises and the job is lost. With `config.solid_queue.buffer_enqueues_on_database_error = true`, Solid Queue instead holds jobs in memory when enqueuing fails because the database can't be reached (`ActiveRecord::ConnectionNotEstablished` or `ActiveRecord::ConnectionFailed`), reports them as enqueued, and enqueues them from a background thread once the database is back, retrying with increasing waits of up to 30 seconds. This applies to `perform_later` and `perform_all_later` alike. Any other error raises as usual.

Keep in mind:
- Held jobs live in the process that enqueued them. If that process crashes, they're lost. On a normal exit, Solid Queue tries once more to enqueue them.
- Each process holds at most `enqueue_buffer_size` jobs. Once that's reached, enqueuing raises `SolidQueue::Job::EnqueueError` as it does without buffering.
- A held job has no `provider_job_id` until it's actually enqueued.

## Concurrency controls

Solid Queue extends Active Job with concurrency controls, that allows you to limit how many jobs of a certain type or with certain arguments can run at the same time. When limited in this way, **by default, jobs will be blocked from running**, and they'll stay blocked until another job finishes and unblocks them, or after the set expiry time (concurrency limit's _duration_) elapses.
Expand Down
5 changes: 5 additions & 0 deletions app/models/solid_queue/job.rb
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,9 @@ def enqueue_all(active_jobs)
end

active_jobs.count(&:successfully_enqueued?)
rescue *EnqueueBuffer::CONNECTION_ERRORS => error
raise unless EnqueueBuffer.hold(active_jobs, error)
active_jobs.size
end

def enqueue(active_job, scheduled_at: Time.current)
Expand All @@ -48,6 +51,8 @@ def enqueue(active_job, scheduled_at: Time.current)
active_job.provider_job_id = enqueued_job.id if enqueued_job.persisted?
active_job.successfully_enqueued = enqueued_job.persisted?
end
rescue EnqueueError => error
raise unless EnqueueBuffer.hold([ active_job ], error.cause)
end

private
Expand Down
3 changes: 3 additions & 0 deletions lib/solid_queue.rb
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,9 @@ module SolidQueue
mattr_accessor :clear_finished_jobs_after, default: 1.day
mattr_accessor :default_concurrency_control_period, default: 3.minutes

mattr_accessor :buffer_enqueues_on_database_error, default: false
mattr_accessor :enqueue_buffer_size, default: 1_000

mattr_reader :time_zone

def time_zone=(zone)
Expand Down
119 changes: 119 additions & 0 deletions lib/solid_queue/enqueue_buffer.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,119 @@
# frozen_string_literal: true

module SolidQueue
# Holds jobs in memory while the queue database can't be reached, and enqueues
# them once it's back, instead of raising and losing them. Opt-in with
# +SolidQueue.buffer_enqueues_on_database_error+.
#
# The buffer lives in the enqueuing process: jobs held when that process
# crashes are lost. A normal exit tries once more to enqueue what's left.
module EnqueueBuffer
extend self
extend AppExecutor

CONNECTION_ERRORS = [ ActiveRecord::ConnectionNotEstablished, ActiveRecord::ConnectionFailed ].freeze
MIN_RETRY_INTERVAL = 1.second
MAX_RETRY_INTERVAL = 30.seconds

# Held jobs belong to the process that held them. A forked child starts
# empty, so it can't enqueue its parent's jobs a second time, and gets its
# own flusher, since threads don't survive a fork.
def reset
@jobs = []
@mutex = Mutex.new
@flusher = nil
end

reset
ActiveSupport::ForkTracker.after_fork { reset }

# Holds +active_jobs+ when +error+ means the database can't be reached,
# buffering is enabled and the buffer has room for all of them. Held jobs
# count as enqueued.
def hold(active_jobs, error)
return false unless holdable?(error)

held = mutex.synchronize do
if jobs.size + active_jobs.size <= SolidQueue.enqueue_buffer_size
jobs.concat(active_jobs)
start_flushing
true
end
end

SolidQueue.instrument(:buffer_enqueue, size: active_jobs.size, held: !!held, error: error)
active_jobs.each { |job| job.successfully_enqueued = true } if held

!!held
end

# Tries to enqueue every held job. Returns false when the database still
# can't be reached, putting the jobs back to try again later.
def flush
batch = mutex.synchronize { jobs.shift(jobs.size) }
return true if batch.empty?

SolidQueue.instrument(:flush_enqueue_buffer, size: batch.size) do
flushing { wrap_in_app_executor { Job.enqueue_all(batch) } }
end
true
rescue *CONNECTION_ERRORS
mutex.synchronize { jobs.unshift(*batch) }
false
rescue StandardError => error
# Anything but an unreachable database won't fix itself by retrying:
# report it, as other Solid Queue thread errors are, and drop the batch.
handle_thread_error(error)
true
end

def size
mutex.synchronize { jobs.size }
end

private
def holdable?(error)
SolidQueue.buffer_enqueues_on_database_error &&
CONNECTION_ERRORS.any? { |error_class| error.is_a?(error_class) } &&
!Thread.current[:solid_queue_flushing_enqueue_buffer]
end

def flushing
Thread.current[:solid_queue_flushing_enqueue_buffer] = true
yield
ensure
Thread.current[:solid_queue_flushing_enqueue_buffer] = false
end

# Called with the mutex held.
def start_flushing
@flusher ||= create_thread { flush_until_empty }
@exit_hook ||= at_exit { flush_on_exit }
end

def flush_until_empty
interval = MIN_RETRY_INTERVAL

loop do
sleep interval
interval = flush ? MIN_RETRY_INTERVAL : [ interval * 2, MAX_RETRY_INTERVAL ].min
break if stop_flushing_if_empty
end
end

def flush_on_exit
SolidQueue.instrument(:lose_held_jobs, size: size) unless flush
end

# Decided under the mutex, so a job held right after the check starts a
# new flusher instead of waiting for one that is about to exit.
def stop_flushing_if_empty
mutex.synchronize do
@flusher = nil if jobs.empty?
@flusher.nil?
end
end

attr_reader :jobs, :mutex
end
end
18 changes: 18 additions & 0 deletions lib/solid_queue/log_subscriber.rb
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,24 @@ def deregister_process(event)
end
end

def buffer_enqueue(event)
attributes = event.payload.slice(:size, :held).merge(error: formatted_error(event.payload[:error]))

if event.payload[:held]
warn formatted_event(event, action: "Hold jobs until the database is reachable", **attributes)
else
error formatted_event(event, action: "Enqueue buffer full, jobs not held", **attributes)
end
end

def flush_enqueue_buffer(event)
info formatted_event(event, action: "Enqueue held jobs", **event.payload.slice(:size))
end

def lose_held_jobs(event)
error formatted_event(event, action: "Lose held jobs on exit, the database is still unreachable", **event.payload.slice(:size))
end

def prune_processes(event)
debug formatted_event(event, action: "Prune dead processes", **event.payload.slice(:size))
end
Expand Down
136 changes: 136 additions & 0 deletions test/unit/enqueue_buffer_test.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,136 @@
# frozen_string_literal: true

require "test_helper"

class EnqueueBufferTest < ActiveSupport::TestCase
self.use_transactional_tests = false

setup do
SolidQueue.buffer_enqueues_on_database_error = true
# The flusher thread is exercised on its own; elsewhere flush by hand.
SolidQueue::EnqueueBuffer.stubs(:start_flushing)
end

teardown do
SolidQueue.buffer_enqueues_on_database_error = false
SolidQueue.enqueue_buffer_size = 1_000
SolidQueue::EnqueueBuffer.reset
end

test "an enqueue that can't reach the database is held and reported as enqueued" do
database_unreachable_for(:create!)

job = AddToBufferJob.set(queue: :critical, priority: 7).perform_later(42)

assert job.successfully_enqueued?
assert_nil job.provider_job_id
assert_equal 1, SolidQueue::EnqueueBuffer.size
assert_equal 0, SolidQueue::Job.count
end

test "held jobs are enqueued as they were once the database is back" do
database_unreachable_for(:create!)
job = AddToBufferJob.set(queue: :critical, priority: 7, wait: 5.minutes).perform_later(42)
SolidQueue::Job.unstub(:create!)

assert SolidQueue::EnqueueBuffer.flush

enqueued = SolidQueue::Job.find_by!(active_job_id: job.job_id)
assert_equal [ "critical", 7 ], [ enqueued.queue_name, enqueued.priority ]
assert_in_delta job.scheduled_at, enqueued.scheduled_at, 1.second
assert enqueued.scheduled?
assert_equal enqueued.id, job.provider_job_id
assert_equal 0, SolidQueue::EnqueueBuffer.size
end

test "a bulk enqueue that can't reach the database is held as a whole" do
database_unreachable_for(:insert_all, ActiveRecord::ConnectionFailed)
jobs = [ AddToBufferJob.new(1), AddToBufferJob.new(2) ]

assert_equal 2, SolidQueue::Job.enqueue_all(jobs)

assert jobs.all?(&:successfully_enqueued?)
assert_equal 2, SolidQueue::EnqueueBuffer.size
end

test "held jobs stay held while the database is still unreachable" do
database_unreachable_for(:create!)
AddToBufferJob.perform_later(42)
database_unreachable_for(:insert_all)

assert_not SolidQueue::EnqueueBuffer.flush
assert_equal 1, SolidQueue::EnqueueBuffer.size
end

test "enqueue errors raise as before when buffering is off" do
SolidQueue.buffer_enqueues_on_database_error = false
database_unreachable_for(:create!)

assert_raises(SolidQueue::Job::EnqueueError) { AddToBufferJob.perform_later(42) }
assert_equal 0, SolidQueue::EnqueueBuffer.size
end

test "errors other than an unreachable database still raise" do
SolidQueue::Job.stubs(:create!).raises(ActiveRecord::StatementInvalid)

assert_raises(SolidQueue::Job::EnqueueError) { AddToBufferJob.perform_later(42) }
assert_equal 0, SolidQueue::EnqueueBuffer.size
end

test "enqueues raise once the buffer is full" do
SolidQueue.enqueue_buffer_size = 1
database_unreachable_for(:create!)
AddToBufferJob.perform_later(1)

assert_raises(SolidQueue::Job::EnqueueError) { AddToBufferJob.perform_later(2) }
assert_equal 1, SolidQueue::EnqueueBuffer.size
end

test "the flusher enqueues held jobs in the background and then stops" do
SolidQueue::EnqueueBuffer.unstub(:start_flushing)
SolidQueue::EnqueueBuffer.stubs(:sleep)
database_unreachable_for(:create!) # the flusher enqueues through insert_all, which works

job = AddToBufferJob.perform_later(42)

wait_for(timeout: 5.seconds) { SolidQueue::EnqueueBuffer.instance_variable_get(:@flusher).nil? }
assert SolidQueue::Job.exists?(active_job_id: job.job_id)
assert_equal 0, SolidQueue::EnqueueBuffer.size
end

test "held jobs still unenqueued at exit are reported as lost" do
database_unreachable_for(:create!)
AddToBufferJob.perform_later(42)
database_unreachable_for(:insert_all)
lost = []
subscriber = ActiveSupport::Notifications.subscribe("lose_held_jobs.solid_queue") { |event| lost << event.payload[:size] }

SolidQueue::EnqueueBuffer.send(:flush_on_exit)

assert_equal [ 1 ], lost
ensure
ActiveSupport::Notifications.unsubscribe(subscriber)
end

test "a forked child doesn't inherit held jobs" do
database_unreachable_for(:create!)
AddToBufferJob.perform_later(42)

reader, writer = IO.pipe
pid = fork do
reader.close
writer.puts SolidQueue::EnqueueBuffer.size
writer.close
end
writer.close
Process.wait(pid)

assert_equal "0", reader.read.strip
assert_equal 1, SolidQueue::EnqueueBuffer.size
end

private
def database_unreachable_for(method, error = ActiveRecord::ConnectionNotEstablished)
SolidQueue::Job.stubs(method).raises(error)
end
end
28 changes: 28 additions & 0 deletions test/unit/log_subscriber_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,34 @@ def set_logger(logger)
assert_match_logged :debug, "Release blocked job", "job_id: 42, concurrency_key: \"foo/1\", released: true"
end

test "jobs held while the database is unreachable" do
attach_log_subscriber
instrument "buffer_enqueue.solid_queue", size: 2, held: true, error: ActiveRecord::ConnectionNotEstablished.new("connection refused")

assert_match_logged :warn, "Hold jobs until the database is reachable", "size: 2, held: true, error: \"ActiveRecord::ConnectionNotEstablished connection refused\""
end

test "jobs not held because the enqueue buffer is full" do
attach_log_subscriber
instrument "buffer_enqueue.solid_queue", size: 1, held: false, error: ActiveRecord::ConnectionNotEstablished.new("connection refused")

assert_match_logged :error, "Enqueue buffer full, jobs not held", "size: 1, held: false"
end

test "held jobs enqueued" do
attach_log_subscriber
instrument "flush_enqueue_buffer.solid_queue", size: 3

assert_match_logged :info, "Enqueue held jobs", "size: 3"
end

test "held jobs lost on exit" do
attach_log_subscriber
instrument "lose_held_jobs.solid_queue", size: 4

assert_match_logged :error, "Lose held jobs on exit, the database is still unreachable", "size: 4"
end

test "unblock many jobs" do
attach_log_subscriber
instrument "release_many_blocked.solid_queue", limit: 42, size: 10
Expand Down
Loading