Class: Gritz::Testing::Cluster
- Inherits:
-
Object
- Object
- Gritz::Testing::Cluster
- 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
- #pid ⇒ Object readonly
Class Method Summary collapse
Instance Method Summary collapse
-
#initialize(config_path:, env: {}, command: nil) ⇒ Cluster
constructor
A new instance of Cluster.
- #logs ⇒ Object
- #master_pid ⇒ Object
- #signal(name, pid: @pid) ⇒ Object
- #start ⇒ Object
- #status ⇒ Object
- #stop(timeout: 10) ⇒ Object
- #wait(timeout: 10) ⇒ Object
-
#wait_until(state: "ready", workers: nil, timeout: 10) ⇒ Object
A predicate can inspect snapshots without fixed sleeps or dependence on log wording.
- #workers ⇒ Object
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.(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
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
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, []) |