Class: DataMigration::Job

Inherits:
ActiveJob::Base
  • Object
show all
Defined in:
lib/data_migration/job.rb

Instance Method Summary collapse

Instance Method Details

#perform(task_id, *job_args, **job_kwargs) ⇒ Object



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