Class: CapistranoMulticonfigParallel::RakeWorker
- Inherits:
-
Object
- Object
- CapistranoMulticonfigParallel::RakeWorker
show all
- Includes:
- ApplicationHelper, Celluloid, Celluloid::Logger
- Defined in:
- lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb
Overview
class that handles the rake task and waits for approval from the celluloid worker
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
Instance Attribute Details
#action ⇒ Object
Returns the value of attribute action.
10
11
12
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 10
def action
@action
end
|
#client ⇒ Object
Returns the value of attribute client.
10
11
12
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 10
def client
@client
end
|
#env ⇒ Object
Returns the value of attribute env.
10
11
12
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 10
def env
@env
end
|
#job_id ⇒ Object
Returns the value of attribute job_id.
10
11
12
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 10
def job_id
@job_id
end
|
#publisher_channel ⇒ Object
Returns the value of attribute publisher_channel.
10
11
12
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 10
def publisher_channel
@publisher_channel
end
|
#stdin_result ⇒ Object
Returns the value of attribute stdin_result.
10
11
12
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 10
def stdin_result
@stdin_result
end
|
#subscription_channel ⇒ Object
Returns the value of attribute subscription_channel.
10
11
12
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 10
def subscription_channel
@subscription_channel
end
|
#successfull_subscription ⇒ Object
Returns the value of attribute successfull_subscription.
10
11
12
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 10
def successfull_subscription
@successfull_subscription
end
|
#task ⇒ Object
Returns the value of attribute task.
10
11
12
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 10
def task
@task
end
|
#task_approved ⇒ Object
Returns the value of attribute task_approved.
10
11
12
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 10
def task_approved
@task_approved
end
|
Instance Method Details
#custom_attributes ⇒ Object
22
23
24
25
26
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 22
def custom_attributes
@publisher_channel = "worker_#{@job_id}"
@action = @options['actor_id'].include?('_count') ? 'count' : 'invoke'
@task = @options['task']
end
|
#default_settings ⇒ Object
45
46
47
48
49
50
51
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 45
def default_settings
@stdin_result = nil
@job_id = @options['job_id']
@subscription_channel = @options['actor_id']
@task_approved = false
@successfull_subscription = false
end
|
#get_question_details(data) ⇒ Object
136
137
138
139
140
141
142
143
144
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 136
def get_question_details(data)
question = ''
default = nil
if data =~ /(.*)\?*\s*\:*\s*(\([^)]*\))*/m
question = Regexp.last_match(1)
default = Regexp.last_match(2)
end
question.present? ? [question, default] : nil
end
|
#initialize_subscription ⇒ Object
53
54
55
56
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 53
def initialize_subscription
return if defined?(@client) && @client.present?
@client = CelluloidPubsub::Client.connect(actor: Actor.current, enable_debug: debug_websocket?, channel: @subscription_channel)
end
|
#log_debug(action, message) ⇒ Object
87
88
89
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 87
def log_debug(action, message)
log_to_file("Rake worker #{@job_id} received after #{action}: #{message}")
end
|
#msg_for_stdin?(message) ⇒ Boolean
91
92
93
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 91
def msg_for_stdin?(message)
message['action'] == 'stdin'
end
|
#msg_for_task?(message) ⇒ Boolean
95
96
97
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 95
def msg_for_task?(message)
message['task'].present?
end
|
#on_close(code, reason) ⇒ Object
131
132
133
134
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 131
def on_close(code, reason)
log_to_file("websocket connection closed: #{code.inspect}, #{reason.inspect}")
terminate
end
|
#on_message(message) ⇒ Object
74
75
76
77
78
79
80
81
82
83
84
85
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 74
def on_message(message)
return unless message.present?
log_debug('on_message', message)
if @client.succesfull_subscription?(message)
publish_subscription_successfull(message)
elsif msg_for_task?(message) || msg_for_stdin?(message)
task_approval(message)
stdin_approval(message)
else
show_warning "unknown action: #{message.inspect}"
end
end
|
#printing_question?(data) ⇒ Boolean
146
147
148
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 146
def printing_question?(data)
get_question_details(data).present?
end
|
#publish_new_work(env, new_options = {}) ⇒ Object
28
29
30
31
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 28
def publish_new_work(env, new_options = {})
work(env, @options.merge(new_options))
publish_to_worker(task_data)
end
|
#publish_subscription_successfull(message) ⇒ Object
99
100
101
102
103
104
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 99
def publish_subscription_successfull(message)
return unless @client.succesfull_subscription?(message)
log_debug('publish_subscription_successfull', message)
@successfull_subscription = true
publish_to_worker(task_data)
end
|
#publish_to_worker(data) ⇒ Object
70
71
72
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 70
def publish_to_worker(data)
@client.publish(@publisher_channel, data)
end
|
#stdin_approval(message) ⇒ Object
113
114
115
116
117
118
119
120
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 113
def stdin_approval(message)
return unless msg_for_stdin?(message)
if @job_id.to_i == message['job_id'].to_i && message['result'].present?
@stdin_result = message['result']
else
show_warning "unknown invocation #{message.inspect}"
end
end
|
#task_approval(message) ⇒ Object
122
123
124
125
126
127
128
129
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 122
def task_approval(message)
return unless msg_for_task?(message)
if @job_id.to_i == message['job_id'].to_i && message['task'] == task_name && message['approved'] == 'yes'
@task_approved = true
else
show_warning "unknown invocation #{message.inspect}"
end
end
|
#task_data ⇒ Object
62
63
64
65
66
67
68
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 62
def task_data
{
action: @action,
task: task_name,
job_id: @job_id
}
end
|
#task_name ⇒ Object
58
59
60
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 58
def task_name
@task.name
end
|
#user_prompt_needed?(data) ⇒ Boolean
150
151
152
153
154
155
156
157
158
159
160
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 150
def user_prompt_needed?(data)
return if !printing_question?(data) || @action != 'invoke'
details = get_question_details(data)
default = details.second.present? ? details.second : nil
publish_to_worker(action: 'stdout',
question: details.first,
default: default.delete('()'),
job_id: @job_id)
wait_for_stdin_input
end
|
#wait_execution(name = task_name, time = 0.1) ⇒ Object
33
34
35
36
37
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 33
def wait_execution(name = task_name, time = 0.1)
Actor.current.wait_for(name, time)
end
|
#wait_for(_name, time) ⇒ Object
39
40
41
42
43
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 39
def wait_for(_name, time)
sleep time
end
|
106
107
108
109
110
111
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 106
def wait_for_stdin_input
wait_execution until @stdin_result.present?
output = @stdin_result.clone
Actor.current.stdin_result = nil
output
end
|
#work(env, options = {}) ⇒ Object
14
15
16
17
18
19
20
|
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 14
def work(env, options = {})
@options = options.stringify_keys
@env = env
default_settings
custom_attributes
initialize_subscription
end
|