Class: Hive::Worker

Inherits:
Object
  • Object
show all
Defined in:
lib/hive/worker.rb,
lib/hive/worker/shell.rb

Overview

The generic worker class

Direct Known Subclasses

Shell

Defined Under Namespace

Classes: DeviceNotReady, InvalidJobReservationError, NoPortsAvailable, Shell

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(options) ⇒ Worker

The main worker process loop



29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
# File 'lib/hive/worker.rb', line 29

def initialize(options)
  @options = options
  @parent_pid = @options['parent_pid']
  @device_id = @options['id']
  @hive_id = @options['hive_id']
  @default_component ||= self.class.to_s
  @current_job_start_time = nil
  @hive_mind ||= mind_meld_klass.new(
    url: Chamber.env.network.hive_mind? ? Chamber.env.network.hive_mind : nil,
    pem: Chamber.env.network.cert ? Chamber.env.network.cert : nil,
    ca_file: Chamber.env.network.cafile ? Chamber.env.network.cafile : nil,
    verify_mode: Chamber.env.network.verify_mode ? Chamber.env.network.verify_mode : nil,
    device: hive_mind_device_identifiers
  )
  @device_identity = @options['device_identity'] || 'unknown-device'
  pid = Process.pid
  $PROGRAM_NAME = "#{@options['name_stub'] || 'WORKER'}.#{pid}"
  @log = Hive::Log.new
  @log.add_logger(
    "#{LOG_DIRECTORY}/#{pid}.#{@device_identity}.log",
    Hive.config.logging.worker_level || 'INFO'
  )
  @log.hive_mind = @hive_mind
  @log.default_progname = @default_component

  self.update_queues

  @port_allocator = (@options.has_key?('port_allocator') ? @options['port_allocator'] : Hive::PortAllocator.new(ports: []))
  
  platform = self.class.to_s.scan(/[^:][^:]*/)[2].downcase
  @diagnostic_runner = Hive::DiagnosticRunner.new(@options, Hive.config.diagnostics, platform, @hive_mind) if Hive.config.diagnostics? && Hive.config.diagnostics[platform]

  Hive::Messages.configure do |config|
    config.base_path = Hive.config.network.scheduler
    config.pem_file = Hive.config.network.cert
    config.ssl_verify_mode = OpenSSL::SSL::VERIFY_NONE
  end

  Signal.trap('TERM') do
    @log.info("Worker terminated")
    exit
  end

  @log.info('Starting worker')
  while keep_worker_running?
    begin
      @log.clear({component: @default_component, level: Hive.config.logging.hm_logs_to_delete})
      update_queues
      poll_queue if diagnostics
    rescue DeviceNotReady => e
      @log.warn("#{e.message}\n");
    rescue StandardError => e
      @log.warn("Worker loop aborted: #{e.message}\n  : #{e.backtrace.join("\n  : ")}")
    end
    sleep Hive.config.timings.worker_loop_interval
  end
  @log.info('Exiting worker')
end

Instance Attribute Details

#device_api ⇒ Object

Device API Object for device associated with this worker



26
27
28
# File 'lib/hive/worker.rb', line 26

def device_api
  @device_api
end

#queues ⇒ Object

Device API Object for device associated with this worker



26
27
28
# File 'lib/hive/worker.rb', line 26

def queues
  @queues
end

Instance Method Details

#after_error(job, file_system, script) ⇒ Object

Any tasks to do after a script has terminated with an error



456
457
# File 'lib/hive/worker.rb', line 456

def after_error(job, file_system, script)
end

#allocate_port ⇒ Object

Allocate a port



475
476
477
478
479
# File 'lib/hive/worker.rb', line 475

def allocate_port
  @log.warn("Using deprecated 'Hive::Worker.allocate_port' method")
  @log.warn("Use @port_allocator.allocate_port instead")
  @port_allocator.allocate_port
end

#autogenerated_queues ⇒ Object

