Class: CapistranoMulticonfigParallel::CelluloidWorker

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

Overview

rubocop:disable ClassLength worker that will spawn a child process in order to execute a capistrano job and monitor that process

Defined Under Namespace

Classes: TaskFailed

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, 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

Instance Attribute Details

#action_name ⇒ Object

Returns the value of attribute action_name.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def action_name
  @action_name
end

#app_name ⇒ Object

Returns the value of attribute app_name.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def app_name
  @app_name
end

#client ⇒ Object

Returns the value of attribute client.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def client
  @client
end

#current_task_number ⇒ Object

Returns the value of attribute current_task_number.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def current_task_number
  @current_task_number
end

#env_name ⇒ Object

Returns the value of attribute env_name.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def env_name
  @env_name
end

#env_options ⇒ Object

Returns the value of attribute env_options.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def env_options
  @env_options
end

#filename ⇒ Object

Returns the value of attribute filename.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def filename
  @filename
end

#invocation_chain ⇒ Object

Returns the value of attribute invocation_chain.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def invocation_chain
  @invocation_chain
end

#job ⇒ Hash

Returns options used for executing capistrano task.

Returns:

  • (Hash) —

    options used for executing capistrano task



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
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 20

class CelluloidWorker
  include Celluloid
  include Celluloid::Notifications
  include Celluloid::Logger
  include CapistranoMulticonfigParallel::ApplicationHelper
  class TaskFailed < StandardError; end

  attr_accessor :job, :manager, :job_id, :app_name, :env_name, :action_name, :env_options, :machine, :client, :task_argv,
                :rake_tasks, :current_task_number, # tracking tasks
                :successfull_subscription, :subscription_channel, :publisher_channel, # for subscriptions and publishing events
                :job_termination_condition, :worker_state, :invocation_chain, :filename, :worker_log

  def work(job, manager)
    @job = job
    @worker_state = 'started'
    @manager = manager
    @job_confirmation_conditions = []
    process_job(job) if job.present?
    log_to_file("worker #{@job_id} received #{job.inspect}")
    @subscription_channel = "worker_#{@job_id}"
    @machine = CapistranoMulticonfigParallel::StateMachine.new(job, Actor.current)
    manager.register_worker_for_job(job, Actor.current)
  end

  def start_task
    @manager.setup_worker_conditions(Actor.current)
    log_to_file("exec worker #{@job_id} starts task with #{@job.inspect}")
    @client = CelluloidPubsub::Client.connect(actor: Actor.current, enable_debug: debug_websocket?, channel: subscription_channel)
  end

  def publish_rake_event(data)
    @client.publish(rake_actor_id(data), data)
  end

  def rake_actor_id(data)
    data['action'].present? && data['action'] == 'count' ? "rake_worker_#{@job_id}_count" : "rake_worker_#{@job_id}"
  end

  def on_message(message)
    log_to_file("worker #{@job_id} received:  #{message.inspect}")
    if @client.succesfull_subscription?(message)
      @successfull_subscription = true
      execute_after_succesfull_subscription
    else
      handle_subscription(message)
    end
  end

  def execute_after_succesfull_subscription
    async.execute_deploy
  end

  def rake_tasks
    @rake_tasks ||= []
  end

  def invocation_chain
    @invocation_chain ||= []
  end

  def cd_working_directory
    "cd #{detect_root}"
  end

  # def generate_command_new
  #   <<-CMD
  #     bundle exec ruby -e "require 'bundler' ;   Bundler.with_clean_env { %x[cd #{cd_working_directory} && bundle install && RAILS_ENV=#{@env_name} bundle exec cap #{@task_argv.join(' ')}] } "
  #   CMD
  # end

  def generate_command
    <<-CMD
    #{cd_working_directory} && RAILS_ENV=#{@env_name} bundle exec multi_cap #{@task_argv.join(' ')}
    CMD
  end

  def execute_deploy
    log_to_file("invocation chain #{@job_id} is : #{@rake_tasks.inspect}")
    check_child_proces
    setup_task_arguments
    log_to_file("worker #{@job_id} executes: #{generate_command}")
    @child_process.async.work(generate_command, actor: Actor.current, silent: true)
    @manager.wait_task_confirmations_worker(Actor.current)
  end

  def check_child_proces
    if !defined?(@child_process) || @child_process.nil?
      @child_process = CapistranoMulticonfigParallel::ChildProcess.new
      Actor.current.link @child_process
    else
      @client.unsubscribe("rake_worker_#{@job_id}_count")
      @child_process.exit_status = nil
    end
  end

  def on_close(code, reason)
    log_to_file("worker #{@job_id} websocket connection closed: #{code.inspect}, #{reason.inspect}")
  end

  def check_gitflow
    return if @env_name != 'staging' || !@manager.can_tag_staging? || !executed_task?(CapistranoMulticonfigParallel::GITFLOW_TAG_STAGING_TASK)
    @manager.dispatch_new_job(@job.merge('env' => 'production'))
  end

  def handle_subscription(message)
    if message_is_about_a_task?(message)
      check_gitflow
      save_tasks_to_be_executed(message)
      update_machine_state(message['task']) # if message['action'] == 'invoke'
      log_to_file("worker #{@job_id} state is #{@machine.state}")
      task_approval(message)
    elsif message_is_for_stdout?(message)
      result = Celluloid::Actor[:terminal_server].show_confirmation(message['question'], message['default'])
      publish_rake_event(message.merge('action' => 'stdin', 'result' => result, 'client_action' => 'stdin'))
    else
      log_to_file(message, @job_id)
    end
  end

  def message_is_for_stdout?(message)
    message.present? && message.is_a?(Hash) && message['action'].present? && message['job_id'].present? && message['action'] == 'stdout'
  end

  def message_is_about_a_task?(message)
    message.present? && message.is_a?(Hash) && message['action'].present? && message['job_id'].present? && message['task'].present?
  end

  def executed_task?(task)
    rake_tasks.present? && rake_tasks.index(task.to_s).present?
  end

  def task_approval(message)
    job_conditions = @manager.job_to_condition[@job_id]
    if job_conditions.present? && app_configuration.task_confirmations.include?(message['task']) && message['action'] == 'invoke'
      task_confirmation = job_conditions[message['task']]
      task_confirmation[:status] = 'confirmed'
      task_confirmation[:condition].signal(message['task'])
    else
      publish_rake_event(message.merge('approved' => 'yes'))
    end
  end

  def save_tasks_to_be_executed(message)
    log_to_file("worler #{@job_id} current invocation chain : #{rake_tasks.inspect}")
    rake_tasks << message['task'] if rake_tasks.last != message['task']
    invocation_chain << message['task'] if invocation_chain.last != message['task']
  end

  def update_machine_state(name)
    log_to_file("worker #{@job_id} triest to transition from #{@machine.state} to  #{name}")
    @machine.transitions.on(name.to_s, @machine.state => name.to_s)
    @machine.go_to_transition(name.to_s)
    abort(CapistranoMulticonfigParallel::CelluloidWorker::TaskFailed, "task #{@action} failed ") if name == 'deploy:failed' # force worker to rollback
  end

  def setup_command_line(*options)
    @task_argv = []
    options.each do |option|
      @task_argv << option
    end
    @task_argv
  end

  def worker_stage
    @app_name.present? ? "#{@app_name}:#{@env_name}" : "#{@env_name}"
  end

  def worker_action
    "#{@action_name}[#{@task_arguments.join(',')}]"
  end

  def setup_task_arguments(*args)
    #   stage = "#{@app_name}:#{@env_name} #{@action_name}"
    array_options = []
    @env_options.each do |key, value|
      array_options << "#{key}=#{value}" if value.present?
    end
    array_options << '--trace' if app_debug_enabled?
    args.each do |arg|
      array_options << arg
    end
    @manager.jobs[@job_id]['job_argv'] = array_options.clone
    array_options.unshift("#{worker_action}")
    array_options.unshift("#{worker_stage}")
    setup_command_line(*array_options)
  end

  def send_msg(channel, message = nil)
    publish channel, message.present? && message.is_a?(Hash) ? { job_id: @job_id }.merge(message) : { job_id: @job_id, time: Time.now }
  end

  def process_job(job)
    processed_job = @manager.process_job(job)
    @job_id = processed_job['job_id']
    @app_name = processed_job['app_name']
    @env_name = processed_job['env_name']
    @action_name = processed_job['action_name']
    @env_options = processed_job['env_options']
    @task_arguments = processed_job['task_arguments']
  end

  def crashed?
    @action_name == 'deploy:rollback' || @action_name == 'deploy:failed' || @manager.job_failed?(@job)
  end

  def finish_worker
    @manager.mark_completed_remaining_tasks(Actor.current)
    @manager.jobs[@job_id]['worker_action'] = 'finished'
    @manager.workers_terminated.signal('completed') if @manager.all_workers_finished?
  end

  def worker_finshed?
    @manager.jobs[@job_id]['worker_action'] == 'finished'
  end

  def notify_finished(exit_status)
    if exit_status.exitstatus != 0
      log_to_file("worker #{job_id} tries to terminate")
      abort(CapistranoMulticonfigParallel::CelluloidWorker::TaskFailed, "task  failed with exit status #{exit_status.inspect} ") # force worker to rollback
    else
      update_machine_state('FINISHED')
      log_to_file("worker #{job_id} notifies manager has finished")
      finish_worker
    end
  end
