Class: CapistranoMulticonfigParallel::CelluloidManager
- Inherits:
-
Object
- Object
- CapistranoMulticonfigParallel::CelluloidManager
- 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
-
#job_to_condition ⇒ Object
Returns the value of attribute job_to_condition.
-
#job_to_worker ⇒ Object
Returns the value of attribute job_to_worker.
-
#jobs ⇒ Object
Returns the value of attribute jobs.
-
#mutex ⇒ Object
Returns the value of attribute mutex.
-
#registration_complete ⇒ Object
Returns the value of attribute registration_complete.
-
#worker_supervisor ⇒ Object
readonly
Returns the value of attribute worker_supervisor.
-
#worker_to_job ⇒ Object
Returns the value of attribute worker_to_job.
-
#workers ⇒ Object
readonly
Returns the value of attribute workers.
-
#workers_terminated ⇒ Object
Returns the value of attribute workers_terminated.
Instance Method Summary collapse
- #all_workers_finished? ⇒ Boolean
- #apply_confirmation_for_job(job) ⇒ Object
- #apply_confirmations? ⇒ Boolean
- #can_tag_staging? ⇒ Boolean
- #confirm_task_approval(result, task, processed_job = nil) ⇒ Object
-
#delegate(job) ⇒ Object
call to send an actor a job.
- #dispatch_new_job(job, options = {}) ⇒ Object
-
#get_job_status(job) ⇒ Object
lookup status of job by asking actor running it.
- #get_worker_for_job(job) ⇒ Object
-
#initialize(job_manager) ⇒ CelluloidManager
constructor
A new instance of CelluloidManager.
- #job_crashed?(job) ⇒ Boolean
- #job_failed?(job) ⇒ Boolean
- #mark_completed_remaining_tasks(job) ⇒ Object
- #print_confirm_task_approvall(result, task, job) ⇒ Object
- #process_jobs ⇒ Object
-
#register_worker_for_job(job, worker) ⇒ Object
call back from actor once it has received it's job actor should do this asap.
- #setup_worker_conditions(job) ⇒ Object
- #syncronized_confirmation? ⇒ Boolean
- #wait_condition_for_task(job_id, task) ⇒ Object
- #wait_task_confirmations ⇒ Object
- #wait_task_confirmations_worker(job) ⇒ Object
- #worker_died(worker, reason) ⇒ 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, 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
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
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
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, = {}) env_opts = @job_manager.(job.app, job.stage) job. = .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
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
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 |
#print_confirm_task_approvall(result, task, job) ⇒ Object
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) = "Do you want to continue the deployment and execute #{task.upcase}" += " for JOB #{job.id}" if job.present? += '?' apps_symlink_confirmation = Celluloid::Actor[:terminal_server].show_confirmation(, '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
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 |