Class: CapistranoMulticonfigParallel::CelluloidManager
- Inherits:
-
Object
- Object
- CapistranoMulticonfigParallel::CelluloidManager
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
parse_task_string, strip_characters_from_string
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
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
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
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 18
def initialize(job_manager)
@worker_supervisor = Celluloid::SupervisionGroup.run!
@job_manager = job_manager
@registration_complete = false
@mutex = Mutex.new
@workers = @worker_supervisor.pool(CapistranoMulticonfigParallel::CelluloidWorker, as: :workers, size: 10)
Actor.current.link @workers
@conditions = []
@jobs = {}
@job_to_worker = {}
@worker_to_job = {}
@job_to_condition = {}
@worker_supervisor.supervise_as(:terminal_server, CapistranoMulticonfigParallel::TerminalTable, Actor.current, @job_manager)
@worker_supervisor.supervise_as(:web_server, CapistranoMulticonfigParallel::WebServer, websocket_config)
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
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
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
79
80
81
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 79
def all_workers_finished?
@job_to_worker.all? { |job_id, _worker| @jobs[job_id]['worker_action'] == 'finished' }
end
|
#apply_confirmation_for_worker(worker) ⇒ Object
108
109
110
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 108
def apply_confirmation_for_worker(worker)
worker.alive? && app_configuration.apply_stage_confirmation.include?(worker.env_name) && apply_confirmations?
end
|
#apply_confirmations? ⇒ Boolean
99
100
101
102
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 99
def apply_confirmations?
confirmations = app_configuration.task_confirmations
confirmations.is_a?(Array) && confirmations.present?
end
|
#can_tag_staging? ⇒ Boolean
199
200
201
202
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 199
def can_tag_staging?
@job_manager.can_tag_staging? &&
@jobs.find { |_job_id, job| job['env'] == 'production' }.blank?
end
|
#confirm_task_approval(result, task, worker = nil) ⇒ Object
172
173
174
175
176
177
178
179
180
181
182
183
184
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 172
def confirm_task_approval(result, task, worker = nil)
return unless result.present?
result = print_confirm_task_approvall(result, task, worker = nil)
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
47
48
49
50
51
52
53
54
55
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 47
def delegate(job)
job = job.stringify_keys
job['id'] = generate_job_id(job) unless job_failed?(job)
@jobs[job['id']] = job
job['env_options'][CapistranoMulticonfigParallel::ENV_KEY_JOB_ID] = job['id']
@workers.async.work(job, Actor.current)
end
|
#dispatch_new_job(job) ⇒ Object
204
205
206
207
208
209
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 204
def dispatch_new_job(job)
original_env = job['env_options']
env_opts = @job_manager.get_app_additional_env_options(job['app_name'], job['env'])
job['env_options'] = original_env.merge(env_opts)
async.delegate(job)
end
|
#filtered_env_keys ⇒ Object
211
212
213
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 211
def filtered_env_keys
%w(STAGES ACTION)
end
|
#generate_job_id(job) ⇒ Object
40
41
42
43
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 40
def generate_job_id(job)
@jobs[job['id']] = job
job['id']
end
|
#get_job_status(job) ⇒ Object
lookup status of job by asking actor running it
237
238
239
240
241
242
243
244
245
246
247
248
249
250
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 237
def get_job_status(job)
status = nil
if job.present?
if job.is_a?(Hash)
job = job.stringify_keys
actor = @registered_jobs[job['id']]
status = actor.status
else
actor = @registered_jobs[job.to_i]
status = actor.status
end
end
status
end
|
#get_worker_for_job(job) ⇒ Object
186
187
188
189
190
191
192
193
194
195
196
197
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 186
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.to_i]
end
else
return nil
end
end
|
#job_failed?(job) ⇒ Boolean
252
253
254
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 252
def job_failed?(job)
job['worker_action'].present? && job['worker_action'] == 'worker_died'
end
|
#mark_completed_remaining_tasks(worker) ⇒ Object
121
122
123
124
125
126
127
128
129
130
131
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 121
def mark_completed_remaining_tasks(worker)
return unless apply_confirmation_for_worker(worker)
app_configuration.task_confirmations.each_with_index do |task, _index|
fake_result = proc { |sum| sum }
task_confirmation = @job_to_condition[worker.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, worker = nil) ⇒ Object
160
161
162
163
164
165
166
167
168
169
170
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 160
def print_confirm_task_approvall(result, task, worker = nil)
return if result.is_a?(Proc)
message = "Do you want to continue the deployment and execute #{task.upcase}"
message += " for JOB #{worker.job_id}" if worker.present?
message += '?'
apps_symlink_confirmation = Celluloid::Actor[:terminal_server].show_confirmation(message, 'Y/N')
until apps_symlink_confirmation.present?
sleep(0.1)
end
apps_symlink_confirmation
end
|
#process_job(job) ⇒ Object
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 215
def process_job(job)
if job['processed']
@jobs[job['job_id']]
else
env_options = {}
job['env_options'].each do |key, value|
env_options[key] = value if value.present? && !filtered_env_keys.include?(key)
end
{
'job_id' => job['id'],
'app_name' => job['app'],
'env_name' => job['env'],
'action_name' => job['action'],
'env_options' => env_options,
'task_arguments' => job['task_arguments'],
'job_argv' => job.fetch('job_argv', []),
'processed' => true
}
end
end
|
#process_jobs ⇒ Object
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 83
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)
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
59
60
61
62
63
64
65
66
67
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 59
def register_worker_for_job(job, worker)
job = job.stringify_keys
if job['id'].blank?
log_to_file("job id not found. delegating again the job #{job.inspect}")
delegate(job)
else
start_worker(job, worker)
end
end
|
#setup_worker_conditions(worker) ⇒ Object
112
113
114
115
116
117
118
119
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 112
def setup_worker_conditions(worker)
return unless apply_confirmation_for_worker(worker)
hash_conditions = {}
app_configuration.task_confirmations.each do |task|
hash_conditions[task] = { condition: Celluloid::Condition.new, status: 'unconfirmed' }
end
@job_to_condition[worker.job_id] = hash_conditions
end
|
#start_worker(job, worker) ⇒ Object
69
70
71
72
73
74
75
76
77
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 69
def start_worker(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
|
#syncronized_confirmation? ⇒ Boolean
104
105
106
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 104
def syncronized_confirmation?
!@job_manager.can_tag_staging?
end
|
#wait_condition_for_task(job_id, task) ⇒ Object
141
142
143
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 141
def wait_condition_for_task(job_id, task)
@job_to_condition[job_id][task][:condition].wait
end
|
#wait_task_confirmations ⇒ Object
145
146
147
148
149
150
151
152
153
154
155
156
157
158
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 145
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(worker) ⇒ Object
133
134
135
136
137
138
139
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 133
def wait_task_confirmations_worker(worker)
return if !apply_confirmation_for_worker(worker) || !syncronized_confirmation?
app_configuration.task_confirmations.each_with_index do |task, _index|
result = wait_condition_for_task(worker.job_id, task)
confirm_task_approval(result, task, worker) if result.present?
end
end
|
#worker_died(worker, reason) ⇒ Object
256
257
258
259
260
261
262
263
264
265
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/celluloid_manager.rb', line 256
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_failed?(job)
return unless job['action_name'] == 'deploy'
log_to_file "restarting #{job} on new worker"
job = job.merge(:action => 'deploy:rollback', 'worker_action' => 'worker_died')
delegate(job)
end
|