Skip to content
Draft
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
30 changes: 27 additions & 3 deletions lib/rage/configuration.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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?
Expand Down Expand Up @@ -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.
Expand Down
21 changes: 21 additions & 0 deletions lib/rage/deferred/dead_tasks_queue_backends/disk.rb
Original file line number Diff line number Diff line change
@@ -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
15 changes: 15 additions & 0 deletions lib/rage/deferred/dead_tasks_queue_backends/nil.rb
Original file line number Diff line number Diff line change
@@ -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
12 changes: 11 additions & 1 deletion lib/rage/deferred/deferred.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -97,6 +102,9 @@ def self.__initialize
module Backends
end

module DeadTasksQueueBackends
end

class PushTimeout < StandardError
end
end
Expand All @@ -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
Expand Down
13 changes: 11 additions & 2 deletions lib/rage/deferred/queue.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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

Expand Down
Loading