List of autogenerated queues for the worker



276
277
278
# File 'lib/hive/worker.rb', line 276

def autogenerated_queues
  []
end

#checkout_code(repository, checkout_directory, branch) ⇒ Object

Get a checkout of the repository



402
403
404
# File 'lib/hive/worker.rb', line 402

def checkout_code(repository, checkout_directory, branch)
  CodeCache.repo(repository).checkout(:head, checkout_directory, branch) or raise "Unable to checkout repository #{repository}"
end

#cleanup ⇒ Object

Do whatever device cleanup is required



471
472
# File 'lib/hive/worker.rb', line 471

def cleanup
end

#detect_res_file(results_dir) ⇒ Object



376
377
378
# File 'lib/hive/worker.rb', line 376

def detect_res_file(results_dir)
  Dir.glob( "#{results_dir}/*.res" ).first
end

#device_status ⇒ Object

Current state of the device This method should be replaced in child classes, as appropriate



265
266
267
# File 'lib/hive/worker.rb', line 265

def device_status
  @device_status ||= 'happy'
end

#diagnostics ⇒ Object

Diagnostics function to be extended in child class, as required

Raises:



251
252
253
254
255
256
257
258
259
260
261
# File 'lib/hive/worker.rb', line 251

def diagnostics
  retn = true
  protect
  retn = @diagnostic_runner.run if !@diagnostic_runner.nil?
  unprotect
  @log.info('Diagnostics failed') if not retn
  status = device_status
  status = set_device_status('happy') if status == 'busy'
  raise DeviceNotReady.new("Current device status: '#{status}'") if status != 'happy'
  retn
end

#exceeded_time_limit? ⇒ Boolean

Returns:

  • (Boolean)


427
428
429
430
431
432
433
434
435
436
437
438
439
# File 'lib/hive/worker.rb', line 427

def exceeded_time_limit?
  if @job && !@job.nil?
    if max_time = @job.execution_variables.job_timeout rescue nil
      elapsed = (Time.now - @current_job_start_time).to_i
      @log.debug("Elapsed = #{elapsed} seconds, Max = #{max_time} minutes")
      if elapsed > max_time.to_i * 60          
        @log.warn("Job has exceeded max time of #{max_time} minutes")
        return true 
      end
    end
  end
  false
end

#execute_job ⇒ Object

Execute a job



131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
# File 'lib/hive/worker.rb', line 131

