diff --git a/lib/rage/configuration.rb b/lib/rage/configuration.rb index c4d3e17e..3bee2dc6 100644 --- a/lib/rage/configuration.rb +++ b/lib/rage/configuration.rb @@ -873,17 +873,29 @@ def backend=(config) [config, {}] end - @backend_class = case backend_id + @backend_class, @dead_tasks_queue_backend_class = case backend_id when :disk @backend_options = parse_disk_backend_options(opts) - Rage::Deferred::Backends::Disk + @dead_tasks_queue_backend_options = dead_tasks_queue_disk_backend_options(@backend_options) + [Rage::Deferred::Backends::Disk, Rage::Deferred::DeadTasksQueueBackends::Disk] when nil - Rage::Deferred::Backends::Nil + @dead_tasks_queue_backend_options = {} + [Rage::Deferred::Backends::Nil, Rage::Deferred::DeadTasksQueueBackends::Nil] else raise ArgumentError, "unsupported backend value; supported keys are `:disk` and `nil`" end end + def dead_tasks_queue_backend + unless @dead_tasks_queue_backend_class + @dead_tasks_queue_backend_class = Rage::Deferred::DeadTasksQueueBackends::Disk + default_disk_opts = parse_disk_backend_options({}) + @dead_tasks_queue_backend_options = dead_tasks_queue_disk_backend_options(default_disk_opts) + end + + @dead_tasks_queue_backend_class.new(**@dead_tasks_queue_backend_options) + end + class Backpressure attr_reader :high_water_mark, :low_water_mark, :timeout, :sleep_interval, :timeout_iterations @@ -1001,6 +1013,10 @@ def default_disk_storage_prefix "deferred-" end + def default_disk_dead_tasks_queue_storage_prefix + "dead_tasks_queue-" + end + # @private def has_default_disk_storage? default_disk_storage_path.glob("#{default_disk_storage_prefix}*").any? @@ -1040,6 +1056,14 @@ def parse_disk_backend_options(opts) parsed_options end + + def dead_tasks_queue_disk_backend_options(disk_opts) + { + path: disk_opts[:path], + fsync_frequency: disk_opts[:fsync_frequency], + prefix: default_disk_dead_tasks_queue_storage_prefix + } + end end # The class allows configuring telemetry handlers. See {MiddlewareRegistry} for details on available methods. diff --git a/lib/rage/deferred/dead_tasks_queue_backends/disk.rb b/lib/rage/deferred/dead_tasks_queue_backends/disk.rb new file mode 100644 index 00000000..24e3d30f --- /dev/null +++ b/lib/rage/deferred/dead_tasks_queue_backends/disk.rb @@ -0,0 +1,21 @@ +# frozen_string_literal: true + +class Rage::Deferred::DeadTasksQueueBackends::Disk + def initialize(path:, prefix:, fsync_frequency:) + @storage_path = path + @prefix = prefix + @fsync_frequency = fsync_frequency + + # TODO: implement initializer + # @storage_path.mkpath + end + + def add(_, **) + end + + def remove(_) + end + + def retry(_, **) + end +end diff --git a/lib/rage/deferred/dead_tasks_queue_backends/nil.rb b/lib/rage/deferred/dead_tasks_queue_backends/nil.rb new file mode 100644 index 00000000..d4504d39 --- /dev/null +++ b/lib/rage/deferred/dead_tasks_queue_backends/nil.rb @@ -0,0 +1,15 @@ +# frozen_string_literal: true + +class Rage::Deferred::DeadTasksQueueBackends::Nil + def initialize(**) + end + + def add(_, **) + end + + def remove(_) + end + + def retry(_, **) + end +end diff --git a/lib/rage/deferred/deferred.rb b/lib/rage/deferred/deferred.rb index 274eba34..3a726258 100644 --- a/lib/rage/deferred/deferred.rb +++ b/lib/rage/deferred/deferred.rb @@ -56,9 +56,14 @@ def self.__backend @__backend ||= Rage.config.deferred.backend end + # @private + def self.__dead_tasks_queue_backend + @__dead_tasks_queue_backend ||= Rage.config.deferred.dead_tasks_queue_backend + end + # @private def self.__queue - @__queue ||= Rage::Deferred::Queue.new(__backend) + @__queue ||= Rage::Deferred::Queue.new(__backend, __dead_tasks_queue_backend) end # @private @@ -97,6 +102,9 @@ def self.__initialize module Backends end + module DeadTasksQueueBackends + end + class PushTimeout < StandardError end end @@ -110,6 +118,8 @@ class PushTimeout < StandardError require_relative "middleware_chain" require_relative "backends/disk" require_relative "backends/nil" +require_relative "dead_tasks_queue_backends/disk" +require_relative "dead_tasks_queue_backends/nil" if Iodine.running? Rage::Deferred.__initialize diff --git a/lib/rage/deferred/queue.rb b/lib/rage/deferred/queue.rb index 1912c71d..f7ee1eeb 100644 --- a/lib/rage/deferred/queue.rb +++ b/lib/rage/deferred/queue.rb @@ -3,8 +3,9 @@ class Rage::Deferred::Queue attr_reader :backlog_size - def initialize(backend) + def initialize(backend, dead_tasks_queue_backend) @backend = backend + @dead_tasks_queue_backend = dead_tasks_queue_backend @backlog_size = 0 @backpressure = Rage.config.deferred.backpressure end @@ -48,7 +49,15 @@ def schedule(task_id, context, publish_in: nil) if retry_in enqueue(context, delay: retry_in, task_id:) else - @backend.remove(task_id) + added_to_dead_tasks_queue = begin + @dead_tasks_queue_backend.add(context, exception: result, task_id:) + true + rescue => e + Rage.logger.error("Could not add task #{task_id} to the dead tasks queue: #{e.class} (#{e.message})") + false + end + + @backend.remove(task_id) if added_to_dead_tasks_queue end end