Class: Foruiman::Engine
- Inherits:
-
Object
- Object
- Foruiman::Engine
- 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
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
-
#env ⇒ Object
readonly
Returns the value of attribute env.
-
#exit_on ⇒ Object
readonly
Returns the value of attribute exit_on.
-
#logs ⇒ Object
readonly
Returns the value of attribute logs.
-
#processes ⇒ Object
readonly
Returns the value of attribute processes.
-
#procfile_path ⇒ Object
readonly
Returns the value of attribute procfile_path.
-
#root ⇒ Object
readonly
Returns the value of attribute root.
-
#term_timeout ⇒ Object
readonly
Returns the value of attribute term_timeout.
Instance Method Summary collapse
- #close ⇒ Object
- #exit_code ⇒ Object
- #finished? ⇒ Boolean
-
#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
constructor
A new instance of Engine.
- #interrupt_process(name) ⇒ Object
- #load_procfile(filename) ⇒ Object
-
#manage_input! ⇒ Object
Give every child a dedicated pseudo-terminal for stdin.
- #on_event(&listener) ⇒ Object
- #process(name) ⇒ Object
- #process_names ⇒ Object
- #register(name, command) ⇒ Object
- #resize_inputs(rows, columns) ⇒ Object
- #restart(name) ⇒ Object
- #run(keep_open: false) ⇒ Object
- #select(name) ⇒ Object
- #shutdown(explicit: true) ⇒ Object
- #shutting_down? ⇒ Boolean
- #start(name = nil) ⇒ Object (also: #start_all)
- #state(name) ⇒ Object
- #stop(name) ⇒ Object
-
#tick(timeout: 0.03) ⇒ Object
Embedders may drive this method directly; start/restart/stop are called on that same thread.
- #validate_ports!(entries = processes) ⇒ Object
- #write_input(name, bytes) ⇒ Object
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.(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
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.(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.
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
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
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
215 216 217 |
# File 'lib/foruiman/engine.rb', line 215 def shutting_down? @shutdown end |
#start(name = nil) ⇒ Object Also known as: start_all
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 |