Class: Gritz::Transport::Async

Inherits:
Object
  • Object
show all
Defined in:
/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async.rb,
/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async/body.rb,
/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async/bridge.rb,
/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async/call.rb,
/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async/input.rb,
/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async/server.rb

Overview

Runs RPC controllers in fibers without the grpc C extension.

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(config:, dispatcher:, logger:) ⇒ Async

Returns a new instance of Async.



16
17
18
19
20
21
22
23
# File '/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async.rb', line 16

def initialize(config:, dispatcher:, logger:)
  @config = config
  @dispatcher = dispatcher
  @logger = logger
  @lock = Mutex.new
  @inflight = {}.compare_by_identity
  @requests_total = @rejected_total = 0
end

Class Method Details

.capabilities ⇒ Object



14
# File '/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async.rb', line 14

def self.capabilities = Set[:unary, :client_streaming, :server_streaming, :bidi, :reuseport, :inherited_fd, :fiber].freeze

Instance Method Details

#bind(listener_spec = @config.bind) ⇒ Integer

The caller retains ownership of an inherited socket; only its duplicate is closed.

Returns:

  • (Integer) —

    the selected TCP port

Raises:

  • (ArgumentError)


27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
# File '/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async.rb', line 27

def bind(listener_spec = @config.bind)
  raise ArgumentError, "transport is already bound" if @listener
  raise ArgumentError, "at least one controller must be registered" if @dispatcher.router.routes.empty?
  raise ConfigurationError, "Async reflection is not supported" if @config.reflection
  raise ConfigurationError, "Async TLS is not supported" unless @config.tls.empty?

  %i[max_connection_age max_connection_age_grace keepalive_time keepalive_permit_without_calls].each do |setting|
    unless @config.public_send(setting) == Configuration::DEFAULTS.fetch(setting)
      raise ConfigurationError, "#{setting} is only supported by the Native adapter"
    end
  end

  if listener_spec.respond_to?(:accept)
    @listener = listener_spec.dup
  else
    @listener = Supervisor::Listener.bind(listener_spec, reuseport: @config.listener_strategy == :reuseport)
  end
  @listener.local_address.ip_port
end

#drain! ⇒ Object



107
108
109
110
# File '/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async.rb', line 107

def drain!
  @draining = true
  self
end

#kill ⇒ Object



105
# File '/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async.rb', line 105

def kill = stop(deadline: Time.now)

#refresh_health ⇒ Object



114
115
116
117
118
119
120
121
# File '/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async.rb', line 114

def refresh_health
  checks = @config.health_checks.transform_values do |check|
    check.call ? true : false
  rescue StandardError
    false
  end
  update_health(ready: running?, checks:)
end

#running? ⇒ Boolean

Returns:

  • (Boolean)


86
# File '/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async.rb', line 86

def running? = @thread&.alive? && !@stopped

#start ⇒ Object

Raises:

  • (ArgumentError)


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
78
79
80
81
82
83
84
# File '/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async.rb', line 47

def start
  raise ArgumentError, "transport is stopped" if @stopped
  raise ArgumentError, "bind must be called before start" unless @listener
  raise ArgumentError, "transport is already started" if @thread

  @control_r, @control_w = IO.pipe
  ready = Queue.new
  @thread = Thread.new do
    Async do |root|
      bridge = Bridge.new(self, @dispatcher.router.routes, config: @config, logger: @logger)
      @server = Server.new(bridge, @logger)
      acceptor = root.async do
        loop do
          peer, address = @listener.accept
          root.async(peer, address) do |_task, client, remote|
            @server.accept(client, remote)
          ensure
            client.close unless client.closed?
          end
        end
      end
      ready << true
      @control_r.read(1)
      acceptor.stop
      @listener.close
      @server.stop(deadline: @stop_deadline)
    ensure
      root.children.to_a.each(&:stop)
    end
  rescue StandardError => e
    ready << e
    raise
  end
  result = ready.pop(timeout: 5)
  raise(result || "Async server did not start within 5 seconds") unless result == true

  self
end

#stats ⇒ Object



123
124
125
126
127
128
129
130
# File '/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async.rb', line 123

def stats
  @lock.synchronize do
    oldest = @inflight.values.min
    { inflight: @inflight.size, busy: @inflight.size, capacity: @config.threads + @config.max_waiting_requests,
      rejected_total: @rejected_total, requests_total: @requests_total,
      oldest_inflight_age: oldest ? Process.clock_gettime(Process::CLOCK_MONOTONIC) - oldest : 0 }
  end
end

#stop(deadline:) ⇒ Object



89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
# File '/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async.rb', line 89

def stop(deadline:)
  drain!
  return if @stopped

  @stopped = true
  @stop_deadline = deadline
  if @thread&.alive?
    @control_w.write(".")
    wait
  end
ensure
  @listener&.close unless @listener&.closed?
  @control_r&.close unless @control_r&.closed?
  @control_w&.close unless @control_w&.closed?
end

#update_health(ready:, checks: {}) ⇒ Object



112
# File '/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async.rb', line 112

def update_health(ready:, checks: {}) = ready && checks.values.all? && !@draining

#wait ⇒ Object



87
# File '/home/runner/work/gritz/gritz/sources/gritz-async/lib/gritz/transport/async.rb', line 87

def wait = @thread&.value