end

#job_id ⇒ Object

Returns the value of attribute job_id.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def job_id
  @job_id
end

#job_termination_condition ⇒ Object

Returns the value of attribute job_termination_condition.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def job_termination_condition
  @job_termination_condition
end

#machine ⇒ Object

Returns the value of attribute machine.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def machine
  @machine
end

#manager ⇒ CapistranoMulticonfigParallel::CelluloidManager

Returns the instance of the manager that delegated the job to this worker.

Returns:



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
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 20

class CelluloidWorker
  include Celluloid
  include Celluloid::Notifications
  include Celluloid::Logger
  include CapistranoMulticonfigParallel::ApplicationHelper
  class TaskFailed < StandardError; end

  attr_accessor :job, :manager, :job_id, :app_name, :env_name, :action_name, :env_options, :machine, :client, :task_argv,
                :rake_tasks, :current_task_number, # tracking tasks
                :successfull_subscription, :subscription_channel, :publisher_channel, # for subscriptions and publishing events
                :job_termination_condition, :worker_state, :invocation_chain, :filename, :worker_log

  def work(job, manager)
    @job = job
    @worker_state = 'started'
    @manager = manager
    @job_confirmation_conditions = []
    process_job(job) if job.present?
    log_to_file("worker #{@job_id} received #{job.inspect}")
    @subscription_channel = "worker_#{@job_id}"
    @machine = CapistranoMulticonfigParallel::StateMachine.new(job, Actor.current)
    manager.register_worker_for_job(job, Actor.current)
  end

  def start_task
    @manager.setup_worker_conditions(Actor.current)
    log_to_file("exec worker #{@job_id} starts task with #{@job.inspect}")
    @client = CelluloidPubsub::Client.connect(actor: Actor.current, enable_debug: debug_websocket?, channel: subscription_channel)
  end

  def publish_rake_event(data)
    @client.publish(rake_actor_id(data), data)
  end

  def rake_actor_id(data)
    data['action'].present? && data['action'] == 'count' ? "rake_worker_#{@job_id}_count" : "rake_worker_#{@job_id}"
  end

  def on_message(message)
    log_to_file("worker #{@job_id} received:  #{message.inspect}")
    if @client.succesfull_subscription?(message)
      @successfull_subscription = true
      execute_after_succesfull_subscription
    else
      handle_subscription(message)
    end
  end

  def execute_after_succesfull_subscription
    async.execute_deploy
  end

  def rake_tasks
    @rake_tasks ||= []
  end

  def invocation_chain
    @invocation_chain ||= []
  end

  def cd_working_directory
    "cd #{detect_root}"
  end

  # def generate_command_new
  #   <<-CMD
  #     bundle exec ruby -e "require 'bundler' ;   Bundler.with_clean_env { %x[cd #{cd_working_directory} && bundle install && RAILS_ENV=#{@env_name} bundle exec cap #{@task_argv.join(' ')}] } "
  #   CMD
  # end

  def generate_command
    <<-CMD
    #{cd_working_directory} && RAILS_ENV=#{@env_name} bundle exec multi_cap #{@task_argv.join(' ')}
    CMD
  end

  def execute_deploy
    log_to_file("invocation chain #{@job_id} is : #{@rake_tasks.inspect}")
    check_child_proces
    setup_task_arguments
    log_to_file("worker #{@job_id} executes: #{generate_command}")
    @child_process.async.work(generate_command, actor: Actor.current, silent: true)
    @manager.wait_task_confirmations_worker(Actor.current)
  end

  def check_child_proces
    if !defined?(@child_process) || @child_process.nil?
      @child_process = CapistranoMulticonfigParallel::ChildProcess.new
      Actor.current.link @child_process
    else
      @client.unsubscribe("rake_worker_#{@job_id}_count")
      @child_process.exit_status = nil
    end
  end

  def on_close(code, reason)
    log_to_file("worker #{@job_id} websocket connection closed: #{code.inspect}, #{reason.inspect}")
  end

  def check_gitflow
    return if @env_name != 'staging' || !@manager.can_tag_staging? || !executed_task?(CapistranoMulticonfigParallel::GITFLOW_TAG_STAGING_TASK)
    @manager.dispatch_new_job(@job.merge('env' => 'production'))
  end

  def handle_subscription(message)
    if message_is_about_a_task?(message)
      check_gitflow
      save_tasks_to_be_executed(message)
      update_machine_state(message['task']) # if message['action'] == 'invoke'
      log_to_file("worker #{@job_id} state is #{@machine.state}")
      task_approval(message)
    elsif message_is_for_stdout?(message)
      result = Celluloid::Actor[:terminal_server].show_confirmation(message['question'], message['default'])
      publish_rake_event(message.merge('action' => 'stdin', 'result' => result, 'client_action' => 'stdin'))
    else
      log_to_file(message, @job_id)
    end
  end

  def message_is_for_stdout?(message)
    message.present? && message.is_a?(Hash) && message['action'].present? && message['job_id'].present? && message['action'] == 'stdout'
  end

  def message_is_about_a_task?(message)
    message.present? && message.is_a?(Hash) && message['action'].present? && message['job_id'].present? && message['task'].present?
  end

  def executed_task?(task)
    rake_tasks.present? && rake_tasks.index(task.to_s).present?
  end

  def task_approval(message)
    job_conditions = @manager.job_to_condition[@job_id]
    if job_conditions.present? && app_configuration.task_confirmations.include?(message['task']) && message['action'] == 'invoke'
      task_confirmation = job_conditions[message['task']]
      task_confirmation[:status] = 'confirmed'
      task_confirmation[:condition].signal(message['task'])
    else
      publish_rake_event(message.merge('approved' => 'yes'))
    end
  end

  def save_tasks_to_be_executed(message)
    log_to_file("worler #{@job_id} current invocation chain : #{rake_tasks.inspect}")
    rake_tasks << message['task'] if rake_tasks.last != message['task']
    invocation_chain << message['task'] if invocation_chain.last != message['task']
  end

  def update_machine_state(name)
    log_to_file("worker #{@job_id} triest to transition from #{@machine.state} to  #{name}")
    @machine.transitions.on(name.to_s, @machine.state => name.to_s)
    @machine.go_to_transition(name.to_s)
    abort(CapistranoMulticonfigParallel::CelluloidWorker::TaskFailed, "task #{@action} failed ") if name == 'deploy:failed' # force worker to rollback
  end

  def setup_command_line(*options)
    @task_argv = []
    options.each do |option|
      @task_argv << option
    end
    @task_argv
  end

  def worker_stage
    @app_name.present? ? "#{@app_name}:#{@env_name}" : "#{@env_name}"
  end

  def worker_action
    "#{@action_name}[#{@task_arguments.join(',')}]"
  end

  def setup_task_arguments(*args)
    #   stage = "#{@app_name}:#{@env_name} #{@action_name}"
    array_options = []
    @env_options.each do |key, value|
      array_options << "#{key}=#{value}" if value.present?
    end
    array_options << '--trace' if app_debug_enabled?
    args.each do |arg|
      array_options << arg
    end
    @manager.jobs[@job_id]['job_argv'] = array_options.clone
    array_options.unshift("#{worker_action}")
    array_options.unshift("#{worker_stage}")
    setup_command_line(*array_options)
  end

  def send_msg(channel, message = nil)
    publish channel, message.present? && message.is_a?(Hash) ? { job_id: @job_id }.merge(message) : { job_id: @job_id, time: Time.now }
  end

  def process_job(job)
    processed_job = @manager.process_job(job)
    @job_id = processed_job['job_id']
    @app_name = processed_job['app_name']
    @env_name = processed_job['env_name']
    @action_name = processed_job['action_name']
    @env_options = processed_job['env_options']
    @task_arguments = processed_job['task_arguments']
  end

  def crashed?
    @action_name == 'deploy:rollback' || @action_name == 'deploy:failed' || @manager.job_failed?(@job)
  end

  def finish_worker
    @manager.mark_completed_remaining_tasks(Actor.current)
    @manager.jobs[@job_id]['worker_action'] = 'finished'
    @manager.workers_terminated.signal('completed') if @manager.all_workers_finished?
  end

  def worker_finshed?
    @manager.jobs[@job_id]['worker_action'] == 'finished'
  end

  def notify_finished(exit_status)
    if exit_status.exitstatus != 0
      log_to_file("worker #{job_id} tries to terminate")
      abort(CapistranoMulticonfigParallel::CelluloidWorker::TaskFailed, "task  failed with exit status #{exit_status.inspect} ") # force worker to rollback
    else
      update_machine_state('FINISHED')
      log_to_file("worker #{job_id} notifies manager has finished")
      finish_worker
    end
  end
