From e1442967b0268bddf5a6b0559248d7c71446ae16 Mon Sep 17 00:00:00 2001 From: Oleksandr Rohachev Date: Thu, 13 Aug 2026 16:06:15 +0300 Subject: [PATCH 1/2] Added DLQ configuration, interface for DLQ backends and logic to save dead task into DLQ after all retries got exhausted --- lib/rage/configuration.rb | 29 +++++++++++++++++++++++--- lib/rage/deferred/deferred.rb | 9 +++++++- lib/rage/deferred/dlq_backends/disk.rb | 21 +++++++++++++++++++ lib/rage/deferred/dlq_backends/nil.rb | 15 +++++++++++++ lib/rage/deferred/queue.rb | 13 ++++++++++-- 5 files changed, 81 insertions(+), 6 deletions(-) create mode 100644 lib/rage/deferred/dlq_backends/disk.rb create mode 100644 lib/rage/deferred/dlq_backends/nil.rb diff --git a/lib/rage/configuration.rb b/lib/rage/configuration.rb index c4d3e17e..90c5d1c5 100644 --- a/lib/rage/configuration.rb +++ b/lib/rage/configuration.rb @@ -873,17 +873,28 @@ def backend=(config) [config, {}] end - @backend_class = case backend_id + @backend_class, @dlq_backend_class = case backend_id when :disk @backend_options = parse_disk_backend_options(opts) - Rage::Deferred::Backends::Disk + @dlq_backend_options = dlq_disk_backend_options(@backend_options) + [Rage::Deferred::Backends::Disk, Rage::Deferred::DLQBackends::Disk] when nil - Rage::Deferred::Backends::Nil + [Rage::Deferred::Backends::Nil, Rage::Deferred::DLQBackends::Nil] else raise ArgumentError, "unsupported backend value; supported keys are `:disk` and `nil`" end end + def dlq_backend + unless @dlq_backend_class + @dlq_backend_class = Rage::Deferred::DLQBackends::Disk + default_disk_opts = parse_disk_backend_options({}) + @dlq_backend_options = dlq_disk_backend_options(default_disk_opts) + end + + @dlq_backend_class.new(**@dlq_backend_options) + end + class Backpressure attr_reader :high_water_mark, :low_water_mark, :timeout, :sleep_interval, :timeout_iterations @@ -1001,6 +1012,10 @@ def default_disk_storage_prefix "deferred-" end + def default_disk_dlq_storage_prefix + "dlq-" + end + # @private def has_default_disk_storage? default_disk_storage_path.glob("#{default_disk_storage_prefix}*").any? @@ -1040,6 +1055,14 @@ def parse_disk_backend_options(opts) parsed_options end + + def dlq_disk_backend_options(disk_opts) + { + path: disk_opts[:path], + fsync_frequency: disk_opts[:fsync_frequency], + prefix: default_disk_dlq_storage_prefix + } + end end # The class allows configuring telemetry handlers. See {MiddlewareRegistry} for details on available methods. diff --git a/lib/rage/deferred/deferred.rb b/lib/rage/deferred/deferred.rb index 274eba34..68e16caa 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.__dlq_backend + @__dlq_backend ||= Rage.config.deferred.dlq_backend + end + # @private def self.__queue - @__queue ||= Rage::Deferred::Queue.new(__backend) + @__queue ||= Rage::Deferred::Queue.new(__backend, __dlq_backend) end # @private @@ -110,6 +115,8 @@ class PushTimeout < StandardError require_relative "middleware_chain" require_relative "backends/disk" require_relative "backends/nil" +require_relative "dlq_backends/disk" +require_relative "dlq_backends/nil" if Iodine.running? Rage::Deferred.__initialize diff --git a/lib/rage/deferred/dlq_backends/disk.rb b/lib/rage/deferred/dlq_backends/disk.rb new file mode 100644 index 00000000..d25966b0 --- /dev/null +++ b/lib/rage/deferred/dlq_backends/disk.rb @@ -0,0 +1,21 @@ +# frozen_string_literal: true + +class Rage::Deferred::DLQBackends::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/dlq_backends/nil.rb b/lib/rage/deferred/dlq_backends/nil.rb new file mode 100644 index 00000000..00af11ed --- /dev/null +++ b/lib/rage/deferred/dlq_backends/nil.rb @@ -0,0 +1,15 @@ +# frozen_string_literal: true + +class Rage::Deferred::DLQBackends::Nil + def initialize + end + + def add(_, **) + end + + def remove(_) + end + + def retry(_, **) + end +end diff --git a/lib/rage/deferred/queue.rb b/lib/rage/deferred/queue.rb index 1912c71d..c81a653f 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, dlq_backend) @backend = backend + @dlq_backend = dlq_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) + dead_tasked = begin + @dlq_backend.add(context, exception: result, task_id:) + true + rescue => e + Rage.logger.error("Could not dead-task task #{task_id}: #{e.class} (#{e.message})") + false + end + + @backend.remove(task_id) if dead_tasked end end From 8979dd53eccc819051ce7ec3002bbcf216bce70f Mon Sep 17 00:00:00 2001 From: Oleksandr Rohachev Date: Thu, 13 Aug 2026 17:16:13 +0300 Subject: [PATCH 2/2] Rename DLQ to Dead Tasks Queue --- lib/rage/configuration.rb | 29 ++++++++++--------- .../disk.rb | 2 +- .../nil.rb | 4 +-- lib/rage/deferred/deferred.rb | 15 ++++++---- lib/rage/deferred/queue.rb | 12 ++++---- 5 files changed, 33 insertions(+), 29 deletions(-) rename lib/rage/deferred/{dlq_backends => dead_tasks_queue_backends}/disk.rb (85%) rename lib/rage/deferred/{dlq_backends => dead_tasks_queue_backends}/nil.rb (61%) diff --git a/lib/rage/configuration.rb b/lib/rage/configuration.rb index 90c5d1c5..3bee2dc6 100644 --- a/lib/rage/configuration.rb +++ b/lib/rage/configuration.rb @@ -873,26 +873,27 @@ def backend=(config) [config, {}] end - @backend_class, @dlq_backend_class = case backend_id + @backend_class, @dead_tasks_queue_backend_class = case backend_id when :disk @backend_options = parse_disk_backend_options(opts) - @dlq_backend_options = dlq_disk_backend_options(@backend_options) - [Rage::Deferred::Backends::Disk, Rage::Deferred::DLQBackends::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, Rage::Deferred::DLQBackends::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 dlq_backend - unless @dlq_backend_class - @dlq_backend_class = Rage::Deferred::DLQBackends::Disk + 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({}) - @dlq_backend_options = dlq_disk_backend_options(default_disk_opts) + @dead_tasks_queue_backend_options = dead_tasks_queue_disk_backend_options(default_disk_opts) end - - @dlq_backend_class.new(**@dlq_backend_options) + + @dead_tasks_queue_backend_class.new(**@dead_tasks_queue_backend_options) end class Backpressure @@ -1012,8 +1013,8 @@ def default_disk_storage_prefix "deferred-" end - def default_disk_dlq_storage_prefix - "dlq-" + def default_disk_dead_tasks_queue_storage_prefix + "dead_tasks_queue-" end # @private @@ -1056,11 +1057,11 @@ def parse_disk_backend_options(opts) parsed_options end - def dlq_disk_backend_options(disk_opts) + def dead_tasks_queue_disk_backend_options(disk_opts) { path: disk_opts[:path], fsync_frequency: disk_opts[:fsync_frequency], - prefix: default_disk_dlq_storage_prefix + prefix: default_disk_dead_tasks_queue_storage_prefix } end end diff --git a/lib/rage/deferred/dlq_backends/disk.rb b/lib/rage/deferred/dead_tasks_queue_backends/disk.rb similarity index 85% rename from lib/rage/deferred/dlq_backends/disk.rb rename to lib/rage/deferred/dead_tasks_queue_backends/disk.rb index d25966b0..24e3d30f 100644 --- a/lib/rage/deferred/dlq_backends/disk.rb +++ b/lib/rage/deferred/dead_tasks_queue_backends/disk.rb @@ -1,6 +1,6 @@ # frozen_string_literal: true -class Rage::Deferred::DLQBackends::Disk +class Rage::Deferred::DeadTasksQueueBackends::Disk def initialize(path:, prefix:, fsync_frequency:) @storage_path = path @prefix = prefix diff --git a/lib/rage/deferred/dlq_backends/nil.rb b/lib/rage/deferred/dead_tasks_queue_backends/nil.rb similarity index 61% rename from lib/rage/deferred/dlq_backends/nil.rb rename to lib/rage/deferred/dead_tasks_queue_backends/nil.rb index 00af11ed..d4504d39 100644 --- a/lib/rage/deferred/dlq_backends/nil.rb +++ b/lib/rage/deferred/dead_tasks_queue_backends/nil.rb @@ -1,7 +1,7 @@ # frozen_string_literal: true -class Rage::Deferred::DLQBackends::Nil - def initialize +class Rage::Deferred::DeadTasksQueueBackends::Nil + def initialize(**) end def add(_, **) diff --git a/lib/rage/deferred/deferred.rb b/lib/rage/deferred/deferred.rb index 68e16caa..3a726258 100644 --- a/lib/rage/deferred/deferred.rb +++ b/lib/rage/deferred/deferred.rb @@ -56,14 +56,14 @@ def self.__backend @__backend ||= Rage.config.deferred.backend end - # @private - def self.__dlq_backend - @__dlq_backend ||= Rage.config.deferred.dlq_backend + # @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, __dlq_backend) + @__queue ||= Rage::Deferred::Queue.new(__backend, __dead_tasks_queue_backend) end # @private @@ -102,6 +102,9 @@ def self.__initialize module Backends end + module DeadTasksQueueBackends + end + class PushTimeout < StandardError end end @@ -115,8 +118,8 @@ class PushTimeout < StandardError require_relative "middleware_chain" require_relative "backends/disk" require_relative "backends/nil" -require_relative "dlq_backends/disk" -require_relative "dlq_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 c81a653f..f7ee1eeb 100644 --- a/lib/rage/deferred/queue.rb +++ b/lib/rage/deferred/queue.rb @@ -3,9 +3,9 @@ class Rage::Deferred::Queue attr_reader :backlog_size - def initialize(backend, dlq_backend) + def initialize(backend, dead_tasks_queue_backend) @backend = backend - @dlq_backend = dlq_backend + @dead_tasks_queue_backend = dead_tasks_queue_backend @backlog_size = 0 @backpressure = Rage.config.deferred.backpressure end @@ -49,15 +49,15 @@ def schedule(task_id, context, publish_in: nil) if retry_in enqueue(context, delay: retry_in, task_id:) else - dead_tasked = begin - @dlq_backend.add(context, exception: result, 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 dead-task task #{task_id}: #{e.class} (#{e.message})") + 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 dead_tasked + @backend.remove(task_id) if added_to_dead_tasks_queue end end