9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
|
# File 'lib/data_migration/job.rb', line 9
def perform(task_id, *job_args, **job_kwargs)
checked_in = false
failed = false
migration_started = false
migration_class = nil
task = DataMigration::Task.find(task_id)
DataMigration.config.monitoring_context.call(task)
migration_name = task.name
migration_path = task.file_path
unless task.file_exists?
DataMigration.notify("#{migration_name} not found")
task.update_columns(status: task.class.statuses.fetch("failed"), updated_at: Time.current)
return
end
task.job_check_in!(job_id, job_args: job_args, job_kwargs: job_kwargs)
checked_in = true
migration_started = true
require migration_path
klass_name = migration_name.gsub(/^[0-9_]+/, "").camelize
migration_class = klass_name.safe_constantize
raise "Data migration class #{klass_name} not found" unless migration_class.is_a?(Class)
raise "Data migration class #{klass_name} must implement `perform` method" unless migration_class.method_defined?(:perform)
if task.started_at.blank?
task.update!(started_at: Time.current, status: :started)
end
if task.requires_pause?
DataMigration::Job.set(wait: task.pause_minutes.minutes).perform_later(task_id, *job_args, **job_kwargs)
task.update!(status: :paused)
return
end
Thread.current[:data_migration_enqueue_called] ||= {}
Thread.current[:data_migration_enqueue_kwargs] ||= {}
migration_class.define_method(:enqueue) do |**enqueue_kwargs|
Thread.current[:data_migration_enqueue_called][migration_class.name] = true
Thread.current[:data_migration_enqueue_kwargs][migration_class.name] = enqueue_kwargs
end
task.update!(status: :performing, pause_minutes: 0)
migration_class.new.perform(**job_kwargs)
task.job_check_out!(job_id)
checked_in = false
enqueue_called = Thread.current[:data_migration_enqueue_called].delete(migration_class.name)
enqueue_kwargs = Thread.current[:data_migration_enqueue_kwargs].delete(migration_class.name)
if enqueue_called
if enqueue_kwargs[:background] == false
self.class.new.perform(task_id, *job_args, **enqueue_kwargs)
else
DataMigration::Job.perform_later(task_id, *job_args, **enqueue_kwargs)
end
else
task.update!(completed_at: Time.current, status: :completed)
end
rescue
failed = migration_started
raise
ensure
task.job_check_out!(job_id, status: (failed ? :failed : nil)) if checked_in
task.update_columns(status: task.class.statuses.fetch("failed"), updated_at: Time.current) if failed && !checked_in
if migration_class
Thread.current[:data_migration_enqueue_called]&.delete(migration_class.name)
Thread.current[:data_migration_enqueue_kwargs]&.delete(migration_class.name)
end
end
|