end

#publisher_channel ⇒ Object

Returns the value of attribute publisher_channel.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def publisher_channel
  @publisher_channel
end

#rake_tasks ⇒ Object

Returns the value of attribute rake_tasks.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def rake_tasks
  @rake_tasks
end

#subscription_channel ⇒ Object

Returns the value of attribute subscription_channel.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def subscription_channel
  @subscription_channel
end

#successfull_subscription ⇒ Object

Returns the value of attribute successfull_subscription.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def successfull_subscription
  @successfull_subscription
end

#task_argv ⇒ Object

Returns the value of attribute task_argv.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def task_argv
  @task_argv
end

#worker_log ⇒ Object

Returns the value of attribute worker_log.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def worker_log
  @worker_log
end

#worker_state ⇒ Object

Returns the value of attribute worker_state.



27
28
29
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 27

def worker_state
  @worker_state
end

Instance Method Details

#cd_working_directory ⇒ Object



80
81
82
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 80

def cd_working_directory
  "cd #{detect_root}"
end

#check_child_proces ⇒ Object



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

def check_child_proces
  if !defined?(@child_process) || @child_process.nil?
    @child_process = CapistranoMulticonfigParallel::ChildProcess.new
    Actor.current.link @child_process
  else
    @client.unsubscribe("rake_worker_#{@job_id}_count")
    @child_process.exit_status = nil
  end
