Class: Gritz::Testing::Cluster

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

Overview

Starts a supervisor in a fresh interpreter, avoiding gRPC state inherited from a test process.

Constant Summary collapse

MAX_LOG_BYTES =
64 * 1024

Instance Attribute Summary collapse

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(config_path:, env: {}, command: nil) ⇒ Cluster

Returns a new instance of Cluster.



28
29
30
31
32
33
# File '/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/testing/cluster.rb', line 28

def initialize(config_path:, env: {}, command: nil)
  @config_path = File.expand_path(config_path)
  @env = env
  @command = command
  @status = { state: "starting", workers: [] }
end

Instance Attribute Details

#pid ⇒ Object (readonly)



15
16
17
# File '/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/testing/cluster.rb', line 15

def pid
  @pid
end

Class Method Details

.start ⇒ Object



17
18
19
20
21
22
23
24
25
26
# File '/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/testing/cluster.rb', line 17

def self.start(**)
  cluster = new(**).start
  return cluster unless block_given?

  begin
    yield cluster
  ensure
    cluster.stop
  end
end

Instance Method Details

#logs ⇒ Object



67
68
69
70
71
72
73
74
75
# File '/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/testing/cluster.rb', line 67

def logs
  return @logs || "" unless @log && !@log.closed?

  @log.flush
  File.open(@log.path) do |file|
    file.seek([file.size - MAX_LOG_BYTES, 0].max)
    file.read
  end
end

#master_pid ⇒ Object



65
# File '/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/testing/cluster.rb', line 65

def master_pid = status[:pid] || @pid

#signal(name, pid: @pid) ⇒ Object

Raises:

  • (ArgumentError)


77
78
79
80
81
82
# File '/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/testing/cluster.rb', line 77

def signal(name, pid: @pid)
  raise ArgumentError, "cluster is not started" unless pid

  Process.kill(name, pid)
  self
end

#start ⇒ Object



35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
# File '/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/testing/cluster.rb', line 35

def start
  raise ArgumentError, "cluster is already started" if @pid

  Supervisor::Launcher.enable_subreaper!
  @reader, writer = IO.pipe
  @channel = Supervisor::StatusChannel.new(@reader)
  @log = Tempfile.new(["gritz-cluster", ".log"])
  command = @command || [
    RbConfig.ruby, "-I", $LOAD_PATH.join(File::PATH_SEPARATOR), "-rgritz/core", "-e",
    "exit Gritz::CLI.new(status_io: IO.for_fd(3), launch: true).run(ARGV)", "--", "start", "-C", @config_path
  ]
  @pid = Process.spawn(@env, *command, 3 => writer, out: @log, err: @log, pgroup: true)
  self
rescue StandardError
  unless @pid
    @channel&.close
    @log&.close!
  end
  raise
ensure
  writer&.close
end

#status ⇒ Object



58
59
60
61
# File '/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/testing/cluster.rb', line 58

def status
  @channel&.read&.each { |row| @status = row }
  @status
end

#stop(timeout: 10) ⇒ Object



120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
# File '/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/testing/cluster.rb', line 120

def stop(timeout: 10)
  return self unless @pid

  signal("TERM") unless exited?
  wait(timeout:)
  self
rescue Timeout::Error
  signal("QUIT") if status[:owner_pid] && !exited?
  begin
    wait(timeout: 0.5)
  rescue Timeout::Error
    owned_groups.each { |pid| kill_group(pid) }
    wait(timeout: 5)
  end
  self
ensure
  cleanup_groups if @pid && @exit_status
  @channel&.close
  @logs = logs
  @log&.close!
end

#wait(timeout: 10) ⇒ Object

Raises:

  • (ArgumentError)


108
109
110
111
112
113
114
115
116
117
118
# File '/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/testing/cluster.rb', line 108

def wait(timeout: 10)
  raise ArgumentError, "cluster is not started" unless @pid

  deadline = monotonic + timeout
  until exited?
    raise Timeout::Error, "cluster did not exit\n#{logs}" if monotonic >= deadline

    sleep((deadline - monotonic).clamp(0, 0.01))
  end
  @exit_status
end

#wait_until(state: "ready", workers: nil, timeout: 10) ⇒ Object

A predicate can inspect snapshots without fixed sleeps or dependence on log wording.



85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
# File '/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/testing/cluster.rb', line 85

def wait_until(state: "ready", workers: nil, timeout: 10)
  deadline = monotonic + timeout
  loop do
    snapshot = status
    active_workers = snapshot.fetch(:workers, []).reject { |worker| worker[:retiring] }
    ready = if state == "ready"
              snapshot[:state] == "running" && active_workers.any? && active_workers.all? { |worker| worker[:state] == "ready" }
            else
              state.nil? || snapshot[:state] == state
            end
    count = workers.nil? || active_workers.size == workers
    return self if ready && count && (!block_given? || yield(snapshot))

    if exited?
      raise "cluster exited #{@exit_status.inspect} before reaching #{state.inspect}\n#{logs}"
    end
    raise Timeout::Error, "cluster did not reach #{state.inspect}: #{snapshot.inspect}\n#{logs}" if monotonic >= deadline

    pause = (deadline - monotonic).clamp(0, 0.05)
    @reader.closed? ? sleep(pause) : @reader.wait_readable(pause)
  end
end

#workers ⇒ Object



63
# File '/home/runner/work/gritz/gritz/sources/gritz-core/lib/gritz/testing/cluster.rb', line 63

def workers = status.fetch(:workers, [])