def execute_job
  # Ensure that a killed worker cleans up correctly
  Signal.trap('TERM') do |s|
    Signal.trap('TERM') {} # Prevent retry signals
    @log.info "Caught TERM signal"
    @log.info "Terminating script, if running"
    @script.terminate if @script
    @log.info "Post-execution cleanup"
    signal_safe_post_script(@job, @file_system, @script)

    # Upload results
    @file_system.finalise_results_directory
    upload_files(@job, @file_system.results_path, @file_system.logs_path)
    set_job_state_to :completed
    @job.error('Worker killed')
    @log.info "Worker terminated"
    exit
  end

  @log.info('Job starting')
  @job.prepare(@hive_mind.id)
  exception = nil
  begin
    @log.info "Setting job paths"
    @file_system = Hive::FileSystem.new(@job.job_id, Hive.config.logging.home, @log)
    set_job_state_to :preparing

    if ! @job.repository.to_s.empty?
      @log.info "Checking out the repository"
      @log.debug "  #{@job.repository}"
      @log.debug "  #{@file_system.testbed_path}"
      checkout_code(@job.repository, @file_system.testbed_path, @job.branch)
    end

    @log.info "Initialising execution script"
    @script = Hive::ExecutionScript.new(
      file_system: @file_system,
      log: @log,
      keep_running: ->() { self.keep_script_running? }
    )
    @script.append_bash_cmd "mkdir -p #{@file_system.testbed_path}/#{@job.execution_directory}"
    @script.append_bash_cmd "cd #{@file_system.testbed_path}/#{@job.execution_directory}"

    @log.info "Setting the execution variables in the environment"
    @script.set_env 'HIVE_RESULTS', @file_system.results_path
    @script.set_env 'HIVE_SCRIPT_ERRORS', @file_system.script_errors_file
    @job.execution_variables.to_h.each_pair do |var, val|
      @script.set_env "HIVE_#{var.to_s}".upcase, val if ! val.kind_of?(Array)
    end
    if @job.execution_variables.retry_urns && !@job.execution_variables.retry_urns.empty?
      @script.set_env "RETRY_URNS", @job.execution_variables.retry_urns
    end
    if @job.execution_variables.tests && @job.execution_variables.tests != [""]
      @script.set_env "TEST_NAMES", @job.execution_variables.tests
    end
    

    @log.info "Appending test script to execution script"
    @script.append_bash_cmd @job.command

    set_job_state_to :running

    @log.info "Pre-execution setup"
    pre_script(@job, @file_system, @script)

    @job.start
    @log.info "Running execution script"
    exit_value = @script.run
    @job.end(exit_value)
  rescue => e
    exception = e
  end

  begin
    @log.info "Post-execution cleanup"
    set_job_state_to :uploading
    post_script(@job, @file_system, @script)

    # Upload results
    @file_system.finalise_results_directory
    upload_results(@job, "#{@file_system.testbed_path}/#{@job.execution_directory}", @file_system.results_path)
  rescue => e
    @log.error( "Post execution failed: " + e.message)
    @log.error("  : #{e.backtrace.join("\n  : ")}")
  end

  if exception or File.size(@file_system.script_errors_file) > 0
    set_job_state_to :completed
    begin
      after_error(@job, @file_system, @script)
      upload_files(@job, @file_system.results_path, @file_system.logs_path)
    rescue => e
      @log.error("Exception while uploading files: #{e.backtrace.join("\n  : ")}")
    end
    if exception
      @job.error( exception.message )
      raise exception
    else
      @job.error( 'Errors raised by execution script' )
      raise 'See errors file for errors reported in test.'
    end
  else
    @job.complete
    begin
      upload_files(@job, @file_system.results_path, @file_system.logs_path)
    rescue => e
      @log.error("Exception while uploading files: #{e.backtrace.join("\n  : ")}")
    end
  end

  Signal.trap('TERM') do
    @log.info("Worker terminated")
    exit
  end

  set_job_state_to :completed
  exit_value == 0
end

#hive_mind_device_identifiers ⇒ Object

Parameters for uniquely identifying the device



503
504
505
# File 'lib/hive/worker.rb', line 503

def hive_mind_device_identifiers
  { id: @device_id }
end

#job_message_klass ⇒ Object

Get the correct job class This should usually be replaced in the child class



116
117
118
119
# File 'lib/hive/worker.rb', line 116

def job_message_klass
  @log.info 'Generic job class'
  Hive::Messages::Job
end

#keep_script_running? ⇒ Boolean

Keep the execution script running

Returns:

  • (Boolean)


418
419
420
421
422
423
424
425
# File 'lib/hive/worker.rb', line 418

def keep_script_running?
  @log.debug("Keep Running check ")
  if exceeded_time_limit? or parent_process_dead? or File.size(@file_system.script_errors_file) > 0
    return false
  else
    return true
  end
end

#keep_worker_running? ⇒ Boolean

Keep the worker process running

Returns:

  • (Boolean)


407
408
409
410
411
412
413
414
415
# File 'lib/hive/worker.rb', line 407

def keep_worker_running?
  @log.debug("Keep Worker Running check ")
  if parent_process_dead?
    @log.info("Think parent process is dead")
    false
  else
    true
  end
end

#lion_config(checkout) ⇒ Object



397
398
399
# File 'lib/hive/worker.rb', line 397

