Class: Foruiman::Engine

Inherits:
Object
  • Object
show all
Defined in:
lib/foruiman/engine.rb

Overview

Derived from Foreman's registration, process lookup, pipes and self-pipe signal handling. All process state and output now belong to the caller's event loop.

Defined Under Namespace

Classes: Event, State

Constant Summary collapse

FORWARDED_SIGNALS =
%i[USR1 USR2].freeze
HANDLED_SIGNALS =
%i[INT TERM HUP USR1 USR2].freeze
TERM_TIMEOUT =
5.0
READ_CHUNK =
4096
READ_BUDGET =
64 * 1024
INPUT_BUDGET =
64 * 1024

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(procfile: nil, root: Dir.pwd, env: ENV.to_h, input: $stdin, port: 5000, log_lines: 10_000, term_timeout: TERM_TIMEOUT, exit_on: :all) ⇒ Engine

Returns a new instance of Engine.



24
25
26
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
# File 'lib/foruiman/engine.rb', line 24

def initialize(procfile: nil, root: Dir.pwd, env: ENV.to_h, input: $stdin, port: 5000, log_lines: 10_000,
               term_timeout: TERM_TIMEOUT, exit_on: :all)
  raise Foruiman::Error, "port must be an integer in 1..65535" unless port.is_a?(Integer) && (1..65_535).cover?(port)
  raise Foruiman::Error, "log-lines must be a positive integer" unless log_lines.is_a?(Integer) && log_lines.positive?
  unless term_timeout.is_a?(Numeric) && term_timeout.real? && term_timeout.finite? && term_timeout >= 0
    raise Foruiman::Error, "timeout must be a finite nonnegative number"
  end
  raise Foruiman::Error, "exit-on must be all, any, or failure" unless %i[all any failure].include?(exit_on)

  @root = File.expand_path(root)
  raise Foruiman::Error, "working directory does not exist: #{@root}" unless File.directory?(@root)

  @env = env.dup.freeze
  @input = input
  @managed_input = false
  @base_port = port
  @log_lines = log_lines
  @term_timeout = term_timeout
  @exit_on = exit_on
  @processes = []
  @names = {}
  @running = {}
  @readers = {}
  @inputs = {}
  @listeners = []
  @shutdown = false
  @explicit_shutdown = false
  @failed = false
  @closed = false
  @started = false
  @pending_signals = []
  @self_reader, @self_writer = create_pipe
  load_procfile(procfile) if procfile
rescue StandardError
  @self_reader&.close
  @self_writer&.close
  raise
end

Instance Attribute Details

#env ⇒ Object (readonly)

Returns the value of attribute env.



22
23
24
# File 'lib/foruiman/engine.rb', line 22

def env
  @env
end

#exit_on ⇒ Object (readonly)

Returns the value of attribute exit_on.



22
23
24
# File 'lib/foruiman/engine.rb', line 22

def exit_on
  @exit_on
end

#logs ⇒ Object (readonly)

Returns the value of attribute logs.



22
23
24
# File 'lib/foruiman/engine.rb', line 22

def logs
  @logs
end

#processes ⇒ Object (readonly)

Returns the value of attribute processes.



22
23
24
# File 'lib/foruiman/engine.rb', line 22

def processes
  @processes
end

#procfile_path ⇒ Object (readonly)

Returns the value of attribute procfile_path.



22
23
24
# File 'lib/foruiman/engine.rb', line 22

def procfile_path
  @procfile_path
end

#root ⇒ Object (readonly)

Returns the value of attribute root.



22
23
24
# File 'lib/foruiman/engine.rb', line 22

def root
  @root
end

#term_timeout ⇒ Object (readonly)

Returns the value of attribute term_timeout.



22
23
24
# File 'lib/foruiman/engine.rb', line 22

def term_timeout
  @term_timeout
end

Instance Method Details

#close ⇒ Object



262
263
264
265
266
267
268
269
270
271
272
273
274
# File 'lib/foruiman/engine.rb', line 262

