diff --git a/README.md b/README.md index 4367c6b0..db2ddad3 100644 --- a/README.md +++ b/README.md @@ -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) @@ -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 @@ -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. diff --git a/app/models/solid_queue/job.rb b/app/models/solid_queue/job.rb index cb8f25be..a588db86 100644 --- a/app/models/solid_queue/job.rb +++ b/app/models/solid_queue/job.rb @@ -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) @@ -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 diff --git a/lib/solid_queue.rb b/lib/solid_queue.rb index fbc2f5e6..c394f44f 100644 --- a/lib/solid_queue.rb +++ b/lib/solid_queue.rb @@ -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) diff --git a/lib/solid_queue/enqueue_buffer.rb b/lib/solid_queue/enqueue_buffer.rb new file mode 100644 index 00000000..a7ac4b02 --- /dev/null +++ b/lib/solid_queue/enqueue_buffer.rb @@ -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 diff --git a/lib/solid_queue/log_subscriber.rb b/lib/solid_queue/log_subscriber.rb index fd4f54fd..eb4b5777 100644 --- a/lib/solid_queue/log_subscriber.rb +++ b/lib/solid_queue/log_subscriber.rb @@ -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 diff --git a/test/unit/enqueue_buffer_test.rb b/test/unit/enqueue_buffer_test.rb new file mode 100644 index 00000000..c0623b69 --- /dev/null +++ b/test/unit/enqueue_buffer_test.rb @@ -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 diff --git a/test/unit/log_subscriber_test.rb b/test/unit/log_subscriber_test.rb index 9e239923..ab63f334 100644 --- a/test/unit/log_subscriber_test.rb +++ b/test/unit/log_subscriber_test.rb @@ -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