Class: TestQueue::Runner
- Inherits:
-
Object
show all
- Defined in:
- lib/test_queue/runner.rb,
lib/test_queue/runner/rspec.rb,
lib/test_queue/runner/sample.rb,
lib/test_queue/runner/cucumber.rb,
lib/test_queue/runner/minitest.rb,
lib/test_queue/runner/minitest5.rb,
lib/test_queue/runner/puppet_lint.rb,
lib/test_queue/runner/minitest4.rb,
lib/test_queue/runner/testunit.rb
Defined Under Namespace
Classes: Cucumber, MiniTest, PuppetLint, RSpec, Sample, TestUnit
Instance Attribute Summary collapse
Instance Method Summary
collapse
Constructor Details
#initialize(queue, concurrency = nil, socket = nil, relay = nil) ⇒ Runner
Returns a new instance of Runner.
27
28
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
|
# File 'lib/test_queue/runner.rb', line 27
def initialize(queue, concurrency=nil, socket=nil, relay=nil)
raise ArgumentError, 'array required' unless Array === queue
if forced = ENV['TEST_QUEUE_FORCE']
forced = forced.split(/\s*,\s*/)
whitelist = Set.new(forced)
queue = queue.select{ |s| whitelist.include?(s.to_s) }
queue.sort_by!{ |s| forced.index(s.to_s) }
end
@procline = $0
@queue = queue
@suites = queue.inject(Hash.new){ |hash, suite| hash.update suite.to_s => suite }
@workers = {}
@completed = []
@concurrency =
concurrency ||
(ENV['TEST_QUEUE_WORKERS'] && ENV['TEST_QUEUE_WORKERS'].to_i) ||
if File.exists?('/proc/cpuinfo')
File.read('/proc/cpuinfo').split("\n").grep(/processor/).size
elsif RUBY_PLATFORM =~ /darwin/
`/usr/sbin/sysctl -n hw.activecpu`.to_i
else
2
end
@slave_connection_timeout =
(ENV['TEST_QUEUE_RELAY_TIMEOUT'] && ENV['TEST_QUEUE_RELAY_TIMEOUT'].to_i) ||
30
@run_token = ENV['TEST_QUEUE_RELAY_TOKEN'] || SecureRandom.hex(8)
@socket =
socket ||
ENV['TEST_QUEUE_SOCKET'] ||
"/tmp/test_queue_#{$$}_#{object_id}.sock"
@relay =
relay ||
ENV['TEST_QUEUE_RELAY']
@slave_message = ENV["TEST_QUEUE_SLAVE_MESSAGE"] if ENV.has_key?("TEST_QUEUE_SLAVE_MESSAGE")
if @relay == @socket
STDERR.puts "*** Detected TEST_QUEUE_RELAY == TEST_QUEUE_SOCKET. Disabling relay mode."
@relay = nil
elsif @relay
@queue = []
end
end
|
Instance Attribute Details
#concurrency ⇒ Object
Returns the value of attribute concurrency.
25
26
27
|
# File 'lib/test_queue/runner.rb', line 25
def concurrency
@concurrency
end
|
Instance Method Details
#after_fork(num) ⇒ Object
Prepare a worker for executing jobs after a fork.
265
266
|
# File 'lib/test_queue/runner.rb', line 265
def after_fork(num)
end
|
#after_fork_internal(num, iterator) ⇒ Object
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
|
# File 'lib/test_queue/runner.rb', line 237
def after_fork_internal(num, iterator)
srand
output = File.open("/tmp/test_queue_worker_#{$$}_output", 'w')
$stdout.reopen(output)
$stderr.reopen($stdout)
$stdout.sync = $stderr.sync = true
$0 = "test-queue worker [#{num}]"
puts
puts "==> Starting #$0 (#{Process.pid} on #{Socket.gethostname}) - iterating over #{iterator.sock}"
puts
after_fork(num)
end
|
#around_filter(suite) ⇒ Object
260
261
262
|
# File 'lib/test_queue/runner.rb', line 260
def around_filter(suite)
yield
end
|
#cleanup_worker ⇒ Object
281
282
|
# File 'lib/test_queue/runner.rb', line 281
def cleanup_worker
end
|
#connect_to_relay ⇒ Object
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
|
# File 'lib/test_queue/runner.rb', line 368
def connect_to_relay
sock = nil
start = Time.now
puts "Attempting to connect for #{@slave_connection_timeout}s..."
while sock.nil?
begin
sock = TCPSocket.new(*@relay.split(':'))
rescue Errno::ECONNREFUSED => e
raise e if Time.now - start > @slave_connection_timeout
puts "Master not yet available, sleeping..."
sleep 0.5
end
end
sock
end
|
#distribute_queue ⇒ Object
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
|
# File 'lib/test_queue/runner.rb', line 314
def distribute_queue
return if relay?
remote_workers = 0
until @queue.empty? && remote_workers == 0
if IO.select([@server], nil, nil, 0.1).nil?
reap_worker(false) if @workers.any? else
sock = @server.accept
cmd = sock.gets.strip
case cmd
when /^POP/
if obj = @queue.shift
data = Marshal.dump(obj.to_s)
sock.write(data)
end
when /^SLAVE (\d+) ([\w\.-]+) (\w+)(?: (.+))?/
num = $1.to_i
slave = $2
run_token = $3
slave_message = $4
if run_token == @run_token
sock.write("OK\n")
remote_workers += num
else
STDERR.puts "*** Worker from run #{run_token} connected to master for run #{@run_token}; ignoring."
sock.write("WRONG RUN\n")
end
message = "*** #{num} workers connected from #{slave} after #{Time.now-@start_time}s"
message << " " + slave_message if slave_message
STDERR.puts message
when /^WORKER (\d+)/
data = sock.read($1.to_i)
worker = Marshal.load(data)
worker_completed(worker)
remote_workers -= 1
end
sock.close
end
end
ensure
stop_master
until @workers.empty?
reap_worker
end
end
|
#execute ⇒ Object
89
90
91
92
93
94
95
96
97
98
|
# File 'lib/test_queue/runner.rb', line 89
def execute
$stdout.sync = $stderr.sync = true
@start_time = Time.now
@concurrency > 0 ?
execute_parallel :
execute_sequential
ensure
summarize_internal unless $!
end
|
#execute_parallel ⇒ Object
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
|
# File 'lib/test_queue/runner.rb', line 155
def execute_parallel
start_master
prepare(@concurrency)
@prepared_time = Time.now
start_relay if relay?
spawn_workers
distribute_queue
ensure
stop_master
@workers.each do |pid, worker|
Process.kill 'KILL', pid
end
until @workers.empty?
reap_worker
end
end
|
#execute_sequential ⇒ Object
151
152
153
|
# File 'lib/test_queue/runner.rb', line 151
def execute_sequential
exit! run_worker(@queue)
end
|
#prepare(concurrency) ⇒ Object
Run in the master before the fork. Used to create concurrency copies of any databases required by the test workers.
257
258
|
# File 'lib/test_queue/runner.rb', line 257
def prepare(concurrency)
end
|
#reap_worker(blocking = true) ⇒ Object
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
|
# File 'lib/test_queue/runner.rb', line 289
def reap_worker(blocking=true)
if pid = Process.waitpid(-1, blocking ? 0 : Process::WNOHANG) and worker = @workers.delete(pid)
worker.status = $?
worker.end_time = Time.now
if File.exists?(file = "/tmp/test_queue_worker_#{pid}_output")
worker.output = IO.binread(file)
FileUtils.rm(file)
end
if File.exists?(file = "/tmp/test_queue_worker_#{pid}_stats")
worker.stats = Marshal.load(IO.binread(file))
FileUtils.rm(file)
end
relay_to_master(worker) if relay?
worker_completed(worker)
end
end
|
#relay? ⇒ Boolean
364
365
366
|
# File 'lib/test_queue/runner.rb', line 364
def relay?
!!@relay
end
|
#relay_to_master(worker) ⇒ Object
384
385
386
387
388
389
390
391
392
393
|
# File 'lib/test_queue/runner.rb', line 384
def relay_to_master(worker)
worker.host = Socket.gethostname
data = Marshal.dump(worker)
sock = connect_to_relay
sock.puts("WORKER #{data.bytesize}")
sock.write(data)
ensure
sock.close if sock
end
|
#run_worker(iterator) ⇒ Object
Entry point for internal runner implementations. The iterator will yield jobs from the shared queue on the master.
Returns nothing. exits 0 on success. exits N on error, where N is the number of failures.
273
274
275
276
277
278
279
|
# File 'lib/test_queue/runner.rb', line 273
def run_worker(iterator)
iterator.each do |item|
puts " #{item.inspect}"
end
return 0 end
|
#spawn_workers ⇒ Object
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
|
# File 'lib/test_queue/runner.rb', line 219
def spawn_workers
@concurrency.times do |i|
num = i+1
pid = fork do
@server.close if @server
iterator = Iterator.new(relay?? @relay : @socket, @suites, method(:around_filter))
after_fork_internal(num, iterator)
ret = run_worker(iterator) || 0
cleanup_worker
Kernel.exit! ret
end
@workers[pid] = Worker.new(pid, num)
end
end
|
#start_master ⇒ Object
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
|
# File 'lib/test_queue/runner.rb', line 174
def start_master
if !relay?
if @socket =~ /^(?:(.+):)?(\d+)$/
address = $1 || '0.0.0.0'
port = $2.to_i
@socket = "#$1:#$2"
@server = TCPServer.new(address, port)
else
FileUtils.rm(@socket) if File.exists?(@socket)
@server = UNIXServer.new(@socket)
end
end
desc = "test-queue master (#{relay?? "relaying to #{@relay}" : @socket})"
puts "Starting #{desc}"
$0 = "#{desc} - #{@procline}"
end
|
#start_relay ⇒ Object
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
|
# File 'lib/test_queue/runner.rb', line 192
def start_relay
return unless relay?
sock = connect_to_relay
message = @slave_message ? " #{@slave_message}" : ""
message.gsub!(/(\r|\n)/, "") sock.puts("SLAVE #{@concurrency} #{Socket.gethostname} #{@run_token}#{message}")
response = sock.gets.strip
unless response == "OK"
STDERR.puts "*** Got non-OK response from master: #{response}"
sock.close
exit! 1
end
sock.close
rescue Errno::ECONNREFUSED
STDERR.puts "*** Unable to connect to relay #{@relay}. Aborting.."
exit! 1
end
|
#stats ⇒ Object
80
81
82
83
84
85
86
87
|
# File 'lib/test_queue/runner.rb', line 80
def stats
@stats ||=
if File.exists?(file = stats_file)
Marshal.load(IO.binread(file)) || {}
else
{}
end
end
|
#stats_file ⇒ Object
146
147
148
149
|
# File 'lib/test_queue/runner.rb', line 146
def stats_file
ENV['TEST_QUEUE_STATS'] ||
'.test_queue_stats'
end
|
#stop_master ⇒ Object
211
212
213
214
215
216
217
|
# File 'lib/test_queue/runner.rb', line 211
def stop_master
return if relay?
FileUtils.rm_f(@socket) if @socket && @server.is_a?(UNIXServer)
@server.close rescue nil if @server
@socket = @server = nil
end
|
#summarize ⇒ Object
143
144
|
# File 'lib/test_queue/runner.rb', line 143
def summarize
end
|
#summarize_internal ⇒ Object
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
|
# File 'lib/test_queue/runner.rb', line 100
def summarize_internal
puts
puts "==> Summary (#{@completed.size} workers in %.4fs)" % (Time.now-@start_time)
puts
@failures = ''
@completed.each do |worker|
summarize_worker(worker)
@failures << worker.failure_output if worker.failure_output
puts " [%2d] %60s %4d suites in %.4fs (pid %d exit %d%s)" % [
worker.num,
worker.summary,
worker.stats.size,
worker.end_time - worker.start_time,
worker.pid,
worker.status.exitstatus,
worker.host && " on #{worker.host.split('.').first}"
]
end
unless @failures.empty?
puts
puts "==> Failures"
puts
puts @failures
end
puts
if @stats
File.open(stats_file, 'wb') do |f|
f.write Marshal.dump(stats)
end
end
summarize
estatus = @completed.inject(0){ |s, worker| s + worker.status.exitstatus }
estatus = 255 if estatus > 255
exit!(estatus)
end
|
#summarize_worker(worker) ⇒ Object
284
285
286
287
|
# File 'lib/test_queue/runner.rb', line 284
def summarize_worker(worker)
worker.summary = ''
worker.failure_output = ''
end
|
#worker_completed(worker) ⇒ Object
309
310
311
312
|
# File 'lib/test_queue/runner.rb', line 309
def worker_completed(worker)
@completed << worker
puts worker.output if ENV['TEST_QUEUE_VERBOSE'] || worker.status.exitstatus != 0
end
|