Class: CapistranoMulticonfigParallel::CelluloidManager

Inherits:
Object
  • Object
show all
Includes:
ApplicationHelper, Celluloid, Celluloid::Logger, Celluloid::Notifications
Defined in:
lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb

Overview

rubocop:disable ClassLength

Instance Attribute Summary collapse

Instance Method Summary collapse

Methods included from ApplicationHelper

parse_task_string, strip_characters_from_string

Methods included from StagesHelper

check_stage_path, checks_paths, fetch_stages, fetch_stages_paths, sorted_paths, stages_paths

Methods included from CoreHelper

app_configuration, app_debug_enabled?, app_logger, ask_confirm, check_terminal_tty, debug_websocket?, execute_with_rescue, find_config_type, find_loaded_gem, find_worker_log, force_confirmation, format_error, log_error, log_to_file, rescue_error, rescue_interrupt, show_warning, websocket_config, websocket_server_config

Methods included from InternalHelper

config_file, custom_commands, default_internal_config, detect_root, enable_main_log_file, find_env_multi_cap_root, internal_config_directory, internal_config_file, log_directory, main_log_file, root, try_detect_capfile

Constructor Details

#initialize(job_manager) ⇒ CelluloidManager

Returns a new instance of CelluloidManager.



18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 18

def initialize(job_manager)
  # start SupervisionGroup
  @worker_supervisor = Celluloid::SupervisionGroup.run!
  @job_manager = job_manager
  @registration_complete = false
  # Get a handle on the SupervisionGroup::Member
  @mutex = Mutex.new
  # http://rubydoc.info/gems/celluloid/Celluloid/SupervisionGroup/Member
  @workers = @worker_supervisor.pool(CapistranoMulticonfigParallel::CelluloidWorker, as: :workers, size: 10)
  Actor.current.link @workers
  @worker_supervisor.supervise_as(:terminal_server, CapistranoMulticonfigParallel::TerminalTable, Actor.current, @job_manager)
  @worker_supervisor.supervise_as(:web_server, CapistranoMulticonfigParallel::WebServer, websocket_config)

  # Get a handle on the PoolManager
  # http://rubydoc.info/gems/celluloid/Celluloid/PoolManager
  # @workers = workers_pool.actor
  @conditions = []
  @jobs = {}
  @job_to_worker = {}
  @worker_to_job = {}
  @job_to_condition = {}
end

Instance Attribute Details

#job_to_condition ⇒ Object

Returns the value of attribute job_to_condition.



13
14
15
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 13

def job_to_condition
  @job_to_condition
end

#job_to_worker ⇒ Object

Returns the value of attribute job_to_worker.



13
14
15
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 13

def job_to_worker
  @job_to_worker
end

#jobs ⇒ Object

Returns the value of attribute jobs.



13
14
15
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 13

def jobs
  @jobs
end

#mutex ⇒ Object

Returns the value of attribute mutex.



13
14
15
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 13

def mutex
  @mutex
end

#registration_complete ⇒ Object

Returns the value of attribute registration_complete.



13
14
15
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 13

def registration_complete
  @registration_complete
end

#worker_supervisor ⇒ Object (readonly)

Returns the value of attribute worker_supervisor.



15
16
17
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 15

def worker_supervisor
  @worker_supervisor
end

#worker_to_job ⇒ Object

Returns the value of attribute worker_to_job.



13
14
15
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 13

def worker_to_job
  @worker_to_job
end

#workers ⇒ Object (readonly)

Returns the value of attribute workers.



15
16
17
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 15

def workers
  @workers
end

#workers_terminated ⇒ Object

Returns the value of attribute workers_terminated.



13
14
15
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 13

def workers_terminated
  @workers_terminated
end

Instance Method Details

#all_workers_finished? ⇒ Boolean

Returns:

  • (Boolean)


62
63
64
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 62

def all_workers_finished?
  @job_to_worker.all? { |job_id, _worker| @jobs[job_id].finished? }
end

#apply_confirmation_for_job(job) ⇒ Object