end

#check_gitflow ⇒ Object



119
120
121
122
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 119

def check_gitflow
  return if @env_name != 'staging' || !@manager.can_tag_staging? || !executed_task?(CapistranoMulticonfigParallel::GITFLOW_TAG_STAGING_TASK)
  @manager.dispatch_new_job(@job.merge('env' => 'production'))
end

#crashed? ⇒ Boolean

Returns:

  • (Boolean)


221
222
223
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 221

def crashed?
  @action_name == 'deploy:rollback' || @action_name == 'deploy:failed' || @manager.job_failed?(@job)
end

#execute_after_succesfull_subscription ⇒ Object



68
69
70
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 68

def execute_after_succesfull_subscription
  async.execute_deploy
end

#execute_deploy ⇒ Object



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

def execute_deploy
  log_to_file("invocation chain #{@job_id} is : #{@rake_tasks.inspect}")
  check_child_proces
  setup_task_arguments
  log_to_file("worker #{@job_id} executes: #{generate_command}")
  @child_process.async.work(generate_command, actor: Actor.current, silent: true)
  @manager.wait_task_confirmations_worker(Actor.current)
end

#executed_task?(task) ⇒ Boolean

Returns:

  • (Boolean)


147
148
149
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 147

def executed_task?(task)
  rake_tasks.present? && rake_tasks.index(task.to_s).present?