def close
  return if @closed

  # Observers can fail (e.g. a closed stdout). Cleanup must still own the loop.
  @listeners.clear
  shutdown(explicit: false)
  tick(timeout: 0.01) until !@started || finished?
ensure
  restore_signal_handlers
  @self_reader.close unless @self_reader.closed?
  @self_writer.close unless @self_writer.closed?
  @closed = true
end

#exit_code ⇒ Object



223
224
225
226
227
# File 'lib/foruiman/engine.rb', line 223

def exit_code
  return 0 if @explicit_shutdown

  @policy_exit_code || (@failed ? 1 : 0)
end

#finished? ⇒ Boolean

Returns:

  • (Boolean)


219
220
221
# File 'lib/foruiman/engine.rb', line 219

def finished?
  @started && processes.all? { |entry| !entry.pgid && !entry.restart_pending } && @readers.empty?
end

#interrupt_process(name) ⇒ Object



188
189
190
191
192
193
194
# File 'lib/foruiman/engine.rb', line 188

def interrupt_process(name)
  entry = state(name)
  return unless entry.pgid

  signal_group(entry, :INT)
  self
end

#load_procfile(filename) ⇒ Object



77
78
79
80
81
82
83
# File 'lib/foruiman/engine.rb', line 77

def load_procfile(filename)
  parsed = Foruiman::Procfile.new(filename)
  entries = parsed.entries.to_a
  entries.each { |name, command| register(name, command) }
  @procfile_path = File.expand_path(filename).freeze
  self
end

#manage_input! ⇒ Object

Give every child a dedicated pseudo-terminal for stdin. The TUI remains the sole reader of the real terminal and explicitly forwards input to one child.

Raises:



122
123
124
125
126
127
128
# File 'lib/foruiman/engine.rb', line 122

def manage_input!
  raise Foruiman::Error, "cannot change input mode after startup" if @started

  require "pty"
  @managed_input = true
  self
end

#on_event(&listener) ⇒ Object



115
116
117
118
# File 'lib/foruiman/engine.rb', line 115

def on_event(&listener)
  @listeners << listener
  self
end

#process(name) ⇒ Object



98
99
100
# File 'lib/foruiman/engine.rb', line 98

def process(name)
  @names[name]&.process
end

#process_names ⇒ Object



85
86
87
# File 'lib/foruiman/engine.rb', line 85

def process_names
  @names.keys
end

#register(name, command) ⇒ Object

Raises:



63
64
65
66
67
68
69
70
71
72
73
74
75
# File 'lib/foruiman/engine.rb', line 63

def register(name, command)
  raise Foruiman::Error, "cannot register after startup" if @started
  raise Foruiman::Error, "duplicate process: #{name}" if @names.key?(name)

  Foruiman::Procfile.new[name] = command
  process = Foruiman::Process.new(command, cwd: root, env: env)
  state = State.new(name: name.freeze, process: process, port: @base_port + (processes.size * 100),
                    status: :pending, generation: 0, restart_pending: false, reaped: true, group_gone: true,
                    input_buffer: +"".b)
  @names[name] = state
  processes << state
  state
end

#resize_inputs(rows, columns) ⇒ Object



196
197
198
199
200
201
202
# File 'lib/foruiman/engine.rb', line 196

def resize_inputs(rows, columns)
  processes.each do |entry|
    entry.input&.winsize = [rows, columns]
  rescue IOError, SystemCallError
    nil
  end
end

#restart(name) ⇒ Object



148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
# File 'lib/foruiman/engine.rb', line 148

def restart(name)
  return if @shutdown || @closed

  return start(name) unless @started

  entry = state(name)
  return if entry.restart_pending

  validate_ports!([entry])
  entry.restart_pending = true
  entry.status = :restarting
  lifecycle(entry, :restarting, "restarting")
  if entry.pgid
    terminate(entry)
  else
    spawn_process(entry)
  end
end

#run(keep_open: false) ⇒ Object



229
230
231
232
233
234
235
236
237
238
239
240
# File 'lib/foruiman/engine.rb', line 229