91
92
93
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 91

def apply_confirmation_for_job(job)
  app_configuration.apply_stage_confirmation.include?(job.stage) && apply_confirmations?
end

#apply_confirmations? ⇒ Boolean

Returns:

  • (Boolean)


82
83
84
85
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 82

def apply_confirmations?
  confirmations = app_configuration.task_confirmations
  confirmations.is_a?(Array) && confirmations.present?
end

#can_tag_staging? ⇒ Boolean

Returns:

  • (Boolean)


183
184
185
186
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 183

def can_tag_staging?
  @job_manager.can_tag_staging? &&
    @jobs.find { |_job_id, job| job['env'] == 'production' }.blank?
end

#confirm_task_approval(result, task, processed_job = nil) ⇒ Object



156
157
158
159
160
161
162
163
164
165
166
167
168
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 156

def confirm_task_approval(result, task, processed_job = nil)
  return unless result.present?
  result = print_confirm_task_approvall(result, task, processed_job)
  return if result.blank? || result.downcase != 'y'
  @jobs.pmap do |job_id, job|
    worker = get_worker_for_job(job_id)
    worker.publish_rake_event('approved' => 'yes',
                              'action' => 'invoke',
                              'job_id' => job.id,
                              'task' => task
                             )
  end
end

#delegate(job) ⇒ Object

call to send an actor a job



43
44
45
46
47
48
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 43

def delegate(job)
  @jobs[job.id] = job
  # debug(@jobs)
  # start work and send it to the background
  @workers.async.work(job, Actor.current)
end

#dispatch_new_job(job, options = {}) ⇒ Object



188
189
190
191
192
193
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 188

def dispatch_new_job(job, options = {})
  env_opts = @job_manager.get_app_additional_env_options(job.app, job.stage)
  job.env_options = options.merge(env_opts)
  new_job = CapistranoMulticonfigParallel::Job.new(job.to_s)
  async.delegate(new_job)
end

#get_job_status(job) ⇒ Object

lookup status of job by asking actor running it



196
197
198
199
200
201
202
203
204
205
206
207
208
209
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 196

def get_job_status(job)
  status = nil
  if job.present?
    if job.is_a?(Hash)
      job = job.stringify_keys
      actor = @job_to_worker[job.id]
      status = actor.status
    else
      actor = @job_to_worker[job]
      status = actor.status
    end
  end
  status
end

#get_worker_for_job(job) ⇒ Object



170
171
172
173
174
175
176
177
178
179
180
181
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 170

def get_worker_for_job(job)
  if job.present?
    if job.is_a?(Hash)
      job = job.stringify_keys
      @job_to_worker[job.id]
    else
      @job_to_worker[job]
    end
  else
    return nil
  end
end

#job_crashed?(job) ⇒ Boolean

Returns:

  • (Boolean)


211
212
213
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 211

def job_crashed?(job)
  job.action == 'deploy:rollback' || job.action == 'deploy:failed' || job_failed?(job)
end

#job_failed?(job) ⇒ Boolean

Returns:

  • (Boolean)


215
216
217
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 215

def job_failed?(job)
  job.status.present? && job.status == 'worker_died'
end

#mark_completed_remaining_tasks(job) ⇒ Object



104
105
106
107
108
109
110
111
112
113
114
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 104

def mark_completed_remaining_tasks(job)
  return unless apply_confirmation_for_job(job)
  app_configuration.task_confirmations.each_with_index do |task, _index|
    fake_result = proc { |sum| sum }
    task_confirmation = @job_to_condition[job.id][task]
    if task_confirmation[:status] != 'confirmed'
      task_confirmation[:status] = 'confirmed'
      task_confirmation[:condition].signal(fake_result)
    end
  end
end


144
145
146
147
148
149
150
151
152
153
154
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 144