end

#finish_worker ⇒ Object



225
226
227
228
229
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 225

def finish_worker
  @manager.mark_completed_remaining_tasks(Actor.current)
  @manager.jobs[@job_id]['worker_action'] = 'finished'
  @manager.workers_terminated.signal('completed') if @manager.all_workers_finished?
end

#generate_command ⇒ Object

def generate_command_new <<-CMD bundle exec ruby -e "require 'bundler' ; Bundler.with_clean_env { %x[cd ##cd_working_directory && bundle install && RAILS_ENV=#@env_name bundle exec cap #')] } " CMD end



90
91
92
93
94
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 90

def generate_command
  <<-CMD
  #{cd_working_directory} && RAILS_ENV=#{@env_name} bundle exec multi_cap #{@task_argv.join(' ')}
  CMD
end

#handle_subscription(message) ⇒ Object



124
125
126
127
128
129
130
131
132
133
134
135
136
137
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 124

def handle_subscription(message)
  if message_is_about_a_task?(message)
    check_gitflow
    save_tasks_to_be_executed(message)
    update_machine_state(message['task']) # if message['action'] == 'invoke'
    log_to_file("worker #{@job_id} state is #{@machine.state}")
    task_approval(message)
  elsif message_is_for_stdout?(message)
    result = Celluloid::Actor[:terminal_server].show_confirmation(message['question'], message['default'])
    publish_rake_event(message.merge('action' => 'stdin', 'result' => result, 'client_action' => 'stdin'))
  else
    log_to_file(message, @job_id)
  end