def run(keep_open: false)
  register_signal_handlers
  start_all
  loop do
    tick(timeout: keep_open ? 1.0 / 30 : 0.05)
    yield self if block_given?
    break if finished? && (!keep_open || shutting_down?)
  end
  exit_code
ensure
  close
end

#select(name) ⇒ Object

Raises:



106
107
108
109
110
111
112
113
# File 'lib/foruiman/engine.rb', line 106

def select(name)
  raise Foruiman::Error, "cannot select after startup" if @started
  raise Foruiman::Error, "unknown process: #{name}" unless @names.key?(name)

  @names.select! { |key, _entry| key == name }
  processes.select! { |entry| entry.name == name }
  self
end

#shutdown(explicit: true) ⇒ Object



204
205
206
207
208
209
210
211
212
213
# File 'lib/foruiman/engine.rb', line 204

def shutdown(explicit: true)
  @explicit_shutdown ||= explicit
  return if @shutdown

  @shutdown = true
  processes.each do |entry|
    entry.restart_pending = false
    stop(entry.name)
  end
end

#shutting_down? ⇒ Boolean

Returns:

  • (Boolean)


215
216
217
# File 'lib/foruiman/engine.rb', line 215

def shutting_down?
  @shutdown
end

#start(name = nil) ⇒ Object Also known as: start_all

Raises:



130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
# File 'lib/foruiman/engine.rb', line 130

def start(name = nil)
  raise Foruiman::Error, "supervisor is closed" if @closed
  return if @shutdown

  targets = name ? [state(name)] : processes
  validate_ports!(targets)
  unless @started
    @logs = Foruiman::LogStore.new(process_names, capacity: @log_lines) do |record|
      emit(:output, state(record.name), record: record)
    end
    @started = true
  end
  targets.each { |entry| spawn_process(entry) unless entry.pgid }
  self
end

#state(name) ⇒ Object



102
103
104
# File 'lib/foruiman/engine.rb', line 102

def state(name)
  @names.fetch(name)
end

#stop(name) ⇒ Object



167
168
169
170
171
172
173
174
175
# File 'lib/foruiman/engine.rb', line 167

def stop(name)
  entry = state(name)
  entry.restart_pending = false
  return unless entry.pgid

  entry.status = :stopping
  lifecycle(entry, :stopping, "stopping")
  terminate(entry)
end

#tick(timeout: 0.03) ⇒ Object

Embedders may drive this method directly; start/restart/stop are called on that same thread. A bounded round-robin read prevents output starving input.



244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
# File 'lib/foruiman/engine.rb', line 244

def tick(timeout: 0.03)
  return if @closed

  handle_signals
  reap_children
  advance_groups
  writable = processes.filter_map { |entry| entry.input unless entry.input_buffer.empty? }
  ready, writable = IO.select([@self_reader, *@readers.keys], writable, nil, timeout) || [[], []]
  drain_signal_pipe if ready.delete(@self_reader)
  handle_signals
  writable.each { |input| flush_input(@inputs.fetch(input)) if @inputs.key?(input) }
  read_output(ready)
  reap_children
  advance_groups
rescue Errno::EINTR
  # The next tick handles the deferred signal.
end

#validate_ports!(entries = processes) ⇒ Object



89
90
91
92
93
94
95
96
# File 'lib/foruiman/engine.rb', line 89

def validate_ports!(entries = processes)
  entries.each do |entry|
    next if entry.port <= 65_535

    raise Foruiman::Error, "allocated port #{entry.port} exceeds 65535 for #{entry.name}"
  end
  self
end

#write_input(name, bytes) ⇒ Object



177
178
179
180
181
182
183
184
185
186
# File 'lib/foruiman/engine.rb', line 177

def write_input(name, bytes)
  entry = state(name)
  return false unless entry.status == :running && entry.input && !entry.input.closed?
  return false if entry.input_buffer.bytesize + bytes.bytesize > INPUT_BUDGET

  entry.input_buffer << bytes.b
  flush_input(entry)
rescue IOError, SystemCallError
  false
end