Class: MutationTester::ForkRunner

Inherits:
Object
  • Object
show all
Defined in:
lib/mutation_tester/fork_runner.rb

Defined Under Namespace

Classes: InMemoryOutcome

Constant Summary collapse

WORKER_PATH =
File.expand_path('fork_runner/worker.rb', __dir__)
IN_MEMORY_POOL_KEY =
:in_memory
BOOT_TIMEOUT =
60
CLONE_TIMEOUT =
10
RESPONSE_GRACE =
15
SHUTDOWN_GRACE =
2
POLL_INTERVAL =
0.05
HANDSHAKE_POLL =
0.005

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(use_bundle_exec:, framework: :rspec) ⇒ ForkRunner

Returns a new instance of ForkRunner.



140
141
142
143
144
145
146
147
148
149
150
# File 'lib/mutation_tester/fork_runner.rb', line 140

def initialize(use_bundle_exec:, framework: :rspec)
  argv = ['ruby', WORKER_PATH, framework.to_s]
  argv = ['bundle', 'exec', *argv] if use_bundle_exec

  job_reader, job_writer = IO.pipe
  event_reader, event_writer = IO.pipe
  pid = Process.spawn(*argv, pgroup: true, in: job_reader, out: event_writer, err: File::NULL)
  job_reader.close
  event_writer.close
  attach_endpoints(pid, job_writer, event_reader, BOOT_TIMEOUT)
end

Class Method Details

.acquire(use_bundle_exec:, framework: :rspec) ⇒ Object



24
25
26
27
28
29
30
31
# File 'lib/mutation_tester/fork_runner.rb', line 24

def acquire(use_bundle_exec:, framework: :rspec)
  return nil unless available?

  key = [Process.pid, use_bundle_exec, framework]
  return registry[key] if registry.key?(key)

  registry[key] = checkout_pooled(pool_key(use_bundle_exec, framework)) || boot(use_bundle_exec, framework)
end

.attached(pid, job_writer, event_reader) ⇒ Object



95
96
97
98
99
# File 'lib/mutation_tester/fork_runner.rb', line 95

def attached(pid, job_writer, event_reader)
  runner = allocate
  runner.send(:attach_endpoints, pid, job_writer, event_reader, CLONE_TIMEOUT, tolerate_eof: true)
  runner
end

.available?Boolean

Returns:

  • (Boolean)


20
21
22
# File 'lib/mutation_tester/fork_runner.rb', line 20

def available?
  Process.respond_to?(:fork)
end

.checkout_in_memoryObject



52
53
54
55
56
57
58
# File 'lib/mutation_tester/fork_runner.rb', line 52

def checkout_in_memory
  number = parallel_worker_number
  return nil unless number

  entry = pool[IN_MEMORY_POOL_KEY]
  entry && entry[:runners][number]
end

.discard(runner) ⇒ Object



80
81
82
83
84
85
# File 'lib/mutation_tester/fork_runner.rb', line 80

def discard(runner)
  registry.delete_if { |_, value| value.equal?(runner) }
  pool.each_value do |entry|
    entry[:runners].map! { |value| value.equal?(runner) ? nil : value }
  end
end

.in_memory_pool_prepared?Boolean

Returns:

  • (Boolean)


48
49
50
# File 'lib/mutation_tester/fork_runner.rb', line 48

def in_memory_pool_prepared?
  pool.key?(IN_MEMORY_POOL_KEY)
end

.poolObject



91
92
93
# File 'lib/mutation_tester/fork_runner.rb', line 91

def pool
  @pool ||= {}
end

.prepare_in_memory_pool(count, primary) ⇒ Object



42
43
44
45
46
# File 'lib/mutation_tester/fork_runner.rb', line 42

def prepare_in_memory_pool(count, primary)
  return [] unless available? && primary&.ready?

  refill_pool(IN_MEMORY_POOL_KEY, count, primary)
end

.prepare_pool(count, use_bundle_exec:, framework: :rspec, env_for: nil) ⇒ Object



33
34
35
36
37
38
39
40
# File 'lib/mutation_tester/fork_runner.rb', line 33

def prepare_pool(count, use_bundle_exec:, framework: :rspec, env_for: nil)
  return unless available?

  primary = acquire(use_bundle_exec: use_bundle_exec, framework: framework)
  return unless primary

  refill_pool(pool_key(use_bundle_exec, framework), count, primary, env_for: env_for)
end

.registryObject



87
88
89
# File 'lib/mutation_tester/fork_runner.rb', line 87

def registry
  @registry ||= {}
end

.shutdown_allObject



68
69
70
71
72
73
74
75
76
77
78
# File 'lib/mutation_tester/fork_runner.rb', line 68

