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

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

Returns:

  • (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

Returns:

  • (Boolean)


93
94
95
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 93

def msg_for_stdin?(message)
  message['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 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
     publish_to_worker(task_data)
  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

Returns:

  • (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, new_options = {})
  work(env, @options.merge(new_options))
   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(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



118
119
120
121
122
123
124
# File 'lib/capistrano_multiconfig_parallel/celluloid/rake_worker.rb', line 118

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

Returns:

  • (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 = {})
  @options = options.stringify_keys
  @env = env
  default_settings
  custom_attributes
  initialize_subscription
end