def lion_config(checkout)
  Dir.glob( "#{checkout}/.lion.yml" ).first
end

#mind_meld_klass ⇒ Object



121
122
123
# File 'lib/hive/worker.rb', line 121

def mind_meld_klass
  MindMeld::Device
end

#parent_process_dead? ⇒ Boolean

Returns:

  • (Boolean)


441
442
443
444
445
446
447
448
449
# File 'lib/hive/worker.rb', line 441

def parent_process_dead?
  begin
    Process.getpgid(@parent_pid)
    false
  rescue
    @log.warn("Parent process appears to have terminated")
    true
  end
end

#poll_queue ⇒ Object

Check the queues for work



89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
# File 'lib/hive/worker.rb', line 89

def poll_queue
  @job = reserve_job
  if @job.nil?
    @log.info('No job found')
  else
    @log.info('Job starting')
    begin
      @current_job_start_time = Time.now
      execute_job
    rescue => e
      @log.info("Error running test: #{e.message}\n : #{e.backtrace.join("\n :")}")
    end
    cleanup
  end
end

#post_script(job, file_system, script) ⇒ Object

Any device specific steps immediately after the execution script



460
461
462
# File 'lib/hive/worker.rb', line 460

def post_script(job, file_system, script)
  signal_safe_post_script(job, file_system, script)
end

#pre_script(job, file_system, script) ⇒ Object

Any setup required before the execution script



452
453
# File 'lib/hive/worker.rb', line 452

def pre_script(job, file_system, script)
end

#process_xunit_results(results_dir) ⇒ Object



380
381
382
383
384
385
386
387
388
389
390
391
# File 'lib/hive/worker.rb', line 380

def process_xunit_results(results_dir) 
  if !Dir.glob("#{results_dir}/*.xml").empty?
    xunit_output = Res.parse_results(parser: :junit,:file =>  Dir.glob( "#{results_dir}/*.xml" ).first)
    res_output = File.open(xunit_output.io, "rb")
    contents = res_output.read
    res_output.close
    res = File.open("#{results_dir}/xunit.res", "w+")
    res.puts contents   
    res.close   
    res 
  end
end

#release_all_ports ⇒ Object

Release all ports



489
490
491
492
493
# File 'lib/hive/worker.rb', line 489

def release_all_ports
  @log.warn("Using deprecated 'Hive::Worker.release_all_ports' method")
  @log.warn("Use @port_allocator.release_all_ports instead")
  @port_allocator.release_all_ports
end

#release_port(p) ⇒ Object

Release a port



482
483
484
485
486
# File 'lib/hive/worker.rb', line 482

def release_port(p)
  @log.warn("Using deprecated 'Hive::Worker.release_port' method")
  @log.warn("Use @port_allocator.release_port instead")
  @port_allocator.release_port(p)
end

#reservation_details ⇒ Object



125
126
127
128
# File 'lib/hive/worker.rb', line 125

def reservation_details
  @log.debug "Reservations details: hive_id=#{@hive_id}, worker_pid=#{Process.pid}, device_id=#{@hive_mind.id}"
  { hive_id: @hive_id, worker_pid: Process.pid, device_id: @hive_mind.id }
end

#reserve_job ⇒ Object

Try to find and reserve a job



106
107
108
109
110
111
112
# File 'lib/hive/worker.rb', line 106

def reserve_job
  @log.info "Trying to reserve job for queues: #{@queues.join(', ')}"
  job = job_message_klass.reserve(@queues, reservation_details)
  @log.debug "Job: #{job.inspect}"
  raise InvalidJobReservationError.new("Invalid Job Reserved") if ! (job.nil? || job.valid?)
  job
end

#set_device_status(status) ⇒ Object

Set the status of a device This method should be replaced in child classes, as appropriate



271
272
273
# File 'lib/hive/worker.rb', line 271

def set_device_status(status)
  @device_status = status
end

#set_job_state_to(state) ⇒ Object

