Class: Gritz::Supervisor::Master

Inherits:
Object
  • Object
show all
Defined in:
/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/supervisor/master.rb

Overview

Forks workers, monitors their status pipes, and owns every child until reaped.

Constant Summary collapse

SIGNALS =
%w[TERM INT QUIT TTIN TTOU HUP CHLD USR1 USR2].freeze

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(config, logger: Logger.new($stdout), status_io: nil, owner_channel: nil, listener: nil) ⇒ Master

Returns a new instance of Master.



13
14
15
16
17
18
19
20
21
22
23
24
25
# File '/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/supervisor/master.rb', line 13

def initialize(config, logger: Logger.new($stdout), status_io: nil, owner_channel: nil, listener: nil)
  @config = config
  @logger = logger
  @workers = {}
  @desired = config.workers
  @owner_channel = owner_channel
  @listener = listener
  @reports = owner_channel || (StatusChannel.new(status_io) if status_io)
  @metrics = Metrics::Aggregator.new
  @forwarded = []
  @replacement_queue = []
  @exit_status = 0
end

Instance Attribute Details

#workers ⇒ Object (readonly)



11
12
13
# File '/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/supervisor/master.rb', line 11

def workers
  @workers
end

Instance Method Details

#ready? ⇒ Boolean

Returns:

  • (Boolean)


84
85
86
87
88
# File '/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/supervisor/master.rb', line 84

def ready?
  !@shutdown_at && @workers.values.count { |handle|
    handle.state == "ready" && !handle.term_at && handle.stats[:healthy] != false
  } >= @config.min_ready_workers
end

#run ⇒ Object



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
# File '/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/supervisor/master.rb', line 27

def run
  @config.validate_runtime!
  raise ConfigurationError, "Supervisor requires workers > 0" unless @desired.positive?

  require "gritz/#{@config.transport}"
  @listener ||= Listener.bind(@config.bind) if @config.listener_strategy == :inherited_fd
  @signals = SignalQueue.new(signals: SIGNALS)
  unless @owner_channel
    @admin = AdminServer.new(bind: @config.admin_bind, status: -> { status }, ready: -> { ready? },
                             metrics: -> { @metrics.render(workers: @workers.values) }, logger: @logger)
  end
  @guard = ForkGuard.activate(mode: @config.fork_mode == :clean ? @config.fork_guard : :off, logger: @logger)
  @config.preload! if @config.preload_app?
  Process.warmup if @config.preload_app? && Process.respond_to?(:warmup)
  if @config.transport == :native && @config.workers > 1 && !RUBY_PLATFORM.include?("linux")
    @logger.warn("Multiple native workers require Linux for SO_REUSEPORT load balancing; use workers 0 on macOS")
  end
  maintain_worker_count
  loop do
    @signals.drain.each { |signal| handle_signal(signal) }
    read_statuses
    reap_children
    check_timeouts
    sample_memory
    check_recycle
    advance_replacement
    maintain_worker_count unless @shutdown_at
    flush_forwarded
    if @reports && @forwarded.empty? && (!@next_report_at || now >= @next_report_at || @shutdown_at) &&
       @reports.write { @owner_channel ? status.merge(type: "status") : status }
      @next_report_at = now + @config.status_interval
    end
    @admin&.poll
    if @owner_channel&.closed?
      @exit_status = 1
      begin_shutdown(immediate: true)
    end
    if @shutdown_at && @workers.empty?
      break if !@owner_channel || @owner_channel.closed? || (@forwarded.empty? && @owner_channel.flush)
      break if now >= @shutdown_at + @config.drain_delay + @config.shutdown_timeout
    end

    worker_ios = @forwarded.empty? ? @workers.values.filter_map { |handle| handle.channel.io unless handle.channel.closed? } : []
    owner_ios = !@forwarded.empty? && @owner_channel && !@owner_channel.closed? ? [@owner_channel.io] : nil
    IO.select([@signals.io, *@admin&.ios.to_a, *worker_ios], owner_ios, nil, 0.05)
  end
  @exit_status
ensure
  cleanup
  ForkGuard.deactivate if @guard
end

#status ⇒ Object



79
80
81
82
# File '/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/supervisor/master.rb', line 79

def status
  { pid: Process.pid, state: @shutdown_at ? "draining" : "running", desired: @desired,
    phased_restart: !@replacement.nil? || !@replacement_queue.empty?, workers: @workers.values.map(&:to_h) }
end