end

#message_is_about_a_task?(message) ⇒ Boolean

Returns:

  • (Boolean)


143
144
145
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 143

def message_is_about_a_task?(message)
  message.present? && message.is_a?(Hash) && message['action'].present? && message['job_id'].present? && message['task'].present?
end

#message_is_for_stdout?(message) ⇒ Boolean

Returns:

  • (Boolean)


139
140
141
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 139

def message_is_for_stdout?(message)
  message.present? && message.is_a?(Hash) && message['action'].present? && message['job_id'].present? && message['action'] == 'stdout'
end

#notify_finished(exit_status) ⇒ Object



235
236
237
238
239
240
241
242
243
244
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 235

def notify_finished(exit_status)
  if exit_status.exitstatus != 0
    log_to_file("worker #{job_id} tries to terminate")
    abort(CapistranoMulticonfigParallel::CelluloidWorker::TaskFailed, "task  failed with exit status #{exit_status.inspect} ") # force worker to rollback
  else
    update_machine_state('FINISHED')
    log_to_file("worker #{job_id} notifies manager has finished")
    finish_worker
  end
end

#on_close(code, reason) ⇒ Object



115
116
117
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 115

def on_close(code, reason)
  log_to_file("worker #{@job_id} websocket connection closed: #{code.inspect}, #{reason.inspect}")
end

#on_message(message) ⇒ Object



58
59
60
61
62
63
64
65
66
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 58

def on_message(message)
  log_to_file("worker #{@job_id} received:  #{message.inspect}")
  if @client.succesfull_subscription?(message)
    @successfull_subscription = true
    execute_after_succesfull_subscription
  else
    handle_subscription(message)
  end
end

#process_job(job) ⇒ Object



211
212
213
214
215
216
217
218
219
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 211

def process_job(job)
  processed_job = @manager.process_job(job)
  @job_id = processed_job['job_id']
  @app_name = processed_job['app_name']
  @env_name = processed_job['env_name']
  @action_name = processed_job['action_name']
  @env_options = processed_job['env_options']
  @task_arguments = processed_job['task_arguments']
end

#publish_rake_event(data) ⇒ Object



50
51
52
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 50

def publish_rake_event(data)
  @client.publish(rake_actor_id(data), data)
end

