From 2939e3d5d4ad1da93bb5171ba08d17f49815642c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Sa=C3=AFd=20Mimouni?= Date: Tue, 25 Aug 2026 19:34:34 +0200 Subject: [PATCH] Fix rescheduling recreated dynamic tasks --- .../scheduler/recurring_schedule.rb | 70 ++++++++++++++----- test/unit/scheduler_test.rb | 57 +++++++++++++++ 2 files changed, 111 insertions(+), 16 deletions(-) diff --git a/lib/solid_queue/scheduler/recurring_schedule.rb b/lib/solid_queue/scheduler/recurring_schedule.rb index 526a8edcd..7f1f07b7d 100644 --- a/lib/solid_queue/scheduler/recurring_schedule.rb +++ b/lib/solid_queue/scheduler/recurring_schedule.rb @@ -10,7 +10,10 @@ def initialize(static_tasks, dynamic_tasks_enabled: false) @static_tasks = Array(static_tasks).map { |task| RecurringTask.wrap(task) }.select(&:valid?) @dynamic_tasks_enabled = dynamic_tasks_enabled + @schedule_lock = Mutex.new @scheduled_tasks = Concurrent::Hash.new + @scheduled_dynamic_task_ids = {} + @active = false end def configured_tasks @@ -28,18 +31,23 @@ def schedule_tasks reload_dynamic_tasks end - configured_tasks.each do |task| - schedule_task(task) + schedule_lock.synchronize do + @active = true + configured_tasks.each { |task| schedule_task_without_lock(task) } end end def schedule_task(task, run_at: task.next_time) - scheduled_tasks[task.key] = schedule(task, run_at: run_at) + schedule_lock.synchronize { schedule_task_without_lock(task, run_at: run_at) } end def unschedule_tasks - scheduled_tasks.values.each(&:cancel) - scheduled_tasks.clear + schedule_lock.synchronize do + @active = false + scheduled_tasks.values.each(&:cancel) + scheduled_tasks.clear + scheduled_dynamic_task_ids.clear + end end def task_keys @@ -48,14 +56,17 @@ def task_keys def reschedule_dynamic_tasks wrap_in_app_executor do - reload_dynamic_tasks - schedule_created_dynamic_tasks - unschedule_deleted_dynamic_tasks + schedule_lock.synchronize do + reload_dynamic_tasks + schedule_created_dynamic_tasks + reschedule_recreated_dynamic_tasks + unschedule_deleted_dynamic_tasks + end end end private - attr_reader :static_tasks + attr_reader :static_tasks, :schedule_lock, :scheduled_dynamic_task_ids def static_task_keys static_tasks.map(&:key) @@ -70,18 +81,45 @@ def dynamic_tasks_enabled? end def schedule_created_dynamic_tasks - RecurringTask.dynamic.where.not(key: scheduled_tasks.keys).each do |task| - schedule_task(task) + dynamic_tasks.reject { |task| scheduled_tasks.key?(task.key) }.each do |task| + schedule_task_without_lock(task) + end + end + + def reschedule_recreated_dynamic_tasks + dynamic_tasks.each do |task| + next unless scheduled_dynamic_task_ids.key?(task.key) + next if scheduled_dynamic_task_ids[task.key] == task.id + + unschedule_task(task.key) + schedule_task_without_lock(task) end end def unschedule_deleted_dynamic_tasks - (scheduled_tasks.keys - RecurringTask.pluck(:key)).each do |key| - scheduled_tasks[key].cancel - scheduled_tasks.delete(key) + (scheduled_dynamic_task_ids.keys - dynamic_tasks.map(&:key)).each { |key| unschedule_task(key) } + end + + def unschedule_task(key) + scheduled_tasks.delete(key)&.cancel + scheduled_dynamic_task_ids.delete(key) + end + + def schedule_task_without_lock(task, run_at: task.next_time) + scheduled_tasks[task.key] = schedule(task, run_at: run_at) + scheduled_dynamic_task_ids[task.key] = task.id unless task.static? + end + + def schedule_next_task(task, run_at:) + schedule_lock.synchronize do + schedule_task_without_lock(task, run_at: run_at) if @active && current_task?(task) end end + def current_task?(task) + task.static? || scheduled_dynamic_task_ids[task.key] == task.id + end + def persist_static_tasks RecurringTask.static.where.not(key: static_task_keys).delete_all RecurringTask.create_or_update_all static_tasks @@ -102,8 +140,8 @@ def load_dynamic_tasks def schedule(task, run_at: task.next_time) delay = [ (run_at - Time.current).to_f, 0.1 ].max - scheduled_task = Concurrent::ScheduledTask.new(delay, args: [ self, task, run_at ]) do |thread_schedule, thread_task, thread_task_run_at| - thread_schedule.schedule_task(thread_task, run_at: thread_task.next_time_after(thread_task_run_at)) + scheduled_task = Concurrent::ScheduledTask.new(delay, args: [ task, run_at ]) do |thread_task, thread_task_run_at| + schedule_next_task(thread_task, run_at: thread_task.next_time_after(thread_task_run_at)) wrap_in_app_executor do thread_task.enqueue(at: thread_task_run_at) diff --git a/test/unit/scheduler_test.rb b/test/unit/scheduler_test.rb index e914a23ca..35f276980 100644 --- a/test/unit/scheduler_test.rb +++ b/test/unit/scheduler_test.rb @@ -156,4 +156,61 @@ class SchedulerTest < ActiveSupport::TestCase ensure scheduler&.stop end + + test "reschedules a dynamic task recreated with the same key" do + recurring_schedule = nil + travel_to Time.zone.local(2026, 8, 25, 10, 15) do + SolidQueue.schedule_recurring_task( + "dynamic_task", + class: "AddToBufferJob", + schedule: "every hour", + args: [ 42 ] + ) + original_task = SolidQueue::RecurringTask.find_by!(key: "dynamic_task") + recurring_schedule = SolidQueue::Scheduler::RecurringSchedule.new({}, dynamic_tasks_enabled: true) + recurring_schedule.schedule_tasks + original_scheduled_task = recurring_schedule.scheduled_tasks.fetch("dynamic_task") + + SolidQueue.unschedule_recurring_task("dynamic_task") + SolidQueue.schedule_recurring_task( + "dynamic_task", + class: "AddToBufferJob", + schedule: "every day", + args: [ 42 ] + ) + replacement_task = SolidQueue::RecurringTask.find_by!(key: "dynamic_task") + + recurring_schedule.reschedule_dynamic_tasks + replacement_scheduled_task = recurring_schedule.scheduled_tasks.fetch("dynamic_task") + + assert_instance_of Concurrent::CancelledOperationError, original_scheduled_task.reason + assert_not_same original_scheduled_task, replacement_scheduled_task + assert_in_delta replacement_task.next_time - Time.current, replacement_scheduled_task.initial_delay, 0.1 + + recurring_schedule.send(:schedule_next_task, original_task, run_at: original_task.next_time) + assert_same replacement_scheduled_task, recurring_schedule.scheduled_tasks.fetch("dynamic_task") + end + ensure + recurring_schedule&.unschedule_tasks + end + + test "does not reschedule a dynamic task updated in place" do + task = SolidQueue::RecurringTask.create!( + key: "dynamic_task", + static: false, + class_name: "AddToBufferJob", + schedule: "every hour", + arguments: [ 42 ] + ) + recurring_schedule = SolidQueue::Scheduler::RecurringSchedule.new({}, dynamic_tasks_enabled: true) + recurring_schedule.schedule_tasks + scheduled_task = recurring_schedule.scheduled_tasks.fetch("dynamic_task") + + task.update!(schedule: "every day") + recurring_schedule.reschedule_dynamic_tasks + + assert_same scheduled_task, recurring_schedule.scheduled_tasks.fetch("dynamic_task") + ensure + recurring_schedule&.unschedule_tasks + end end