def print_confirm_task_approvall(result, task, job)
  return if result.is_a?(Proc)
  message = "Do you want  to continue the deployment and execute #{task.upcase}"
  message += " for JOB #{job.id}" if job.present?
  message += '?'
  apps_symlink_confirmation = Celluloid::Actor[:terminal_server].show_confirmation(message, 'Y/N')
  until apps_symlink_confirmation.present?
    sleep(0.1) # keep current thread alive
  end
  apps_symlink_confirmation
end

#process_jobs ⇒ Object



66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 66

def process_jobs
  @workers_terminated = Celluloid::Condition.new
  if syncronized_confirmation?
    @job_to_worker.pmap do |_job_id, worker|
      worker.async.start_task
    end
    wait_task_confirmations
  end
  condition = @workers_terminated.wait
  until condition.present?
    sleep(0.1) # keep current thread alive
  end
  log_to_file("all jobs have completed #{condition}")
  Celluloid::Actor[:terminal_server].async.notify_time_change(CapistranoMulticonfigParallel::TerminalTable.topic, type: 'output') if Celluloid::Actor[:terminal_server].alive?
end

#register_worker_for_job(job, worker) ⇒ Object

call back from actor once it has received it's job actor should do this asap



52
53
54
55
56
57
58
59
60
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 52

def register_worker_for_job(job, worker)
  worker.job_id = job.id if worker.job_id.blank?
  @job_to_worker[job.id] = worker
  @worker_to_job[worker.mailbox.address] = job
  log_to_file("worker #{worker.job_id} registed into manager")
  Actor.current.link worker
  worker.async.start_task unless syncronized_confirmation?
  @registration_complete = true if @job_manager.jobs.size == @job_to_worker.size
end

#setup_worker_conditions(job) ⇒ Object



95
96
97
98
99
100
101
102
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 95

def setup_worker_conditions(job)
  return unless apply_confirmation_for_job(job)
  hash_conditions = {}
  app_configuration.task_confirmations.each do |task|
    hash_conditions[task] = { condition: Celluloid::Condition.new, status: 'unconfirmed' }
  end
  @job_to_condition[job.id] = hash_conditions
end

#syncronized_confirmation? ⇒ Boolean

Returns:

  • (Boolean)


87
88
89
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 87

def syncronized_confirmation?
  !@job_manager.can_tag_staging?
end

#wait_condition_for_task(job_id, task) ⇒ Object



125
126
127
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 125

def wait_condition_for_task(job_id, task)
  @job_to_condition[job_id][task][:condition].wait
end

#wait_task_confirmations ⇒ Object



129
130
131
132
133
134
135
136
137
138
139
140
141
142
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 129

def wait_task_confirmations
  stage_apply = app_configuration.apply_stage_confirmation.include?(@job_manager.stage)
  return if !stage_apply || !syncronized_confirmation?
  app_configuration.task_confirmations.each_with_index do |task, _index|
    results = []
    @jobs.pmap do |job_id, _job|
      result = wait_condition_for_task(job_id, task)
      results << result
    end
    if results.size == @jobs.size
      confirm_task_approval(results, task)
    end
  end
end

#wait_task_confirmations_worker(job) ⇒ Object



116
117
118
119
120
121
122
123
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 116

def wait_task_confirmations_worker(job)
  return unless job.finished? || job.exit_status.present?
  return if !apply_confirmation_for_job(job) || !syncronized_confirmation?
  app_configuration.task_confirmations.each_with_index do |task, _index|
    result = wait_condition_for_task(job.id, task)
    confirm_task_approval(result, task, job) if result.present?
  end
end

#worker_died(worker, reason) ⇒ Object



219
220
221
222
223
224
225
226
227
228
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 219

def worker_died(worker, reason)
  job = @worker_to_job[worker.mailbox.address]
  log_to_file("worker job #{job} with mailbox #{worker.mailbox.inspect} died  for reason:  #{reason}")
  @worker_to_job.delete(worker.mailbox.address)
  return if job.blank? || job_crashed?(job)
  return unless job.action == 'deploy'
  log_to_file "restarting #{job} on new worker"
  job.status = 'worker_died'
  dispatch_new_job(job, action: 'deploy:rollback')
end