#rake_actor_id(data) ⇒ Object



54
55
56
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 54

def rake_actor_id(data)
  data['action'].present? && data['action'] == 'count' ? "rake_worker_#{@job_id}_count" : "rake_worker_#{@job_id}"
end

#save_tasks_to_be_executed(message) ⇒ Object



162
163
164
165
166
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 162

def save_tasks_to_be_executed(message)
  log_to_file("worler #{@job_id} current invocation chain : #{rake_tasks.inspect}")
  rake_tasks << message['task'] if rake_tasks.last != message['task']
  invocation_chain << message['task'] if invocation_chain.last != message['task']
end

#send_msg(channel, message = nil) ⇒ Object



207
208
209
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 207

def send_msg(channel, message = nil)
  publish channel, message.present? && message.is_a?(Hash) ? { job_id: @job_id }.merge(message) : { job_id: @job_id, time: Time.now }
end

#setup_command_line(*options) ⇒ Object



175
176
177
178
179
180
181
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 175

def setup_command_line(*options)
  @task_argv = []
  options.each do |option|
    @task_argv << option
  end
  @task_argv
end

#setup_task_arguments(*args) ⇒ Object



191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 191

def setup_task_arguments(*args)
  #   stage = "#{@app_name}:#{@env_name} #{@action_name}"
  array_options = []
  @env_options.each do |key, value|
    array_options << "#{key}=#{value}" if value.present?
  end
  array_options << '--trace' if app_debug_enabled?
  args.each do |arg|
    array_options << arg
  end
  @manager.jobs[@job_id]['job_argv'] = array_options.clone
  array_options.unshift("#{worker_action}")
  array_options.unshift("#{worker_stage}")
  setup_command_line(*array_options)
end

#start_task ⇒ Object



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

def start_task
  @manager.setup_worker_conditions(Actor.current)
  log_to_file("exec worker #{@job_id} starts task with #{@job.inspect}")
  @client = CelluloidPubsub::Client.connect(actor: Actor.current, enable_debug: debug_websocket?, channel: subscription_channel)
end

#task_approval(message) ⇒ Object



151
152
153
154
155
156
157
158
159
160
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 151

def task_approval(message)
  job_conditions = @manager.job_to_condition[@job_id]
  if job_conditions.present? && app_configuration.task_confirmations.include?(message['task']) && message['action'] == 'invoke'
    task_confirmation = job_conditions[message['task']]
    task_confirmation[:status] = 'confirmed'
    task_confirmation[:condition].signal(message['task'])
  else
    publish_rake_event(message.merge('approved' => 'yes'))
  end
end

#update_machine_state(name) ⇒ Object



168
169
170
171
172
173
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 168

def update_machine_state(name)
  log_to_file("worker #{@job_id} triest to transition from #{@machine.state} to  #{name}")
  @machine.transitions.on(name.to_s, @machine.state => name.to_s)
  @machine.go_to_transition(name.to_s)
  abort(CapistranoMulticonfigParallel::CelluloidWorker::TaskFailed, "task #{@action} failed ") if name == 'deploy:failed' # force worker to rollback
end

#work(job, manager) ⇒ Object



32
33
34
35
36
37
38
39
40
41
42
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 32

def work(job, manager)
  @job = job
  @worker_state = 'started'
  @manager = manager
  @job_confirmation_conditions = []
  process_job(job) if job.present?
  log_to_file("worker #{@job_id} received #{job.inspect}")
  @subscription_channel = "worker_#{@job_id}"
  @machine = CapistranoMulticonfigParallel::StateMachine.new(job, Actor.current)
  manager.register_worker_for_job(job, Actor.current)
end

#worker_action ⇒ Object



187
188
189
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 187

def worker_action
  "#{@action_name}[#{@task_arguments.join(',')}]"
end

#worker_finshed? ⇒ Boolean

Returns:

  • (Boolean)


231
232
233
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 231

def worker_finshed?
  @manager.jobs[@job_id]['worker_action'] == 'finished'
end

#worker_stage ⇒ Object



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

def worker_stage
  @app_name.present? ? "#{@app_name}:#{@env_name}" : "#{@env_name}"
end