Set job info file



496
497
498
499
500
# File 'lib/hive/worker.rb', line 496

def set_job_state_to state
  File.open("#{@file_system.home_path}/job_info", 'w') do |f|
    f.puts "#{Process.pid} #{state}"
  end
end

#signal_safe_post_script(job, file_system, script) ⇒ Object

Any device specific steps immediately after the execution script that can be safely run in the a Signal.trap This should be called by post_script



467
468
# File 'lib/hive/worker.rb', line 467

def signal_safe_post_script(job, file_system, script)
end

#testmine_config(checkout) ⇒ Object



393
394
395
# File 'lib/hive/worker.rb', line 393

def testmine_config(checkout)
  Dir.glob( "#{checkout}/.testmi{n,t}e.yml" ).first
end

#update_queue_log ⇒ Object



289
290
291
# File 'lib/hive/worker.rb', line 289

def update_queue_log
  File.open("#{LOG_DIRECTORY}/#{Process.pid}.queues.yml",'w') { |f| f.write @queues.to_yaml}
end

#update_queues ⇒ Object



280
281
282
283
284
285
286
287
# File 'lib/hive/worker.rb', line 280

def update_queues
  # Get Queues from Hive Mind
  @log.debug("Getting queues from Hive Mind")
  @queues = (autogenerated_queues + @hive_mind.hive_queues(true)).uniq
  @log.debug("hive queues: #{@hive_mind.hive_queues}")
  @log.debug("Full list of queues: #{@queues}")
  update_queue_log
end

#upload_files(job, *paths) ⇒ Object

Upload any files from the test



294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
# File 'lib/hive/worker.rb', line 294

def upload_files(job, *paths)
  @log.info("Uploading assets")
  paths.each do |path|
    @log.info("Uploading files from #{path}")
    Dir.foreach(path) do |item|
      @log.info("File: #{item}")
      next if item == '.' or item == '..'
      begin
        artifact = job.report_artifact("#{path}/#{item}")
        @log.info("Artifact uploaded: #{artifact.attributes.to_s}")
      rescue => e
        @log.error("Error uploading artifact #{item}: #{e.message}")
        @log.error("  : #{e.backtrace.join("\n  : ")}")
      end
    end
  end
end

#upload_results(job, checkout, results_dir) ⇒ Object

Update results



313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
# File 'lib/hive/worker.rb', line 313

def upload_results(job, checkout, results_dir)

  res_file = detect_res_file(results_dir) || process_xunit_results(results_dir)
  
  if res_file
    @log.info("Res file found")
  
    begin
      Res.submit_results(
        reporter: :hive,
        ir: res_file,
        job_id: job.job_id
      )
    rescue => e
      @log.warn("Res Hive upload failed #{e.message}")
    end
  
    begin
      if conf_file = testmine_config(checkout)
        Res.submit_results(
          reporter: :testmine,
          ir: res_file,
          config_file: conf_file,
          hive_job_id: job.job_id,
          version: job.execution_variables.version,
          target: job.execution_variables.queue_name,
          cert: Chamber.env.network.cert,
          cacert: Chamber.env.network.cafile,
          ssl_verify_mode: Chamber.env.network.verify_mode
        )
      end
    rescue => e
      @log.warn("Res Testmine upload failed #{e.message}")
    end

    begin
      if conf_file = lion_config(checkout)
        Res.submit_results(
            reporter: :lion,
            ir: res_file,
            config_file: conf_file,
            hive_job_id: job.job_id,
            version: job.execution_variables.version,
            target: job.execution_variables.queue_name,
            cert: Chamber.env.network.cert,
            cacert: Chamber.env.network.cafile,
            ssl_verify_mode: Chamber.env.network.verify_mode
        )
      end
    rescue => e
      @log.warn("Res Lion upload failed #{e.message}")

      end




    # TODO Add in Testrail upload
  
  end
  
end