Class: CapistranoMulticonfigParallel::RakeWorker
- Inherits:
-
Object
- Object
- CapistranoMulticonfigParallel::RakeWorker
- Includes:
- 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
Instance Attribute Summary collapse
-
#action ⇒ Object
Returns the value of attribute action.
-
#client ⇒ Object
Returns the value of attribute client.
-
#env ⇒ Object
Returns the value of attribute env.
-
#job_id ⇒ Object
Returns the value of attribute job_id.
-
#publisher_channel ⇒ Object
Returns the value of attribute publisher_channel.
-
#stdin_result ⇒ Object
Returns the value of attribute stdin_result.
-
#subscription_channel ⇒ Object
Returns the value of attribute subscription_channel.
-
#successfull_subscription ⇒ Object
Returns the value of attribute successfull_subscription.
-
#task ⇒ Object
Returns the value of attribute task.
-
#task_approved ⇒ Object
Returns the value of attribute task_approved.
Instance Method Summary collapse
- #custom_attributes ⇒ Object
- #debug_enabled? ⇒ Boolean
- #default_settings ⇒ Object
- #get_question_details(data) ⇒ Object
- #initialize_subscription ⇒ Object
- #msg_for_stdin?(message) ⇒ Boolean
- #on_close(code, reason) ⇒ Object
- #on_message(message) ⇒ Object
- #printing_question?(data) ⇒ Boolean
- #publish_new_work(env, new_options = {}) ⇒ Object
- #publish_subscription_successfull ⇒ Object
- #publish_to_worker(data) ⇒ Object
- #stdin_approval(message) ⇒ Object
- #task_approval(message) ⇒ Object
- #task_data ⇒ Object
- #task_name ⇒ Object
- #user_prompt_needed?(data) ⇒ Boolean
- #wait_execution(name = task_name, time = 0.1) ⇒ Object
- #wait_for(name, time) ⇒ Object
- #wait_for_stdin_input ⇒ Object
- #work(env, options = {}) ⇒ Object
Instance Attribute Details
#action ⇒ Object
Returns the value of attribute action.
7 8 9 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 7 def action @action end |
#client ⇒ Object
Returns the value of attribute client.
7 8 9 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 7 def client @client end |
#env ⇒ Object
Returns the value of attribute env.
7 8 9 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 7 def env @env end |
#job_id ⇒ Object
Returns the value of attribute job_id.
7 8 9 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 7 def job_id @job_id end |
#publisher_channel ⇒ Object
Returns the value of attribute publisher_channel.
7 8 9 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 7 def publisher_channel @publisher_channel end |
#stdin_result ⇒ Object
Returns the value of attribute stdin_result.
7 8 9 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 7 def stdin_result @stdin_result end |
#subscription_channel ⇒ Object
Returns the value of attribute subscription_channel.
7 8 9 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 7 def subscription_channel @subscription_channel end |
#successfull_subscription ⇒ Object
Returns the value of attribute successfull_subscription.
7 8 9 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 7 def successfull_subscription @successfull_subscription end |
#task ⇒ Object
Returns the value of attribute task.
7 8 9 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 7 def task @task end |
#task_approved ⇒ Object
Returns the value of attribute task_approved.
7 8 9 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 7 def task_approved @task_approved end |
Instance Method Details
#custom_attributes ⇒ Object
21 22 23 24 25 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 21 def custom_attributes @publisher_channel = "worker_#{@job_id}" @action = @options['actor_id'].include?('_count') ? 'count' : 'invoke' @task = @options['task'] end |
#debug_enabled? ⇒ Boolean
57 58 59 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 57 def debug_enabled? CapistranoMulticonfigParallel::CelluloidManager.debug_websocket? end |
#default_settings ⇒ Object
44 45 46 47 48 49 50 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 44 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
131 132 133 134 135 136 137 138 139 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 131 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
52 53 54 55 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 52 def initialize_subscription return if defined?(@client) && @client.present? @client = CelluloidPubsub::Client.connect(actor: Actor.current, enable_debug: debug_enabled?, :channel => @subscription_channel ) end |
#msg_for_stdin?(message) ⇒ Boolean
93 94 95 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 93 def msg_for_stdin?() ['action'] == 'stdin' end |
#on_close(code, reason) ⇒ Object
126 127 128 129 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 126 def on_close(code, reason) debug("websocket connection closed: #{code.inspect}, #{reason.inspect}") if debug_enabled? terminate end |
#on_message(message) ⇒ Object
78 79 80 81 82 83 84 85 86 87 88 89 90 91 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 78 def () debug("Rake worker #{@job_id} received after parse #{}") #if debug_enabled? if @client.succesfull_subscription?() debug("Rake worker #{@job_id} received parse #{}") if debug_enabled? @successfull_subscription = true publish_to_worker(task_data) elsif .present? && ['task'].present? task_approval() elsif .present? && ['action'].present? && ['action'] == 'stdin' stdin_approval() else warn "unknown action: #{.inspect}" if debug_enabled? end end |
#printing_question?(data) ⇒ Boolean
141 142 143 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 141 def printing_question?(data) get_question_details(data).present? end |
#publish_new_work(env, new_options = {}) ⇒ Object
27 28 29 30 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 27 def publish_new_work(env, = {}) work(env, @options.merge()) publish_to_worker(task_data) end |
#publish_subscription_successfull ⇒ Object
97 98 99 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 97 def publish_subscription_successfull end |
#publish_to_worker(data) ⇒ Object
74 75 76 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 74 def publish_to_worker(data) @client.publish(@publisher_channel, data) end |
#stdin_approval(message) ⇒ Object
110 111 112 113 114 115 116 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 110 def stdin_approval() if @job_id.to_i == ['job_id'].to_i && ['result'].present? && ['action'] == 'stdin' @stdin_result = ['result'] else warn "unknown invocation #{.inspect}" if debug_enabled? end end |
#task_approval(message) ⇒ Object
118 119 120 121 122 123 124 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 118 def task_approval() if @job_id.to_i == ['job_id'].to_i && ['task'] == task_name && ['approved'] == 'yes' @task_approved = true else warn "unknown invocation #{.inspect}" if debug_enabled? end end |
#task_data ⇒ Object
65 66 67 68 69 70 71 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 65 def task_data { action: @action, task: task_name, job_id: @job_id } end |
#task_name ⇒ Object
61 62 63 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 61 def task_name @task.name end |
#user_prompt_needed?(data) ⇒ Boolean
145 146 147 148 149 150 151 152 153 154 155 156 157 158 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 145 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
32 33 34 35 36 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 32 def wait_execution(name = task_name, time = 0.1) # info "Before waiting #{name}" Actor.current.wait_for(name, time) # info "After waiting #{name}" end |
#wait_for(name, time) ⇒ Object
38 39 40 41 42 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 38 def wait_for(name, time) # info "waiting for #{time} seconds on #{name}" sleep time # info "done waiting on #{name} " end |
#wait_for_stdin_input ⇒ Object
101 102 103 104 105 106 107 108 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 101 def wait_for_stdin_input until @stdin_result.present? wait_execution end output = @stdin_result.clone Actor.current.stdin_result = nil output end |
#work(env, options = {}) ⇒ Object
13 14 15 16 17 18 19 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 13 def work(env, = {}) @options = .stringify_keys @env = env default_settings custom_attributes initialize_subscription end |