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.
-
#nitialize_options ⇒ Object
Returns the value of attribute nitialize_options.
-
#publisher_channel ⇒ Object
Returns the value of attribute publisher_channel.
-
#questions_prompted ⇒ Object
Returns the value of attribute questions_prompted.
-
#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
- #debug_enabled? ⇒ Boolean
- #get_question_details(data) ⇒ Object
- #has_asked_question?(question) ⇒ Boolean
-
#initialize(options = {}) ⇒ RakeWorker
constructor
A new instance of RakeWorker.
- #initialize_subscription ⇒ Object
- #msg_for_stdin?(message) ⇒ Boolean
- #on_close(code, reason) ⇒ Object
- #on_message(message) ⇒ Object
- #printing_question?(data) ⇒ Boolean
- #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
Constructor Details
#initialize(options = {}) ⇒ RakeWorker
14 15 16 17 18 19 20 21 22 23 24 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 14 def initialize(={}) = .stringify_keys @questions_prompted ||=[] @stdin_result = nil @job_id = ['job_id'] @subscription_channel = ['actor_id'] @publisher_channel = "worker_#{@job_id}" @task_approved = false @successfull_subscription = false initialize_subscription end |
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 |
#nitialize_options ⇒ Object
Returns the value of attribute nitialize_options.
7 8 9 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 7 def 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 |
#questions_prompted ⇒ Object
Returns the value of attribute questions_prompted.
7 8 9 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 7 def questions_prompted @questions_prompted 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
#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 |
#get_question_details(data) ⇒ Object
130 131 132 133 134 135 136 137 138 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 130 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 |
#has_asked_question?(question) ⇒ Boolean
144 145 146 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 144 def has_asked_question?(question) @questions_prompted.include?(question) end |
#initialize_subscription ⇒ Object
51 52 53 54 55 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 51 def initialize_subscription @client = CelluloidPubsub::Client.connect(actor: Actor.current, enable_debug: debug_enabled?) do |ws| ws.subscribe(@subscription_channel) end if !defined?(@client) || @client.nil? end |
#msg_for_stdin?(message) ⇒ Boolean
92 93 94 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 92 def msg_for_stdin?() ['action'] == 'stdin' end |
#on_close(code, reason) ⇒ Object
125 126 127 128 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 125 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 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 78 def () debug("Rake worker #{@job_id} received after parse #{message}") #if debug_enabled? if @client.succesfull_subscription?() debug("Rake worker #{@job_id} received parse #{message}") if debug_enabled? @successfull_subscription = true elsif .present? && ['task'].present? task_approval() elsif .present? && ['action'].present? && ['action'] == 'stdin' stdin_approval() else warn "unknown action: #{message.inspect}" if debug_enabled? end end |
#printing_question?(data) ⇒ Boolean
140 141 142 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 140 def printing_question?(data) get_question_details(data).present? end |
#publish_subscription_successfull ⇒ Object
96 97 98 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 96 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
109 110 111 112 113 114 115 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 109 def stdin_approval() if @job_id.to_i == ['job_id'].to_i && ['result'].present? && ['action'] == 'stdin' @stdin_result = ['result'] else warn "unknown invocation #{message.inspect}" if debug_enabled? end end |
#task_approval(message) ⇒ Object
117 118 119 120 121 122 123 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 117 def task_approval() if @job_id.to_i == ['job_id'].to_i && ['task'] == task_name && ['approved'] == 'yes' @task_approved = true else warn "unknown invocation #{message.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
148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 148 def user_prompt_needed?(data) return if !printing_question?(data) || has_asked_question?(data) 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 }) @questions_prompted << data wait_for_stdin_input end |
#wait_execution(name = task_name, time = 0.1) ⇒ Object
38 39 40 41 42 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 38 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
44 45 46 47 48 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 44 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
100 101 102 103 104 105 106 107 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 100 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
27 28 29 30 31 32 33 34 |
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 27 def work(env, = {}) = .stringify_keys @env = env @action = @subscription_channel.include?('_count') ? 'count' : 'invoke' @task = ['task'] wait_execution until @successfull_subscription == true && defined?(@client) publish_to_worker(task_data) end |