def shutdown_all
  owned = pool.each_value.select { |entry| entry[:owner] == Process.pid }
  owned.flat_map { |entry| entry[:runners] }.compact
       .map { |runner| Thread.new { runner.shutdown } }
       .each(&:join)
  pool.delete_if { |_, entry| entry[:owner] == Process.pid }
  registry.each do |(pid, _), runner|
    runner&.shutdown if pid == Process.pid
  end
  registry.delete_if { |(pid, _), _| pid == Process.pid }
end

.shutdown_in_memory_poolObject



60
61
62
63
64
65
66
# File 'lib/mutation_tester/fork_runner.rb', line 60

def shutdown_in_memory_pool
  entry = pool[IN_MEMORY_POOL_KEY]
  return unless entry && entry[:owner] == Process.pid

  pool.delete(IN_MEMORY_POOL_KEY)
  entry[:runners].compact.map { |runner| Thread.new { runner.shutdown } }.each(&:join)
end

Instance Method Details

#execute(spec_file, timeout: nil, chdir: nil, capture: false, args: [], stop_on_first_failure: false) ⇒ Object



156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
# File 'lib/mutation_tester/fork_runner.rb', line 156

def execute(spec_file, timeout: nil, chdir: nil, capture: false, args: [], stop_on_first_failure: false)
  log = capture ? Tempfile.new(['mutation_tester_fork', '.log']) : nil
  job = {
    spec: spec_file,
    timeout: timeout,
    chdir: chdir,
    log: log&.path,
    args: args,
    stop_on_first_failure: stop_on_first_failure
  }
  @job_writer.puts(JSON.generate(job))
  status = await_result(timeout)['status']
  result = TestCommand::Result.new(status == 'pass', status == 'timeout')
  result.output = File.read(log.path) if log
  result
rescue Errno::EPIPE
  fail_worker
ensure
  log&.close
  log&.unlink
end

#execute_in_memory(source:, path:, timeout: nil, chdir: nil) ⇒ Object



190
191
192
193
194
195
196
197
# File 'lib/mutation_tester/fork_runner.rb', line 190

def execute_in_memory(source:, path:, timeout: nil, chdir: nil)
  job = { in_memory: { source: source, path: path }, timeout: timeout, chdir: chdir }
  @job_writer.puts(JSON.generate(job))
  event = await_result(timeout)
  InMemoryOutcome.new(event['status'], event['message'])
rescue Errno::EPIPE
  fail_worker
end

#fork_clone(env: nil) ⇒ Object



199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
# File 'lib/mutation_tester/fork_runner.rb', line 199

def fork_clone(env: nil)
  return nil unless ready?
  return nil unless File.respond_to?(:mkfifo)

  dir = Dir.mktmpdir('mutation_tester_clone')
  job_path = File.join(dir, 'job')
  events_path = File.join(dir, 'events')
  File.mkfifo(job_path)
  File.mkfifo(events_path)

  clone_request = { 'job' => job_path, 'events' => events_path }
  clone_request['env'] = env if env
  @job_writer.puts(JSON.generate('clone' => clone_request))
  event = read_event(monotonic_time + CLONE_TIMEOUT)
  return nil unless event.is_a?(Hash) && event['event'] == 'cloned'

  adopt_clone(event['pid'], job_path, events_path)
rescue SystemCallError, IOError
  nil
ensure
  FileUtils.remove_entry(dir) if dir
end

#preload(spec_file, chdir: nil, stop_on_first_failure: false) ⇒ Object



178
179
180
181
182
183
184
185
186
187
188
# File 'lib/mutation_tester/fork_runner.rb', line 178

def preload(spec_file, chdir: nil, stop_on_first_failure: false)
  request = { spec: spec_file, chdir: chdir, stop_on_first_failure: stop_on_first_failure }
  @job_writer.puts(JSON.generate(preload: request))
  event = read_event(monotonic_time + BOOT_TIMEOUT)
  return [true, nil] if event.is_a?(Hash) && event['event'] == 'preloaded' && event['status'] == 'ok'

  message = event.is_a?(Hash) && event['message'] ? event['message'] : 'the worker did not confirm the spec preload'
  [false, message]
rescue Errno::EPIPE, IOError
  [false, 'the fork runner worker terminated unexpectedly']
end

#ready?Boolean

Returns:

  • (Boolean)


152
153
154
# File 'lib/mutation_tester/fork_runner.rb', line 152

def ready?
  @ready
end

#shutdownObject



222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
# File 'lib/mutation_tester/fork_runner.rb', line 222

def shutdown
  @shutdown_mutex.synchronize do
    pid = @pid
    next unless pid

    @pid = nil
    kill_group(@current_child) if @current_child
    @current_child = nil
    begin
      Process.kill('TERM', pid)
    rescue Errno::ESRCH
    end
    kill_group(pid) unless reaped_within(pid, SHUTDOWN_GRACE)
    close_pipes
    @ready = false
  end
end