Class: CapistranoMulticonfigParallel::RakeWorker

Inherits:
Object
  • Object
show all
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

Instance Method Summary collapse

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(options={})
 @initialize_options = options.stringify_keys
 @questions_prompted ||=[]
 @stdin_result = nil
 @job_id = @initialize_options['job_id']
 @subscription_channel = @initialize_options['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 nitialize_options
  @nitialize_options
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?(message)
  message['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 on_message(message)
  debug("Rake worker #{@job_id} received after parse #{message}") #if debug_enabled?
  if @client.succesfull_subscription?(message)
   debug("Rake worker #{@job_id} received  parse #{message}") if debug_enabled?
   @successfull_subscription = true
  elsif message.present? && message['task'].present? 
    task_approval(message)
  elsif message.present? && message['action'].present? && message['action'] == 'stdin'
    stdin_approval(message)
  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(message)
  if @job_id.to_i == message['job_id'].to_i && message['result'].present? && message['action'] == 'stdin'
    @stdin_result = message['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(message)
  if @job_id.to_i == message['job_id'].to_i && message['task'] == task_name && message['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, options = {})
  @options = options.stringify_keys
  @env = env
  @action = @subscription_channel.include?('_count') ? 'count' : 'invoke'
  @task = @options['task']
  wait_execution until @successfull_subscription == true && defined?(@client)
  publish_to_worker(task_data)
end