Class: CapistranoMulticonfigParallel::RakeWorker

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

Methods included from ApplicationHelper

parse_task_string, strip_characters_from_string

Methods included from StagesHelper

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

Methods included from InternalHelper

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

Returns:

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

Returns:

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

Returns:

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

Returns:

  • (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)
  #    info "Before waiting #{name}"
  Actor.current.wait_for(name, time)
  #  info "After waiting #{name}"
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)
  # info "waiting for #{time} seconds on #{name}"
  sleep time
  # info "done waiting on #{name} "
end

#wait_for_stdin_input ⇒ Object



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