Class: CapistranoMulticonfigParallel::CelluloidWorker
- Inherits:
-
Object
- Object
- CapistranoMulticonfigParallel::CelluloidWorker
- 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
-
#action_name ⇒ Object
Returns the value of attribute action_name.
-
#app_name ⇒ Object
Returns the value of attribute app_name.
-
#client ⇒ Object
Returns the value of attribute client.
-
#current_task_number ⇒ Object
Returns the value of attribute current_task_number.
-
#env_name ⇒ Object
Returns the value of attribute env_name.
-
#env_options ⇒ Object
Returns the value of attribute env_options.
-
#filename ⇒ Object
Returns the value of attribute filename.
-
#invocation_chain ⇒ Object
Returns the value of attribute invocation_chain.
-
#job ⇒ Hash
Options used for executing capistrano task.
-
#job_id ⇒ Object
Returns the value of attribute job_id.
-
#job_termination_condition ⇒ Object
Returns the value of attribute job_termination_condition.
-
#machine ⇒ Object
Returns the value of attribute machine.
-
#manager ⇒ CapistranoMulticonfigParallel::CelluloidManager
The instance of the manager that delegated the job to this worker.
-
#publisher_channel ⇒ Object
Returns the value of attribute publisher_channel.
-
#rake_tasks ⇒ Object
Returns the value of attribute rake_tasks.
-
#subscription_channel ⇒ Object
Returns the value of attribute subscription_channel.
-
#successfull_subscription ⇒ Object
Returns the value of attribute successfull_subscription.
-
#task_argv ⇒ Object
Returns the value of attribute task_argv.
-
#worker_log ⇒ Object
Returns the value of attribute worker_log.
-
#worker_state ⇒ Object
Returns the value of attribute worker_state.
Instance Method Summary collapse
- #cd_working_directory ⇒ Object
- #check_child_proces ⇒ Object
- #check_gitflow ⇒ Object
- #crashed? ⇒ Boolean
- #execute_after_succesfull_subscription ⇒ Object
- #execute_deploy ⇒ Object
- #executed_task?(task) ⇒ Boolean
- #finish_worker ⇒ Object
-
#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.
- #handle_subscription(message) ⇒ Object
- #message_is_about_a_task?(message) ⇒ Boolean
- #message_is_for_stdout?(message) ⇒ Boolean
- #notify_finished(exit_status) ⇒ Object
- #on_close(code, reason) ⇒ Object
- #on_message(message) ⇒ Object
- #process_job(job) ⇒ Object
- #publish_rake_event(data) ⇒ Object
- #rake_actor_id(data) ⇒ Object
- #save_tasks_to_be_executed(message) ⇒ Object
- #send_msg(channel, message = nil) ⇒ Object
- #setup_command_line(*options) ⇒ Object
- #setup_task_arguments(*args) ⇒ Object
- #start_task ⇒ Object
- #task_approval(message) ⇒ Object
- #update_machine_state(name) ⇒ Object
- #work(job, manager) ⇒ Object
- #worker_action ⇒ Object
- #worker_finshed? ⇒ Boolean
- #worker_stage ⇒ Object
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 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.
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 () log_to_file("worker #{@job_id} received: #{.inspect}") if @client.succesfull_subscription?() @successfull_subscription = true execute_after_succesfull_subscription else handle_subscription() 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() if () check_gitflow save_tasks_to_be_executed() update_machine_state(['task']) # if message['action'] == 'invoke' log_to_file("worker #{@job_id} state is #{@machine.state}") task_approval() elsif () result = Celluloid::Actor[:terminal_server].show_confirmation(['question'], ['default']) publish_rake_event(.merge('action' => 'stdin', 'result' => result, 'client_action' => 'stdin')) else log_to_file(, @job_id) end end def () .present? && .is_a?(Hash) && ['action'].present? && ['job_id'].present? && ['action'] == 'stdout' end def () .present? && .is_a?(Hash) && ['action'].present? && ['job_id'].present? && ['task'].present? end def executed_task?(task) rake_tasks.present? && rake_tasks.index(task.to_s).present? end def task_approval() job_conditions = @manager.job_to_condition[@job_id] if job_conditions.present? && app_configuration.task_confirmations.include?(['task']) && ['action'] == 'invoke' task_confirmation = job_conditions[['task']] task_confirmation[:status] = 'confirmed' task_confirmation[:condition].signal(['task']) else publish_rake_event(.merge('approved' => 'yes')) end end def save_tasks_to_be_executed() log_to_file("worler #{@job_id} current invocation chain : #{rake_tasks.inspect}") rake_tasks << ['task'] if rake_tasks.last != ['task'] invocation_chain << ['task'] if invocation_chain.last != ['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(*) @task_argv = [] .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}" = [] @env_options.each do |key, value| << "#{key}=#{value}" if value.present? end << '--trace' if app_debug_enabled? args.each do |arg| << arg end @manager.jobs[@job_id]['job_argv'] = .clone .unshift("#{worker_action}") .unshift("#{worker_stage}") setup_command_line(*) end def send_msg(channel, = nil) publish channel, .present? && .is_a?(Hash) ? { job_id: @job_id }.merge() : { 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.
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 () log_to_file("worker #{@job_id} received: #{.inspect}") if @client.succesfull_subscription?() @successfull_subscription = true execute_after_succesfull_subscription else handle_subscription() 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() if () check_gitflow save_tasks_to_be_executed() update_machine_state(['task']) # if message['action'] == 'invoke' log_to_file("worker #{@job_id} state is #{@machine.state}") task_approval() elsif () result = Celluloid::Actor[:terminal_server].show_confirmation(['question'], ['default']) publish_rake_event(.merge('action' => 'stdin', 'result' => result, 'client_action' => 'stdin')) else log_to_file(, @job_id) end end def () .present? && .is_a?(Hash) && ['action'].present? && ['job_id'].present? && ['action'] == 'stdout' end def () .present? && .is_a?(Hash) && ['action'].present? && ['job_id'].present? && ['task'].present? end def executed_task?(task) rake_tasks.present? && rake_tasks.index(task.to_s).present? end def task_approval() job_conditions = @manager.job_to_condition[@job_id] if job_conditions.present? && app_configuration.task_confirmations.include?(['task']) && ['action'] == 'invoke' task_confirmation = job_conditions[['task']] task_confirmation[:status] = 'confirmed' task_confirmation[:condition].signal(['task']) else publish_rake_event(.merge('approved' => 'yes')) end end def save_tasks_to_be_executed() log_to_file("worler #{@job_id} current invocation chain : #{rake_tasks.inspect}") rake_tasks << ['task'] if rake_tasks.last != ['task'] invocation_chain << ['task'] if invocation_chain.last != ['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(*) @task_argv = [] .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}" = [] @env_options.each do |key, value| << "#{key}=#{value}" if value.present? end << '--trace' if app_debug_enabled? args.each do |arg| << arg end @manager.jobs[@job_id]['job_argv'] = .clone .unshift("#{worker_action}") .unshift("#{worker_stage}") setup_command_line(*) end def send_msg(channel, = nil) publish channel, .present? && .is_a?(Hash) ? { job_id: @job_id }.merge() : { 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
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
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() if () check_gitflow save_tasks_to_be_executed() update_machine_state(['task']) # if message['action'] == 'invoke' log_to_file("worker #{@job_id} state is #{@machine.state}") task_approval() elsif () result = Celluloid::Actor[:terminal_server].show_confirmation(['question'], ['default']) publish_rake_event(.merge('action' => 'stdin', 'result' => result, 'client_action' => 'stdin')) else log_to_file(, @job_id) end end |
#message_is_about_a_task?(message) ⇒ Boolean
143 144 145 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 143 def () .present? && .is_a?(Hash) && ['action'].present? && ['job_id'].present? && ['task'].present? end |
#message_is_for_stdout?(message) ⇒ Boolean
139 140 141 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_worker.rb', line 139 def () .present? && .is_a?(Hash) && ['action'].present? && ['job_id'].present? && ['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 () log_to_file("worker #{@job_id} received: #{.inspect}") if @client.succesfull_subscription?() @successfull_subscription = true execute_after_succesfull_subscription else handle_subscription() 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() log_to_file("worler #{@job_id} current invocation chain : #{rake_tasks.inspect}") rake_tasks << ['task'] if rake_tasks.last != ['task'] invocation_chain << ['task'] if invocation_chain.last != ['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, = nil) publish channel, .present? && .is_a?(Hash) ? { job_id: @job_id }.merge() : { 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(*) @task_argv = [] .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}" = [] @env_options.each do |key, value| << "#{key}=#{value}" if value.present? end << '--trace' if app_debug_enabled? args.each do |arg| << arg end @manager.jobs[@job_id]['job_argv'] = .clone .unshift("#{worker_action}") .unshift("#{worker_stage}") setup_command_line(*) 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() job_conditions = @manager.job_to_condition[@job_id] if job_conditions.present? && app_configuration.task_confirmations.include?(['task']) && ['action'] == 'invoke' task_confirmation = job_conditions[['task']] task_confirmation[:status] = 'confirmed' task_confirmation[:condition].signal(['task']) else publish_rake_event(.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
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 |